diff --git a/pkg/util/blockingcacheclient/client.go b/pkg/util/blockingcacheclient/client.go index e869ef2ba9..e85abcf81c 100644 --- a/pkg/util/blockingcacheclient/client.go +++ b/pkg/util/blockingcacheclient/client.go @@ -218,8 +218,10 @@ func (c *CacheClient) newEmptyObjectFor(from client.Object) (client.Object, erro // blockApply waits until the applied object appears in the cache with the expected state. // clientObj must be non-nil (caller must have extracted it from the ApplyConfiguration). // preApplyMeta is the object's metadata from a GET before Apply; if nil (e.g. object did not exist), -// we consider the cache updated once the object exists. Otherwise we compare until UID/Generation/ResourceVersion -// differ so the cache has observed the Apply. +// we consider the cache updated once the object exists. Otherwise we wait until +// UID/Generation/ResourceVersion differ so the cache has observed the Apply. +// A no-op SSA leaves those unchanged, so the wait times out; if the object still +// exists, Apply() already succeeded and we treat that as synced. func (c *CacheClient) blockApply(ctx context.Context, obj runtime.ApplyConfiguration, clientObj client.Object, preApplyMeta metav1.Object) error { nn := types.NamespacedName{Namespace: clientObj.GetNamespace(), Name: clientObj.GetName()} newObj, err := c.newEmptyObjectFor(clientObj) @@ -227,7 +229,7 @@ func (c *CacheClient) blockApply(ctx context.Context, obj runtime.ApplyConfigura return err } - return wait.PollUntilContextTimeout(ctx, time.Millisecond*10, time.Second*2, true, func(context.Context) (bool, error) { + err = wait.PollUntilContextTimeout(ctx, time.Millisecond*10, time.Second*2, true, func(context.Context) (bool, error) { err := c.Client.Get(ctx, nn, newObj) if err != nil { if runtime.IsNotRegisteredError(err) { @@ -253,14 +255,17 @@ func (c *CacheClient) blockApply(ctx context.Context, obj runtime.ApplyConfigura return false, err } // Cache has applied state when UID/Generation/ResourceVersion changed from pre-apply. - // Condition 1: UID changed - object was deleted and recreated - // Condition 2: Generation increased - spec was updated - // Condition 3: ResourceVersion changed - any update occurred (metadata, spec, or status) - // If any of these conditions are true, the Apply operation is reflected in the cache. return preApplyMeta.GetUID() != newAccessor.GetUID() || newAccessor.GetGeneration() > preApplyMeta.GetGeneration() || newAccessor.GetResourceVersion() != preApplyMeta.GetResourceVersion(), nil }) + if err == nil { + return nil + } + if getErr := c.Client.Get(ctx, nn, newObj); getErr == nil { + return nil + } + return err } // TODO: implement DeleteAllOf diff --git a/pkg/util/blockingcacheclient/client_test.go b/pkg/util/blockingcacheclient/client_test.go new file mode 100644 index 0000000000..ca5c935352 --- /dev/null +++ b/pkg/util/blockingcacheclient/client_test.go @@ -0,0 +1,66 @@ +package blockingcacheclient + +import ( + "context" + "testing" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" +) + +func TestBlockApplyNoopDoesNotTimeout(t *testing.T) { + obj := &corev1.ConfigMap{ + TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "ConfigMap"}, + ObjectMeta: metav1.ObjectMeta{ + Name: "vcluster-protected-apiservices", + UID: "uid-1", + ResourceVersion: "1", + Generation: 1, + }, + } + + inner := fake.NewClientBuilder().WithObjects(obj.DeepCopy()).Build() + c := &CacheClient{Client: inner, scheme: inner.Scheme()} + + pre := obj.DeepCopy() + if err := c.blockApply(context.Background(), nil, obj, pre); err != nil { + t.Fatalf("no-op apply with unchanged resourceVersion should succeed, got %v", err) + } +} + +func TestBlockApplyWaitsUntilCreated(t *testing.T) { + obj := &corev1.ConfigMap{ + TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "ConfigMap"}, + ObjectMeta: metav1.ObjectMeta{Name: "created"}, + } + + inner := fake.NewClientBuilder().Build() + gets := 0 + wrapped := interceptor.NewClient(inner, interceptor.Funcs{ + Get: func(ctx context.Context, c client.WithWatch, key types.NamespacedName, out client.Object, opts ...client.GetOption) error { + gets++ + if gets >= 3 { + if err := inner.Create(ctx, obj.DeepCopy()); err != nil { + return err + } + } + return inner.Get(ctx, key, out, opts...) + }, + }) + c := &CacheClient{Client: wrapped, scheme: runtime.NewScheme()} + if err := corev1.AddToScheme(c.scheme); err != nil { + t.Fatal(err) + } + + if err := c.blockApply(context.Background(), nil, obj, nil); err != nil { + t.Fatalf("create should succeed once the object appears, got %v", err) + } + if gets < 3 { + t.Fatalf("expected polling before create, got %d gets", gets) + } +}