Repository navigation
ethmonitor: parallel block prefetch for fast chains - #219
Conversation
The monitor fetches each block and its logs serially, which caps ingestion at roughly 1 / (2 x node latency) blocks per second. On chains producing blocks faster than that (robinhood mainnet at ~10 blocks/s) the monitor falls behind and never recovers: node-gateway on dev was serving a head ~14 minutes stale. Add an opt-in prefetcher (Options.PrefetchConcurrency, default 0 = off, and Options.PrefetchWindow, default 4x concurrency). While the monitor trails the head, N workers fetch the blocks past its next block, and their logs, into the cache under the existing cache keys. The serial run loop is unchanged and still the only place the canonical chain is built; it just reads prefetched payloads as cache hits. At the head the prefetcher is idle, and in polling mode it only polls the head while the run loop is catching up, so slow chains see no extra node calls. Reorg safety: - with prefetch on, a block that does not extend the head is confirmed with a direct, uncached node fetch before it is treated as a reorg - a reorg resets the prefetcher and drops its cached blocks above the reverted block; a generation counter discards in-flight stale results Also, affecting all monitors: - block payloads are decoded before caching, and undecodable cache entries are dropped, so one bad node response can't poison a block number until cache expiry - Run owns a child context and joins its publisher, head listeners and prefetcher before returning, so goroutines don't outlive a fatal exit and a later Run - fetchNextBlock reads nextBlockNumber under its lock Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: cf8c4c191a
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| } | ||
|
|
||
| func (p *prefetcher) pollHead(ctx context.Context) { | ||
| tctx, cancel := context.WithTimeout(ctx, 4*time.Second) |
There was a problem hiding this comment.
Honor configured RPC timeout when polling prefetch head
When polling (rather than streaming), prefetching can only learn the scheduling window through this BlockNumber call. This hard-coded 4-second deadline ignores Options.Timeout, so on providers where eth_blockNumber takes more than four seconds but block/log requests succeed within the configured timeout, every head poll fails and latestHead remains zero; the prefetcher never schedules work and the monitor falls back to the slow serial path. Use the monitor's configured timeout here (as the prefetch block and log requests do).
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 987262d697
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| p.gen.Add(1) | ||
| p.cursor, p.highWater = 0, 0 |
There was a problem hiding this comment.
Purge prior-run prefetched entries before restarting
When a monitor is stopped after prefetching head+1, then that block is reorged while stopped, the next Run still reads the old by-number entry from the automatic cache. Its parent hash still matches the retained head, so the parent-mismatch confirmation does not run and the monitor publishes the abandoned block; if the reorged block remains the chain head, it is not corrected until a later block arrives. Resetting highWater here also discards the range needed to evict those entries, so preserve and purge the previous prefetch window (or bypass cached entries on restart) before scheduling the new run.
Useful? React with 👍 / 👎.
TestMonitorPrefetchReorg/prefetched=true raced the monitor: it waited for a block 2-4 past the monitor head to show up in the cache, assuming the monitor (~15ms a block) couldn't reach it first. On slow CI runners (macos-latest) the monitor either drained the fixed chain before the condition was seen, failing the test, or reached the reorg block first, passing without exercising the stale prefetched block at all. Hold block 1060 on the fake chain instead, so the monitor stalls before it while the workers prefetch past it. Reorg from 1062 once it is cached from the old fork, then release. Also assert the reorg was actually exercised: a block from the abandoned fork must be published and removed. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* fix(ethreceipts): harden rollback, filter ownership and cancellation * fix(ethreceipts): preserve retry candidates and filter identity Keep successful canceled retries tied to their learned block generation, reopen canonical adopted blocks after retention, and support stable custom value filter removal. Document supported limits and add regression coverage. * ci: test Go 1.25 and 1.26 * fix(ethreceipts): retry current candidates after stale RPC receipts * refactor(ethreceipts): pass block processing contexts explicitly * refactor(ethreceipts): restore breaker retries with v0.3.2 * test(ethreceipts): organize regressions by behavior * refactor(ethreceipts): clarify pending receipt retries * fix(ethreceipts): harden subscription cleanup and finality * fix(ethreceipts): preserve fetch filter compatibility * fix(ethmonitor): reject invalid empty log responses
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 95199aee60
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| state.generation++ | ||
| state.invalidatedAt = l.reorgRevision | ||
| state.removed = true | ||
| l.blockStates[block.Hash()] = state |
There was a problem hiding this comment.
On a long-running listener that sees recurring reorgs, this inserts one permanent blockState entry for every distinct removed block hash. blockStates is never pruned (including by PurgeHistory), while the monitor itself has bounded block retention, so reorg-prone chains will grow this map and its retained hashes without bound over the lifetime of the process.
Useful? React with 👍 / 👎.
| } | ||
|
|
||
| if len(blocks) == 1 { | ||
| if len(blocks) == 1 && blocks[0].Event != Added { |
There was a problem hiding this comment.
why the && blocks[0].Event != Added ?
There was a problem hiding this comment.
A single Added block needs to go through c.push() now, because that is where the monitor assigns a fresh canonical-state incarnation. Previously every one-block bootstrap took the copy-only shortcut; keeping that shortcut for Added would leave the bootstrap head untracked and its copies unable to observe a later removal.
The Event != Added condition preserves the existing special handling for a single non-Added input while routing canonical additions through the same adoption path as multi-block bootstrap. It does not change the supplied block hash or fetch anything from the provider.
TestBlockCanonicalStateBootstrap covers one and three blocks, both directly and through JSON. Those tests passed under -race in this check.
| // Each adoption owns its state so reusing an input cannot revive old events. | ||
| c.lastIncarnation++ | ||
| block := *nextBlock | ||
| block.canonicalState = &blockCanonicalState{incarnation: c.lastIncarnation} |
There was a problem hiding this comment.
why do we need block.canonicalState ..? will this add more memory overhead then what we already have?
There was a problem hiding this comment.
Yes, there is a small memory and allocation cost. I measured the layout on amd64: Block grows from 48 to 56 bytes, and blockCanonicalState occupies 16 bytes. That is 24 extra bytes per retained adopted block, before allocator/GC overhead, with one small state allocation per adoption. Queued copies share that allocation; each copy carries the additional 8-byte pointer. The adoption also makes a shallow Block copy; it does not duplicate the underlying Ethereum block, transactions, or log data.
The reason is receipt correctness when the receipt listener falls behind monitor retention. After H is removed and readopted, a delayed Added(H) event may arrive after GetBlock(H) no longer finds H. We must distinguish retention eviction from a real removal, and avoid reopening H from an obsolete event when H has been removed again. Shared state preserves that evidence beyond retention; a fresh incarnation on readoption keeps old copies invalid.
This supports the receipt hardening, and applies even with prefetch disabled. There is no additional monitor map retaining every historical block: the state becomes collectible once retained blocks and event copies release it. Retention, Readoption, and concurrent-removal tests passed under -race.
| // OK flag which represents the block is ready for broadcasting | ||
| OK bool | ||
|
|
||
| canonicalState *blockCanonicalState |
There was a problem hiding this comment.
what is canonicalState *blockCanonicalState and why was it added..?
There was a problem hiding this comment.
It is runtime metadata shared by the monitor's copies of one block adoption:
incarnationis a monitor-local ID assigned bypush()for each adoption.canonicalis an atomic flag, set true on adoption and false bypop()on removal.
An Added event describes what happened when it was emitted; that field alone cannot tell a delayed consumer whether the block has since been removed. Sharing this state lets old queued/snapshot copies observe the removal. Retention eviction leaves the flag unchanged, and readopting the same hash creates a new state rather than reviving the old event's state.
ethreceipts.acceptBlock() uses this when deciding whether an ordered Added event can reopen a previously removed block. JSON serialization excludes the state; deserialized/manual blocks remain untracked until adoption, with a retained-lookup fallback for legacy blocks.
One explicit compatibility tradeoff: adding the private field means external positional Block literals no longer compile; callers must use keyed literals such as ethmonitor.Block{Block: b, Event: ethmonitor.Added}. That is the remaining confirmed Codegenie finding and is documented on the type.
| if err == nil { | ||
| // check the logsBloom from the block to check if we should be expecting logs. logsBloom | ||
| // will be included for any indexed logs. | ||
| if len(logs) > 0 || block.Bloom() == (types.Bloom{}) { |
There was a problem hiding this comment.
why did we lose this check..? this was an important feature in the past, please confirm we stlil have this logic. check past commits if necessary
There was a problem hiding this comment.
Confirmed: the bloom protection is still enforced. This check moved into filterLogs() / unmarshalLogs(), so it runs before a response can be cached and also validates cache hits. addLogs() only marks a block ready after that validation succeeds; failures still leave it unready for backfill.
For an unfiltered query with a nonzero bloom, empty [], whitespace arrays, and JSON null are rejected. With a zero bloom, [] and null are valid empty results. The intentional topic-filter exception is that an unrelated log can set the block-wide bloom while the requested topics have no matches; that empty filtered response must be allowed.
I checked the older addLogs guard (b19b6ff) and the cache hardening in b261956 / PR #167. TestMonitorEmptyLogResponses now exercises readiness, retry calls, and caching across 30 cases with prefetch off/on. It passed under -race, along with the poisoned-cache recovery test.
| if err != nil { | ||
| return nil, err | ||
| } | ||
| if blockBloom != (types.Bloom{}) && (len(logsPayload) == 0 || (len(logsPayload) == 2 && logsPayload[0] == '[' && logsPayload[1] == ']')) { |
There was a problem hiding this comment.
this was important logic we added.. to account for different or faulty node behaviour.. please check old commits and make sure we account for this, and we have test cases too
There was a problem hiding this comment.
I checked b261956 / PR #167, which added this guard specifically to prevent a faulty empty origin response from being cached when the bloom says logs exist. That protection remains in the cache getter: we decode and validate the response with unmarshalLogs(..., expectLogs) before returning success to the cache. Invalid existing cache entries are also evicted so backfill can retry the origin.
The check now handles decoded emptiness rather than only the literal two-byte []: unfiltered nonzero-bloom [], [ ], and null all fail. Zero-bloom []/null remain supported. Topic-filtered empty results are allowed because the block bloom also includes unrelated logs. HTTP 503, JSON-RPC errors, an empty HTTP 200 body, and a missing result remain errors and are not cached as valid empty logs.
Coverage is in TestMonitorEmptyLogResponses (30 cases) and TestMonitorRecoversInvalidPrefetchedLogs (bad speculative responses and already-poisoned cache entries, followed by origin recovery and delivery of the expected nonempty logs). Both passed under -race in this check. These changes are already in a43e30f.
Rollback keeps a LimitOne owner's claim on its delivered txn so the same txn can be re-mined. Nothing released that claim once the orphaned delivery was pruned at finality, so a filter whose txn never came back silently dropped every later match: forever with MaxWait 0, or until ErrFilterExhausted otherwise. Release the claim when its delivery is pruned at finality. releaseClaim still keeps it while the owner has queued finality, in-flight, delivered or pending work, so a txn re-mined before then stays selected. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Upgrade goware/channel to v0.6.0 for safe concurrent Send/Close and backed-off queue alerts. Raise the subscriber warning threshold to 10. Release receiptMu before channel sends while deliveryMu continues to serialize publication, rollback, and finality. Recheck cancellation after block validation, accumulate matches across added blocks, and let custom filters manage their own MaxWait exhaustion. Preserve block generations needed by pending retries and late RPCs beyond monitor retention. Cover receipt recovery, orphan rejection, cancellation, blocked sends, batch matching, and custom exhaustion with regression tests. Bound the blocked-send test's startup and cleanup waits. Validation: full ethreceipts and ethmonitor suites with -race; focused regressions repeated 20 times; go build ./... and package vet passed.
Problem
ethmonitorfetches each block and then its logs serially, limiting ingestion to roughly1 / (2 × node latency)blocks/s. Faster chains accumulate lag after fetch failures, reorg pauses, and reconnects. On robinhood mainnet (~10 blocks/s), ingestion managed ~7–8 blocks/s; node-gateway on dev served aneth_blockNumberhead about 14 minutes stale while reporting healthy.Implementation
An opt-in parallel prefetcher fills the existing block/log cache while the monitor trails the head:
Options.PrefetchConcurrencysets the worker count; the default is 0 (off).Options.PrefetchWindowbounds how far beyond the monitor's next block workers fetch; it defaults to 4× concurrency.CacheBackendcreates an in-memory cache.newHeadsfor head information. Polling requestseth_blockNumberonly while catching up, after consecutive successful block fetches; head polling stops after catching up.WithLogs, that block's logs sequentially.Reorg and cache handling
Every monitor with a cache confirms a parent mismatch through a direct, uncached origin fetch before treating it as a reorg, including a shared-cache consumer with local
PrefetchConcurrency=0. Failed cache deletion or immediate stale repopulation cannot supply the confirmation result. Genuine origin-confirmed mismatches still produce canonical removals and replacement additions; uncached monitors retain their normal origin path. Cached monitors also confirm mismatches when their own getter just fetched the block from origin: one extra RPC on a rare mismatch keeps a single confirmation path.Cache errors fall back to origin for next-block fetches, block-by-hash ancestry, and logs while the caller context remains active. This keeps
WithLogspublication working during a shared-cache outage instead of advancing blocks with unavailable logs until the queue fills. The direct fallback deliberately bypasses cache reads and writes, preserves payload validation and RPC timeouts, and does not classify a backend failure itself as an origin miss. Cache reads still wait for the backend's configured timeout before fallback; the automatic memory cache retains its existing eight-second timeout.A reorg resets prefetch scheduling and deletes cached by-number entries above the popped block. A generation counter makes old in-flight workers discard/purge results from before the reset.
A stale prefetched fork that still extends the monitor's head can be accepted until a parent mismatch exposes it, potentially adding up to one prefetch window of old-fork blocks. Normal reorg recovery pauses at least 2 seconds per reverted block. Keep the window small on reorg-prone chains.
Hardening that also applies with prefetch disabled
Runowns a child context and joins its publisher, head listeners, and prefetcher on fatal exit or shutdown, preventing overlap with a later run. Streaming reconnect delay is cancellable. A provider that ignores cancellation can delayRunreturning.fetchNextBlockreadsnextBlockNumberunder its lock.Topic-filtered empty-log trade-off
Unfiltered log queries retain the nonzero-bloom guard: empty results are retried/backfilled when the block bloom indicates logs. With
LogTopics, valid empty results ([]ornull) are accepted even with a nonzero bloom, because the bloom does not establish a match in the first topic position and can produce false positives. This avoids endless retries for legitimately empty filtered queries. Compared withmaster, a lagging node that incorrectly returns an empty filtered result can now have that response cached forCacheExpiryand the block published without its matching logs; that block is not automatically backfilled. Malformed responses and RPC/HTTP errors remain failures.Canonical incarnation evidence
Block.CanonicalState() (incarnation uint64, canonical bool)lets delayed event consumers distinguish retention eviction from an actual removal of the event's accepted incarnation:Blocks.Copyshare the read-only witness. A pop atomically marks that incarnation noncanonical before publishing its Removed event; older Added copies observe the same removal.0, false); JSON decode clears runtime tracking. Added bootstrap inputs, including a single block, receive fresh destination-chain witnesses through normal adoption. The existing single-Removed bootstrap case stays untracked. JSON preserves the existingblock,event,logs, andokschema and carries no runtime witness.The state records this monitor's known acceptance/removal history, not an independent origin-canonicality guarantee. Reorgs beyond retained history keep the monitor's existing limits. The monitor witness lifetime follows existing event/snapshot references; its implementation adds no global registry or unbounded hash history. Receipt processing uses this evidence to authorize canonical same-hash readoption even after retention eviction.
Source migration:
Blockpreviously contained only four exported fields and now has a private state field. External positional literals such asethmonitor.Block{nil, ethmonitor.Added, nil, false}must use keyed literals, for exampleethmonitor.Block{Block: block, Event: ethmonitor.Added, Logs: logs, OK: true}. Structural conversions relying on the old exact four-field layout also change. Exported fields and JSON remain available. Bootstrap stores owned wrappers rather than preserving caller wrapper pointer identity. The exportedBlockGoDoc documents keyed construction.Receipt delivery and channel hardening
goware/channelis upgraded from v0.5.0 to v0.6.0. ItsSendandCloseoperations are safe to race, queue warnings back off as the backlog grows, capacity drops raise an alert, and closed queues drain without spinning or retaining delivered payloads unnecessarily.Receipt publication and finalization hold
receiptMuonly for block validation, releasing it before subscriber channel sends. SubscriberdeliveryMustill serializes delivery, rollback, finalization, and owner retirement. A slow subscriber send therefore no longer holds the listener-wide receipt lock; block processing still waits for its subscriber workers. Publication rechecks cancellation after waiting for block validation.Batch match results accumulate across all Added blocks; rollback notifications do not count as new matches. Custom filters manage their own MaxWait counters and exhaustion signals instead of causing a listener panic. The subscriber queue warning threshold rises from 2 to 10 unread receipts; capacity remains 5000.
Removal and generation state remains retained because pending retries and in-flight receipt RPCs can outlive both monitor retention and block finality. Advancing the chain must neither invalidate a current readopted candidate nor allow a late orphan response into the receipt cache.
Validation
The fake-chain/gomock tests enforce expected provider calls and cover:
Throughput at 100 blocks/s with 15ms calls, polling and streaming with 0/4/8 workers: serial ingestion falls behind while prefetch keeps up.
Prefetched-fork and below-head reorgs with 2/4 workers, checking both retained chain and subscriber event reconstruction.
Slow-chain operation at the head with matching block/log call counts and zero extra
eth_blockNumbercalls, plus idle head polling after backlog recovery.Cache confirmation with local concurrency 0/1 crossed with deletion failure and stale repopulation; genuine reorgs with cached 0/1 and uncached 0 controls.
Invalid block/log recovery, announced-but-not-yet-available heads, fatal exits, panics, restart in both modes, and worker shutdown without goroutine leaks.
Shared-cache outage with local prefetch 0/2: publish 20 consecutive blocks with nonzero blooms and intact logs while the backend is down (queue capacity 10), then publish 20 more and verify caching resumes after recovery.
Cache-error fallback for logs and block-by-hash fetches, including caller cancellation and rejection of malformed origin responses. Next-block timeout/shared-error regressions preserve slow-head waiting and count only actual origin misses.
Empty-log provider compatibility for zero/nonzero blooms, filtered queries,
[],null, malformed payloads, and HTTP/RPC failures.Canonical state across multiple retention depths, value/snapshot copies, pop, same-hash readoption, input reuse, bootstrap/JSON, nil/untracked inputs, and concurrent state observation. Independent probes additionally checked buffered events after alternate-chain eviction and source/destination bootstrap ownership.
Receipt regressions cover recovery of pending same-hash readoptions after chain advancement, rejection and cache exclusion of late orphan RPC responses, and cancellation while block validation waits for
receiptMu.Receipt tests also cover matches in earlier blocks of a batch, rollback-only batches, blocked subscriber sends, and built-in/custom MaxWait behavior. The blocked-send test bounds alert startup and publisher cleanup waits.
Final checks:
ErrQueueFullat head 1010 on62a096awith both local prefetch settings. The new logs/by-hash fallback tests also failed against that production source.-race, 20 repetitions: passed.-race: passed.ethmonitorsuite under-race, including the live fee-history test: passed in 76.329s on Go 1.27.1.ethreceiptssuite under-raceagainst the local Hardhat testchain: passed in 246.050s on Go 1.27.1.-race, 20 repetitions: passed.go vet ./ethreceipts ./ethmonitor,go build ./..., dependency verification, formatting, and diff checks passed.Follow-up
node-gateway can bump ethkit, add a per-network
PrefetchConcurrencysetting, and roll out on robinhood.Bounding receipt block-generation state remains a follow-up. Age-based pruning is deferred because deleting a still-referenced generation can lose canonical pending delivery or accept an orphan RPC result. The existing map-growth concern remains until cleanup accounts for outstanding receipt work.