diff --git a/CHANGELOG.md b/CHANGELOG.md index befb7bfd86..1b673fc484 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,7 @@ Ref: https://keepachangelog.com/en/1.0.0/ # Changelog ## v6.6 sei-chain +* [#3759](https://github.com/sei-protocol/sei-chain/pull/3759) feat(evmrpc): bound `eth_getLogs` peak memory with a matched-log cap. `max_log_no_block` now caps both bounded and open-ended block ranges; exceeding it returns an error instead of silently truncating. * [#3743](https://github.com/sei-protocol/sei-chain/pull/3743) Backport `release/v6.6`: feat(cosmos): range-check Dec conversions and DecCoin validation (CON-369) * [#3732](https://github.com/sei-protocol/sei-chain/pull/3732) Bump version in prep to release v6.6-rc2 * [#3731](https://github.com/sei-protocol/sei-chain/pull/3731) Backport `release/v6.6`: Update v6.6 changelog in prep to cut RC2 diff --git a/evmrpc/config/config.go b/evmrpc/config/config.go index db238922d0..88e3be4eba 100644 --- a/evmrpc/config/config.go +++ b/evmrpc/config/config.go @@ -99,7 +99,9 @@ type Config struct { // Deny list defines list of methods that EVM RPC should fail fast DenyList []string `mapstructure:"deny_list"` - // max number of logs returned if block range is open-ended + // max number of logs a single eth_getLogs query may match before it errors, + // for both bounded and open-ended block ranges (a non-positive value falls + // back to DefaultMaxLogLimit) MaxLogNoBlock int64 `mapstructure:"max_log_no_block"` // max number of blocks to query logs for @@ -695,7 +697,8 @@ enabled_legacy_sei_apis = [ # "sei2_getBlockTransactionCountByNumber", ] -# max number of logs returned if block range is open-ended +# max number of logs a single eth_getLogs query may match before it errors, +# for both bounded and open-ended block ranges max_log_no_block = {{ .EVM.MaxLogNoBlock }} # max number of blocks to query logs for diff --git a/evmrpc/filter.go b/evmrpc/filter.go index 09baf1f2b0..9f2149e848 100644 --- a/evmrpc/filter.go +++ b/evmrpc/filter.go @@ -8,6 +8,7 @@ import ( "fmt" "sort" "sync" + "sync/atomic" "time" "github.com/ethereum/go-ethereum/common" @@ -802,23 +803,31 @@ func (f *LogFetcher) GetLogsByFilters(ctx context.Context, crit filters.FilterCr return nil, 0, fmt.Errorf("block range too large (%d), maximum allowed is %d blocks", blockRange, f.filterConfig.maxBlock) } + // maxLog caps the number of matching logs a single query may return, for + // both bounded and open-ended queries. Exceeding it is an error (not + // silent truncation), so peak memory stays bounded to the cap. NewFilterAPI + // coerces a non-positive value to DefaultMaxLogLimit, so limit is always > 0 + // here; downstream helpers still treat limit <= 0 as uncapped. + limit := f.filterConfig.maxLog + // blockHash queries must use the hash-aware block fetch path below. // Range-query receipt stores only constrain by numeric block range and // do not enforce crit.BlockHash. if crit.BlockHash == nil { // Try efficient range query first // #nosec G115 -- begin and end are validated to be positive block heights above - if logs, rangeErr := f.tryFilterLogsRange(ctx, uint64(begin), uint64(end), crit); rangeErr == nil { + if logs, rangeErr := f.tryFilterLogsRange(ctx, uint64(begin), uint64(end), crit, limit); rangeErr == nil { return logs, end, nil } else if !errors.Is(rangeErr, receipt.ErrRangeQueryNotSupported) { - // If it's a real error (not just unsupported), return it + // A real error (including receipt.ErrTooManyLogs) short-circuits; + // only "unsupported" falls through to the block-by-block path. return nil, 0, rangeErr } } // Fall back to block-by-block querying for backends that don't support range queries bloomIndexes := EncodeFilters(crit.Addresses, crit.Topics) - blocks, end, applyOpenEndedLogLimit, err := f.fetchBlocksByCrit(ctx, crit, lastToHeight, bloomIndexes) + blocks, end, err := f.fetchBlocksByCrit(ctx, crit, lastToHeight, bloomIndexes) if err != nil { return nil, 0, err } @@ -828,6 +837,12 @@ func (f *LogFetcher) GetLogsByFilters(ctx context.Context, crit filters.FilterCr sortedBatches := make([][]*ethtypes.Log, 0) var wg sync.WaitGroup var submitError error + // collected tracks matched logs across workers so the fan-out can stop + // queueing batches once the cap is exceeded. + var collected int64 + capExceeded := func() bool { + return limit > 0 && atomic.LoadInt64(&collected) > limit + } processBatch := func(batch []*coretypes.ResultBlock) { defer wg.Done() @@ -835,7 +850,14 @@ func (f *LogFetcher) GetLogsByFilters(ctx context.Context, crit filters.FilterCr localLogs := f.globalLogSlicePool.Get() for _, block := range batch { + // Cooperative early-abort: once any worker has pushed the matched + // count past the cap, stop materializing further blocks. + if capExceeded() { + break + } + before := len(localLogs) f.GetLogsForBlockPooled(block, crit, &localLogs) + atomic.AddInt64(&collected, int64(len(localLogs)-before)) } // Sort the local batch @@ -855,6 +877,11 @@ func (f *LogFetcher) GetLogsByFilters(ctx context.Context, crit filters.FilterCr // Batch process with fail-fast blockBatch := make([]*coretypes.ResultBlock, 0, evmrpcconfig.WorkerBatchSize) for block := range blocks { + // Stop queueing once the cap is exceeded; in-flight batches finish and + // the overflow is reported after wg.Wait. + if capExceeded() { + break + } blockBatch = append(blockBatch, block) if len(blockBatch) >= evmrpcconfig.WorkerBatchSize { @@ -874,8 +901,8 @@ func (f *LogFetcher) GetLogsByFilters(ctx context.Context, crit filters.FilterCr return nil, 0, submitError } - // Process remaining blocks - if len(blockBatch) > 0 { + // Process remaining blocks, unless the cap is already exceeded + if len(blockBatch) > 0 && !capExceeded() { wg.Add(1) if err := runner.SubmitWithMetrics(func() { processBatch(blockBatch) }); err != nil { wg.Done() @@ -893,13 +920,15 @@ func (f *LogFetcher) GetLogsByFilters(ctx context.Context, crit filters.FilterCr } }() - // Push the cap into the merge so it stops popping once the limit is reached, - // bounding the merged result allocation to O(maxLog) instead of O(total - // matching logs). A limit of 0 means no cap. - var limit int64 - if applyOpenEndedLogLimit { - limit = f.filterConfig.maxLog + // Drain the blocks channel so its producer goroutines are not left blocked + // on a full buffer when we broke out of the loop above on overflow. + for range blocks { + } + + if limit > 0 && collected > limit { + return nil, 0, receipt.NewTooManyLogsError(limit) } + res = f.mergeSortedLogs(sortedBatches, limit) // Ensure we never return nil, always return an array (even if empty) @@ -974,7 +1003,9 @@ func (f *LogFetcher) earliestHeight(ctx context.Context) (int64, error) { // tryFilterLogsRange attempts to use the efficient range query if supported by the backend. // Returns ErrRangeQueryNotSupported if the backend doesn't support range queries. -func (f *LogFetcher) tryFilterLogsRange(ctx context.Context, fromBlock, toBlock uint64, crit filters.FilterCriteria) ([]*ethtypes.Log, error) { +// When limit > 0 the backend aborts with receipt.ErrTooManyLogs once more than +// limit logs match, bounding peak memory; limit <= 0 means no cap. +func (f *LogFetcher) tryFilterLogsRange(ctx context.Context, fromBlock, toBlock uint64, crit filters.FilterCriteria, limit int64) ([]*ethtypes.Log, error) { store := f.k.ReceiptStore() if store == nil { return nil, receipt.ErrRangeQueryNotSupported @@ -984,7 +1015,7 @@ func (f *LogFetcher) tryFilterLogsRange(ctx context.Context, fromBlock, toBlock // #nosec G115 -- toBlock is a block height which fits in int64 sdkCtx := f.ctxProvider(int64(toBlock)) - logs, err := store.FilterLogs(sdkCtx, fromBlock, toBlock, crit) + logs, err := store.FilterLogs(sdkCtx, fromBlock, toBlock, crit, limit) if err != nil { return nil, err } @@ -993,7 +1024,18 @@ func (f *LogFetcher) tryFilterLogsRange(ctx context.Context, fromBlock, toBlock return []*ethtypes.Log{}, nil } - return f.normalizeRangeQueryLogs(ctx, logs, crit) + normalized, err := f.normalizeRangeQueryLogs(ctx, logs, crit) + if err != nil { + return nil, err + } + + // Re-check the cap against the normalized (EVM-visible) count, which the + // store's tag-index count does not match. + if limit > 0 && int64(len(normalized)) > limit { + return nil, receipt.NewTooManyLogsError(limit) + } + + return normalized, nil } // normalizeRangeQueryLogs corrects BlockHash, TxIndex, and LogIndex on logs @@ -1170,7 +1212,7 @@ func MatchesCriteria(log *ethtypes.Log, crit filters.FilterCriteria) bool { } // Optimized fetchBlocksByCrit with batch processing -func (f *LogFetcher) fetchBlocksByCrit(ctx context.Context, crit filters.FilterCriteria, lastToHeight int64, bloomIndexes [][]BloomIndexes) (chan *coretypes.ResultBlock, int64, bool, error) { +func (f *LogFetcher) fetchBlocksByCrit(ctx context.Context, crit filters.FilterCriteria, lastToHeight int64, bloomIndexes [][]BloomIndexes) (chan *coretypes.ResultBlock, int64, error) { if crit.BlockHash != nil { // Check for invalid zero hash zeroHash := common.Hash{} @@ -1178,7 +1220,7 @@ func (f *LogFetcher) fetchBlocksByCrit(ctx context.Context, crit filters.FilterC // For invalid hash, return empty channel instead of error res := make(chan *coretypes.ResultBlock) close(res) - return res, 0, false, nil + return res, 0, nil } block, err := blockByHashRespectingWatermarks(ctx, f.tmClient, f.watermarks, crit.BlockHash[:], 1) @@ -1186,18 +1228,22 @@ func (f *LogFetcher) fetchBlocksByCrit(ctx context.Context, crit filters.FilterC // For non-existent blocks, return empty channel instead of error res := make(chan *coretypes.ResultBlock) close(res) - return res, 0, false, nil + return res, 0, nil } res := make(chan *coretypes.ResultBlock, 1) res <- block close(res) - return res, 0, false, nil + return res, 0, nil } - applyOpenEndedLogLimit := f.filterConfig.maxLog > 0 && (crit.FromBlock == nil || crit.ToBlock == nil) + // Open-ended queries (missing fromBlock or toBlock) have their block range + // windowed down to the most recent maxBlock blocks rather than erroring; a + // bounded query whose range is too large still errors. The matched-log cap + // (maxLog) is enforced separately by the caller. + applyOpenEndedBlockWindow := f.filterConfig.maxLog > 0 && (crit.FromBlock == nil || crit.ToBlock == nil) latest, err := f.watermarks.LatestHeight(ctx) if err != nil { - return nil, 0, false, err + return nil, 0, err } earliest, err := f.watermarks.EarliestHeight(ctx) if err != nil { @@ -1205,22 +1251,22 @@ func (f *LogFetcher) fetchBlocksByCrit(ctx context.Context, crit filters.FilterC } begin, end, err := ComputeBlockBounds(latest, earliest, lastToHeight, crit) if err != nil { - return nil, 0, false, err + return nil, 0, err } blockRange := end - begin + 1 - if applyOpenEndedLogLimit && blockRange > f.filterConfig.maxBlock { + if applyOpenEndedBlockWindow && blockRange > f.filterConfig.maxBlock { begin = end - f.filterConfig.maxBlock + 1 if begin < earliest { begin = earliest } - } else if !applyOpenEndedLogLimit && f.filterConfig.maxBlock > 0 && blockRange > f.filterConfig.maxBlock { + } else if !applyOpenEndedBlockWindow && f.filterConfig.maxBlock > 0 && blockRange > f.filterConfig.maxBlock { // Use consistent error message format - return nil, 0, false, fmt.Errorf("block range too large (%d), maximum allowed is %d blocks", blockRange, f.filterConfig.maxBlock) + return nil, 0, fmt.Errorf("block range too large (%d), maximum allowed is %d blocks", blockRange, f.filterConfig.maxBlock) } if begin > end { - return nil, 0, false, fmt.Errorf("fromBlock %d is after toBlock %d", begin, end) + return nil, 0, fmt.Errorf("fromBlock %d is after toBlock %d", begin, end) } res := make(chan *coretypes.ResultBlock, end-begin+1) @@ -1243,7 +1289,7 @@ func (f *LogFetcher) fetchBlocksByCrit(ctx context.Context, crit filters.FilterC } }(batchStart, batchEnd)); err != nil { wg.Done() - return nil, 0, false, fmt.Errorf("system overloaded, please reduce request frequency: %w", err) + return nil, 0, fmt.Errorf("system overloaded, please reduce request frequency: %w", err) } } @@ -1262,10 +1308,10 @@ func (f *LogFetcher) fetchBlocksByCrit(ctx context.Context, crit filters.FilterC } if firstErr != nil { - return nil, 0, false, firstErr + return nil, 0, firstErr } - return res, end, applyOpenEndedLogLimit, nil + return res, end, nil } // Batch processing function for blocks diff --git a/evmrpc/filter_test.go b/evmrpc/filter_test.go index df5daeadad..76d7e483d5 100644 --- a/evmrpc/filter_test.go +++ b/evmrpc/filter_test.go @@ -299,6 +299,35 @@ func testFilterGetLogs(t *testing.T, namespace string, tests []GetFilterLogTests } } +// TestFilterGetLogsMatchedLogCap exercises the eth_getLogs matched-log cap +// end-to-end over the live JSON-RPC handler. +// This test is intentionally not parallel: it runs during the serial phase, so +// its cache warm-up below completes before the parallel filter tests resume. +func TestFilterGetLogsMatchedLogCap(t *testing.T) { + // Warm the shared block cache with MockHeight8's canonical block via an + // address-less number-range query (no address ⇒ no bloom pre-filter, so the + // block is fetched and cached). + warmCriteria := map[string]interface{}{ + "fromBlock": "0x8", + "toBlock": "0x8", + } + warmRes := sendRequestGood(t, "getLogs", warmCriteria) + _, warmErr := warmRes["error"] + require.False(t, warmErr, "warm-up number query should not error, got: %v", warmRes) + + filterCriteria := map[string]interface{}{ + "blockHash": LogCapBlockHash, + "address": []common.Address{common.HexToAddress(LogCapAddr)}, + } + resObj := sendRequestGood(t, "getLogs", filterCriteria) + + _, hasResult := resObj["result"] + require.False(t, hasResult, "expected no result once the log cap is exceeded, got: %v", resObj["result"]) + errObj, ok := resObj["error"].(map[string]interface{}) + require.True(t, ok, "expected an error object, got: %v", resObj) + require.Contains(t, errObj["message"].(string), "too many logs") +} + func TestFilterGetFilterLogs(t *testing.T) { t.Skip() filterCriteria := map[string]interface{}{ diff --git a/evmrpc/setup_test.go b/evmrpc/setup_test.go index 1f2c093a17..6fec0c2aef 100644 --- a/evmrpc/setup_test.go +++ b/evmrpc/setup_test.go @@ -110,6 +110,16 @@ var multiTxBlockSynthTx *ethtypes.Transaction var Block100NormalTx sdk.Tx var block100NormalTx *ethtypes.Transaction +// Single-tx block whose receipt carries >MaxLogNoBlock matching logs. +var LogCapTx sdk.Tx +var logCapTx *ethtypes.Transaction + +var LogCapBlockHash = "0x00000000000000000000000000000000000000000000000000000000000000ca" + +// LogCapAddr is matched only by the logs in the LogCapBlockHash block, so the +// cap test's filter is unaffected by (and does not affect) any other fixture. +const LogCapAddr = "0x11111111111111111111111111111111111111ca" + var DebugTraceTx sdk.Tx var DebugTracePanicTx sdk.Tx var DebugTraceNonPanicTx sdk.Tx @@ -131,6 +141,27 @@ var MockBlockIDMultiTx = tmtypes.BlockID{ Hash: bytes.HexBytes(mustHexToBytes(MultiTxBlockHash[2:])), } +// mockLogCapBlock returns the single-tx block behind LogCapBlockHash, served +// only via BlockByHash. It reports MockHeight8 to clear the latest-height +// watermark. +func mockLogCapBlock() *coretypes.ResultBlock { + return &coretypes.ResultBlock{ + BlockID: tmtypes.BlockID{Hash: bytes.HexBytes(mustHexToBytes(LogCapBlockHash[2:]))}, + Block: &tmtypes.Block{ + Header: mockBlockHeader(MockHeight8), + Data: tmtypes.Data{ + Txs: []tmtypes.Tx{ + func() []byte { + bz, _ := Encoder(LogCapTx) + return bz + }(), + }, + }, + LastCommit: &tmtypes.Commit{Height: MockHeight8 - 1}, + }, + } +} + var NewHeadsCalled = make(chan struct{}, 1) // NotifierForTest backs the Autobahn-style WS server started on @@ -383,6 +414,9 @@ func (c *MockClient) BlockByHash(_ context.Context, hash bytes.HexBytes) (*coret if hash.String() == MultiTxBlockHash[2:] { return c.mockBlock(MockHeight2), nil } + if strings.EqualFold(hash.String(), LogCapBlockHash[2:]) { + return mockLogCapBlock(), nil + } if strings.ToLower(hash.String()) == "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" { // Match real Tendermint behavior for unknown hashes: ResultBlock with // Block: nil + no error. blockByHashWithRetry wraps this as @@ -813,6 +847,18 @@ func generateTxData() { Data: []byte("abc"), ChainID: chainId, }) + // Build the single tx that holds the over-cap log set for the + // LogCapBlockHash block. + var logCapTxBuilder client.TxBuilder + logCapTxBuilder, logCapTx = buildTx(ethtypes.DynamicFeeTx{ + Nonce: 8, + GasFeeCap: big.NewInt(10), + Gas: 1000, + To: &to, + Value: big.NewInt(1000), + Data: []byte("logcap"), + ChainID: chainId, + }) debugTraceTxBuilder, debugTraceEthTx := buildTx(ethtypes.DynamicFeeTx{ Nonce: 0, GasFeeCap: big.NewInt(1000000000), @@ -856,6 +902,7 @@ func generateTxData() { MultiTxBlockTx4 = txBuilder4.GetTx() MultiTxBlockSynthTx = synthTxBuilder.GetTx() Block100NormalTx = block100TxBuilder.GetTx() + LogCapTx = logCapTxBuilder.GetTx() DebugTraceTx = debugTraceTxBuilder.GetTx() DebugTracePanicTx = debugTracePanicTxBuilder.GetTx() panicEthTx, _ := DebugTracePanicTx.GetMsgs()[0].(*types.MsgEVMTransaction).AsTransaction() @@ -1097,6 +1144,33 @@ func setupLogs() { GasUsed: 21000, EffectiveGasPrice: 100, }) + // LogCapBlockHash block (served only via BlockByHash): a single receipt + // whose 11 logs all match LogCapAddr — one more than the test server's + // MaxLogNoBlock (=10) + capLogs := make([]*types.Log, 0, 11) + ethCapLogs := make([]*ethtypes.Log, 0, 11) + capTopic := common.HexToHash(LogCapBlockHash) + for i := 0; i < 11; i++ { + capLogs = append(capLogs, &types.Log{ + Address: LogCapAddr, + Topics: []string{capTopic.Hex()}, + }) + ethCapLogs = append(ethCapLogs, ðtypes.Log{ + Address: common.HexToAddress(LogCapAddr), + Topics: []common.Hash{capTopic}, + }) + } + bloomCap := ethtypes.CreateBloom(ðtypes.Receipt{Logs: ethCapLogs}) + CtxLogCap := Ctx.WithBlockHeight(MockHeight8) + EVMKeeper.MockReceipt(CtxLogCap, logCapTx.Hash(), &types.Receipt{ + BlockNumber: MockHeight8, + TransactionIndex: 0, + TxHashHex: logCapTx.Hash().Hex(), + LogsBloom: bloomCap[:], + Logs: capLogs, + GasUsed: 21000, + EffectiveGasPrice: 100, + }) CtxMock = Ctx.WithBlockHeight(MockHeight103) EVMKeeper.MockReceipt(CtxMock, common.HexToHash(TestSyntheticTxHash), &types.Receipt{ TxType: types.ShellEVMTxType, diff --git a/evmrpc/watermark_manager_test.go b/evmrpc/watermark_manager_test.go index 0d6c3b6cf1..3c68b07e6b 100644 --- a/evmrpc/watermark_manager_test.go +++ b/evmrpc/watermark_manager_test.go @@ -313,7 +313,7 @@ func (f *fakeReceiptStore) SetReceipts(sdk.Context, []receipt.ReceiptRecord) err return nil } -func (f *fakeReceiptStore) FilterLogs(sdk.Context, uint64, uint64, filters.FilterCriteria) ([]*ethtypes.Log, error) { +func (f *fakeReceiptStore) FilterLogs(sdk.Context, uint64, uint64, filters.FilterCriteria, int64) ([]*ethtypes.Log, error) { return nil, receipt.ErrRangeQueryNotSupported } diff --git a/sei-db/ledger_db/receipt/litt_receipt_store.go b/sei-db/ledger_db/receipt/litt_receipt_store.go index 0526bd5c89..349f0184d8 100644 --- a/sei-db/ledger_db/receipt/litt_receipt_store.go +++ b/sei-db/ledger_db/receipt/litt_receipt_store.go @@ -337,11 +337,11 @@ func (s *littReceiptStore) nextPartIndex(blockNumber uint64) (uint32, error) { // FilterLogs answers eth_getLogs via the tag index. Both bounds inclusive; for // a single block set fromBlock == toBlock. -func (s *littReceiptStore) FilterLogs(_ sdk.Context, fromBlock, toBlock uint64, crit filters.FilterCriteria) ([]*ethtypes.Log, error) { +func (s *littReceiptStore) FilterLogs(_ sdk.Context, fromBlock, toBlock uint64, crit filters.FilterCriteria, limit int64) ([]*ethtypes.Log, error) { if fromBlock > toBlock { return nil, fmt.Errorf("fromBlock (%d) > toBlock (%d)", fromBlock, toBlock) } - return s.filterLogsByTags(fromBlock, toBlock, crit) + return s.filterLogsByTags(fromBlock, toBlock, crit, limit) } // startFlusher bounds litt durability lag to littFlushInterval from a diff --git a/sei-db/ledger_db/receipt/litt_tag_index.go b/sei-db/ledger_db/receipt/litt_tag_index.go index d715d4847c..e42d89000e 100644 --- a/sei-db/ledger_db/receipt/litt_tag_index.go +++ b/sei-db/ledger_db/receipt/litt_tag_index.go @@ -1,9 +1,11 @@ package receipt import ( + "context" "encoding/binary" "fmt" "sort" + "sync/atomic" "github.com/ethereum/go-ethereum/common" ethtypes "github.com/ethereum/go-ethereum/core/types" @@ -193,7 +195,7 @@ func (s *littReceiptStore) stageTagKeys(batch dbtypes.Batch, blockNumber uint64, // This parallelizes both the per-block index scans and the litt body reads, // which dominate wide-range latency. Results are exact (matchLog re-verifies // after decode) and stay in (block, txIndex) order via the indexed buffer. -func (s *littReceiptStore) filterLogsByTags(fromBlock, toBlock uint64, crit filters.FilterCriteria) ([]*ethtypes.Log, error) { +func (s *littReceiptStore) filterLogsByTags(fromBlock, toBlock uint64, crit filters.FilterCriteria, limit int64) ([]*ethtypes.Log, error) { if latest := s.latestVersion.Load(); latest >= 0 && toBlock > uint64(latest) { //nolint:gosec // latest is non-negative toBlock = uint64(latest) //nolint:gosec // latest is non-negative } @@ -211,16 +213,32 @@ func (s *littReceiptStore) filterLogsByTags(fromBlock, toBlock uint64, crit filt // output keeps block order regardless of completion order. errgroup caps // concurrency at the configured limit and hands the next block to whichever // worker frees up, load-balancing across the skewed per-block cost. + // + // When limit > 0 a running counter aborts the fan-out once the matched logs + // exceed the cap: the tripping worker returns ErrTooManyLogs, which cancels + // the group so already-scheduled blocks bail before reading, and the loop + // stops queueing new blocks. Peak memory is thus bounded to the cap plus at + // most one in-flight block per worker, rather than O(total matching logs). results := make([][]*ethtypes.Log, nBlocks) - var eg errgroup.Group + var collected int64 + eg, egCtx := errgroup.WithContext(context.Background()) eg.SetLimit(s.logFilterParallelism) for i := 0; i < nBlocks; i++ { + if egCtx.Err() != nil { + break + } eg.Go(func() error { + if egCtx.Err() != nil { + return nil // group already cancelled; skip without overwriting the cause + } blockLogs, err := s.blockLogs(fromBlock+uint64(i), groups, crit) //nolint:gosec // i < nBlocks if err != nil { return err } results[i] = blockLogs + if limit > 0 && atomic.AddInt64(&collected, int64(len(blockLogs))) > limit { + return NewTooManyLogsError(limit) + } return nil }) } @@ -228,7 +246,7 @@ func (s *littReceiptStore) filterLogsByTags(fromBlock, toBlock uint64, crit filt return nil, err } - var logs []*ethtypes.Log + logs := make([]*ethtypes.Log, 0, collected) for _, blockLogs := range results { logs = append(logs, blockLogs...) } diff --git a/sei-db/ledger_db/receipt/littidx_test.go b/sei-db/ledger_db/receipt/littidx_test.go index 11c2ab568e..aa5909eef8 100644 --- a/sei-db/ledger_db/receipt/littidx_test.go +++ b/sei-db/ledger_db/receipt/littidx_test.go @@ -108,39 +108,72 @@ func TestLittIdxFilterLogs(t *testing.T) { writeLitBlock(t, store, ctx, 3, litReceipt(3, 0, addr1, approve, bob)) // Address OR: either address matches. - logs, err := store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Addresses: []common.Address{addr1, addr2}}) + logs, err := store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Addresses: []common.Address{addr1, addr2}}, 0) require.NoError(t, err) require.Len(t, logs, 4) // AND across topic positions: transfer at 0 AND alice at 1. - logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Topics: [][]common.Hash{{transfer}, {alice}}}) + logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Topics: [][]common.Hash{{transfer}, {alice}}}, 0) require.NoError(t, err) require.Len(t, logs, 1) require.Equal(t, uint64(1), logs[0].BlockNumber) require.Equal(t, uint(0), logs[0].TxIndex) // OR within a topic position: alice or bob at position 1. - logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Topics: [][]common.Hash{{transfer}, {alice, bob}}}) + logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Topics: [][]common.Hash{{transfer}, {alice, bob}}}, 0) require.NoError(t, err) require.Len(t, logs, 2) // Wildcard position 0 (empty), bob at position 1. - logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Topics: [][]common.Hash{{}, {bob}}}) + logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Topics: [][]common.Hash{{}, {bob}}}, 0) require.NoError(t, err) require.Len(t, logs, 2) // Range bound: only block 2. - logs, err = store.FilterLogs(ctx, 2, 2, filters.FilterCriteria{}) + logs, err = store.FilterLogs(ctx, 2, 2, filters.FilterCriteria{}, 0) require.NoError(t, err) require.Len(t, logs, 1) require.Equal(t, uint64(2), logs[0].BlockNumber) // No match. - logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Addresses: []common.Address{common.HexToAddress("0x9")}}) + logs, err = store.FilterLogs(ctx, 1, 3, filters.FilterCriteria{Addresses: []common.Address{common.HexToAddress("0x9")}}, 0) require.NoError(t, err) require.Empty(t, logs) } +// TestLittIdxFilterLogsLimit: FilterLogs aborts with ErrTooManyLogs once the +// matched logs exceed a positive limit, and returns the full result when the +// match count is at or below it. +func TestLittIdxFilterLogsLimit(t *testing.T) { + store, ctx := setupLittIdx(t, t.TempDir()) + defer func() { _ = store.Close() }() + + addr := common.HexToAddress("0x4444444444444444444444444444444444444444") + topic := common.HexToHash("0x1111aa") + // 3 blocks x 2 logs = 6 matching logs. + for block := uint64(1); block <= 3; block++ { + writeLitBlock(t, store, ctx, block, + litReceipt(block, 0, addr, topic), + litReceipt(block, 1, addr, topic), + ) + } + crit := filters.FilterCriteria{Addresses: []common.Address{addr}} + + // Exactly at the limit: no error, all logs returned. + logs, err := store.FilterLogs(ctx, 1, 3, crit, 6) + require.NoError(t, err) + require.Len(t, logs, 6) + + // Over the limit: ErrTooManyLogs. + _, err = store.FilterLogs(ctx, 1, 3, crit, 5) + require.ErrorIs(t, err, receipt.ErrTooManyLogs) + + // limit <= 0 disables the cap. + logs, err = store.FilterLogs(ctx, 1, 3, crit, 0) + require.NoError(t, err) + require.Len(t, logs, 6) +} + // TestLittIdxMultiLogReceipt: a receipt with several logs is numbered with // contiguous block-wide log indices. func TestLittIdxMultiLogReceipt(t *testing.T) { @@ -162,7 +195,7 @@ func TestLittIdxMultiLogReceipt(t *testing.T) { } writeLitBlock(t, store, ctx, 7, rec) - logs, err := store.FilterLogs(ctx, 7, 7, filters.FilterCriteria{Addresses: []common.Address{addr}}) + logs, err := store.FilterLogs(ctx, 7, 7, filters.FilterCriteria{Addresses: []common.Address{addr}}, 0) require.NoError(t, err) require.Len(t, logs, 3) for i, lg := range logs { @@ -192,7 +225,7 @@ func TestLittIdxBlockWideLogIndex(t *testing.T) { // tx0: 2 logs, tx1: 1 log, tx2: 3 logs -> contiguous block-wide indices 0..5. writeLitBlock(t, store, ctx, 8, mk(0, 2), mk(1, 1), mk(2, 3)) - logs, err := store.FilterLogs(ctx, 8, 8, filters.FilterCriteria{Addresses: []common.Address{addr}}) + logs, err := store.FilterLogs(ctx, 8, 8, filters.FilterCriteria{Addresses: []common.Address{addr}}, 0) require.NoError(t, err) require.Len(t, logs, 6) for i, lg := range logs { @@ -224,7 +257,7 @@ func TestLittIdxMultiPart(t *testing.T) { require.Equal(t, txIndex, rcpt.TransactionIndex) } - logs, err := store.FilterLogs(ctx, 4, 4, filters.FilterCriteria{Addresses: []common.Address{addr}}) + logs, err := store.FilterLogs(ctx, 4, 4, filters.FilterCriteria{Addresses: []common.Address{addr}}, 0) require.NoError(t, err) require.Len(t, logs, 2) } @@ -277,7 +310,7 @@ func TestLittIdxReopen(t *testing.T) { require.NoError(t, err) require.Equal(t, uint64(3), rcpt.BlockNumber) - logs, err := store.FilterLogs(ctx, 1, 5, filters.FilterCriteria{Addresses: []common.Address{addr}}) + logs, err := store.FilterLogs(ctx, 1, 5, filters.FilterCriteria{Addresses: []common.Address{addr}}, 0) require.NoError(t, err) require.Len(t, logs, 5) } @@ -301,7 +334,7 @@ func TestLittIdxPrune(t *testing.T) { _, err := store.GetReceiptFromStore(ctx, litTxHash(block, 0)) require.ErrorIs(t, err, receipt.ErrNotFound, "block %d should be pruned", block) } - logs, err := store.FilterLogs(ctx, 1, 5, filters.FilterCriteria{Addresses: []common.Address{addr}}) + logs, err := store.FilterLogs(ctx, 1, 5, filters.FilterCriteria{Addresses: []common.Address{addr}}, 0) require.NoError(t, err) require.Empty(t, logs) @@ -311,7 +344,7 @@ func TestLittIdxPrune(t *testing.T) { require.NoError(t, err) require.Equal(t, block, rcpt.BlockNumber) } - logs, err = store.FilterLogs(ctx, 1, 10, filters.FilterCriteria{Addresses: []common.Address{addr}}) + logs, err = store.FilterLogs(ctx, 1, 10, filters.FilterCriteria{Addresses: []common.Address{addr}}, 0) require.NoError(t, err) require.Len(t, logs, 5) } @@ -341,7 +374,7 @@ func TestLittIdxFilterLogsParallelOrder(t *testing.T) { writeLitBlock(t, store, ctx, b, recs...) } - logs, err := store.FilterLogs(ctx, 1, blocks, filters.FilterCriteria{Addresses: []common.Address{addr}}) + logs, err := store.FilterLogs(ctx, 1, blocks, filters.FilterCriteria{Addresses: []common.Address{addr}}, 0) require.NoError(t, err) got := make([]string, len(logs)) for i, lg := range logs { diff --git a/sei-db/ledger_db/receipt/receipt_bench_read_test.go b/sei-db/ledger_db/receipt/receipt_bench_read_test.go index e2143acf0f..3ea6877bdf 100644 --- a/sei-db/ledger_db/receipt/receipt_bench_read_test.go +++ b/sei-db/ledger_db/receipt/receipt_bench_read_test.go @@ -216,7 +216,7 @@ func execFilterLogs(env *readBenchEnv, fromBlock, toBlock uint64, crit filters.F if env.isPebble { return pebbleFilterLogs(env.store, env.ctx, fromBlock, toBlock, env.idx, crit) } - return env.store.FilterLogs(env.ctx, fromBlock, toBlock, crit) + return env.store.FilterLogs(env.ctx, fromBlock, toBlock, crit, 0) } // readBenchEnv holds everything produced by the write phase that the read diff --git a/sei-db/ledger_db/receipt/receipt_store.go b/sei-db/ledger_db/receipt/receipt_store.go index e4da39aaa8..63697f7017 100644 --- a/sei-db/ledger_db/receipt/receipt_store.go +++ b/sei-db/ledger_db/receipt/receipt_store.go @@ -30,8 +30,16 @@ var ( ErrNotFound = errors.New("receipt not found") ErrNotConfigured = errors.New("receipt store not configured") ErrRangeQueryNotSupported = errors.New("range query not supported by this backend") + // ErrTooManyLogs is returned by FilterLogs when a query matches more logs + // than the caller-supplied limit. It lets callers cap peak memory by + // aborting a query instead of materializing an unbounded result set. + ErrTooManyLogs = errors.New("query matches too many logs") ) +func NewTooManyLogsError(limit int64) error { + return fmt.Errorf("%w: result exceeds the maximum of %d logs; narrow the block range or filter criteria", ErrTooManyLogs, limit) +} + // ReceiptStore exposes receipt-specific operations without leaking the StateStore interface. type ReceiptStore interface { LatestVersion() int64 @@ -43,7 +51,9 @@ type ReceiptStore interface { SetReceipts(ctx sdk.Context, receipts []ReceiptRecord) error // FilterLogs queries logs across a range of blocks. // For single-block queries, set fromBlock == toBlock. - FilterLogs(ctx sdk.Context, fromBlock, toBlock uint64, crit filters.FilterCriteria) ([]*ethtypes.Log, error) + // When limit > 0 the query aborts and returns ErrTooManyLogs once more than + // limit logs match, bounding peak memory; limit <= 0 means no cap. + FilterLogs(ctx sdk.Context, fromBlock, toBlock uint64, crit filters.FilterCriteria, limit int64) ([]*ethtypes.Log, error) Close() error } @@ -265,7 +275,7 @@ func (s *receiptStore) SetReceipts(ctx sdk.Context, receipts []ReceiptRecord) er // FilterLogs is not efficiently supported by the pebble backend since receipts // are indexed by tx hash, not by block number. Returns ErrRangeQueryNotSupported. // Callers should fall back to fetching receipts individually via GetReceipt. -func (s *receiptStore) FilterLogs(_ sdk.Context, _, _ uint64, _ filters.FilterCriteria) ([]*ethtypes.Log, error) { +func (s *receiptStore) FilterLogs(_ sdk.Context, _, _ uint64, _ filters.FilterCriteria, _ int64) ([]*ethtypes.Log, error) { return nil, ErrRangeQueryNotSupported } diff --git a/sei-db/ledger_db/receipt/receipt_store_test.go b/sei-db/ledger_db/receipt/receipt_store_test.go index 26c65a24b6..b061bc7c22 100644 --- a/sei-db/ledger_db/receipt/receipt_store_test.go +++ b/sei-db/ledger_db/receipt/receipt_store_test.go @@ -170,7 +170,7 @@ func TestReceiptStorePebbleBackendBasic(t *testing.T) { logs, err := store.FilterLogs(ctx, 1, 1, filters.FilterCriteria{ Addresses: []common.Address{addr}, Topics: [][]common.Hash{{topic}}, - }) + }, 0) require.ErrorIs(t, err, receipt.ErrRangeQueryNotSupported) require.Empty(t, logs) } @@ -178,7 +178,7 @@ func TestReceiptStorePebbleBackendBasic(t *testing.T) { func TestFilterLogsRangeQueryNotSupported(t *testing.T) { store, ctx, _ := setupReceiptStore(t) // Pebble backend does not support range queries, so FilterLogs returns ErrRangeQueryNotSupported. - _, err := store.FilterLogs(ctx, 1, 10, filters.FilterCriteria{}) + _, err := store.FilterLogs(ctx, 1, 10, filters.FilterCriteria{}, 0) require.ErrorIs(t, err, receipt.ErrRangeQueryNotSupported) } diff --git a/sei-db/state_db/bench/cryptosim/reciept_store_simulator.go b/sei-db/state_db/bench/cryptosim/reciept_store_simulator.go index 7056d013fe..5ede2d08e9 100644 --- a/sei-db/state_db/bench/cryptosim/reciept_store_simulator.go +++ b/sei-db/state_db/bench/cryptosim/reciept_store_simulator.go @@ -411,7 +411,7 @@ func (r *RecieptStoreSimulator) executeLogFilterRead(crand *crand.CannedRandom) sdkCtx := sdk.NewContext(nil, tmproto.Header{}, false) start := time.Now() - logs, err := r.store.FilterLogs(sdkCtx, fromBlock, toBlock, crit) + logs, err := r.store.FilterLogs(sdkCtx, fromBlock, toBlock, crit, 0) r.metrics.RecordReceiptLogFilterDuration(time.Since(start).Seconds()) r.metrics.RecordLogFilterLogsReturned(int64(len(logs)))