diff --git a/apps/hub/src/index.ts b/apps/hub/src/index.ts index 66db0cc6d..810956806 100644 --- a/apps/hub/src/index.ts +++ b/apps/hub/src/index.ts @@ -67,7 +67,7 @@ import { import { AGENT_SECTION_MODE, - CHAT_TURN_TIMEOUT_MS, + DEFAULT_TURN_CLAIM_TTL_MS, createArtifactDeliveryHandler, createDrizzleAgentTurnStore, createInMemoryTurnClaimStore, @@ -1310,7 +1310,7 @@ export async function createHub(config: HubConfig) { // still serializes against the others rather than each queue only // seeing its own slice of the traffic. const turnQueue = createWorkbenchTurnQueue({ - claims: createInMemoryTurnClaimStore({ ttlMs: CHAT_TURN_TIMEOUT_MS }), + claims: createInMemoryTurnClaimStore({ ttlMs: DEFAULT_TURN_CLAIM_TTL_MS }), publish: workbenchSubscribers.publish, }); // The room timeline store (CL-6327): a workbench's own messages, held @@ -1532,7 +1532,7 @@ export async function createHub(config: HubConfig) { conditionRegistry: chatConditionRegistry, }), isInvitableDefinition: isPickerListableDefinition, - turnTimeoutMs: CHAT_TURN_TIMEOUT_MS, + turnTimeoutMs: DEFAULT_TURN_CLAIM_TTL_MS, resolvePrincipalName: async (_tenantId, principalId) => { const principalRow = await db.query.principal.findFirst({ where: (p, { eq: equals }) => equals(p.id, principalId), diff --git a/packages/chat/src/index.ts b/packages/chat/src/index.ts index fd9692b23..098acb821 100644 --- a/packages/chat/src/index.ts +++ b/packages/chat/src/index.ts @@ -128,7 +128,7 @@ export { CHAT_TURN_TIMEOUT_MS, createInMemoryTurnClaimStore, } from "./turn-claims"; -export type { TurnClaim, TurnClaimStore } from "./turn-claims"; +export type { TurnClaim, TurnClaimStore, TurnClaimToken } from "./turn-claims"; export { AGENT_SECTION_MODE, workbenchLaunchPersistExtra, @@ -227,6 +227,9 @@ export { dispatchTurn, DEFAULT_TURN_DISPATCH_TIMEOUT_MS, turnDispatchTimeoutMessage, + DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS, + DEFAULT_TURN_CLAIM_TTL_MS, + waitUntilFreeTimeoutMessage, launchAndJoinAgent, postCannedGreeting, cannedGreeting, diff --git a/packages/chat/src/routes.ts b/packages/chat/src/routes.ts index 196da3ab5..a469b9cef 100644 --- a/packages/chat/src/routes.ts +++ b/packages/chat/src/routes.ts @@ -146,7 +146,14 @@ export type CreateChatRoutesDeps = { * schedulable workflows masquerade as chat partners. */ isInvitableDefinition: (definition: InvitableDefinitionRecord) => boolean; - /** Per-turn timeout, the default write-claim TTL. */ + /** + * The default turn-claim TTL — see + * `./workbench-service.ts`'s `DEFAULT_TURN_CLAIM_TTL_MS` for the + * production default and the arithmetic it has to clear + * (`waitUntilFreeTimeoutMs` + `turnDispatchTimeoutMs`, with margin) + * to stay an unreachable backstop rather than a bound that fires on + * a still-legitimate dispatch. + */ turnTimeoutMs: number; /** * CL-6644's turn-level deadline: see `SendWorkbenchMessageDeps`'s field @@ -154,6 +161,12 @@ export type CreateChatRoutesDeps = { * `dispatchTurnBatch` uses `DEFAULT_TURN_DISPATCH_TIMEOUT_MS`. */ turnDispatchTimeoutMs?: number; + /** + * CL-7129's bound on the CL-6670 wait: see `SendWorkbenchMessageDeps`'s + * field of the same name in `./workbench-service.ts`. Omitted, + * `dispatchTurnBatch` uses `DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS`. + */ + waitUntilFreeTimeoutMs?: number; /** * Resolves a principal to the display name a greeting can use. The * hub wires this to its user table; omitted, the canned greeting @@ -2162,6 +2175,9 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { ...(deps.turnDispatchTimeoutMs !== undefined ? { turnDispatchTimeoutMs: deps.turnDispatchTimeoutMs } : {}), + ...(deps.waitUntilFreeTimeoutMs !== undefined + ? { waitUntilFreeTimeoutMs: deps.waitUntilFreeTimeoutMs } + : {}), }, { tenantId: ownerTenantId, @@ -2308,6 +2324,9 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { ...(deps.turnDispatchTimeoutMs !== undefined ? { turnDispatchTimeoutMs: deps.turnDispatchTimeoutMs } : {}), + ...(deps.waitUntilFreeTimeoutMs !== undefined + ? { waitUntilFreeTimeoutMs: deps.waitUntilFreeTimeoutMs } + : {}), }, { tenantId: ownerTenantId, diff --git a/packages/chat/src/turn-claims.test.ts b/packages/chat/src/turn-claims.test.ts index 59985e578..c2b01e952 100644 --- a/packages/chat/src/turn-claims.test.ts +++ b/packages/chat/src/turn-claims.test.ts @@ -2,9 +2,10 @@ import { describe, expect, test } from "bun:test"; import { createInMemoryTurnClaimStore } from "./turn-claims"; describe("createInMemoryTurnClaimStore", () => { - test("the first tryClaim for a workbench wins", async () => { + test("the first tryClaim for a workbench wins, returning a token", async () => { const store = createInMemoryTurnClaimStore({ ttlMs: 60_000 }); - await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.toBe(true); + const token = await store.tryClaim({ workbenchId: "wb_1" }); + expect(typeof token).toBe("string"); }); test("a second tryClaim for the same workbench loses while the first is held", async () => { @@ -15,22 +16,52 @@ describe("createInMemoryTurnClaimStore", () => { test("claims on different workbenches never contend", async () => { const store = createInMemoryTurnClaimStore({ ttlMs: 60_000 }); - await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.toBe(true); - await expect(store.tryClaim({ workbenchId: "wb_2" })).resolves.toBe(true); + await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.not.toBe( + false, + ); + await expect(store.tryClaim({ workbenchId: "wb_2" })).resolves.not.toBe( + false, + ); }); test("release frees the claim for a fresh tryClaim", async () => { const store = createInMemoryTurnClaimStore({ ttlMs: 60_000 }); - await store.tryClaim({ workbenchId: "wb_1" }); - await store.release({ workbenchId: "wb_1" }); - await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.toBe(true); + const token = await store.tryClaim({ workbenchId: "wb_1" }); + await expect( + store.release({ workbenchId: "wb_1" }, token as string), + ).resolves.toBe(true); + await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.not.toBe( + false, + ); }); test("releasing a claim nobody holds is a harmless no-op", async () => { const store = createInMemoryTurnClaimStore({ ttlMs: 60_000 }); await expect( - store.release({ workbenchId: "wb_never_claimed" }), - ).resolves.toBeUndefined(); + store.release({ workbenchId: "wb_never_claimed" }, "some_token"), + ).resolves.toBe(false); + }); + + test("releasing with a stale token — one the TTL already reassigned — is a no-op", async () => { + let now = 0; + const store = createInMemoryTurnClaimStore({ + ttlMs: 1_000, + now: () => now, + }); + const firstToken = await store.tryClaim({ workbenchId: "wb_1" }); + now += 1_000; // the TTL elapses while the first holder is still "in flight" + const secondToken = await store.tryClaim({ workbenchId: "wb_1" }); + expect(secondToken).not.toBe(false); + expect(secondToken).not.toBe(firstToken); + + // The first holder's `finally` releasing its now-stale token must + // never evict the second holder's live claim. + await expect( + store.release({ workbenchId: "wb_1" }, firstToken as string), + ).resolves.toBe(false); + await expect( + store.holds({ workbenchId: "wb_1" }, secondToken as string), + ).resolves.toBe(true); }); test("a claim older than the TTL is reclaimable even without release — the crash/hang backstop", async () => { @@ -43,6 +74,38 @@ describe("createInMemoryTurnClaimStore", () => { now += 999; await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.toBe(false); now += 2; - await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.toBe(true); + await expect(store.tryClaim({ workbenchId: "wb_1" })).resolves.not.toBe( + false, + ); + }); + + describe("holds", () => { + test("is true for the current, unexpired holder's own token", async () => { + const store = createInMemoryTurnClaimStore({ ttlMs: 60_000 }); + const token = await store.tryClaim({ workbenchId: "wb_1" }); + await expect( + store.holds({ workbenchId: "wb_1" }, token as string), + ).resolves.toBe(true); + }); + + test("is false for a token nobody ever held", async () => { + const store = createInMemoryTurnClaimStore({ ttlMs: 60_000 }); + await expect( + store.holds({ workbenchId: "wb_1" }, "never_claimed"), + ).resolves.toBe(false); + }); + + test("is false once the TTL expires, even without a fresh tryClaim", async () => { + let now = 0; + const store = createInMemoryTurnClaimStore({ + ttlMs: 1_000, + now: () => now, + }); + const token = await store.tryClaim({ workbenchId: "wb_1" }); + now += 1_000; + await expect( + store.holds({ workbenchId: "wb_1" }, token as string), + ).resolves.toBe(false); + }); }); }); diff --git a/packages/chat/src/turn-claims.ts b/packages/chat/src/turn-claims.ts index 272c71cdf..9dabfacc8 100644 --- a/packages/chat/src/turn-claims.ts +++ b/packages/chat/src/turn-claims.ts @@ -21,31 +21,57 @@ export type TurnClaim = { readonly workbenchId: string; }; +/** + * Proof of having won a `tryClaim` call — opaque to callers, meaningful + * only back to the store that issued it. Required by `release` (and by + * `turn-queue.ts`'s own `holds` check) so a claim the TTL already + * reassigned to a second winner can never be released, or mistaken for + * still-held, by the first winner's stale reference (CL-7129). + */ +export type TurnClaimToken = string; + export type TurnClaimStore = { /** - * Atomically claims `workbenchId`: `true` if this call won the claim - * (the caller should run the next turn), `false` if a turn is - * already in flight for this workbench (the caller must queue - * instead, never dispatch). + * Atomically claims `workbenchId`: a token if this call won the claim + * (the caller should run the next turn, and must present this token + * to `release`/`holds`), `false` if a turn is already in flight for + * this workbench (the caller must queue instead, never dispatch). */ - tryClaim(claim: TurnClaim): Promise; + tryClaim(claim: TurnClaim): Promise; /** * Un-claims `workbenchId` — called once the turn that won the claim * has settled, success or failure alike, so the next queued batch (or - * a fresh message) can run. Never left un-called on a path that can - * still complete; the TTL below is only the backstop for a dispatch - * that never settles at all. + * a fresh message) can run. A no-op (returns `false`) if `token` is + * not the current holder — e.g. the TTL already reassigned the claim + * to a second winner, in which case releasing here would free a + * claim this caller no longer owns. Never left un-called on a path + * that can still complete; the TTL below is only the backstop for a + * dispatch that never settles at all. + */ + release(claim: TurnClaim, token: TurnClaimToken): Promise; + /** + * Whether `token` is still the current, unexpired holder of + * `workbenchId`'s claim — how a long-running holder notices the TTL + * reassigned its claim to someone else without itself calling + * `tryClaim` again (which would attempt to *take* the claim, not just + * check it). */ - release(claim: TurnClaim): Promise; + holds(claim: TurnClaim, token: TurnClaimToken): Promise; }; /** - * How long one agent turn may run before the room stops waiting on it. - * One number, two enforcement points that must never drift: the section - * body's own per-occurrence `timeout` (CL-6329, pinned into a room - * agent's deployed definition by `./platform-adapter.ts`) and the claim - * TTL that stops a workbench wedging behind a dispatch that never - * settles. + * How long one agent turn may run before the room stops waiting on it — + * the section body's own per-occurrence `timeout` (CL-6329, pinned into + * a room agent's deployed definition by `./platform-adapter.ts`). + * + * CL-7129: this is also the floor `./workbench-service.ts`'s + * `DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS` is built from — that wait has to + * tolerate a prior turn legitimately running the full length this + * constant allows. The claim TTL is no longer this same number (it + * used to be); it now has to be strictly longer than that wait bound + * plus the dispatch deadline, so see + * `./workbench-service.ts`'s `DEFAULT_TURN_CLAIM_TTL_MS` for the actual + * TTL callers should use. */ export const CHAT_TURN_TIMEOUT_MS = 5 * 60 * 1000; @@ -66,18 +92,38 @@ export function createInMemoryTurnClaimStore(options: { readonly now?: () => number; }): TurnClaimStore { const now = options.now ?? Date.now; - const claimedAt = new Map(); + const holders = new Map(); + let nextToken = 0; + + function isExpired(holder: { claimedAt: number }): boolean { + return now() - holder.claimedAt >= options.ttlMs; + } + return { async tryClaim(claim) { - const existing = claimedAt.get(claim.workbenchId); - if (existing !== undefined && now() - existing < options.ttlMs) { + const existing = holders.get(claim.workbenchId); + if (existing !== undefined && !isExpired(existing)) { + return false; + } + const token = String(++nextToken); + holders.set(claim.workbenchId, { token, claimedAt: now() }); + return token; + }, + async release(claim, token) { + const existing = holders.get(claim.workbenchId); + if (existing === undefined || existing.token !== token) { return false; } - claimedAt.set(claim.workbenchId, now()); + holders.delete(claim.workbenchId); return true; }, - async release(claim) { - claimedAt.delete(claim.workbenchId); + async holds(claim, token) { + const existing = holders.get(claim.workbenchId); + return ( + existing !== undefined && + existing.token === token && + !isExpired(existing) + ); }, }; } diff --git a/packages/chat/src/turn-queue.test.ts b/packages/chat/src/turn-queue.test.ts index 119de24a1..846eda26c 100644 --- a/packages/chat/src/turn-queue.test.ts +++ b/packages/chat/src/turn-queue.test.ts @@ -125,6 +125,73 @@ describe("createWorkbenchTurnQueue", () => { expect(dispatched).toEqual([[turn("msg_2", "two")]]); }); + // 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 + // queues normally behind an in-flight turn, the TTL then elapses + // while that turn is still open, a second `run()` wins a fresh claim + // via the backstop and starts its own loop — and only ONE of the two + // loops may ever go on to pop and dispatch the queued message. + test("a claim's TTL expiring mid-dispatch stops the old loop from also draining what queued behind it", async () => { + let now = 0; + const claims = createInMemoryTurnClaimStore({ + ttlMs: 1_000, + now: () => now, + }); + const queue = createWorkbenchTurnQueue({ + claims, + publish: () => undefined, + }); + + const calls: { batch: readonly QueuedTurn[]; resolve: () => void }[] = []; + const dispatch = (batch: readonly QueuedTurn[]) => + new Promise((resolve) => { + calls.push({ batch, resolve }); + }); + + // The first run() wins the claim and starts a turn that outlives + // the claim TTL while still "in flight" — the crash/hang shape the + // TTL backstop exists for. + const firstRun = queue.run("wb_1", turn("msg_1", "a"), dispatch); + await Promise.resolve(); + expect(calls).toHaveLength(1); + + // A second message arrives before the TTL elapses: the first loop + // is still the valid holder, so this queues normally. + await queue.run("wb_1", turn("msg_2", "b"), dispatch); + + // The claim's TTL elapses while the first turn is still open. + now += 1_000; + + // A third, unrelated run() for the same workbench now wins a fresh + // claim via the TTL backstop and starts its own loop, dispatching + // its own turn first. + const thirdRun = queue.run("wb_1", turn("msg_3", "c"), dispatch); + await Promise.resolve(); + expect(calls).toHaveLength(2); + expect(calls[1]?.batch).toEqual([turn("msg_3", "c")]); + + // The first loop's own dispatch finally settles. Before this fix, + // it would go on to pop and dispatch "b" itself here — a second, + // concurrent drain of the very queue the third loop now owns. The + // fix's `holds` check makes it notice it lost the claim and stop. + calls[0]?.resolve(); + await firstRun; + expect(calls).toHaveLength(2); + + // The third loop's own dispatch settles next, picks up "b" (still + // sitting in the shared queue, untouched by the first loop) exactly + // once, and drains it under its own, still-valid claim. + calls[1]?.resolve(); + await Promise.resolve(); + await Promise.resolve(); + expect(calls).toHaveLength(3); + expect(calls[2]?.batch).toEqual([turn("msg_2", "b")]); + + calls[2]?.resolve(); + await thirdRun; + }); + test("workbenches never contend with each other", async () => { const queue = createWorkbenchTurnQueue({ claims: createInMemoryTurnClaimStore({ ttlMs: 60_000 }), diff --git a/packages/chat/src/turn-queue.ts b/packages/chat/src/turn-queue.ts index 786cef40b..e362b6301 100644 --- a/packages/chat/src/turn-queue.ts +++ b/packages/chat/src/turn-queue.ts @@ -86,8 +86,8 @@ export function createWorkbenchTurnQueue( return { async run(workbenchId, turn, dispatch) { - const claimed = await deps.claims.tryClaim({ workbenchId }); - if (!claimed) { + const token = await deps.claims.tryClaim({ workbenchId }); + if (token === false) { const queue = pendingByWorkbench.get(workbenchId) ?? []; queue.push(turn); pendingByWorkbench.set(workbenchId, queue); @@ -110,17 +110,27 @@ export function createWorkbenchTurnQueue( // 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 }); + await deps.claims.release({ workbenchId }, token); } }, }; diff --git a/packages/chat/src/workbench-service.timeouts.test.ts b/packages/chat/src/workbench-service.timeouts.test.ts new file mode 100644 index 000000000..73791cc83 --- /dev/null +++ b/packages/chat/src/workbench-service.timeouts.test.ts @@ -0,0 +1,39 @@ +// CL-7129: the constants a claim's backstop TTL and the CL-6670 wait +// bound are built from have to clear specific arithmetic, not just +// exist — a value that merely narrows the old unbounded wait can still +// undercut `CHAT_TURN_TIMEOUT_MS`, the documented max length a +// legitimate prior turn is allowed to run (see `./turn-claims.test.ts` +// for the store's own TTL-reclaim mechanics; this file is about the +// numbers those mechanics are configured with). +import { describe, expect, test } from "bun:test"; + +import { CHAT_TURN_TIMEOUT_MS } from "./turn-claims"; +import { + DEFAULT_TURN_DISPATCH_TIMEOUT_MS, + DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS, + DEFAULT_TURN_CLAIM_TTL_MS, +} from "./workbench-service"; + +describe("CL-7129's timeout constants", () => { + test("the CL-6670 wait bound covers the longest a legitimate prior turn may run", () => { + // A prior turn is allowed to run the full `CHAT_TURN_TIMEOUT_MS` + // (the section body's own per-occurrence timeout). A wait bound at + // or below that would time out a message queued behind such a + // turn and drop it as undelivered instead of waiting it out. + expect(DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS).toBeGreaterThan( + CHAT_TURN_TIMEOUT_MS, + ); + }); + + test("the claim TTL exceeds the wait bound plus the dispatch deadline, so it never fires on a well-behaved dispatch", () => { + // A well-behaved dispatch spends, in the worst case, the wait bound + // waiting on a prior turn and then the dispatch deadline making its + // own call — back to back, inside one claim's lifetime. The TTL has + // to clear that sum (with margin) to stay an unreachable backstop + // rather than a bound that reassigns a claim a live dispatch still + // holds. + expect(DEFAULT_TURN_CLAIM_TTL_MS).toBeGreaterThan( + DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS + DEFAULT_TURN_DISPATCH_TIMEOUT_MS, + ); + }); +}); diff --git a/packages/chat/src/workbench-service.ts b/packages/chat/src/workbench-service.ts index ee38afc31..025e91712 100644 --- a/packages/chat/src/workbench-service.ts +++ b/packages/chat/src/workbench-service.ts @@ -44,6 +44,7 @@ import type { import { postRoomMessage, type RoomMessageStore } from "./room-messages"; import type { WorkbenchSubscriberRegistry } from "./workbench-events"; import type { QueuedTurn, WorkbenchTurnQueue } from "./turn-queue"; +import { CHAT_TURN_TIMEOUT_MS } from "./turn-claims"; import type { WorkbenchTenancyStore } from "./workbench-tenancy"; import type { ChatStore } from "./store"; import { withTimeout } from "./with-timeout"; @@ -910,6 +911,18 @@ export type SendWorkbenchMessageDeps = { * the production default. */ readonly turnDispatchTimeoutMs?: number; + /** + * CL-7129's bound on the CL-6670 wait: `dispatchTurnBatch` wraps + * `agentTurns.waitUntilFree` in this budget before it ever starts the + * `turnDispatchTimeoutMs` clock below, defaulting to + * `DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS`. Left unbounded, a slow prior + * turn for the same agent could hold `turn-queue.ts`'s claim past its + * TTL while this `await` was still open — the exact gap that let a + * second `run()` win the claim and start a second concurrent drain. + * Injectable so tests exercise the bound in milliseconds instead of + * the production default. + */ + readonly waitUntilFreeTimeoutMs?: number; }; /** CL-6644's default turn-level deadline: generous enough to cover a @@ -917,6 +930,46 @@ export type SendWorkbenchMessageDeps = { * `DEFAULT_WAKE_TIMEOUT_MS` (30s) uses for the wake step alone. */ export const DEFAULT_TURN_DISPATCH_TIMEOUT_MS = 120_000; +/** + * CL-7129's default bound on `agentTurns.waitUntilFree`: must cover the + * longest a legitimate prior turn is allowed to run — + * `CHAT_TURN_TIMEOUT_MS`, the same bound the section body's own + * per-occurrence timeout enforces — plus a grace margin. A message + * queued behind a prior turn that is merely slow, not hung, has to be + * delivered once that turn frees up rather than time out here and land + * as an undelivered notice (CL-6670's exact case; a bound narrower than + * `CHAT_TURN_TIMEOUT_MS` reintroduces the bug CL-6670 fixed for any + * prior turn that runs past it). + */ +export const DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS = CHAT_TURN_TIMEOUT_MS + 30_000; + +/** + * CL-7129's default turn-claim TTL: the backstop for a workbench whose + * dispatch loop crashed or hung without ever calling `release` (see + * `./turn-claims.ts`'s `createInMemoryTurnClaimStore`). A well-behaved + * dispatch runs `DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS`'s wait and + * `DEFAULT_TURN_DISPATCH_TIMEOUT_MS`'s dispatch call back-to-back + * inside one claim's lifetime, so the TTL has to clear the sum of both + * with margin to spare — otherwise the TTL fires on a claim a legitimate + * dispatch is still using, letting a second `run()` win it and drain + * the same workbench's queue concurrently (the CL-7129 bug this whole + * store change exists to close). Kept strictly greater by a 30s margin + * so it is unreachable in any well-behaved case. + */ +export const DEFAULT_TURN_CLAIM_TTL_MS = + DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS + + DEFAULT_TURN_DISPATCH_TIMEOUT_MS + + 30_000; + +/** `waitUntilFree`'s own timeout message: names the agent address the + * wait was blocked on and the budget it exceeded. */ +export function waitUntilFreeTimeoutMessage( + agentAddress: string, + timeoutMs: number, +): string { + return `waiting for "${agentAddress}"'s prior turn to close did not settle within ${String(timeoutMs)}ms`; +} + /** The turn-level deadline's own rejection message: names the turn's * run address (the recipient every dispatch failure is already reported * and notified against) and the budget it exceeded, so a person reading @@ -1224,6 +1277,7 @@ async function dispatchTurnBatch( | "roomMessages" | "publish" | "turnDispatchTimeoutMs" + | "waitUntilFreeTimeoutMs" | "agentTurns" >, tenantId: string, @@ -1232,6 +1286,8 @@ async function dispatchTurnBatch( ): Promise { const turnDispatchTimeoutMs = deps.turnDispatchTimeoutMs ?? DEFAULT_TURN_DISPATCH_TIMEOUT_MS; + const waitUntilFreeTimeoutMs = + deps.waitUntilFreeTimeoutMs ?? DEFAULT_WAIT_UNTIL_FREE_TIMEOUT_MS; const recipientSet = new Set(); for (const turn of batch) { for (const agentAddress of turn.recipients) recipientSet.add(agentAddress); @@ -1259,24 +1315,38 @@ async function dispatchTurnBatch( await Promise.all( recipients.map(async (agentAddress) => { try { - // CL-6670: wait OUTSIDE the per-hop deadline below. An agent - // that already has a turn running must never be handed a - // second occurrence while the first is still generating — the - // sidecar's `agent.event` stream carries only the agent's - // address, so `chat-orchestrator.ts`'s reply path cannot tell - // two simultaneously-`running` turns for the same agent apart - // (see `AgentTurnStore.findRunningTurn`'s own doc comment) and - // one of the two replies would land stamped onto the wrong - // turn, or as the wrong turn's own drop notice. Waiting here - // serializes the SAME agent's turns into arrival order, in - // this recipient's own concurrent branch only — a different - // agent named in the same batch (`recipients.map` above) is a + // CL-6670: wait under its own deadline, separate from the + // per-hop deadline below. An agent that already has a turn + // running must never be handed a second occurrence while the + // first is still generating — the sidecar's `agent.event` + // stream carries only the agent's address, so + // `chat-orchestrator.ts`'s reply path cannot tell two + // simultaneously-`running` turns for the same agent apart (see + // `AgentTurnStore.findRunningTurn`'s own doc comment) and one + // of the two replies would land stamped onto the wrong turn, or + // as the wrong turn's own drop notice. Waiting here serializes + // the SAME agent's turns into arrival order, in this + // recipient's own concurrent branch only — a different agent + // named in the same batch (`recipients.map` above) is a // different key and proceeds immediately, unaffected. - await deps.agentTurns?.waitUntilFree({ - tenantId, - workbenchId, - agentAddress, - }); + // + // CL-7129: this wait is bounded — `waitUntilFreeTimeoutMs`, kept + // below the turn-claim TTL alongside `turnDispatchTimeoutMs` + // below — because it runs while `turn-queue.ts` still holds + // this workbench's claim; left unbounded, a slow prior turn + // could hold that claim past its TTL and let a second `run()` + // start a second, concurrent drain of the same queue. + if (deps.agentTurns !== undefined) { + await withTimeout( + deps.agentTurns.waitUntilFree({ + tenantId, + workbenchId, + agentAddress, + }), + waitUntilFreeTimeoutMs, + waitUntilFreeTimeoutMessage(agentAddress, waitUntilFreeTimeoutMs), + ); + } // CL-6644: one deadline around the whole turn, not another // per-hop bound. `dispatchTurn` only ever reaches "the mail was // handed to the agent's mailbox" (see `./turn-queue.ts`'s own diff --git a/packages/chat/test/agent-turn-dispatch.test.ts b/packages/chat/test/agent-turn-dispatch.test.ts index bf0c3d2e6..7e5b6195c 100644 --- a/packages/chat/test/agent-turn-dispatch.test.ts +++ b/packages/chat/test/agent-turn-dispatch.test.ts @@ -15,6 +15,7 @@ import { sendText, settleFanout, TENANT, + timelineOf, } from "./test-support"; async function roomWithAgent( @@ -245,6 +246,60 @@ describe("overlapping turns across agents and messages (CL-6670)", () => { (deps.platform as ReturnType).sentMail, ).toHaveLength(2); }); + + // CL-7129: the wait this queues behind has to tolerate a prior turn + // that is merely slow, not hung — up to the full window a legitimate + // turn is allowed to run. A bound narrower than that reintroduces + // exactly the drop CL-6670 fixed: the queued send times out and lands + // as an undelivered notice instead of dispatching once the agent + // frees up. Modeled at test scale: `waitUntilFreeTimeoutMs` stands in + // for `CHAT_TURN_TIMEOUT_MS`, and the first turn is held open until + // just under that injected budget before closing. + test("a prior turn that runs almost the full wait budget still lets the queued send through, not dropped as undelivered", async () => { + const { app, deps, workbenchId, agentTurns } = await roomWithAgent({ + platform: fakePlatform({ invitable: [{ id: "wfd_echo", name: "echo" }] }), + waitUntilFreeTimeoutMs: 200, + }); + + await sendText(app, workbenchId, "one"); + const [firstTurn] = await agentTurns.listTurns({ + tenantId: TENANT.id, + workbenchId, + }); + + const secondSend = sendText(app, workbenchId, "two"); + + // Closes the first turn just under the injected wait budget, well + // past where a narrower (e.g. halved) bound would already have + // timed the wait out. + await new Promise((resolve) => setTimeout(resolve, 150)); + await agentTurns.finishTurn({ + tenantId: TENANT.id, + turnId: firstTurn?.id ?? "", + status: "completed", + replyMessageId: "msg_reply1", + }); + + await secondSend; + await settleFanout(); + + const timeline = await timelineOf(deps, workbenchId); + const undelivered = timeline.find((message) => + message.parts.some( + (part) => part.kind === "text" && part.turnFailed === true, + ), + ); + expect(undelivered).toBeUndefined(); + + const turns = await agentTurns.listTurns({ + tenantId: TENANT.id, + workbenchId, + }); + expect(turns.map((turn) => turn.childRunId).sort()).toEqual([ + "turn__0", + "turn__1", + ]); + }); }); describe("the turns routes", () => {