diff --git a/pkg/kube/statuswait.go b/pkg/kube/statuswait.go index 91eec30ef..8d5e96f97 100644 --- a/pkg/kube/statuswait.go +++ b/pkg/kube/statuswait.go @@ -33,7 +33,9 @@ import ( "github.com/fluxcd/cli-utils/pkg/kstatus/watcher" "github.com/fluxcd/cli-utils/pkg/object" appsv1 "k8s.io/api/apps/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/client-go/dynamic" watchtools "k8s.io/client-go/tools/watch" @@ -154,7 +156,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, w.deleteStatusObserver(cancelCtx, cancel)) <-done if statusCollector.Error != nil { @@ -245,11 +247,6 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, logger *sl if rs == nil { continue } - // If a resource is already deleted before waiting has started, it will show as unknown. - // This check ensures we don't wait forever for a resource that is already deleted. - if rs.Status == status.UnknownStatus && desired == status.NotFoundStatus { - continue - } // 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 { @@ -267,17 +264,97 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, logger *sl return } - if len(nonDesiredResources) > 0 { - // Log a single resource so the user knows what they're waiting for without an overwhelming amount of output - sort.Slice(nonDesiredResources, func(i, j int) bool { - return nonDesiredResources[i].Identifier.Name < nonDesiredResources[j].Identifier.Name - }) - first := nonDesiredResources[0] - logger.Debug("waiting for resource", "namespace", first.Identifier.Namespace, "name", first.Identifier.Name, "kind", first.Identifier.GroupKind.Kind, "expectedStatus", desired, "actualStatus", first.Status) + logFirstNonDesiredResource(logger, desired, nonDesiredResources) + } +} + +// deleteStatusObserver returns an observer for delete waits, where the desired +// status is NotFound. +// +// UnknownStatus is ambiguous on this path. While the status watcher initializes +// its informer caches, every watched resource briefly reports Unknown, so an +// Unknown-only set must not complete the wait: resources that still exist would +// be reported as deleted before a single real status event arrived (#32261). +// But a resource deleted before the watch started also stays Unknown forever, +// because the watcher never emits an event for an object absent from its +// initial LIST; waiting for a real status event would hang until the timeout +// for resources that are already gone (#32214). +// +// The observer therefore makes no completion decision until the watcher +// delivers its Sync event, which marks the informer caches as populated. From +// then on, any resource still reporting Unknown is confirmed with a live +// lookup: NotFound confirms the deletion, while anything else keeps the wait +// running until the watcher reports a real status for it. +func (w *statusWaiter) deleteStatusObserver(ctx context.Context, cancel context.CancelFunc) collector.ObserverFunc { + desired := status.NotFoundStatus + synced := false + confirmedGone := map[object.ObjMetadata]bool{} + return func(statusCollector *collector.ResourceStatusCollector, e event.Event) { + if e.Type == event.SyncEvent { + synced = true } + if !synced { + return + } + var rss []*event.ResourceStatus + var nonDesiredResources []*event.ResourceStatus + for _, rs := range statusCollector.ResourceStatuses { + if rs == nil { + continue + } + if rs.Status == status.UnknownStatus { + if !confirmedGone[rs.Identifier] && w.isResourceGone(ctx, rs.Identifier) { + confirmedGone[rs.Identifier] = true + } + if confirmedGone[rs.Identifier] { + continue + } + } + rss = append(rss, rs) + if rs.Status != desired { + nonDesiredResources = append(nonDesiredResources, rs) + } + } + + if aggregator.AggregateStatus(rss, desired) == desired { + w.Logger().Debug("all resources achieved desired status", "desiredStatus", desired, "resourceCount", len(rss)) + cancel() + return + } + + logFirstNonDesiredResource(w.Logger(), desired, nonDesiredResources) } } +// isResourceGone reports whether the resource is confirmed absent from the +// cluster by a live lookup. Any error (including transient API errors) reports +// false so the wait keeps running; the lookup is retried on the next event. +func (w *statusWaiter) isResourceGone(ctx context.Context, id object.ObjMetadata) bool { + mapping, err := w.restMapper.RESTMapping(id.GroupKind) + if err != nil { + w.Logger().Debug("unable to map resource to confirm deletion", "namespace", id.Namespace, "name", id.Name, "kind", id.GroupKind.Kind, "error", err) + return false + } + _, err = w.client.Resource(mapping.Resource).Namespace(id.Namespace).Get(ctx, id.Name, metav1.GetOptions{}) + if err != nil && !apierrors.IsNotFound(err) { + w.Logger().Debug("unable to confirm resource deletion", "namespace", id.Namespace, "name", id.Name, "kind", id.GroupKind.Kind, "error", err) + } + return apierrors.IsNotFound(err) +} + +// logFirstNonDesiredResource logs a single resource so the user knows what +// they're waiting for without an overwhelming amount of output +func logFirstNonDesiredResource(logger *slog.Logger, desired status.Status, nonDesiredResources []*event.ResourceStatus) { + if len(nonDesiredResources) == 0 { + return + } + sort.Slice(nonDesiredResources, func(i, j int) bool { + return nonDesiredResources[i].Identifier.Name < nonDesiredResources[j].Identifier.Name + }) + first := nonDesiredResources[0] + logger.Debug("waiting for resource", "namespace", first.Identifier.Namespace, "name", first.Identifier.Name, "kind", first.Identifier.GroupKind.Kind, "expectedStatus", desired, "actualStatus", first.Status) +} + type hookOnlyWaiter struct { sw *statusWaiter } diff --git a/pkg/kube/statuswait_test.go b/pkg/kube/statuswait_test.go index 5f5f5d051..d69ca3a7c 100644 --- a/pkg/kube/statuswait_test.go +++ b/pkg/kube/statuswait_test.go @@ -29,6 +29,7 @@ import ( "github.com/fluxcd/cli-utils/pkg/kstatus/polling/engine" "github.com/fluxcd/cli-utils/pkg/kstatus/polling/event" "github.com/fluxcd/cli-utils/pkg/kstatus/status" + "github.com/fluxcd/cli-utils/pkg/kstatus/watcher" "github.com/fluxcd/cli-utils/pkg/object" "github.com/fluxcd/cli-utils/pkg/testutil" "github.com/stretchr/testify/assert" @@ -259,6 +260,17 @@ metadata: name: test-namespace ` +var hookPodManifest = ` +apiVersion: v1 +kind: Pod +metadata: + name: pre-upgrade-hook + namespace: ns + annotations: + "helm.sh/hook": pre-upgrade + "helm.sh/hook-delete-policy": before-hook-creation +` + func getGVR(t *testing.T, mapper meta.RESTMapper, obj *unstructured.Unstructured) schema.GroupVersionResource { t.Helper() gvk := obj.GroupVersionKind() @@ -375,6 +387,179 @@ func TestStatusWaitForDeleteNonExistentObject(t *testing.T) { assert.NoError(t, statusWaiter.WaitForDelete(resourceList, timeout)) } +// scriptedStatusWatcher emits a fixed sequence of events, then keeps the event +// channel open until the watch context is cancelled. It pins down orderings +// that are racy with the real DefaultStatusWatcher, such as the informer +// initialization window where the Sync event is delivered while every watched +// resource still reports UnknownStatus. +type scriptedStatusWatcher struct { + events []event.Event +} + +func (s *scriptedStatusWatcher) Watch(ctx context.Context, _ object.ObjMetadataSet, _ watcher.Options) <-chan event.Event { + ch := make(chan event.Event) + go func() { + defer close(ch) + for _, e := range s.events { + ch <- e + } + <-ctx.Done() + }() + return ch +} + +// TestStatusWaitForDeleteInformerSync deterministically covers the informer +// initialization window at the root of #32261: every watched resource reports +// UnknownStatus until the watcher's caches sync, and the Sync event can be +// delivered before any real status event. An Unknown-only (or empty, once +// Unknown is filtered) status set must not complete a delete wait for +// resources that still exist, while a resource that is genuinely absent must +// still complete promptly (#32214). +func TestStatusWaitForDeleteInformerSync(t *testing.T) { + t.Parallel() + timeout := time.Second + current := func(id object.ObjMetadata) event.Event { + return event.Event{ + Type: event.ResourceUpdateEvent, + Resource: &event.ResourceStatus{ + Identifier: id, + Status: status.CurrentStatus, + Message: "Resource is current", + }, + } + } + tests := []struct { + name string + createObject bool + failGets bool + events func(id object.ObjMetadata) []event.Event + expectErrs []string + expectPrompt bool + }{ + { + name: "existing resource with all statuses Unknown at sync does not complete", + createObject: true, + events: func(_ object.ObjMetadata) []event.Event { + return []event.Event{{Type: event.SyncEvent}} + }, + expectErrs: []string{"context deadline exceeded"}, + }, + { + name: "status event delivered after sync does not complete", + createObject: true, + events: func(id object.ObjMetadata) []event.Event { + return []event.Event{{Type: event.SyncEvent}, current(id)} + }, + expectErrs: []string{"resource Pod/ns/current-pod still exists. status: Current", "context deadline exceeded"}, + }, + { + name: "resource absent from the cluster completes promptly at sync", + events: func(_ object.ObjMetadata) []event.Event { + return []event.Event{{Type: event.SyncEvent}} + }, + expectPrompt: true, + }, + { + name: "lookup errors do not confirm deletion", + failGets: true, + events: func(_ object.ObjMetadata) []event.Event { + return []event.Event{{Type: event.SyncEvent}} + }, + expectErrs: []string{"context deadline exceeded"}, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + c := newTestClient(t) + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper(v1.SchemeGroupVersion.WithKind("Pod")) + statusWaiter := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + statusWaiter.SetLogger(slog.Default().Handler()) + objs := getRuntimeObjFromManifests(t, []string{podCurrentManifest}) + u := objs[0].(*unstructured.Unstructured) + if tt.createObject { + gvr := getGVR(t, fakeMapper, u) + require.NoError(t, fakeClient.Tracker().Create(gvr, u, u.GetNamespace())) + } + if tt.failGets { + fakeClient.PrependReactor("get", "pods", func(_ clienttesting.Action) (bool, runtime.Object, error) { + return true, nil, errors.New("transient apiserver error") + }) + } + id, err := object.RuntimeToObjMeta(u) + require.NoError(t, err) + sw := &scriptedStatusWatcher{events: tt.events(id)} + resourceList := getResourceListFromRuntimeObjs(t, c, objs) + + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + start := time.Now() + err = statusWaiter.waitForDelete(ctx, resourceList, sw) + elapsed := time.Since(start) + + if tt.expectErrs != nil { + require.Error(t, err) + for _, expectedErrStr := range tt.expectErrs { + require.ErrorContains(t, err, expectedErrStr) + } + // The wait must run until the deadline instead of completing + // on the Unknown-only status set observed at sync. + assert.GreaterOrEqual(t, elapsed, timeout/2, "delete wait completed before any real status event") + } else { + require.NoError(t, err) + } + if tt.expectPrompt { + assert.Less(t, elapsed, timeout/2, "delete wait for an absent resource should complete well before the timeout") + } + }) + } +} + +// TestStatusWaitForDeleteAlreadyDeletedHookReturnsPromptly covers the #32214 +// acceptance criterion for fixing #32261: a hook object that was already +// deleted (e.g. via the before-hook-creation delete policy) never receives a +// status event, because the watcher never emits one for an object absent from +// its initial LIST. The wait must still complete promptly rather than block +// until the timeout, which is the regression that got the previous fix +// (#32081) reverted. +func TestStatusWaitForDeleteAlreadyDeletedHookReturnsPromptly(t *testing.T) { + t.Parallel() + c := newTestClient(t) + // The timeout is deliberately generous: the reverted #32081 made this + // scenario block for the entire timeout, so completion well before this + // timeout is what separates correct behavior from the regression. + timeout := time.Second * 60 + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper( + v1.SchemeGroupVersion.WithKind("Pod"), + ) + statusWaiter := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + statusWaiter.SetLogger(slog.Default().Handler()) + // The hook object is intentionally never created: it was already deleted + // before the wait started, so it is absent from the watcher's initial LIST + // and stays UnknownStatus for the entire watch. + objManifest := getRuntimeObjFromManifests(t, []string{hookPodManifest}) + resourceList := getResourceListFromRuntimeObjs(t, c, objManifest) + + errCh := make(chan error, 1) + go func() { + errCh <- statusWaiter.WaitForDelete(resourceList, timeout) + }() + select { + case err := <-errCh: + require.NoError(t, err) + case <-time.After(time.Second * 10): + t.Fatal("WaitForDelete blocked on an already-deleted hook object; it must return promptly instead of waiting out the timeout") + } +} + func TestStatusWait(t *testing.T) { t.Parallel() tests := []struct {