From a2ed474d09feffed28cafc5125a4f6dead9c35a2 Mon Sep 17 00:00:00 2001 From: nvsriram Date: Tue, 8 Sep 2026 06:52:24 -0700 Subject: [PATCH 1/4] fix: pass unbounded queries through without bisection --- .../pkg/accessors/evm/evm_source_reader.go | 21 +++----- .../accessors/evm/evm_source_reader_test.go | 51 ++++--------------- 2 files changed, 16 insertions(+), 56 deletions(-) diff --git a/integration/pkg/accessors/evm/evm_source_reader.go b/integration/pkg/accessors/evm/evm_source_reader.go index 7e0f49553..32e173cde 100644 --- a/integration/pkg/accessors/evm/evm_source_reader.go +++ b/integration/pkg/accessors/evm/evm_source_reader.go @@ -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 } diff --git a/integration/pkg/accessors/evm/evm_source_reader_test.go b/integration/pkg/accessors/evm/evm_source_reader_test.go index 87409a01d..41897ac52 100644 --- a/integration/pkg/accessors/evm/evm_source_reader_test.go +++ b/integration/pkg/accessors/evm/evm_source_reader_test.go @@ -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), @@ -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) { From 1c5adbc6c8b65507c8155f304847b9321d9da93a Mon Sep 17 00:00:00 2001 From: nvsriram Date: Wed, 9 Sep 2026 07:12:40 -0700 Subject: [PATCH 2/4] Revert "fix(verifier): add adaptive log fetching and return early on no progress (#1370)" This reverts commit 92ca56b6fa987433550d6d85004d574bd0536fb9. --- .../pkg/accessors/evm/evm_source_reader.go | 140 --------------- .../accessors/evm/evm_source_reader_test.go | 160 +----------------- verifier/pkg/sourcereader/service.go | 5 - verifier/pkg/sourcereader/service_test.go | 35 ++-- 4 files changed, 22 insertions(+), 318 deletions(-) diff --git a/integration/pkg/accessors/evm/evm_source_reader.go b/integration/pkg/accessors/evm/evm_source_reader.go index 32e173cde..d526edd14 100644 --- a/integration/pkg/accessors/evm/evm_source_reader.go +++ b/integration/pkg/accessors/evm/evm_source_reader.go @@ -7,8 +7,6 @@ import ( "fmt" "maps" "math/big" - "strings" - "sync/atomic" "time" "github.com/ethereum/go-ethereum" @@ -36,78 +34,6 @@ var ( _ chainaccess.CriticalSourceInvariantCallbackSetter = (*SourceReader)(nil) ) -// rangeLimitErrorSubstrings are common RPC provider error messages indicating -// the requested eth_getLogs block range exceeds the provider's limit. -// -// Sources: -// -// - geth: "exceed maximum block range %d" -// https://github.com/ethereum/go-ethereum/blob/master/eth/filters/filter.go#L148-L149 -// -// - erigon: "query block range exceeds server limit, narrow your filter" -// and "query returns too many logs, narrow your filter" -// https://github.com/erigontech/erigon/blob/main/rpc/jsonrpc/eth_receipts.go#L49-L54 -// -// - reth: "query exceeds max block range %d" -// and "query exceeds max results ..." -// https://github.com/paradigmxyz/reth/blob/main/crates/rpc/rpc/src/eth/filter.rs#L976-L981 -// -// - nethermind: "Block range ... exceeds the maximum of ... blocks per logs request." -// https://github.com/NethermindEth/nethermind/blob/master/src/Nethermind/Nethermind.JsonRpc/Modules/Eth/EthRpcModule.cs#L1278-L1291 -// -// - infura: "query returned more than ... results. Try with this block range ..." -// https://github.com/smartcontractkit/chainlink-evm/blob/develop/pkg/client/errors.go#L697-L699 -// -// - alchemy: "Log response size exceeded. You can make eth_getLogs requests with up to ..." -// https://github.com/smartcontractkit/chainlink-evm/blob/develop/pkg/client/errors.go#L701-L703 -// -// - quicknode: "eth_getLogs and eth_newFilter are limited to a 10,000 blocks range" -// https://support.quicknode.com/articles/3261121056 -// "eth_getLogs is limited to a ... range" -// https://github.com/smartcontractkit/chainlink-evm/blob/develop/pkg/client/errors.go#L705-L707 -// -// - simplyvc: "too wide blocks range, the limit is ..." -// https://github.com/smartcontractkit/chainlink-evm/blob/develop/pkg/client/errors.go#L709-L711 -// -// - dRPC: "requested too many blocks from ... to ..., maximum is set to ..." -// https://github.com/smartcontractkit/chainlink-evm/blob/develop/pkg/client/errors.go#L713-L715 -// -// - hyperliquid: "query exceeds max block range" -// https://github.com/smartcontractkit/chainlink-evm/blob/develop/pkg/client/errors.go#L717-L720 -// -// - jovay: "Exceeded max range limit for eth_getLogs: 1000" -// Observed from the Jovay RPC -var rangeLimitErrorSubstrings = []string{ - "exceed maximum block range", // geth - "query block range exceeds server limit", // erigon - "query returns too many logs", // erigon - "query exceeds max block range", // reth, hyperliquid - "query exceeds max results", // reth - "blocks per logs request", // nethermind - "query returned more than", // infura - "log response size exceeded", // alchemy - "eth_getlogs and eth_newfilter are limited to", // quicknode - "eth_getlogs is limited to", // quicknode - "too wide blocks range", // simplyvc - "requested too many blocks", // dRPC - "exceeded max range limit", // jovay -} - -// isRangeLimitError reports whether err looks like an RPC rejection due to the -// requested block range being too large, based on common provider wording. -func isRangeLimitError(err error) bool { - if err == nil { - return false - } - msg := strings.ToLower(err.Error()) - for _, substr := range rangeLimitErrorSubstrings { - if strings.Contains(msg, substr) { - return true - } - } - return false -} - type SourceReader struct { chainClient client.Client headTracker heads.Tracker @@ -121,7 +47,6 @@ type SourceReader struct { lggr logger.Logger onRampABI *abi.ABI // Cached ABI to avoid re-parsing onCriticalInvariant func(context.Context) - maxFilterBlockRange *atomic.Uint64 // Single eth_getLogs query block span sourceReaderHeaderFetchBatchSize int } @@ -219,7 +144,6 @@ func NewEVMSourceReader( chainSelector: chainSelector, lggr: lggr, onRampABI: onRampABI, - maxFilterBlockRange: new(atomic.Uint64), sourceReaderHeaderFetchBatchSize: sourceReaderHeaderFetchBatchSize(headerFetchBatchSize), } reader.SetCriticalSourceInvariantCallback(onCriticalInvariant) @@ -321,71 +245,7 @@ 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 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) { - if toBlock == nil { - return r.fetchMessageSentEventsRange(ctx, fromBlock, nil) - } - - from := fromBlock.Uint64() - to := toBlock.Uint64() - if to < from { - return nil, nil - } - - maxRange := r.maxFilterBlockRange.Load() - var allEvents []protocol.MessageSentEvent - - for from <= to { - end := to - if maxRange > 0 { - end = min(from+maxRange-1, to) - } - - events, err := r.fetchMessageSentEventsRange(ctx, - new(big.Int).SetUint64(from), - new(big.Int).SetUint64(end), - ) - if err == nil { - allEvents = append(allEvents, events...) - if end >= to { - break - } - from = end + 1 - continue - } - - // Error handling - if ctx.Err() != nil { - return allEvents, err - } - if !isRangeLimitError(err) { - return allEvents, err - } - - querySize := end - from + 1 - if querySize <= 1 { - return allEvents, err - } - - newMaxRange := querySize / 2 - r.lggr.Warnw( - "Log query rejected as range too large, retrying with smaller block range", - "error", err, - "fromBlock", from, - "querySize", querySize, - "nextQuerySize", newMaxRange, - ) - // Persist range for future calls - maxRange = newMaxRange - r.maxFilterBlockRange.Store(maxRange) - } - return allEvents, nil -} - -// fetchMessageSentEventsRange queries a single eth_getLogs range without chunking. -func (r *SourceReader) fetchMessageSentEventsRange(ctx context.Context, fromBlock, toBlock *big.Int) ([]protocol.MessageSentEvent, error) { rangeQuery := ethereum.FilterQuery{ FromBlock: fromBlock, ToBlock: toBlock, diff --git a/integration/pkg/accessors/evm/evm_source_reader_test.go b/integration/pkg/accessors/evm/evm_source_reader_test.go index 41897ac52..eba9cac55 100644 --- a/integration/pkg/accessors/evm/evm_source_reader_test.go +++ b/integration/pkg/accessors/evm/evm_source_reader_test.go @@ -3,7 +3,6 @@ package evm import ( "context" "errors" - "fmt" "math/big" "strconv" "sync/atomic" @@ -13,7 +12,6 @@ import ( "github.com/ethereum/go-ethereum" "github.com/ethereum/go-ethereum/accounts/abi/bind" "github.com/ethereum/go-ethereum/common" - "github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/rpc" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" @@ -30,156 +28,6 @@ import ( "github.com/smartcontractkit/chainlink-ccv/protocol" ) -// mockFilterLogsClient embeds evmclient.Client and overrides FilterLogs to -// simulate RPC range-limit rejections and successes. -type mockFilterLogsClient struct { - evmclient.Client - filterLogsFunc func(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error) -} - -func (m *mockFilterLogsClient) FilterLogs(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error) { - return m.filterLogsFunc(ctx, q) -} - -func newTestSourceReader(t *testing.T, chainClient evmclient.Client) *SourceReader { - t.Helper() - return &SourceReader{ - chainClient: chainClient, - headTracker: heads.NullTracker, - onRampAddress: common.HexToAddress("0x1234"), - ccipMessageSentTopic: common.Hash{}.Hex(), - chainSelector: protocol.ChainSelector(1337), - lggr: logger.Test(t), - maxFilterBlockRange: new(atomic.Uint64), - } -} - -func rangeLimitError() error { - return fmt.Errorf("RPC call failed: exceeded max range limit for eth_getLogs") -} - -func TestFetchMessageSentEvents_BoundedQueryShrinksAndRetries(t *testing.T) { - t.Parallel() - - var queriedRanges [][2]uint64 // track [from, to] pairs - callCount := 0 - - 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}) - callCount++ - - // Reject ranges > 500 blocks - if to-from+1 > 500 { - return nil, rangeLimitError() - } - return []types.Log{}, nil - }, - } - - reader := newTestSourceReader(t, client) - - // Query [0, 999] — 1000 blocks, should be rejected, halved to 500, then succeed - events, err := reader.FetchMessageSentEvents(context.Background(), big.NewInt(0), big.NewInt(999)) - require.NoError(t, err) - require.Empty(t, events) - - // 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]) - - // The shrunk limit should persist - require.Equal(t, uint64(500), reader.maxFilterBlockRange.Load()) -} - -func TestFetchMessageSentEvents_ShrunkLimitPersistsForNextCall(t *testing.T) { - t.Parallel() - - callCount := 0 - client := &mockFilterLogsClient{ - filterLogsFunc: func(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error) { - callCount++ - from := q.FromBlock.Uint64() - to := q.ToBlock.Uint64() - if to-from+1 > 500 { - return nil, rangeLimitError() - } - return []types.Log{}, nil - }, - } - - reader := newTestSourceReader(t, client) - - // First call: [0, 999] — triggers shrink to 500 - _, err := reader.FetchMessageSentEvents(context.Background(), big.NewInt(0), big.NewInt(999)) - require.NoError(t, err) - require.Equal(t, uint64(500), reader.maxFilterBlockRange.Load()) - - // Second call: [1000, 1999] — should use the persisted 500 limit - var secondCallRanges [][2]uint64 - client.filterLogsFunc = func(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error) { - from := q.FromBlock.Uint64() - to := q.ToBlock.Uint64() - secondCallRanges = append(secondCallRanges, [2]uint64{from, to}) - return []types.Log{}, nil - } - - _, err = reader.FetchMessageSentEvents(context.Background(), big.NewInt(1000), big.NewInt(1999)) - require.NoError(t, err) - - // Should have chunked into [1000,1499], [1500,1999] - require.Len(t, secondCallRanges, 2) - require.Equal(t, [2]uint64{1000, 1499}, secondCallRanges[0]) - require.Equal(t, [2]uint64{1500, 1999}, secondCallRanges[1]) -} - -func TestFetchMessageSentEvents_UnboundedQueryPassesThrough(t *testing.T) { - t.Parallel() - - var queriedWithNil bool - client := &mockFilterLogsClient{ - filterLogsFunc: func(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error) { - queriedWithNil = q.ToBlock == nil - return []types.Log{}, nil - }, - } - - reader := newTestSourceReader(t, client) - - // 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") - - // No limit should be set - require.Equal(t, uint64(0), reader.maxFilterBlockRange.Load()) -} - -func TestFetchMessageSentEvents_GenericErrorDoesNotShrink(t *testing.T) { - t.Parallel() - - genericErr := fmt.Errorf("connection refused") - client := &mockFilterLogsClient{ - filterLogsFunc: func(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error) { - return nil, genericErr - }, - } - - reader := newTestSourceReader(t, client) - - events, err := reader.FetchMessageSentEvents(context.Background(), big.NewInt(0), big.NewInt(999)) - require.ErrorIs(t, err, genericErr) - require.Empty(t, events) - - // Limit should NOT be set - require.Equal(t, uint64(0), reader.maxFilterBlockRange.Load()) -} - type stubOnRampStaticConfigGetter struct { cfg onramp.OnRampStaticConfig err error @@ -280,6 +128,14 @@ func blockNumFromArg(arg any) int64 { return n } +func newTestSourceReader(t *testing.T, cc evmclient.Client) *SourceReader { + t.Helper() + return &SourceReader{ + chainClient: cc, + lggr: logger.Test(t), + } +} + func TestGetBlocksHeaders_BatchesAndChunks(t *testing.T) { t.Parallel() diff --git a/verifier/pkg/sourcereader/service.go b/verifier/pkg/sourcereader/service.go index 1f9701187..2629d8107 100644 --- a/verifier/pkg/sourcereader/service.go +++ b/verifier/pkg/sourcereader/service.go @@ -350,11 +350,6 @@ func (r *Service) processEventCycle(ctx context.Context, latest, finalized *prot r.logger.Warnw("Error when querying logs", "error", err, "fromBlock", fromBlock.String(), "toBlock", "latest") - - // Only return early when no progress was made - if lastQueriedBlock.Cmp(fromBlock) == 0 { - return false - } } tasks := make([]verifier.VerificationTask, 0, len(events)) diff --git a/verifier/pkg/sourcereader/service_test.go b/verifier/pkg/sourcereader/service_test.go index aa0823034..54609f136 100644 --- a/verifier/pkg/sourcereader/service_test.go +++ b/verifier/pkg/sourcereader/service_test.go @@ -1324,8 +1324,7 @@ func TestSRS_FailureRetriesNextTick(t *testing.T) { reader := mocks.NewMockSourceReader(t) - // latest == fromBlock (99) collapses the range to be non-bisectable - latest := &protocol.BlockHeader{Number: 99} + latest := &protocol.BlockHeader{Number: 1000} finalized := &protocol.BlockHeader{Number: 900} reader.EXPECT(). @@ -1444,8 +1443,7 @@ func TestSRS_FailureDoesNotDeleteExistingTasks(t *testing.T) { reader := mocks.NewMockSourceReader(t) - // latest == fromBlock (99) collapses the range to be non-bisectable - latest := &protocol.BlockHeader{Number: 99} + latest := &protocol.BlockHeader{Number: 1000} finalized := &protocol.BlockHeader{Number: 900} reader.EXPECT(). @@ -1796,9 +1794,8 @@ func TestSRS_PartialRead_EventsFromSuccessfulChunksQueued(t *testing.T) { reader := mocks.NewMockSourceReader(t) - // maxBlockRange=500 splits [100, 601) into two chunks: [100,600] and [601,nil]. - // latest == 601 makes the second chunk non-bisectable - latest := &protocol.BlockHeader{Number: 601} + // maxBlockRange=500 splits [100, 700) into two chunks: [100,600] and [601,nil]. + latest := &protocol.BlockHeader{Number: 700} finalized := &protocol.BlockHeader{Number: 600} reader.EXPECT().LatestAndFinalizedBlock(mock.Anything).Return(latest, finalized, nil).Maybe() @@ -1849,9 +1846,8 @@ func TestSRS_PartialRead_ProgressAdvancesToLastSuccessfulChunkBound(t *testing.T reader := mocks.NewMockSourceReader(t) - // maxBlockRange=500 splits [100, 601) into [100,600] and [601,nil]. - // latest == 601 makes the second chunk non-bisectable - latest := &protocol.BlockHeader{Number: 601} + // maxBlockRange=500 splits [100, 700) into [100,600] and [601,nil]. + latest := &protocol.BlockHeader{Number: 700} finalized := &protocol.BlockHeader{Number: 600} reader.EXPECT().LatestAndFinalizedBlock(mock.Anything).Return(latest, finalized, nil).Maybe() @@ -1898,10 +1894,9 @@ func TestSRS_PartialRead_MultipleChunksSucceedBeforeFailure(t *testing.T) { reader := mocks.NewMockSourceReader(t) - // maxBlockRange=300 splits [100, 702) into three chunks: + // maxBlockRange=300 splits [100, 1000) into three chunks: // [100,400], [401,701], [702,nil]. - // latest == 702 makes the third chunk non-bisectable - latest := &protocol.BlockHeader{Number: 702} + latest := &protocol.BlockHeader{Number: 1000} finalized := &protocol.BlockHeader{Number: 800} reader.EXPECT().LatestAndFinalizedBlock(mock.Anything).Return(latest, finalized, nil).Maybe() @@ -1963,15 +1958,13 @@ func TestSRS_PartialRead_TotalFailureDoesNotAdvanceProgress(t *testing.T) { reader := mocks.NewMockSourceReader(t) - // latest == fromBlock (100) collapses the range to be non-bisectable - latest := &protocol.BlockHeader{Number: 100} - finalized := &protocol.BlockHeader{Number: 100} + latest := &protocol.BlockHeader{Number: 700} + finalized := &protocol.BlockHeader{Number: 600} reader.EXPECT().LatestAndFinalizedBlock(mock.Anything).Return(latest, finalized, nil).Maybe() // Chunk 1 fails immediately — no events, no subsequent chunk calls. - nilBigInt := mock.MatchedBy(func(b *big.Int) bool { return b == nil }) reader.EXPECT(). - FetchMessageSentEvents(mock.Anything, big.NewInt(100), nilBigInt). + FetchMessageSentEvents(mock.Anything, big.NewInt(100), big.NewInt(600)). Return(nil, assert.AnError). Once() // Chunk 2 must NOT be called: testify will fail the test on any unexpected call. @@ -2052,9 +2045,9 @@ func TestSRS_PartialRead_ProgressCapsAtFinalizedWhenChunkBoundExceedsIt(t *testi reader := mocks.NewMockSourceReader(t) - // maxBlockRange=400 splits [100, 902) into chunks [100,500], [501,901], [902,nil]. - // latest == 902 makes the third chunk non-bisectable - latest := &protocol.BlockHeader{Number: 902} + // maxBlockRange=400 splits [100, 1000) into chunks [100,500], [501,901], [902,nil]. + // finalized is 600, so the successful chunk 2 boundary (901) exceeds finalized. + latest := &protocol.BlockHeader{Number: 1000} finalized := &protocol.BlockHeader{Number: 600} reader.EXPECT().LatestAndFinalizedBlock(mock.Anything).Return(latest, finalized, nil).Maybe() From 75803793e33b27995d2a80c9d8bf6428504e09df Mon Sep 17 00:00:00 2001 From: nvsriram Date: Wed, 9 Sep 2026 07:17:14 -0700 Subject: [PATCH 3/4] chore: update DefaultMaxBlockRange to 100 --- verifier/pkg/sourcereader/service.go | 2 +- verifier/pkg/sourcereader/service_test.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/verifier/pkg/sourcereader/service.go b/verifier/pkg/sourcereader/service.go index 2629d8107..da3ca2278 100644 --- a/verifier/pkg/sourcereader/service.go +++ b/verifier/pkg/sourcereader/service.go @@ -30,7 +30,7 @@ import ( const ( DefaultPollInterval = 2100 * time.Millisecond DefaultPollTimeout = 10 * time.Second - DefaultMaxBlockRange = 1500 + DefaultMaxBlockRange = 100 ) type blockRange struct { diff --git a/verifier/pkg/sourcereader/service_test.go b/verifier/pkg/sourcereader/service_test.go index 54609f136..17e742195 100644 --- a/verifier/pkg/sourcereader/service_test.go +++ b/verifier/pkg/sourcereader/service_test.go @@ -1194,7 +1194,7 @@ func TestSRS_LargeRangeChunkedInSingleCycle(t *testing.T) { curseDetector.EXPECT().Start(mock.Anything).Return(nil).Maybe() curseDetector.EXPECT().Close().Return(nil).Maybe() - srs, _, _ := newTestSRS(t, chain, reader, chainStatusMgr, curseDetector, 10*time.Millisecond, 0) + srs, _, _ := newTestSRS(t, chain, reader, chainStatusMgr, curseDetector, 10*time.Millisecond, 1500) srs.lastProcessedFinalizedBlock.Store(big.NewInt(99)) srs.processEventCycle(ctx, latest, finalized) From eb8d82ca3520e4f80cdc62ee0df6c4673747f62f Mon Sep 17 00:00:00 2001 From: nvsriram Date: Wed, 9 Sep 2026 07:23:16 -0700 Subject: [PATCH 4/4] feat: readd early return on no progress --- verifier/pkg/sourcereader/service.go | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/verifier/pkg/sourcereader/service.go b/verifier/pkg/sourcereader/service.go index da3ca2278..1fb80dbee 100644 --- a/verifier/pkg/sourcereader/service.go +++ b/verifier/pkg/sourcereader/service.go @@ -350,6 +350,11 @@ func (r *Service) processEventCycle(ctx context.Context, latest, finalized *prot r.logger.Warnw("Error when querying logs", "error", err, "fromBlock", fromBlock.String(), "toBlock", "latest") + + // Only return early when no progress was made + if lastQueriedBlock.Cmp(fromBlock) == 0 { + return false + } } tasks := make([]verifier.VerificationTask, 0, len(events))