diff --git a/internal/controllers/watch/kind.go b/internal/controllers/watch/kind.go index 3b385543..522ddd89 100644 --- a/internal/controllers/watch/kind.go +++ b/internal/controllers/watch/kind.go @@ -5,11 +5,11 @@ import ( "fmt" "math/rand" "path" - "reflect" "slices" "time" apiv1 "github.com/Azure/eno/api/v1" + "github.com/Azure/eno/internal/flowcontrol" "github.com/Azure/eno/internal/manager" "github.com/go-logr/logr" "golang.org/x/time/rate" @@ -17,7 +17,6 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/types" - "k8s.io/client-go/util/retry" "k8s.io/client-go/util/workqueue" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" @@ -36,6 +35,7 @@ func SetKindWatchRateLimit(rps int) { type KindWatchController struct { client client.Client + buffer *flowcontrol.CompositionInputRevisionWriteBuffer gvk schema.GroupVersionKind cancel context.CancelFunc } @@ -45,6 +45,7 @@ func NewKindWatchController(ctx context.Context, parent *WatchController, resour k := &KindWatchController{ client: parent.mgr.GetClient(), + buffer: parent.buffer, gvk: schema.GroupVersionKind{ Group: resource.Group, Version: resource.Version, @@ -219,10 +220,7 @@ func (k *KindWatchController) Reconcile(ctx context.Context, req ctrl.Request) ( logger.Error(err, "failed to list compositions for synthesizer", "synthesizerName", synth.Name) return ctrl.Result{}, fmt.Errorf("listing compositions: %w", err) } - modified, err := k.updateCompositions(ctx, logger, &synth, meta, list, isDeleted) - if modified || err != nil { - return ctrl.Result{}, err - } + k.updateCompositions(logger, &synth, meta, list, isDeleted) } list := &apiv1.CompositionList{} @@ -233,61 +231,42 @@ func (k *KindWatchController) Reconcile(ctx context.Context, req ctrl.Request) ( logger.Error(err, "failed to list compositions by binding", "synthesizerName", synth.Name) return ctrl.Result{}, fmt.Errorf("listing compositions: %w", err) } - modified, err := k.updateCompositions(ctx, logger, &synth, meta, list, isDeleted) - if modified || err != nil { - return ctrl.Result{}, err - } + k.updateCompositions(logger, &synth, meta, list, isDeleted) } logger.Info("finished reconciling watched resource") return ctrl.Result{}, nil } -func (k *KindWatchController) updateCompositions(ctx context.Context, logger logr.Logger, synth *apiv1.Synthesizer, meta *metav1.PartialObjectMetadata, list *apiv1.CompositionList, isDeleted bool) (bool, error) { - for _, comp := range list.Items { +// updateCompositions enqueues an input-revision change for every composition bound to the +// watched resource. Writes are coalesced per composition by the write buffer, so multiple +// inputs changing for the same composition in a short window collapse into a single +// Status().Update instead of one write per input. +func (k *KindWatchController) updateCompositions(logger logr.Logger, synth *apiv1.Synthesizer, meta *metav1.PartialObjectMetadata, list *apiv1.CompositionList, isDeleted bool) { + for i := range list.Items { + comp := list.Items[i] key := findRefKey(&comp, synth, meta) if key == "" { continue } compKey := client.ObjectKeyFromObject(&comp) - var modified bool - - err := retry.RetryOnConflict(retry.DefaultBackoff, func() error { - // Re-fetch composition to get the latest version - if err := k.client.Get(ctx, compKey, &comp); err != nil { - return err - } - - if isDeleted { - // Only remove InputRevisions for optional refs - // Required refs should trigger MissingInputs status instead - if isOptionalRef(synth, key) { - modified = removeInputRevision(&comp, key) - } - // For required refs, do nothing - the composition controller will handle MissingInputs - } else { - revs := apiv1.NewInputRevisions(meta, key) - modified = setInputRevisions(&comp, revs) - } - if !modified { - return nil + if isDeleted { + // Only remove InputRevisions for optional refs. + // Required refs should trigger MissingInputs status instead - the + // composition controller handles that case. + if isOptionalRef(synth, key) { + k.buffer.RemoveInputRevisionAsync(compKey, key) + logger.V(1).Info("buffered input revision removal", "compositionName", comp.Name, "compositionNamespace", comp.Namespace, "ref", key) } - - return k.client.Status().Update(ctx, &comp) - }) - if err != nil { - return false, fmt.Errorf("updating input revisions: %w", err) + continue } - if modified { - logger.V(1).Info("noticed input resource change", "compositionName", comp.Name, "compositionNamespace", comp.Namespace, "ref", key) - return true, nil // wait for requeue - } + revs := apiv1.NewInputRevisions(meta, key) + k.buffer.PatchInputRevisionAsync(compKey, revs) + logger.V(1).Info("buffered input resource change", "compositionName", comp.Name, "compositionNamespace", comp.Namespace, "ref", key) } - - return false, nil } func findRefKey(comp *apiv1.Composition, synth *apiv1.Synthesizer, meta *metav1.PartialObjectMetadata) string { @@ -320,28 +299,3 @@ func isOptionalRef(synth *apiv1.Synthesizer, key string) bool { } return false } - -func removeInputRevision(comp *apiv1.Composition, key string) bool { - for i, ir := range comp.Status.InputRevisions { - if ir.Key == key { - comp.Status.InputRevisions = append(comp.Status.InputRevisions[:i], comp.Status.InputRevisions[i+1:]...) - return true - } - } - return false -} - -func setInputRevisions(comp *apiv1.Composition, revs *apiv1.InputRevisions) bool { - for i, ir := range comp.Status.InputRevisions { - if ir.Key != revs.Key { - continue - } - if reflect.DeepEqual(ir, *revs) { - return false - } - comp.Status.InputRevisions[i] = *revs - return true - } - comp.Status.InputRevisions = append(comp.Status.InputRevisions, *revs) - return true -} diff --git a/internal/controllers/watch/kind_test.go b/internal/controllers/watch/kind_test.go index 093b7f60..5500bea9 100644 --- a/internal/controllers/watch/kind_test.go +++ b/internal/controllers/watch/kind_test.go @@ -6,122 +6,9 @@ import ( apiv1 "github.com/Azure/eno/api/v1" "github.com/stretchr/testify/assert" "k8s.io/apimachinery/pkg/types" - "k8s.io/utils/ptr" "sigs.k8s.io/controller-runtime/pkg/reconcile" ) -func TestSetInputRevisions(t *testing.T) { - tests := []struct { - name string - comp *apiv1.Composition - revs *apiv1.InputRevisions - expected bool - finalRevs []apiv1.InputRevisions - }{ - { - name: "add new revision when key is not found", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - }, - }, - }, - revs: &apiv1.InputRevisions{ - Key: "rev2", - Revision: ptr.To(2), - }, - expected: true, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - {Key: "rev2", Revision: ptr.To(2)}, - }, - }, - { - name: "update existing revision", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - }, - }, - }, - revs: &apiv1.InputRevisions{ - Key: "rev1", - Revision: ptr.To(2), - }, - expected: true, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(2)}, - }, - }, - { - name: "no update if revision is identical", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - }, - }, - }, - revs: &apiv1.InputRevisions{ - Key: "rev1", - Revision: ptr.To(1), - }, - expected: false, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - }, - }, - { - name: "no update if revision is identical and synth generation is set", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))}, - }, - }, - }, - revs: &apiv1.InputRevisions{ - Key: "rev1", - Revision: ptr.To(1), - SynthesizerGeneration: ptr.To(int64(3)), - }, - expected: false, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))}, - }, - }, - { - name: "update if revision is identical but synth generation is not", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))}, - }, - }, - }, - revs: &apiv1.InputRevisions{ - Key: "rev1", - Revision: ptr.To(1), - SynthesizerGeneration: ptr.To(int64(5)), - }, - expected: true, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(5))}, - }, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - result := setInputRevisions(tt.comp, tt.revs) - assert.Equal(t, tt.expected, result) - assert.Equal(t, tt.finalRevs, tt.comp.Status.InputRevisions) - }) - } -} - func TestBuildRequests(t *testing.T) { tests := []struct { name string @@ -266,87 +153,6 @@ func TestBuildRequests(t *testing.T) { } } -func TestRemoveInputRevision(t *testing.T) { - tests := []struct { - name string - comp *apiv1.Composition - key string - expected bool - finalRevs []apiv1.InputRevisions - }{ - { - name: "remove existing revision", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - {Key: "rev2", Revision: ptr.To(2)}, - }, - }, - }, - key: "rev1", - expected: true, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev2", Revision: ptr.To(2)}, - }, - }, - { - name: "remove last revision", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - }, - }, - }, - key: "rev1", - expected: true, - finalRevs: []apiv1.InputRevisions{}, - }, - { - name: "no removal if key not found", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - }, - }, - }, - key: "rev2", - expected: false, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - }, - }, - { - name: "remove from middle of list", - comp: &apiv1.Composition{ - Status: apiv1.CompositionStatus{ - InputRevisions: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - {Key: "rev2", Revision: ptr.To(2)}, - {Key: "rev3", Revision: ptr.To(3)}, - }, - }, - }, - key: "rev2", - expected: true, - finalRevs: []apiv1.InputRevisions{ - {Key: "rev1", Revision: ptr.To(1)}, - {Key: "rev3", Revision: ptr.To(3)}, - }, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - result := removeInputRevision(tt.comp, tt.key) - assert.Equal(t, tt.expected, result) - assert.Equal(t, tt.finalRevs, tt.comp.Status.InputRevisions) - }) - } -} - func TestIsOptionalRef(t *testing.T) { tests := []struct { name string diff --git a/internal/controllers/watch/watch.go b/internal/controllers/watch/watch.go index 3144af41..a7e11948 100644 --- a/internal/controllers/watch/watch.go +++ b/internal/controllers/watch/watch.go @@ -9,23 +9,31 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" apiv1 "github.com/Azure/eno/api/v1" + "github.com/Azure/eno/internal/flowcontrol" "github.com/Azure/eno/internal/manager" ) type WatchController struct { mgr ctrl.Manager client client.Client + buffer *flowcontrol.CompositionInputRevisionWriteBuffer refControllers map[apiv1.ResourceRef]*KindWatchController } func NewController(mgr ctrl.Manager) error { - err := ctrl.NewControllerManagedBy(mgr). + buffer, err := flowcontrol.NewCompositionInputRevisionWriteBufferForManager(mgr) + if err != nil { + return err + } + + err = ctrl.NewControllerManagedBy(mgr). Named("watchControllerController"). Watches(&apiv1.Synthesizer{}, manager.SingleEventHandler()). WithLogConstructor(manager.NewLogConstructor(mgr, "watchController")). Complete(&WatchController{ mgr: mgr, client: mgr.GetClient(), + buffer: buffer, refControllers: map[apiv1.ResourceRef]*KindWatchController{}, }) if err != nil { diff --git a/internal/flowcontrol/inputrevbuffer.go b/internal/flowcontrol/inputrevbuffer.go new file mode 100644 index 00000000..e5063e35 --- /dev/null +++ b/internal/flowcontrol/inputrevbuffer.go @@ -0,0 +1,272 @@ +package flowcontrol + +import ( + "context" + "reflect" + "sync" + "time" + + "github.com/go-logr/logr" + k8serrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/util/workqueue" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + + apiv1 "github.com/Azure/eno/api/v1" +) + +// inputRevisionOp is a single pending mutation to one input-revision key of a composition. +// When remove is true the key should be dropped; otherwise revs is written (last-write-wins). +type inputRevisionOp struct { + revs *apiv1.InputRevisions + remove bool +} + +// CompositionInputRevisionWriteBuffer reduces load on etcd/apiserver by collecting +// input-revision updates for each Composition over a short window and applying them +// in a single scoped status patch per Composition. +// +// Multiple input kinds are watched by independent KindWatchControllers, so a burst of +// input changes affecting the same Composition would otherwise produce one status write +// per change. Coalescing them per Composition (last-write-wins per input key) collapses +// those into a single mutating request while preserving the latest revision for every key. +type CompositionInputRevisionWriteBuffer struct { + client client.Client + + // queue items are per-composition. + // the state map collects multiple per-key updates per composition to be dispatched together. + mut sync.Mutex + state map[types.NamespacedName]map[string]inputRevisionOp + insertionTime map[types.NamespacedName]time.Time + // queued tracks whether a composition currently has a queue item awaiting + // processing. Coupling it with state under mut guarantees that any update + // enqueued while a flush is in flight re-queues the composition, so an update + // is never stranded in state without a pending queue item. + queued map[types.NamespacedName]bool + queue workqueue.TypedRateLimitingInterface[types.NamespacedName] +} + +func NewCompositionInputRevisionWriteBufferForManager(mgr ctrl.Manager) (*CompositionInputRevisionWriteBuffer, error) { + r := NewCompositionInputRevisionWriteBuffer(mgr.GetClient()) + if err := mgr.Add(r); err != nil { + return nil, err + } + return r, nil +} + +func NewCompositionInputRevisionWriteBuffer(cli client.Client) *CompositionInputRevisionWriteBuffer { + q := workqueue.NewTypedRateLimitingQueueWithConfig( + workqueue.NewTypedItemExponentialFailureRateLimiter[types.NamespacedName](time.Millisecond*100, 8*time.Second), + workqueue.TypedRateLimitingQueueConfig[types.NamespacedName]{ + Name: "compositionInputRevisionWriteBuffer", + }) + setInputRevisionBufferLenSource(q.Len) + return &CompositionInputRevisionWriteBuffer{ + client: cli, + state: make(map[types.NamespacedName]map[string]inputRevisionOp), + insertionTime: make(map[types.NamespacedName]time.Time), + queued: make(map[types.NamespacedName]bool), + queue: q, + } +} + +// PatchInputRevisionAsync enqueues an input-revision update for the given composition. +// The update is coalesced last-write-wins per input key and eventually flushed, or dropped +// only if the composition is deleted. +func (w *CompositionInputRevisionWriteBuffer) PatchInputRevisionAsync(comp types.NamespacedName, revs *apiv1.InputRevisions) { + w.enqueue(comp, revs.Key, inputRevisionOp{revs: revs}) +} + +// RemoveInputRevisionAsync enqueues the removal of a single input-revision key for the given +// composition. Coalesced last-write-wins with any pending update for the same key. +func (w *CompositionInputRevisionWriteBuffer) RemoveInputRevisionAsync(comp types.NamespacedName, key string) { + w.enqueue(comp, key, inputRevisionOp{remove: true}) +} + +func (w *CompositionInputRevisionWriteBuffer) enqueue(comp types.NamespacedName, key string, op inputRevisionOp) { + w.mut.Lock() + defer w.mut.Unlock() + + ops := w.state[comp] + if ops == nil { + ops = make(map[string]inputRevisionOp) + w.state[comp] = ops + } + ops[key] = op // last write wins + + if _, found := w.insertionTime[comp]; !found { + w.insertionTime[comp] = time.Now() + } + if !w.queued[comp] { + w.queued[comp] = true + w.queue.AddRateLimited(comp) + } +} + +func (w *CompositionInputRevisionWriteBuffer) Start(ctx context.Context) error { + go func() { + <-ctx.Done() + w.queue.ShutDown() + }() + for w.processQueueItem(ctx) { + } + return nil +} + +func (w *CompositionInputRevisionWriteBuffer) processQueueItem(ctx context.Context) bool { + comp, shutdown := w.queue.Get() + if shutdown { + return false + } + defer w.queue.Done(comp) + + logger := logr.FromContextOrDiscard(ctx).WithValues("compositionName", comp.Name, "compositionNamespace", comp.Namespace, "controller", "compositionInputRevisionWriteBuffer") + ctx = logr.NewContext(ctx, logger) + + w.mut.Lock() + // Mark not-queued up front (under the same lock that guards state) so that any + // update enqueued from here on re-queues the composition instead of being stranded. + delete(w.queued, comp) + insertionTime := w.insertionTime[comp] + ops := w.state[comp] + delete(w.state, comp) + w.mut.Unlock() + + if len(ops) == 0 { + w.settle(comp) + return true + } + + if w.updateComposition(ctx, insertionTime, comp, ops) { + w.settle(comp) + return true + } + + // Put the updates back in the buffer to retry, without clobbering newer updates + // that arrived while we were flushing, and ensure the composition is re-queued. + w.mut.Lock() + pending := w.state[comp] + if pending == nil { + w.state[comp] = ops + } else { + for key, op := range ops { + if _, ok := pending[key]; !ok { + pending[key] = op + } + } + } + if !w.queued[comp] { + w.queued[comp] = true + w.queue.AddRateLimited(comp) + } + w.mut.Unlock() + return true +} + +// settle is called after a flush succeeds (including a no-op flush, where there was +// nothing to do). If a producer enqueued new ops for comp while the flush was in flight, +// enqueue has already set w.queued[comp] and re-added it to the queue - in that case we +// leave the rate limiter alone so a composition under sustained churn keeps backing off +// exponentially instead of resetting to the fast path on every flush, and we keep +// insertionTime so the next flush reports an accurate latency instead of the zero time. +// Only once a flush finds nothing pending do we Forget (reset to the fast path) and clear +// insertionTime, since the composition is genuinely idle. +func (w *CompositionInputRevisionWriteBuffer) settle(comp types.NamespacedName) { + w.mut.Lock() + defer w.mut.Unlock() + if w.queued[comp] { + return + } + w.queue.Forget(comp) + delete(w.insertionTime, comp) +} + +// updateComposition applies all pending per-key ops to the composition via a merge patch +// carrying an optimistic lock. Returns true when the work is done (including no-op and +// composition-deleted), false when it should be retried. +func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Context, insertionTime time.Time, comp types.NamespacedName, ops map[string]inputRevisionOp) (success bool) { + logger := logr.FromContextOrDiscard(ctx) + + obj := &apiv1.Composition{} + err := w.client.Get(ctx, comp, obj) + if k8serrors.IsNotFound(err) { + logger.V(1).Info("composition deleted - dropping buffered input revision updates") + return true + } + if err != nil { + logger.Error(err, "unable to get composition") + inputRevisionBufferFlushErrors.WithLabelValues("get").Inc() + return false + } + + original := obj.DeepCopy() + + modified := false + for key, op := range ops { + if op.remove { + if removeInputRevision(obj, key) { + modified = true + } + continue + } + if setInputRevisions(obj, op.revs) { + modified = true + } + } + if !modified { + return true + } + + // A merge patch with an optimistic lock (resourceVersion) merges into whatever is + // currently stored, handling an absent status and an emptied list correctly, and still + // conflicts (and retries) against a concurrent writer, whether to inputRevisions itself + // (e.g. pruning.go) or, more conservatively, to an unrelated status field. + err = w.client.Status().Patch(ctx, obj, client.MergeFromWithOptions(original, client.MergeFromWithOptimisticLock{})) + if k8serrors.IsNotFound(err) { + logger.V(1).Info("composition deleted - dropping buffered input revision updates") + return true + } + if k8serrors.IsConflict(err) { + logger.V(1).Info("conflict updating composition input revisions - will retry") + return false + } + if err != nil { + logger.Error(err, "unable to update composition input revisions") + inputRevisionBufferFlushErrors.WithLabelValues("patch").Inc() + return false + } + + inputRevisionBufferFlushes.Inc() + logger.V(1).Info("flushed input revision updates to composition", "keyCount", len(ops), "latencyMs", time.Since(insertionTime).Abs().Milliseconds()) + return true +} + +// setInputRevisions sets or replaces the revision for revs.Key on the composition. +// Returns true when the composition's status was changed. +func setInputRevisions(comp *apiv1.Composition, revs *apiv1.InputRevisions) bool { + for i, ir := range comp.Status.InputRevisions { + if ir.Key != revs.Key { + continue + } + if reflect.DeepEqual(ir, *revs) { + return false + } + comp.Status.InputRevisions[i] = *revs + return true + } + comp.Status.InputRevisions = append(comp.Status.InputRevisions, *revs) + return true +} + +// removeInputRevision drops the revision for key from the composition. +// Returns true when the composition's status was changed. +func removeInputRevision(comp *apiv1.Composition, key string) bool { + for i, ir := range comp.Status.InputRevisions { + if ir.Key == key { + comp.Status.InputRevisions = append(comp.Status.InputRevisions[:i], comp.Status.InputRevisions[i+1:]...) + return true + } + } + return false +} diff --git a/internal/flowcontrol/inputrevbuffer_test.go b/internal/flowcontrol/inputrevbuffer_test.go new file mode 100644 index 00000000..1f571b92 --- /dev/null +++ b/internal/flowcontrol/inputrevbuffer_test.go @@ -0,0 +1,806 @@ +package flowcontrol + +import ( + "context" + "errors" + "fmt" + "math/rand/v2" + "strconv" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + k8serrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/utils/ptr" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" + + apiv1 "github.com/Azure/eno/api/v1" + "github.com/Azure/eno/internal/testutil" +) + +func newTestComposition(t *testing.T, cli client.Client, name string) *apiv1.Composition { + t.Helper() + comp := &apiv1.Composition{} + comp.Name = name + comp.Namespace = "default" + comp.Spec.Synthesizer.Name = "test-synth" + require.NoError(t, cli.Create(testutil.NewContext(t), comp)) + return comp +} + +func TestInputRevisionWriteBufferBasics(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + + w.processQueueItem(ctx) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + require.Len(t, comp.Status.InputRevisions, 1) + assert.Equal(t, "foo", comp.Status.InputRevisions[0].Key) + assert.Equal(t, "1", comp.Status.InputRevisions[0].ResourceVersion) + + // state fully flushed + assert.Len(t, w.state, 0) + assert.Equal(t, 0, w.queue.Len()) +} + +func TestInputRevisionWriteBufferCoalescesMultipleKeys(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + + // Three different input keys enqueued before any flush. + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "a", ResourceVersion: "1"}) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "b", ResourceVersion: "1"}) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "c", ResourceVersion: "1"}) + + // A single dequeue should write all three at once. + rvBefore := getResourceVersion(t, cli, nsn) + w.processQueueItem(ctx) + rvAfter := getResourceVersion(t, cli, nsn) + assert.NotEqual(t, rvBefore, rvAfter, "expected exactly one status update") + + require.NoError(t, cli.Get(ctx, nsn, comp)) + keys := inputRevisionKeys(comp) + assert.ElementsMatch(t, []string{"a", "b", "c"}, keys) +} + +func TestInputRevisionWriteBufferLastWriteWinsPerKey(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "2"}) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "3"}) + + w.processQueueItem(ctx) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + require.Len(t, comp.Status.InputRevisions, 1) + assert.Equal(t, "3", comp.Status.InputRevisions[0].ResourceVersion) +} + +func TestInputRevisionWriteBufferRemove(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + comp.Status.InputRevisions = []apiv1.InputRevisions{ + {Key: "keep", ResourceVersion: "1"}, + {Key: "drop", ResourceVersion: "1"}, + } + require.NoError(t, cli.Status().Update(ctx, comp)) + + w.RemoveInputRevisionAsync(nsn, "drop") + w.processQueueItem(ctx) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.ElementsMatch(t, []string{"keep"}, inputRevisionKeys(comp)) +} + +func TestInputRevisionWriteBufferSetThenRemoveSameKeyCoalesces(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + + // set then remove the same key before any flush: the remove wins. + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.RemoveInputRevisionAsync(nsn, "foo") + w.processQueueItem(ctx) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.Empty(t, comp.Status.InputRevisions) +} + +func TestInputRevisionWriteBufferNoOpDoesNotError(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + comp.Status.InputRevisions = []apiv1.InputRevisions{{Key: "foo", ResourceVersion: "1"}} + require.NoError(t, cli.Status().Update(ctx, comp)) + rvBefore := getResourceVersion(t, cli, nsn) + + // Enqueue an identical revision - should be a no-op, no status write. + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) + + assert.Equal(t, rvBefore, getResourceVersion(t, cli, nsn), "no-op update must not write status") + assert.Len(t, w.state, 0) +} + +func TestInputRevisionWriteBufferDeletedCompositionDropped(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + nsn := client.ObjectKey{Namespace: "default", Name: "does-not-exist"} + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) + + // Nothing to assert other than the buffer drained cleanly without hanging/erroring. + assert.Len(t, w.state, 0) + assert.Equal(t, 0, w.queue.Len()) +} + +func getResourceVersion(t *testing.T, cli client.Client, nsn client.ObjectKey) string { + t.Helper() + comp := &apiv1.Composition{} + require.NoError(t, cli.Get(testutil.NewContext(t), nsn, comp)) + return comp.ResourceVersion +} + +func inputRevisionKeys(comp *apiv1.Composition) []string { + keys := make([]string, 0, len(comp.Status.InputRevisions)) + for _, ir := range comp.Status.InputRevisions { + keys = append(keys, ir.Key) + } + return keys +} + +func TestSetInputRevisions(t *testing.T) { + tests := []struct { + name string + comp *apiv1.Composition + revs *apiv1.InputRevisions + expected bool + finalRevs []apiv1.InputRevisions + }{ + { + name: "add new revision when key is not found", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + }, + }, + }, + revs: &apiv1.InputRevisions{ + Key: "rev2", + Revision: ptr.To(2), + }, + expected: true, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + {Key: "rev2", Revision: ptr.To(2)}, + }, + }, + { + name: "update existing revision", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + }, + }, + }, + revs: &apiv1.InputRevisions{ + Key: "rev1", + Revision: ptr.To(2), + }, + expected: true, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(2)}, + }, + }, + { + name: "no update if revision is identical", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + }, + }, + }, + revs: &apiv1.InputRevisions{ + Key: "rev1", + Revision: ptr.To(1), + }, + expected: false, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + }, + }, + { + name: "no update if revision is identical and synth generation is set", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))}, + }, + }, + }, + revs: &apiv1.InputRevisions{ + Key: "rev1", + Revision: ptr.To(1), + SynthesizerGeneration: ptr.To(int64(3)), + }, + expected: false, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))}, + }, + }, + { + name: "update if revision is identical but synth generation is not", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))}, + }, + }, + }, + revs: &apiv1.InputRevisions{ + Key: "rev1", + Revision: ptr.To(1), + SynthesizerGeneration: ptr.To(int64(5)), + }, + expected: true, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(5))}, + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + result := setInputRevisions(tt.comp, tt.revs) + assert.Equal(t, tt.expected, result) + assert.Equal(t, tt.finalRevs, tt.comp.Status.InputRevisions) + }) + } +} + +func TestRemoveInputRevision(t *testing.T) { + tests := []struct { + name string + comp *apiv1.Composition + key string + expected bool + finalRevs []apiv1.InputRevisions + }{ + { + name: "remove existing revision", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + {Key: "rev2", Revision: ptr.To(2)}, + }, + }, + }, + key: "rev1", + expected: true, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev2", Revision: ptr.To(2)}, + }, + }, + { + name: "remove last revision", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + }, + }, + }, + key: "rev1", + expected: true, + finalRevs: []apiv1.InputRevisions{}, + }, + { + name: "no removal if key not found", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + }, + }, + }, + key: "rev2", + expected: false, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + }, + }, + { + name: "remove from middle of list", + comp: &apiv1.Composition{ + Status: apiv1.CompositionStatus{ + InputRevisions: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + {Key: "rev2", Revision: ptr.To(2)}, + {Key: "rev3", Revision: ptr.To(3)}, + }, + }, + }, + key: "rev2", + expected: true, + finalRevs: []apiv1.InputRevisions{ + {Key: "rev1", Revision: ptr.To(1)}, + {Key: "rev3", Revision: ptr.To(3)}, + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + result := removeInputRevision(tt.comp, tt.key) + assert.Equal(t, tt.expected, result) + assert.Equal(t, tt.finalRevs, tt.comp.Status.InputRevisions) + }) + } +} + +// TestInputRevisionWriteBufferMapsDoNotLeak asserts no per-composition bookkeeping survives once a composition is created, flushed, and deleted. +func TestInputRevisionWriteBufferMapsDoNotLeak(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) // create + flush + + require.NoError(t, cli.Delete(ctx, comp)) + w.RemoveInputRevisionAsync(nsn, "foo") + w.processQueueItem(ctx) // composition is gone, buffered update is dropped + + w.mut.Lock() + defer w.mut.Unlock() + assert.Empty(t, w.state, "state must not retain an entry once settled") + assert.Empty(t, w.insertionTime, "insertionTime must not retain an entry once settled") + assert.Empty(t, w.queued, "queued must not retain an entry once settled") +} + +// TestInputRevisionWriteBufferInsertionTimePreservedDuringRace asserts insertionTime survives a producer racing in a new update mid-flush, instead of resetting to the zero time. +func TestInputRevisionWriteBufferInsertionTimePreservedDuringRace(t *testing.T) { + ctx := testutil.NewContext(t) + var w *CompositionInputRevisionWriteBuffer + var nsn client.ObjectKey + raced := false + cli := testutil.NewClientWithInterceptors(t, &interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, c client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + err := c.SubResource(subResourceName).Patch(ctx, obj, patch, opts...) + if !raced { + raced = true + // Simulate a producer racing in a new update while this flush is still + // in flight, i.e. before the success path decides whether to settle. + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "bar", ResourceVersion: "2"}) + } + return err + }, + }) + w = NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn = client.ObjectKeyFromObject(comp) + + before := time.Now() + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) // flushes "foo"; races in "bar" mid-flush + + w.mut.Lock() + insertionTime := w.insertionTime[nsn] + w.mut.Unlock() + require.False(t, insertionTime.IsZero(), "insertionTime must not be cleared while an update is still pending") + assert.False(t, insertionTime.Before(before), "insertionTime should reflect the racing enqueue, not a stale or zero value") + + latency := time.Since(insertionTime) + assert.GreaterOrEqual(t, latency, time.Duration(0), "latency must be non-negative") + assert.Less(t, latency, 10*time.Second, "latency must reflect a real timestamp, not the zero time") + + // The racing update is still buffered and flushes cleanly on the next pass. + w.processQueueItem(ctx) + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.ElementsMatch(t, []string{"foo", "bar"}, inputRevisionKeys(comp)) + assert.Empty(t, w.state) + assert.Empty(t, w.insertionTime) +} + +// TestInputRevisionWriteBufferBackoffGrowsUnderSustainedLoad asserts the rate limiter backs off exponentially while a composition keeps churning. +func TestInputRevisionWriteBufferBackoffGrowsUnderSustainedLoad(t *testing.T) { + ctx := testutil.NewContext(t) + var w *CompositionInputRevisionWriteBuffer + var nsn client.ObjectKey + const churnCycles = 3 + patchCalls := 0 + cli := testutil.NewClientWithInterceptors(t, &interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, c client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + err := c.SubResource(subResourceName).Patch(ctx, obj, patch, opts...) + patchCalls++ + if patchCalls <= churnCycles { + // Simulate another input changing while this flush is still in flight. + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: strconv.Itoa(patchCalls + 1)}) + } + return err + }, + }) + w = NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn = client.ObjectKeyFromObject(comp) + + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + + prevRequeues := 0 + for i := 0; i < churnCycles; i++ { + w.processQueueItem(ctx) + requeues := w.queue.NumRequeues(nsn) + assert.Greater(t, requeues, prevRequeues, "rate limit must not be forgotten while the composition is still churning") + prevRequeues = requeues + } + + // One more flush finds nothing pending: the composition is idle, so the rate limiter resets. + w.processQueueItem(ctx) + assert.Equal(t, 0, w.queue.NumRequeues(nsn), "rate limit should reset once the composition goes idle") +} + +// TestInputRevisionWriteBufferBackoffStaysBaseForIsolatedComposition asserts a single, isolated update resets the rate limiter to the fast path. +func TestInputRevisionWriteBufferBackoffStaysBaseForIsolatedComposition(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) + + assert.Equal(t, 0, w.queue.NumRequeues(nsn), "an isolated update should reset the rate limit to the fast path") +} + +func TestInputRevisionWriteBufferRequeuesFailedFlushWithoutClobberingNewerUpdates(t *testing.T) { + ctx := testutil.NewContext(t) + var w *CompositionInputRevisionWriteBuffer + var nsn client.ObjectKey + patchCalls := 0 + cli := testutil.NewClientWithInterceptors(t, &interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, c client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + patchCalls++ + if patchCalls == 1 { + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "2"}) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "baz", ResourceVersion: "1"}) + return errors.New("transient patch failure") + } + return c.SubResource(subResourceName).Patch(ctx, obj, patch, opts...) + }, + }) + w = NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn = client.ObjectKeyFromObject(comp) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "bar", ResourceVersion: "1"}) + + w.processQueueItem(ctx) + + w.mut.Lock() + pending := w.state[nsn] + queued := w.queued[nsn] + w.mut.Unlock() + require.Len(t, pending, 3) + assert.Equal(t, "2", pending["foo"].revs.ResourceVersion, "the newer update must win") + assert.Equal(t, "1", pending["bar"].revs.ResourceVersion, "the failed update must be restored") + assert.Equal(t, "1", pending["baz"].revs.ResourceVersion, "the concurrent update must be preserved") + assert.True(t, queued, "the composition must remain queued after a failed flush") + assert.Greater(t, w.queue.NumRequeues(nsn), 0) + + w.processQueueItem(ctx) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.Equal(t, "2", revisionFor(comp, "foo")) + assert.Equal(t, "1", revisionFor(comp, "bar")) + assert.Equal(t, "1", revisionFor(comp, "baz")) + assert.Empty(t, w.state) + assert.False(t, w.queued[nsn]) + assert.Equal(t, 2, patchCalls) +} + +func TestInputRevisionWriteBufferPatchErrors(t *testing.T) { + resource := schema.GroupResource{Group: apiv1.SchemeGroupVersion.Group, Resource: "compositions"} + tests := []struct { + name string + patchErr error + wantRetry bool + }{ + { + name: "not found drops update", + patchErr: k8serrors.NewNotFound(resource, "test-comp-1"), + }, + { + name: "conflict retries update", + patchErr: k8serrors.NewConflict(resource, "test-comp-1", errors.New("conflict")), + wantRetry: true, + }, + { + name: "generic error retries update", + patchErr: errors.New("transient patch failure"), + wantRetry: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctx := testutil.NewContext(t) + patchCalls := 0 + cli := testutil.NewClientWithInterceptors(t, &interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, c client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + patchCalls++ + if patchCalls == 1 { + return tt.patchErr + } + return c.SubResource(subResourceName).Patch(ctx, obj, patch, opts...) + }, + }) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) + + if !tt.wantRetry { + assert.Empty(t, w.state) + assert.False(t, w.queued[nsn]) + assert.Equal(t, 0, w.queue.NumRequeues(nsn)) + assert.Equal(t, 1, patchCalls) + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.Empty(t, comp.Status.InputRevisions) + return + } + + require.Len(t, w.state[nsn], 1) + assert.True(t, w.queued[nsn]) + assert.Greater(t, w.queue.NumRequeues(nsn), 0) + + w.processQueueItem(ctx) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.Equal(t, "1", revisionFor(comp, "foo")) + assert.Empty(t, w.state) + assert.False(t, w.queued[nsn]) + assert.Equal(t, 2, patchCalls) + }) + } +} + +// TestInputRevisionWriteBufferRetriesAfterConflictWithUnrelatedStatusWrite asserts a concurrent write to an unrelated status field conflicts, retries, and is preserved rather than clobbered. +func TestInputRevisionWriteBufferRetriesAfterConflictWithUnrelatedStatusWrite(t *testing.T) { + ctx := testutil.NewContext(t) + raced := false + cli := testutil.NewClientWithInterceptors(t, &interceptor.Funcs{ + Get: func(ctx context.Context, c client.WithWatch, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + err := c.Get(ctx, key, obj, opts...) + if _, ok := obj.(*apiv1.Composition); ok && !raced && err == nil { + raced = true + // Simulate a concurrent writer (e.g. the synthesis controller) updating an + // unrelated status field between our Get and our Patch. + concurrent := &apiv1.Composition{} + require.NoError(t, c.Get(ctx, key, concurrent)) + concurrent.Status.CurrentSynthesis = &apiv1.Synthesis{UUID: "concurrent-write"} + require.NoError(t, c.Status().Update(ctx, concurrent)) + } + return err + }, + }) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) // our resourceVersion is now stale, conflicts, retries + w.processQueueItem(ctx) // retry against a fresh read: succeeds + + assert.Empty(t, w.state) + assert.Equal(t, 0, w.queue.NumRequeues(nsn)) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + require.Len(t, comp.Status.InputRevisions, 1) + assert.Equal(t, "foo", comp.Status.InputRevisions[0].Key) + require.NotNil(t, comp.Status.CurrentSynthesis) + + assert.Equal(t, "concurrent-write", comp.Status.CurrentSynthesis.UUID, "unrelated status field must be preserved") +} + +// TestInputRevisionWriteBufferDoesNotResurrectConcurrentlyPrunedKey asserts a concurrent prune (internal/controllers/watch/pruning.go) of a different key is not resurrected by a racing flush. +func TestInputRevisionWriteBufferDoesNotResurrectConcurrentlyPrunedKey(t *testing.T) { + ctx := testutil.NewContext(t) + raced := false + cli := testutil.NewClientWithInterceptors(t, &interceptor.Funcs{ + Get: func(ctx context.Context, c client.WithWatch, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + err := c.Get(ctx, key, obj, opts...) + if _, ok := obj.(*apiv1.Composition); ok && !raced && err == nil { + raced = true + // Simulate the pruning controller's own Get -> mutate -> Status().Update + // removing "baz", interleaved between our Get and our Patch. + pruned := &apiv1.Composition{} + require.NoError(t, c.Get(ctx, key, pruned)) + pruned.Status.InputRevisions = []apiv1.InputRevisions{{Key: "bar", ResourceVersion: "1"}} + require.NoError(t, c.Status().Update(ctx, pruned)) + } + return err + }, + }) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + comp.Status.InputRevisions = []apiv1.InputRevisions{ + {Key: "bar", ResourceVersion: "1"}, + {Key: "baz", ResourceVersion: "1"}, + } + require.NoError(t, cli.Status().Update(ctx, comp)) + + // The buffer flushes an update to an unrelated key based on its now-stale read, which + // still includes "baz". + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + w.processQueueItem(ctx) // first attempt: resourceVersion is stale, pruner's write conflicts, retries + w.processQueueItem(ctx) // retry against a fresh read: succeeds + + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.ElementsMatch(t, []string{"bar", "foo"}, inputRevisionKeys(comp), "the pruned key must not be resurrected") + assert.Empty(t, w.state) +} + +func revisionFor(comp *apiv1.Composition, key string) string { + for _, ir := range comp.Status.InputRevisions { + if ir.Key == key { + return ir.ResourceVersion + } + } + return "" +} + +// TestInputRevisionWriteBufferConcurrentWrites runs many goroutines against a live Start +// loop (run with -race); each (composition, key) pair has exactly one writer goroutine, so +// the final value it wrote is unambiguous. +func TestInputRevisionWriteBufferConcurrentWrites(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + go w.Start(ctx) + + const compCount, keysPerComp, iterations = 3, 2, 15 + + nsns := make([]client.ObjectKey, compCount) + for i := range nsns { + comp := newTestComposition(t, cli, fmt.Sprintf("test-comp-%d", i)) + nsns[i] = client.ObjectKeyFromObject(comp) + } + + var mu sync.Mutex + want := make(map[client.ObjectKey]map[string]string, compCount) // "" means removed + for _, nsn := range nsns { + want[nsn] = make(map[string]string, keysPerComp) + } + + var wg sync.WaitGroup + for c := 0; c < compCount; c++ { + for k := 0; k < keysPerComp; k++ { + nsn, key := nsns[c], fmt.Sprintf("key-%d", k) + wg.Add(1) + go func() { + defer wg.Done() + var last string + for i := 0; i < iterations; i++ { + if rand.IntN(4) == 0 { + w.RemoveInputRevisionAsync(nsn, key) + last = "" + } else { + last = strconv.Itoa(i) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: key, ResourceVersion: last}) + } + } + mu.Lock() + want[nsn][key] = last + mu.Unlock() + }() + } + } + wg.Wait() + + testutil.Eventually(t, func() bool { + w.mut.Lock() + defer w.mut.Unlock() + return len(w.state) == 0 && w.queue.Len() == 0 + }) + + for _, nsn := range nsns { + comp := &apiv1.Composition{} + require.NoError(t, cli.Get(ctx, nsn, comp)) + for key, rv := range want[nsn] { + assert.Equal(t, rv, revisionFor(comp, key), "key %s on %s", key, nsn.Name) + } + } +} + +// TestInputRevisionWriteBufferIntegrationStorageEdgeCases runs against a real apiserver +// (envtest), since a fresh composition with no status key yet, and an emptied list, are +// storage-layer behaviors the fake client can't reproduce. +func TestInputRevisionWriteBufferIntegrationStorageEdgeCases(t *testing.T) { + ctx := testutil.NewContext(t) + mgr := testutil.NewManager(t) + cli := mgr.GetClient() + w, err := NewCompositionInputRevisionWriteBufferForManager(mgr.Manager) + require.NoError(t, err) + mgr.Start(t) + + comp := &apiv1.Composition{} + comp.Name = "test-comp" + comp.Namespace = "default" + comp.Spec.Synthesizer.Name = "test-synth" + require.NoError(t, cli.Create(ctx, comp)) + nsn := client.ObjectKeyFromObject(comp) + + // (a) First-ever status write on a fresh composition: the apiserver has never + // initialized status for this object, so it may be entirely absent server-side. + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) + testutil.Eventually(t, func() bool { + if err := cli.Get(ctx, nsn, comp); err != nil { + return false + } + return len(comp.Status.InputRevisions) == 1 && comp.Status.InputRevisions[0].Key == "foo" + }) + + // (b) Remove the only input revision, then add a different one: the field must not + // get stuck as a stored [] that a later flush can never write past. + w.RemoveInputRevisionAsync(nsn, "foo") + testutil.Eventually(t, func() bool { + cli.Get(ctx, nsn, comp) + return len(comp.Status.InputRevisions) == 0 + }) + + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "bar", ResourceVersion: "1"}) + testutil.Eventually(t, func() bool { + cli.Get(ctx, nsn, comp) + return len(comp.Status.InputRevisions) == 1 && comp.Status.InputRevisions[0].Key == "bar" + }) +} diff --git a/internal/flowcontrol/metrics.go b/internal/flowcontrol/metrics.go index 7cb22de1..347b0765 100644 --- a/internal/flowcontrol/metrics.go +++ b/internal/flowcontrol/metrics.go @@ -45,10 +45,50 @@ var ( Help: "Errors encountered while flushing resource slice status updates, partitioned by operation", }, []string{"op"}, ) + + // Depth of the composition input-revision write buffer's internal queue. A + // persistently non-zero value means input-revision patches are accumulating + // faster than the buffer can flush them to apiserver, which directly delays + // downstream resynthesis. + inputRevisionBufferDepth = prometheus.NewGaugeFunc( + prometheus.GaugeOpts{ + Name: "eno_input_revision_write_buffer_depth", + Help: "Current depth of the composition input-revision status write buffer queue", + }, + func() float64 { + fn := inputRevisionBufferLenFn.Load() + if fn == nil { + return 0 + } + return float64((*fn)()) + }, + ) + // Closure that returns the current input-revision write buffer queue depth. + // Installed by NewCompositionInputRevisionWriteBuffer; nil before the buffer is + // constructed. atomic.Pointer makes the swap safe against the scrape goroutine. + inputRevisionBufferLenFn atomic.Pointer[func() int] + + // Errors hit while flushing input-revision patches, partitioned by op (get/patch). + // These silently retry today, so without this counter they're invisible. + inputRevisionBufferFlushErrors = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "eno_input_revision_write_buffer_flush_errors_total", + Help: "Errors encountered while flushing composition input revision updates, partitioned by operation", + }, []string{"op"}, + ) + + // Count of successful flushes of composition input-revision updates. + inputRevisionBufferFlushes = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "eno_input_revision_write_buffer_flush_total", + Help: "Count of successful flushes of composition input revision updates", + }, + ) ) func init() { - metrics.Registry.MustRegister(sliceStatusUpdates, writeBufferDepth, writeBufferStatusUpdateErrors) + metrics.Registry.MustRegister(sliceStatusUpdates, writeBufferDepth, writeBufferStatusUpdateErrors, + inputRevisionBufferDepth, inputRevisionBufferFlushErrors, inputRevisionBufferFlushes) } // setWriteBufferLenSource installs the queue.Len reader used by the depth gauge. @@ -56,3 +96,9 @@ func init() { func setWriteBufferLenSource(fn func() int) { writeBufferLenFn.Store(&fn) } + +// setInputRevisionBufferLenSource installs the queue.Len reader used by the input-revision +// buffer's depth gauge. Called once during NewCompositionInputRevisionWriteBuffer. +func setInputRevisionBufferLenSource(fn func() int) { + inputRevisionBufferLenFn.Store(&fn) +} diff --git a/internal/flowcontrol/metrics_test.go b/internal/flowcontrol/metrics_test.go index 14f981a8..8401f95f 100644 --- a/internal/flowcontrol/metrics_test.go +++ b/internal/flowcontrol/metrics_test.go @@ -50,3 +50,45 @@ func TestWriteBufferDepthGauge(t *testing.T) { setWriteBufferLenSource(func() int { return 12 }) assert.Equal(t, 12.0, prometheustestutil.ToFloat64(writeBufferDepth)) } + +func TestInputRevisionBufferFlushErrorsLabels(t *testing.T) { + inputRevisionBufferFlushErrors.Reset() + + inputRevisionBufferFlushErrors.WithLabelValues("get").Inc() + inputRevisionBufferFlushErrors.WithLabelValues("get").Inc() + inputRevisionBufferFlushErrors.WithLabelValues("patch").Inc() + + assert.Equal(t, 2.0, prometheustestutil.ToFloat64(inputRevisionBufferFlushErrors.WithLabelValues("get"))) + assert.Equal(t, 1.0, prometheustestutil.ToFloat64(inputRevisionBufferFlushErrors.WithLabelValues("patch"))) + // Untouched labels stay at zero. + assert.Equal(t, 0.0, prometheustestutil.ToFloat64(inputRevisionBufferFlushErrors.WithLabelValues("other"))) +} + +func TestInputRevisionBufferDepthGauge(t *testing.T) { + prev := inputRevisionBufferLenFn.Load() + t.Cleanup(func() { + if prev == nil { + inputRevisionBufferLenFn.Store(nil) + } else { + inputRevisionBufferLenFn.Store(prev) + } + }) + + // Before installation the gauge reports 0. + inputRevisionBufferLenFn.Store(nil) + assert.Equal(t, 0.0, prometheustestutil.ToFloat64(inputRevisionBufferDepth)) + + // After installation the gauge reflects the closure's return value live. + depth := 0 + setInputRevisionBufferLenSource(func() int { return depth }) + + depth = 4 + assert.Equal(t, 4.0, prometheustestutil.ToFloat64(inputRevisionBufferDepth)) + + depth = 0 + assert.Equal(t, 0.0, prometheustestutil.ToFloat64(inputRevisionBufferDepth)) + + // A second call replaces the source. + setInputRevisionBufferLenSource(func() int { return 7 }) + assert.Equal(t, 7.0, prometheustestutil.ToFloat64(inputRevisionBufferDepth)) +}