pull/32586/merge
Chadi Dridi 2 days ago committed by GitHub
commit 59f26e147d
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -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,7 +97,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)
// 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 {
@ -154,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, w.Logger()))
done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.NotFoundStatus, nil, &observedResources{}, w.Logger()))
<-done
if statusCollector.Error != nil {
@ -180,6 +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, nil)
}
// 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{}
@ -198,7 +233,8 @@ 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()))
observed := &observedResources{}
done := statusCollector.ListenWithObserver(eventCh, statusObserver(cancel, status.CurrentStatus, deletedIsDone, observed, w.Logger()))
<-done
if statusCollector.Error != nil {
@ -211,6 +247,9 @@ func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw w
if rs.Status == status.CurrentStatus {
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))
}
@ -237,7 +276,60 @@ 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 {
// 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) {
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 (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
var nonDesiredResources []*event.ResourceStatus
@ -250,11 +342,29 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, logger *sl
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
}
// 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 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 {
nonDesiredResources = append(nonDesiredResources, rs)

@ -1777,3 +1777,203 @@ func TestWatchUntilReadyWithCustomReaders(t *testing.T) {
})
}
}
var jobTTLNoStatusManifest = `
apiVersion: batch/v1
kind: Job
metadata:
name: test
namespace: qual
generation: 1
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
// hook that goes away, and a TTL Job that was never seen running, still fail.
func TestWatchUntilReadyHookDeleted(t *testing.T) {
t.Parallel()
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",
},
},
}
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())
// 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()
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
// 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")
}

Loading…
Cancel
Save