From 2b5c9d96e30bf98c919d13591d0fecf9fa523a98 Mon Sep 17 00:00:00 2001 From: cawthorne Date: Wed, 7 Oct 2026 00:58:43 +0100 Subject: [PATCH] fix(chipingress): empty PublishBatch results with nil error counts as delivered MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A nil-error PublishBatch response carrying no per-event results is reported to every event's callback as ErrCodeResultsMismatch since dc8a1dc3 removed the len(resp.Results) > 0 guard in sendBatch. Servers that accept a batch without per-event detail do exist: local-cre chip-router v1.x forwards every event to Kafka but returns an empty results array. With the durable emitter on top this is not a cosmetic mismatch: every event is treated as undelivered and retransmitted after 60s, forever — the ack never comes, so the node re-sends its whole event stream endlessly. On a chainlink v2.67.2-rc.0 node against chip-router v1.0.1 this produced 301k 'failed to deliver event. Relying on retransmit.' warnings on a single node within minutes of startup, a ~20x duplicated Kafka topic stream, and metric-event delivery latency gaps long enough to break downstream freshness gates. Restore the pre-dc8a1dc3 semantics: a nil-error response with an empty results array resolves every callback with nil (delivered); a non-empty array is still correlated per event, and a partial array still mismatches for the missing tail (unchanged, covered by the existing test). Adds a regression test for the zero-results case. --- pkg/chipingress/batch/client.go | 13 ++++++++-- pkg/chipingress/batch/client_test.go | 37 ++++++++++++++++++++++++++++ 2 files changed, 48 insertions(+), 2 deletions(-) diff --git a/pkg/chipingress/batch/client.go b/pkg/chipingress/batch/client.go index 3313a301b7..d99ff158d6 100644 --- a/pkg/chipingress/batch/client.go +++ b/pkg/chipingress/batch/client.go @@ -355,10 +355,19 @@ func (b *Client) sendBatch(ctx context.Context, messages []*messageWithCallback) if err != nil { b.log.Errorw("failed to publish batch", "error", err) b.completeBatchCallbacks(batchMessages, err) - } else if !b.transactionEnabled { - // always call, even when no results + } else if !b.transactionEnabled && resp != nil && len(resp.Results) > 0 { + // A nil-error response with a non-empty results array is correlated + // per event below; partial arrays still mismatch for the missing tail. b.completeBatchCallbacksFromResults(batchMessages, resp.Results) } else { + // A nil-error response with NO per-event results means the server + // accepted the batch without per-event detail — e.g. local-cre + // chip-router v1.x forwards every event to Kafka but returns an + // empty results array. Reporting RESULTS_MISMATCH here makes the + // durable emitter retransmit every event forever (it never sees an + // ack), duplicating the whole stream and flooding the topic, while + // the events were in fact delivered. Treat empty results as + // delivered. b.completeBatchCallbacks(batchMessages, nil) } } diff --git a/pkg/chipingress/batch/client_test.go b/pkg/chipingress/batch/client_test.go index a93b58dd8e..27558ccf04 100644 --- a/pkg/chipingress/batch/client_test.go +++ b/pkg/chipingress/batch/client_test.go @@ -2138,6 +2138,43 @@ func TestTransactionEnabledEdgeCases(t *testing.T) { assert.Contains(t, results["e3"].Error(), "server returned 1 results for 3 events") }) + t.Run("no results with nil error counts as delivered, no retransmit", func(t *testing.T) { + // Servers that accept a batch without per-event detail (e.g. local-cre + // chip-router v1.x forwards everything to Kafka but returns an empty + // results array) must not trip RESULTS_MISMATCH: the durable emitter + // would retransmit every event forever, duplicating the whole stream. + mockClient := mocks.NewClient(t) + mockClient.EXPECT().Close().Return(nil).Maybe() + mockClient. + On("PublishBatch", mock.Anything, mock.Anything). + Return(&chipingress.PublishResponse{ + Results: []*chipingress.PublishResult{}, + }, nil) + + client, err := NewBatchClient(mockClient, WithTransactionEnabled(false)) + require.NoError(t, err) + + var mu sync.Mutex + results := make(map[string]error) + messages := []*messageWithCallback{ + {event: &chipingress.CloudEventPb{Id: "e1", Source: "s", SpecVersion: "1.0", Type: "t"}, callback: func(err error) { mu.Lock(); results["e1"] = err; mu.Unlock() }}, + {event: &chipingress.CloudEventPb{Id: "e2", Source: "s", SpecVersion: "1.0", Type: "t"}, callback: func(err error) { mu.Lock(); results["e2"] = err; mu.Unlock() }}, + {event: &chipingress.CloudEventPb{Id: "e3", Source: "s", SpecVersion: "1.0", Type: "t"}, callback: func(err error) { mu.Lock(); results["e3"] = err; mu.Unlock() }}, + } + + client.sendBatch(t.Context(), messages) + // Wait for send goroutine to finish (acquire+release semaphore slot). + client.maxConcurrentSends <- struct{}{} + <-client.maxConcurrentSends + client.callbackWg.Wait() + + mu.Lock() + defer mu.Unlock() + require.NoError(t, results["e1"]) + require.NoError(t, results["e2"]) + require.NoError(t, results["e3"]) + }) + t.Run("more results than messages does not panic", func(t *testing.T) { mockClient := mocks.NewClient(t) mockClient.EXPECT().Close().Return(nil).Maybe()