From ac12c51a0d19c345a308f6a64cc8897909a0f41a Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 23:35:27 -0700 Subject: [PATCH 1/5] Add tests for turn queue enqueue/release race and dispatch rejections 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. --- packages/chat/src/turn-queue.test.ts | 143 +++++++++++++++++++++++++-- 1 file changed, 137 insertions(+), 6 deletions(-) diff --git a/packages/chat/src/turn-queue.test.ts b/packages/chat/src/turn-queue.test.ts index 846eda26c..60ad8ac5b 100644 --- a/packages/chat/src/turn-queue.test.ts +++ b/packages/chat/src/turn-queue.test.ts @@ -1,5 +1,6 @@ import { describe, expect, test } from "bun:test"; import { createInMemoryTurnClaimStore } from "./turn-claims"; +import type { TurnClaimStore } from "./turn-claims"; import { createWorkbenchTurnQueue, type QueuedTurn } from "./turn-queue"; import type { ChatWorkbenchEvent } from "./platform-port"; @@ -105,17 +106,18 @@ describe("createWorkbenchTurnQueue", () => { await firstRun; }); - test("a failed dispatch still releases the claim", async () => { + test("a failed dispatch still releases the claim and does not reject run()", async () => { const queue = createWorkbenchTurnQueue({ claims: createInMemoryTurnClaimStore({ ttlMs: 60_000 }), publish: () => undefined, }); - await expect( - queue.run("wb_1", turn("msg_1", "one"), async () => { - throw new Error("dispatch blew up"); - }), - ).rejects.toThrow("dispatch blew up"); + // `dispatch` must never reject per its own documented contract; a + // dispatch that breaks it is this queue's failure to contain, not + // something that should propagate out of `run()`. + await queue.run("wb_1", turn("msg_1", "one"), async () => { + throw new Error("dispatch blew up"); + }); // The claim was released despite the throw — a fresh turn can run. const dispatched: (readonly QueuedTurn[])[] = []; @@ -125,6 +127,34 @@ describe("createWorkbenchTurnQueue", () => { expect(dispatched).toEqual([[turn("msg_2", "two")]]); }); + test("a rejecting dispatch does not strand what queued behind it", async () => { + const queue = createWorkbenchTurnQueue({ + claims: createInMemoryTurnClaimStore({ ttlMs: 60_000 }), + publish: () => undefined, + }); + + let rejectFirst: ((err: Error) => void) | undefined; + const dispatched: string[] = []; + const dispatch = (batch: readonly QueuedTurn[]): Promise => { + if (batch[0]?.messageId === "msg_1") { + return new Promise((_resolve, reject) => { + rejectFirst = reject; + }); + } + for (const t of batch) dispatched.push(t.messageId); + return Promise.resolve(); + }; + + const firstRun = queue.run("wb_1", turn("msg_1", "one"), dispatch); + // msg_2 queues behind msg_1's still-open (and later rejecting) dispatch. + await queue.run("wb_1", turn("msg_2", "two"), dispatch); + + rejectFirst?.(new Error("dispatch blew up")); + await firstRun; + + expect(dispatched).toEqual(["msg_2"]); + }); + // CL-7129: the TTL is a crash/hang backstop for a dispatch that never // settles at all, not a license for two loops to drain the same // workbench's queue at once. Reproduces the reported shape: a message @@ -211,3 +241,104 @@ describe("createWorkbenchTurnQueue", () => { await firstRun; }); }); + +describe("turn queue drains everything that was enqueued", () => { + test("does not strand a turn enqueued while the holder is releasing", async () => { + const holders = new Map(); + let n = 0; + let onReleaseStart: (() => Promise) | undefined; + + const claims: TurnClaimStore = { + async tryClaim(c) { + if (holders.has(c.workbenchId)) return false; + const t = String(++n); + holders.set(c.workbenchId, t); + return t; + }, + async release(c, t) { + // A durable claim store awaits I/O here; anything can arrive first. + const hook = onReleaseStart; + onReleaseStart = undefined; + if (hook) await hook(); + if (holders.get(c.workbenchId) !== t) return false; + holders.delete(c.workbenchId); + return true; + }, + async holds(c, t) { + return holders.get(c.workbenchId) === t; + }, + }; + + const queue = createWorkbenchTurnQueue({ claims, publish: () => undefined }); + const dispatched: string[] = []; + const dispatch = async (batch: readonly QueuedTurn[]) => { + for (const t of batch) dispatched.push(t.messageId); + }; + + // m2 arrives after the drain loop saw an empty queue, before release lands. + onReleaseStart = async () => { + await queue.run("wb", turn("m2", "two"), dispatch); + }; + + await queue.run("wb", turn("m1", "one"), dispatch); + await new Promise((r) => setTimeout(r, 20)); + + expect(dispatched).toEqual(["m1", "m2"]); + }); + + // The symmetric half of the same race: the holder's `release()` can + // resolve and its own post-release recheck can find nothing pending + // *before* a concurrent enqueuer's failed `tryClaim` call ever + // resolves — the enqueuer's push then lands only after the holder is + // already gone and has already returned. Closing this requires the + // enqueue path to reclaim on its own, not just the release path. + test("does not strand a turn whose push lands only after the holder already returned", async () => { + const holders = new Map(); + let n = 0; + let releaseBlockedClaim: (() => void) | undefined; + + const claims: TurnClaimStore = { + async tryClaim(c) { + if (holders.has(c.workbenchId)) { + // Simulate a claim read that is still in flight (e.g. a slow + // durable store) when the holder's own release already lands. + await new Promise((resolve) => { + releaseBlockedClaim = resolve; + }); + return false; + } + const t = String(++n); + holders.set(c.workbenchId, t); + return t; + }, + async release(c, t) { + if (holders.get(c.workbenchId) !== t) return false; + holders.delete(c.workbenchId); + return true; + }, + async holds(c, t) { + return holders.get(c.workbenchId) === t; + }, + }; + + const queue = createWorkbenchTurnQueue({ claims, publish: () => undefined }); + const dispatched: string[] = []; + const dispatch = async (batch: readonly QueuedTurn[]) => { + for (const t of batch) dispatched.push(t.messageId); + }; + + const firstRun = queue.run("wb", turn("m1", "one"), dispatch); + const secondRun = queue.run("wb", turn("m2", "two"), dispatch); + + // m1 fully dispatches, sees nothing pending (m2's push hasn't + // landed yet — its tryClaim call is still parked), and releases. + await firstRun; + expect(dispatched).toEqual(["m1"]); + + // Only now does m2's tryClaim resolve `false` and its push execute. + releaseBlockedClaim?.(); + await secondRun; + + expect(dispatched).toEqual(["m1", "m2"]); + }); +}); From a35271e160f1f1df5edcadff2ad814a652967ad0 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 23:36:43 -0700 Subject: [PATCH 2/5] Close the turn queue's enqueue/release race and dispatch strand MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- packages/chat/src/turn-queue.ts | 185 ++++++++++++++++++++++++-------- 1 file changed, 138 insertions(+), 47 deletions(-) diff --git a/packages/chat/src/turn-queue.ts b/packages/chat/src/turn-queue.ts index d5b921fbf..469a17baa 100644 --- a/packages/chat/src/turn-queue.ts +++ b/packages/chat/src/turn-queue.ts @@ -14,12 +14,16 @@ // "release on completion" below means "release once this seam's own // dispatch call settles", not "once the agent finished replying" — // exactly the gap `./turn-claims.ts`'s TTL exists to backstop. +import { reportError } from "@corbits/error-sink"; +import { getLogger } from "@intx/log"; import { type } from "arktype"; import type { Part as PartType } from "./parts"; -import type { TurnClaimStore } from "./turn-claims"; +import type { TurnClaimStore, TurnClaimToken } from "./turn-claims"; import type { WorkbenchSubscriberRegistry } from "./workbench-events"; +const log = getLogger(["chat", "turn-queue"]); + /** * The non-persisted stream event a room's client renders as the queued * strip. Never written to the timeline — a queued message's own @@ -57,11 +61,11 @@ export type WorkbenchTurnQueueDeps = { * never reject — exactly like `routeMessage`'s own contract in * `workbench-service.ts`, which this wraps: a recipient that fails is * `dispatch`'s own job to report (an undelivered notice on the - * timeline), never something this queue has to catch. A `dispatch` that - * broke this contract would strand whatever queued behind it — the - * claim still releases (see `run`'s `finally`), but nothing would be - * left to notice and drain the leftover queue until the next unrelated - * message arrived. + * timeline), never something this queue has to catch. `run` below + * still guards this contract itself (a rejection is reported via + * `reportError` and treated as a settled, empty turn) so a `dispatch` + * that breaks it degrades to a logged defect rather than stranding + * whatever queued behind it. */ export type DispatchTurnBatch = (batch: readonly QueuedTurn[]) => Promise; @@ -83,59 +87,146 @@ export type WorkbenchTurnQueue = { ): Promise; }; +function reportDispatchFailure( + workbenchId: string, + batch: readonly QueuedTurn[], + err: unknown, +): void { + const messageIds = batch.map((t) => t.messageId); + const refId = reportError(err, { + operation: "chat.turnQueue.dispatch", + roomId: workbenchId, + extra: { messageIds }, + }); + // `dispatch` documents "must never reject" (see `DispatchTurnBatch`); + // reaching here means some caller broke that contract. + log.error( + "turn queue: dispatch rejected for workbench {workbenchId}, message(s) {messageIds} (ref {refId}), violating its \"never reject\" contract: {err}", + { workbenchId, messageIds, refId, err }, + ); +} + export function createWorkbenchTurnQueue( deps: WorkbenchTurnQueueDeps, ): WorkbenchTurnQueue { + // Process-local, paired one-to-one with `deps.claims`: a multi-replica + // hub would need every workbench's turns routed to the same replica, + // or replica B enqueues here into a `Map` replica A never reads and a + // turn queued on B is never dispatched at all. This queue is + // single-replica-only until the pending list moves into whatever + // store backs `deps.claims`. const pendingByWorkbench = new Map(); + // Runs `firstBatch` (already claimed under `token`) to completion, + // draining whatever else queues behind it under the same claim, then + // releases. Holds the claim across every batch rather than releasing + // and re-claiming between them: releasing early would open a window + // (each `await` is a yield point) where a fresh, unrelated call could + // win the claim ahead of a batch that was already queued and waiting + // — breaking the ordering this queue exists to guarantee. + // + // `dispatch` is the one `await` in this loop the TTL backstop + // (CL-7129) can fire during: a dispatch that outlives `ttlMs` lets a + // second, unrelated `run()` call win a fresh claim and start its own + // drain of this same `pendingByWorkbench` entry. The `holds` check + // right after each `dispatch` is how this loop notices that happened + // and stops — leaving whatever is still queued for the new holder to + // pick up — instead of popping and dispatching a batch the new holder + // is already about to dispatch itself. + async function drain( + workbenchId: string, + firstToken: TurnClaimToken, + firstBatch: readonly QueuedTurn[], + dispatch: DispatchTurnBatch, + ): Promise { + let token = firstToken; + let batch = firstBatch; + for (;;) { + try { + await dispatch(batch); + } catch (err) { + reportDispatchFailure(workbenchId, batch, err); + } + + if (!(await deps.claims.holds({ workbenchId }, token))) { + // The TTL backstop already reassigned this claim to a fresh + // `run()` call; that call's own loop now owns draining + // whatever is pending here. + return; + } + + const queued = pendingByWorkbench.get(workbenchId); + if (queued !== undefined && queued.length > 0) { + pendingByWorkbench.delete(workbenchId); + batch = queued; + continue; + } + + const released = await deps.claims.release({ workbenchId }, token); + if (!released) return; + + // `release` can itself await (any durable store reopens exactly + // this window — see `./turn-claims.ts`'s own doc comment): a + // concurrent `run()` whose `tryClaim` read us as still holding, + // before this `release` took effect, pushed onto + // `pendingByWorkbench` with nobody left holding the claim to + // drain it. Recheck and reclaim rather than stranding it — `run`'s + // enqueue path below runs the symmetric dance from its own side, + // so between the two, whichever notices the claim is free first + // picks up what's pending. + const late = pendingByWorkbench.get(workbenchId); + if (late === undefined || late.length === 0) return; + + const reclaimed = await deps.claims.tryClaim({ workbenchId }); + if (reclaimed === false) return; // the new holder's own loop will see it + + const afterReclaim = pendingByWorkbench.get(workbenchId); + if (afterReclaim === undefined || afterReclaim.length === 0) { + // Drained by whoever we raced against for the reclaim itself. + await deps.claims.release({ workbenchId }, reclaimed); + return; + } + pendingByWorkbench.delete(workbenchId); + token = reclaimed; + batch = afterReclaim; + } + } + return { async run(workbenchId, turn, dispatch) { const token = await deps.claims.tryClaim({ workbenchId }); - if (token === false) { - const queue = pendingByWorkbench.get(workbenchId) ?? []; - queue.push(turn); - pendingByWorkbench.set(workbenchId, queue); - deps.publish(workbenchId, { - type: "chat.turn-queued", - data: { - workbenchId, - messageId: turn.messageId, - queueLength: queue.length, - }, - }); + if (token !== false) { + await drain(workbenchId, token, [turn], dispatch); return; } - // Holds the claim across every batch it dispatches below, rather - // than releasing and re-claiming between them: releasing early - // would open a window (each `await` is a yield point) where a - // fresh, unrelated call could win the claim ahead of a batch that - // was already queued and waiting — breaking the ordering this - // queue exists to guarantee. The pending-queue check between - // batches runs with no `await` in between, so nothing can slip in - // during it. - // - // `dispatch` is the one `await` in this loop the TTL backstop can - // fire during (CL-7129): a dispatch that outlives `ttlMs` lets a - // second, unrelated `run()` call win a fresh claim and start its - // own drain of this same `pendingByWorkbench` entry. The `holds` - // check right after each `dispatch` is how this loop notices that - // happened and stops — leaving whatever is still queued for the - // new holder to pick up — instead of popping and dispatching a - // batch the new holder is already about to dispatch itself. - let batch: readonly QueuedTurn[] = [turn]; - try { - for (;;) { - await dispatch(batch); - if (!(await deps.claims.holds({ workbenchId }, token))) break; - const next = pendingByWorkbench.get(workbenchId); - if (next === undefined || next.length === 0) break; - pendingByWorkbench.delete(workbenchId); - batch = next; - } - } finally { - await deps.claims.release({ workbenchId }, token); + const queue = pendingByWorkbench.get(workbenchId) ?? []; + queue.push(turn); + pendingByWorkbench.set(workbenchId, queue); + deps.publish(workbenchId, { + type: "chat.turn-queued", + data: { + workbenchId, + messageId: turn.messageId, + queueLength: queue.length, + }, + }); + + // Symmetric to `drain`'s own reclaim above: the holder we lost + // `tryClaim` to may already have released — its own emptiness + // check ran before this push landed — leaving nobody holding the + // claim to drain what was just enqueued. Reclaim in that case + // rather than returning with the turn stranded. + const reclaimed = await deps.claims.tryClaim({ workbenchId }); + if (reclaimed === false) return; + + const batch = pendingByWorkbench.get(workbenchId); + pendingByWorkbench.delete(workbenchId); + if (batch === undefined || batch.length === 0) { + await deps.claims.release({ workbenchId }, reclaimed); + return; } + await drain(workbenchId, reclaimed, batch, dispatch); }, }; } From ab71166c319e29b58bb3cdbd382b2363deac0a90 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:55:16 -0700 Subject: [PATCH 3/5] Format turn-queue sources with prettier --- packages/chat/src/turn-queue.test.ts | 10 ++++++++-- packages/chat/src/turn-queue.ts | 2 +- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/packages/chat/src/turn-queue.test.ts b/packages/chat/src/turn-queue.test.ts index 60ad8ac5b..043919348 100644 --- a/packages/chat/src/turn-queue.test.ts +++ b/packages/chat/src/turn-queue.test.ts @@ -269,7 +269,10 @@ describe("turn queue drains everything that was enqueued", () => { }, }; - const queue = createWorkbenchTurnQueue({ claims, publish: () => undefined }); + const queue = createWorkbenchTurnQueue({ + claims, + publish: () => undefined, + }); const dispatched: string[] = []; const dispatch = async (batch: readonly QueuedTurn[]) => { for (const t of batch) dispatched.push(t.messageId); @@ -321,7 +324,10 @@ describe("turn queue drains everything that was enqueued", () => { }, }; - const queue = createWorkbenchTurnQueue({ claims, publish: () => undefined }); + const queue = createWorkbenchTurnQueue({ + claims, + publish: () => undefined, + }); const dispatched: string[] = []; const dispatch = async (batch: readonly QueuedTurn[]) => { for (const t of batch) dispatched.push(t.messageId); diff --git a/packages/chat/src/turn-queue.ts b/packages/chat/src/turn-queue.ts index 469a17baa..1bba31960 100644 --- a/packages/chat/src/turn-queue.ts +++ b/packages/chat/src/turn-queue.ts @@ -101,7 +101,7 @@ function reportDispatchFailure( // `dispatch` documents "must never reject" (see `DispatchTurnBatch`); // reaching here means some caller broke that contract. log.error( - "turn queue: dispatch rejected for workbench {workbenchId}, message(s) {messageIds} (ref {refId}), violating its \"never reject\" contract: {err}", + 'turn queue: dispatch rejected for workbench {workbenchId}, message(s) {messageIds} (ref {refId}), violating its "never reject" contract: {err}', { workbenchId, messageIds, refId, err }, ); } From deee7bb17607aa71681f3d05b7a3bdb8a6f8c3a8 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 05:30:12 -0700 Subject: [PATCH 4/5] Add a test for the claim store rejecting mid-drain 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. --- packages/chat/src/turn-queue.test.ts | 48 ++++++++++++++++++++++++++++ 1 file changed, 48 insertions(+) diff --git a/packages/chat/src/turn-queue.test.ts b/packages/chat/src/turn-queue.test.ts index 043919348..bdf4c6b9f 100644 --- a/packages/chat/src/turn-queue.test.ts +++ b/packages/chat/src/turn-queue.test.ts @@ -347,4 +347,52 @@ describe("turn queue drains everything that was enqueued", () => { expect(dispatched).toEqual(["m1", "m2"]); }); + + // CL-7195's own reasoning ("a durable store reopens the window") + // implies these calls can also reject outright, not just resolve + // `false` — a real possibility for I/O, unlike the in-memory store. + // `run()` documents "never rejects", so the drain loop's claim-store + // calls must not leak that rejection, and must still attempt to free + // the claim rather than leaving it held until the TTL backstop. + test("still releases the claim and does not reject run() when the claim store throws mid-drain", async () => { + const holders = new Map(); + let n = 0; + let releaseCalls = 0; + + const claims: TurnClaimStore = { + async tryClaim(c) { + if (holders.has(c.workbenchId)) return false; + const t = String(++n); + holders.set(c.workbenchId, t); + return t; + }, + async release(c, t) { + releaseCalls++; + if (holders.get(c.workbenchId) !== t) return false; + holders.delete(c.workbenchId); + return true; + }, + async holds() { + throw new Error("claim store connection reset"); + }, + }; + + const queue = createWorkbenchTurnQueue({ + claims, + publish: () => undefined, + }); + const dispatched: string[] = []; + const dispatch = async (batch: readonly QueuedTurn[]) => { + for (const t of batch) dispatched.push(t.messageId); + }; + + // Does not reject, despite `holds` throwing. + await queue.run("wb", turn("m1", "one"), dispatch); + + expect(dispatched).toEqual(["m1"]); + expect(releaseCalls).toBeGreaterThan(0); + // The claim was freed despite the throw — a fresh turn can claim it + // directly rather than only being queued. + expect(holders.has("wb")).toBe(false); + }); }); From 9a968315e396679208d05955781ddaa90b0587ea Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 05:30:21 -0700 Subject: [PATCH 5/5] Keep the claim store's release contract when its own calls reject 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. --- packages/chat/src/turn-queue.ts | 152 ++++++++++++++++++++------------ 1 file changed, 98 insertions(+), 54 deletions(-) diff --git a/packages/chat/src/turn-queue.ts b/packages/chat/src/turn-queue.ts index 1bba31960..399d64d1f 100644 --- a/packages/chat/src/turn-queue.ts +++ b/packages/chat/src/turn-queue.ts @@ -106,6 +106,21 @@ function reportDispatchFailure( ); } +function reportClaimStoreFailure( + operation: string, + workbenchId: string, + err: unknown, +): void { + const refId = reportError(err, { + operation, + roomId: workbenchId, + }); + log.error( + "turn queue: claim store call failed for workbench {workbenchId} (ref {refId}): {err}", + { workbenchId, refId, err }, + ); +} + export function createWorkbenchTurnQueue( deps: WorkbenchTurnQueueDeps, ): WorkbenchTurnQueue { @@ -141,54 +156,72 @@ export function createWorkbenchTurnQueue( ): Promise { let token = firstToken; let batch = firstBatch; - for (;;) { - try { - await dispatch(batch); - } catch (err) { - reportDispatchFailure(workbenchId, batch, err); - } + // The claim store's own contract (`./turn-claims.ts`) requires a won + // claim to be "never left un-called on a path that can still + // complete" — the outer try/catch is what makes that hold even when + // `holds`/`release`/`tryClaim` themselves reject (a durable store + // doing I/O can, unlike the in-memory one), not just when they + // resolve normally. `dispatch` gets its own inner try/catch because + // a rejecting dispatch is an expected, contract-violating case this + // loop must keep draining through, not stop for. + try { + for (;;) { + try { + await dispatch(batch); + } catch (err) { + reportDispatchFailure(workbenchId, batch, err); + } - if (!(await deps.claims.holds({ workbenchId }, token))) { - // The TTL backstop already reassigned this claim to a fresh - // `run()` call; that call's own loop now owns draining - // whatever is pending here. - return; - } + if (!(await deps.claims.holds({ workbenchId }, token))) { + // The TTL backstop already reassigned this claim to a fresh + // `run()` call; that call's own loop now owns draining + // whatever is pending here. + return; + } - const queued = pendingByWorkbench.get(workbenchId); - if (queued !== undefined && queued.length > 0) { - pendingByWorkbench.delete(workbenchId); - batch = queued; - continue; - } + const queued = pendingByWorkbench.get(workbenchId); + if (queued !== undefined && queued.length > 0) { + pendingByWorkbench.delete(workbenchId); + batch = queued; + continue; + } - const released = await deps.claims.release({ workbenchId }, token); - if (!released) return; - - // `release` can itself await (any durable store reopens exactly - // this window — see `./turn-claims.ts`'s own doc comment): a - // concurrent `run()` whose `tryClaim` read us as still holding, - // before this `release` took effect, pushed onto - // `pendingByWorkbench` with nobody left holding the claim to - // drain it. Recheck and reclaim rather than stranding it — `run`'s - // enqueue path below runs the symmetric dance from its own side, - // so between the two, whichever notices the claim is free first - // picks up what's pending. - const late = pendingByWorkbench.get(workbenchId); - if (late === undefined || late.length === 0) return; - - const reclaimed = await deps.claims.tryClaim({ workbenchId }); - if (reclaimed === false) return; // the new holder's own loop will see it - - const afterReclaim = pendingByWorkbench.get(workbenchId); - if (afterReclaim === undefined || afterReclaim.length === 0) { - // Drained by whoever we raced against for the reclaim itself. - await deps.claims.release({ workbenchId }, reclaimed); - return; + const released = await deps.claims.release({ workbenchId }, token); + if (!released) return; + + // `release` can itself await (any durable store reopens exactly + // this window — see `./turn-claims.ts`'s own doc comment): a + // concurrent `run()` whose `tryClaim` read us as still holding, + // before this `release` took effect, pushed onto + // `pendingByWorkbench` with nobody left holding the claim to + // drain it. Recheck and reclaim rather than stranding it — `run`'s + // enqueue path below runs the symmetric dance from its own side, + // so between the two, whichever notices the claim is free first + // picks up what's pending. + const late = pendingByWorkbench.get(workbenchId); + if (late === undefined || late.length === 0) return; + + const reclaimed = await deps.claims.tryClaim({ workbenchId }); + if (reclaimed === false) return; // the new holder's own loop will see it + + const afterReclaim = pendingByWorkbench.get(workbenchId); + if (afterReclaim === undefined || afterReclaim.length === 0) { + // Drained by whoever we raced against for the reclaim itself. + await deps.claims.release({ workbenchId }, reclaimed); + return; + } + pendingByWorkbench.delete(workbenchId); + token = reclaimed; + batch = afterReclaim; } - pendingByWorkbench.delete(workbenchId); - token = reclaimed; - batch = afterReclaim; + } catch (err) { + reportClaimStoreFailure("chat.turnQueue.drain", workbenchId, err); + // Best-effort: `token` may already be released (this call then + // no-ops per the store's contract) or may be the one this loop + // never got to release because the failure happened first either + // way, attempting it is strictly better than leaking it until the + // TTL backstop expires. + await deps.claims.release({ workbenchId }, token).catch(() => undefined); } } @@ -216,17 +249,28 @@ export function createWorkbenchTurnQueue( // `tryClaim` to may already have released — its own emptiness // check ran before this push landed — leaving nobody holding the // claim to drain what was just enqueued. Reclaim in that case - // rather than returning with the turn stranded. - const reclaimed = await deps.claims.tryClaim({ workbenchId }); - if (reclaimed === false) return; - - const batch = pendingByWorkbench.get(workbenchId); - pendingByWorkbench.delete(workbenchId); - if (batch === undefined || batch.length === 0) { - await deps.claims.release({ workbenchId }, reclaimed); - return; + // rather than returning with the turn stranded. `run` documents + // "never rejects", so a claim store that throws here (rather than + // resolving `false`) is caught and reported, same as `drain`'s own + // guard, instead of propagating out. + try { + const reclaimed = await deps.claims.tryClaim({ workbenchId }); + if (reclaimed === false) return; + + const batch = pendingByWorkbench.get(workbenchId); + pendingByWorkbench.delete(workbenchId); + if (batch === undefined || batch.length === 0) { + await deps.claims.release({ workbenchId }, reclaimed); + return; + } + await drain(workbenchId, reclaimed, batch, dispatch); + } catch (err) { + reportClaimStoreFailure( + "chat.turnQueue.enqueueReclaim", + workbenchId, + err, + ); } - await drain(workbenchId, reclaimed, batch, dispatch); }, }; }