Skip to content

Commit 8c26ebb

Browse files
authored
Code for requeue event once even loop exit, registry cleanup an… (#1281)
* Code added for requeue event once even loop exit, registry cleanup and orphaned registry detection * Unit and integration tests added for requeueEvent, registry cleanup and orphaned registry detection * Apply gofumpt formatting * Removing parallel run of VSS registry tests * Code added to close dead WebSocket before removal in getOrCreateWebSocket * Code added to replace infinite retry with errorCount threshold in eventLoop * Code modified to provide threshold for connection retry and recreate shared websocket after reconnect failure threshold is reached * Requeue event, registry cleanup and orphan entry code added in VDS+VSS along with testcases * Test added for vaultDynamicSecret in event_watcher_registry_test.go * Comments updated * Fixed WebSocket stop-notification ordering, atomicity, and delivery; eliminated subscriber double-delete race; narrowed the Client interface; deduplicated subscriber code; closed orphan detection gaps; and cleaned up tests.
1 parent d762f62 commit 8c26ebb

10 files changed

Lines changed: 1454 additions & 159 deletions

controllers/event_watcher_registry_test.go

Lines changed: 364 additions & 30 deletions
Large diffs are not rendered by default.

controllers/vaultdynamicsecret_controller.go

Lines changed: 91 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -1065,6 +1065,49 @@ func (r *VaultDynamicSecretReconciler) ensureEventWatcher(
10651065
logger := log.FromContext(ctx).WithName("ensureEventWatcher")
10661066
name := client.ObjectKeyFromObject(o)
10671067

1068+
// onStop is shared across all Subscriber structs for this resource so that
1069+
// only the first OnStop call deletes the registry entry. Without this, a
1070+
// second or third independent closure could delete the entry that a
1071+
// reconciler already re-created after the first OnStop fired.
1072+
var onStopOnce sync.Once
1073+
onStop := func() {
1074+
onStopOnce.Do(func() {
1075+
r.eventWatcherRegistry.Delete(name)
1076+
})
1077+
}
1078+
1079+
// newDVSSubscriber is a local factory that stamps out a *vault.Subscriber
1080+
// with all fields common to every subscription for this resource. Only
1081+
// vaultNS and vaultPath vary between the engine-events subscriber (Step 4)
1082+
// and the lease-events subscribers (Steps 1 and 5).
1083+
//
1084+
// NOTE: lease subscribers intentionally pass vaultNS="" — lease events are
1085+
// routed by lease ID alone, so setting a namespace would produce a
1086+
// "<namespace>/<leaseID>" key that never matches the "<leaseID>" lookup.
1087+
newDVSSubscriber := func(vaultNS, vaultPath string) *vault.Subscriber {
1088+
return &vault.Subscriber{
1089+
ResourceKey: name,
1090+
VaultNS: vaultNS,
1091+
VaultPath: vaultPath,
1092+
ResourceType: vault.ResourceTypeVaultDynamicSecret,
1093+
ReconcileCh: r.SourceCh,
1094+
PendingVaultIndex: &r.pendingVaultIndex,
1095+
// OnStop cleans up the registry entry when the WebSocket dies.
1096+
// The log is emitted by notifySubscribersOfStop via ws.logger,
1097+
// which is always valid (unlike the reconcile-context logger
1098+
// captured here, which goes stale after Reconcile returns).
1099+
OnStop: onStop,
1100+
NewObject: func() client.Object {
1101+
return &secretsv1beta1.VaultDynamicSecret{
1102+
ObjectMeta: metav1.ObjectMeta{
1103+
Namespace: name.Namespace,
1104+
Name: name.Name,
1105+
},
1106+
}
1107+
},
1108+
}
1109+
}
1110+
10681111
currentLeaseID := o.Status.SecretLease.ID
10691112
meta, hasMeta := r.eventWatcherRegistry.Get(name)
10701113

@@ -1103,31 +1146,40 @@ func (r *VaultDynamicSecretReconciler) ensureEventWatcher(
11031146

11041147
// Lease events are mount-type-independent; still subscribe if applicable.
11051148
if !o.Spec.AllowStaticCreds && currentLeaseID != "" {
1106-
if !(hasMeta && meta.LastLeaseID == currentLeaseID && meta.LastClientID == c.ID()) {
1107-
leaseSubscriber := &vault.Subscriber{
1108-
ResourceKey: name,
1109-
VaultPath: currentLeaseID,
1110-
ResourceType: "VaultDynamicSecret",
1111-
ReconcileCh: r.SourceCh,
1112-
PendingVaultIndex: &r.pendingVaultIndex,
1113-
}
1114-
if err := c.SubscribeToEvents(ctx, vault.EventTypeLease, leaseSubscriber); err != nil {
1115-
// Surface this too, instead of only logging it, so a compound
1116-
// failure (mount type AND lease subscribe both broken) is visible.
1117-
resultErr = errors.Join(resultErr, fmt.Errorf("failed to subscribe to lease events: %w", err))
1118-
} else {
1119-
// Record what we did establish so the next reconcile can detect
1120-
// "lease-only subscription already active" and skip re-subscribing
1121-
// on every cycle while the mount-type lookup keeps failing.
1122-
// LastEventType is left empty so Step 2's staleness check still
1123-
// forces a full re-subscribe once GetMountType succeeds.
1124-
r.eventWatcherRegistry.Register(name, &eventWatcherMeta{
1125-
LastClientID: c.ID(),
1126-
LastGeneration: o.GetGeneration(),
1127-
LastLeaseID: currentLeaseID,
1128-
LastEventType: "",
1129-
})
1130-
}
1149+
// Skip re-subscription only when the existing lease watcher is
1150+
// confirmed alive. A metadata-only match is not sufficient: if the
1151+
// lease WebSocket died (e.g. reconnect threshold exceeded) but
1152+
// OnStop has not yet cleaned the registry, the entry looks valid
1153+
// but the resource is actually unsubscribed. The liveness check
1154+
// catches this orphaned state and falls through to re-subscribe.
1155+
leaseAlive := hasMeta &&
1156+
meta.LastLeaseID == currentLeaseID &&
1157+
meta.LastClientID == c.ID() &&
1158+
c.IsWebSocketHealthy(vault.EventTypeLease)
1159+
if leaseAlive {
1160+
return resultErr
1161+
}
1162+
// Dead or missing — clear any stale entry before re-subscribing.
1163+
if hasMeta {
1164+
r.eventWatcherRegistry.Delete(name)
1165+
}
1166+
leaseSubscriber := newDVSSubscriber("", currentLeaseID)
1167+
if err := c.SubscribeToEvents(ctx, vault.EventTypeLease, leaseSubscriber); err != nil {
1168+
// Surface this too, instead of only logging it, so a compound
1169+
// failure (mount type AND lease subscribe both broken) is visible.
1170+
resultErr = errors.Join(resultErr, fmt.Errorf("failed to subscribe to lease events: %w", err))
1171+
} else {
1172+
// Record what we did establish so the next reconcile can detect
1173+
// "lease-only subscription already active" and skip re-subscribing
1174+
// on every cycle while the mount-type lookup keeps failing.
1175+
// LastEventType is left empty so Step 2's staleness check still
1176+
// forces a full re-subscribe once GetMountType succeeds.
1177+
r.eventWatcherRegistry.Register(name, &eventWatcherMeta{
1178+
LastClientID: c.ID(),
1179+
LastGeneration: o.GetGeneration(),
1180+
LastLeaseID: currentLeaseID,
1181+
LastEventType: "",
1182+
})
11311183
}
11321184
}
11331185
return resultErr
@@ -1141,9 +1193,19 @@ func (r *VaultDynamicSecretReconciler) ensureEventWatcher(
11411193
meta.LastClientID == c.ID() &&
11421194
meta.LastLeaseID == currentLeaseID &&
11431195
meta.LastEventType == eventType {
1144-
logger.V(consts.LogLevelDebug).Info("Event subscription already active",
1196+
// Orphaned entry detection: verify the WebSocket is actually alive.
1197+
// This is a safety net for cases where OnStop did not fire (e.g. operator
1198+
// restart) or there was a race between WebSocket death and OnStop cleanup.
1199+
if c.IsWebSocketHealthy(eventType) {
1200+
logger.V(consts.LogLevelDebug).Info("Event subscription already active",
1201+
"namespace", o.Namespace, "name", o.Name)
1202+
return nil
1203+
}
1204+
// WebSocket is dead or missing — orphaned registry entry detected.
1205+
logger.Info("Detected orphaned registry entry (WebSocket is dead or missing), cleaning up",
11451206
"namespace", o.Namespace, "name", o.Name)
1146-
return nil
1207+
r.eventWatcherRegistry.Delete(name)
1208+
hasMeta = false
11471209
}
11481210

11491211
// Step 3: tear down the stale subscription now that we have a valid type.
@@ -1155,14 +1217,7 @@ func (r *VaultDynamicSecretReconciler) ensureEventWatcher(
11551217

11561218
// Step 4: subscribe to engine events.
11571219
vaultPath := buildVaultEventKey(o)
1158-
subscriber := &vault.Subscriber{
1159-
ResourceKey: name,
1160-
VaultNS: o.Spec.Namespace,
1161-
VaultPath: vaultPath,
1162-
ResourceType: "VaultDynamicSecret",
1163-
ReconcileCh: r.SourceCh,
1164-
PendingVaultIndex: &r.pendingVaultIndex,
1165-
}
1220+
subscriber := newDVSSubscriber(o.Spec.Namespace, vaultPath)
11661221
if err := c.SubscribeToEvents(ctx, eventType, subscriber); err != nil {
11671222
return fmt.Errorf("failed to subscribe to %s events: %w", eventType, err)
11681223
}
@@ -1176,13 +1231,7 @@ func (r *VaultDynamicSecretReconciler) ensureEventWatcher(
11761231
// unWatchEventsWithLeaseID), so setting VaultNS here would key the subscriber
11771232
// as "<namespace>/<leaseID>" and never match the "<leaseID>" lookup.
11781233
if !o.Spec.AllowStaticCreds && o.Status.SecretLease.ID != "" {
1179-
leaseSubscriber := &vault.Subscriber{
1180-
ResourceKey: name,
1181-
VaultPath: currentLeaseID,
1182-
ResourceType: "VaultDynamicSecret",
1183-
ReconcileCh: r.SourceCh,
1184-
PendingVaultIndex: &r.pendingVaultIndex,
1185-
}
1234+
leaseSubscriber := newDVSSubscriber("", currentLeaseID)
11861235
if err := c.SubscribeToEvents(ctx, vault.EventTypeLease, leaseSubscriber); err != nil {
11871236
// Non-fatal: the database/LDAP subscription above still provides
11881237
// event-driven updates. Surface the failure as a warning event so

controllers/vaultdynamicsecret_controller_test.go

Lines changed: 110 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2443,14 +2443,22 @@ type mockEnsureClient struct {
24432443
mountTypeErr error
24442444
subscribed []vault.EventType
24452445
seen []vault.EventType // from UnsubscribeFromEvents
2446+
// webSocketHealthy is returned by IsWebSocketHealthy.
2447+
webSocketHealthy bool
2448+
// onSubscribe is an optional hook called on every SubscribeToEvents call,
2449+
// allowing tests to capture Subscriber fields such as OnStop.
2450+
onSubscribe func(vault.EventType, *vault.Subscriber)
24462451
}
24472452

24482453
func (m *mockEnsureClient) GetMountType(_ context.Context, _ string) (string, error) {
24492454
return m.mountTypeResult, m.mountTypeErr
24502455
}
24512456

2452-
func (m *mockEnsureClient) SubscribeToEvents(_ context.Context, et vault.EventType, _ *vault.Subscriber) error {
2457+
func (m *mockEnsureClient) SubscribeToEvents(_ context.Context, et vault.EventType, sub *vault.Subscriber) error {
24532458
m.subscribed = append(m.subscribed, et)
2459+
if m.onSubscribe != nil {
2460+
m.onSubscribe(et, sub)
2461+
}
24542462
return nil
24552463
}
24562464

@@ -2461,6 +2469,10 @@ func (m *mockEnsureClient) UnsubscribeFromEvents(et vault.EventType, _ vault.Sub
24612469

24622470
func (m *mockEnsureClient) ID() string { return "test-client" }
24632471

2472+
func (m *mockEnsureClient) IsWebSocketHealthy(_ vault.EventType) bool {
2473+
return m.webSocketHealthy
2474+
}
2475+
24642476
// Test_ensureEventWatcher_GetMountTypeError_ReusesPriorEventType verifies
24652477
// that when GetMountType fails but a prior LastEventType is stored, that prior
24662478
// type is reused — keeping the subscription on the correct event stream instead
@@ -2694,3 +2706,100 @@ func TestVaultDynamicSecretReconciler_syncSecret_vaultIndex(t *testing.T) {
26942706
})
26952707
}
26962708
}
2709+
2710+
// Test_ensureEventWatcher_OrphanedEntry_NilWebSocket verifies that when the
2711+
// registry has a matching entry but IsWebSocketHealthy returns false (e.g. after
2712+
// an operator restart), ensureEventWatcher detects the orphaned entry, clears it,
2713+
// and re-subscribes rather than returning nil.
2714+
func Test_ensureEventWatcher_OrphanedEntry_NilWebSocket(t *testing.T) {
2715+
ch := make(chan event.GenericEvent, 10)
2716+
r := &VaultDynamicSecretReconciler{
2717+
eventWatcherRegistry: newEventWatcherRegistry(),
2718+
SourceCh: ch,
2719+
Recorder: record.NewFakeRecorder(10),
2720+
}
2721+
2722+
o := &secretsv1beta1.VaultDynamicSecret{
2723+
ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: "db-secret", Generation: 1},
2724+
Spec: secretsv1beta1.VaultDynamicSecretSpec{
2725+
Mount: "database",
2726+
Path: "creds/my-role",
2727+
},
2728+
}
2729+
key := client.ObjectKeyFromObject(o)
2730+
2731+
// Pre-populate registry with fully-matching metadata so the staleness check
2732+
// would normally return nil — the only thing that should trigger re-subscribe
2733+
// is the nil WebSocket.
2734+
r.eventWatcherRegistry.Register(key, &eventWatcherMeta{
2735+
LastClientID: "test-client",
2736+
LastGeneration: 1,
2737+
LastLeaseID: "",
2738+
LastEventType: vault.EventTypeDatabase,
2739+
})
2740+
2741+
// IsWebSocketHealthy returns false — simulates operator restart with no live WebSocket.
2742+
m := &mockEnsureClient{
2743+
mountTypeResult: "database",
2744+
webSocketHealthy: false,
2745+
}
2746+
2747+
err := r.ensureEventWatcher(context.Background(), o, m)
2748+
2749+
require.NoError(t, err)
2750+
assert.Contains(t, m.subscribed, vault.EventTypeDatabase,
2751+
"must re-subscribe when WebSocket is nil (orphaned entry)")
2752+
2753+
// Registry must be refreshed with new metadata after re-subscription.
2754+
meta, ok := r.eventWatcherRegistry.Get(key)
2755+
require.True(t, ok, "registry must have an entry after re-subscription")
2756+
assert.Equal(t, vault.EventTypeDatabase, meta.LastEventType)
2757+
}
2758+
2759+
// Test_VDS_ensureEventWatcher_OnStop_CleansRegistry verifies that the OnStop
2760+
// callback set on the VDS engine-events Subscriber deletes the registry entry
2761+
// when invoked, so the next reconcile falls through to re-subscribe instead of
2762+
// returning early because it sees a stale registry entry with matching metadata.
2763+
func Test_VDS_ensureEventWatcher_OnStop_CleansRegistry(t *testing.T) {
2764+
ch := make(chan event.GenericEvent, 10)
2765+
r := &VaultDynamicSecretReconciler{
2766+
eventWatcherRegistry: newEventWatcherRegistry(),
2767+
SourceCh: ch,
2768+
Recorder: record.NewFakeRecorder(10),
2769+
}
2770+
2771+
o := &secretsv1beta1.VaultDynamicSecret{
2772+
ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: "db-secret", Generation: 1},
2773+
Spec: secretsv1beta1.VaultDynamicSecretSpec{
2774+
Mount: "database",
2775+
Path: "creds/my-role",
2776+
},
2777+
}
2778+
key := client.ObjectKeyFromObject(o)
2779+
2780+
// captureOnStop intercepts the first OnStop set on any SubscribeToEvents call.
2781+
var capturedOnStop func()
2782+
m := &mockEnsureClient{
2783+
mountTypeResult: "database",
2784+
onSubscribe: func(_ vault.EventType, sub *vault.Subscriber) {
2785+
if capturedOnStop == nil && sub.OnStop != nil {
2786+
capturedOnStop = sub.OnStop
2787+
}
2788+
},
2789+
}
2790+
2791+
err := r.ensureEventWatcher(context.Background(), o, m)
2792+
require.NoError(t, err)
2793+
2794+
// Registry must have an entry after subscription.
2795+
_, ok := r.eventWatcherRegistry.Get(key)
2796+
require.True(t, ok, "registry must have entry after ensureEventWatcher")
2797+
2798+
// Simulate WebSocket death by invoking the captured OnStop callback.
2799+
require.NotNil(t, capturedOnStop, "OnStop must be set on the engine-events subscriber")
2800+
capturedOnStop()
2801+
2802+
// Registry entry must be gone — next reconcile will re-subscribe.
2803+
_, ok = r.eventWatcherRegistry.Get(key)
2804+
assert.False(t, ok, "registry entry must be deleted after OnStop fires")
2805+
}

controllers/vaultstaticsecret_controller.go

Lines changed: 48 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -373,19 +373,46 @@ func (r *VaultStaticSecretReconciler) ensureEventWatcher(ctx context.Context, o
373373
logger := log.FromContext(ctx).WithName("ensureEventWatcher")
374374
name := client.ObjectKeyFromObject(o)
375375

376+
// onStop is shared across all Subscriber structs for this resource so that
377+
// only the first OnStop call deletes the registry entry. A single subscriber
378+
// is registered today, but the Once guard makes the pattern consistent with
379+
// VaultDynamicSecret and safe if additional subscribers are added in future.
380+
var onStopOnce sync.Once
381+
onStop := func() {
382+
onStopOnce.Do(func() {
383+
r.eventWatcherRegistry.Delete(name)
384+
})
385+
}
386+
387+
// vssEventType is the single event type VSS subscribes to. Declaring it once
388+
// here keeps the orphan check and the subscribe call in sync — a future change
389+
// to the VSS event type only needs to be made in one place.
390+
const vssEventType = vault.EventTypeKV
391+
376392
meta, ok := r.eventWatcherRegistry.Get(name)
377393
if ok {
378-
// The subscription is active, and if the VSS object has not been updated,
379-
// and the client ID is the same, just return
380-
if meta.LastGeneration == o.GetGeneration() && meta.LastClientID == c.ID() {
381-
logger.V(consts.LogLevelDebug).Info("Event subscription already active",
394+
// Check if the WebSocket is actually healthy (Orphaned Entry Detection)
395+
// This is a safety net in case OnStop callback didn't run or there was a race condition
396+
if c.IsWebSocketHealthy(vssEventType) {
397+
// WebSocket is healthy, check if metadata matches
398+
if meta.LastGeneration == o.GetGeneration() && meta.LastClientID == c.ID() {
399+
// The subscription is active, and if the VSS object has not been updated,
400+
// and the client ID is the same, just return
401+
logger.V(consts.LogLevelDebug).Info("Event subscription already active",
402+
"namespace", o.Namespace, "name", o.Name)
403+
return nil
404+
}
405+
// The subscription exists but metadata or vault client has changed, unsubscribe first
406+
logger.V(consts.LogLevelDebug).Info("Unsubscribing due to metadata or client change",
407+
"namespace", o.Namespace, "name", o.Name)
408+
r.unWatchEvents(o, c)
409+
} else {
410+
// WebSocket is dead or missing - orphaned registry entry detected
411+
logger.Info("Detected orphaned registry entry (WebSocket is dead or missing), cleaning up",
382412
"namespace", o.Namespace, "name", o.Name)
383-
return nil
413+
r.eventWatcherRegistry.Delete(name)
414+
// Fall through to create a new WebSocket subscription
384415
}
385-
// The subscription exists but metadata or vault client has changed, unsubscribe first
386-
logger.V(consts.LogLevelDebug).Info("Unsubscribing due to metadata or client change",
387-
"namespace", o.Namespace, "name", o.Name)
388-
r.unWatchEvents(o, c)
389416
}
390417

391418
// Build the vault path for subscription
@@ -396,12 +423,22 @@ func (r *VaultStaticSecretReconciler) ensureEventWatcher(ctx context.Context, o
396423
ResourceKey: name,
397424
VaultNS: o.Spec.Namespace,
398425
VaultPath: vaultPath,
399-
ResourceType: "VaultStaticSecret",
426+
ResourceType: vault.ResourceTypeVaultStaticSecret,
400427
ReconcileCh: r.SourceCh,
401428
PendingVaultIndex: &r.pendingVaultIndex,
429+
// OnStop callback cleans up registry when WebSocket dies
430+
OnStop: onStop,
431+
NewObject: func() client.Object {
432+
return &secretsv1beta1.VaultStaticSecret{
433+
ObjectMeta: metav1.ObjectMeta{
434+
Namespace: name.Namespace,
435+
Name: name.Name,
436+
},
437+
}
438+
},
402439
}
403440

404-
if err := c.SubscribeToEvents(ctx, vault.EventTypeKV, subscriber); err != nil {
441+
if err := c.SubscribeToEvents(ctx, vssEventType, subscriber); err != nil {
405442
return fmt.Errorf("failed to subscribe to events: %w", err)
406443
}
407444

0 commit comments

Comments
 (0)