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
6 changes: 3 additions & 3 deletions apps/hub/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ import {

import {
AGENT_SECTION_MODE,
CHAT_TURN_TIMEOUT_MS,
DEFAULT_TURN_CLAIM_TTL_MS,
createArtifactDeliveryHandler,
createDrizzleAgentTurnStore,
createInMemoryTurnClaimStore,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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),
Expand Down
5 changes: 4 additions & 1 deletion packages/chat/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
21 changes: 20 additions & 1 deletion packages/chat/src/routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -146,14 +146,27 @@ 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
* of the same name in `./workbench-service.ts`. Omitted,
* `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
Expand Down Expand Up @@ -2162,6 +2175,9 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono<TenantEnv> {
...(deps.turnDispatchTimeoutMs !== undefined
? { turnDispatchTimeoutMs: deps.turnDispatchTimeoutMs }
: {}),
...(deps.waitUntilFreeTimeoutMs !== undefined
? { waitUntilFreeTimeoutMs: deps.waitUntilFreeTimeoutMs }
: {}),
},
{
tenantId: ownerTenantId,
Expand Down Expand Up @@ -2308,6 +2324,9 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono<TenantEnv> {
...(deps.turnDispatchTimeoutMs !== undefined
? { turnDispatchTimeoutMs: deps.turnDispatchTimeoutMs }
: {}),
...(deps.waitUntilFreeTimeoutMs !== undefined
? { waitUntilFreeTimeoutMs: deps.waitUntilFreeTimeoutMs }
: {}),
},
{
tenantId: ownerTenantId,
Expand Down
83 changes: 73 additions & 10 deletions packages/chat/src/turn-claims.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand All @@ -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);
});
});
});
88 changes: 67 additions & 21 deletions packages/chat/src/turn-claims.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<boolean>;
tryClaim(claim: TurnClaim): Promise<TurnClaimToken | false>;
/**
* 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<boolean>;
/**
* 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<void>;
holds(claim: TurnClaim, token: TurnClaimToken): Promise<boolean>;
};

/**
* 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;

Expand All @@ -66,18 +92,38 @@ export function createInMemoryTurnClaimStore(options: {
readonly now?: () => number;
}): TurnClaimStore {
const now = options.now ?? Date.now;
const claimedAt = new Map<string, number>();
const holders = new Map<string, { token: string; claimedAt: number }>();
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)
);
},
};
}
67 changes: 67 additions & 0 deletions packages/chat/src/turn-queue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((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 }),
Expand Down
Loading
Loading