Skip to content

fix: concurrentPool buffered-delivery stall and unbounded prefetch - #399

Merged
ppeeou merged 3 commits into
mainfrom
fix/concurrent-pool-backpressure
Aug 8, 2026
Merged

ppeeou merged 3 commits into
mainfrom
fix/concurrent-pool-backpressure

Conversation

@ppeeou

@ppeeou ppeeou commented Aug 8, 2026

Copy link
Copy Markdown
Member

Summary

Fixes three empirically reproduced defects in concurrentPool (regression tests committed first, all failing on the previous implementation):

  • Buffered-delivery stall: next() only relied on a future fetch settling to trigger delivery, so once the pull chain drained, already-buffered results were never handed out while the following upstream fetch hung. next() now serves the buffer directly.
  • Unbounded prefetch with a slow consumer: every next() launched length new fetches regardless of buffered results (upstream pulls == consumed × pool size). Prefetch is now capped at the pool size when no consumer is waiting; under active demand the pool stays saturated, which is what order-preserving pools (p-map included) must do past a slow head item.
  • Post-error pulls (the long-standing TODO): the source is no longer pulled after it has thrown, and once nothing is in flight and the buffer is drained, later consumers receive done.

Throughput parity re-verified after the fix: p-map 455ms vs concurrentPool 455ms on the uneven-duration benchmark (concurrency 2) — the migration guide's claims still hold.

Test plan

  • 3 new regression tests (hang via Promise.race timeout, produced ≤ consumed + pool, zero post-throw pulls), all red before / green after
  • Full jest suite passes (140 suites, 1258 tests); lint 0; compile:check 0
  • p-map parity probe on the built output: 455ms == 455ms

🤖 Generated with Claude Code

ppeeou and others added 2 commits August 8, 2026 17:48
…ulls

Three failing regressions:
- buffered items are not delivered once the pull chain drains - next()
  hangs on a buffered result while the following upstream fetch hangs
- every consumer next() launches `length` new fetches regardless of
  buffered results, so a slow consumer grows the buffer unboundedly
  (produced == consumed x pool size)
- the source is pulled again after it has thrown (the known TODO)

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…currentPool

- next() now serves already-buffered results immediately instead of
  depending on a future fetch settling, fixing a hang where buffered
  items were never delivered once the pull chain drained while the
  next upstream fetch was pending
- prefetch is capped at the pool size when no consumer is waiting, so
  a slow consumer no longer grows the buffer unboundedly (was
  consumed x pool-size upstream pulls); under active demand the pool
  stays saturated, preserving p-map-equivalent throughput (verified
  455ms == 455ms on the uneven-duration benchmark)
- the source is no longer pulled after it has thrown (resolves the
  long-standing TODO), and post-error consumers receive done once
  nothing is in flight and the buffer is drained

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@ppeeou ppeeou self-assigned this Aug 8, 2026
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@ppeeou
ppeeou merged commit d6518e2 into main Aug 8, 2026
4 checks passed
@ppeeou
ppeeou deleted the fix/concurrent-pool-backpressure branch August 8, 2026 13:57
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