Skip to content

Commit ef57dbd

Browse files
committed
feat(node-drainer): replace event informer in node-drainer with eventrecorder
Signed-off-by: Ajay Mishra <ajmishra@nvidia.com>
1 parent c9c73b5 commit ef57dbd

8 files changed

Lines changed: 74 additions & 121 deletions

File tree

distros/kubernetes/nvsentinel/charts/node-drainer/templates/clusterrole.yaml

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -53,10 +53,8 @@ rules:
5353
resources:
5454
- events
5555
verbs:
56-
- list
5756
- create
58-
- update
59-
- watch
57+
- patch
6058
- apiGroups:
6159
- ""
6260
resources:

docs/configuration/node-drainer.md

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -254,7 +254,7 @@ kubectl get events -n default \
254254
--field-selector "involvedObject.kind=Node,involvedObject.name=$NODE"
255255
```
256256

257-
Node Drainer events are created in the **`default`** namespace (`Type=NodeDraining`,
257+
Node Drainer events are created in the **`default`** namespace (`Type=Normal`,
258258
`Source=nvsentinel-node-drainer`). Example output for a pod `training-job-0` in namespace
259259
`batch-jobs` on node `gpu-node-42`:
260260

@@ -283,7 +283,7 @@ database postgres-0 1/1 Running 0 2h
283283
284284
# kubectl get events -n default --field-selector "involvedObject.name=gpu-node-42"
285285
LAST SEEN TYPE REASON OBJECT MESSAGE
286-
2m NodeDraining AwaitingPodCompletion Node/gpu-node-42 Waiting for following pods to finish: [database/postgres-0]
286+
2m Normal AwaitingPodCompletion Node/gpu-node-42 Waiting for following pods to finish: [database/postgres-0]
287287
```
288288

289289
#### `DeleteAfterTimeout`
@@ -295,7 +295,7 @@ batch-jobs training-job-0 1/1 Running 0 2h
295295
296296
# kubectl get events -n default --field-selector "involvedObject.name=gpu-node-42"
297297
LAST SEEN TYPE REASON OBJECT MESSAGE
298-
1m NodeDraining WaitingBeforeForceDelete Node/gpu-node-42 Waiting for following pods to finish: [training-job-0] in namespace: [batch-jobs] or they will be force deleted on: 2026-07-01 18:30:00 +0000 UTC
298+
1m Normal WaitingBeforeForceDelete Node/gpu-node-42 Waiting for following pods to finish: [training-job-0] in namespace: [batch-jobs] or they will be force deleted on: 2026-07-01 18:30:00 +0000 UTC
299299
```
300300

301301
After the deadline, the pod is force-deleted and `kubectl get pods` shows it gone.

node-drainer/pkg/informers/informers.go

Lines changed: 28 additions & 103 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,10 @@ import (
3434
"k8s.io/apimachinery/pkg/types"
3535
"k8s.io/client-go/informers"
3636
"k8s.io/client-go/kubernetes"
37+
"k8s.io/client-go/kubernetes/scheme"
38+
typedcorev1 "k8s.io/client-go/kubernetes/typed/core/v1"
3739
"k8s.io/client-go/tools/cache"
40+
"k8s.io/client-go/tools/record"
3841
"k8s.io/utils/ptr"
3942

4043
"github.com/nvidia/nvsentinel/data-models/pkg/model"
@@ -44,20 +47,19 @@ import (
4447
)
4548

4649
const (
47-
NodeIndex = "node"
48-
NamespaceNodeIndex = "namespace-node"
49-
NodeEventReasonIndex = "node-event-reason"
50+
NodeIndex = "node"
51+
NamespaceNodeIndex = "namespace-node"
5052
)
5153

5254
type Informers struct {
5355
podInformer cache.SharedIndexInformer
54-
eventInformer cache.SharedIndexInformer
5556
nodeInformer cache.SharedIndexInformer
57+
eventBroadcaster record.EventBroadcaster
58+
eventRecorder record.EventRecorder
5659
clientset kubernetes.Interface
5760
notReadyTimeoutMinutes *int
5861
drainGPUPods bool
5962
dryRunMode []string
60-
namespace string
6163
}
6264

6365
func NewInformers(clientset kubernetes.Interface, resyncPeriod time.Duration,
@@ -87,41 +89,30 @@ func NewInformers(clientset kubernetes.Interface, resyncPeriod time.Duration,
8789
return nil, fmt.Errorf("failed to add indexer: %w", err)
8890
}
8991

90-
eventInformerFactory := informers.NewSharedInformerFactoryWithOptions(
91-
clientset,
92-
resyncPeriod,
93-
informers.WithNamespace(metav1.NamespaceDefault),
94-
)
95-
96-
eventInformer := eventInformerFactory.Core().V1().Events().Informer()
97-
98-
err = eventInformer.GetIndexer().AddIndexers(
99-
cache.Indexers{
100-
NodeEventReasonIndex: NodeEventReasonIndexFunc,
101-
})
102-
if err != nil {
103-
return nil, fmt.Errorf("failed to add event indexer: %w", err)
104-
}
105-
10692
nodeInformer := informerFactory.Core().V1().Nodes().Informer()
10793
if err := nodeInformer.SetTransform(nodeTransform); err != nil {
10894
return nil, fmt.Errorf("failed to set node informer transform: %w", err)
10995
}
11096

97+
eventBroadcaster := record.NewBroadcaster()
98+
11199
dryRunMode := []string{}
112100
if dryRun {
113101
dryRunMode = []string{metav1.DryRunAll}
114102
}
115103

116104
return &Informers{
117-
clientset: clientset,
118-
podInformer: podInformer,
119-
eventInformer: eventInformer,
120-
nodeInformer: nodeInformer,
105+
clientset: clientset,
106+
podInformer: podInformer,
107+
nodeInformer: nodeInformer,
108+
eventBroadcaster: eventBroadcaster,
109+
eventRecorder: eventBroadcaster.NewRecorder(
110+
scheme.Scheme,
111+
v1.EventSource{Component: "nvsentinel-node-drainer"},
112+
),
121113
notReadyTimeoutMinutes: notReadyTimeoutMinutes,
122114
drainGPUPods: drainGPUPods,
123115
dryRunMode: dryRunMode,
124-
namespace: metav1.NamespaceDefault,
125116
}, nil
126117
}
127118

@@ -273,7 +264,7 @@ func isDaemonSetOwned(ownerReferences []metav1.OwnerReference) bool {
273264
}
274265

275266
func (i *Informers) HasSynced() bool {
276-
return i.podInformer.HasSynced() && i.eventInformer.HasSynced() && i.nodeInformer.HasSynced()
267+
return i.podInformer.HasSynced() && i.nodeInformer.HasSynced()
277268
}
278269

279270
func NodeIndexFunc(obj any) ([]string, error) {
@@ -304,25 +295,17 @@ func NamespaceNodeIndexFunc(obj any) ([]string, error) {
304295
return []string{compositeKey}, nil
305296
}
306297

307-
func NodeEventReasonIndexFunc(obj any) ([]string, error) {
308-
event, ok := obj.(*v1.Event)
309-
if !ok {
310-
return []string{}, nil
311-
}
312-
313-
if event.InvolvedObject.Kind != "Node" {
314-
return []string{}, nil
315-
}
316-
317-
compositeKey := fmt.Sprintf("%s/%s", event.InvolvedObject.Name, event.Reason)
318-
319-
return []string{compositeKey}, nil
320-
}
321-
322298
func (i *Informers) Run(ctx context.Context) error {
299+
i.eventBroadcaster.StartRecordingToSink(&typedcorev1.EventSinkImpl{
300+
Interface: i.clientset.CoreV1().Events(metav1.NamespaceDefault),
301+
})
302+
323303
go i.podInformer.Run(ctx.Done())
324-
go i.eventInformer.Run(ctx.Done())
325304
go i.nodeInformer.Run(ctx.Done())
305+
go func() {
306+
<-ctx.Done()
307+
i.eventBroadcaster.Shutdown()
308+
}()
326309

327310
if ok := cache.WaitForCacheSync(ctx.Done(),
328311
i.HasSynced); !ok {
@@ -714,71 +697,13 @@ func (i *Informers) sendEvictionRequestForPod(ctx context.Context, namespace str
714697
return nil
715698
}
716699

717-
func (i *Informers) UpdateNodeEvent(ctx context.Context, nodeName string, reason string, message string) error {
718-
compositeKey := fmt.Sprintf("%s/%s", nodeName, reason)
719-
720-
cachedEvents, err := i.eventInformer.GetIndexer().ByIndex(NodeEventReasonIndex, compositeKey)
721-
if err != nil {
722-
slog.ErrorContext(ctx, "Failed to query event cache", "error", err)
723-
return fmt.Errorf("error querying event cache: %w", err)
724-
}
725-
726-
now := metav1.NewTime(time.Now())
727-
eventsClient := i.clientset.CoreV1().Events(i.namespace)
728-
729-
for _, obj := range cachedEvents {
730-
existingEvent, ok := obj.(*v1.Event)
731-
if !ok {
732-
continue
733-
}
734-
735-
if existingEvent.Message == message {
736-
eventCopy := existingEvent.DeepCopy()
737-
eventCopy.Count++
738-
eventCopy.LastTimestamp = now
739-
740-
_, err = eventsClient.Update(ctx, eventCopy, metav1.UpdateOptions{})
741-
if err != nil {
742-
slog.ErrorContext(ctx, "Failed to update event occurrence count", "error", err)
743-
return fmt.Errorf("error in updating event occurrence count: %w", err)
744-
}
745-
746-
return nil
747-
}
748-
}
749-
750-
// Get node from informer cache to retrieve its UID for proper event association
700+
func (i *Informers) UpdateNodeEvent(_ context.Context, nodeName string, reason string, message string) error {
751701
node, err := i.GetNode(nodeName)
752702
if err != nil {
753703
return fmt.Errorf("error getting node from cache %s: %w", nodeName, err)
754704
}
755705

756-
newEvent := &v1.Event{
757-
ObjectMeta: metav1.ObjectMeta{
758-
GenerateName: nodeName + "-",
759-
Namespace: i.namespace,
760-
},
761-
InvolvedObject: v1.ObjectReference{
762-
Kind: "Node",
763-
Name: nodeName,
764-
UID: node.UID,
765-
APIVersion: "v1",
766-
},
767-
Reason: reason,
768-
Message: message,
769-
Type: "NodeDraining",
770-
Source: v1.EventSource{Component: "nvsentinel-node-drainer"},
771-
FirstTimestamp: now,
772-
LastTimestamp: now,
773-
Count: 1,
774-
}
775-
776-
_, err = eventsClient.Create(ctx, newEvent, metav1.CreateOptions{})
777-
if err != nil {
778-
slog.ErrorContext(ctx, "Failed to create event", "error", err, "node", nodeName, "reason", reason)
779-
return fmt.Errorf("error in creating event: %w", err)
780-
}
781-
706+
i.eventRecorder.Event(node, v1.EventTypeNormal, reason, message)
782707
return nil
783708
}
784709

node-drainer/pkg/informers/informers_transform_test.go

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -245,10 +245,19 @@ func TestInformerTransformsIntegrateWithIndexesAndNodeEvents(t *testing.T) {
245245
assert.Equal(t, node.Annotations, cachedNode.Annotations)
246246

247247
require.NoError(t, informers.UpdateNodeEvent(ctx, node.Name, "AwaitingPodCompletion", "waiting"))
248-
events, err := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, metav1.ListOptions{})
249-
require.NoError(t, err)
250-
require.Len(t, events.Items, 1)
251-
assert.Equal(t, node.UID, events.Items[0].InvolvedObject.UID)
248+
require.NoError(t, informers.UpdateNodeEvent(ctx, node.Name, "AwaitingPodCompletion", "waiting"))
249+
require.EventuallyWithT(t, func(collect *assert.CollectT) {
250+
events, listErr := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, metav1.ListOptions{})
251+
if !assert.NoError(collect, listErr) || !assert.Len(collect, events.Items, 1) {
252+
return
253+
}
254+
255+
event := events.Items[0]
256+
assert.Equal(collect, node.UID, event.InvolvedObject.UID)
257+
assert.Equal(collect, v1.EventTypeNormal, event.Type)
258+
assert.Equal(collect, "nvsentinel-node-drainer", event.Source.Component)
259+
assert.Equal(collect, int32(2), event.Count)
260+
}, 5*time.Second, 50*time.Millisecond)
252261
}
253262

254263
func richDrainEligiblePod(namespace, name, nodeName string) *v1.Pod {

node-drainer/pkg/reconciler/reconciler_integration_test.go

Lines changed: 24 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -504,9 +504,16 @@ func TestReconciler_ProcessEvent(t *testing.T) {
504504
assert.Error(t, err)
505505
assert.Contains(t, err.Error(), "waiting for pods to complete")
506506

507-
nodeEvents, err := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, metav1.ListOptions{FieldSelector: fmt.Sprintf("involvedObject.name=%s,involvedObject.kind=Node", nodeName)})
507+
eventListOptions := metav1.ListOptions{FieldSelector: fmt.Sprintf("involvedObject.name=%s,involvedObject.kind=Node", nodeName)}
508+
require.Eventually(t, func() bool {
509+
nodeEvents, listErr := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, eventListOptions)
510+
return listErr == nil && len(nodeEvents.Items) == 1
511+
}, 5*time.Second, 50*time.Millisecond, "node event should be recorded asynchronously")
512+
513+
nodeEvents, err := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, eventListOptions)
508514
require.NoError(t, err)
509515
require.Len(t, nodeEvents.Items, 1, "only one event should be created despite multiple reconciliations")
516+
require.Equal(t, v1.EventTypeNormal, nodeEvents.Items[0].Type)
510517
require.Equal(t, nodeEvents.Items[0].Reason, "AwaitingPodCompletion")
511518
expectedMessage := "Waiting for following pods to finish: [completion-test/running-pod-1]"
512519
require.Equal(t, expectedMessage, nodeEvents.Items[0].Message, "only expected pods should be drained")
@@ -661,9 +668,16 @@ func TestReconciler_ProcessEvent(t *testing.T) {
661668
assert.Error(t, err)
662669
assert.Contains(t, err.Error(), "failed timeout eviction for node test-node: waiting for 1 pods to complete or timeout")
663670

664-
nodeEvents, err := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, metav1.ListOptions{FieldSelector: fmt.Sprintf("involvedObject.name=%s,involvedObject.kind=Node", nodeName)})
671+
eventListOptions := metav1.ListOptions{FieldSelector: fmt.Sprintf("involvedObject.name=%s,involvedObject.kind=Node", nodeName)}
672+
require.Eventually(t, func() bool {
673+
nodeEvents, listErr := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, eventListOptions)
674+
return listErr == nil && len(nodeEvents.Items) == 1
675+
}, 5*time.Second, 50*time.Millisecond, "node event should be recorded asynchronously")
676+
677+
nodeEvents, err := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, eventListOptions)
665678
require.NoError(t, err)
666679
require.Len(t, nodeEvents.Items, 1, "only one event should be created despite multiple reconciliations")
680+
require.Equal(t, v1.EventTypeNormal, nodeEvents.Items[0].Type)
667681
require.Equal(t, nodeEvents.Items[0].Reason, "WaitingBeforeForceDelete")
668682
expectedMessage := "Waiting for following pods to finish: [pod-1] in namespace: [timeout-test] or they will be force deleted on:"
669683
require.Contains(t, nodeEvents.Items[0].Message, expectedMessage, "only expected pods should be drained")
@@ -892,9 +906,16 @@ func TestReconciler_ProcessEvent(t *testing.T) {
892906
assert.Error(t, err)
893907
assert.Contains(t, err.Error(), "waiting for pods to complete")
894908

895-
nodeEvents, err := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, metav1.ListOptions{FieldSelector: fmt.Sprintf("involvedObject.name=%s,involvedObject.kind=Node", nodeName)})
909+
eventListOptions := metav1.ListOptions{FieldSelector: fmt.Sprintf("involvedObject.name=%s,involvedObject.kind=Node", nodeName)}
910+
require.Eventually(t, func() bool {
911+
nodeEvents, listErr := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, eventListOptions)
912+
return listErr == nil && len(nodeEvents.Items) == 1
913+
}, 5*time.Second, 50*time.Millisecond, "node event should be recorded asynchronously")
914+
915+
nodeEvents, err := client.CoreV1().Events(metav1.NamespaceDefault).List(ctx, eventListOptions)
896916
require.NoError(t, err)
897917
require.Len(t, nodeEvents.Items, 1, "only one event should be created despite multiple reconciliations")
918+
require.Equal(t, v1.EventTypeNormal, nodeEvents.Items[0].Type)
898919
require.Equal(t, nodeEvents.Items[0].Reason, "AwaitingPodCompletion")
899920
expectedMessage := "Waiting for following pods to finish: [completion-test/running-pod-1 completion-test/running-pod-2 completion-test/running-pod-3]"
900921
require.Equal(t, expectedMessage, nodeEvents.Items[0].Message, "pod list should be in sorted order")

tests/fault_management_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -346,7 +346,7 @@ func TestNodeRecoveryDuringDrain(t *testing.T) {
346346
require.NoError(t, err)
347347

348348
require.Eventually(t, func() bool {
349-
found, _ := helpers.CheckNodeEventExists(ctx, client, testCtx.NodeName, "NodeDraining", "WaitingBeforeForceDelete", time.Time{})
349+
found, _ := helpers.CheckNodeEventExists(ctx, client, testCtx.NodeName, v1.EventTypeNormal, "WaitingBeforeForceDelete", time.Time{})
350350
return found
351351
}, helpers.EventuallyWaitTimeout, helpers.WaitInterval, "WaitingBeforeForceDelete event should be created")
352352

@@ -446,7 +446,7 @@ func TestManualUncordonPropagation(t *testing.T) {
446446
require.NoError(t, err)
447447

448448
require.Eventually(t, func() bool {
449-
found, _ := helpers.CheckNodeEventExists(ctx, client, testCtx.NodeName, "NodeDraining", "WaitingBeforeForceDelete", time.Time{})
449+
found, _ := helpers.CheckNodeEventExists(ctx, client, testCtx.NodeName, v1.EventTypeNormal, "WaitingBeforeForceDelete", time.Time{})
450450
return found
451451
}, helpers.EventuallyWaitTimeout, helpers.WaitInterval, "WaitingBeforeForceDelete event should be created")
452452

tests/node_drainer_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -169,7 +169,7 @@ func TestNodeDrainerEvictionModes(t *testing.T) {
169169
require.NoError(t, err)
170170

171171
require.Eventually(t, func() bool {
172-
found, event := helpers.CheckNodeEventExists(ctx, client, testCtx.NodeName, "NodeDraining", "", restartTime)
172+
found, event := helpers.CheckNodeEventExists(ctx, client, testCtx.NodeName, v1.EventTypeNormal, "", restartTime)
173173
if found {
174174
t.Logf("Found event after restart: %s", event.Reason)
175175
}

tests/smoke_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,7 @@ func TestFatalHealthEvent(t *testing.T) {
130130
helpers.WaitForNodesWithLabel(ctx, t, client, []string{nodeName}, statemanager.NVSentinelStateLabelKey, string(statemanager.DrainingLabelValue))
131131

132132
expectedDrainingNodeEvent := v1.Event{
133-
Type: "NodeDraining",
133+
Type: v1.EventTypeNormal,
134134
Reason: "AwaitingPodCompletion",
135135
}
136136
helpers.WaitForNodeEvent(ctx, t, client, nodeName, expectedDrainingNodeEvent)
@@ -321,7 +321,7 @@ func TestFatalUnsupportedHealthEvent(t *testing.T) {
321321
helpers.WaitForNodesWithLabel(ctx, t, client, []string{nodeName}, statemanager.NVSentinelStateLabelKey, string(statemanager.DrainingLabelValue))
322322

323323
expectedDrainingNodeEvent := v1.Event{
324-
Type: "NodeDraining",
324+
Type: v1.EventTypeNormal,
325325
Reason: "AwaitingPodCompletion",
326326
}
327327
helpers.WaitForNodeEvent(ctx, t, client, nodeName, expectedDrainingNodeEvent)

0 commit comments

Comments
 (0)