pull/32222/merge
Carter Pewterschmidt 2 days ago committed by GitHub
commit 0eae3a86bd
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -33,7 +33,9 @@ import (
"github.com/fluxcd/cli-utils/pkg/kstatus/watcher" "github.com/fluxcd/cli-utils/pkg/kstatus/watcher"
"github.com/fluxcd/cli-utils/pkg/object" "github.com/fluxcd/cli-utils/pkg/object"
appsv1 "k8s.io/api/apps/v1" appsv1 "k8s.io/api/apps/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta" "k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/client-go/dynamic" "k8s.io/client-go/dynamic"
watchtools "k8s.io/client-go/tools/watch" watchtools "k8s.io/client-go/tools/watch"
@ -62,6 +64,10 @@ type statusWaiter struct {
// when they don't set a timeout. // when they don't set a timeout.
var DefaultStatusWatcherTimeout = 30 * time.Second 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) { func alwaysReady(_ *unstructured.Unstructured) (*status.Result, error) {
return &status.Result{ return &status.Result{
Status: status.CurrentStatus, Status: status.CurrentStatus,
@ -162,16 +168,50 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL
} }
errs := []error{} 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 { for _, id := range resources {
rs := statusCollector.ResourceStatuses[id] 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 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", 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)) rs.Identifier.GroupKind.Kind, rs.Identifier.Namespace, rs.Identifier.Name, rs.Status, rs.Message))
} }
if err := ctx.Err(); err != nil { 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 { if len(errs) > 0 {
return errors.Join(errs...) return errors.Join(errs...)
@ -179,6 +219,21 @@ func (w *statusWaiter) waitForDelete(ctx context.Context, resourceList ResourceL
return nil 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 { func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw watcher.StatusWatcher) error {
cancelCtx, cancel := context.WithCancel(ctx) cancelCtx, cancel := context.WithCancel(ctx)
defer cancel() defer cancel()

@ -22,6 +22,7 @@ import (
"fmt" "fmt"
"log/slog" "log/slog"
"strings" "strings"
"sync"
"sync/atomic" "sync/atomic"
"testing" "testing"
"time" "time"
@ -375,6 +376,158 @@ func TestStatusWaitForDeleteNonExistentObject(t *testing.T) {
assert.NoError(t, statusWaiter.WaitForDelete(resourceList, timeout)) 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) { func TestStatusWait(t *testing.T) {
t.Parallel() t.Parallel()
tests := []struct { tests := []struct {

Loading…
Cancel
Save