From 6648e50891ca25aa5c86ded86fb73725d9fb88b6 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:46:56 -0700 Subject: [PATCH 1/2] Add tests for cancellable dispatch timeouts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit withTimeout races a timer against a promise with no way to signal the loser that it lost. dispatchTurnBatch wraps waitUntilFree and dispatchTurn in it, so a timed-out call keeps running unbounded — for waitUntilFree's uncancellable for(;;) loop, and for dispatchTurn's turn row, which can stay `running` long enough for a later reply to land behind the undelivered notice already posted for it. Underneath, waitUntilFree's wait primitive registers a resolve into a Set and schedules a backstop timer notify() never clears; when the backstop fires first, the resolve is never removed either, so a key that keeps timing out accumulates waiters forever. These tests assert the fix ahead of the implementation: withTimeout must pass work an AbortSignal it can react to; the waiter Set must never grow across repeated timeouts; finishTurn must be compare-and-set on status so a late reply can never clobber a turn a timeout already closed; and a dispatch deadline must close its turn row the instant it fires, not whenever the abandoned send eventually settles. --- packages/chat/src/agent-turns.test.ts | 127 ++++++++++++++++++ packages/chat/src/with-timeout.test.ts | 75 +++++++++++ .../chat/test/turn-dispatch-deadline.test.ts | 37 +++++ 3 files changed, 239 insertions(+) create mode 100644 packages/chat/src/with-timeout.test.ts diff --git a/packages/chat/src/agent-turns.test.ts b/packages/chat/src/agent-turns.test.ts index efab69370..34044fd2c 100644 --- a/packages/chat/src/agent-turns.test.ts +++ b/packages/chat/src/agent-turns.test.ts @@ -3,6 +3,7 @@ import { describe, expect, test } from "bun:test"; import { AGENT_TURN_STALE_MS, createInMemoryAgentTurnStore, + createTurnFreedSignal, } from "./agent-turns"; const BASE = { @@ -223,5 +224,131 @@ describe("createInMemoryAgentTurnStore", () => { await waiting; expect(freed).toBe(true); }); + + // CL-7193: a key that keeps timing out (nothing ever calls `notify`) + // used to leave one abandoned resolve in the waiter Set per attempt + // forever. An `AbortSignal` now lets a timed-out wait remove itself + // immediately instead of lingering until its own backstop fires. + test("an aborted wait throws instead of reporting the agent free", async () => { + const store = createInMemoryAgentTurnStore(); + await store.startTurn(BASE); + const controller = new AbortController(); + + const waiting = store.waitUntilFree(BASE, controller.signal); + controller.abort(new Error("deadline")); + + await expect(waiting).rejects.toThrow("deadline"); + }); + + test("an already-aborted signal throws immediately, without polling", async () => { + const store = createInMemoryAgentTurnStore(); + await store.startTurn(BASE); + const controller = new AbortController(); + controller.abort(new Error("already gone")); + + await expect( + store.waitUntilFree(BASE, controller.signal), + ).rejects.toThrow("already gone"); + }); + }); + + // CL-7193: `finishTurn` used to overwrite a turn's row unconditionally. + // A dispatch deadline closing a turn as `failed` and the turn's own + // late reply closing it as `completed` could race — whichever landed + // last silently clobbered the other, discarding its `replyMessageId`. + // `finishTurn` is now compare-and-set on `status === "running"`, so + // exactly one of two racing closes ever applies. + describe("finishTurn is compare-and-set on status (CL-7193)", () => { + test("a second finishTurn call on an already-finished turn is a no-op", async () => { + const store = createInMemoryAgentTurnStore(); + const opened = await store.startTurn(BASE); + + const first = await store.finishTurn({ + tenantId: BASE.tenantId, + turnId: opened.id, + status: "failed", + error: "turn dispatch timed out", + }); + expect(first?.status).toBe("failed"); + + const second = await store.finishTurn({ + tenantId: BASE.tenantId, + turnId: opened.id, + status: "completed", + replyMessageId: "msg_late_reply", + }); + expect(second).toBeUndefined(); + + const stored = await store.getTurn({ + tenantId: BASE.tenantId, + turnId: opened.id, + }); + expect(stored?.status).toBe("failed"); + expect(stored?.replyMessageId).toBeNull(); + }); + }); +}); + +// CL-7193: the process-local pubsub `waitUntilFree` is built on. Exercised +// directly because the defect it fixes -- a backstop firing before +// `notify` leaves its waiter in the Set forever -- is about the waiter +// bookkeeping itself, not any one store's behavior. +describe("createTurnFreedSignal (CL-7193)", () => { + test("a backstop that fires before notify removes its own waiter from the Set", async () => { + const freed = createTurnFreedSignal(); + + await freed.wait("key", 1); + + expect(freed.waiterCount("key")).toBe(0); + }); + + test("repeated timed-out waits never grow the waiter Set", async () => { + const freed = createTurnFreedSignal(); + + for (let i = 0; i < 25; i++) { + await freed.wait("key", 1); + } + + expect(freed.waiterCount("key")).toBe(0); + }); + + test("notify clears the backstop timer so it never fires again after resolving", async () => { + const freed = createTurnFreedSignal(); + const waiting = freed.wait("key", 10_000); + + freed.notify("key"); + await waiting; + + expect(freed.waiterCount("key")).toBe(0); + }); + + test("an aborted wait removes itself from the Set immediately", async () => { + const freed = createTurnFreedSignal(); + const controller = new AbortController(); + + const waiting = freed.wait("key", 10_000, controller.signal); + expect(freed.waiterCount("key")).toBe(1); + + controller.abort(); + await waiting; + + expect(freed.waiterCount("key")).toBe(0); + }); + + test("a different key's waiters are never touched by another key's notify", async () => { + const freed = createTurnFreedSignal(); + let otherFreed = false; + const waiting = freed.wait("other", 10_000).then(() => { + otherFreed = true; + }); + + freed.notify("key"); + await Promise.resolve(); + expect(otherFreed).toBe(false); + expect(freed.waiterCount("other")).toBe(1); + + freed.notify("other"); + await waiting; + expect(otherFreed).toBe(true); }); }); diff --git a/packages/chat/src/with-timeout.test.ts b/packages/chat/src/with-timeout.test.ts new file mode 100644 index 000000000..d2fd66017 --- /dev/null +++ b/packages/chat/src/with-timeout.test.ts @@ -0,0 +1,75 @@ +import { describe, expect, test } from "bun:test"; + +import { withTimeout } from "./with-timeout"; + +describe("withTimeout", () => { + test("resolves with the value once work settles before the deadline", async () => { + const result = await withTimeout(async () => "ok", 50, "timed out"); + expect(result).toBe("ok"); + }); + + test("propagates a rejection work produces on its own", async () => { + await expect( + withTimeout( + async () => { + throw new Error("boom"); + }, + 50, + "timed out", + ), + ).rejects.toThrow("boom"); + }); + + test("rejects with the timeout message when work never settles", async () => { + await expect( + withTimeout( + () => new Promise(() => {}), + 10, + "did not settle within 10ms", + ), + ).rejects.toThrow("did not settle within 10ms"); + }); + + // CL-7193: the losing promise used to be abandoned outright, with + // nothing telling it the caller had already given up. `work` now + // receives the same signal the timeout fires the moment it wins. + test("aborts the signal work received the moment the deadline wins", async () => { + let aborted = false; + let reason: unknown; + + await expect( + withTimeout( + (signal) => { + signal.addEventListener("abort", () => { + aborted = true; + reason = signal.reason; + }); + return new Promise(() => {}); + }, + 10, + "did not settle within 10ms", + ), + ).rejects.toThrow("did not settle within 10ms"); + + expect(aborted).toBe(true); + expect(reason).toBeInstanceOf(Error); + }); + + test("never aborts the signal when work settles before the deadline", async () => { + let sawAbort = false; + + await withTimeout( + (signal) => { + signal.addEventListener("abort", () => { + sawAbort = true; + }); + return Promise.resolve("done"); + }, + 50, + "timed out", + ); + + await Bun.sleep(60); + expect(sawAbort).toBe(false); + }); +}); diff --git a/packages/chat/test/turn-dispatch-deadline.test.ts b/packages/chat/test/turn-dispatch-deadline.test.ts index 290f2af05..769f33572 100644 --- a/packages/chat/test/turn-dispatch-deadline.test.ts +++ b/packages/chat/test/turn-dispatch-deadline.test.ts @@ -9,6 +9,7 @@ // bounds stay as diagnostics; this is the backstop. import { describe, expect, test } from "bun:test"; +import { createInMemoryAgentTurnStore } from "../src/agent-turns"; import { createChatRoutes } from "../src/routes"; import { buildDeps, @@ -16,6 +17,7 @@ import { fakePlatform, mountAs, settleFanout, + TENANT, timelineOf, } from "./test-support"; @@ -168,6 +170,41 @@ describe("dispatchTurnBatch's turn-level deadline (CL-6644)", () => { ); }); + // CL-7193: the timed-out `dispatchTurn` call used to be abandoned along + // with `sendMail` -- its turn row stayed `running` until the agent's + // real reply (if `sendMail` ever landed) closed it later, well behind + // the undelivered notice already on the timeline. The deadline now + // closes the row itself the instant it fires, before `sendMail` ever + // settles. + test("a timed-out dispatch closes its turn row immediately, not whenever sendMail eventually settles", async () => { + const platform = fakePlatform({ + invitable: [{ id: "wfd_echo", name: "echo" }], + }); + platform.sendMail = () => new Promise(() => {}); + const agentTurns = createInMemoryAgentTurnStore(); + + const deps = buildDeps({ platform, agentTurns, turnDispatchTimeoutMs: 20 }); + const app = mountAs(createChatRoutes(deps), "prn_alice"); + const { body: workbench } = await createWorkbench(app, { + kind: "chat", + definitionId: "wfd_echo", + }); + + await app.request(`/workbenches/${workbench.id}/messages`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ parts: [{ kind: "text", text: "hello" }] }), + }); + await settleFanout(); + + const turns = await agentTurns.listTurns({ + tenantId: TENANT.id, + workbenchId: workbench.id, + }); + expect(turns).toHaveLength(1); + expect(turns[0]?.status).toBe("failed"); + }); + test("the timeout message names the turn's run address and its elapsed budget", async () => { const { turnDispatchTimeoutMessage } = await import("../src/workbench-service"); From d9e705990f0ffa613cd57eac572dac0ebb93483e Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:58:37 -0700 Subject: [PATCH 2/2] Give a timed-out dispatch a real cancellation path withTimeout now threads an AbortSignal into the work it wraps, firing the moment its own timer wins the race, instead of rejecting and walking away from a promise that keeps running unobserved. waitUntilFree's for(;;) loop accepts that signal and throws once it aborts, rather than polling until its own backstop timer -- up to AGENT_TURN_STALE_MS later -- or silently resolving as though the agent were free. The waiter pubsub underneath (createTurnFreedSignal) gives every waiter its own teardown: notify(), a fired backstop, and an aborted signal all funnel through one settle() that clears the backstop timer and removes the waiter, so whichever reaches it first is the only one that ever runs -- a key that keeps timing out no longer accumulates one abandoned waiter per attempt. dispatchTurn takes the same signal and closes its turn row `failed` the instant it aborts, rather than leaving the row `running` until (or unless) the abandoned sendMail eventually settles -- sendMail itself has no cancellable primitive, so this closes the bookkeeping around the send rather than the send itself. finishTurn is now compare-and-set on status === "running", so this abort-driven close and a late reply's own close can race safely: exactly one applies, and the loser can never clobber a completed turn's replyMessageId. The two withTimeout call sites in platform-adapter.ts (wakeByAddress, sendFoldedMail) have no cancellable primitive to hook a signal into, so they keep today's abandon-on-timeout behavior -- only their call shape changes to match the new signature. With a turn closed on timeout, a late reply with no correlation hint across multiple room memberships now has no running turn to attach to; chat-orchestrator.ts's existing single-membership fallback still posts it untied, but a multi-membership agent's late reply in that situation is dropped rather than landing behind the undelivered notice -- trading a confusing double-post for a rarer silent drop. --- packages/chat/src/agent-turns.ts | 111 ++++++++++++++++++++----- packages/chat/src/platform-adapter.ts | 11 ++- packages/chat/src/with-timeout.ts | 26 ++++-- packages/chat/src/workbench-service.ts | 98 +++++++++++++++++----- 4 files changed, 192 insertions(+), 54 deletions(-) diff --git a/packages/chat/src/agent-turns.ts b/packages/chat/src/agent-turns.ts index ba547aeee..a5080384b 100644 --- a/packages/chat/src/agent-turns.ts +++ b/packages/chat/src/agent-turns.ts @@ -68,7 +68,15 @@ export interface AgentTurnStore { * turn is on the projection from its first moment. */ startTurn(input: StartAgentTurnInput): Promise; - /** Closes a turn. Undefined for a turn id this store does not hold. */ + /** + * Closes a turn — compare-and-set on `status === "running"` (CL-7193), + * so two closes racing the same turn (a dispatch deadline closing it + * `failed` the instant it fires, the turn's own late reply closing it + * `completed` once it lands) can never clobber each other: exactly one + * applies, and the loser reads back `undefined` rather than + * overwriting the winner's `replyMessageId`. Also undefined for a + * turn id this store does not hold. + */ finishTurn(input: FinishAgentTurnInput): Promise; /** * The newest turn still `running` for (workbench, agent) — what the @@ -105,12 +113,20 @@ export interface AgentTurnStore { * dispatch: this is scoped to one (workbench, agent) pair, so two * different agents mentioned in the same message still turn * concurrently. + * + * CL-7193: an aborted `signal` makes the wait stop polling and throw + * — never resolve as though the agent were free — so a caller whose + * own deadline already fired never mistakes an abandoned wait for a + * green light to dispatch. */ - waitUntilFree(input: { - readonly tenantId: string; - readonly workbenchId: string; - readonly agentAddress: string; - }): Promise; + waitUntilFree( + input: { + readonly tenantId: string; + readonly workbenchId: string; + readonly agentAddress: string; + }, + signal?: AbortSignal, + ): Promise; listTurns(input: { readonly tenantId: string; readonly workbenchId: string; @@ -168,27 +184,66 @@ function sectionKey(input: { * in-memory one by construction, the Drizzle one because this hub runs * as a single replica — see `AGENT_TURN_STALE_MS`'s own doc comment for * the ceiling this backstop mirrors). + * + * CL-7193: each waiter used to be a bare `resolve` — `notify` deleted + * the whole per-key Set and called every one, but never cleared their + * individual backstop timers, and a backstop firing first (the key kept + * timing out, `notify` never came) never removed its own entry either. + * A key that kept timing out accumulated one abandoned waiter per + * attempt forever. Every waiter now owns its full teardown — clearing + * its own timer, removing itself from the Set, detaching its abort + * listener — so whichever of {`notify`, the backstop, an aborted + * `signal`} reaches it first is the only one that ever runs it. */ -function createTurnFreedSignal(): { +export function createTurnFreedSignal(): { notify(key: string): void; - wait(key: string, backstopMs: number): Promise; + wait(key: string, backstopMs: number, signal?: AbortSignal): Promise; + /** Waiters still pending for `key` — diagnostic, and what proves the + * Set never grows unboundedly across repeated timeouts. */ + waiterCount(key: string): number; } { - const waiters = new Map void>>(); + type Waiter = { + readonly settle: () => void; + readonly timer: ReturnType; + }; + // A waiter's own settle() needs to remove that same waiter from its + // Set -- an id minted before the waiter object exists gives it a way + // to name itself without the object having to reference its own + // not-yet-created binding. + const waiters = new Map>(); + let nextWaiterId = 0; + + function forget(key: string, id: number): void { + const byId = waiters.get(key); + if (byId === undefined) return; + byId.delete(id); + if (byId.size === 0) waiters.delete(key); + } + return { notify(key) { - const set = waiters.get(key); - if (set === undefined) return; - waiters.delete(key); - for (const resolve of set) resolve(); + const byId = waiters.get(key); + if (byId === undefined) return; + for (const waiter of [...byId.values()]) waiter.settle(); }, - wait(key, backstopMs) { + wait(key, backstopMs, signal) { return new Promise((resolve) => { - const set = waiters.get(key) ?? new Set(); - set.add(resolve); - waiters.set(key, set); - setTimeout(resolve, backstopMs); + const byId = waiters.get(key) ?? new Map(); + waiters.set(key, byId); + const id = nextWaiterId++; + const settle = () => { + clearTimeout(byId.get(id)?.timer); + forget(key, id); + signal?.removeEventListener("abort", settle); + resolve(); + }; + byId.set(id, { settle, timer: setTimeout(settle, backstopMs) }); + signal?.addEventListener("abort", settle, { once: true }); }); }, + waiterCount(key) { + return waiters.get(key)?.size ?? 0; + }, }; } @@ -263,6 +318,12 @@ export function createInMemoryAgentTurnStore( if (existing === undefined || existing.tenantId !== input.tenantId) { return undefined; } + // CL-7193: compare-and-set on status -- a turn already closed + // (by a dispatch deadline's abort close, or an earlier finishTurn) + // never gets clobbered by a second close racing in behind it. + if (existing.status !== "running") { + return undefined; + } const finished: AgentTurn = { id: existing.id, tenantId: existing.tenantId, @@ -287,13 +348,14 @@ export function createInMemoryAgentTurnStore( return runningTurn(input); }, - async waitUntilFree(input) { + async waitUntilFree(input, signal) { for (;;) { + if (signal?.aborted) throw signal.reason; const running = runningTurn(input); if (running === undefined) return; const staleAt = Date.parse(running.startedAt) + AGENT_TURN_STALE_MS; const backstopMs = Math.max(0, staleAt - now()) + 50; - await freed.wait(sectionKey(input), backstopMs); + await freed.wait(sectionKey(input), backstopMs, signal); } }, @@ -455,6 +517,9 @@ export function createDrizzleAgentTurnStore< }, async finishTurn(input) { + // CL-7193: compare-and-set on status -- a turn already closed + // (by a dispatch deadline's abort close, or an earlier finishTurn) + // never gets clobbered by a second close racing in behind it. const [row] = await db .update(agentTurns) .set({ @@ -472,6 +537,7 @@ export function createDrizzleAgentTurnStore< and( eq(agentTurns.id, input.turnId), eq(agentTurns.tenantId, input.tenantId), + eq(agentTurns.status, "running"), ), ) .returning(); @@ -485,13 +551,14 @@ export function createDrizzleAgentTurnStore< return resolveRunningTurn(input); }, - async waitUntilFree(input) { + async waitUntilFree(input, signal) { for (;;) { + if (signal?.aborted) throw signal.reason; const running = await resolveRunningTurn(input); if (running === undefined) return; const staleAt = Date.parse(running.startedAt) + AGENT_TURN_STALE_MS; const backstopMs = Math.max(0, staleAt - now()) + 50; - await freed.wait(sectionKey(input), backstopMs); + await freed.wait(sectionKey(input), backstopMs, signal); } }, diff --git a/packages/chat/src/platform-adapter.ts b/packages/chat/src/platform-adapter.ts index b09d305fa..5c47fbbf4 100644 --- a/packages/chat/src/platform-adapter.ts +++ b/packages/chat/src/platform-adapter.ts @@ -697,8 +697,11 @@ export function createHubChatPlatform( */ async function wakeByAddressBounded(address: string): Promise { if (lifecycle === undefined) { + // CL-7193: `wakeByAddress` has no cancellable primitive to hook a + // signal into, so a timeout here still abandons the underlying wake + // exactly as before — the signal parameter is unused on purpose. await withTimeout( - wakeByAddress(address), + () => wakeByAddress(address), DEFAULT_WAKE_TIMEOUT_MS, `wake for "${address}" did not settle within ${String(DEFAULT_WAKE_TIMEOUT_MS)}ms`, ); @@ -877,8 +880,12 @@ export function createHubChatPlatform( let loggedRetryStart = false; for (let attempt = 0; ; attempt++) { try { + // CL-7193: `sendFoldedMail` does a DB write plus a sidecar + // delivery that shouldn't be half-cancelled, and has no signal + // to accept regardless — the signal parameter is unused here, + // same abandon-on-timeout behavior as before. return await withTimeout( - sendFoldedMail(foldedRunsDeps, params), + () => sendFoldedMail(foldedRunsDeps, params), MAIL_DELIVERY_TIMEOUT_MS, `mail to ${params.agentAddress} did not settle within ${String(MAIL_DELIVERY_TIMEOUT_MS)}ms`, ); diff --git a/packages/chat/src/with-timeout.ts b/packages/chat/src/with-timeout.ts index 1a3f4b4fd..a162b3469 100644 --- a/packages/chat/src/with-timeout.ts +++ b/packages/chat/src/with-timeout.ts @@ -1,17 +1,29 @@ -// A wall-clock bound on a promise that may never settle at all — the -// shape #316 first wrote inline in `platform-adapter.ts` to catch a -// wedged post-deploy mail ack, and CL-6644's turn-level deadline in -// `workbench-service.ts` reuses rather than duplicates. `promise` +// A wall-clock bound on work that may never settle at all — the shape +// #316 first wrote inline in `platform-adapter.ts` to catch a wedged +// post-deploy mail ack, and CL-6644's turn-level deadline in +// `workbench-service.ts` reuses rather than duplicates. `work` // rejecting or resolving on its own always wins the race; the timer // only fires when neither ever happens. +// +// CL-7193: the losing side used to be abandoned outright — nothing told +// it the caller had already given up, so it kept running uncancelled. +// `work` now receives the `AbortSignal` this fires the moment the +// timeout wins, so a caller with something cancellable to do (stop a +// poll loop, close a row it opened) can react instead of running to +// completion unobserved. A `work` that ignores the signal keeps +// today's abandon-on-timeout behavior. export function withTimeout( - promise: Promise, + work: (signal: AbortSignal) => Promise, ms: number, message: string, ): Promise { + const controller = new AbortController(); return new Promise((resolve, reject) => { - const timer = setTimeout(() => reject(new Error(message)), ms); - promise.then( + const timer = setTimeout(() => { + controller.abort(new Error(message)); + reject(new Error(message)); + }, ms); + work(controller.signal).then( (value) => { clearTimeout(timer); resolve(value); diff --git a/packages/chat/src/workbench-service.ts b/packages/chat/src/workbench-service.ts index 8dde9bda0..4904eb723 100644 --- a/packages/chat/src/workbench-service.ts +++ b/packages/chat/src/workbench-service.ts @@ -1714,13 +1714,14 @@ async function dispatchTurnBatch( // 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) { + const agentTurns = deps.agentTurns; + if (agentTurns !== undefined) { await withTimeout( - deps.agentTurns.waitUntilFree({ - tenantId, - workbenchId, - agentAddress, - }), + (signal) => + agentTurns.waitUntilFree( + { tenantId, workbenchId, agentAddress }, + signal, + ), waitUntilFreeTimeoutMs, waitUntilFreeTimeoutMessage(agentAddress, waitUntilFreeTimeoutMs), ); @@ -1733,24 +1734,34 @@ async function dispatchTurnBatch( // `chat-orchestrator.ts`'s independent sidecar-event // subscription. This deadline therefore can never cut off a // reply in progress: nothing it awaits is the reply. + // + // CL-7193: `dispatchTurn` gets the deadline's own `AbortSignal` + // so it can close its turn row the instant the deadline fires, + // instead of leaving it `running` until (or unless) the + // abandoned `sendMail` below eventually settles. await withTimeout( - dispatchTurn(deps, { - tenantId, - workbenchId, - principalId: last.principalId, - agentAddress, - parts, - requestMessageIds: messageIds, - // A batch concatenating more than one queued message's parts - // still stamps the whole combined body as the answer when any - // one of them carries a correlationId — acceptable (the user - // did answer), but a batch mixing the actual answer with an - // unrelated follow-up hands the gate the whole blob, not just - // the answer. - ...(batchCorrelationId !== undefined - ? { correlationId: batchCorrelationId } - : {}), - }), + (signal) => + dispatchTurn( + deps, + { + tenantId, + workbenchId, + principalId: last.principalId, + agentAddress, + parts, + requestMessageIds: messageIds, + // A batch concatenating more than one queued message's parts + // still stamps the whole combined body as the answer when any + // one of them carries a correlationId — acceptable (the user + // did answer), but a batch mixing the actual answer with an + // unrelated follow-up hands the gate the whole blob, not just + // the answer. + ...(batchCorrelationId !== undefined + ? { correlationId: batchCorrelationId } + : {}), + }, + signal, + ), turnDispatchTimeoutMs, turnDispatchTimeoutMessage(agentAddress, turnDispatchTimeoutMs), ); @@ -1817,10 +1828,21 @@ export type DispatchTurnInput = { * will carry is already allocated. A trigger that never lands closes the * row `failed` and rethrows, leaving the caller to post the undelivered * notice it always has. + * + * CL-7193: `sendMail` itself has no cancellable primitive — a caller's + * `signal` firing can't stop the send in flight, only close the + * bookkeeping around it. The moment it aborts, this closes the turn row + * `failed` immediately rather than waiting for `sendMail` to eventually + * settle, so a real reply that later lands (behind the undelivered + * notice the caller already posted for this timeout) finds no `running` + * row to attach to — `finishTurn`'s compare-and-set means whichever of + * the abort close and the late reply's own close reaches the row first + * is the only one that applies. */ export async function dispatchTurn( deps: Pick, input: DispatchTurnInput, + signal?: AbortSignal, ): Promise { const turn = await deps.agentTurns?.startTurn({ tenantId: input.tenantId, @@ -1828,6 +1850,34 @@ export async function dispatchTurn( agentAddress: input.agentAddress, requestMessageIds: input.requestMessageIds, }); + + const closeAsTimedOut = () => { + if (turn === undefined) return; + deps.agentTurns + ?.finishTurn({ + tenantId: input.tenantId, + turnId: turn.id, + status: "failed", + error: + signal?.reason instanceof Error + ? signal.reason.message + : "turn dispatch timed out", + }) + .catch((err: unknown) => { + reportError(err, { + operation: "chat.dispatchTurn.closeAsTimedOut", + tenantId: input.tenantId, + roomId: input.workbenchId, + agentId: input.agentAddress, + }); + }); + }; + if (signal?.aborted) { + closeAsTimedOut(); + } else { + signal?.addEventListener("abort", closeAsTimedOut, { once: true }); + } + try { await deps.platform.sendMail({ tenantId: input.tenantId, @@ -1849,6 +1899,8 @@ export async function dispatchTurn( }); } throw err; + } finally { + signal?.removeEventListener("abort", closeAsTimedOut); } }