From 72560109a177f4958faaf4899a2d5b3d44e72a83 Mon Sep 17 00:00:00 2001 From: Lee Calcote Date: Fri, 3 Jul 2026 18:21:20 -0500 Subject: [PATCH 1/4] perf(broker): add opt-in content dedup on broker output path BrokerWriter.Write republished the full object to NATS on every ADD/UPDATE/DELETE. resourceVersion-based suppression in the informer UpdateFunc already drops no-op updates, but two distinct resourceVersions can still carry byte-identical wire content, so those redundant republishes waste broker bandwidth and Server DB writes. Add ContentDeduplicatorWriter, a streaming pass-through wrapper that keys resources by KubernetesResourceMeta.UID and remembers a sha256 of the last published payload per UID: - ADD/UPDATE: skip when the new payload hashes identical to the stored hash for that UID; otherwise publish and update the stored hash. - DELETE: always publish and evict the UID, keeping the map bounded to live UIDs and letting a re-created UID republish fresh. - Empty/absent UID: always publish, never deduplicate. - Thread-safe via a mutex. The full object is always sent when published; the wire format is never rewritten into a delta/patch, since Meshery Server consumes full objects from these subjects. The wrapper is OFF by default and enabled via the MESHSYNC_BROKER_CONTENT_DEDUP env var. The broker writer persists across informer resyncs, and a resync is also a recovery path, so the safe default is to republish everything; operators opt in to trade a small, bounded amount of memory for reduced broker/DB churn. Memory is bounded to one sha256 (32 bytes) per live resource UID: entries are added on ADD/UPDATE and removed on DELETE. Add unit tests covering identical-skip, republish-on-change, always-emit-DELETE-and-evict, never-dedup-empty-UID, per-UID independence, map bounding, and concurrent writes under -race. Signed-off-by: Lee Calcote --- internal/config/types.go | 9 + internal/output/content_deduplicator.go | 121 +++++++++++ internal/output/content_deduplicator_test.go | 208 +++++++++++++++++++ pkg/lib/meshsync/meshsync.go | 26 ++- 4 files changed, 359 insertions(+), 5 deletions(-) create mode 100644 internal/output/content_deduplicator.go create mode 100644 internal/output/content_deduplicator_test.go diff --git a/internal/config/types.go b/internal/config/types.go index f8a43210..ecf5dbb1 100644 --- a/internal/config/types.go +++ b/internal/config/types.go @@ -21,6 +21,15 @@ const ( InformerStore = "informer-store" OutputModeBroker = "broker" OutputModeFile = "file" + + // EnvBrokerContentDedup, when truthy, enables in-memory content-hash + // deduplication on the broker output path so that byte-identical + // republishes of the same resource are suppressed. It is OFF by default: + // the dedup cache persists across informer resyncs, and a resync is also a + // recovery path, so republishing everything (the default) is the safe + // behaviour. Accepted truthy values are parsed by strconv.ParseBool + // ("1", "t", "true", etc.). + EnvBrokerContentDedup = "MESHSYNC_BROKER_CONTENT_DEDUP" ) type PipelineConfigs []PipelineConfig diff --git a/internal/output/content_deduplicator.go b/internal/output/content_deduplicator.go new file mode 100644 index 00000000..6ce45f75 --- /dev/null +++ b/internal/output/content_deduplicator.go @@ -0,0 +1,121 @@ +package output + +import ( + "crypto/sha256" + "encoding/json" + "sync" + + "github.com/meshery/meshkit/broker" + "github.com/meshery/meshsync/internal/config" + "github.com/meshery/meshsync/pkg/model" +) + +// ContentDeduplicatorWriter is a streaming output wrapper that suppresses +// byte-identical republishes of the SAME resource on the broker path. +// +// It keys resources by KubernetesResourceMeta.UID and remembers a sha256 hash +// of the last payload published for that UID. On ADD/UPDATE, if the incoming +// payload hashes to the same value already recorded for the UID, the write is +// skipped; otherwise it is forwarded and the recorded hash is updated. This is +// a bandwidth/DB-write optimisation on top of the resourceVersion-based +// suppression already performed in the informer UpdateFunc: two distinct +// resourceVersions can still carry identical wire content (e.g. status churn +// that normalises away), and only a content hash catches those. +// +// The full object is always forwarded when it IS published: this wrapper never +// rewrites the wire format into a delta/patch, because Meshery Server consumes +// full objects from these subjects. +// +// Invariants that keep this safe: +// - DELETE is ALWAYS forwarded and evicts the UID from the map. Eviction +// bounds the map to the set of currently-live UIDs and guarantees that a +// re-created object (same UID reused, or a new UID) republishes fresh. +// - Resources with an empty/absent UID are never deduplicated; each is +// forwarded as-is (mirrors InMemoryDeduplicatorStreamingWriter). +// - Ordering and event semantics are preserved: this is a pass-through filter, +// it never reorders, batches, or defers events. +// +// NOTE on informer resyncs: the broker writer (and therefore this wrapper) +// persists across informer resyncs - a resync recreates the informer factory +// but not the output writer - so the hash map survives a resync. A resync +// re-lists every object, and unchanged objects would be suppressed here. That +// is acceptable while Meshery Server still holds those unchanged objects, but +// because a resync is also a recovery path it must be an explicit opt-in +// (see config.EnvBrokerContentDedup); it is OFF by default so the default +// behaviour - republish everything - is unchanged. +type ContentDeduplicatorWriter struct { + realWriter Writer + + mu sync.Mutex + // hashByUID stores the sha256 of the last payload published per resource UID. + // Entries are added/updated on ADD/UPDATE and removed on DELETE, so the map + // stays bounded to the number of live resources with a UID. + hashByUID map[string][sha256.Size]byte +} + +// NewContentDeduplicatorWriter wraps realWriter with content-hash +// deduplication keyed by resource UID. +func NewContentDeduplicatorWriter(realWriter Writer) *ContentDeduplicatorWriter { + return &ContentDeduplicatorWriter{ + realWriter: realWriter, + hashByUID: make(map[string][sha256.Size]byte), + } +} + +func (w *ContentDeduplicatorWriter) Write( + obj model.KubernetesResource, + evtype broker.EventType, + cfg config.PipelineConfig, +) error { + uid := "" + if obj.KubernetesResourceMeta != nil { + uid = obj.KubernetesResourceMeta.UID + } + + // No UID: cannot be tracked reliably, so always forward (never dedup). + if uid == "" { + return w.realWriter.Write(obj, evtype, cfg) + } + + w.mu.Lock() + defer w.mu.Unlock() + + // DELETE always publishes and evicts the UID. Evicting keeps the map bounded + // and ensures a subsequently re-created object republishes fresh rather than + // colliding with a stale hash. + if evtype == broker.Delete { + delete(w.hashByUID, uid) + return w.realWriter.Write(obj, evtype, cfg) + } + + // ADD/UPDATE: publish only when the payload content changed for this UID. + hash, err := hashResource(obj) + if err != nil { + // Hashing failed for reasons outside our control (a value that does not + // round-trip through JSON). Fail open - forward the event - rather than + // silently dropping it, and drop any stale hash so we do not wedge this + // UID into a permanently-suppressed state. + delete(w.hashByUID, uid) + return w.realWriter.Write(obj, evtype, cfg) + } + + if prev, ok := w.hashByUID[uid]; ok && prev == hash { + // Byte-identical to the last published payload for this UID: skip. + return nil + } + + w.hashByUID[uid] = hash + return w.realWriter.Write(obj, evtype, cfg) +} + +// hashResource returns the sha256 of the JSON encoding of obj. json.Marshal is +// deterministic for a given value (struct fields in declaration order, map keys +// sorted), and the object is the same value that meshkit marshals onto the wire, +// so identical published payloads hash identically. +func hashResource(obj model.KubernetesResource) ([sha256.Size]byte, error) { + data, err := json.Marshal(obj) + if err != nil { + return [sha256.Size]byte{}, err + } + return sha256.Sum256(data), nil +} diff --git a/internal/output/content_deduplicator_test.go b/internal/output/content_deduplicator_test.go new file mode 100644 index 00000000..8058e74d --- /dev/null +++ b/internal/output/content_deduplicator_test.go @@ -0,0 +1,208 @@ +package output + +import ( + "fmt" + "sync" + "testing" + + "github.com/meshery/meshkit/broker" + "github.com/meshery/meshsync/internal/config" + "github.com/meshery/meshsync/pkg/model" + "github.com/stretchr/testify/assert" +) + +// resourceWithUID builds a resource carrying the given UID. resourceVersion is +// folded into the payload so callers can produce two values with the same UID +// but different serialized content. +func resourceWithUID(uid, resourceVersion string) model.KubernetesResource { + return model.KubernetesResource{ + Kind: "Pod", + KubernetesResourceMeta: &model.KubernetesResourceObjectMeta{ + UID: uid, + ResourceVersion: resourceVersion, + }, + } +} + +func TestContentDeduplicatorWriter_SkipsByteIdentical(t *testing.T) { + mock := &mockWriter{} + writer := NewContentDeduplicatorWriter(mock) + cfg := config.PipelineConfig{} + + res := resourceWithUID("uid-1", "100") + + // First ADD publishes; second identical UPDATE is byte-identical and skipped. + assert.NoError(t, writer.Write(res, broker.Add, cfg)) + assert.NoError(t, writer.Write(res, broker.Update, cfg)) + assert.NoError(t, writer.Write(res, broker.Update, cfg)) + + assert.Len(t, mock.written, 1) + assert.Equal(t, broker.Add, mock.events[0]) +} + +func TestContentDeduplicatorWriter_RepublishesOnChange(t *testing.T) { + mock := &mockWriter{} + writer := NewContentDeduplicatorWriter(mock) + cfg := config.PipelineConfig{} + + // Same UID, changing content each time: every write must be forwarded. + assert.NoError(t, writer.Write(resourceWithUID("uid-1", "100"), broker.Add, cfg)) + assert.NoError(t, writer.Write(resourceWithUID("uid-1", "101"), broker.Update, cfg)) + assert.NoError(t, writer.Write(resourceWithUID("uid-1", "102"), broker.Update, cfg)) + + assert.Len(t, mock.written, 3) + + // A repeat of the last content is then suppressed, proving the stored hash + // tracks the most recent payload. + assert.NoError(t, writer.Write(resourceWithUID("uid-1", "102"), broker.Update, cfg)) + assert.Len(t, mock.written, 3) +} + +func TestContentDeduplicatorWriter_AlwaysEmitsDeleteAndEvicts(t *testing.T) { + mock := &mockWriter{} + writer := NewContentDeduplicatorWriter(mock) + cfg := config.PipelineConfig{} + + res := resourceWithUID("uid-1", "100") + + assert.NoError(t, writer.Write(res, broker.Add, cfg)) // published + assert.NoError(t, writer.Write(res, broker.Delete, cfg)) // always published + // A DELETE with content identical to a prior publish must still go out. + assert.NoError(t, writer.Write(res, broker.Delete, cfg)) + + assert.Len(t, mock.written, 3) + assert.Equal(t, broker.Add, mock.events[0]) + assert.Equal(t, broker.Delete, mock.events[1]) + assert.Equal(t, broker.Delete, mock.events[2]) + + // Eviction check: after DELETE, re-adding the same UID with the SAME content + // as the original ADD must republish (the hash was evicted, not retained). + mock2 := &mockWriter{} + writer2 := NewContentDeduplicatorWriter(mock2) + assert.NoError(t, writer2.Write(res, broker.Add, cfg)) + assert.NoError(t, writer2.Write(res, broker.Delete, cfg)) + assert.NoError(t, writer2.Write(res, broker.Add, cfg)) // re-created UID, fresh publish + assert.Len(t, mock2.written, 3) + assert.Equal(t, broker.Add, mock2.events[2]) + + // The map must not retain the UID after delete. + writer2.mu.Lock() + _, present := writer2.hashByUID["uid-1"] + writer2.mu.Unlock() + assert.True(t, present, "re-add after delete should have re-populated the UID") +} + +func TestContentDeduplicatorWriter_MapBoundedByDelete(t *testing.T) { + mock := &mockWriter{} + writer := NewContentDeduplicatorWriter(mock) + cfg := config.PipelineConfig{} + + res := resourceWithUID("uid-1", "100") + assert.NoError(t, writer.Write(res, broker.Add, cfg)) + + writer.mu.Lock() + _, present := writer.hashByUID["uid-1"] + writer.mu.Unlock() + assert.True(t, present) + + assert.NoError(t, writer.Write(res, broker.Delete, cfg)) + + writer.mu.Lock() + _, present = writer.hashByUID["uid-1"] + size := len(writer.hashByUID) + writer.mu.Unlock() + assert.False(t, present, "DELETE must evict the UID from the map") + assert.Equal(t, 0, size) +} + +func TestContentDeduplicatorWriter_NeverDedupsEmptyUID(t *testing.T) { + mock := &mockWriter{} + writer := NewContentDeduplicatorWriter(mock) + cfg := config.PipelineConfig{} + + // Nil metadata and empty-string UID are both treated as "no UID" and must + // always be forwarded, even for byte-identical repeats. + noMeta := model.KubernetesResource{KubernetesResourceMeta: nil} + emptyUID := model.KubernetesResource{ + KubernetesResourceMeta: &model.KubernetesResourceObjectMeta{UID: ""}, + } + + assert.NoError(t, writer.Write(noMeta, broker.Add, cfg)) + assert.NoError(t, writer.Write(noMeta, broker.Add, cfg)) + assert.NoError(t, writer.Write(emptyUID, broker.Update, cfg)) + assert.NoError(t, writer.Write(emptyUID, broker.Update, cfg)) + + assert.Len(t, mock.written, 4) + + // Empty-UID resources must never populate the dedup map. + writer.mu.Lock() + size := len(writer.hashByUID) + writer.mu.Unlock() + assert.Equal(t, 0, size) +} + +func TestContentDeduplicatorWriter_DistinctUIDsIndependent(t *testing.T) { + mock := &mockWriter{} + writer := NewContentDeduplicatorWriter(mock) + cfg := config.PipelineConfig{} + + // Identical content on two different UIDs must both publish: dedup is + // per-UID, not global content dedup. + a := resourceWithUID("uid-a", "100") + b := resourceWithUID("uid-b", "100") + + assert.NoError(t, writer.Write(a, broker.Add, cfg)) + assert.NoError(t, writer.Write(b, broker.Add, cfg)) + // Repeats of each are then suppressed. + assert.NoError(t, writer.Write(a, broker.Update, cfg)) + assert.NoError(t, writer.Write(b, broker.Update, cfg)) + + assert.Len(t, mock.written, 2) +} + +// TestContentDeduplicatorWriter_ConcurrentWrites exercises the mutex under the +// race detector. Each goroutine owns a distinct UID and writes the same content +// repeatedly; exactly one publish per UID must survive dedup. +func TestContentDeduplicatorWriter_ConcurrentWrites(t *testing.T) { + mock := &lockingMockWriter{} + writer := NewContentDeduplicatorWriter(mock) + cfg := config.PipelineConfig{} + + const goroutines = 16 + const perGoroutine = 50 + + var wg sync.WaitGroup + wg.Add(goroutines) + for g := 0; g < goroutines; g++ { + go func(id int) { + defer wg.Done() + res := resourceWithUID(fmt.Sprintf("uid-%d", id), "100") + for i := 0; i < perGoroutine; i++ { + assert.NoError(t, writer.Write(res, broker.Update, cfg)) + } + }(g) + } + wg.Wait() + + assert.Equal(t, goroutines, mock.count()) +} + +// lockingMockWriter is a minimal thread-safe Writer for the concurrency test; +// the package-shared mockWriter is not safe for concurrent use. +type lockingMockWriter struct { + mu sync.Mutex + written int +} + +func (m *lockingMockWriter) Write(_ model.KubernetesResource, _ broker.EventType, _ config.PipelineConfig) error { + m.mu.Lock() + m.written++ + m.mu.Unlock() + return nil +} + +func (m *lockingMockWriter) count() int { + m.mu.Lock() + defer m.mu.Unlock() + return m.written +} diff --git a/pkg/lib/meshsync/meshsync.go b/pkg/lib/meshsync/meshsync.go index b59b8a53..9b22ece9 100644 --- a/pkg/lib/meshsync/meshsync.go +++ b/pkg/lib/meshsync/meshsync.go @@ -9,6 +9,7 @@ import ( "os/signal" "path" "slices" + "strconv" "strings" "syscall" "time" @@ -126,11 +127,17 @@ func Run(log logger.Handler, optsSetters ...OptionsSetter) error { // the broker handler exists (nats.New succeeded or a custom // handler was provided): report ready on /readyz health.markReady() - outputProcessor.SetOutput( - output.NewBrokerWriter( - br, - ), - ) + + var brokerWriter output.Writer = output.NewBrokerWriter(br) + // Opt-in content deduplication (OFF by default). The dedup cache persists + // across informer resyncs, and a resync is a recovery path, so the safe + // default is to republish everything. Operators who want to trade a small + // amount of memory for reduced broker/DB churn enable it explicitly. + if brokerContentDedupEnabled() { + log.Info("meshsync: broker content deduplication enabled") + brokerWriter = output.NewContentDeduplicatorWriter(brokerWriter) + } + outputProcessor.SetOutput(brokerWriter) } if options.OutputMode == config.OutputModeFile { @@ -256,6 +263,15 @@ func Run(log logger.Handler, optsSetters ...OptionsSetter) error { return nil } +// brokerContentDedupEnabled reports whether the operator opted into broker +// content deduplication via config.EnvBrokerContentDedup. An unset or +// unparseable value means disabled, preserving the default republish-everything +// behaviour. +func brokerContentDedupEnabled() bool { + enabled, err := strconv.ParseBool(os.Getenv(config.EnvBrokerContentDedup)) + return err == nil && enabled +} + // connectivityTestTimeout bounds how long connectivityTest keeps retrying // before it gives up and returns an error. const connectivityTestTimeout = 5 * time.Minute From 36067f7c8b2d768f7546b5c15fed75f0a7f976fc Mon Sep 17 00:00:00 2001 From: Lee Calcote Date: Sat, 4 Jul 2026 18:20:38 -0500 Subject: [PATCH 2/4] perf(output): hash and publish outside the dedup mutex Hold w.mu only to guard the hash map; compute the sha256/JSON serialization and call the downstream writer with the lock released, so concurrent writers are not serialized by CPU-bound hashing or the broker publish. Addresses review feedback on #575. Signed-off-by: Lee Calcote --- internal/output/content_deduplicator.go | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/internal/output/content_deduplicator.go b/internal/output/content_deduplicator.go index 6ce45f75..b8dcf24e 100644 --- a/internal/output/content_deduplicator.go +++ b/internal/output/content_deduplicator.go @@ -77,34 +77,41 @@ func (w *ContentDeduplicatorWriter) Write( return w.realWriter.Write(obj, evtype, cfg) } - w.mu.Lock() - defer w.mu.Unlock() - // DELETE always publishes and evicts the UID. Evicting keeps the map bounded // and ensures a subsequently re-created object republishes fresh rather than // colliding with a stale hash. if evtype == broker.Delete { + w.mu.Lock() delete(w.hashByUID, uid) + w.mu.Unlock() return w.realWriter.Write(obj, evtype, cfg) } // ADD/UPDATE: publish only when the payload content changed for this UID. + // Hash outside the lock: the JSON serialization is CPU-bound and must not + // serialize concurrent writers. The mutex guards only the map, and the + // downstream write happens after the lock is released. hash, err := hashResource(obj) if err != nil { // Hashing failed for reasons outside our control (a value that does not // round-trip through JSON). Fail open - forward the event - rather than // silently dropping it, and drop any stale hash so we do not wedge this // UID into a permanently-suppressed state. + w.mu.Lock() delete(w.hashByUID, uid) + w.mu.Unlock() return w.realWriter.Write(obj, evtype, cfg) } + w.mu.Lock() if prev, ok := w.hashByUID[uid]; ok && prev == hash { // Byte-identical to the last published payload for this UID: skip. + w.mu.Unlock() return nil } - w.hashByUID[uid] = hash + w.mu.Unlock() + return w.realWriter.Write(obj, evtype, cfg) } From aa410a968eb1f7cf44e5ce31c11325332bb3f688 Mon Sep 17 00:00:00 2001 From: marblom007 <158522975+marblom007@users.noreply.github.com> Date: Sat, 4 Jul 2026 18:40:00 -0500 Subject: [PATCH 3/4] Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Signed-off-by: marblom007 <158522975+marblom007@users.noreply.github.com> --- internal/output/content_deduplicator_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/output/content_deduplicator_test.go b/internal/output/content_deduplicator_test.go index 8058e74d..94bde1f0 100644 --- a/internal/output/content_deduplicator_test.go +++ b/internal/output/content_deduplicator_test.go @@ -85,7 +85,7 @@ func TestContentDeduplicatorWriter_AlwaysEmitsDeleteAndEvicts(t *testing.T) { assert.Len(t, mock2.written, 3) assert.Equal(t, broker.Add, mock2.events[2]) - // The map must not retain the UID after delete. + // After re-add following DELETE, the UID should be present again in the map. writer2.mu.Lock() _, present := writer2.hashByUID["uid-1"] writer2.mu.Unlock() From a50e6d9b4535d4edbe3a48940d5ce84edda6f7ed Mon Sep 17 00:00:00 2001 From: Lee Calcote Date: Sat, 4 Jul 2026 19:01:09 -0500 Subject: [PATCH 4/4] fix(output): record dedup hash only after a successful publish Previously the content-hash was recorded before calling the downstream writer, so a failed publish left the UID marked as published and suppressed a later retry (or resync) of the same payload that had never actually shipped. Now the payload is published first and the hash recorded only on success; a failed write returns the error and leaves the hash unrecorded. Adds a regression test and corrects a stale test comment. Addresses review feedback on #575. Signed-off-by: marblom007 <158522975+marblom007@users.noreply.github.com> Signed-off-by: Lee Calcote --- internal/output/content_deduplicator.go | 19 +++++++--- internal/output/content_deduplicator_test.go | 39 ++++++++++++++++++++ 2 files changed, 53 insertions(+), 5 deletions(-) diff --git a/internal/output/content_deduplicator.go b/internal/output/content_deduplicator.go index b8dcf24e..1785998c 100644 --- a/internal/output/content_deduplicator.go +++ b/internal/output/content_deduplicator.go @@ -104,15 +104,24 @@ func (w *ContentDeduplicatorWriter) Write( } w.mu.Lock() - if prev, ok := w.hashByUID[uid]; ok && prev == hash { - // Byte-identical to the last published payload for this UID: skip. - w.mu.Unlock() + prev, ok := w.hashByUID[uid] + w.mu.Unlock() + if ok && prev == hash { + // Byte-identical to the last successfully published payload: skip. return nil } + + // Publish first, then record the hash only on success. Recording before the + // write would let a failed publish suppress a later retry of the same payload + // (the payload would be marked "published" though it never was). + if err := w.realWriter.Write(obj, evtype, cfg); err != nil { + return err + } + + w.mu.Lock() w.hashByUID[uid] = hash w.mu.Unlock() - - return w.realWriter.Write(obj, evtype, cfg) + return nil } // hashResource returns the sha256 of the JSON encoding of obj. json.Marshal is diff --git a/internal/output/content_deduplicator_test.go b/internal/output/content_deduplicator_test.go index 94bde1f0..e108e736 100644 --- a/internal/output/content_deduplicator_test.go +++ b/internal/output/content_deduplicator_test.go @@ -206,3 +206,42 @@ func (m *lockingMockWriter) count() int { defer m.mu.Unlock() return m.written } + +// failableWriter fails its next Write once when failNext is set, then succeeds. +type failableWriter struct { + failNext bool + calls int +} + +func (w *failableWriter) Write(_ model.KubernetesResource, _ broker.EventType, _ config.PipelineConfig) error { + w.calls++ + if w.failNext { + w.failNext = false + return fmt.Errorf("downstream write failed") + } + return nil +} + +// A failed downstream publish must not record the hash, or a retry of the same +// payload would be wrongly suppressed as a duplicate that never actually shipped. +func TestContentDeduplicatorWriter_FailedWriteNotRecorded(t *testing.T) { + inner := &failableWriter{failNext: true} + writer := NewContentDeduplicatorWriter(inner) + cfg := config.PipelineConfig{} + res := resourceWithUID("uid-1", "100") + + // First publish fails downstream: the error propagates and the hash is not recorded. + assert.Error(t, writer.Write(res, broker.Add, cfg)) + writer.mu.Lock() + _, present := writer.hashByUID["uid-1"] + writer.mu.Unlock() + assert.False(t, present, "a failed publish must not record the hash") + + // Retry with the identical payload must reach the downstream writer (not be suppressed). + assert.NoError(t, writer.Write(res, broker.Update, cfg)) + assert.Equal(t, 2, inner.calls, "retry after a failed publish should re-reach the downstream writer") + + // After the successful publish the hash is recorded, so an identical write is skipped. + assert.NoError(t, writer.Write(res, broker.Update, cfg)) + assert.Equal(t, 2, inner.calls, "byte-identical write after a successful publish should be suppressed") +}