Close the turn queue's enqueue/release race and dispatch strand - #474
TheGreatAxios merged 5 commits into
Conversation
|
greybeard · reviewed design pre-implementation Reviewed the drain-loop/reclaim design before implementation, not the diff cold. The initial plan only fixed the releaser-side race (recheck pending after
|
|
critique - verified empirically (checked out new tests against pre-fix code: all 4 fail deterministically; post-fix: 9/9 pass, tsc clean, package suite green aside from one unrelated pre-existing timeout in blocks.test.ts). Traced the reclaim dance against TurnClaimStore's contract by hand: no drop scenario found even with 3+ concurrent racers. Two non-blocking observations judged non-issues: an apparently-unreachable defensive branch in the afterReclaim checks, and an unchanged pre-existing publish-before-reclaim ordering. All four CL-7195 acceptance criteria met and test-covered. |
|
Reviewed the drain-loop/reclaim design and the diff directly (not just the description). The symmetric two-sided reclaim holds up: traced the interleavings on both the releaser side ( One gap found and fixed: Verified: |
a84807b to
82398b3
Compare
Reproduces CL-7195: the drain loop's emptiness check and its claim release are not atomic against a concurrent run() whose tryClaim fails and pushes to the pending map in that window. Covers both the releaser-side interleaving (release() awaits past a concurrent push, written during the audit) and the symmetric pusher-side interleaving (a failed tryClaim resolves only after the holder already released and returned). Also covers a rejecting dispatch stranding whatever queued behind it, and updates the existing failed-dispatch test for the "dispatch never rejects run()" contract the fix will enforce. These fail against the current implementation.
The drain loop broke on an empty pending list and released the claim in a finally, but the check and the release were not atomic against a concurrent run() that failed tryClaim and pushed in that window — the in-memory claim store's release only looked safe because it resolved before its first await; a durable store reopens the window and the pushed turn is never drained by anyone. Both sides of the race now reclaim: the drain loop rechecks the pending list after release() settles and reclaims if something landed during it, and run()'s enqueue path does the same after a failed tryClaim. Whichever side notices the claim is free first picks up what's pending, so nothing queued is ever left with no holder to drain it. dispatch is now guarded by a try/catch: a rejection (which breaks its own documented "never reject" contract) is reported via reportError and logged instead of propagating out of run() and stranding whatever queued behind it. pendingByWorkbench is documented as process-local and single-replica- only, paired one-to-one with the claim store it backs — a multi- replica hub needs every workbench's turns routed to the same replica until this pending list moves into whatever store backs the claims.
The drain loop's holds/release/tryClaim calls were only proven safe against resolving false; a durable store doing I/O can also reject outright, and run() documents that it never rejects.
drain()'s loop wrapped only the dispatch call in try/catch; a holds, release, or tryClaim call that rejected (rather than resolving false) propagated straight out of run(), skipping the claim release the store's own contract requires on every path that can still complete and breaking run()'s own "never rejects" guarantee. Both drain() and run()'s symmetric enqueue-side reclaim now catch these rejections, report them, and still attempt to free the claim.
82398b3 to
9a96831
Compare
Summary
Fixes CL-7195:
packages/chat/src/turn-queue.ts's drain loop checked its pending map for emptiness and then released the workbench's turn claim in afinally, non-atomically. A concurrentrun()call that failedtryClaim(the claim was still held) and pushed a turn into the pending map during the window beforerelease()actually took effect would strand that turn forever — nobody would be left holding the claim to drain it. The message lands in the timeline and is never answered, with no error.This was invisible only because the in-memory
TurnClaimStore'sreleaseresolves synchronously before its first await.turn-claims.ts's own interface explicitly anticipates a durable implementation, and any store whosereleaseawaits (e.g. Postgres-backed) reopens the window.Changes
release()settles, it rechecks the pending map; if a turn landed during that await, it reclaims viatryClaimand keeps draining. Symmetrically,run()'s enqueue path (whentryClaiminitially fails and the turn is pushed) also retriestryClaimimmediately after pushing, in case the holder had already released before the push landed. Whichever side notices the claim is free first drains what's pending — nothing queued is ever left with no holder.dispatchcan no longer strand the queue.dispatch(batch)is wrapped in try/catch: a rejection (which breaks its own documented "never reject" contract) is reported viareportErrorand logged, instead of propagating out ofrun()and stopping the drain loop mid-batch.pendingByWorkbenchis a process-localMappaired one-to-one with the claim store; a multi-replica hub would need every workbench's turns routed to the same replica, or a turn queued on one replica is never read by another. Documented as single-replica-only (matching the existing precedent inagent-turns.ts), rather than building actual cross-replica correctness — out of scope for this ticket.Tests
TurnClaimStorewith an injected hook to force the exact interleaving: one on the release side (from the CL-7195 audit), one on the symmetric enqueue side. Both fail deterministically against the pre-fix code and pass after.dispatchno longer strands what queued behind it.run()now resolves instead of rejecting whendispatchthrows, matching its own documented "never rejects" contract.Review
Design reviewed with Greybeard before implementation (closed a residual pusher-side race in the initial one-sided plan). Reviewed by Critique after landing — verified the fix empirically (checked out the new tests against pre-fix code, confirmed all four fail deterministically; confirmed the full pass post-fix,
tsc --noEmitclean, and the rest of the package's test suite passing aside from one unrelated pre-existing timeout inblocks.test.ts). No verified issues; two informational observations (a defensive branch that appears unreachable given the store's mutual-exclusion contract, and an unchanged pre-existing event-publish ordering) were judged non-issues.Acceptance criteria (CL-7195)
releaseawaitsNote on the "guard" wording: no runtime guard was added for the multi-replica case. There's no reliable in-process signal for "am I replica N of M" here, and a guard gated on a config flag or env var nobody is required to set would be a safety net that lies. This repo's existing precedent (
agent-turns.ts'screateTurnFreedSignal) is a doc comment stating the single-replica assumption, not a runtime check — this PR follows that precedent. Happy to add an explicit guard if there's a concrete mechanism preferred here.Interchange defects found
None. This ticket only touches
packages/chatin this repo; no Interchange code was read or modified.