diff --git a/pkg/kube/statuswait.go b/pkg/kube/statuswait.go index 91eec30ef..28b7766d6 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" @@ -62,6 +64,10 @@ type statusWaiter struct { // when they don't set a timeout. var DefaultStatusWatcherTimeout = 30 * time.Second +// deleteVerificationTimeout bounds the live GETs used to confirm deletions that +// the status watcher did not observe before WaitForDelete gave up. +const deleteVerificationTimeout = 10 * time.Second + func alwaysReady(_ *unstructured.Unstructured) (*status.Result, error) { return &status.Result{ Status: status.CurrentStatus, @@ -162,16 +168,50 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL } errs := []error{} + // A cancelled caller context signals an external interruption, so only + // fall back to live GETs when the watcher stopped on its own or timed out. + verifyWithGet := !errors.Is(ctx.Err(), context.Canceled) + // ctx has usually expired by now, so keep its values but not its deadline, + // and bound all fallback GETs together by a single grace period. + getCtx, cancelGet := context.WithTimeout(context.WithoutCancel(ctx), deleteVerificationTimeout) + defer cancelGet() + unknown, confirmedGone := 0, 0 for _, id := range resources { rs := statusCollector.ResourceStatuses[id] - if rs.Status == status.NotFoundStatus || rs.Status == status.UnknownStatus { + if rs.Status == status.NotFoundStatus { + continue + } + if rs.Status == status.UnknownStatus { + unknown++ continue } + // The watcher may have missed the deletion event (e.g. due to a + // connection drop or informer lag). Verify with a live GET before + // reporting the resource as still existing. + if verifyWithGet { + gone, err := w.isResourceGone(getCtx, id) + if err != nil { + errs = append(errs, fmt.Errorf("resource %s/%s/%s may still exist. status: %s, message: %s: unable to verify deletion: %w", + rs.Identifier.GroupKind.Kind, rs.Identifier.Namespace, rs.Identifier.Name, rs.Status, rs.Message, err)) + continue + } + if gone { + w.Logger().Debug("watcher reported resource as existing but live GET confirms deletion", + "kind", id.GroupKind.Kind, "namespace", id.Namespace, "name", id.Name) + confirmedGone++ + continue + } + } errs = append(errs, fmt.Errorf("resource %s/%s/%s still exists. status: %s, message: %s", rs.Identifier.GroupKind.Kind, rs.Identifier.Namespace, rs.Identifier.Name, rs.Status, rs.Message)) } if err := ctx.Err(); err != nil { - errs = append(errs, err) + // The watcher timing out is not a failure when every resource it did not + // see deleted was confirmed gone by a live GET. + missedDeletesOnly := errors.Is(err, context.DeadlineExceeded) && confirmedGone > 0 && unknown == 0 && len(errs) == 0 + if !missedDeletesOnly { + errs = append(errs, err) + } } if len(errs) > 0 { return errors.Join(errs...) @@ -179,6 +219,21 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL return nil } +// isResourceGone reports whether a live GET returns NotFound for the resource. +// Any other GET error is returned so that callers do not mistake an inability to +// verify deletion for the resource still existing. +func (w *statusWaiter) isResourceGone(ctx context.Context, id object.ObjMetadata) (bool, error) { + mapping, err := w.restMapper.RESTMapping(id.GroupKind) + if err != nil { + return false, err + } + _, err = w.client.Resource(mapping.Resource).Namespace(id.Namespace).Get(ctx, id.Name, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + return true, nil + } + return false, err +} + func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw watcher.StatusWatcher) error { cancelCtx, cancel := context.WithCancel(ctx) defer cancel() diff --git a/pkg/kube/statuswait_test.go b/pkg/kube/statuswait_test.go index 5f5f5d051..cd7015fdc 100644 --- a/pkg/kube/statuswait_test.go +++ b/pkg/kube/statuswait_test.go @@ -22,6 +22,7 @@ import ( "fmt" "log/slog" "strings" + "sync" "sync/atomic" "testing" "time" @@ -375,6 +376,158 @@ func TestStatusWaitForDeleteNonExistentObject(t *testing.T) { assert.NoError(t, statusWaiter.WaitForDelete(resourceList, timeout)) } +func TestWaitForDeleteWithMissedWatchEvent(t *testing.T) { + t.Parallel() + c := newTestClient(t) + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper( + v1.SchemeGroupVersion.WithKind("Pod"), + ) + // Return a watcher with no events to simulate a missed deletion notification. + // The watch starts after the initial list, so signal that the watcher has + // seen the resource before deleting it. + watchStarted := make(chan struct{}) + var once sync.Once + fakeClient.PrependWatchReactor("pods", func(_ clienttesting.Action) (bool, watch.Interface, error) { + once.Do(func() { close(watchStarted) }) + return true, watch.NewFake(), nil + }) + sw := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + sw.SetLogger(slog.Default().Handler()) + objs := getRuntimeObjFromManifests(t, []string{podCurrentManifest}) + gvrs := make([]schema.GroupVersionResource, 0, len(objs)) + for _, obj := range objs { + u := obj.(*unstructured.Unstructured) + gvr := getGVR(t, fakeMapper, u) + err := fakeClient.Tracker().Create(gvr, u, u.GetNamespace()) + require.NoError(t, err) + gvrs = append(gvrs, gvr) + } + // Delete the resource once the watch has started. The watcher will miss + // this deletion because watch events are suppressed, but a live GET should + // confirm the resource is gone. + deleteErrs := make(chan error, 1) + go func() { + <-watchStarted + var errs []error + for i, obj := range objs { + u := obj.(*unstructured.Unstructured) + errs = append(errs, fakeClient.Tracker().Delete(gvrs[i], u.GetNamespace(), u.GetName())) + } + deleteErrs <- errors.Join(errs...) + }() + resourceList := getResourceListFromRuntimeObjs(t, c, objs) + err := sw.WaitForDelete(resourceList, 500*time.Millisecond) + require.NoError(t, <-deleteErrs) + assert.NoError(t, err) +} + +func TestWaitForDeleteWithMissedWatchEventGetError(t *testing.T) { + t.Parallel() + c := newTestClient(t) + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper( + v1.SchemeGroupVersion.WithKind("Pod"), + ) + fakeClient.PrependWatchReactor("pods", func(_ clienttesting.Action) (bool, watch.Interface, error) { + return true, watch.NewFake(), nil + }) + // The live GET cannot verify the deletion, so the resource must not be + // reported as simply still existing, nor as deleted. + fakeClient.PrependReactor("get", "pods", func(_ clienttesting.Action) (bool, runtime.Object, error) { + return true, nil, apierrors.NewForbidden(schema.GroupResource{Resource: "pods"}, "current-pod", errors.New("get not allowed")) + }) + sw := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + sw.SetLogger(slog.Default().Handler()) + objs := getRuntimeObjFromManifests(t, []string{podCurrentManifest}) + for _, obj := range objs { + u := obj.(*unstructured.Unstructured) + gvr := getGVR(t, fakeMapper, u) + err := fakeClient.Tracker().Create(gvr, u, u.GetNamespace()) + require.NoError(t, err) + } + resourceList := getResourceListFromRuntimeObjs(t, c, objs) + err := sw.WaitForDelete(resourceList, 500*time.Millisecond) + require.Error(t, err) + assert.True(t, apierrors.IsForbidden(err), "expected the GET error to be wrapped, got: %v", err) + require.ErrorIs(t, err, context.DeadlineExceeded) + assert.Contains(t, err.Error(), "resource Pod/ns/current-pod may still exist. status: Current") + assert.Contains(t, err.Error(), "unable to verify deletion") +} + +func TestIsResourceGone(t *testing.T) { + t.Parallel() + tests := []struct { + name string + getErr error + exists bool + expected bool + expectErr func(error) bool + }{ + { + name: "not found", + expected: true, + }, + { + name: "still exists", + exists: true, + }, + { + name: "forbidden", + getErr: apierrors.NewForbidden(schema.GroupResource{Resource: "pods"}, "current-pod", errors.New("get not allowed")), + expectErr: apierrors.IsForbidden, + }, + { + name: "timeout", + getErr: context.DeadlineExceeded, + expectErr: func(err error) bool { return errors.Is(err, context.DeadlineExceeded) }, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + fakeClient := dynamicfake.NewSimpleDynamicClient(scheme.Scheme) + fakeMapper := testutil.NewFakeRESTMapper( + v1.SchemeGroupVersion.WithKind("Pod"), + ) + if tt.getErr != nil { + fakeClient.PrependReactor("get", "pods", func(_ clienttesting.Action) (bool, runtime.Object, error) { + return true, nil, tt.getErr + }) + } + sw := statusWaiter{ + restMapper: fakeMapper, + client: fakeClient, + } + sw.SetLogger(slog.Default().Handler()) + if tt.exists { + u := getRuntimeObjFromManifests(t, []string{podCurrentManifest})[0].(*unstructured.Unstructured) + err := fakeClient.Tracker().Create(getGVR(t, fakeMapper, u), u, u.GetNamespace()) + require.NoError(t, err) + } + id := object.ObjMetadata{ + GroupKind: v1.SchemeGroupVersion.WithKind("Pod").GroupKind(), + Namespace: "ns", + Name: "current-pod", + } + gone, err := sw.isResourceGone(t.Context(), id) + if tt.expectErr != nil { + require.Error(t, err) + assert.True(t, tt.expectErr(err), "unexpected error: %v", err) + } else { + require.NoError(t, err) + } + assert.Equal(t, tt.expected, gone) + }) + } +} + func TestStatusWait(t *testing.T) { t.Parallel() tests := []struct {