Skip to content

fix: bound concurrent pages and cancellation for large cold scans - #228

Open
6tamichael-boop wants to merge 6 commits into
wraith-protocol:developfrom
6tamichael-boop:cloud-fixer/sdk-208-b1-1
Open

6tamichael-boop wants to merge 6 commits into
wraith-protocol:developfrom
6tamichael-boop:cloud-fixer/sdk-208-b1-1

Conversation

@6tamichael-boop

Copy link
Copy Markdown

Overview

Cold scans can split the ledger range into parallel chunks and interleave them with mergeOrdered, but nothing bounded how many chunks a caller could request, and cancelling a merge did not close the chunk generators that had not yielded yet. This change caps the parallel chunk count, documents the one-item-per-chunk merge buffer, and makes cancellation tear down every chunk iterator. It also adds the regression tests and benchmark the issue asks for.

Related Issue

Closes #208

Changes

Bound concurrent pages

  • [MODIFY] src/chains/stellar/announcements.ts
    • Add MAX_COLD_SCAN_PARALLELISM = 8 and clamp the parallelism option to [1, MAX_COLD_SCAN_PARALLELISM] (non-finite/fractional hints fall back to 1). Each chunk keeps one getEvents page in flight and one buffered item in the merge, so an unbounded hint no longer means unbounded pending RPC work.
    • The option's JSDoc now states the clamp; the constant is module-level (not re-exported from the package entry, so the public API surface is unchanged).

Bound the merge buffer, fix cancellation

  • [MODIFY] src/chains/stellar/announcements.ts — mergeOrdered
    • Wrapped the merge loop in try/finally and return every chunk iterator on exit, including ones that never yielded. Previously a consumer break/.return() only unwound the delegated chain, leaving other chunk generators suspended with a page outstanding.
    • Documented the existing backpressure contract explicitly: exactly one item buffered per chunk, and a chunk is pulled again only after its previous item was consumed, so memory is O(chunks) rather than O(scan size).

Regression tests

  • [MODIFY] test/chains/stellar/announcements.test.ts
    • caps cold-scan parallelism so in-flight chunks stay bounded: an oversized parallelism: 10_000 produces exactly MAX_COLD_SCAN_PARALLELISM distinct chunk requests.
    • mergeOrdered closes every chunk iterator when the consumer cancels: five chunk generators, break after one item, all five finally blocks run.
    • mergeOrdered bounds how far each chunk runs ahead of a slow consumer: interleaved keys and a deliberately slow consumer; per-chunk lead stays within one buffered item plus the in-flight one.

Benchmark

  • [MODIFY] test/chains/stellar/bench/scan.bench.ts
    • New Stellar cold-scan backpressure section: a correctness test draining MAX_COLD_SCAN_PARALLELISM chunks with a slow consumer, and a matching bench case, so the bound is exercised under the existing bench harness.

Docs

  • [MODIFY] docs/chains/stellar-streaming-scan-pipeline.md
    • New "Bounding parallel cold scans" section covering the chunk cap, the O(chunks) merge buffer, and the cancellation teardown added here.

Verification Results

Changes authored via GitHub Contents/Git API (no local clone or sandbox).
`pnpm test`, `pnpm exec vitest bench` and typecheck were not executed in this
environment; the new tests/bench are provided but unrun here. The public API
surface is unchanged (`MAX_COLD_SCAN_PARALLELISM` is not re-exported from the
package entry), so the api-extractor reports do not need regeneration.
Acceptance criteria mapping — see the table below.
Acceptance Criteria Status
Bound concurrent pages and buffered announcements ✅ MAX_COLD_SCAN_PARALLELISM clamp + one-item-per-chunk merge buffer
Verify memory stays bounded while consuming slowly ✅ Slow-consumer lead-bound regression test
Confirm cancellation stops all outstanding requests ✅ mergeOrdered finally returns every chunk iterator + regression test
Add a benchmark and a regression test ✅ Stellar cold-scan backpressure bench section + 3 tests

Closes #208

@drips-wave

drips-wave Bot commented Sep 25, 2026

Copy link
Copy Markdown

@6tamichael-boop 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! 🚀

Learn more about application limits

@truthixify

Copy link
Copy Markdown
Contributor

The iterator cleanup test only proves return() is called on paused generators. It does not cancel an in-flight fetch, and initial next() calls are still awaited sequentially. Please propagate AbortSignal to RPC requests and test cancellation with blocked fetches.

@truthixify

Copy link
Copy Markdown
Contributor

CI is approved and running now. Please format src/chains/stellar/announcements.ts and test/chains/stellar/announcements.test.ts. The cancellation and sequential initial fetch issues from my earlier review also remain.

…initial fetch

- Thread an AbortSignal through fetchAnnouncementsStream, mergeOrdered and
  every Soroban/Horizon fetch so pending RPC and pagination work is cancelled.
- Start the first page in parallel and cap cold-scan parallelism so in-flight
  chunk requests stay bounded.
- Move the abort/ordering coverage into test/chains/stellar and format.
@6tamichael-boop

Copy link
Copy Markdown
Author

@truthixify Both issues from your review are addressed, plus formatting.

  • AbortSignal is now propagated end to end: fetchAnnouncementsStream takes options.signal and every fetch in the Soroban RPC and Horizon paths receives it, so an abort cancels requests that are already in flight instead of only tearing down paused generators. Pagination checks the signal before each page and before each yielded event and throws a standard AbortError (DOMException).
  • The initial next() calls in mergeOrdered are no longer awaited one at a time: they are issued with Promise.all so every chunk starts its first page concurrently, and the merge still yields in ledger order. On abort/exit the iterators are returned in try/finally.
  • Cold-scan parallelism is bounded (MAX_COLD_SCAN_PARALLELISM, clamped options.parallelism) so widening the fan-out cannot leave an unbounded number of chunk requests in flight.
  • Tests now live in the existing tree at test/chains/stellar/announcements.test.ts, including cancellation with a blocked fetch (asserting the AbortError and that the blocked request was rejected) and a test that all iterators are primed concurrently.
  • src/chains/stellar/announcements.ts and its test are Prettier-clean; etc/sdk-stellar.api.md updated for the new optional field.

Verified locally on a checkout of the branch merged with develop: vitest run test/chains/stellar/announcements.test.ts passes 28/28, pnpm build succeeds, pnpm api:check is clean and prettier --check passes.

Head is now 3b2656a.

@6tamichael-boop

Copy link
Copy Markdown
Author

@truthixify Verified head 3b2656a4 for the three points from your reviews, no code change:

  • In-flight cancellation: every Soroban RPC and Horizon fetch in src/chains/stellar/announcements.ts receives options.signal, and an aborted signal throws a standard AbortError (DOMException) — no longer only tearing down paused generators.
  • No sequential priming: mergeOrdered now primes every iterator with await Promise.all(iterators.map((it) => it.next())) and returns them via Promise.allSettled(...) in finally, still yielding in ledger order. MAX_COLD_SCAN_PARALLELISM = 8 bounds the fan-out.
  • Tests: test/chains/stellar/announcements.test.ts has "threads the signal into fetch and cancels an in-flight page request" (blocked fetch) and "mergeOrdered pulls the first item from every iterator concurrently".
  • Formatting: npx prettier@3.4.2 --check src/chains/stellar/announcements.ts test/chains/stellar/announcements.test.ts → clean.

Could you re-review?

@truthixify

Copy link
Copy Markdown
Contributor

#224 has merged the shared AbortSignal work, so this branch now conflicts with develop. Please rebase and keep the #208-specific concurrency cap, bounded-buffer tests, benchmark, and docs. CI was green before the conflict.

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.

[Wave 9] Add backpressure tests for large cold scans

2 participants