Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 12 additions & 7 deletions pkg/util/blockingcacheclient/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -218,16 +218,18 @@ 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)
if err != nil {
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) {
Expand All @@ -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
Expand Down
66 changes: 66 additions & 0 deletions pkg/util/blockingcacheclient/client_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}