From 1c0d9be589164f6795d3f33d16d32a08feed17ef Mon Sep 17 00:00:00 2001 From: Erik Miller Date: Fri, 24 Jul 2026 12:55:06 -0700 Subject: [PATCH] fix(controller): retry async callback status update on conflict The async operation callback reads the managed resource, sets the conditions describing the result of the operation and writes them with a single, un-retried status update. Because the callback runs in its own goroutine, that write races the managed reconciler's own writes to the same object and regularly loses the optimistic concurrency check with "the object has been modified; please apply your changes to the latest version and try again". When that happens the result of the async operation is silently dropped: the provider error never reaches the LastAsyncOperation and Synced conditions, so the only durable record of why an async create, update or delete failed is lost and the resource just reports a stale status. Wrap the read and the status update in RetryOnConflict so that a conflict is retried against a freshly read object, re-applying the conditions on each attempt. The reconcile request and the returned errors are unchanged, so a status update that keeps conflicting is still reported. Signed-off-by: Erik Miller --- pkg/controller/api.go | 63 ++++++++++++------- pkg/controller/api_test.go | 126 +++++++++++++++++++++++++++++++++++++ 2 files changed, 165 insertions(+), 24 deletions(-) diff --git a/pkg/controller/api.go b/pkg/controller/api.go index 6113e1a1..b0db8d8d 100644 --- a/pkg/controller/api.go +++ b/pkg/controller/api.go @@ -13,6 +13,7 @@ import ( v1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/util/retry" "sigs.k8s.io/controller-runtime/pkg/client" ctrl "sigs.k8s.io/controller-runtime/pkg/manager" @@ -124,33 +125,47 @@ type APICallbacks struct { func (ac *APICallbacks) callbackFn(nn types.NamespacedName, op asyncOperation, requestReconcile bool) terraform.CallbackFn { //nolint:gocyclo // for better readability return func(err error, ctx context.Context) error { tr := ac.newTerraformed() - if kErr := ac.kube.Get(ctx, nn, tr); kErr != nil { - return errors.Wrapf(kErr, errGetFmt, tr.GetObjectKind().GroupVersionKind().String(), nn, op) - } - // For the no-fork architecture, we will need to be able to report - // reconciliation errors. The proper place is the `Synced` - // status condition but we need changes in the managed reconciler - // to do so. So we keep the `LastAsyncOperation` condition. - // TODO: move this to the `Synced` condition. - tr.SetConditions(resource.LastAsyncOperationCondition(err)) - if err != nil { - wrapMsg := "" - switch op { - case opCreate: - wrapMsg = errXPReconcileCreate - case opUpdate: - wrapMsg = errXPReconcileUpdate - case opDestroy: - wrapMsg = errXPReconcileDelete + setConditions := func() { + // For the no-fork architecture, we will need to be able to report + // reconciliation errors. The proper place is the `Synced` + // status condition but we need changes in the managed reconciler + // to do so. So we keep the `LastAsyncOperation` condition. + // TODO: move this to the `Synced` condition. + tr.SetConditions(resource.LastAsyncOperationCondition(err)) + if err != nil { + wrapMsg := "" + switch op { + case opCreate: + wrapMsg = errXPReconcileCreate + case opUpdate: + wrapMsg = errXPReconcileUpdate + case opDestroy: + wrapMsg = errXPReconcileDelete + } + tr.SetConditions(xpv2.ReconcileError(errors.Wrap(err, wrapMsg))) + } else { + tr.SetConditions(xpv2.ReconcileSuccess()) + } + if ac.enableStatusUpdates { + tr.SetConditions(resource.AsyncOperationFinishedCondition()) } - tr.SetConditions(xpv2.ReconcileError(errors.Wrap(err, wrapMsg))) - } else { - tr.SetConditions(xpv2.ReconcileSuccess()) } - if ac.enableStatusUpdates { - tr.SetConditions(resource.AsyncOperationFinishedCondition()) + // The managed reconciler concurrently updates the status of the same + // resource, so this status update can lose an optimistic concurrency + // race. Re-read the resource and set the conditions again on conflict, + // otherwise the result of the async operation is lost. + var kErr error + sErr := retry.RetryOnConflict(retry.DefaultRetry, func() error { + if kErr = ac.kube.Get(ctx, nn, tr); kErr != nil { + return kErr + } + setConditions() + return ac.kube.Status().Update(ctx, tr) + }) + if kErr != nil { + return errors.Wrapf(kErr, errGetFmt, tr.GetObjectKind().GroupVersionKind().String(), nn, op) } - uErr := errors.Wrapf(ac.kube.Status().Update(ctx, tr), errUpdateStatusFmt, tr.GetObjectKind().GroupVersionKind().String(), nn, op) + uErr := errors.Wrapf(sErr, errUpdateStatusFmt, tr.GetObjectKind().GroupVersionKind().String(), nn, op) if ac.eventHandler != nil && requestReconcile { rateLimiter := handler.NoRateLimiter switch { diff --git a/pkg/controller/api_test.go b/pkg/controller/api_test.go index f2c90eb3..92c99124 100644 --- a/pkg/controller/api_test.go +++ b/pkg/controller/api_test.go @@ -12,8 +12,11 @@ import ( xpresource "github.com/crossplane/crossplane-runtime/v2/pkg/resource" xpfake "github.com/crossplane/crossplane-runtime/v2/pkg/resource/fake" "github.com/crossplane/crossplane-runtime/v2/pkg/test" + xpv2 "github.com/crossplane/crossplane/apis/v2/core/v2" "github.com/google/go-cmp/cmp" "github.com/pkg/errors" + kerrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/types" "sigs.k8s.io/controller-runtime/pkg/client" ctrl "sigs.k8s.io/controller-runtime/pkg/manager" @@ -569,6 +572,129 @@ func newUnprimedEventHandler() *handler.EventHandler { return handler.NewEventHandler(handler.WithLogger(logging.NewNopLogger())) } +func TestAPICallbacksStatusUpdateConflict(t *testing.T) { + errConflict := kerrors.NewConflict(schema.GroupResource{Resource: "terraformeds"}, "name", errors.New("the object has been modified")) + asyncErr := tjerrors.NewApplyFailed(nil) + type args struct { + // statusUpdate returns the error of the calls'th status update. + statusUpdate func(calls int) error + withEventHandler bool + } + type want struct { + err error + // minGetCalls is the minimum number of times the resource must have + // been read. It is not an exact count because the number of retries + // depends on retry.DefaultRetry. + minGetCalls int + // maxGetCalls, when non-zero, is the maximum number of times the + // resource may have been read. + maxGetCalls int + } + cases := map[string]struct { + reason string + args + want + }{ + "ConflictThenSuccess": { + reason: "It should re-read the resource and set the conditions again if the status update conflicts", + args: args{ + statusUpdate: func(calls int) error { + if calls == 1 { + return errConflict + } + return nil + }, + }, + want: want{ + minGetCalls: 2, + }, + }, + "PersistentConflict": { + reason: "It should return the wrapped status update error if the conflict cannot be resolved", + args: args{ + statusUpdate: func(_ int) error { + return errConflict + }, + }, + want: want{ + err: errors.Wrapf(errConflict, errUpdateStatusFmt, "", ", Kind=//name", opUpdate), + minGetCalls: 2, + }, + }, + "PersistentConflictWithEventHandler": { + reason: "It should still request a reconciliation if the conflict cannot be resolved; with an uninitialized queue this surfaces as an errReconcileRequestFmt error.", + args: args{ + statusUpdate: func(_ int) error { + return errConflict + }, + withEventHandler: true, + }, + want: want{ + err: errors.Errorf(errReconcileRequestFmt, "", ", Kind=//name", opUpdate), + minGetCalls: 2, + }, + }, + "NonConflictError": { + reason: "It should not retry a status update error other than a conflict", + args: args{ + statusUpdate: func(_ int) error { + return errBoom + }, + }, + want: want{ + err: errors.Wrapf(errBoom, errUpdateStatusFmt, "", ", Kind=//name", opUpdate), + minGetCalls: 1, + maxGetCalls: 1, + }, + }, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + getCalls, updateCalls := 0, 0 + mgr := &xpfake.Manager{ + Client: &test.MockClient{ + MockGet: func(_ context.Context, _ client.ObjectKey, obj client.Object) error { + getCalls++ + // The API client decodes into a zeroed target, so the + // conditions of a previous attempt are not observed + // here. + obj.(*fake.Terraformed).ConditionedStatus = xpv2.ConditionedStatus{} + return nil + }, + MockStatusUpdate: func(_ context.Context, obj client.Object, _ ...client.SubResourceUpdateOption) error { + updateCalls++ + err := tc.args.statusUpdate(updateCalls) + if err != nil { + return err + } + got := obj.(resource.Terraformed).GetCondition(resource.TypeLastAsyncOperation) + if diff := cmp.Diff(resource.LastAsyncOperationCondition(asyncErr), got); diff != "" { + t.Errorf("\nUpdate(...): -want condition, +got condition:\n%s", diff) + } + return nil + }, + }, + Scheme: xpfake.SchemeWith(&fake.Terraformed{}), + } + var opts []APICallbacksOption + if tc.args.withEventHandler { + opts = append(opts, WithEventHandler(newUnprimedEventHandler())) + } + e := NewAPICallbacks(mgr, xpresource.ManagedKind(xpfake.GVK(&fake.Terraformed{})), opts...) + err := e.Update(types.NamespacedName{Name: "name"}, true)(asyncErr, context.TODO()) + if diff := cmp.Diff(tc.want.err, err, test.EquateErrors()); diff != "" { + t.Errorf("\n%s\nUpdate(...): -want error, +got error:\n%s", tc.reason, diff) + } + if getCalls < tc.want.minGetCalls { + t.Errorf("\n%s\nUpdate(...): expected the resource to be read at least %d times, got %d", tc.reason, tc.want.minGetCalls, getCalls) + } + if tc.want.maxGetCalls > 0 && getCalls > tc.want.maxGetCalls { + t.Errorf("\n%s\nUpdate(...): expected the resource to be read at most %d times, got %d", tc.reason, tc.want.maxGetCalls, getCalls) + } + }) + } +} + func TestAPICallbacksCreateRequestReconcile(t *testing.T) { type args struct { err error