From e15d9e572165879bdfd22f73ee5849f5699ecc84 Mon Sep 17 00:00:00 2001 From: cawthorne Date: Wed, 7 Oct 2026 01:58:08 +0100 Subject: [PATCH 1/5] fix(chip-router): populate per-event PublishBatch results MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PublishBatch always returned an empty PublishResponse regardless of the batch's outcome. The node-side batch client (chainlink-common/pkg/chipingress/batch) requires one PublishResult per event to resolve delivery: since chainlink-common#2326 an empty results array is reported to every event's callback as ErrCodeResultsMismatch, so a chainlink v2.67 node's durable emitter against this router treated every event as undelivered and retransmitted the whole stream every 60s, forever — measured at 301k 'failed to deliver event. Relying on retransmit.' warnings on a single node within minutes of startup, a ~20x duplicated Kafka topic, and metric delivery-latency gaps long enough to break downstream freshness gates. - Ack each event (PublishResult{EventId}, nil error) when the batch was handed to at least one subscriber; the events were forwarded in that case, so the caller can resolve delivery and stop retransmitting. - Report a per-event PublishError when no subscriber accepted the batch, so the caller retains and retries (at-least-once preserved). - Return Unavailable when no subscribers are registered instead of an empty success — the previous behavior acked events that went nowhere, letting a durable store delete them undelivered (the data-loss shape flagged by the audit behind chainlink-common#2326). Bumps the module's chainlink-common/pkg/chipingress pin (Dec 2025 -> Sep 2026) — the old pin's pb predates PublishResult.Error. --- .../chiprouter/cmd/chip-router/main.go | 30 +++++- .../chiprouter/cmd/chip-router/main_test.go | 97 +++++++++++++++++++ framework/components/chiprouter/go.mod | 6 +- framework/components/chiprouter/go.sum | 8 +- 4 files changed, 132 insertions(+), 9 deletions(-) create mode 100644 framework/components/chiprouter/cmd/chip-router/main_test.go diff --git a/framework/components/chiprouter/cmd/chip-router/main.go b/framework/components/chiprouter/cmd/chip-router/main.go index 516d5b71f..2cb04dc61 100644 --- a/framework/components/chiprouter/cmd/chip-router/main.go +++ b/framework/components/chiprouter/cmd/chip-router/main.go @@ -11,6 +11,7 @@ import ( "os/signal" "strings" "sync" + "sync/atomic" "syscall" "time" @@ -19,7 +20,9 @@ import ( chippb "github.com/smartcontractkit/chainlink-common/pkg/chipingress/pb" "golang.org/x/sync/errgroup" "google.golang.org/grpc" + "google.golang.org/grpc/codes" "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/status" "github.com/smartcontractkit/chainlink-testing-framework/framework" ) @@ -164,9 +167,12 @@ func (r *router) Publish(ctx context.Context, event *cepb.CloudEvent) (*chippb.P func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch) (*chippb.PublishResponse, error) { snapshot := r.snapshotSubscribers() if len(snapshot) == 0 { - return &chippb.PublishResponse{}, nil + // Nothing will carry these events. Reporting an empty success here + // would let the caller's durable store delete them undelivered. + return nil, status.Error(codes.Unavailable, "no subscribers registered") } + var forwarded atomic.Bool var group errgroup.Group group.SetLimit(forwardParallel) for _, sub := range snapshot { @@ -177,13 +183,33 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch _, err := sub.client.PublishBatch(forwardCtx, batch) if err != nil { framework.L.Error().Msgf("chip router failed to forward batch to subscriber id=%s name=%s endpoint=%s err=%v", sub.id, sub.name, sub.endpoint, err) + } else { + forwarded.Store(true) } framework.L.Debug().Msgf("chip router forwarded batch to subscriber id=%s", sub.id) return nil }) } _ = group.Wait() - return &chippb.PublishResponse{}, nil + + // PublishBatch callers (the node-side batch client driving the durable + // emitter) require one result per event to resolve delivery: an empty + // results array reads as RESULTS_MISMATCH and retransmits the whole + // stream forever. A batch the router handed to at least one subscriber + // counts as accepted for every event; a batch no subscriber accepted is + // reported failed per event so the caller retains and retries. + results := make([]*chippb.PublishResult, 0, len(batch.GetEvents())) + for _, ev := range batch.GetEvents() { + result := &chippb.PublishResult{EventId: ev.GetId()} + if !forwarded.Load() { + result.Error = &chippb.PublishError{ + ErrorCode: chippb.PublishErrorCode_PUBLISH_ERROR_CODE_UNKNOWN, + Reason: "chip router could not forward the batch to any subscriber", + } + } + results = append(results, result) + } + return &chippb.PublishResponse{Results: results}, nil } func (r *router) handleHealth(w http.ResponseWriter, req *http.Request) { diff --git a/framework/components/chiprouter/cmd/chip-router/main_test.go b/framework/components/chiprouter/cmd/chip-router/main_test.go new file mode 100644 index 000000000..25dc6957d --- /dev/null +++ b/framework/components/chiprouter/cmd/chip-router/main_test.go @@ -0,0 +1,97 @@ +package main + +import ( + "context" + "sync" + "testing" + + cepb "github.com/cloudevents/sdk-go/binding/format/protobuf/v2/pb" + chippb "github.com/smartcontractkit/chainlink-common/pkg/chipingress/pb" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +type stubPublishClient struct { + mu sync.Mutex + batches int + err error +} + +func (s *stubPublishClient) Publish(context.Context, *cepb.CloudEvent, ...grpc.CallOption) (*chippb.PublishResponse, error) { + return &chippb.PublishResponse{}, nil +} + +func (s *stubPublishClient) PublishBatch(context.Context, *chippb.CloudEventBatch, ...grpc.CallOption) (*chippb.PublishResponse, error) { + s.mu.Lock() + defer s.mu.Unlock() + s.batches++ + return &chippb.PublishResponse{}, s.err +} + +func (s *stubPublishClient) batchCount() int { + s.mu.Lock() + defer s.mu.Unlock() + return s.batches +} + +func batchOf(ids ...string) *chippb.CloudEventBatch { + events := make([]*cepb.CloudEvent, 0, len(ids)) + for _, id := range ids { + events = append(events, &cepb.CloudEvent{Id: id, Source: "s", SpecVersion: "1.0", Type: "t"}) + } + return &chippb.CloudEventBatch{Events: events} +} + +func TestPublishBatchPopulatesPerEventResults(t *testing.T) { + r := &router{subscribers: map[string]*subscriber{ + "sink": {id: "sink", client: &stubPublishClient{}}, + }} + + resp, err := r.PublishBatch(t.Context(), batchOf("e1", "e2", "e3")) + if err != nil { + t.Fatalf("PublishBatch: %v", err) + } + if len(resp.Results) != 3 { + t.Fatalf("want 3 results, got %d", len(resp.Results)) + } + for i, want := range []string{"e1", "e2", "e3"} { + if resp.Results[i].EventId != want { + t.Errorf("results[%d].EventId = %q, want %q", i, resp.Results[i].EventId, want) + } + if resp.Results[i].Error != nil { + t.Errorf("results[%d].Error = %v, want nil", i, resp.Results[i].Error) + } + } +} + +func TestPublishBatchNoSubscribersIsUnavailable(t *testing.T) { + r := &router{subscribers: map[string]*subscriber{}} + + _, err := r.PublishBatch(t.Context(), batchOf("e1")) + if err == nil { + t.Fatal("want error with no subscribers, got nil") + } + if status.Code(err) != codes.Unavailable { + t.Errorf("want Unavailable, got %v", status.Code(err)) + } +} + +func TestPublishBatchAllForwardsFailedReportsPerEventErrors(t *testing.T) { + r := &router{subscribers: map[string]*subscriber{ + "sink": {id: "sink", client: &stubPublishClient{err: context.DeadlineExceeded}}, + }} + + resp, err := r.PublishBatch(t.Context(), batchOf("e1", "e2")) + if err != nil { + t.Fatalf("PublishBatch: %v", err) + } + if len(resp.Results) != 2 { + t.Fatalf("want 2 results, got %d", len(resp.Results)) + } + for i, res := range resp.Results { + if res.Error == nil { + t.Errorf("results[%d].Error = nil, want a forward-failure error", i) + } + } +} diff --git a/framework/components/chiprouter/go.mod b/framework/components/chiprouter/go.mod index c31193136..18347708e 100644 --- a/framework/components/chiprouter/go.mod +++ b/framework/components/chiprouter/go.mod @@ -1,6 +1,6 @@ module github.com/smartcontractkit/chainlink-testing-framework/framework/components/chiprouter -go 1.26.5 +go 1.26.6 replace github.com/smartcontractkit/chainlink-testing-framework/framework => ../../ @@ -9,11 +9,11 @@ require ( github.com/google/uuid v1.6.0 github.com/moby/moby/api v1.54.1 github.com/pkg/errors v0.9.1 - github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20251211140724-319861e514c4 + github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20260915184316-2730f1867c92 github.com/smartcontractkit/chainlink-testing-framework/framework v0.15.8 github.com/testcontainers/testcontainers-go v0.42.0 golang.org/x/sync v0.22.0 - google.golang.org/grpc v1.83.1 + google.golang.org/grpc v1.83.2 ) require ( diff --git a/framework/components/chiprouter/go.sum b/framework/components/chiprouter/go.sum index 52871d8a1..38bf00c87 100644 --- a/framework/components/chiprouter/go.sum +++ b/framework/components/chiprouter/go.sum @@ -205,8 +205,8 @@ github.com/shirou/gopsutil/v4 v4.26.3 h1:2ESdQt90yU3oXF/CdOlRCJxrP+Am1aBYubTMTfx github.com/shirou/gopsutil/v4 v4.26.3/go.mod h1:LZ6ewCSkBqUpvSOf+LsTGnRinC6iaNUNMGBtDkJBaLQ= github.com/sirupsen/logrus v1.9.4 h1:TsZE7l11zFCLZnZ+teH4Umoq5BhEIfIzfRDZ1Uzql2w= github.com/sirupsen/logrus v1.9.4/go.mod h1:ftWc9WdOfJ0a92nsE2jF5u5ZwH8Bv2zdeOC42RjbV2g= -github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20251211140724-319861e514c4 h1:NOUsjsMzNecbjiPWUQGlRSRAutEvCFrqqyETDJeh5q4= -github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20251211140724-319861e514c4/go.mod h1:Zpvul9sTcZNAZOVzt5vBl1XZGNvQebFpnpn3/KOQvOQ= +github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20260915184316-2730f1867c92 h1:Z1OIGukDhJD+Gknr5G07Vtml24E7BwKaIQoWzV/SdqA= +github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20260915184316-2730f1867c92/go.mod h1:Scyk7O80fde5rJF3Ty3fEE3xiS6j84eWbu8/2gl3pIY= github.com/spf13/pflag v1.0.6 h1:jFzHGLGAlb3ruxLB8MhbI6A8+AQX/2eW4qeyNZXNp2o= github.com/spf13/pflag v1.0.6/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4= @@ -278,8 +278,8 @@ gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= google.golang.org/genproto/googleapis/rpc v0.0.0-20260819154853-08b0e4226688 h1:cYNAzI2sUwhmCcoj9TxvihSrqsxt6uIkj3rDRhSDmW4= google.golang.org/genproto/googleapis/rpc v0.0.0-20260819154853-08b0e4226688/go.mod h1:DjtHYE8FKJLivXcBEjGwndXfIC23G0VpXiXKqG179uA= -google.golang.org/grpc v1.83.1 h1:HIO0+BEtBP6soyqvqC8sNUjZ7bTs+0hFQuFF+RAy++Y= -google.golang.org/grpc v1.83.1/go.mod h1:kDyl6SKsiHKt0uylY5gtn5cEjkrIOhQOGDgIc4JGwzQ= +google.golang.org/grpc v1.83.2 h1:EManeRomTObA0BU7I8vXgg/78uE5MJ9M8B39EX2WscU= +google.golang.org/grpc v1.83.2/go.mod h1:YPI1hK3kDked6iHvgX3tR0y+nX/qpMFKhPgFsokw1S8= google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc= google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/evanphx/json-patch.v4 v4.12.0 h1:n6jtcsulIzXPJaxegRbvFNNrZDjbij7ny3gmSPG+6V4= From 066433d0105d8ebb7ccd2a5fd5d49ad6ee8924b0 Mon Sep 17 00:00:00 2001 From: cawthorne Date: Wed, 7 Oct 2026 10:51:12 +0100 Subject: [PATCH 2/5] fix(chip-router): aggregate downstream per-event outcomes; fail the RPC when nothing accepts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses the review findings: - A nil downstream RPC error does not mean every event was accepted — the protocol supports per-event rejections in a successful response. The router now aggregates the subscribers' per-event results and passes a rejection through only while no subscriber has accepted the event. Subscribers that accept without per-event detail (nil-error response, empty results array) count as accepting the whole batch — the router trusts a subscriber's RPC success at the forwarding boundary, the same as Publish does. - When no subscriber accepts the batch, return Unavailable instead of a nil-error response with per-event errors: transactional callers resolve the whole batch from the RPC outcome alone, so per-event errors would still acknowledge everything. An RPC error makes every caller mode retain and retry. --- .../chiprouter/cmd/chip-router/main.go | 63 +++++++++++---- .../chiprouter/cmd/chip-router/main_test.go | 78 ++++++++++++++++--- 2 files changed, 115 insertions(+), 26 deletions(-) diff --git a/framework/components/chiprouter/cmd/chip-router/main.go b/framework/components/chiprouter/cmd/chip-router/main.go index 2cb04dc61..be6772163 100644 --- a/framework/components/chiprouter/cmd/chip-router/main.go +++ b/framework/components/chiprouter/cmd/chip-router/main.go @@ -11,7 +11,6 @@ import ( "os/signal" "strings" "sync" - "sync/atomic" "syscall" "time" @@ -172,7 +171,19 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch return nil, status.Error(codes.Unavailable, "no subscribers registered") } - var forwarded atomic.Bool + // Aggregate the downstream outcomes so the per-event results reflect what + // actually happened to each event: + // - anyAccepted: at least one subscriber's PublishBatch RPC succeeded. + // - detail: merged per-event results from the subscribers that provide + // them. An event is reported rejected only if every detail-providing + // subscriber rejected it; subscribers that accept without per-event + // detail (a nil-error response with an empty results array) count as + // accepting the whole batch — the router trusts a subscriber's RPC + // success at the forwarding boundary, the same as Publish does. + var mu sync.Mutex + anyAccepted := false + var detail []*chippb.PublishResult + var group errgroup.Group group.SetLimit(forwardParallel) for _, sub := range snapshot { @@ -180,32 +191,50 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch framework.L.Debug().Msgf("chip router forwarding batch to subscriber id=%s name=%s endpoint=%s", sub.id, sub.name, sub.endpoint) forwardCtx, cancel := context.WithTimeout(ctx, forwardTimeout) defer cancel() - _, err := sub.client.PublishBatch(forwardCtx, batch) + resp, err := sub.client.PublishBatch(forwardCtx, batch) if err != nil { framework.L.Error().Msgf("chip router failed to forward batch to subscriber id=%s name=%s endpoint=%s err=%v", sub.id, sub.name, sub.endpoint, err) - } else { - forwarded.Store(true) + return nil } framework.L.Debug().Msgf("chip router forwarded batch to subscriber id=%s", sub.id) + mu.Lock() + defer mu.Unlock() + anyAccepted = true + if results := resp.GetResults(); len(results) > 0 { + if detail == nil { + detail = results + return nil + } + // Merge: an accepting result overrides a rejecting one for the + // same event index; a rejection sticks only while no subscriber + // has accepted the event. + for i := range detail { + if i >= len(results) { + break + } + if detail[i].GetError() != nil && results[i].GetError() == nil { + detail[i] = results[i] + } + } + } return nil }) } _ = group.Wait() - // PublishBatch callers (the node-side batch client driving the durable - // emitter) require one result per event to resolve delivery: an empty - // results array reads as RESULTS_MISMATCH and retransmits the whole - // stream forever. A batch the router handed to at least one subscriber - // counts as accepted for every event; a batch no subscriber accepted is - // reported failed per event so the caller retains and retries. + // A nil RPC error does not imply the batch was accepted: every forward + // can fail (or a transactional caller resolves the whole batch from the + // RPC outcome alone), so fail the RPC and let the caller retain and + // retry instead of acknowledging with per-event errors. + if !anyAccepted { + return nil, status.Error(codes.Unavailable, fmt.Sprintf("all %d subscribers failed to accept the batch", len(snapshot))) + } + results := make([]*chippb.PublishResult, 0, len(batch.GetEvents())) - for _, ev := range batch.GetEvents() { + for i, ev := range batch.GetEvents() { result := &chippb.PublishResult{EventId: ev.GetId()} - if !forwarded.Load() { - result.Error = &chippb.PublishError{ - ErrorCode: chippb.PublishErrorCode_PUBLISH_ERROR_CODE_UNKNOWN, - Reason: "chip router could not forward the batch to any subscriber", - } + if i < len(detail) && detail[i].GetError() != nil { + result.Error = detail[i].GetError() } results = append(results, result) } diff --git a/framework/components/chiprouter/cmd/chip-router/main_test.go b/framework/components/chiprouter/cmd/chip-router/main_test.go index 25dc6957d..19e59ecbb 100644 --- a/framework/components/chiprouter/cmd/chip-router/main_test.go +++ b/framework/components/chiprouter/cmd/chip-router/main_test.go @@ -13,9 +13,10 @@ import ( ) type stubPublishClient struct { - mu sync.Mutex - batches int - err error + mu sync.Mutex + batches int + err error + response *chippb.PublishResponse } func (s *stubPublishClient) Publish(context.Context, *cepb.CloudEvent, ...grpc.CallOption) (*chippb.PublishResponse, error) { @@ -26,7 +27,13 @@ func (s *stubPublishClient) PublishBatch(context.Context, *chippb.CloudEventBatc s.mu.Lock() defer s.mu.Unlock() s.batches++ - return &chippb.PublishResponse{}, s.err + if s.err != nil { + return nil, s.err + } + if s.response != nil { + return s.response, nil + } + return &chippb.PublishResponse{}, nil } func (s *stubPublishClient) batchCount() int { @@ -77,11 +84,38 @@ func TestPublishBatchNoSubscribersIsUnavailable(t *testing.T) { } } -func TestPublishBatchAllForwardsFailedReportsPerEventErrors(t *testing.T) { +func TestPublishBatchAllForwardsFailedIsUnavailable(t *testing.T) { r := &router{subscribers: map[string]*subscriber{ "sink": {id: "sink", client: &stubPublishClient{err: context.DeadlineExceeded}}, }} + // A nil RPC error with per-event errors would still acknowledge the whole + // batch for transactional callers (they resolve delivery from the RPC + // outcome alone), so the router must fail the RPC instead. + _, err := r.PublishBatch(t.Context(), batchOf("e1", "e2")) + if err == nil { + t.Fatal("want error when every forward fails, got nil") + } + if status.Code(err) != codes.Unavailable { + t.Errorf("want Unavailable, got %v", status.Code(err)) + } +} + +func TestPublishBatchAggregatesDownstreamRejections(t *testing.T) { + // One subscriber rejects e2 per-event (nil RPC error); the router must + // pass that rejection through rather than acknowledging the whole batch. + r := &router{subscribers: map[string]*subscriber{ + "sink": {id: "sink", client: &stubPublishClient{response: &chippb.PublishResponse{ + Results: []*chippb.PublishResult{ + {EventId: "e1"}, + {EventId: "e2", Error: &chippb.PublishError{ + ErrorCode: chippb.PublishErrorCode_PUBLISH_ERROR_CODE_VALIDATION_FAILED, + Reason: "invalid event", + }}, + }, + }}}, + }} + resp, err := r.PublishBatch(t.Context(), batchOf("e1", "e2")) if err != nil { t.Fatalf("PublishBatch: %v", err) @@ -89,9 +123,35 @@ func TestPublishBatchAllForwardsFailedReportsPerEventErrors(t *testing.T) { if len(resp.Results) != 2 { t.Fatalf("want 2 results, got %d", len(resp.Results)) } - for i, res := range resp.Results { - if res.Error == nil { - t.Errorf("results[%d].Error = nil, want a forward-failure error", i) - } + if resp.Results[0].Error != nil { + t.Errorf("results[0].Error = %v, want nil (accepted)", resp.Results[0].Error) + } + if resp.Results[1].Error == nil { + t.Error("results[1].Error = nil, want the downstream rejection passed through") + } +} + +func TestPublishBatchAnAcceptingSubscriberOverridesADetaillessRejection(t *testing.T) { + // Two detail-providing subscribers disagree: the event is rejected by one + // and accepted by the other — acceptance wins. + r := &router{subscribers: map[string]*subscriber{ + "a": {id: "a", client: &stubPublishClient{response: &chippb.PublishResponse{ + Results: []*chippb.PublishResult{ + {EventId: "e1", Error: &chippb.PublishError{Reason: "rejected by a"}}, + }, + }}}, + "b": {id: "b", client: &stubPublishClient{response: &chippb.PublishResponse{ + Results: []*chippb.PublishResult{ + {EventId: "e1"}, + }, + }}}, + }} + + resp, err := r.PublishBatch(t.Context(), batchOf("e1")) + if err != nil { + t.Fatalf("PublishBatch: %v", err) + } + if len(resp.Results) != 1 || resp.Results[0].Error != nil { + t.Errorf("want e1 accepted, got results=%+v", resp.Results) } } From e8fb57235b7d12fe0323f5faaf801014d09bcb35 Mon Sep 17 00:00:00 2001 From: cawthorne Date: Wed, 7 Oct 2026 11:04:43 +0100 Subject: [PATCH 3/5] chore: bump the chip_router example's go directive to 1.26.6 (gomods tidy) --- framework/examples/chip_router/go.mod | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/framework/examples/chip_router/go.mod b/framework/examples/chip_router/go.mod index 8a652b397..474e75575 100644 --- a/framework/examples/chip_router/go.mod +++ b/framework/examples/chip_router/go.mod @@ -1,6 +1,6 @@ module github.com/smartcontractkit/chainlink-testing-framework/framework/examples/chiprouter -go 1.26.5 +go 1.26.6 replace github.com/smartcontractkit/chainlink-testing-framework/framework => ../../ From 3c74ee1f968170f17a6410dc17e4f2195d42cfb8 Mon Sep 17 00:00:00 2001 From: cawthorne Date: Wed, 7 Oct 2026 11:29:53 +0100 Subject: [PATCH 4/5] docs: trim comments --- .../chiprouter/cmd/chip-router/main.go | 26 ++++++------------- .../chiprouter/cmd/chip-router/main_test.go | 4 +-- 2 files changed, 9 insertions(+), 21 deletions(-) diff --git a/framework/components/chiprouter/cmd/chip-router/main.go b/framework/components/chiprouter/cmd/chip-router/main.go index be6772163..b02f92e2b 100644 --- a/framework/components/chiprouter/cmd/chip-router/main.go +++ b/framework/components/chiprouter/cmd/chip-router/main.go @@ -166,20 +166,14 @@ func (r *router) Publish(ctx context.Context, event *cepb.CloudEvent) (*chippb.P func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch) (*chippb.PublishResponse, error) { snapshot := r.snapshotSubscribers() if len(snapshot) == 0 { - // Nothing will carry these events. Reporting an empty success here - // would let the caller's durable store delete them undelivered. + // Acking events nothing carried would let a durable store delete them undelivered. return nil, status.Error(codes.Unavailable, "no subscribers registered") } - // Aggregate the downstream outcomes so the per-event results reflect what - // actually happened to each event: - // - anyAccepted: at least one subscriber's PublishBatch RPC succeeded. - // - detail: merged per-event results from the subscribers that provide - // them. An event is reported rejected only if every detail-providing - // subscriber rejected it; subscribers that accept without per-event - // detail (a nil-error response with an empty results array) count as - // accepting the whole batch — the router trusts a subscriber's RPC - // success at the forwarding boundary, the same as Publish does. + // anyAccepted: at least one subscriber's RPC succeeded. detail: merged + // per-event results — an event is rejected only if every detail-providing + // subscriber rejected it; a subscriber accepting without per-event detail + // counts as accepting the whole batch. var mu sync.Mutex anyAccepted := false var detail []*chippb.PublishResult @@ -205,9 +199,7 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch detail = results return nil } - // Merge: an accepting result overrides a rejecting one for the - // same event index; a rejection sticks only while no subscriber - // has accepted the event. + // An accepting result overrides a rejection for the same event. for i := range detail { if i >= len(results) { break @@ -222,10 +214,8 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch } _ = group.Wait() - // A nil RPC error does not imply the batch was accepted: every forward - // can fail (or a transactional caller resolves the whole batch from the - // RPC outcome alone), so fail the RPC and let the caller retain and - // retry instead of acknowledging with per-event errors. + // Fail the RPC when nothing accepted: transactional callers resolve the + // whole batch from the RPC outcome alone. if !anyAccepted { return nil, status.Error(codes.Unavailable, fmt.Sprintf("all %d subscribers failed to accept the batch", len(snapshot))) } diff --git a/framework/components/chiprouter/cmd/chip-router/main_test.go b/framework/components/chiprouter/cmd/chip-router/main_test.go index 19e59ecbb..69bf248e8 100644 --- a/framework/components/chiprouter/cmd/chip-router/main_test.go +++ b/framework/components/chiprouter/cmd/chip-router/main_test.go @@ -89,9 +89,7 @@ func TestPublishBatchAllForwardsFailedIsUnavailable(t *testing.T) { "sink": {id: "sink", client: &stubPublishClient{err: context.DeadlineExceeded}}, }} - // A nil RPC error with per-event errors would still acknowledge the whole - // batch for transactional callers (they resolve delivery from the RPC - // outcome alone), so the router must fail the RPC instead. + // Transactional callers resolve the whole batch from the RPC outcome alone. _, err := r.PublishBatch(t.Context(), batchOf("e1", "e2")) if err == nil { t.Fatal("want error when every forward fails, got nil") From c69acae8c2faa2c6a8f33d22eb1d760d910c84d3 Mon Sep 17 00:00:00 2001 From: cawthorne Date: Wed, 7 Oct 2026 11:34:02 +0100 Subject: [PATCH 5/5] fix(chip-router): honor transaction_enabled; a downstream all-rejection fails the RPC MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A subscriber's nil RPC error does not mean its events were accepted: the pinned protocol reports per-event rejections in a successful response. Renamed the flag to anyForwarded (what an RPC success actually tells the router) and gated the response on the batch's own PublishOptions: - transaction_enabled: any unaccepted event fails the RPC (Unavailable) — transactional callers resolve the whole batch from the RPC outcome alone, so per-event errors would still acknowledge everything, including events no subscriber produced. - otherwise: per-event results as before (rejections passed through, acceptance wins over rejection across subscribers). --- .../chiprouter/cmd/chip-router/main.go | 22 ++++-- .../chiprouter/cmd/chip-router/main_test.go | 71 +++++++++++++++++++ 2 files changed, 87 insertions(+), 6 deletions(-) diff --git a/framework/components/chiprouter/cmd/chip-router/main.go b/framework/components/chiprouter/cmd/chip-router/main.go index b02f92e2b..e680b6a86 100644 --- a/framework/components/chiprouter/cmd/chip-router/main.go +++ b/framework/components/chiprouter/cmd/chip-router/main.go @@ -170,12 +170,12 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch return nil, status.Error(codes.Unavailable, "no subscribers registered") } - // anyAccepted: at least one subscriber's RPC succeeded. detail: merged + // anyForwarded: at least one subscriber's RPC succeeded. detail: merged // per-event results — an event is rejected only if every detail-providing // subscriber rejected it; a subscriber accepting without per-event detail // counts as accepting the whole batch. var mu sync.Mutex - anyAccepted := false + anyForwarded := false var detail []*chippb.PublishResult var group errgroup.Group @@ -193,7 +193,7 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch framework.L.Debug().Msgf("chip router forwarded batch to subscriber id=%s", sub.id) mu.Lock() defer mu.Unlock() - anyAccepted = true + anyForwarded = true if results := resp.GetResults(); len(results) > 0 { if detail == nil { detail = results @@ -214,9 +214,9 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch } _ = group.Wait() - // Fail the RPC when nothing accepted: transactional callers resolve the - // whole batch from the RPC outcome alone. - if !anyAccepted { + // No subscriber accepted the batch: fail the RPC so every caller mode + // retains and retries. + if !anyForwarded { return nil, status.Error(codes.Unavailable, fmt.Sprintf("all %d subscribers failed to accept the batch", len(snapshot))) } @@ -228,6 +228,16 @@ func (r *router) PublishBatch(ctx context.Context, batch *chippb.CloudEventBatch } results = append(results, result) } + + // Transactional callers resolve the whole batch from the RPC outcome + // alone, so any unaccepted event must fail the RPC (all-or-nothing). + if batch.GetOptions().GetTransactionEnabled() { + for _, result := range results { + if result.Error != nil { + return nil, status.Error(codes.Unavailable, "not all events were accepted by a subscriber") + } + } + } return &chippb.PublishResponse{Results: results}, nil } diff --git a/framework/components/chiprouter/cmd/chip-router/main_test.go b/framework/components/chiprouter/cmd/chip-router/main_test.go index 69bf248e8..3ce630356 100644 --- a/framework/components/chiprouter/cmd/chip-router/main_test.go +++ b/framework/components/chiprouter/cmd/chip-router/main_test.go @@ -42,6 +42,8 @@ func (s *stubPublishClient) batchCount() int { return s.batches } +func ptr(b bool) *bool { return &b } + func batchOf(ids ...string) *chippb.CloudEventBatch { events := make([]*cepb.CloudEvent, 0, len(ids)) for _, id := range ids { @@ -129,6 +131,75 @@ func TestPublishBatchAggregatesDownstreamRejections(t *testing.T) { } } +func TestPublishBatchTransactionalPartialRejectionFailsRPC(t *testing.T) { + // Transactional callers resolve the whole batch from the RPC outcome + // alone, so a partial rejection must fail the RPC, not ack everything. + r := &router{subscribers: map[string]*subscriber{ + "sink": {id: "sink", client: &stubPublishClient{response: &chippb.PublishResponse{ + Results: []*chippb.PublishResult{ + {EventId: "e1"}, + {EventId: "e2", Error: &chippb.PublishError{Reason: "rejected"}}, + }, + }}}, + }} + + batch := batchOf("e1", "e2") + batch.Options = &chippb.PublishOptions{TransactionEnabled: ptr(true)} + _, err := r.PublishBatch(t.Context(), batch) + if err == nil { + t.Fatal("want error for a partial rejection under transaction_enabled, got nil") + } + if status.Code(err) != codes.Unavailable { + t.Errorf("want Unavailable, got %v", status.Code(err)) + } +} + +func TestPublishBatchTransactionalAllRejectedByDownstreamFailsRPC(t *testing.T) { + // A nil downstream RPC error with every event rejected must NOT count as + // accepted: a transactional caller would otherwise delete them undelivered. + r := &router{subscribers: map[string]*subscriber{ + "sink": {id: "sink", client: &stubPublishClient{response: &chippb.PublishResponse{ + Results: []*chippb.PublishResult{ + {EventId: "e1", Error: &chippb.PublishError{Reason: "rejected"}}, + {EventId: "e2", Error: &chippb.PublishError{Reason: "rejected"}}, + }, + }}}, + }} + + batch := batchOf("e1", "e2") + batch.Options = &chippb.PublishOptions{TransactionEnabled: ptr(true)} + _, err := r.PublishBatch(t.Context(), batch) + if err == nil { + t.Fatal("want error when the subscriber rejected every event, got nil") + } + if status.Code(err) != codes.Unavailable { + t.Errorf("want Unavailable, got %v", status.Code(err)) + } +} + +func TestPublishBatchTransactionalAllAcceptedPasses(t *testing.T) { + r := &router{subscribers: map[string]*subscriber{ + "sink": {id: "sink", client: &stubPublishClient{response: &chippb.PublishResponse{ + Results: []*chippb.PublishResult{ + {EventId: "e1"}, + {EventId: "e2"}, + }, + }}}, + }} + + batch := batchOf("e1", "e2") + batch.Options = &chippb.PublishOptions{TransactionEnabled: ptr(true)} + resp, err := r.PublishBatch(t.Context(), batch) + if err != nil { + t.Fatalf("PublishBatch: %v", err) + } + for i, res := range resp.Results { + if res.Error != nil { + t.Errorf("results[%d].Error = %v, want nil", i, res.Error) + } + } +} + func TestPublishBatchAnAcceptingSubscriberOverridesADetaillessRejection(t *testing.T) { // Two detail-providing subscribers disagree: the event is rejected by one // and accepted by the other — acceptance wins.