Skip to content
Draft
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
61 changes: 58 additions & 3 deletions framework/components/chiprouter/cmd/chip-router/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down Expand Up @@ -164,26 +166,79 @@ 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 {
group.Go(func() error {
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) {
Expand Down
226 changes: 226 additions & 0 deletions framework/components/chiprouter/cmd/chip-router/main_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
6 changes: 3 additions & 3 deletions framework/components/chiprouter/go.mod
Original file line number Diff line number Diff line change
@@ -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 => ../../

Expand All @@ -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 (
Expand Down
8 changes: 4 additions & 4 deletions framework/components/chiprouter/go.sum

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion framework/examples/chip_router/go.mod
Original file line number Diff line number Diff line change
@@ -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 => ../../

Expand Down
Loading