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/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.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/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); } } 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");