Skip to content

Commit ec0a413

Browse files
Merge pull request #474 from corbitsdev/cl-7195-stop-the-turn-queue-from-silently-stranding-a-user-message
Close the turn queue's enqueue/release race and dispatch strand
2 parents baf38a7 + 9a96831 commit ec0a413

2 files changed

Lines changed: 371 additions & 51 deletions

File tree

‎packages/chat/src/turn-queue.test.ts‎

Lines changed: 191 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { describe, expect, test } from "bun:test";
22
import { createInMemoryTurnClaimStore } from "./turn-claims";
3+
import type { TurnClaimStore } from "./turn-claims";
34
import { createWorkbenchTurnQueue, type QueuedTurn } from "./turn-queue";
45
import type { ChatWorkbenchEvent } from "./platform-port";
56

@@ -105,17 +106,18 @@ describe("createWorkbenchTurnQueue", () => {
105106
await firstRun;
106107
});
107108

108-
test("a failed dispatch still releases the claim", async () => {
109+
test("a failed dispatch still releases the claim and does not reject run()", async () => {
109110
const queue = createWorkbenchTurnQueue({
110111
claims: createInMemoryTurnClaimStore({ ttlMs: 60_000 }),
111112
publish: () => undefined,
112113
});
113114

114-
await expect(
115-
queue.run("wb_1", turn("msg_1", "one"), async () => {
116-
throw new Error("dispatch blew up");
117-
}),
118-
).rejects.toThrow("dispatch blew up");
115+
// `dispatch` must never reject per its own documented contract; a
116+
// dispatch that breaks it is this queue's failure to contain, not
117+
// something that should propagate out of `run()`.
118+
await queue.run("wb_1", turn("msg_1", "one"), async () => {
119+
throw new Error("dispatch blew up");
120+
});
119121

120122
// The claim was released despite the throw — a fresh turn can run.
121123
const dispatched: (readonly QueuedTurn[])[] = [];
@@ -125,6 +127,34 @@ describe("createWorkbenchTurnQueue", () => {
125127
expect(dispatched).toEqual([[turn("msg_2", "two")]]);
126128
});
127129

130+
test("a rejecting dispatch does not strand what queued behind it", async () => {
131+
const queue = createWorkbenchTurnQueue({
132+
claims: createInMemoryTurnClaimStore({ ttlMs: 60_000 }),
133+
publish: () => undefined,
134+
});
135+
136+
let rejectFirst: ((err: Error) => void) | undefined;
137+
const dispatched: string[] = [];
138+
const dispatch = (batch: readonly QueuedTurn[]): Promise<void> => {
139+
if (batch[0]?.messageId === "msg_1") {
140+
return new Promise((_resolve, reject) => {
141+
rejectFirst = reject;
142+
});
143+
}
144+
for (const t of batch) dispatched.push(t.messageId);
145+
return Promise.resolve();
146+
};
147+
148+
const firstRun = queue.run("wb_1", turn("msg_1", "one"), dispatch);
149+
// msg_2 queues behind msg_1's still-open (and later rejecting) dispatch.
150+
await queue.run("wb_1", turn("msg_2", "two"), dispatch);
151+
152+
rejectFirst?.(new Error("dispatch blew up"));
153+
await firstRun;
154+
155+
expect(dispatched).toEqual(["msg_2"]);
156+
});
157+
128158
// CL-7129: the TTL is a crash/hang backstop for a dispatch that never
129159
// settles at all, not a license for two loops to drain the same
130160
// workbench's queue at once. Reproduces the reported shape: a message
@@ -211,3 +241,158 @@ describe("createWorkbenchTurnQueue", () => {
211241
await firstRun;
212242
});
213243
});
244+
245+
describe("turn queue drains everything that was enqueued", () => {
246+
test("does not strand a turn enqueued while the holder is releasing", async () => {
247+
const holders = new Map<string, string>();
248+
let n = 0;
249+
let onReleaseStart: (() => Promise<void>) | undefined;
250+
251+
const claims: TurnClaimStore = {
252+
async tryClaim(c) {
253+
if (holders.has(c.workbenchId)) return false;
254+
const t = String(++n);
255+
holders.set(c.workbenchId, t);
256+
return t;
257+
},
258+
async release(c, t) {
259+
// A durable claim store awaits I/O here; anything can arrive first.
260+
const hook = onReleaseStart;
261+
onReleaseStart = undefined;
262+
if (hook) await hook();
263+
if (holders.get(c.workbenchId) !== t) return false;
264+
holders.delete(c.workbenchId);
265+
return true;
266+
},
267+
async holds(c, t) {
268+
return holders.get(c.workbenchId) === t;
269+
},
270+
};
271+
272+
const queue = createWorkbenchTurnQueue({
273+
claims,
274+
publish: () => undefined,
275+
});
276+
const dispatched: string[] = [];
277+
const dispatch = async (batch: readonly QueuedTurn[]) => {
278+
for (const t of batch) dispatched.push(t.messageId);
279+
};
280+
281+
// m2 arrives after the drain loop saw an empty queue, before release lands.
282+
onReleaseStart = async () => {
283+
await queue.run("wb", turn("m2", "two"), dispatch);
284+
};
285+
286+
await queue.run("wb", turn("m1", "one"), dispatch);
287+
await new Promise((r) => setTimeout(r, 20));
288+
289+
expect(dispatched).toEqual(["m1", "m2"]);
290+
});
291+
292+
// The symmetric half of the same race: the holder's `release()` can
293+
// resolve and its own post-release recheck can find nothing pending
294+
// *before* a concurrent enqueuer's failed `tryClaim` call ever
295+
// resolves — the enqueuer's push then lands only after the holder is
296+
// already gone and has already returned. Closing this requires the
297+
// enqueue path to reclaim on its own, not just the release path.
298+
test("does not strand a turn whose push lands only after the holder already returned", async () => {
299+
const holders = new Map<string, string>();
300+
let n = 0;
301+
let releaseBlockedClaim: (() => void) | undefined;
302+
303+
const claims: TurnClaimStore = {
304+
async tryClaim(c) {
305+
if (holders.has(c.workbenchId)) {
306+
// Simulate a claim read that is still in flight (e.g. a slow
307+
// durable store) when the holder's own release already lands.
308+
await new Promise<void>((resolve) => {
309+
releaseBlockedClaim = resolve;
310+
});
311+
return false;
312+
}
313+
const t = String(++n);
314+
holders.set(c.workbenchId, t);
315+
return t;
316+
},
317+
async release(c, t) {
318+
if (holders.get(c.workbenchId) !== t) return false;
319+
holders.delete(c.workbenchId);
320+
return true;
321+
},
322+
async holds(c, t) {
323+
return holders.get(c.workbenchId) === t;
324+
},
325+
};
326+
327+
const queue = createWorkbenchTurnQueue({
328+
claims,
329+
publish: () => undefined,
330+
});
331+
const dispatched: string[] = [];
332+
const dispatch = async (batch: readonly QueuedTurn[]) => {
333+
for (const t of batch) dispatched.push(t.messageId);
334+
};
335+
336+
const firstRun = queue.run("wb", turn("m1", "one"), dispatch);
337+
const secondRun = queue.run("wb", turn("m2", "two"), dispatch);
338+
339+
// m1 fully dispatches, sees nothing pending (m2's push hasn't
340+
// landed yet — its tryClaim call is still parked), and releases.
341+
await firstRun;
342+
expect(dispatched).toEqual(["m1"]);
343+
344+
// Only now does m2's tryClaim resolve `false` and its push execute.
345+
releaseBlockedClaim?.();
346+
await secondRun;
347+
348+
expect(dispatched).toEqual(["m1", "m2"]);
349+
});
350+
351+
// CL-7195's own reasoning ("a durable store reopens the window")
352+
// implies these calls can also reject outright, not just resolve
353+
// `false` — a real possibility for I/O, unlike the in-memory store.
354+
// `run()` documents "never rejects", so the drain loop's claim-store
355+
// calls must not leak that rejection, and must still attempt to free
356+
// the claim rather than leaving it held until the TTL backstop.
357+
test("still releases the claim and does not reject run() when the claim store throws mid-drain", async () => {
358+
const holders = new Map<string, string>();
359+
let n = 0;
360+
let releaseCalls = 0;
361+
362+
const claims: TurnClaimStore = {
363+
async tryClaim(c) {
364+
if (holders.has(c.workbenchId)) return false;
365+
const t = String(++n);
366+
holders.set(c.workbenchId, t);
367+
return t;
368+
},
369+
async release(c, t) {
370+
releaseCalls++;
371+
if (holders.get(c.workbenchId) !== t) return false;
372+
holders.delete(c.workbenchId);
373+
return true;
374+
},
375+
async holds() {
376+
throw new Error("claim store connection reset");
377+
},
378+
};
379+
380+
const queue = createWorkbenchTurnQueue({
381+
claims,
382+
publish: () => undefined,
383+
});
384+
const dispatched: string[] = [];
385+
const dispatch = async (batch: readonly QueuedTurn[]) => {
386+
for (const t of batch) dispatched.push(t.messageId);
387+
};
388+
389+
// Does not reject, despite `holds` throwing.
390+
await queue.run("wb", turn("m1", "one"), dispatch);
391+
392+
expect(dispatched).toEqual(["m1"]);
393+
expect(releaseCalls).toBeGreaterThan(0);
394+
// The claim was freed despite the throw — a fresh turn can claim it
395+
// directly rather than only being queued.
396+
expect(holders.has("wb")).toBe(false);
397+
});
398+
});

0 commit comments

Comments
 (0)