Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
127 changes: 127 additions & 0 deletions packages/chat/src/agent-turns.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import { describe, expect, test } from "bun:test";
import {
AGENT_TURN_STALE_MS,
createInMemoryAgentTurnStore,
createTurnFreedSignal,
} from "./agent-turns";

const BASE = {
Expand Down Expand Up @@ -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);
});
});
111 changes: 89 additions & 22 deletions packages/chat/src/agent-turns.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,15 @@ export interface AgentTurnStore {
* turn is on the projection from its first moment.
*/
startTurn(input: StartAgentTurnInput): Promise<AgentTurn>;
/** 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<AgentTurn | undefined>;
/**
* The newest turn still `running` for (workbench, agent) — what the
Expand Down Expand Up @@ -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<void>;
waitUntilFree(
input: {
readonly tenantId: string;
readonly workbenchId: string;
readonly agentAddress: string;
},
signal?: AbortSignal,
): Promise<void>;
listTurns(input: {
readonly tenantId: string;
readonly workbenchId: string;
Expand Down Expand Up @@ -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<void>;
wait(key: string, backstopMs: number, signal?: AbortSignal): Promise<void>;
/** 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<string, Set<() => void>>();
type Waiter = {
readonly settle: () => void;
readonly timer: ReturnType<typeof setTimeout>;
};
// 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<string, Map<number, Waiter>>();
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<number, Waiter>();
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;
},
};
}

Expand Down Expand Up @@ -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,
Expand All @@ -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);
}
},

Expand Down Expand Up @@ -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({
Expand All @@ -472,6 +537,7 @@ export function createDrizzleAgentTurnStore<
and(
eq(agentTurns.id, input.turnId),
eq(agentTurns.tenantId, input.tenantId),
eq(agentTurns.status, "running"),
),
)
.returning();
Expand All @@ -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);
}
},

Expand Down
11 changes: 9 additions & 2 deletions packages/chat/src/platform-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -697,8 +697,11 @@ export function createHubChatPlatform(
*/
async function wakeByAddressBounded(address: string): Promise<void> {
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`,
);
Expand Down Expand Up @@ -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`,
);
Expand Down
Loading
Loading