Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 11 additions & 2 deletions pkg/chipingress/batch/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
Expand Down
37 changes: 37 additions & 0 deletions pkg/chipingress/batch/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
Loading