Repository navigation
Harden ethreceipts rollback, filter ownership and cancellation - #220
Conversation
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 375c403208
ℹ️ 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".
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.
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. |
| // | ||
| fmt.Println("listening for txns..") | ||
|
|
||
| sub := receiptsListener.Subscribe( |
There was a problem hiding this comment.
for the deleted lines in this hunk, I want to ensure we are still properly testing these filter options..?
There was a problem hiding this comment.
Yes. I compared the old filter test with the current assertions. TestReceiptsListenerFilters still covers:
FilterFromwithLimitOne: exactly the first matching transaction.FilterTo: all three matching senders.- Transaction-hash filters and their distinct public IDs.
- Adding a filter to an existing subscription and later finding an already-mined transaction through cache search.
- Removing the matching owner after its mined event, with no final event for that removed owner.
- An unmatched filter exhausting after
MaxWait(4). - Mined and final events on the other two subscriptions.
The old secondary-subscription loops mostly printed receipts. The current expected-delivery map asserts exact subscription/hash/filter/finality combinations, original filter identity, success status, absence of reorg flags and duplicates, and the final active-filter set. MaxWait(0) is explicit on the long-lived From/To filters so their separate expiry does not make those assertions timing-dependent; the dedicated MaxWait(4) expiry case remains. TestFetchTransactionReceiptBlast also still waits for mined and final success for every requested transaction, now propagating worker errors to the test goroutine.
|
|
||
| // Add a bunch of subscribers that will each send a txn and wait for it | ||
| for i := 0; i < 50; i++ { | ||
| txr := ðtxn.TransactionRequest{ |
There was a problem hiding this comment.
Please confirm we're testing all these cases properly.
There was a problem hiding this comment.
Yes: this subtest retains a healthy monitor and puts the receipt listener's RPC provider into 100% failure for 20s before restoring it. Transaction submission still uses the separate healthy provider. The common helper sends 50 distinct transactions and requires the matching successful receipt for each subscription, with worker errors and timeouts propagated to the test goroutine.
The following fourth subtest also remains: monitor and receipt listener both use the failing provider, while submission uses the healthy provider. It uses the same 20s restoration and 50-receipt assertions. All four combinations (healthy/healthy, failing monitor, failing receipt provider, both failing) are still present and executed.
| blockNum *big.Int | ||
| blockHash common.Hash | ||
| generation uint64 | ||
| owner *filterOwner |
There was a problem hiding this comment.
Why do we care about the filter owner? What does it do and who does it represent?
There was a problem hiding this comment.
This is the internal token for the exact filter registration that produced the receipt. Every Subscribe/AddFilter registration gets its own token, even if the caller reuses the same public filter object. Receipt.Filter still exposes that original public filter.
Keeping the token on the receipt lets later finality and rollback processing find the exact registration that owns the work. Public IDs can repeat, custom filter values may be non-comparable, and the same query can be removed and registered again, so neither a public ID nor the filter value alone can safely identify that lifetime. The field is private and is not a new public handle. Owner-isolation and remove/re-add tests exercise this distinction.
|
|
||
| // Every registration owns a comparable token; the public filter may contain | ||
| // slices, maps or functions and may be reused in a later registration. | ||
| type filterOwner struct{ Filterer } |
There was a problem hiding this comment.
It seems the subscription dot go files change substantially with the concept of a filter owner What is that for we do things like retire and release and all these different methods? What's the general direction of the new subscription logic and why is it necessary for us to do this?
There was a problem hiding this comment.
The direction is to give each filter registration a precise lifetime and serialize publication against cancellation/rollback. filterOwner is the internal identity of one registration within a subscription; the caller still receives the original filter value.
The key operations are:
reserve: claim a matching transaction before the RPC wait, so concurrent cache search/live processing cannot select two transactions for oneLimitOneregistration or start duplicate work.releaseReservation/releaseClaim: remove the temporary claim after an attempt, retaining it only while pending work, a delivered receipt, or queued finality still needs it.retireOwner: remove a completed/canceled registration and its work. Explicit removal clears its finality; normal exhaustion keeps already-mined receipts eligible for finality.publish: commit delivery only while that exact registration and block generation remain valid. RPC waits happen outside the delivery lock.
Previously, pending/finality identity could collide across filters, LimitOne was decided after concurrent work had started, and removing/re-adding a query could let an old worker deliver into its new lifetime. The owner token and shared receipt key close those gaps together. Regressions cover overlapping owners with duplicate public IDs, concurrent LimitOne, remove/clear during in-flight work, remove/re-add, rollback/re-mining, and expiry preserving queued finals. The amount of code reflects the cleanup needed across these stages; the public subscription API remains the same.
| @@ -0,0 +1,26 @@ | |||
| # Receipt listener fixes | |||
There was a problem hiding this comment.
Please delete the fixes.md file. We don't need it.
There was a problem hiding this comment.
Removed FIXES.md in b33c95faeee305c80fedf563f60b8e71f187b6a1 and removed its reference from the PR description. The implementation reasoning and validation remain in the PR description and these review replies.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: b33c95faee
ℹ️ 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".
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: d711f5c0d0
ℹ️ 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".
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 584f585fae
ℹ️ 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".
| if ok && !l.validFetchedBlock(receipt.BlockHash, expected, started) { | ||
| l.receiptMu.Unlock() | ||
| errCh <- ethereum.NotFound | ||
| return |
There was a problem hiding this comment.
Re-fetch mismatched cached receipts for a live block
When the past-receipts cache contains a transaction from old block A whose removal fell outside this listener's observed retention, and the monitor later sees the same transaction re-mined in block B, fetchReceipt supplies B as expectedBlock. This branch rejects cached A as NotFound; matchFiltersAndPublish treats that error as terminal and does not enqueue a retry, so the receipt for live block B is silently never delivered until cache eviction or manual purge. For a cache mismatch while the expected block is still current, bypass/delete the cached entry and fetch the provider instead.
Useful? React with 👍 / 👎.
* ethmonitor: parallel block prefetch for fast chains 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> * test(ethmonitor): avoid timing-dependent polling counts * fix(ethmonitor): reject malformed logs before caching * ethmonitor: make the prefetched reorg test deterministic 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(ethmonitor): confirm cached parents and track canonical incarnations * Harden ethreceipts rollback, filter ownership and cancellation (#220) * 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 * fix(ethmonitor): preserve empty-log provider compatibility * fix(ethmonitor): retry origin after cache fetch failures * fix(ethmonitor): keep log publication alive during cache outages * fix(ethreceipts): release orphaned LimitOne selection at finality 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> * fix(ethreceipts): harden delivery and channel handling 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. --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Receipt listeners could lose canonical deliveries across reorgs, overlapping filters and canceled retries, or publish/finalize work whose block or registration had already been invalidated. This change tracks receipts by transaction, canonical block generation and filter registration so rollback, re-mining and finality affect the correct work.
parallel-monitor.Public filter values and custom Match behavior remain intact. RemoveFilter uses ordinary Go equality for comparable values and pointers. Callback-bearing values otherwise use a stable nonnil Exhausted channel and the same concrete type as base identity; shared-base values are aliases. Removal selects the first matching active registration, or its first queued finality owner if inactive. Distinct bases or pointers provide independent removal; nil or unstable signals require pointer registration.
Tests are organized by behavior in ownership, pending-receipt, receipt-fetch, finalizer, lifecycle, reorg and stale-receipt files, with shared fixtures in
helpers_test.go. Existing local-chain integration tests propagate service errors, cancel and join workers, and verify mined/final deliveries. The Makefile serializes package binaries that share that chain. CI targets Go 1.25.x and 1.26.x on Ubuntu and macOS.Validation:
go test -race ./ethreceipts -count=1 -timeout=10mpassed against the local chain in 238.415s, including all regression and integration tests.go vet ./ethreceipts,go build ./..., formatting and diff checks pass.Fetch helpers preserve legacy FilterQuery builder conversion, including builders resolved by the optional Finalize step, before snapshotting the resolved filter. Mined and final receipts expose the original public Filterer type/value while helper options, counters and exhaustion remain independent. Five compatibility cases cover built-in, custom pointer/value and builder filters; all failed before this fix and pass afterward, including repeated race checks and both CI Go versions.
Test startup cleanup is registered before waiting, cancels Run directly and uses a reusable completion signal. Prefetch panic-test teardown bounds its cancellation wait. Forced-failure probes confirm startup work is canceled/joined and an unresponsive prefetch test returns after its bounded teardown rather than reaching the global test timeout.