Skip to content

Close the turn queue's enqueue/release race and dispatch strand - #474

Merged
TheGreatAxios merged 5 commits into
mainfrom
cl-7195-stop-the-turn-queue-from-silently-stranding-a-user-message
Aug 30, 2026
Merged

TheGreatAxios merged 5 commits into
mainfrom
cl-7195-stop-the-turn-queue-from-silently-stranding-a-user-message

Conversation

@TheGreatAxios

Copy link
Copy Markdown
Contributor

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 a finally, non-atomically. A concurrent run() call that failed tryClaim (the claim was still held) and pushed a turn into the pending map during the window before release() 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's release resolves synchronously before its first await. turn-claims.ts's own interface explicitly anticipates a durable implementation, and any store whose release awaits (e.g. Postgres-backed) reopens the window.

Changes

  • Symmetric two-sided reclaim. After the drain loop's release() settles, it rechecks the pending map; if a turn landed during that await, it reclaims via tryClaim and keeps draining. Symmetrically, run()'s enqueue path (when tryClaim initially fails and the turn is pushed) also retries tryClaim immediately 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.
  • dispatch can no longer strand the queue. dispatch(batch) is wrapped in 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 stopping the drain loop mid-batch.
  • Multi-replica behavior documented. pendingByWorkbench is a process-local Map paired 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 in agent-turns.ts), rather than building actual cross-replica correctness — out of scope for this ticket.

Tests

  • Two new regression tests reproduce the race directly, using a fake TurnClaimStore with 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.
  • A new test proves a rejecting dispatch no longer strands what queued behind it.
  • The existing "failed dispatch still releases the claim" test is updated: run() now resolves instead of rejecting when dispatch throws, 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 --noEmit clean, and the rest of the package's test suite passing aside from one unrelated pre-existing timeout in blocks.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)

  • Claim release and the emptiness check are atomic, so no enqueue can be lost between them
  • The dispatch call is wrapped so a rejection cannot strand the rest of the queue
  • The invariant "every enqueued turn is eventually dispatched" holds against a claim store whose release awaits
  • The audit's failing strand test is committed
  • Multi-replica behaviour is either made correct or explicitly documented as single-replica-only with a guard (documented; no runtime guard — see note below)

Note 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's createTurnFreedSignal) 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/chat in this repo; no Interchange code was read or modified.

@TheGreatAxios

Copy link
Copy Markdown
Contributor Author

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 release() settles, reclaim if something landed). That's insufficient alone: the enqueue path (tryClaim fails → push) has its own yield point, so a concurrent pusher whose tryClaim resolves false only after the holder has already released-and-returned still strands its turn. Closing the ticket's stated invariant needs both sides retrying tryClaim — releaser after release(), pusher after its own push — so whichever side notices the claim is free first drains what's pending. The shipped diff implements the symmetric version; the second regression test (does not strand a turn whose push lands only after the holder already returned) is exactly the case the one-sided version misses.

.catch/try-catch around dispatch with reportError is the right call and matches workbench-service.ts:1738's existing pattern. Single-replica documentation over a runtime guard is correct — there's no reliable in-process signal for replica count here, and agent-turns.ts already sets that precedent for this codebase.

@TheGreatAxios

Copy link
Copy Markdown
Contributor Author

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.

@TheGreatAxios

Copy link
Copy Markdown
Contributor Author

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 (drain's post-release() recheck) and the enqueue side (run's post-push tryClaim), and — matching what the two new regression tests already cover — whichever side notices the claim is free first drains what's pending; no interleaving leaves a turn in pendingByWorkbench with nobody holding the claim. Single-replica-only documentation without a runtime guard is the right call here (no in-process signal for replica count exists, and it matches agent-turns.ts's existing precedent).

One gap found and fixed: drain()'s loop only wrapped the dispatch call in try/catch. holds/release/tryClaim were unguarded, so if any of them rejected (rather than resolving false — a real possibility for the durable, I/O-backed store this ticket is about) the claim was never released and the rejection propagated out of run(), breaking its own documented "never rejects" contract and leaking the claim until the TTL backstop. Fixed in packages/chat/src/turn-queue.ts: both drain()'s loop and run()'s symmetric enqueue-side reclaim now catch these, report via reportError, and still make a best-effort release attempt. Added a regression test (still releases the claim and does not reject run() when the claim store throws mid-drain) that fails against the pre-fix code and passes now.

Verified: bunx prettier --check, WORKBENCH_CHECK_SINCE=origin/main bun run typecheck, eslint on the touched files, and bun test src/turn-queue.test.ts (10/10 pass) all clean.

@TheGreatAxios
TheGreatAxios force-pushed the cl-7195-stop-the-turn-queue-from-silently-stranding-a-user-message branch from a84807b to 82398b3 Compare August 30, 2026 21:36
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.
@TheGreatAxios
TheGreatAxios force-pushed the cl-7195-stop-the-turn-queue-from-silently-stranding-a-user-message branch from 82398b3 to 9a96831 Compare August 30, 2026 21:45
@TheGreatAxios
TheGreatAxios merged commit ec0a413 into main Aug 30, 2026
7 checks passed
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