Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
21 changes: 7 additions & 14 deletions integration/pkg/accessors/evm/evm_source_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -321,22 +321,15 @@ func (r *SourceReader) fetchHeadBatch(ctx context.Context, blockNumbers []*big.I

// FetchMessageSentEvents returns MessageSentEvents in the given block range.
// The toBlock parameter can be nil to query up to the latest block.
// When an RPC provider rejects a query as too large, the range is halved and
// retried. The shrunk limit persists for subsequent calls on this instance.
// When an RPC provider rejects a bounded query as too large, the range is halved
// and retried. The shrunk limit persists for subsequent calls on this instance.
func (r *SourceReader) FetchMessageSentEvents(ctx context.Context, fromBlock, toBlock *big.Int) ([]protocol.MessageSentEvent, error) {
from := fromBlock.Uint64()

var to uint64
if toBlock != nil {
to = toBlock.Uint64()
} else {
latest, _, err := r.headTracker.LatestAndFinalizedBlock(ctx)
if err != nil || latest == nil || latest.Number < 0 {
// Can't resolve upper bound — pass through unbounded
return r.fetchMessageSentEventsRange(ctx, fromBlock, nil)
}
to = uint64(latest.Number)
if toBlock == nil {
return r.fetchMessageSentEventsRange(ctx, fromBlock, nil)
}

from := fromBlock.Uint64()
to := toBlock.Uint64()
if to < from {
return nil, nil
}
Expand Down
51 changes: 9 additions & 42 deletions integration/pkg/accessors/evm/evm_source_reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,28 +41,11 @@ func (m *mockFilterLogsClient) FilterLogs(ctx context.Context, q ethereum.Filter
return m.filterLogsFunc(ctx, q)
}

// mockHeadTracker wraps heads.NullTracker and overrides LatestAndFinalizedBlock.
type mockHeadTracker struct {
heads.Tracker
latest *evmtypes.Head
finalized *evmtypes.Head
err error
}

func (m *mockHeadTracker) LatestAndFinalizedBlock(ctx context.Context) (*evmtypes.Head, *evmtypes.Head, error) {
return m.latest, m.finalized, m.err
}

func newTestSourceReader(t *testing.T, chainClient evmclient.Client) *SourceReader {
t.Helper()
return newTestSourceReaderWithTracker(t, chainClient, heads.NullTracker)
}

func newTestSourceReaderWithTracker(t *testing.T, chainClient evmclient.Client, tracker heads.Tracker) *SourceReader {
t.Helper()
return &SourceReader{
chainClient: chainClient,
headTracker: tracker,
headTracker: heads.NullTracker,
onRampAddress: common.HexToAddress("0x1234"),
ccipMessageSentTopic: common.Hash{}.Hex(),
chainSelector: protocol.ChainSelector(1337),
Expand Down Expand Up @@ -154,43 +137,27 @@ func TestFetchMessageSentEvents_ShrunkLimitPersistsForNextCall(t *testing.T) {
require.Equal(t, [2]uint64{1500, 1999}, secondCallRanges[1])
}

func TestFetchMessageSentEvents_UnboundedQueryShrinksAndRetries(t *testing.T) {
func TestFetchMessageSentEvents_UnboundedQueryPassesThrough(t *testing.T) {
t.Parallel()

latestHead := &evmtypes.Head{
Number: 999,
Hash: common.BigToHash(big.NewInt(999)),
Timestamp: time.Now(),
}
tracker := &mockHeadTracker{latest: latestHead}

var queriedRanges [][2]uint64
var queriedWithNil bool
client := &mockFilterLogsClient{
filterLogsFunc: func(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error) {
from := q.FromBlock.Uint64()
to := q.ToBlock.Uint64()
queriedRanges = append(queriedRanges, [2]uint64{from, to})
if to-from+1 > 500 {
return nil, rangeLimitError()
}
queriedWithNil = q.ToBlock == nil
return []types.Log{}, nil
},
}

reader := newTestSourceReaderWithTracker(t, client, tracker)
reader := newTestSourceReader(t, client)

// Unbounded query [0, nil] — resolves latest=999, span=1000, halved to 500
// Unbounded query should pass through directly without bisection
events, err := reader.FetchMessageSentEvents(context.Background(), big.NewInt(0), nil)
require.NoError(t, err)
require.Empty(t, events)
require.True(t, queriedWithNil, "unbounded query should pass nil toBlock to FilterLogs")

// Should have made 3 calls: [0,999] rejected, [0,499] ok, [500,999] ok
require.Len(t, queriedRanges, 3)
require.Equal(t, [2]uint64{0, 999}, queriedRanges[0])
require.Equal(t, [2]uint64{0, 499}, queriedRanges[1])
require.Equal(t, [2]uint64{500, 999}, queriedRanges[2])

require.Equal(t, uint64(500), reader.maxFilterBlockRange.Load())
// No limit should be set
require.Equal(t, uint64(0), reader.maxFilterBlockRange.Load())
}

func TestFetchMessageSentEvents_GenericErrorDoesNotShrink(t *testing.T) {
Expand Down
Loading