From 36f0c62fe105c7c037905fb98bb6f6710a41048f Mon Sep 17 00:00:00 2001 From: ChadiDridi Date: Thu, 27 Aug 2026 23:22:18 +0100 Subject: [PATCH 1/3] fix(kube): finish hook waits when the hook is deleted The legacy waiter ended a hook wait as soon as it saw a delete event for the watched resource. The kstatus watcher, which is the default in Helm 4, has no equivalent rule: it waits for the Current status, so a hook that is removed while Helm is waiting is reported as NotFound and blocks until the timeout expires, failing the release. That is the normal lifecycle of a Job hook that sets .spec.ttlSecondsAfterFinished: the TTL controller removes the Job as soon as it completes. The ingress-nginx chart, for example, sets that field on its admission webhook patch Jobs, so upgrades hang for the whole hook timeout and then fail. WatchUntilReady is only used for hooks, so it now treats a resource that is not found as done. Wait and WaitWithJobs are unchanged: for regular chart resources a disappearing resource is still an error. Closes #31786 Signed-off-by: ChadiDridi --- pkg/kube/statuswait.go | 28 +++++++++++++--- pkg/kube/statuswait_test.go | 67 +++++++++++++++++++++++++++++++++++++ 2 files changed, 91 insertions(+), 4 deletions(-) diff --git a/pkg/kube/statuswait.go b/pkg/kube/statuswait.go index 91eec30ef..dda1d2c40 100644 --- a/pkg/kube/statuswait.go +++ b/pkg/kube/statuswait.go @@ -95,7 +95,11 @@ func (w *statusWaiter) WatchUntilReady(resourceList ResourceList, timeout time.D StatusReaders: append(w.readers, jobSR, podSR, genericSR), } sw.StatusReader = sr - return w.wait(ctx, resourceList, sw) + // Hook resources may legitimately disappear while Helm is waiting for them. + // A Job hook with .spec.ttlSecondsAfterFinished is removed by the TTL + // controller as soon as it completes, so a hook that is gone is treated as + // done rather than as a resource that never became ready. + return w.waitFor(ctx, resourceList, sw, true) } func (w *statusWaiter) Wait(resourceList ResourceList, timeout time.Duration) error { @@ -154,7 +158,7 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL RESTScopeStrategy: watcher.RESTScopeNamespace, }) statusCollector := collector.NewResourceStatusCollector(resources) - done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.NotFoundStatus, w.Logger())) + done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.NotFoundStatus, false, w.Logger())) <-done if statusCollector.Error != nil { @@ -180,6 +184,13 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL } func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw watcher.StatusWatcher) error { + return w.waitFor(ctx, resourceList, sw, false) +} + +// waitFor waits until every resource reaches the current status. When +// deletedIsDone is set, a resource that is not found is considered done +// instead of blocking until the timeout expires. +func (w *statusWaiter) waitFor(ctx context.Context, resourceList ResourceList, sw watcher.StatusWatcher, deletedIsDone bool) error { cancelCtx, cancel := context.WithCancel(ctx) defer cancel() resources := []object.ObjMetadata{} @@ -198,7 +209,7 @@ func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw w RESTScopeStrategy: watcher.RESTScopeNamespace, }) statusCollector := collector.NewResourceStatusCollector(resources) - done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.CurrentStatus, w.Logger())) + done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.CurrentStatus, deletedIsDone, w.Logger())) <-done if statusCollector.Error != nil { @@ -211,6 +222,9 @@ func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw w if rs.Status == status.CurrentStatus { continue } + if deletedIsDone && rs.Status == status.NotFoundStatus { + continue + } errs = append(errs, fmt.Errorf("resource %s/%s/%s not ready. status: %s, message: %s", rs.Identifier.GroupKind.Kind, rs.Identifier.Namespace, rs.Identifier.Name, rs.Status, rs.Message)) } @@ -237,7 +251,7 @@ func contextWithTimeout(ctx context.Context, timeout time.Duration) (context.Con return watchtools.ContextWithOptionalTimeout(ctx, timeout) } -func statusObserver(cancel context.CancelFunc, desired status.Status, logger *slog.Logger) collector.ObserverFunc { +func statusObserver(cancel context.CancelFunc, desired status.Status, deletedIsDone bool, logger *slog.Logger) collector.ObserverFunc { return func(statusCollector *collector.ResourceStatusCollector, _ event.Event) { var rss []*event.ResourceStatus var nonDesiredResources []*event.ResourceStatus @@ -255,6 +269,12 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, logger *sl if rs.Status == status.FailedStatus && desired == status.CurrentStatus { continue } + // A hook that is gone has finished its job: the TTL controller + // removes completed Jobs that set .spec.ttlSecondsAfterFinished, and + // the legacy waiter treated a delete event as the end of the wait. + if deletedIsDone && rs.Status == status.NotFoundStatus { + continue + } rss = append(rss, rs) if rs.Status != desired { nonDesiredResources = append(nonDesiredResources, rs) diff --git a/pkg/kube/statuswait_test.go b/pkg/kube/statuswait_test.go index 5f5f5d051..6d2851ee5 100644 --- a/pkg/kube/statuswait_test.go +++ b/pkg/kube/statuswait_test.go @@ -1777,3 +1777,70 @@ func TestWatchUntilReadyWithCustomReaders(t *testing.T) { }) } } + +// TestWatchUntilReadyHookDeletedWhileWaiting covers a Job hook that sets +// .spec.ttlSecondsAfterFinished: the TTL controller removes the Job as soon as +// it completes, so the hook can disappear while Helm is still waiting for it. +// The wait has to end there instead of running until the timeout expires. +func TestWatchUntilReadyHookDeletedWhileWaiting(t *testing.T) { + t.Parallel() + c := newTestClient(t) + timeout := 3 * time.Second + timeUntilJobDelete := 500 * time.Millisecond + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper( + batchv1.SchemeGroupVersion.WithKind("Job"), + ) + statusWaiter := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + statusWaiter.SetLogger(slog.Default().Handler()) + + // The Job never reports completion, so only its deletion can end the wait. + objs := getRuntimeObjFromManifests(t, []string{jobNoStatusManifest}) + u := objs[0].(*unstructured.Unstructured) + gvr := getGVR(t, fakeMapper, u) + require.NoError(t, fakeClient.Tracker().Create(gvr, u, u.GetNamespace())) + + go func() { + time.Sleep(timeUntilJobDelete) + assert.NoError(t, fakeClient.Tracker().Delete(gvr, u.GetNamespace(), u.GetName())) + }() + + resourceList := getResourceListFromRuntimeObjs(t, c, objs) + start := time.Now() + require.NoError(t, statusWaiter.WatchUntilReady(resourceList, timeout)) + assert.Less(t, time.Since(start), timeout, "wait should end when the hook is deleted, not on timeout") +} + +// TestStatusWaitDeletedResourceStillFails makes sure the hook behaviour above +// does not leak into Wait, where a resource that goes away is still an error. +func TestStatusWaitDeletedResourceStillFails(t *testing.T) { + t.Parallel() + c := newTestClient(t) + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper( + batchv1.SchemeGroupVersion.WithKind("Job"), + ) + statusWaiter := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + statusWaiter.SetLogger(slog.Default().Handler()) + + objs := getRuntimeObjFromManifests(t, []string{jobNoStatusManifest}) + u := objs[0].(*unstructured.Unstructured) + gvr := getGVR(t, fakeMapper, u) + require.NoError(t, fakeClient.Tracker().Create(gvr, u, u.GetNamespace())) + + go func() { + time.Sleep(200 * time.Millisecond) + assert.NoError(t, fakeClient.Tracker().Delete(gvr, u.GetNamespace(), u.GetName())) + }() + + resourceList := getResourceListFromRuntimeObjs(t, c, objs) + err := statusWaiter.Wait(resourceList, time.Second) + require.Error(t, err) + require.ErrorContains(t, err, "resource Job/qual/test not ready. status: NotFound") +} From e1da86fe7390ea761bce8a0e795eafe50912cfbb Mon Sep 17 00:00:00 2001 From: ChadiDridi Date: Sun, 6 Sep 2026 22:22:16 +0100 Subject: [PATCH 2/3] fix(kube): scope the hook wait exception to TTL Jobs The first version accepted NotFound for any resource passed to WatchUntilReady, which also accepted a Pod hook, a Job without a TTL, or a hook deleted by an operator before it finished, and accepted a hook that was already missing when the wait started. The exception is now limited to the Jobs in the wait that set .spec.ttlSecondsAfterFinished, which are the ones the TTL controller removes when they complete, and only once the hook has actually been observed on the cluster. Everything else that disappears, and a TTL Job that was never seen running, still fail the wait. Signed-off-by: ChadiDridi --- pkg/kube/statuswait.go | 98 ++++++++++++++++++++++------ pkg/kube/statuswait_test.go | 124 +++++++++++++++++++++++++++--------- 2 files changed, 173 insertions(+), 49 deletions(-) diff --git a/pkg/kube/statuswait.go b/pkg/kube/statuswait.go index dda1d2c40..f1ac06601 100644 --- a/pkg/kube/statuswait.go +++ b/pkg/kube/statuswait.go @@ -22,6 +22,7 @@ import ( "fmt" "log/slog" "sort" + "sync" "time" "github.com/fluxcd/cli-utils/pkg/kstatus/polling/aggregator" @@ -33,6 +34,7 @@ import ( "github.com/fluxcd/cli-utils/pkg/kstatus/watcher" "github.com/fluxcd/cli-utils/pkg/object" appsv1 "k8s.io/api/apps/v1" + batchv1 "k8s.io/api/batch/v1" "k8s.io/apimachinery/pkg/api/meta" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/client-go/dynamic" @@ -95,11 +97,11 @@ func (w *statusWaiter) WatchUntilReady(resourceList ResourceList, timeout time.D StatusReaders: append(w.readers, jobSR, podSR, genericSR), } sw.StatusReader = sr - // Hook resources may legitimately disappear while Helm is waiting for them. - // A Job hook with .spec.ttlSecondsAfterFinished is removed by the TTL - // controller as soon as it completes, so a hook that is gone is treated as - // done rather than as a resource that never became ready. - return w.waitFor(ctx, resourceList, sw, true) + // A Job hook that sets .spec.ttlSecondsAfterFinished is removed by the TTL + // controller as soon as it completes, so it can disappear while Helm is + // still waiting for it. Only those hooks may report NotFound instead of + // completing; every other hook that goes away is still an error. + return w.waitFor(ctx, resourceList, sw, ttlJobs(resourceList)) } func (w *statusWaiter) Wait(resourceList ResourceList, timeout time.Duration) error { @@ -158,7 +160,7 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL RESTScopeStrategy: watcher.RESTScopeNamespace, }) statusCollector := collector.NewResourceStatusCollector(resources) - done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.NotFoundStatus, false, w.Logger())) + done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.NotFoundStatus, nil, &observedResources{}, w.Logger())) <-done if statusCollector.Error != nil { @@ -184,13 +186,35 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL } func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw watcher.StatusWatcher) error { - return w.waitFor(ctx, resourceList, sw, false) + return w.waitFor(ctx, resourceList, sw, nil) } -// waitFor waits until every resource reaches the current status. When -// deletedIsDone is set, a resource that is not found is considered done -// instead of blocking until the timeout expires. -func (w *statusWaiter) waitFor(ctx context.Context, resourceList ResourceList, sw watcher.StatusWatcher, deletedIsDone bool) error { +// ttlJobs returns the identities of the Jobs in resourceList that set +// .spec.ttlSecondsAfterFinished, which the TTL controller deletes once they +// finish. +func ttlJobs(resourceList ResourceList) map[object.ObjMetadata]struct{} { + var ttl map[object.ObjMetadata]struct{} + for _, r := range resourceList { + job, ok := AsVersioned(r).(*batchv1.Job) + if !ok || job.Spec.TTLSecondsAfterFinished == nil { + continue + } + obj, err := object.RuntimeToObjMeta(r.Object) + if err != nil { + continue + } + if ttl == nil { + ttl = map[object.ObjMetadata]struct{}{} + } + ttl[obj] = struct{}{} + } + return ttl +} + +// waitFor waits until every resource reaches the current status. Resources in +// deletedIsDone are considered done when they are not found, instead of +// blocking until the timeout expires. +func (w *statusWaiter) waitFor(ctx context.Context, resourceList ResourceList, sw watcher.StatusWatcher, deletedIsDone map[object.ObjMetadata]struct{}) error { cancelCtx, cancel := context.WithCancel(ctx) defer cancel() resources := []object.ObjMetadata{} @@ -209,7 +233,8 @@ func (w *statusWaiter) waitFor(ctx context.Context, resourceList ResourceList, s RESTScopeStrategy: watcher.RESTScopeNamespace, }) statusCollector := collector.NewResourceStatusCollector(resources) - done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.CurrentStatus, deletedIsDone, w.Logger())) + observed := &observedResources{} + done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.CurrentStatus, deletedIsDone, observed, w.Logger())) <-done if statusCollector.Error != nil { @@ -222,8 +247,10 @@ func (w *statusWaiter) waitFor(ctx context.Context, resourceList ResourceList, s if rs.Status == status.CurrentStatus { continue } - if deletedIsDone && rs.Status == status.NotFoundStatus { - continue + if rs.Status == status.NotFoundStatus && observed.wasPresent(id) { + if _, ok := deletedIsDone[id]; ok { + continue + } } errs = append(errs, fmt.Errorf("resource %s/%s/%s not ready. status: %s, message: %s", rs.Identifier.GroupKind.Kind, rs.Identifier.Namespace, rs.Identifier.Name, rs.Status, rs.Message)) @@ -251,7 +278,31 @@ func contextWithTimeout(ctx context.Context, timeout time.Duration) (context.Con return watchtools.ContextWithOptionalTimeout(ctx, timeout) } -func statusObserver(cancel context.CancelFunc, desired status.Status, deletedIsDone bool, logger *slog.Logger) collector.ObserverFunc { +// observedResources records the resources that were seen on the cluster while +// waiting, so that a resource which disappears can be told apart from one that +// was never there. +type observedResources struct { + mu sync.Mutex + present map[object.ObjMetadata]struct{} +} + +func (o *observedResources) markPresent(id object.ObjMetadata) { + o.mu.Lock() + defer o.mu.Unlock() + if o.present == nil { + o.present = map[object.ObjMetadata]struct{}{} + } + o.present[id] = struct{}{} +} + +func (o *observedResources) wasPresent(id object.ObjMetadata) bool { + o.mu.Lock() + defer o.mu.Unlock() + _, ok := o.present[id] + return ok +} + +func statusObserver(cancel context.CancelFunc, desired status.Status, deletedIsDone map[object.ObjMetadata]struct{}, observed *observedResources, logger *slog.Logger) collector.ObserverFunc { return func(statusCollector *collector.ResourceStatusCollector, _ event.Event) { var rss []*event.ResourceStatus var nonDesiredResources []*event.ResourceStatus @@ -269,11 +320,18 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, deletedIsD if rs.Status == status.FailedStatus && desired == status.CurrentStatus { continue } - // A hook that is gone has finished its job: the TTL controller - // removes completed Jobs that set .spec.ttlSecondsAfterFinished, and - // the legacy waiter treated a delete event as the end of the wait. - if deletedIsDone && rs.Status == status.NotFoundStatus { - continue + if rs.Status != status.NotFoundStatus { + observed.markPresent(rs.Identifier) + } + // A Job hook that sets .spec.ttlSecondsAfterFinished is deleted by + // the TTL controller once it completes, so its disappearance ends + // the wait rather than blocking it. This only applies to a hook that + // was seen on the cluster first: one that is already gone when the + // wait starts was never observed running and is still an error. + if rs.Status == status.NotFoundStatus && observed.wasPresent(rs.Identifier) { + if _, ok := deletedIsDone[rs.Identifier]; ok { + continue + } } rss = append(rss, rs) if rs.Status != desired { diff --git a/pkg/kube/statuswait_test.go b/pkg/kube/statuswait_test.go index 6d2851ee5..cbfbfc2f4 100644 --- a/pkg/kube/statuswait_test.go +++ b/pkg/kube/statuswait_test.go @@ -1778,40 +1778,106 @@ func TestWatchUntilReadyWithCustomReaders(t *testing.T) { } } -// TestWatchUntilReadyHookDeletedWhileWaiting covers a Job hook that sets -// .spec.ttlSecondsAfterFinished: the TTL controller removes the Job as soon as -// it completes, so the hook can disappear while Helm is still waiting for it. -// The wait has to end there instead of running until the timeout expires. -func TestWatchUntilReadyHookDeletedWhileWaiting(t *testing.T) { +var jobTTLNoStatusManifest = ` +apiVersion: batch/v1 +kind: Job +metadata: + name: test + namespace: qual + generation: 1 +spec: + ttlSecondsAfterFinished: 0 +` + +// TestWatchUntilReadyHookDeleted covers hooks that disappear while Helm waits. +// A Job that sets .spec.ttlSecondsAfterFinished is removed by the TTL +// controller as soon as it completes, so its deletion ends the wait. Any other +// hook that goes away, and a TTL Job that was never seen running, still fail. +func TestWatchUntilReadyHookDeleted(t *testing.T) { t.Parallel() - c := newTestClient(t) - timeout := 3 * time.Second - timeUntilJobDelete := 500 * time.Millisecond - fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) - fakeMapper := testutil.NewFakeRESTMapper( - batchv1.SchemeGroupVersion.WithKind("Job"), - ) - statusWaiter := statusWaiter{ - restMapper: fakeMapper, - client: fakeClient, + tests := []struct { + name string + manifest string + create bool + expectErrStrs []string + }{ + { + name: "TTL Job deleted while waiting is done", + manifest: jobTTLNoStatusManifest, + create: true, + }, + { + name: "Job without TTL deleted while waiting fails", + manifest: jobNoStatusManifest, + create: true, + expectErrStrs: []string{ + "resource Job/qual/test not ready. status: NotFound", + "context deadline exceeded", + }, + }, + { + name: "Pod hook deleted while waiting fails", + manifest: podNoStatusManifest, + create: true, + expectErrStrs: []string{ + "resource Pod/ns/in-progress-pod not ready. status: NotFound", + "context deadline exceeded", + }, + }, + { + name: "TTL Job that was never running fails", + manifest: jobTTLNoStatusManifest, + create: false, + expectErrStrs: []string{ + "resource Job/qual/test not ready.", + "context deadline exceeded", + }, + }, } - statusWaiter.SetLogger(slog.Default().Handler()) - // The Job never reports completion, so only its deletion can end the wait. - objs := getRuntimeObjFromManifests(t, []string{jobNoStatusManifest}) - u := objs[0].(*unstructured.Unstructured) - gvr := getGVR(t, fakeMapper, u) - require.NoError(t, fakeClient.Tracker().Create(gvr, u, u.GetNamespace())) + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + c := newTestClient(t) + timeout := 3 * time.Second + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper( + batchv1.SchemeGroupVersion.WithKind("Job"), + v1.SchemeGroupVersion.WithKind("Pod"), + ) + statusWaiter := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + statusWaiter.SetLogger(slog.Default().Handler()) - go func() { - time.Sleep(timeUntilJobDelete) - assert.NoError(t, fakeClient.Tracker().Delete(gvr, u.GetNamespace(), u.GetName())) - }() + // The hook never reports completion, so only its deletion can end + // the wait. + objs := getRuntimeObjFromManifests(t, []string{tt.manifest}) + u := objs[0].(*unstructured.Unstructured) + gvr := getGVR(t, fakeMapper, u) + if tt.create { + require.NoError(t, fakeClient.Tracker().Create(gvr, u, u.GetNamespace())) + go func() { + time.Sleep(500 * time.Millisecond) + assert.NoError(t, fakeClient.Tracker().Delete(gvr, u.GetNamespace(), u.GetName())) + }() + } - resourceList := getResourceListFromRuntimeObjs(t, c, objs) - start := time.Now() - require.NoError(t, statusWaiter.WatchUntilReady(resourceList, timeout)) - assert.Less(t, time.Since(start), timeout, "wait should end when the hook is deleted, not on timeout") + resourceList := getResourceListFromRuntimeObjs(t, c, objs) + start := time.Now() + err := statusWaiter.WatchUntilReady(resourceList, timeout) + if tt.expectErrStrs != nil { + require.Error(t, err) + for _, expectedErrStr := range tt.expectErrStrs { + require.ErrorContains(t, err, expectedErrStr) + } + return + } + require.NoError(t, err) + assert.Less(t, time.Since(start), timeout, "wait should end when the hook is deleted, not on timeout") + }) + } } // TestStatusWaitDeletedResourceStillFails makes sure the hook behaviour above From 8a2cd379dbee4357f7d7d76775d973a51aa3a784 Mon Sep 17 00:00:00 2001 From: ChadiDridi Date: Thu, 1 Oct 2026 18:24:34 +0100 Subject: [PATCH 3/3] fix(kube): do not let a deleted hook hide its failure A failed Job is deleted by the TTL controller too, and the delete event replaces Failed with NotFound in the status collector. The deletion exception then matched and the failed hook passed the wait. Failures are now recorded as they are seen, and a resource that is gone only ends the wait when it was seen running and did not fail while it was. Signed-off-by: ChadiDridi --- pkg/kube/statuswait.go | 64 ++++++++++++++++++++++++++--------- pkg/kube/statuswait_test.go | 67 +++++++++++++++++++++++++++++++++++++ 2 files changed, 115 insertions(+), 16 deletions(-) diff --git a/pkg/kube/statuswait.go b/pkg/kube/statuswait.go index f1ac06601..634faff3d 100644 --- a/pkg/kube/statuswait.go +++ b/pkg/kube/statuswait.go @@ -247,10 +247,8 @@ func (w *statusWaiter) waitFor(ctx context.Context, resourceList ResourceList, s if rs.Status == status.CurrentStatus { continue } - if rs.Status == status.NotFoundStatus && observed.wasPresent(id) { - if _, ok := deletedIsDone[id]; ok { - continue - } + if rs.Status == status.NotFoundStatus && observed.deletionIsSuccess(id, deletedIsDone) { + continue } errs = append(errs, fmt.Errorf("resource %s/%s/%s not ready. status: %s, message: %s", rs.Identifier.GroupKind.Kind, rs.Identifier.Namespace, rs.Identifier.Name, rs.Status, rs.Message)) @@ -278,12 +276,14 @@ func contextWithTimeout(ctx context.Context, timeout time.Duration) (context.Con return watchtools.ContextWithOptionalTimeout(ctx, timeout) } -// observedResources records the resources that were seen on the cluster while -// waiting, so that a resource which disappears can be told apart from one that -// was never there. +// observedResources records what was seen on the cluster while waiting, so that +// a resource which disappears can be told apart from one that was never there, +// and so that a failure is not forgotten when the resource is deleted +// afterwards. type observedResources struct { mu sync.Mutex present map[object.ObjMetadata]struct{} + failed map[object.ObjMetadata]struct{} } func (o *observedResources) markPresent(id object.ObjMetadata) { @@ -302,6 +302,33 @@ func (o *observedResources) wasPresent(id object.ObjMetadata) bool { return ok } +func (o *observedResources) markFailed(id object.ObjMetadata) { + o.mu.Lock() + defer o.mu.Unlock() + if o.failed == nil { + o.failed = map[object.ObjMetadata]struct{}{} + } + o.failed[id] = struct{}{} +} + +func (o *observedResources) hasFailed(id object.ObjMetadata) bool { + o.mu.Lock() + defer o.mu.Unlock() + _, ok := o.failed[id] + return ok +} + +// deletionIsSuccess reports whether a resource that is no longer found should +// end the wait successfully: it has to be one of the resources deletion is +// expected for, it has to have been seen running, and it must not have failed +// while it was. +func (o *observedResources) deletionIsSuccess(id object.ObjMetadata, deletedIsDone map[object.ObjMetadata]struct{}) bool { + if _, ok := deletedIsDone[id]; !ok { + return false + } + return o.wasPresent(id) && !o.hasFailed(id) +} + func statusObserver(cancel context.CancelFunc, desired status.Status, deletedIsDone map[object.ObjMetadata]struct{}, observed *observedResources, logger *slog.Logger) collector.ObserverFunc { return func(statusCollector *collector.ResourceStatusCollector, _ event.Event) { var rss []*event.ResourceStatus @@ -315,23 +342,28 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, deletedIsD if rs.Status == status.UnknownStatus && desired == status.NotFoundStatus { continue } + if rs.Status != status.NotFoundStatus { + observed.markPresent(rs.Identifier) + } + // Remember the failure: the resource may be deleted shortly after, + // and the delete event would otherwise replace Failed with NotFound + // and hide it. + if rs.Status == status.FailedStatus { + observed.markFailed(rs.Identifier) + } // Failed is a terminal state. This check ensures we don't wait forever for a resource // that has already failed, as intervention is required to resolve the failure. if rs.Status == status.FailedStatus && desired == status.CurrentStatus { continue } - if rs.Status != status.NotFoundStatus { - observed.markPresent(rs.Identifier) - } // A Job hook that sets .spec.ttlSecondsAfterFinished is deleted by // the TTL controller once it completes, so its disappearance ends // the wait rather than blocking it. This only applies to a hook that - // was seen on the cluster first: one that is already gone when the - // wait starts was never observed running and is still an error. - if rs.Status == status.NotFoundStatus && observed.wasPresent(rs.Identifier) { - if _, ok := deletedIsDone[rs.Identifier]; ok { - continue - } + // was seen running on the cluster and did not fail: one that is + // already gone when the wait starts, or that failed before being + // deleted, is still an error. + if rs.Status == status.NotFoundStatus && observed.deletionIsSuccess(rs.Identifier, deletedIsDone) { + continue } rss = append(rss, rs) if rs.Status != desired { diff --git a/pkg/kube/statuswait_test.go b/pkg/kube/statuswait_test.go index cbfbfc2f4..8df0d2029 100644 --- a/pkg/kube/statuswait_test.go +++ b/pkg/kube/statuswait_test.go @@ -1789,6 +1789,73 @@ spec: ttlSecondsAfterFinished: 0 ` +// TestDeletionIsSuccess covers which disappearing resources may end a wait +// successfully. A failed hook must not pass just because the TTL controller +// removed it afterwards: the delete event replaces Failed with NotFound in the +// collector, so the failure has to be remembered. +func TestDeletionIsSuccess(t *testing.T) { + t.Parallel() + id := object.ObjMetadata{ + GroupKind: batchv1.SchemeGroupVersion.WithKind("Job").GroupKind(), + Namespace: "qual", + Name: "test", + } + other := id + other.Name = "other" + + tests := []struct { + name string + present bool + failed bool + ttl bool + expected bool + }{ + { + name: "TTL Job seen running and then deleted", + present: true, + ttl: true, + expected: true, + }, + { + name: "TTL Job that failed before being deleted", + present: true, + failed: true, + ttl: true, + expected: false, + }, + { + name: "TTL Job that was never seen running", + ttl: true, + expected: false, + }, + { + name: "resource deletion is not expected for", + present: true, + expected: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + observed := &observedResources{} + if tt.present { + observed.markPresent(id) + } + if tt.failed { + observed.markFailed(id) + } + deletedIsDone := map[object.ObjMetadata]struct{}{} + if tt.ttl { + deletedIsDone[id] = struct{}{} + } + assert.Equal(t, tt.expected, observed.deletionIsSuccess(id, deletedIsDone)) + // An unrelated resource is never covered by the exception. + assert.False(t, observed.deletionIsSuccess(other, deletedIsDone)) + }) + } +} + // TestWatchUntilReadyHookDeleted covers hooks that disappear while Helm waits. // A Job that sets .spec.ttlSecondsAfterFinished is removed by the TTL // controller as soon as it completes, so its deletion ends the wait. Any other