From 99a9ed081cbaaf67522e7bdca5931b6e62a7888d Mon Sep 17 00:00:00 2001 From: MrJack <36191829+biagiopietro@users.noreply.github.com> Date: Wed, 2 Sep 2026 14:08:38 +0200 Subject: [PATCH 1/2] fix(kube): don't hang waiting on CRDs that are still being established Signed-off-by: MrJack <36191829+biagiopietro@users.noreply.github.com> --- pkg/kube/statuswait.go | 112 +++++++++++++++++++++ pkg/kube/statuswait_crd_test.go | 167 ++++++++++++++++++++++++++++++++ 2 files changed, 279 insertions(+) create mode 100644 pkg/kube/statuswait_crd_test.go diff --git a/pkg/kube/statuswait.go b/pkg/kube/statuswait.go index 91eec30ef..5a865201b 100644 --- a/pkg/kube/statuswait.go +++ b/pkg/kube/statuswait.go @@ -33,10 +33,14 @@ 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" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/dynamic" watchtools "k8s.io/client-go/tools/watch" + "sigs.k8s.io/controller-runtime/pkg/client/apiutil" "helm.sh/helm/v4/internal/logging" helmStatusReaders "helm.sh/helm/v4/internal/statusreaders" @@ -194,6 +198,15 @@ func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw w resources = append(resources, obj) } + // Ensure the RESTMapper can resolve every GroupKind before the status + // watcher starts its informers. Otherwise a CRD that is being created or + // established asynchronously would cause the watcher to permanently + // abandon the informer for that GroupKind, leaving the affected resources + // stuck in the Unknown status until the wait times out. + if err := w.ensureResourceMappingsReady(cancelCtx, resources); err != nil { + return err + } + eventCh := sw.Watch(cancelCtx, resources, watcher.Options{ RESTScopeStrategy: watcher.RESTScopeNamespace, }) @@ -237,6 +250,105 @@ func contextWithTimeout(ctx context.Context, timeout time.Duration) (context.Con return watchtools.ContextWithOptionalTimeout(ctx, timeout) } +// ensureResourceMappingsReady blocks until the RESTMapper can resolve every +// GroupKind present in resources, or until the context is done. +// +// The status watcher starts one informer per GroupKind. If a GroupKind cannot +// be mapped when its informer starts (e.g. because the backing CRD was created +// by the same release and is not established yet), the underlying watcher gives +// up on that informer permanently, leaving the affected resources stuck in the +// Unknown status until the wait times out. Resolving the mappings up front +// avoids that race so informers always start successfully. +func (w *statusWaiter) ensureResourceMappingsReady(ctx context.Context, resources []object.ObjMetadata) error { + // Do not do any work if the wait has already been cancelled. + if err := ctx.Err(); err != nil { + return err + } + groupKinds := make([]schema.GroupKind, 0, len(resources)) + seen := make(map[schema.GroupKind]struct{}, len(resources)) + for _, r := range resources { + if r.GroupKind.Kind == "" { + // A resource without a type cannot be resolved; leave it to the + // watcher to report as not ready. + continue + } + if _, ok := seen[r.GroupKind]; ok { + continue + } + seen[r.GroupKind] = struct{}{} + groupKinds = append(groupKinds, r.GroupKind) + } + if len(groupKinds) == 0 { + return nil + } + + // Fast path: if every GroupKind is already known there is nothing to wait + // for and the status watcher can start immediately. + if mappingsResolved(w.restMapper, groupKinds) { + return nil + } + w.Logger().Debug("waiting for resource types to become available", + "groupKinds", groupKindNames(groupKinds)) + + err := wait.PollUntilContextCancel(ctx, time.Second, true, func(_ context.Context) (bool, error) { + for _, gk := range groupKinds { + if _, err := w.restMapper.RESTMapping(gk); err != nil { + if isRetryableMappingError(err) { + return false, nil + } + return false, err + } + } + return true, nil + }) + if err != nil && !wait.Interrupted(err) { + return err + } + return nil +} + +// mappingsResolved returns true when the RESTMapper can resolve every GroupKind. +func mappingsResolved(mapper meta.RESTMapper, groupKinds []schema.GroupKind) bool { + for _, gk := range groupKinds { + if _, err := mapper.RESTMapping(gk); err != nil { + return false + } + } + return true +} + +// groupKindNames returns the GroupKinds in a human-readable form. +func groupKindNames(groupKinds []schema.GroupKind) []string { + names := make([]string, 0, len(groupKinds)) + for _, gk := range groupKinds { + if gk.Group == "" { + names = append(names, gk.Kind) + } else { + names = append(names, gk.Group+"/"+gk.Kind) + } + } + return names +} + +// isRetryableMappingError returns true if a RESTMapper lookup failed because a +// resource type is not (yet) known to the API server, which is expected while a +// CRD is being created or established. +func isRetryableMappingError(err error) bool { + if meta.IsNoMatchError(err) { + return true + } + discoveryFailed, ok := errors.AsType[*apiutil.ErrResourceDiscoveryFailed](err) + if !ok { + return false + } + for _, gvErr := range *discoveryFailed { + if meta.IsNoMatchError(gvErr) || apierrors.IsNotFound(gvErr) { + return true + } + } + return false +} + func statusObserver(cancel context.CancelFunc, desired status.Status, logger *slog.Logger) collector.ObserverFunc { return func(statusCollector *collector.ResourceStatusCollector, _ event.Event) { var rss []*event.ResourceStatus diff --git a/pkg/kube/statuswait_crd_test.go b/pkg/kube/statuswait_crd_test.go new file mode 100644 index 000000000..3a8525128 --- /dev/null +++ b/pkg/kube/statuswait_crd_test.go @@ -0,0 +1,167 @@ +/* +Copyright The Helm Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package kube + +import ( + "context" + "log/slog" + "sync/atomic" + "testing" + "time" + + "k8s.io/apimachinery/pkg/api/meta" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/cli-runtime/pkg/resource" + dynamicfake "k8s.io/client-go/dynamic/fake" + "k8s.io/kubectl/pkg/scheme" +) + +const rolloutManifest = ` +apiVersion: argoproj.io/v1alpha1 +kind: Rollout +metadata: + name: app + namespace: default + generation: 1 +status: + observedGeneration: 1 + conditions: + - type: Promoted + status: "True" +` + +var rolloutGVR = schema.GroupVersionResource{ + Group: "argoproj.io", + Version: "v1alpha1", + Resource: "rollouts", +} + +var rolloutGK = schema.GroupKind{Group: "argoproj.io", Kind: "Rollout"} + +// delayedMapper simulates a CRD whose GroupKind is not registered in the +// RESTMapper until a short delay has passed (e.g. the CRD is created by the +// same Helm release and takes a moment to be established). Until then +// RESTMapping returns a NoMatchError. +type delayedMapper struct { + meta.RESTMapper + rolloutGK schema.GroupKind + ready atomic.Bool +} + +func (d *delayedMapper) RESTMapping(gk schema.GroupKind, versions ...string) (*meta.RESTMapping, error) { + if gk == d.rolloutGK && !d.ready.Load() { + return nil, &meta.NoResourceMatchError{ + PartialResource: schema.GroupVersionResource{Group: gk.Group, Resource: "rollouts"}, + } + } + return d.RESTMapper.RESTMapping(gk, versions...) +} + +func newDelayedMapper(ready bool) *delayedMapper { + delegate := meta.NewDefaultRESTMapper([]schema.GroupVersion{ + {Group: "argoproj.io", Version: "v1alpha1"}, + }) + delegate.Add(schema.GroupVersionKind{Group: "argoproj.io", Version: "v1alpha1", Kind: "Rollout"}, meta.RESTScopeNamespace) + mapper := &delayedMapper{ + RESTMapper: delegate, + rolloutGK: rolloutGK, + } + mapper.ready.Store(ready) + return mapper +} + +func newRolloutStatusWaiter(t *testing.T, mapper meta.RESTMapper) (*statusWaiter, ResourceList) { + t.Helper() + rollout := getRuntimeObjFromManifests(t, []string{rolloutManifest})[0].(*unstructured.Unstructured) + + fakeClient := dynamicfake.NewSimpleDynamicClientWithCustomListKinds( + scheme.Scheme, + map[schema.GroupVersionResource]string{rolloutGVR: "RolloutList"}, + rollout, + ) + + sw := &statusWaiter{ + client: fakeClient, + restMapper: mapper, + } + sw.SetLogger(slog.Default().Handler()) + + resourceList := ResourceList{ + &resource.Info{ + Object: rollout, + Namespace: rollout.GetNamespace(), + Name: rollout.GetName(), + }, + } + return sw, resourceList +} + +// TestStatusWaitCustomResource ensures that waiting on a custom resource (such +// as an Argo Rollout) does not hang when its CRD is not yet registered in the +// RESTMapper when the wait starts. The wait must succeed as soon as the CRD +// becomes available, instead of leaving the resource in the Unknown status +// until the timeout. +func TestStatusWaitCustomResource(t *testing.T) { + t.Parallel() + tests := []struct { + name string + initReady bool + readyAfter time.Duration + }{ + { + name: "CRD is already registered", + initReady: true, + }, + { + // Simulate the CRD becoming established shortly after the wait + // starts. + name: "CRD is registered while waiting", + readyAfter: 500 * time.Millisecond, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + mapper := newDelayedMapper(tt.initReady) + if tt.readyAfter > 0 { + time.AfterFunc(tt.readyAfter, func() { + mapper.ready.Store(true) + }) + } + + sw, resourceList := newRolloutStatusWaiter(t, mapper) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + done := make(chan error, 1) + go func() { + done <- sw.Wait(resourceList, 5*time.Second) + }() + + select { + case err := <-done: + if err != nil { + t.Fatalf("Wait failed: %v", err) + } + case <-ctx.Done(): + t.Fatal("Wait hung: custom resource remained in the Unknown status despite its CRD becoming available") + } + }) + } +} From 32f33419aa191407bfed660946967aea7924bfc5 Mon Sep 17 00:00:00 2001 From: MrJack <36191829+biagiopietro@users.noreply.github.com> Date: Wed, 2 Sep 2026 15:31:04 +0200 Subject: [PATCH 2/2] fix(kube): don't hang waiting on resources whose status cannot be computed Argo Rollout stores status.observedGeneration as a string, unlike the int64 convention kstatus expects. Helm v4's status wait runs status.Compute for CRDs, which errors on the string value; the error is attached to the ResourceStatus as "Unknown" instead of being surfaced, so the Rollout stays Unknown forever and --wait times out with "waiting for resource ... kind=Rollout ... actualStatus=Unknown". Treat a resource whose status could not be computed (Unknown with an attached error) as ready when waiting for Current, matching Helm 3's behavior of considering kinds Helm cannot evaluate as ready. Also prime the RESTMapper before starting the status watcher's informers so a CRD that is still being established does not cause an informer to be permanently abandoned. Signed-off-by: MrJack <36191829+biagiopietro@users.noreply.github.com> --- pkg/kube/statuswait.go | 26 ++++++++++++++++++++++---- pkg/kube/statuswait_crd_test.go | 30 +++++++++++++++++++++++++++++- 2 files changed, 51 insertions(+), 5 deletions(-) diff --git a/pkg/kube/statuswait.go b/pkg/kube/statuswait.go index 5a865201b..6097521c1 100644 --- a/pkg/kube/statuswait.go +++ b/pkg/kube/statuswait.go @@ -221,7 +221,7 @@ func (w *statusWaiter) wait(ctx context.Context, resourceList ResourceList, sw w errs := []error{} for _, id := range resources { rs := statusCollector.ResourceStatuses[id] - if rs.Status == status.CurrentStatus { + if resourceStatusSatisfied(rs, status.CurrentStatus) { continue } errs = append(errs, fmt.Errorf("resource %s/%s/%s not ready. status: %s, message: %s", @@ -367,10 +367,16 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, logger *sl if rs.Status == status.FailedStatus && desired == status.CurrentStatus { continue } - rss = append(rss, rs) - if rs.Status != desired { - nonDesiredResources = append(nonDesiredResources, rs) + // A resource whose status could not be computed (Unknown with an attached + // error) is treated as having reached the desired state. Helm 3 considered + // all kinds it could not evaluate (e.g. CRDs such as Argo Rollout whose + // status fields deviate from Kubernetes conventions) as ready, so don't + // block the wait on them. + if resourceStatusSatisfied(rs, desired) { + continue } + rss = append(rss, rs) + nonDesiredResources = append(nonDesiredResources, rs) } if aggregator.AggregateStatus(rss, desired) == desired { @@ -390,6 +396,18 @@ func statusObserver(cancel context.CancelFunc, desired status.Status, logger *sl } } +// resourceStatusSatisfied reports whether the given resource status satisfies the +// desired status. A resource whose status could not be computed (Unknown with an +// attached error, e.g. a CRD such as an Argo Rollout whose status fields do not +// follow Kubernetes conventions) is treated as satisfied for the Current status, +// matching Helm 3's behavior of considering such kinds ready. +func resourceStatusSatisfied(rs *event.ResourceStatus, desired status.Status) bool { + if rs.Status == desired { + return true + } + return desired == status.CurrentStatus && rs.Status == status.UnknownStatus && rs.Error != nil +} + type hookOnlyWaiter struct { sw *statusWaiter } diff --git a/pkg/kube/statuswait_crd_test.go b/pkg/kube/statuswait_crd_test.go index 3a8525128..3732a2b9e 100644 --- a/pkg/kube/statuswait_crd_test.go +++ b/pkg/kube/statuswait_crd_test.go @@ -39,10 +39,11 @@ metadata: namespace: default generation: 1 status: - observedGeneration: 1 + observedGeneration: "1" conditions: - type: Promoted status: "True" + lastTransitionTime: "2026-01-01T00:00:00Z" ` var rolloutGVR = schema.GroupVersionResource{ @@ -165,3 +166,30 @@ func TestStatusWaitCustomResource(t *testing.T) { }) } } + +// TestStatusWaitCustomResourceUncomputableStatus ensures that waiting on a custom +// resource whose status cannot be computed by the kstatus library (such as an +// Argo Rollout, which stores status.observedGeneration as a string) does not +// hang. Helm 3 considered kinds it could not evaluate as ready, so the wait must +// succeed instead of leaving the resource in the Unknown status until the timeout. +func TestStatusWaitCustomResourceUncomputableStatus(t *testing.T) { + t.Parallel() + mapper := newDelayedMapper(true) + sw, resourceList := newRolloutStatusWaiter(t, mapper) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + done := make(chan error, 1) + go func() { + done <- sw.Wait(resourceList, 5*time.Second) + }() + + select { + case err := <-done: + if err != nil { + t.Fatalf("Wait failed: %v", err) + } + case <-ctx.Done(): + t.Fatal("Wait hung: custom resource remained in the Unknown status because its status could not be computed") + } +}