diff --git a/packages/chat/src/turn-queue.test.ts b/packages/chat/src/turn-queue.test.ts index 846eda26c..bdf4c6b9f 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,158 @@ 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"]); + }); + + // 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); + }); +}); diff --git a/packages/chat/src/turn-queue.ts b/packages/chat/src/turn-queue.ts index d5b921fbf..399d64d1f 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,58 +87,189 @@ 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 }, + ); +} + +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 { + // 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; + // 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; + } + + 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; + } + } 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); + } + } + 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]; + 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. `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 { - 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; + 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; } - } finally { - await deps.claims.release({ workbenchId }, token); + await drain(workbenchId, reclaimed, batch, dispatch); + } catch (err) { + reportClaimStoreFailure( + "chat.turnQueue.enqueueReclaim", + workbenchId, + err, + ); } }, };