Skip to content

Commit 4c100ce

Browse files
committed
fix(sdk): keep idle Stop from discarding the next response
1 parent 75ec579 commit 4c100ce

4 files changed

Lines changed: 307 additions & 71 deletions

File tree

‎.changeset/chat-stop-successor-boundary.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
---
44

55
Keep new chat responses intact after Stop, including slow Stop acknowledgments and page reloads.
6+
Stop on an idle hydrated session no longer discards the next response. Early resumed Stop retains its protection.
7+
68
Sequence-free replies after Stop require a transcript reload before further messages.
79
Loading a fresh transcript through `useLoadTranscript` restores blocked sessions only after its saved input cursor covers the stopped turn.
810
Transcript recovery reports missing cursor evidence and empty output polls. An empty recovery poll keeps the accepted message available for reconnect.

‎packages/trigger-sdk/src/v3/chat-stop.test.ts‎

Lines changed: 196 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -78,12 +78,14 @@ describe("Stop with a successor response", () => {
7878
let outputHeaders: { peek: boolean; timeout: number }[];
7979
let inputSeq: number;
8080
let holdStop: boolean;
81+
let holdMessages: boolean;
8182
let stopStatus: number;
8283
let settled: boolean;
8384
let resumeAfterStoppedCheckpoint: boolean;
8485
let emptyRecoveredOutput: boolean;
8586
let includeSequence: boolean;
8687
let pendingStop: { response: ServerResponse; seq: number } | undefined;
88+
let pendingMessages: { response: ServerResponse; seq: number }[];
8789
let saved: ChatSessionPersistedState | null;
8890

8991
function createTransport(
@@ -113,12 +115,14 @@ describe("Stop with a successor response", () => {
113115
outputHeaders = [];
114116
inputSeq = 10;
115117
holdStop = false;
118+
holdMessages = false;
116119
stopStatus = 200;
117120
settled = false;
118121
resumeAfterStoppedCheckpoint = false;
119122
emptyRecoveredOutput = false;
120123
includeSequence = true;
121124
pendingStop = undefined;
125+
pendingMessages = [];
122126
saved = null;
123127
server = createServer(async (request, response) => {
124128
if (request.method === "POST") {
@@ -130,6 +134,8 @@ describe("Stop with a successor response", () => {
130134
const seq = inputSeq++;
131135
if (isStop && holdStop) {
132136
pendingStop = { response, seq };
137+
} else if (!isStop && holdMessages) {
138+
pendingMessages.push({ response, seq });
133139
} else {
134140
appendResponse(response, seq, isStop ? stopStatus : 200);
135141
}
@@ -288,6 +294,162 @@ describe("Stop with a successor response", () => {
288294
}
289295
);
290296

297+
it.each([
298+
["constructor", undefined],
299+
["constructor", "1"],
300+
["setSession", undefined],
301+
["setSession", "1"],
302+
] as const)("does not gate idle %s hydration with cursor %s", async (hydrate, lastEventId) => {
303+
const session = { publicAccessToken: "test-token", lastEventId };
304+
if (hydrate === "constructor") {
305+
transport.dispose();
306+
transport = createTransport(session);
307+
} else {
308+
transport.setSession("chat", session);
309+
}
310+
expect(await transport.stopGeneration("chat")).toBe(true);
311+
expect(transport.getSession("chat")?.skipToTurnComplete).not.toBe(true);
312+
const next = await send();
313+
emit([...reply(2), complete(7, 11)]);
314+
await expect(readText(next)).resolves.toBe("New response");
315+
expect(inputSeq).toBe(12);
316+
});
317+
318+
it.each([
319+
["message", "idle"],
320+
["action", "idle"],
321+
["message", "repeated Stop"],
322+
["action", "repeated Stop"],
323+
] as const)("retains Stop during a pending %s append (%s)", async (kind, state) => {
324+
holdMessages = true;
325+
const first = kind === "message" ? send() : transport.sendAction("chat", { type: "undo" });
326+
await vi.waitFor(() => expect(pendingMessages).toHaveLength(1));
327+
expect(await transport.stopGeneration("chat")).toBe(true);
328+
if (state === "repeated Stop") expect(await transport.stopGeneration("chat")).toBe(true);
329+
expect(transport.getSession("chat")).toMatchObject({
330+
skipToTurnComplete: true,
331+
transcriptRecoveryInputSeq: 11,
332+
});
333+
expect(transport.getSession("chat")).not.toHaveProperty("pendingInputCount");
334+
holdMessages = false;
335+
const pending = pendingMessages[0]!;
336+
appendResponse(pending.response, pending.seq);
337+
const firstStream = await first;
338+
await vi.waitFor(() => expect(outputs).toHaveLength(1));
339+
const firstResult = readText(firstStream);
340+
const nextInput = inputSeq;
341+
const next = await send();
342+
await expect(firstResult).resolves.toBe("");
343+
emit(oldTailAndReply(10, nextInput));
344+
await expect(readText(next)).resolves.toBe("New response");
345+
});
346+
347+
it.each(["message", "action"] as const)(
348+
"retains Stop from the first %s acknowledgment event",
349+
async (kind) => {
350+
let stopping: Promise<boolean> | undefined;
351+
let stopFirst = true;
352+
transport.setOnEvent((event) => {
353+
if (event.type === "message-sent" && event.source !== "stop" && stopFirst) {
354+
stopFirst = false;
355+
queueMicrotask(() => {
356+
stopping = transport.stopGeneration("chat");
357+
});
358+
}
359+
});
360+
const first = await (kind === "message"
361+
? send()
362+
: transport.sendAction("chat", { type: "undo" }));
363+
await vi.waitFor(() => expect(outputs).toHaveLength(1));
364+
await expect(stopping).resolves.toBe(true);
365+
expect(transport.getSession("chat")?.skipToTurnComplete).toBe(true);
366+
const stopped = readText(first);
367+
const next = await send();
368+
await expect(stopped).resolves.toBe("");
369+
emit(oldTailAndReply(10, 12));
370+
await expect(readText(next)).resolves.toBe("New response");
371+
}
372+
);
373+
374+
it.each(["message", "action"] as const)(
375+
"clears pending activity after a failed %s append",
376+
async (kind) => {
377+
holdMessages = true;
378+
const first = kind === "message" ? send() : transport.sendAction("chat", { type: "undo" });
379+
const failed = expect(first).rejects.toThrow();
380+
await vi.waitFor(() => expect(pendingMessages).toHaveLength(1));
381+
const pending = pendingMessages[0]!;
382+
appendResponse(pending.response, pending.seq, 400);
383+
await failed;
384+
holdMessages = false;
385+
await transport.stopGeneration("chat");
386+
expect(transport.getSession("chat")?.skipToTurnComplete).not.toBe(true);
387+
const next = await send();
388+
emit([...reply(2), complete(7, 12)]);
389+
await expect(readText(next)).resolves.toBe("New response");
390+
}
391+
);
392+
393+
it.each([
394+
["message", "idle"],
395+
["action", "idle"],
396+
["message", "abandoned"],
397+
["action", "abandoned"],
398+
] as const)("keeps a rejected %s append idle after Stop (%s)", async (kind, state) => {
399+
transport.setSession("chat", { publicAccessToken: "test-token", isStreaming: false });
400+
if (state === "abandoned") transport.clearSupersedeGate("chat");
401+
holdMessages = true;
402+
const first = kind === "message" ? send() : transport.sendAction("chat", { type: "undo" });
403+
const failed = expect(first).rejects.toThrow();
404+
await vi.waitFor(() => expect(pendingMessages).toHaveLength(1));
405+
await transport.stopGeneration("chat");
406+
const pending = pendingMessages[0]!;
407+
appendResponse(pending.response, pending.seq, 400);
408+
await failed;
409+
holdMessages = false;
410+
expect(transport.getSession("chat")?.skipToTurnComplete).not.toBe(true);
411+
const next = await send();
412+
emit([...reply(2), complete(7, 12)]);
413+
await expect(readText(next)).resolves.toBe("New response");
414+
});
415+
416+
it("retains pending activity until every overlapping append finishes", async () => {
417+
holdMessages = true;
418+
const first = transport.sendAction("chat", { type: "first" });
419+
const firstFailure = expect(first).rejects.toThrow();
420+
const second = transport.sendAction("chat", { type: "second" });
421+
await vi.waitFor(() => expect(pendingMessages).toHaveLength(2));
422+
const rejected = pendingMessages[0]!;
423+
appendResponse(rejected.response, rejected.seq, 400);
424+
await firstFailure;
425+
await transport.stopGeneration("chat");
426+
expect(transport.getSession("chat")?.skipToTurnComplete).toBe(true);
427+
holdMessages = false;
428+
const accepted = pendingMessages[1]!;
429+
appendResponse(accepted.response, accepted.seq);
430+
const stopped = readText(await second);
431+
await vi.waitFor(() => expect(outputs).toHaveLength(1));
432+
const next = await send();
433+
await expect(stopped).resolves.toBe("");
434+
emit(oldTailAndReply(11, 13));
435+
await expect(readText(next)).resolves.toBe("New response");
436+
});
437+
438+
it("does not mark an already-canceled reconnect as an outstanding turn", async () => {
439+
const abort = new AbortController();
440+
abort.abort();
441+
const resumed = await transport.reconnectToStream({
442+
chatId: "chat",
443+
abortSignal: abort.signal,
444+
});
445+
if (!resumed) throw new Error("Expected a resumed stream");
446+
await expect(readText(resumed)).resolves.toBe("");
447+
expect(await transport.stopGeneration("chat")).toBe(true);
448+
const next = await send();
449+
emit([...reply(2), complete(7, 11)]);
450+
await expect(readText(next)).resolves.toBe("New response");
451+
});
452+
291453
it("does not gate a response after an empty settled resume", async () => {
292454
settled = true;
293455
transport.setSession("chat", { publicAccessToken: "test-token", lastEventId: "1" });
@@ -303,6 +465,25 @@ describe("Stop with a successor response", () => {
303465
await expect(readText(next)).resolves.toBe("New response");
304466
});
305467

468+
it("does not transfer an unknown resumed turn to a replacement idle session", async () => {
469+
const abort = new AbortController();
470+
const resumed = await transport.reconnectToStream({
471+
chatId: "chat",
472+
abortSignal: abort.signal,
473+
});
474+
if (!resumed) throw new Error("Expected a resumed stream");
475+
await vi.waitFor(() => expect(outputs).toHaveLength(1));
476+
abort.abort();
477+
await expect(readText(resumed)).resolves.toBe("");
478+
const session = transport.getSession("chat")!;
479+
expect(session).not.toHaveProperty("resumedUnknownTurn");
480+
transport.setSession("chat", session);
481+
await transport.stopGeneration("chat");
482+
const next = await send();
483+
emit([...reply(2), complete(7, 11)]);
484+
await expect(readText(next)).resolves.toBe("New response");
485+
});
486+
306487
it("does not gate a response after Stop on a known idle watch", async () => {
307488
transport.dispose();
308489
transport = createTransport(
@@ -560,6 +741,8 @@ describe("Stop with a successor response", () => {
560741
it.each([false, true])(
561742
"retains the first Stop input through repeated Stop (delayed first acknowledgment: %s)",
562743
async (delayed) => {
744+
await transport.reconnectToStream({ chatId: "chat" });
745+
await vi.waitFor(() => expect(outputs).toHaveLength(1));
563746
holdStop = delayed;
564747
const firstStop = transport.stopGeneration("chat");
565748
if (delayed) await vi.waitFor(() => expect(pendingStop).toBeDefined());
@@ -770,6 +953,19 @@ describe("Stop with a successor response", () => {
770953
expect(transport.getSession("chat")).toMatchObject({ requiresTranscriptReload: true });
771954
});
772955

956+
it("retains the stopped boundary for recovered output before reconnect", async () => {
957+
await hydrateBlockedSession("constructor");
958+
expect(
959+
transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "11", lastInEventId: "11" })
960+
).toBe(true);
961+
await transport.stopGeneration("chat");
962+
expect(transport.getSession("chat")?.skipToTurnComplete).toBe(true);
963+
const next = await send();
964+
emit([complete(12, 12), ...reply(13), complete(18, 14)]);
965+
await expect(readText(next)).resolves.toBe("New response");
966+
expect(inputSeq).toBe(15);
967+
});
968+
773969
it.each(["abort", "stop"] as const)("closes recovery quietly after %s", async (operation) => {
774970
await hydrateBlockedSession("constructor");
775971
expect(

‎packages/trigger-sdk/src/v3/chat.test.ts‎

Lines changed: 13 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1435,18 +1435,21 @@ describe("TriggerChatTransport", () => {
14351435
expect(await drainChunks(next)).toEqual(sampleChunks);
14361436
});
14371437

1438-
it("does not gate a stop with no turn outstanding", async () => {
1439-
mockFetch([() => defaultSseResponse()]);
1440-
1441-
const transport = await armedGate("chat-idle-stop", {
1442-
publicAccessToken: "p",
1443-
isStreaming: false,
1444-
});
1438+
it.each([undefined, false])(
1439+
"does not gate an idle stop with isStreaming %s",
1440+
async (isStreaming) => {
1441+
mockFetch([() => defaultSseResponse()]);
1442+
1443+
const transport = await armedGate("chat-idle-stop", {
1444+
publicAccessToken: "p",
1445+
isStreaming,
1446+
});
14451447

1446-
const stream = await send(transport, "chat-idle-stop");
1448+
const stream = await send(transport, "chat-idle-stop");
14471449

1448-
expect(await drainChunks(stream)).toEqual(sampleChunks);
1449-
});
1450+
expect(await drainChunks(stream)).toEqual(sampleChunks);
1451+
}
1452+
);
14501453

14511454
it("does not arm or write a stop when the abort lands after the boundary", async () => {
14521455
const bodies: string[] = [];

0 commit comments

Comments
 (0)