From 23a6231f045efa005fdbc141724b08b622be4927 Mon Sep 17 00:00:00 2001 From: Ruinan Liu Date: Tue, 28 Jul 2026 21:27:15 +0000 Subject: [PATCH 1/7] reduce input buffer --- internal/controllers/watch/kind.go | 66 ++---- internal/controllers/watch/watch.go | 3 + internal/flowcontrol/inputrevbuffer.go | 248 ++++++++++++++++++++ internal/flowcontrol/inputrevbuffer_test.go | 172 ++++++++++++++ 4 files changed, 446 insertions(+), 43 deletions(-) create mode 100644 internal/flowcontrol/inputrevbuffer.go create mode 100644 internal/flowcontrol/inputrevbuffer_test.go diff --git a/internal/controllers/watch/kind.go b/internal/controllers/watch/kind.go index 3b385543..c12bb714 100644 --- a/internal/controllers/watch/kind.go +++ b/internal/controllers/watch/kind.go @@ -10,6 +10,7 @@ import ( "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 +18,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 +36,7 @@ func SetKindWatchRateLimit(rps int) { type KindWatchController struct { client client.Client + buffer *flowcontrol.CompositionInputRevisionWriteBuffer gvk schema.GroupVersionKind cancel context.CancelFunc } @@ -45,6 +46,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 +221,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 +232,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 { diff --git a/internal/controllers/watch/watch.go b/internal/controllers/watch/watch.go index 3144af41..d8ada56c 100644 --- a/internal/controllers/watch/watch.go +++ b/internal/controllers/watch/watch.go @@ -9,12 +9,14 @@ 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 } @@ -26,6 +28,7 @@ func NewController(mgr ctrl.Manager) error { Complete(&WatchController{ mgr: mgr, client: mgr.GetClient(), + buffer: flowcontrol.NewCompositionInputRevisionWriteBufferForManager(mgr), 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..ed87b7f8 --- /dev/null +++ b/internal/flowcontrol/inputrevbuffer.go @@ -0,0 +1,248 @@ +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 Status().Update 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 { + r := NewCompositionInputRevisionWriteBuffer(mgr.GetClient()) + mgr.Add(r) + return r +} + +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", + }) + 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. + w.queued[comp] = false + insertionTime := w.insertionTime[comp] + ops := w.state[comp] + delete(w.state, comp) + w.mut.Unlock() + + if len(ops) == 0 { + w.queue.Forget(comp) + w.mut.Lock() + delete(w.insertionTime, comp) + w.mut.Unlock() + return true + } + + if w.updateComposition(ctx, insertionTime, comp, ops) { + w.queue.Forget(comp) + w.mut.Lock() + delete(w.insertionTime, comp) + w.mut.Unlock() + 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 +} + +// updateComposition applies all pending per-key ops to the composition in a single +// Status().Update. 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") + return false + } + + 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 + } + + err = w.client.Status().Update(ctx, obj) + 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") + return false + } + + 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..87778eed --- /dev/null +++ b/internal/flowcontrol/inputrevbuffer_test.go @@ -0,0 +1,172 @@ +package flowcontrol + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "sigs.k8s.io/controller-runtime/pkg/client" + + 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 +} From f4a1d7e8fbed86c183c9d8eb94e3b6fc7942b010 Mon Sep 17 00:00:00 2001 From: Ruinan Liu Date: Wed, 29 Jul 2026 22:16:51 +0000 Subject: [PATCH 2/7] flowcontrol: add metrics for composition input-revision write buffer Adds a queue-depth gauge (eno_input_revision_write_buffer_depth), a flush-error counter partitioned by op (eno_input_revision_write_buffer_flush_errors_total), and a successful-flush counter (eno_input_revision_write_buffer_flush_total), mirroring the existing ResourceSliceWriteBuffer metrics. These errors retry silently today, so without this counter they're invisible. Addresses review item 5 on PR #634. --- internal/flowcontrol/metrics.go | 49 +++++++++++++++++++++++++++- internal/flowcontrol/metrics_test.go | 43 ++++++++++++++++++++++++ 2 files changed, 91 insertions(+), 1 deletion(-) diff --git a/internal/flowcontrol/metrics.go b/internal/flowcontrol/metrics.go index 7cb22de1..42455cf7 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,10 @@ 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..94ba7043 100644 --- a/internal/flowcontrol/metrics_test.go +++ b/internal/flowcontrol/metrics_test.go @@ -50,3 +50,46 @@ 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)) +} + From e5b8c7234915f5fea2965d1d106e70cc3f111dba Mon Sep 17 00:00:00 2001 From: Ruinan Liu Date: Wed, 29 Jul 2026 22:20:09 +0000 Subject: [PATCH 3/7] resolving comments --- internal/controllers/watch/kind.go | 25 -- internal/controllers/watch/kind_test.go | 194 ---------- internal/controllers/watch/watch.go | 9 +- internal/flowcontrol/inputrevbuffer.go | 61 +++- internal/flowcontrol/inputrevbuffer_test.go | 373 ++++++++++++++++++++ 5 files changed, 424 insertions(+), 238 deletions(-) diff --git a/internal/controllers/watch/kind.go b/internal/controllers/watch/kind.go index c12bb714..a0e02e03 100644 --- a/internal/controllers/watch/kind.go +++ b/internal/controllers/watch/kind.go @@ -5,7 +5,6 @@ import ( "fmt" "math/rand" "path" - "reflect" "slices" "time" @@ -301,27 +300,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 d8ada56c..a7e11948 100644 --- a/internal/controllers/watch/watch.go +++ b/internal/controllers/watch/watch.go @@ -21,14 +21,19 @@ type WatchController struct { } 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: flowcontrol.NewCompositionInputRevisionWriteBufferForManager(mgr), + buffer: buffer, refControllers: map[apiv1.ResourceRef]*KindWatchController{}, }) if err != nil { diff --git a/internal/flowcontrol/inputrevbuffer.go b/internal/flowcontrol/inputrevbuffer.go index ed87b7f8..1099bdd4 100644 --- a/internal/flowcontrol/inputrevbuffer.go +++ b/internal/flowcontrol/inputrevbuffer.go @@ -25,7 +25,7 @@ type inputRevisionOp struct { // CompositionInputRevisionWriteBuffer reduces load on etcd/apiserver by collecting // input-revision updates for each Composition over a short window and applying them -// in a single Status().Update per Composition. +// 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 @@ -47,10 +47,12 @@ type CompositionInputRevisionWriteBuffer struct { queue workqueue.TypedRateLimitingInterface[types.NamespacedName] } -func NewCompositionInputRevisionWriteBufferForManager(mgr ctrl.Manager) *CompositionInputRevisionWriteBuffer { +func NewCompositionInputRevisionWriteBufferForManager(mgr ctrl.Manager) (*CompositionInputRevisionWriteBuffer, error) { r := NewCompositionInputRevisionWriteBuffer(mgr.GetClient()) - mgr.Add(r) - return r + if err := mgr.Add(r); err != nil { + return nil, err + } + return r, nil } func NewCompositionInputRevisionWriteBuffer(cli client.Client) *CompositionInputRevisionWriteBuffer { @@ -59,6 +61,7 @@ func NewCompositionInputRevisionWriteBuffer(cli client.Client) *CompositionInput workqueue.TypedRateLimitingQueueConfig[types.NamespacedName]{ Name: "compositionInputRevisionWriteBuffer", }) + setInputRevisionBufferLenSource(q.Len) return &CompositionInputRevisionWriteBuffer{ client: cli, state: make(map[types.NamespacedName]map[string]inputRevisionOp), @@ -124,25 +127,19 @@ func (w *CompositionInputRevisionWriteBuffer) processQueueItem(ctx context.Conte 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. - w.queued[comp] = false + delete(w.queued, comp) insertionTime := w.insertionTime[comp] ops := w.state[comp] delete(w.state, comp) w.mut.Unlock() if len(ops) == 0 { - w.queue.Forget(comp) - w.mut.Lock() - delete(w.insertionTime, comp) - w.mut.Unlock() + w.settle(comp) return true } if w.updateComposition(ctx, insertionTime, comp, ops) { - w.queue.Forget(comp) - w.mut.Lock() - delete(w.insertionTime, comp) - w.mut.Unlock() + w.settle(comp) return true } @@ -167,9 +164,28 @@ func (w *CompositionInputRevisionWriteBuffer) processQueueItem(ctx context.Conte return true } -// updateComposition applies all pending per-key ops to the composition in a single -// Status().Update. Returns true when the work is done (including no-op and composition-deleted), -// false when it should be retried. +// 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 +// scoped to whatever fields we actually changed (status.inputRevisions). 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) @@ -181,9 +197,12 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont } 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 { @@ -200,7 +219,12 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont return true } - err = w.client.Status().Update(ctx, obj) + // client.MergeFrom diffs against the unmodified copy, so the resulting patch only + // contains status.inputRevisions - the one field we mutated - and applies cleanly + // even if status.inputRevisions didn't exist yet. Unlike a full Status().Update, a + // concurrent writer to another status field (e.g. currentSynthesis) is never + // clobbered and never causes a conflict here. + err = w.client.Status().Patch(ctx, obj, client.MergeFrom(original)) if k8serrors.IsNotFound(err) { logger.V(1).Info("composition deleted - dropping buffered input revision updates") return true @@ -211,9 +235,11 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont } 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 } @@ -246,3 +272,4 @@ func removeInputRevision(comp *apiv1.Composition, key string) bool { } return false } + diff --git a/internal/flowcontrol/inputrevbuffer_test.go b/internal/flowcontrol/inputrevbuffer_test.go index 87778eed..6ecdfeb8 100644 --- a/internal/flowcontrol/inputrevbuffer_test.go +++ b/internal/flowcontrol/inputrevbuffer_test.go @@ -1,11 +1,16 @@ package flowcontrol import ( + "context" + "strconv" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "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" @@ -170,3 +175,371 @@ func inputRevisionKeys(comp *apiv1.Composition) []string { } 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 that once a composition is created, +// flushed, and deleted, the buffer retains no per-composition bookkeeping - guarding +// against the `queued` map (previously only ever set to false, never deleted) growing +// without bound on a cluster with composition churn. +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 simulates a producer +// enqueuing a new update for the same composition while a flush is still in flight (inside +// the window between reading state and settling). Before the fix, the success path +// unconditionally deleted insertionTime, so the next flush would compute latency against +// the zero time (~62 billion ms). The fix must keep insertionTime intact whenever a +// racing update is still pending. +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 proves that the rate limiter +// backs off exponentially for a composition whose inputs keep churning: Forget must only be +// called once a flush finds nothing pending for that composition. +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 proves that a +// composition with a single, isolated update is not penalized by the exponential backoff - +// it flushes and its rate limiter resets 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") +} + +// TestInputRevisionWriteBufferPatchDoesNotConflictWithUnrelatedStatusWrite proves that +// scoping the flush to a JSON patch on status.inputRevisions avoids the conflict a full +// Status().Update would hit when another controller concurrently writes a different status +// field (e.g. CurrentSynthesis) between our Get and our write. +func TestInputRevisionWriteBufferPatchDoesNotConflictWithUnrelatedStatusWrite(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) + + // The scoped patch succeeded on the first attempt despite the concurrent write. + 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") +} + From 27befafcfd7f99428ea59dce456757564d29bdd6 Mon Sep 17 00:00:00 2001 From: Ruinan Liu Date: Thu, 30 Jul 2026 00:16:56 +0000 Subject: [PATCH 4/7] Fix tests --- internal/controllers/watch/kind.go | 1 - internal/flowcontrol/inputrevbuffer.go | 49 +++-- internal/flowcontrol/inputrevbuffer_test.go | 219 ++++++++++++++++++++ internal/flowcontrol/metrics.go | 3 +- internal/flowcontrol/metrics_test.go | 3 +- 5 files changed, 259 insertions(+), 16 deletions(-) diff --git a/internal/controllers/watch/kind.go b/internal/controllers/watch/kind.go index a0e02e03..522ddd89 100644 --- a/internal/controllers/watch/kind.go +++ b/internal/controllers/watch/kind.go @@ -299,4 +299,3 @@ func isOptionalRef(synth *apiv1.Synthesizer, key string) bool { } return false } - diff --git a/internal/flowcontrol/inputrevbuffer.go b/internal/flowcontrol/inputrevbuffer.go index 1099bdd4..5e237dd9 100644 --- a/internal/flowcontrol/inputrevbuffer.go +++ b/internal/flowcontrol/inputrevbuffer.go @@ -2,6 +2,7 @@ package flowcontrol import ( "context" + "encoding/json" "reflect" "sync" "time" @@ -182,10 +183,9 @@ func (w *CompositionInputRevisionWriteBuffer) settle(comp types.NamespacedName) delete(w.insertionTime, comp) } -// updateComposition applies all pending per-key ops to the composition via a merge patch -// scoped to whatever fields we actually changed (status.inputRevisions). Returns true when -// the work is done (including no-op and composition-deleted), false when it should be -// retried. +// updateComposition applies all pending per-key ops to the composition via a JSON patch +// scoped to status.inputRevisions. 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) @@ -201,7 +201,9 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont return false } - original := obj.DeepCopy() + // Snapshot before mutating obj in place - not a slice alias, since setInputRevisions + // may overwrite obj.Status.InputRevisions' backing array. + before := append([]apiv1.InputRevisions(nil), obj.Status.InputRevisions...) modified := false for key, op := range ops { @@ -219,12 +221,20 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont return true } - // client.MergeFrom diffs against the unmodified copy, so the resulting patch only - // contains status.inputRevisions - the one field we mutated - and applies cleanly - // even if status.inputRevisions didn't exist yet. Unlike a full Status().Update, a - // concurrent writer to another status field (e.g. currentSynthesis) is never - // clobbered and never causes a conflict here. - err = w.client.Status().Patch(ctx, obj, client.MergeFrom(original)) + // Scoped to status.inputRevisions with a test op against the value we just read: a + // concurrent writer to another status field (e.g. currentSynthesis) neither conflicts + // nor gets clobbered, while a concurrent writer to inputRevisions itself - the pruning + // controller, or another flush racing ahead of our (possibly stale) cached read - fails + // the test op, so we retry against a fresh read instead of silently overwriting it. + patch := buildInputRevisionPatch(before, obj.Status.InputRevisions) + patchJson, err := json.Marshal(&patch) + if err != nil { + logger.Error(err, "unable to encode input revision patch") + inputRevisionBufferFlushErrors.WithLabelValues("marshal").Inc() + return false + } + + err = w.client.Status().Patch(ctx, obj, client.RawPatch(types.JSONPatchType, patchJson)) if k8serrors.IsNotFound(err) { logger.V(1).Info("composition deleted - dropping buffered input revision updates") return true @@ -244,6 +254,22 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont return true } +// buildInputRevisionPatch returns a JSON patch scoped to status.inputRevisions: a test op +// asserting the value is still what we read, followed by an add (path not yet present) or +// replace (path present) with the new value. reuses the jsonPatch type from writebuffer.go. +func buildInputRevisionPatch(before, after []apiv1.InputRevisions) []*jsonPatch { + if len(before) == 0 { + return []*jsonPatch{ + {Op: "test", Path: "/status/inputRevisions", Value: nil}, + {Op: "add", Path: "/status/inputRevisions", Value: after}, + } + } + return []*jsonPatch{ + {Op: "test", Path: "/status/inputRevisions", Value: before}, + {Op: "replace", Path: "/status/inputRevisions", Value: after}, + } +} + // 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 { @@ -272,4 +298,3 @@ func removeInputRevision(comp *apiv1.Composition, key string) bool { } return false } - diff --git a/internal/flowcontrol/inputrevbuffer_test.go b/internal/flowcontrol/inputrevbuffer_test.go index 6ecdfeb8..092258ed 100644 --- a/internal/flowcontrol/inputrevbuffer_test.go +++ b/internal/flowcontrol/inputrevbuffer_test.go @@ -2,12 +2,18 @@ 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" @@ -543,3 +549,216 @@ func TestInputRevisionWriteBufferPatchDoesNotConflictWithUnrelatedStatusWrite(t assert.Equal(t, "concurrent-write", comp.Status.CurrentSynthesis.UUID, "unrelated status field must be preserved") } +// TestInputRevisionWriteBufferDoesNotResurrectConcurrentlyPrunedKey proves the test op +// protects against a second writer to status.inputRevisions itself (e.g. the pruning +// controller in internal/controllers/watch/pruning.go), unlike an unlocked merge patch, +// which would silently resurrect whatever the buffer's stale read still had for the key +// the pruner just removed. +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: test op fails against the pruner's write, 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) +} + +// TestInputRevisionWriteBufferStaleCacheReadDoesNotRevertPreviousFlush simulates the +// flush-N / flush-N+1 sequence where the informer cache backing w.client.Get lags behind +// flush N's own write: flush N+1 (for a different key) reads a stale pre-flush-N snapshot. +// A JSON merge patch would blindly replace the whole array with that stale view, reverting +// flush N's key; the test op here must instead fail and force a retry against a fresh read. +func TestInputRevisionWriteBufferStaleCacheReadDoesNotRevertPreviousFlush(t *testing.T) { + ctx := testutil.NewContext(t) + var stale *apiv1.Composition + returnStale := 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 { + if returnStale { + if comp, ok := obj.(*apiv1.Composition); ok { + returnStale = false // the buffer's retry then reads a fresh value + stale.DeepCopyInto(comp) + return nil + } + } + return c.Get(ctx, key, obj, opts...) + }, + }) + w := NewCompositionInputRevisionWriteBuffer(cli) + + comp := newTestComposition(t, cli, "test-comp-1") + nsn := client.ObjectKeyFromObject(comp) + comp.Status.InputRevisions = []apiv1.InputRevisions{ + {Key: "A", ResourceVersion: "1"}, + {Key: "B", ResourceVersion: "1"}, + } + require.NoError(t, cli.Status().Update(ctx, comp)) + stale = comp.DeepCopy() // the pre-flush-N snapshot an informer cache would still be serving + + // Flush N: patch A to rv2. + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "A", ResourceVersion: "2"}) + w.processQueueItem(ctx) + + require.NoError(t, cli.Get(ctx, nsn, comp)) + require.Equal(t, "2", revisionFor(comp, "A"), "flush N must have applied") + + // Flush N+1: patch B to rv2, but the buffer's next Get returns the stale pre-flush-N + // snapshot, as if the informer cache had not caught up with flush N's write yet. + returnStale = true + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "B", ResourceVersion: "2"}) + w.processQueueItem(ctx) // first attempt against the stale read: test op fails, retries + w.processQueueItem(ctx) // retry against a fresh read: succeeds + + require.NoError(t, cli.Get(ctx, nsn, comp)) + assert.Equal(t, "2", revisionFor(comp, "A"), "A must not revert to rv1") + assert.Equal(t, "2", revisionFor(comp, "B"), "B must be updated") +} + +// TestInputRevisionWriteBufferPatchConflictIsRetried proves the k8serrors.IsConflict branch +// in updateComposition is live: a 409 from the patch call is logged as a conflict and the +// buffered update is retried rather than dropped. +func TestInputRevisionWriteBufferPatchConflictIsRetried(t *testing.T) { + ctx := testutil.NewContext(t) + failNext := true + 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 { + if failNext { + failNext = false + return k8serrors.NewConflict(schema.GroupResource{Group: "eno.azure.io", Resource: "compositions"}, obj.GetName(), errors.New("simulated conflict")) + } + 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) // rejected with a conflict, buffered update is kept for retry + w.processQueueItem(ctx) // retry succeeds + + require.NoError(t, cli.Get(ctx, nsn, comp)) + require.Len(t, comp.Status.InputRevisions, 1) + assert.Equal(t, "foo", comp.Status.InputRevisions[0].Key) +} + +func revisionFor(comp *apiv1.Composition, key string) string { + for _, ir := range comp.Status.InputRevisions { + if ir.Key == key { + return ir.ResourceVersion + } + } + return "" +} + +// TestInputRevisionWriteBufferConcurrentWrites drives many goroutines that concurrently call +// PatchInputRevisionAsync/RemoveInputRevisionAsync while the buffer's real Start loop flushes +// in the background, and asserts every composition converges to the expected final state. +// Each (composition, key) pair is owned by exactly one goroutine, so there's no ambiguity +// about which write should win - this isolates the assertion from scheduling nondeterminism +// while still exercising the buffer's locking, coalescing, and workqueue under genuine +// concurrency (run with -race). +func TestInputRevisionWriteBufferConcurrentWrites(t *testing.T) { + ctx := testutil.NewContext(t) + cli := testutil.NewClient(t) + w := NewCompositionInputRevisionWriteBuffer(cli) + go w.Start(ctx) + + const compCount = 6 + const keysPerComp = 4 + const iterations = 40 + + nsns := make([]client.ObjectKey, compCount) + for i := range nsns { + comp := newTestComposition(t, cli, fmt.Sprintf("test-comp-%d", i)) + nsns[i] = client.ObjectKeyFromObject(comp) + } + + type finalState struct { + revision string + removed bool + } + var mu sync.Mutex + want := make(map[client.ObjectKey]map[string]finalState, compCount) + for _, nsn := range nsns { + want[nsn] = make(map[string]finalState, 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 finalState + for i := 0; i < iterations; i++ { + if rand.IntN(4) == 0 { + w.RemoveInputRevisionAsync(nsn, key) + last = finalState{removed: true} + } else { + rv := strconv.Itoa(i) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: key, ResourceVersion: rv}) + last = finalState{revision: rv} + } + time.Sleep(time.Duration(rand.IntN(2)) * time.Millisecond) + } + mu.Lock() + want[nsn][key] = last + mu.Unlock() + }() + } + } + wg.Wait() + + // Wait for the buffer to fully drain before asserting the final state. + 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, exp := range want[nsn] { + rv := revisionFor(comp, key) + if exp.removed { + assert.Empty(t, rv, "key %s on %s should have been removed", key, nsn.Name) + } else { + assert.Equal(t, exp.revision, rv, "key %s on %s has wrong final revision", key, nsn.Name) + } + } + } +} diff --git a/internal/flowcontrol/metrics.go b/internal/flowcontrol/metrics.go index 42455cf7..9d1bae9d 100644 --- a/internal/flowcontrol/metrics.go +++ b/internal/flowcontrol/metrics.go @@ -68,7 +68,7 @@ var ( // 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). + // Errors hit while flushing input-revision patches, partitioned by op (get/patch/marshal). // These silently retry today, so without this counter they're invisible. inputRevisionBufferFlushErrors = prometheus.NewCounterVec( prometheus.CounterOpts{ @@ -102,4 +102,3 @@ func setWriteBufferLenSource(fn func() int) { func setInputRevisionBufferLenSource(fn func() int) { inputRevisionBufferLenFn.Store(&fn) } - diff --git a/internal/flowcontrol/metrics_test.go b/internal/flowcontrol/metrics_test.go index 94ba7043..00b85a22 100644 --- a/internal/flowcontrol/metrics_test.go +++ b/internal/flowcontrol/metrics_test.go @@ -56,9 +56,11 @@ func TestInputRevisionBufferFlushErrorsLabels(t *testing.T) { inputRevisionBufferFlushErrors.WithLabelValues("get").Inc() inputRevisionBufferFlushErrors.WithLabelValues("get").Inc() + inputRevisionBufferFlushErrors.WithLabelValues("marshal").Inc() inputRevisionBufferFlushErrors.WithLabelValues("patch").Inc() assert.Equal(t, 2.0, prometheustestutil.ToFloat64(inputRevisionBufferFlushErrors.WithLabelValues("get"))) + assert.Equal(t, 1.0, prometheustestutil.ToFloat64(inputRevisionBufferFlushErrors.WithLabelValues("marshal"))) 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"))) @@ -92,4 +94,3 @@ func TestInputRevisionBufferDepthGauge(t *testing.T) { setInputRevisionBufferLenSource(func() int { return 7 }) assert.Equal(t, 7.0, prometheustestutil.ToFloat64(inputRevisionBufferDepth)) } - From e2ae64ff9286427fe817ca97e18561d48837d67e Mon Sep 17 00:00:00 2001 From: Ruinan Liu Date: Thu, 30 Jul 2026 01:49:46 +0000 Subject: [PATCH 5/7] fixing comments --- internal/flowcontrol/inputrevbuffer.go | 44 +---- internal/flowcontrol/inputrevbuffer_test.go | 208 +++++++------------- internal/flowcontrol/metrics.go | 2 +- internal/flowcontrol/metrics_test.go | 2 - 4 files changed, 75 insertions(+), 181 deletions(-) diff --git a/internal/flowcontrol/inputrevbuffer.go b/internal/flowcontrol/inputrevbuffer.go index 5e237dd9..e5063e35 100644 --- a/internal/flowcontrol/inputrevbuffer.go +++ b/internal/flowcontrol/inputrevbuffer.go @@ -2,7 +2,6 @@ package flowcontrol import ( "context" - "encoding/json" "reflect" "sync" "time" @@ -183,8 +182,8 @@ func (w *CompositionInputRevisionWriteBuffer) settle(comp types.NamespacedName) delete(w.insertionTime, comp) } -// updateComposition applies all pending per-key ops to the composition via a JSON patch -// scoped to status.inputRevisions. Returns true when the work is done (including no-op and +// 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) @@ -201,9 +200,7 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont return false } - // Snapshot before mutating obj in place - not a slice alias, since setInputRevisions - // may overwrite obj.Status.InputRevisions' backing array. - before := append([]apiv1.InputRevisions(nil), obj.Status.InputRevisions...) + original := obj.DeepCopy() modified := false for key, op := range ops { @@ -221,20 +218,11 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont return true } - // Scoped to status.inputRevisions with a test op against the value we just read: a - // concurrent writer to another status field (e.g. currentSynthesis) neither conflicts - // nor gets clobbered, while a concurrent writer to inputRevisions itself - the pruning - // controller, or another flush racing ahead of our (possibly stale) cached read - fails - // the test op, so we retry against a fresh read instead of silently overwriting it. - patch := buildInputRevisionPatch(before, obj.Status.InputRevisions) - patchJson, err := json.Marshal(&patch) - if err != nil { - logger.Error(err, "unable to encode input revision patch") - inputRevisionBufferFlushErrors.WithLabelValues("marshal").Inc() - return false - } - - err = w.client.Status().Patch(ctx, obj, client.RawPatch(types.JSONPatchType, patchJson)) + // 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 @@ -254,22 +242,6 @@ func (w *CompositionInputRevisionWriteBuffer) updateComposition(ctx context.Cont return true } -// buildInputRevisionPatch returns a JSON patch scoped to status.inputRevisions: a test op -// asserting the value is still what we read, followed by an add (path not yet present) or -// replace (path present) with the new value. reuses the jsonPatch type from writebuffer.go. -func buildInputRevisionPatch(before, after []apiv1.InputRevisions) []*jsonPatch { - if len(before) == 0 { - return []*jsonPatch{ - {Op: "test", Path: "/status/inputRevisions", Value: nil}, - {Op: "add", Path: "/status/inputRevisions", Value: after}, - } - } - return []*jsonPatch{ - {Op: "test", Path: "/status/inputRevisions", Value: before}, - {Op: "replace", Path: "/status/inputRevisions", Value: after}, - } -} - // 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 { diff --git a/internal/flowcontrol/inputrevbuffer_test.go b/internal/flowcontrol/inputrevbuffer_test.go index 092258ed..188e8fdb 100644 --- a/internal/flowcontrol/inputrevbuffer_test.go +++ b/internal/flowcontrol/inputrevbuffer_test.go @@ -2,7 +2,6 @@ package flowcontrol import ( "context" - "errors" "fmt" "math/rand/v2" "strconv" @@ -12,8 +11,6 @@ import ( "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" @@ -375,10 +372,7 @@ func TestRemoveInputRevision(t *testing.T) { } } -// TestInputRevisionWriteBufferMapsDoNotLeak asserts that once a composition is created, -// flushed, and deleted, the buffer retains no per-composition bookkeeping - guarding -// against the `queued` map (previously only ever set to false, never deleted) growing -// without bound on a cluster with composition churn. +// 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) @@ -401,12 +395,7 @@ func TestInputRevisionWriteBufferMapsDoNotLeak(t *testing.T) { assert.Empty(t, w.queued, "queued must not retain an entry once settled") } -// TestInputRevisionWriteBufferInsertionTimePreservedDuringRace simulates a producer -// enqueuing a new update for the same composition while a flush is still in flight (inside -// the window between reading state and settling). Before the fix, the success path -// unconditionally deleted insertionTime, so the next flush would compute latency against -// the zero time (~62 billion ms). The fix must keep insertionTime intact whenever a -// racing update is still pending. +// 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 @@ -451,9 +440,7 @@ func TestInputRevisionWriteBufferInsertionTimePreservedDuringRace(t *testing.T) assert.Empty(t, w.insertionTime) } -// TestInputRevisionWriteBufferBackoffGrowsUnderSustainedLoad proves that the rate limiter -// backs off exponentially for a composition whose inputs keep churning: Forget must only be -// called once a flush finds nothing pending for that composition. +// 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 @@ -491,9 +478,7 @@ func TestInputRevisionWriteBufferBackoffGrowsUnderSustainedLoad(t *testing.T) { assert.Equal(t, 0, w.queue.NumRequeues(nsn), "rate limit should reset once the composition goes idle") } -// TestInputRevisionWriteBufferBackoffStaysBaseForIsolatedComposition proves that a -// composition with a single, isolated update is not penalized by the exponential backoff - -// it flushes and its rate limiter resets to the fast path. +// 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) @@ -508,11 +493,8 @@ func TestInputRevisionWriteBufferBackoffStaysBaseForIsolatedComposition(t *testi assert.Equal(t, 0, w.queue.NumRequeues(nsn), "an isolated update should reset the rate limit to the fast path") } -// TestInputRevisionWriteBufferPatchDoesNotConflictWithUnrelatedStatusWrite proves that -// scoping the flush to a JSON patch on status.inputRevisions avoids the conflict a full -// Status().Update would hit when another controller concurrently writes a different status -// field (e.g. CurrentSynthesis) between our Get and our write. -func TestInputRevisionWriteBufferPatchDoesNotConflictWithUnrelatedStatusWrite(t *testing.T) { +// 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{ @@ -536,9 +518,9 @@ func TestInputRevisionWriteBufferPatchDoesNotConflictWithUnrelatedStatusWrite(t nsn := client.ObjectKeyFromObject(comp) w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "foo", ResourceVersion: "1"}) - w.processQueueItem(ctx) + w.processQueueItem(ctx) // our resourceVersion is now stale, conflicts, retries + w.processQueueItem(ctx) // retry against a fresh read: succeeds - // The scoped patch succeeded on the first attempt despite the concurrent write. assert.Empty(t, w.state) assert.Equal(t, 0, w.queue.NumRequeues(nsn)) @@ -546,14 +528,11 @@ func TestInputRevisionWriteBufferPatchDoesNotConflictWithUnrelatedStatusWrite(t 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 proves the test op -// protects against a second writer to status.inputRevisions itself (e.g. the pruning -// controller in internal/controllers/watch/pruning.go), unlike an unlocked merge patch, -// which would silently resurrect whatever the buffer's stale read still had for the key -// the pruner just removed. +// 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 @@ -585,7 +564,7 @@ func TestInputRevisionWriteBufferDoesNotResurrectConcurrentlyPrunedKey(t *testin // 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: test op fails against the pruner's write, retries + 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)) @@ -593,86 +572,6 @@ func TestInputRevisionWriteBufferDoesNotResurrectConcurrentlyPrunedKey(t *testin assert.Empty(t, w.state) } -// TestInputRevisionWriteBufferStaleCacheReadDoesNotRevertPreviousFlush simulates the -// flush-N / flush-N+1 sequence where the informer cache backing w.client.Get lags behind -// flush N's own write: flush N+1 (for a different key) reads a stale pre-flush-N snapshot. -// A JSON merge patch would blindly replace the whole array with that stale view, reverting -// flush N's key; the test op here must instead fail and force a retry against a fresh read. -func TestInputRevisionWriteBufferStaleCacheReadDoesNotRevertPreviousFlush(t *testing.T) { - ctx := testutil.NewContext(t) - var stale *apiv1.Composition - returnStale := 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 { - if returnStale { - if comp, ok := obj.(*apiv1.Composition); ok { - returnStale = false // the buffer's retry then reads a fresh value - stale.DeepCopyInto(comp) - return nil - } - } - return c.Get(ctx, key, obj, opts...) - }, - }) - w := NewCompositionInputRevisionWriteBuffer(cli) - - comp := newTestComposition(t, cli, "test-comp-1") - nsn := client.ObjectKeyFromObject(comp) - comp.Status.InputRevisions = []apiv1.InputRevisions{ - {Key: "A", ResourceVersion: "1"}, - {Key: "B", ResourceVersion: "1"}, - } - require.NoError(t, cli.Status().Update(ctx, comp)) - stale = comp.DeepCopy() // the pre-flush-N snapshot an informer cache would still be serving - - // Flush N: patch A to rv2. - w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "A", ResourceVersion: "2"}) - w.processQueueItem(ctx) - - require.NoError(t, cli.Get(ctx, nsn, comp)) - require.Equal(t, "2", revisionFor(comp, "A"), "flush N must have applied") - - // Flush N+1: patch B to rv2, but the buffer's next Get returns the stale pre-flush-N - // snapshot, as if the informer cache had not caught up with flush N's write yet. - returnStale = true - w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: "B", ResourceVersion: "2"}) - w.processQueueItem(ctx) // first attempt against the stale read: test op fails, retries - w.processQueueItem(ctx) // retry against a fresh read: succeeds - - require.NoError(t, cli.Get(ctx, nsn, comp)) - assert.Equal(t, "2", revisionFor(comp, "A"), "A must not revert to rv1") - assert.Equal(t, "2", revisionFor(comp, "B"), "B must be updated") -} - -// TestInputRevisionWriteBufferPatchConflictIsRetried proves the k8serrors.IsConflict branch -// in updateComposition is live: a 409 from the patch call is logged as a conflict and the -// buffered update is retried rather than dropped. -func TestInputRevisionWriteBufferPatchConflictIsRetried(t *testing.T) { - ctx := testutil.NewContext(t) - failNext := true - 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 { - if failNext { - failNext = false - return k8serrors.NewConflict(schema.GroupResource{Group: "eno.azure.io", Resource: "compositions"}, obj.GetName(), errors.New("simulated conflict")) - } - 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) // rejected with a conflict, buffered update is kept for retry - w.processQueueItem(ctx) // retry succeeds - - require.NoError(t, cli.Get(ctx, nsn, comp)) - require.Len(t, comp.Status.InputRevisions, 1) - assert.Equal(t, "foo", comp.Status.InputRevisions[0].Key) -} - func revisionFor(comp *apiv1.Composition, key string) string { for _, ir := range comp.Status.InputRevisions { if ir.Key == key { @@ -682,22 +581,16 @@ func revisionFor(comp *apiv1.Composition, key string) string { return "" } -// TestInputRevisionWriteBufferConcurrentWrites drives many goroutines that concurrently call -// PatchInputRevisionAsync/RemoveInputRevisionAsync while the buffer's real Start loop flushes -// in the background, and asserts every composition converges to the expected final state. -// Each (composition, key) pair is owned by exactly one goroutine, so there's no ambiguity -// about which write should win - this isolates the assertion from scheduling nondeterminism -// while still exercising the buffer's locking, coalescing, and workqueue under genuine -// concurrency (run with -race). +// 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 = 6 - const keysPerComp = 4 - const iterations = 40 + const compCount, keysPerComp, iterations = 3, 2, 15 nsns := make([]client.ObjectKey, compCount) for i := range nsns { @@ -705,14 +598,10 @@ func TestInputRevisionWriteBufferConcurrentWrites(t *testing.T) { nsns[i] = client.ObjectKeyFromObject(comp) } - type finalState struct { - revision string - removed bool - } var mu sync.Mutex - want := make(map[client.ObjectKey]map[string]finalState, compCount) + want := make(map[client.ObjectKey]map[string]string, compCount) // "" means removed for _, nsn := range nsns { - want[nsn] = make(map[string]finalState, keysPerComp) + want[nsn] = make(map[string]string, keysPerComp) } var wg sync.WaitGroup @@ -722,17 +611,15 @@ func TestInputRevisionWriteBufferConcurrentWrites(t *testing.T) { wg.Add(1) go func() { defer wg.Done() - var last finalState + var last string for i := 0; i < iterations; i++ { if rand.IntN(4) == 0 { w.RemoveInputRevisionAsync(nsn, key) - last = finalState{removed: true} + last = "" } else { - rv := strconv.Itoa(i) - w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: key, ResourceVersion: rv}) - last = finalState{revision: rv} + last = strconv.Itoa(i) + w.PatchInputRevisionAsync(nsn, &apiv1.InputRevisions{Key: key, ResourceVersion: last}) } - time.Sleep(time.Duration(rand.IntN(2)) * time.Millisecond) } mu.Lock() want[nsn][key] = last @@ -742,7 +629,6 @@ func TestInputRevisionWriteBufferConcurrentWrites(t *testing.T) { } wg.Wait() - // Wait for the buffer to fully drain before asserting the final state. testutil.Eventually(t, func() bool { w.mut.Lock() defer w.mut.Unlock() @@ -752,13 +638,51 @@ func TestInputRevisionWriteBufferConcurrentWrites(t *testing.T) { for _, nsn := range nsns { comp := &apiv1.Composition{} require.NoError(t, cli.Get(ctx, nsn, comp)) - for key, exp := range want[nsn] { - rv := revisionFor(comp, key) - if exp.removed { - assert.Empty(t, rv, "key %s on %s should have been removed", key, nsn.Name) - } else { - assert.Equal(t, exp.revision, rv, "key %s on %s has wrong final revision", key, nsn.Name) - } + 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 9d1bae9d..347b0765 100644 --- a/internal/flowcontrol/metrics.go +++ b/internal/flowcontrol/metrics.go @@ -68,7 +68,7 @@ var ( // 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/marshal). + // 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{ diff --git a/internal/flowcontrol/metrics_test.go b/internal/flowcontrol/metrics_test.go index 00b85a22..8401f95f 100644 --- a/internal/flowcontrol/metrics_test.go +++ b/internal/flowcontrol/metrics_test.go @@ -56,11 +56,9 @@ func TestInputRevisionBufferFlushErrorsLabels(t *testing.T) { inputRevisionBufferFlushErrors.WithLabelValues("get").Inc() inputRevisionBufferFlushErrors.WithLabelValues("get").Inc() - inputRevisionBufferFlushErrors.WithLabelValues("marshal").Inc() inputRevisionBufferFlushErrors.WithLabelValues("patch").Inc() assert.Equal(t, 2.0, prometheustestutil.ToFloat64(inputRevisionBufferFlushErrors.WithLabelValues("get"))) - assert.Equal(t, 1.0, prometheustestutil.ToFloat64(inputRevisionBufferFlushErrors.WithLabelValues("marshal"))) 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"))) From bafbb3eae9a27b49dfe3f8735ae6b3e2c1aac2c3 Mon Sep 17 00:00:00 2001 From: Ruinan Liu Date: Thu, 6 Aug 2026 20:40:59 +0000 Subject: [PATCH 6/7] REsolving comments --- internal/flowcontrol/inputrevbuffer_test.go | 118 ++++++++++++++++++++ 1 file changed, 118 insertions(+) diff --git a/internal/flowcontrol/inputrevbuffer_test.go b/internal/flowcontrol/inputrevbuffer_test.go index 188e8fdb..1f571b92 100644 --- a/internal/flowcontrol/inputrevbuffer_test.go +++ b/internal/flowcontrol/inputrevbuffer_test.go @@ -2,6 +2,7 @@ package flowcontrol import ( "context" + "errors" "fmt" "math/rand/v2" "strconv" @@ -11,6 +12,8 @@ import ( "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" @@ -493,6 +496,121 @@ func TestInputRevisionWriteBufferBackoffStaysBaseForIsolatedComposition(t *testi 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) From cad9a1f2230c5daa9121ae524258dedb08f6a228 Mon Sep 17 00:00:00 2001 From: Ruinan Liu Date: Fri, 7 Aug 2026 18:26:11 +0000 Subject: [PATCH 7/7] Trigger CI