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
38 changes: 33 additions & 5 deletions pkg/platform/oap/install/apply.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/controller-runtime/pkg/client"

"github.com/authzed/openagentprimitives/pkg/apis/v1alpha1"
Expand Down Expand Up @@ -569,11 +570,38 @@ func secretUnstructured(s SecretSpec, ns string) (*unstructured.Unstructured, er
func ssaApply(ctx context.Context, c client.Client, obj *unstructured.Unstructured, identity objectIdentity) error {
// UID prevents a same-name replacement from being seized. ResourceVersion
// makes the apply conditional on the exact state Prepare approved (or Create
// returned), so disappearance, replacement, or an intervening update fails
// closed instead of turning SSA into an untracked create/overwrite.
obj.SetUID(identity.uid)
obj.SetResourceVersion(identity.resourceVersion)
if err := c.Patch(ctx, obj, client.Apply, client.FieldOwner(FieldManager), client.ForceOwnership); err != nil {
// returned), so disappearance or replacement fails closed instead of
// turning SSA into an untracked create/overwrite.
//
// A bare resourceVersion conflict, though, is not by itself foul play: the
// object's own controller watches these kinds and routinely writes to a
// fresh object (a status condition, a finalizer) between that observation
// and this apply. Losing that race used to abort the whole install — a
// recurring e2e flake. So a conflict re-reads the object: the same UID
// proves it is still the one this run approved, and the apply retries on
// the current resourceVersion; a different UID (or a failed re-read) is a
// mid-install replacement or disappearance and still fails closed.
rv := identity.resourceVersion
err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
applied := obj.DeepCopy()
applied.SetUID(identity.uid)
applied.SetResourceVersion(rv)
patchErr := c.Patch(ctx, applied, client.Apply, client.FieldOwner(FieldManager), client.ForceOwnership)
if patchErr == nil || !apierrors.IsConflict(patchErr) {
return patchErr
}
current := &unstructured.Unstructured{}
current.SetGroupVersionKind(obj.GroupVersionKind())
if getErr := c.Get(ctx, client.ObjectKeyFromObject(obj), current); getErr != nil {
return fmt.Errorf("re-read after apply conflict: %w", getErr)
}
if current.GetUID() != identity.uid {
return fmt.Errorf("object was replaced during install (uid %s, expected %s)", current.GetUID(), identity.uid)
}
rv = current.GetResourceVersion()
return patchErr
})
if err != nil {
return fmt.Errorf("ssa-apply %s/%s: %w", obj.GetKind(), obj.GetName(), err)
}
return nil
Expand Down
139 changes: 139 additions & 0 deletions pkg/platform/oap/install/apply_conflict_retry_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,139 @@
package install_test

import (
"context"
"errors"
"strings"
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
"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"

"github.com/authzed/openagentprimitives/pkg/platform/oap"
"github.com/authzed/openagentprimitives/pkg/platform/oap/install"
)

// agentClassConflict builds the optimistic-concurrency error the apiserver
// returns when an SSA apply carries a resourceVersion that a concurrent write
// has already advanced past.
func agentClassConflict(name string) error {
return apierrors.NewConflict(
schema.GroupResource{Group: "agentprimitives.authzed.com", Resource: "agentclasses"},
name,
// The apiserver's registry.OptimisticLockErrorMsg, verbatim.
errors.New("the object has been modified; please apply your changes to the latest version and try again"))
}

// TestInstall_ControllerWriteBetweenCreateAndApply_RetriesOnCurrentVersion pins
// the fix for a recurring e2e flake: applyPreparedObject Creates an object and
// then SSA-applies it conditional on the Create-time resourceVersion, but the
// operator's own controllers watch these kinds and routinely write to the fresh
// object (status, finalizers) in that window. The resulting conflict is benign
// — the object is still the one this run created — so Install must re-read the
// current resourceVersion and retry, not abort.
//
// Fake-client unit test for the same reason as the adopt-hint tests above: a
// real apiserver cannot be made to lose this race deterministically.
func TestInstall_ControllerWriteBetweenCreateAndApply_RetriesOnCurrentVersion(t *testing.T) {
const ns = "test-ns"

patches := 0
c := fake.NewClientBuilder().
WithScheme(runtime.NewScheme()).
WithInterceptorFuncs(interceptor.Funcs{
Patch: func(ctx context.Context, cl client.WithWatch, obj client.Object, patch client.Patch, opts ...client.PatchOption) error {
patches++
if patches == 1 {
// The controller's write landed first; the conditional
// apply loses the race exactly once.
return agentClassConflict(obj.GetName())
}
return cl.Patch(ctx, obj, patch, opts...)
},
}).
Build()

result, err := install.Install(context.Background(), c, bundleWithSoleAgentClass("fresh-agent"), oap.Answers{}, nil, install.InstallOpts{
Namespace: ns,
})

require.NoError(t, err, "a lost race against the object's own controller must be retried, not surfaced")
require.NotNil(t, result)
assert.GreaterOrEqual(t, patches, 2, "the apply must have been attempted again after the conflict")
}

// TestInstall_ObjectReplacedDuringApply_FailsClosed is the safety property the
// retry must not relax: a conflict whose re-read comes back with a DIFFERENT
// uid means the object this run created was deleted and replaced mid-install.
// That is the foreign-seizure case the conditional apply exists to refuse —
// retrying on the replacement's resourceVersion would silently overwrite an
// object this run does not own.
func TestInstall_ObjectReplacedDuringApply_FailsClosed(t *testing.T) {
const ns = "test-ns"

patches := 0
c := fake.NewClientBuilder().
WithScheme(runtime.NewScheme()).
WithInterceptorFuncs(interceptor.Funcs{
Patch: func(ctx context.Context, cl client.WithWatch, obj client.Object, patch client.Patch, opts ...client.PatchOption) error {
patches++
return agentClassConflict(obj.GetName())
},
Get: func(ctx context.Context, cl client.WithWatch, key types.NamespacedName, obj client.Object, opts ...client.GetOption) error {
if err := cl.Get(ctx, key, obj, opts...); err != nil {
return err
}
if u, ok := obj.(*unstructured.Unstructured); ok && u.GetKind() == "AgentClass" {
// The object the conflict re-read finds is a same-name
// replacement, not the one this run created.
u.SetUID("replacement-uid")
}
return nil
},
}).
Build()

_, err := install.Install(context.Background(), c, bundleWithSoleAgentClass("fresh-agent"), oap.Answers{}, nil, install.InstallOpts{
Namespace: ns,
})

require.Error(t, err, "a replaced object must abort Install")
assert.Contains(t, err.Error(), "AgentClass/fresh-agent", "the error must name the object")
assert.Contains(t, strings.ToLower(err.Error()), "replaced", "the error must say the object was replaced, not report a bare conflict")
assert.Equal(t, 1, patches, "a replacement is not a retriable condition; no second apply may be attempted")
}

// TestInstall_PersistentConflict_GivesUpWithConflictError bounds the retry: a
// conflict that never resolves (every attempt loses) must eventually surface
// as the apply failure it is, not loop forever.
func TestInstall_PersistentConflict_GivesUpWithConflictError(t *testing.T) {
const ns = "test-ns"

patches := 0
c := fake.NewClientBuilder().
WithScheme(runtime.NewScheme()).
WithInterceptorFuncs(interceptor.Funcs{
Patch: func(ctx context.Context, cl client.WithWatch, obj client.Object, patch client.Patch, opts ...client.PatchOption) error {
patches++
return agentClassConflict(obj.GetName())
},
}).
Build()

_, err := install.Install(context.Background(), c, bundleWithSoleAgentClass("fresh-agent"), oap.Answers{}, nil, install.InstallOpts{
Namespace: ns,
})

require.Error(t, err, "a conflict that never resolves must fail the Install")
assert.Contains(t, err.Error(), "AgentClass/fresh-agent", "the error must name the object")
assert.Greater(t, patches, 1, "the conflict must have been retried before giving up")
assert.Less(t, patches, 10, "the retry must be bounded")
}
Loading