diff --git a/framework/components/chiprouter/cmd/chip-router/main.go b/framework/components/chiprouter/cmd/chip-router/main.go index 516d5b71f..e680b6a86 100644 --- a/framework/components/chiprouter/cmd/chip-router/main.go +++ b/framework/components/chiprouter/cmd/chip-router/main.go @@ -19,7 +19,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 +166,18 @@ 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 + // Acking events nothing carried would let a durable store delete them undelivered. + return nil, status.Error(codes.Unavailable, "no subscribers registered") } + // 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 + anyForwarded := false + var detail []*chippb.PublishResult + var group errgroup.Group group.SetLimit(forwardParallel) for _, sub := range snapshot { @@ -174,16 +185,60 @@ 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) + return nil } framework.L.Debug().Msgf("chip router forwarded batch to subscriber id=%s", sub.id) + mu.Lock() + defer mu.Unlock() + anyForwarded = true + if results := resp.GetResults(); len(results) > 0 { + if detail == nil { + detail = results + return nil + } + // An accepting result overrides a rejection for the same 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() - return &chippb.PublishResponse{}, nil + + // 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))) + } + + results := make([]*chippb.PublishResult, 0, len(batch.GetEvents())) + for i, ev := range batch.GetEvents() { + result := &chippb.PublishResult{EventId: ev.GetId()} + if i < len(detail) && detail[i].GetError() != nil { + result.Error = detail[i].GetError() + } + 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 } 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..3ce630356 --- /dev/null +++ b/framework/components/chiprouter/cmd/chip-router/main_test.go @@ -0,0 +1,226 @@ +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 + response *chippb.PublishResponse +} + +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++ + 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 { + s.mu.Lock() + defer s.mu.Unlock() + 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 { + 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 TestPublishBatchAllForwardsFailedIsUnavailable(t *testing.T) { + r := &router{subscribers: map[string]*subscriber{ + "sink": {id: "sink", client: &stubPublishClient{err: context.DeadlineExceeded}}, + }} + + // 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") + } + 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) + } + if len(resp.Results) != 2 { + t.Fatalf("want 2 results, got %d", len(resp.Results)) + } + 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 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. + 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) + } +} 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= 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 => ../../