feat(listener): add event processing checkpoints (#783) - #820
Open
Jessepriase wants to merge 2 commits into
Open
Jessepriase wants to merge 2 commits into
Jessepriase wants to merge 2 commits into
Conversation
Persist the latest successfully processed ledger position so the
listener can resume from a known checkpoint after interruption,
rather than replaying from the beginning on every restart.
Changes
-------
event-subscriber.ts
- Add restoreCheckpoints() private method: loops over all configured
contract addresses, calls deduplicationService.getLastCursor() for
each, and populates this.lastCursors with the stored cursor string.
Per-contract failures are caught and logged as warnings so a single
bad DB row never prevents startup. Logs per-contract and a summary.
- Call await this.restoreCheckpoints() from start() before queue
start and the poll loop — the only moment restoration is needed.
- Add private backfillStartLedger field (was referenced by
resolveBackfillStartLedger but never declared).
- Remove duplicate processableEvents re-declaration and duplicate
request re-declaration in getContractEvents (pre-existing merge
artifacts that caused SyntaxError: Identifier already declared).
index.ts
- Add deduplicationService = new EventDeduplicationService(db)
immediately after initializeDatabase(). The let declaration and
import already existed; the missing assignment meant every
if (this.deduplicationService) guard inside EventSubscriber was
always false at runtime — cursor persistence and checkpoint
restoration were dead code.
request-id.ts
- Fix missing closing } on generateCorrelationId(): the JSDoc comment
for isValidRequestId had leaked inside the function body, causing a
parse error that prevented the module from loading.
event-subscriber-checkpoint.test.ts (new)
- 16 tests across three describe blocks using in-memory SQLite.
Acceptance criteria met
-----------------------
Restarting the listener does not require processing the entire history
restoreCheckpoints() loads the persisted cursor from polling_cursors
into this.lastCursors before the first poll; getContractEvents() uses
the cursor branch (not startLedger) on the first call after restart.
Checkpoints are only advanced after successful processing
updatePollingCursor() is only called inside the response.cursor block
at the end of checkForEvents(), after all events in the batch have
been processed. A poll that throws never reaches that block.
Recovery behaviour is covered by tests
checkpoint persistence: cursor written after success, advances per
poll, absent when no cursor returned, stable when RPC throws,
independent per contract.
checkpoint restoration: cursor loaded into lastCursors, used in
first RPC call, cold start uses startLedger when no record, multi-
contract restore, partial restore, no-op with null dedup service,
logging verified.
end-to-end: second subscriber instance resumes from the first's
cursor; checkpoint only advances on success across failure cycles.
|
@Jessepriase Great news! 🎉 Based on an automated assessment of this PR, the linked Wave issue(s) no longer count against your application limits. You can now already apply to more issues while waiting for a review of this PR. Keep up the great work! 🚀 |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Persist the latest successfully processed ledger position so the listener can resume from a known checkpoint after interruption, rather than replaying from the beginning on every restart.
Changes
event-subscriber.ts
index.ts
request-id.ts
event-subscriber-checkpoint.test.ts (new)
Acceptance criteria met
Restarting the listener does not require processing the entire history
restoreCheckpoints() loads the persisted cursor from polling_cursors
into this.lastCursors before the first poll; getContractEvents() uses
the cursor branch (not startLedger) on the first call after restart.
Checkpoints are only advanced after successful processing
updatePollingCursor() is only called inside the response.cursor block
at the end of checkForEvents(), after all events in the batch have
been processed. A poll that throws never reaches that block.
Recovery behaviour is covered by tests
checkpoint persistence: cursor written after success, advances per
poll, absent when no cursor returned, stable when RPC throws,
independent per contract.
checkpoint restoration: cursor loaded into lastCursors, used in
first RPC call, cold start uses startLedger when no record, multi-
contract restore, partial restore, no-op with null dedup service,
logging verified.
end-to-end: second subscriber instance resumes from the first's
cursor; checkpoint only advances on success across failure cycles.
Overview
Related Issue
Closes #
Changes
Verification
How to Test
Checklist
maincargo fmt --allrun (if Rust changes)npm run lintpasses (if TypeScript changes)closes Add Event Processing Checkpoints #783