Improve chain selection for unordered batch creation - #683
Conversation
WalkthroughRename FetchState.queueSize → bufferSize, add Array.lastUnsafe and FetchState helpers for unordered-batch selection/sorting, switch ChainManager to prepare/filter/sort fetch-states for single-pass unordered batching, and update tests/comments to reference bufferSize. Changes
Sequence Diagram(s)sequenceDiagram
participant CM as ChainManager
participant FS as FetchState
participant M as Metrics
participant B as BatchConsumer
CM->>FS: filterAndSortForUnorderedBatch(fetchStates)
FS-->>CM: preparedFetchStates (filtered & sorted by last-item timestamp)
loop single-pass over preparedFetchStates until maxBatchSize
CM->>FS: pop items from fetchState.queue
CM->>M: record per-chain metrics using fetchState.chainId and bufferSize
end
CM-->>B: emit unordered batch
Estimated code review effort🎯 4 (Complex) | ⏱️ ~40 minutes Possibly related PRs
Suggested reviewers
Poem
✨ Finishing Touches🧪 Generate unit tests
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. 🪧 TipsChatThere are 3 ways to chat with CodeRabbit:
SupportNeed help? Create a ticket on our support page for assistance with any issues or questions. CodeRabbit Commands (Invoked using PR/Issue comments)Type Other keywords and placeholders
CodeRabbit Configuration File (
|
There was a problem hiding this comment.
Actionable comments posted: 2
🧹 Nitpick comments (2)
codegenerator/cli/npm/envio/src/Utils.res (1)
253-254: Document the non-empty precondition for Array.lastUnsafe (and consider fail-fast).This helper is intentionally unsafe and will read undefined for empty arrays (index -1). Add a concise doc comment to make the contract explicit and avoid accidental misuse.
Apply this diff to add documentation:
let last = (arr: array<'a>): option<'a> => arr->Belt.Array.get(arr->Array.length - 1) + /** + UNSAFE: Expects a non-empty array. + Prefer `last` when the array may be empty. + */ let lastUnsafe = (arr: array<'a>): 'a => arr->Belt.Array.getUnsafe(arr->Array.length - 1)codegenerator/cli/npm/envio/src/FetchState.res (1)
1237-1241: Add deterministic tie-breakers in comparator for equal timestamps.Current comparator returns 0 when timestamps are equal, which can lead to non-deterministic ordering across runs. Consider tie-breaking on blockNumber, then logIndex, then chainId for stability.
Apply this diff to strengthen ordering:
-let compareUnorderedBatchChainPriority = (a: t, b: t) => { - // Use unsafe since we filtered out all queues without batch items - (a.queue->Utils.Array.lastUnsafe).timestamp - (b.queue->Utils.Array.lastUnsafe).timestamp -} +let compareUnorderedBatchChainPriority = (a: t, b: t) => { + // Use unsafe since we filtered out all queues without batch items + let ai = a.queue->Utils.Array.lastUnsafe + let bi = b.queue->Utils.Array.lastUnsafe + if ai.timestamp !== bi.timestamp { + ai.timestamp - bi.timestamp + } else if ai.blockNumber !== bi.blockNumber { + ai.blockNumber - bi.blockNumber + } else if ai.logIndex !== bi.logIndex { + ai.logIndex - bi.logIndex + } else { + a.chainId - b.chainId + } +}
📜 Review details
Configuration used: CodeRabbit UI
Review profile: CHILL
Plan: Pro
📒 Files selected for processing (8)
codegenerator/cli/npm/envio/src/FetchState.res(4 hunks)codegenerator/cli/npm/envio/src/Utils.res(1 hunks)codegenerator/cli/templates/static/codegen/src/eventFetching/ChainFetcher.res(1 hunks)codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res(2 hunks)scenarios/erc20_multichain_factory/test/RollbackDynamicContract_test.res(1 hunks)scenarios/erc20_multichain_factory/test/RollbackMultichain_test.res(1 hunks)scenarios/test_codegen/test/ChainManager_test.res(1 hunks)scenarios/test_codegen/test/rollback/Rollback_test.res(2 hunks)
🧰 Additional context used
📓 Path-based instructions (4)
**/*.{res,resi}
📄 CodeRabbit Inference Engine (.cursor/rules/rescript.mdc)
**/*.{res,resi}: Never use[| item |]to create an array. Use[ item ]instead.
Must always use=for setting value to a field. Use:=only for ref values created usingreffunction.
ReScript has record types which require a type definition before hand. You can access record fields by dot likefoo.myField.
It's also possible to define an inline object, it'll have quoted fields in this case.
Use records when working with structured data, and objects to conveniently pass payload data between functions.
Never use %raw to access object fields if you know the type.
Files:
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainFetcher.resscenarios/test_codegen/test/ChainManager_test.resscenarios/test_codegen/test/rollback/Rollback_test.resscenarios/erc20_multichain_factory/test/RollbackMultichain_test.rescodegenerator/cli/npm/envio/src/Utils.resscenarios/erc20_multichain_factory/test/RollbackDynamicContract_test.rescodegenerator/cli/npm/envio/src/FetchState.rescodegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res
codegenerator/cli/templates/{dynamic/**/*.hbs,static/**}
📄 CodeRabbit Inference Engine (.cursor/rules/navigation.mdc)
Templates live under codegenerator/cli/templates: dynamic/ for Handlebars (.hbs), static/ for raw Rescript files copied verbatim.
Files:
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainFetcher.rescodegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res
{**/generated/src/**/*.res,codegenerator/cli/templates/static/codegen/src/**/*.res,codegenerator/cli/templates/dynamic/codegen/src/**/*.res}
📄 CodeRabbit Inference Engine (.cursor/rules/navigation.mdc)
Runtime code lives in each project’s generated/src, but template versions (good for editing) are under codegenerator/cli/templates/static/codegen/src or codegenerator/cli/templates/dynamic/codegen/src.
Files:
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainFetcher.rescodegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res
codegenerator/cli/npm/envio/**
📄 CodeRabbit Inference Engine (.cursor/rules/navigation.mdc)
Library-fied runtime shared across indexers lives in codegenerator/cli/npm/envio.
Files:
codegenerator/cli/npm/envio/src/Utils.rescodegenerator/cli/npm/envio/src/FetchState.res
⏰ Context from checks skipped due to timeout of 90000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (1)
- GitHub Check: build_and_test
🔇 Additional comments (13)
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainFetcher.res (1)
500-501: LGTM: rename to bufferSize is correct and consistent with FetchState API.hasNoMoreEventsToProcess now correctly uses FetchState.bufferSize.
scenarios/erc20_multichain_factory/test/RollbackDynamicContract_test.res (1)
182-183: LGTM: comment-only rename to bufferSize.Matches the updated FetchState API and keeps tests consistent with the new semantics.
scenarios/test_codegen/test/ChainManager_test.res (1)
210-211: LGTM: bufferSize aggregation reflects the new API.Summing bufferSize across fetch states is the correct replacement for queueSize.
scenarios/test_codegen/test/rollback/Rollback_test.res (1)
193-194: LGTM: assertions updated to bufferSize.Both assertions correctly align with FetchState.bufferSize after the API rename.
Also applies to: 259-260
scenarios/erc20_multichain_factory/test/RollbackMultichain_test.res (1)
269-276: Rename alignment to bufferSize looks correct.The sum now reflects FetchState.bufferSize as intended by the PR-wide rename.
codegenerator/cli/npm/envio/src/FetchState.res (5)
1089-1089: Public helper rename to bufferSize is consistent.This mirrors the internal change and aligns external call sites with the new naming.
1205-1214: End-block active-indexing check now correctly depends on buffered items.Using fetchState->bufferSize > 0 when past endBlock is a sensible behavior change to drain remaining items before declaring inactivity.
1229-1236: hasBatchItem aligns with getEarliestEvent semantics.Filtering only states whose earliest item is not ahead of latestFullyFetchedBlock is the right precondition for unordered batching.
1242-1246: Filtering and sorting pipeline is correct and efficient.In-place sort on a fresh values array avoids mutating chain state, and the pre-filter prevents unsafe access.
261-267: Metrics update verified:IndexingBufferSizeexists and is used consistently
- Defined in
codegenerator/cli/npm/envio/src/Prometheus.res(moduleIndexingBufferSizewith gauge name"envio_indexing_buffer_size")- Invoked in
FetchState.resat lines 266 and 1064- No remaining references to the legacy
IndexingQueueSizePlease ensure dashboards and alerts have been updated to use the new
envio_indexing_buffer_sizemetric.codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res (3)
185-189: Pre-filtering/sorting fetch states for unordered batch creation is a solid simplification.This prepares a single pass over relevant chains and aligns with the new priority strategy favoring behind chains.
218-226: Per-chain metrics keyed by fetchState.chainId are correct.Using chainId directly avoids reliance on chain objects in this path and matches the new FetchState API. The targetBlockNumber derivation is safe here since the last item just appended belongs to the same chain.
270-273: Batch buffer reporting updated to bufferSize.The logged buffers now reflect the renamed API and will match Prometheus metrics for consistency.
There was a problem hiding this comment.
Actionable comments posted: 0
♻️ Duplicate comments (1)
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res (1)
198-199: Follow-up: previous 'sortedFetchStates' identifier is resolvedThe earlier undefined reference is now correctly replaced with preparedFetchStates. No further action needed.
🧹 Nitpick comments (4)
scenarios/test_codegen/test/lib_tests/FetchState_test.res (2)
2689-2693: Prefer helper for consistent block timestampsFor consistency with the rest of the tests, consider using getBlockData(~blockNumber=latestBlock) to set latestFetchedBlock (timestamp aligned with blockNumber*15), instead of manually mirroring blockNumber into blockTimestamp.
Apply:
- ~latestFetchedBlock={blockNumber: latestBlock, blockTimestamp: latestBlock}, + ~latestFetchedBlock=getBlockData(~blockNumber=latestBlock),
2706-2710: Use lastUnsafe to avoid Option plumbing in testSince the array is guaranteed non-empty in this test, you can simplify by using a dedicated unsafe accessor instead of Option.getUnsafe.
Apply:
- prepared->Array.map(fs => (fs.queue->Utils.Array.last->Option.getUnsafe).blockNumber), + prepared->Array.map(fs => (fs.queue->Utils.Array.lastUnsafe).blockNumber),codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res (2)
197-215: Make per-chain targetBlockNumber independent of global items arraySmall robustness tweak: track the last popped blockNumber for the current chain locally instead of relying on items->Utils.Array.last. This decouples metrics from global state and avoids Option.getUnsafe usage here.
Apply:
- while batchSize.contents < maxBatchSize && idx.contents < preparedNumber { + while batchSize.contents < maxBatchSize && idx.contents < preparedNumber { let fetchState = preparedFetchStates->Array.getUnsafe(idx.contents) let batchSizeBeforeTheChain = batchSize.contents + let lastBlockNumberForChain = ref(-1) let rec loop = () => if batchSize.contents < maxBatchSize { let earliestEvent = fetchState->FetchState.getEarliestEvent switch earliestEvent { | NoItem(_) => () | Item({item, popItemOffQueue}) => { popItemOffQueue() items->Js.Array2.push(item)->ignore batchSize := batchSize.contents + 1 + lastBlockNumberForChain := item.blockNumber loop() } } } loop() let chainBatchSize = batchSize.contents - batchSizeBeforeTheChain if chainBatchSize > 0 { mutProcessingMetricsByChainId->Js.Dict.set( fetchState.chainId->Int.toString, { batchSize: chainBatchSize, // If there's the chainBatchSize, // then it's guaranteed that the last item belongs to the chain - targetBlockNumber: (items->Utils.Array.last->Option.getUnsafe).blockNumber, + targetBlockNumber: lastBlockNumberForChain.contents, }, ) } idx := idx.contents + 1 }
194-197: Polish comment grammar for clarityMinor wording improvement.
Apply:
- // Accumulate items for all actively indexing chains - // the way to group as many items from a single chain as possible - // This way the loaders optimisations will hit more often + // Accumulate items across prepared chains, grouping as many items per chain as possible. + // This increases the hit rate of downstream loader optimizations.
📜 Review details
Configuration used: CodeRabbit UI
Review profile: CHILL
Plan: Pro
📒 Files selected for processing (2)
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res(2 hunks)scenarios/test_codegen/test/lib_tests/FetchState_test.res(1 hunks)
🧰 Additional context used
📓 Path-based instructions (3)
**/*.{res,resi}
📄 CodeRabbit Inference Engine (.cursor/rules/rescript.mdc)
**/*.{res,resi}: Never use[| item |]to create an array. Use[ item ]instead.
Must always use=for setting value to a field. Use:=only for ref values created usingreffunction.
ReScript has record types which require a type definition before hand. You can access record fields by dot likefoo.myField.
It's also possible to define an inline object, it'll have quoted fields in this case.
Use records when working with structured data, and objects to conveniently pass payload data between functions.
Never use %raw to access object fields if you know the type.
Files:
scenarios/test_codegen/test/lib_tests/FetchState_test.rescodegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res
codegenerator/cli/templates/{dynamic/**/*.hbs,static/**}
📄 CodeRabbit Inference Engine (.cursor/rules/navigation.mdc)
Templates live under codegenerator/cli/templates: dynamic/ for Handlebars (.hbs), static/ for raw Rescript files copied verbatim.
Files:
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res
{**/generated/src/**/*.res,codegenerator/cli/templates/static/codegen/src/**/*.res,codegenerator/cli/templates/dynamic/codegen/src/**/*.res}
📄 CodeRabbit Inference Engine (.cursor/rules/navigation.mdc)
Runtime code lives in each project’s generated/src, but template versions (good for editing) are under codegenerator/cli/templates/static/codegen/src or codegenerator/cli/templates/dynamic/codegen/src.
Files:
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res
⏰ Context from checks skipped due to timeout of 90000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (1)
- GitHub Check: build_and_test
🔇 Additional comments (4)
scenarios/test_codegen/test/lib_tests/FetchState_test.res (1)
2669-2712: Good coverage for unordered-batch pre-filter/sort behaviorThe test validates both exclusion (>latestFullyFetchedBlock) and ordering by earliest-eligible last-queue block. It exercises the public API path end-to-end. LGTM.
codegenerator/cli/templates/static/codegen/src/eventFetching/ChainManager.res (3)
185-193: Correct: pre-filter and pre-sort fetch-states for unordered batchesUsing FetchState.filterAndSortForUnorderedBatch over ChainMap.values is the right entry-point and prevents wasted scans. This aligns with the PR’s goal to prioritize chains that are furthest behind.
219-226: Ensure metrics keys are consistent across ordered/unordered pathsIn unordered, the metrics key uses fetchState.chainId->Int.toString; in ordered, it uses item.chain->toChainId->Int.toString. They should resolve to the same string, but please confirm consumers don’t rely on one format or the other. If needed, normalize to a shared helper for chainId-string generation.
Would you like me to search and list all consumers of processingMetricsByChainId to confirm consistent key handling?
272-272: LGTM: bufferSize migrationSwitching to FetchState.bufferSize in buffers reporting is correct and consistent with the rename throughout the codebase.
|
|
||
| let hasBatchItem = ({queue, latestFullyFetchedBlock}: t) => { | ||
| switch queue->Utils.Array.last { | ||
| | Some(item) => item.blockNumber <= latestFullyFetchedBlock.blockNumber |
There was a problem hiding this comment.
In what case is this false? Is it a defensive check?
There was a problem hiding this comment.
After refactoring a few months ago, we now have in a queue all items we fetched and even items for some partitions which are further ahead than the latest fully fetched block.
There was a problem hiding this comment.
Ah right because of the dynamic contracts?
There was a problem hiding this comment.
Yes, or wildcard also creates a new partition
| let compareUnorderedBatchChainPriority = (a: t, b: t) => { | ||
| // Use unsafe since we filtered out all queues without batch items |
There was a problem hiding this comment.
| let compareUnorderedBatchChainPriority = (a: t, b: t) => { | |
| // Use unsafe since we filtered out all queues without batch items | |
| /** FetchStates must have items in their queues or there will be a runtime exn */ | |
| let compareUnorderedBatchChainPriority = (a: t, b: t) => { | |
| // Use unsafe since we filtered out all queues without batch items |
There was a problem hiding this comment.
Sweet, yeah I see now its internal to the function
There was a problem hiding this comment.
Actionable comments posted: 0
🧹 Nitpick comments (1)
codegenerator/cli/npm/envio/src/FetchState.res (1)
1244-1247: Add deterministic tie-breakers to comparator for stability across equal timestampsWhen timestamps are equal, Array.sort may rely on engine stability. Add secondary keys to keep ordering deterministic and avoid churn between runs.
Apply this diff:
- // Use unsafe since we filtered out all queues without batch items - (a.queue->Utils.Array.lastUnsafe).timestamp - (b.queue->Utils.Array.lastUnsafe).timestamp + // Use unsafe since we filtered out all queues without batch items + { + let aItem = a.queue->Utils.Array.lastUnsafe + let bItem = b.queue->Utils.Array.lastUnsafe + let dt = aItem.timestamp - bItem.timestamp + if dt !== 0 { + dt + } else { + let db = aItem.blockNumber - bItem.blockNumber + if db !== 0 { + db + } else { + let dl = aItem.logIndex - bItem.logIndex + if dl !== 0 { + dl + } else { + a.chainId - b.chainId + } + } + } + }
📜 Review details
Configuration used: CodeRabbit UI
Review profile: CHILL
Plan: Pro
📒 Files selected for processing (1)
codegenerator/cli/npm/envio/src/FetchState.res(4 hunks)
🧰 Additional context used
📓 Path-based instructions (2)
**/*.{res,resi}
📄 CodeRabbit Inference Engine (.cursor/rules/rescript.mdc)
**/*.{res,resi}: Never use[| item |]to create an array. Use[ item ]instead.
Must always use=for setting value to a field. Use:=only for ref values created usingreffunction.
ReScript has record types which require a type definition before hand. You can access record fields by dot likefoo.myField.
It's also possible to define an inline object, it'll have quoted fields in this case.
Use records when working with structured data, and objects to conveniently pass payload data between functions.
Never use %raw to access object fields if you know the type.
Files:
codegenerator/cli/npm/envio/src/FetchState.res
codegenerator/cli/npm/envio/**
📄 CodeRabbit Inference Engine (.cursor/rules/navigation.mdc)
Library-fied runtime shared across indexers lives in codegenerator/cli/npm/envio.
Files:
codegenerator/cli/npm/envio/src/FetchState.res
⏰ Context from checks skipped due to timeout of 90000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (1)
- GitHub Check: build_and_test
🔇 Additional comments (5)
codegenerator/cli/npm/envio/src/FetchState.res (5)
261-267: Metric rename and buffer size computation look correctSwitching to bufferSize and wiring Prometheus.IndexingBufferSize.set(~bufferSize, ...) is consistent with the new terminology and intent.
1089-1089: New bufferSize helper is clear and reusableGood extraction. Centralizing this logic will reduce drift and keep naming consistent across the codebase.
1210-1210: Correct gating of “actively indexing” when past endBlockUsing fetchState->bufferSize > 0 once past endBlock matches the expected lifecycle and prevents premature idling.
1231-1236: Filtering ensures lastUnsafe is safehasBatchItem correctly guarantees that queues are non-empty and contain an eligible batch item before sorting, making lastUnsafe safe to use later.
1238-1242: Comparator’s use oftimestamponInternal.eventItemis correct
Confirmed that incodegenerator/cli/npm/envio/src/Internal.res(lines 136–142) theeventItemtype declares atimestamp: intfield. There is noblockTimestamponInternal.eventItem, so the comparator at FetchState.res:1238–1242 may safely use.timestamp. No changes required.
Previously, when we created a batch, we tried to get items from a new chain every time.
The updated logic now prefers chains which are most behind. This is especially nice for indexers like UNIV4 where after indexing for a few hours, all chains besides base and optimism catch up to head, and we are left only with two chains left. We lose the HyperSync concurrency optimization. With this change, Base and optimism will have higher priority and progress with all the other changes simultaneously. Always guaranteeing that we have events and buffers to process (processing a full batch)
So there's no drop from 10k events per second to 4k events per second.
Summary by CodeRabbit
New Features
Behavior Changes
Performance
Tests