Skip to content

Harden ethreceipts rollback, filter ownership and cancellation - #220

Merged
pkieltyka merged 12 commits into
parallel-monitorfrom
ethreceipts-hardening
Oct 6, 2026
Merged

pkieltyka merged 12 commits into
parallel-monitorfrom
ethreceipts-hardening

Conversation

@pkieltyka

@pkieltyka pkieltyka commented Oct 5, 2026 •

Copy link
Copy Markdown
Member

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.

  • Ownership is isolated per registration, including duplicate public filter IDs and LimitOne matching across snapshots and concurrent processing. Remove, Clear, completion and exhausted rollback release the exact owner's work. Ordinary exhaustion preserves queued finality; re-registration rejects old workers.
  • Pending retries preserve valid candidates until publication commits. A canceled attempt remains retryable without resurrecting removed owners or rolled-back blocks. Partial QueryOnChain results retain their completed block/generation. A lagging RPC response from orphan block A is rejected while a still-current block B candidate remains eligible for bounded retries; obsolete generations remain terminal.
  • Canonical re-adoption uses the monitor incarnation evidence from PR #219. Delayed live Added events can reopen the correct same-hash block after retention eviction. Cached snapshots, removed incarnations and earlier RPC generations cannot revive invalidated work. This PR is based directly on parallel-monitor.
  • Fetch helpers own their query state and reuse complete QueryOnChain receipts. Stale requests preserve fresh canonical cache entries. Signed transaction From/To matching and finality initialization for subscriptions created before Run are corrected. Run/Stop and RPC, semaphore, head and finality waits honor cancellation. ChainID retries use goware/breaker v0.3.2 with cancellable backoff.
  • Unsubscribe is safe to repeat or call concurrently. Fetch helpers handle exhaustion once and continue waiting for queued finality after mining. Finality notifications use the same inclusive block threshold as the listener. The helper uses its existing completion channel for cleanup and a local match flag. Dead commented-out code and misleading option/deprecation comments are cleaned up.

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=10m passed against the local chain in 238.415s, including all regression and integration tests.
  • The unsubscribe, finality-boundary and post-mining-exhaustion regressions fail before their fixes and pass afterward. Focused race checks passed five repetitions and on Go 1.25.6 and 1.26.6.
  • The real-monitor shallow-reorg/exhaustion reproduction passes with the fixes.
  • 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.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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".

Comment thread ethreceipts/subscription.go Outdated
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.
@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Oct 5, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-10-06T16:44:26.213273Z 584f585 New commits
ℹ️ About Codex in GitHub

Your team has set up Codex to 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" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

Comment thread ethreceipts/ethreceipts.go Outdated
Comment thread ethreceipts/ethreceipts.go Outdated
Comment thread ethreceipts/ethreceipts.go Outdated
Comment thread ethreceipts/ethreceipts.go
//
fmt.Println("listening for txns..")

sub := receiptsListener.Subscribe(

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

for the deleted lines in this hunk, I want to ensure we are still properly testing these filter options..?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. I compared the old filter test with the current assertions. TestReceiptsListenerFilters still covers:

  • FilterFrom with LimitOne: 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 := &ethtxn.TransactionRequest{

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please confirm we're testing all these cases properly.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread ethreceipts/finalizer.go
Comment thread ethreceipts/receipt.go
blockNum *big.Int
blockHash common.Hash
generation uint64
owner *filterOwner

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why do we care about the filter owner? What does it do and who does it represent?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 }

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 one LimitOne registration 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.

Comment thread FIXES.md Outdated
@@ -0,0 +1,26 @@
# Receipt listener fixes

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please delete the fixes.md file. We don't need it.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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".

Comment thread ethreceipts/ethreceipts.go

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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".

Comment thread ethreceipts/ethreceipts.go

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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".

Comment on lines +523 to +526
if ok && !l.validFetchedBlock(receipt.BlockHash, expected, started) {
l.receiptMu.Unlock()
errCh <- ethereum.NotFound
return

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

@pkieltyka
pkieltyka merged commit 95199ae into parallel-monitor Oct 6, 2026
8 checks passed
@pkieltyka
pkieltyka deleted the ethreceipts-hardening branch October 6, 2026 16:56
pkieltyka added a commit that referenced this pull request Oct 7, 2026
* 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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant