Files
Ronak Nathani d7ba03a966 Improves leader election so that we don't lose events (#153)
Currently, when a replica loses its leadership, a new leader isn't elected until leaseDuration seconds.
Here, that is 15s. The max time till we get a new leader is leaseDuration (15s) + retryPeriod (2s) = 17s.

This commit updates the shutdown process such that if the leader replica is sent a shutdown signal,
it sleeps for leaseDuration seconds. This allows the leader replica to continue to export events until
a new leader is elected. And a new leader is elected only if lease hasn't been renewed and leaseDuration expires.

In addition to this, leader election now uses the leases object instead of configMaps and leases. The clusterRole
is also updated to allow writing to the leases object.

For use cases where no event loss is tolerable, users should use maxEventAgeSeconds to > 1.
2024-02-23 08:26:06 +03:00

155 lines
4.3 KiB
Go

package kube
import (
"sync"
"time"
"github.com/resmoio/kubernetes-event-exporter/pkg/metrics"
"github.com/rs/zerolog/log"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/cache"
)
var startUpTime = time.Now()
type EventHandler func(event *EnhancedEvent)
type EventWatcher struct {
wg sync.WaitGroup
informer cache.SharedInformer
stopper chan struct{}
objectMetadataCache ObjectMetadataProvider
omitLookup bool
fn EventHandler
maxEventAgeSeconds time.Duration
metricsStore *metrics.Store
dynamicClient *dynamic.DynamicClient
clientset *kubernetes.Clientset
}
func NewEventWatcher(config *rest.Config, namespace string, MaxEventAgeSeconds int64, metricsStore *metrics.Store, fn EventHandler, omitLookup bool, cacheSize int) *EventWatcher {
clientset := kubernetes.NewForConfigOrDie(config)
factory := informers.NewSharedInformerFactoryWithOptions(clientset, 0, informers.WithNamespace(namespace))
informer := factory.Core().V1().Events().Informer()
watcher := &EventWatcher{
informer: informer,
stopper: make(chan struct{}),
objectMetadataCache: NewObjectMetadataProvider(cacheSize),
omitLookup: omitLookup,
fn: fn,
maxEventAgeSeconds: time.Second * time.Duration(MaxEventAgeSeconds),
metricsStore: metricsStore,
dynamicClient: dynamic.NewForConfigOrDie(config),
clientset: clientset,
}
informer.AddEventHandler(watcher)
informer.SetWatchErrorHandler(func(r *cache.Reflector, err error) {
watcher.metricsStore.WatchErrors.Inc()
})
return watcher
}
func (e *EventWatcher) OnAdd(obj interface{}) {
event := obj.(*corev1.Event)
e.onEvent(event)
}
func (e *EventWatcher) OnUpdate(oldObj, newObj interface{}) {
// Ignore updates
}
// Ignore events older than the maxEventAgeSeconds
func (e *EventWatcher) isEventDiscarded(event *corev1.Event) bool {
timestamp := event.LastTimestamp.Time
if timestamp.IsZero() {
timestamp = event.EventTime.Time
}
eventAge := time.Since(timestamp)
if eventAge > e.maxEventAgeSeconds {
// Log discarded events if they were created after the watcher started
// (to suppres warnings from initial synchrnization)
if timestamp.After(startUpTime) {
log.Warn().
Str("event age", eventAge.String()).
Str("event namespace", event.Namespace).
Str("event name", event.Name).
Msg("Event discarded as being older then maxEventAgeSeconds")
e.metricsStore.EventsDiscarded.Inc()
}
return true
}
return false
}
func (e *EventWatcher) onEvent(event *corev1.Event) {
if e.isEventDiscarded(event) {
return
}
log.Debug().
Str("msg", event.Message).
Str("namespace", event.Namespace).
Str("reason", event.Reason).
Str("involvedObject", event.InvolvedObject.Name).
Msg("Received event")
e.metricsStore.EventsProcessed.Inc()
ev := &EnhancedEvent{
Event: *event.DeepCopy(),
}
ev.Event.ManagedFields = nil
if e.omitLookup {
ev.InvolvedObject.ObjectReference = *event.InvolvedObject.DeepCopy()
} else {
objectMetadata, err := e.objectMetadataCache.GetObjectMetadata(&event.InvolvedObject, e.clientset, e.dynamicClient, e.metricsStore)
if err != nil {
if errors.IsNotFound(err) {
ev.InvolvedObject.Deleted = true
log.Error().Err(err).Msg("Object not found, likely deleted")
} else {
log.Error().Err(err).Msg("Failed to get object metadata")
}
ev.InvolvedObject.ObjectReference = *event.InvolvedObject.DeepCopy()
} else {
ev.InvolvedObject.Labels = objectMetadata.Labels
ev.InvolvedObject.Annotations = objectMetadata.Annotations
ev.InvolvedObject.OwnerReferences = objectMetadata.OwnerReferences
ev.InvolvedObject.ObjectReference = *event.InvolvedObject.DeepCopy()
ev.InvolvedObject.Deleted = objectMetadata.Deleted
}
}
e.fn(ev)
}
func (e *EventWatcher) OnDelete(obj interface{}) {
// Ignore deletes
}
func (e *EventWatcher) Start() {
e.wg.Add(1)
go func() {
defer e.wg.Done()
e.informer.Run(e.stopper)
}()
}
func (e *EventWatcher) Stop() {
close(e.stopper)
e.wg.Wait()
}
func (e *EventWatcher) setStartUpTime(time time.Time) {
startUpTime = time
}