Skip to content

Commit ad4e2ef

Browse files
committed
fix(sdk): require input evidence for stopped transcript recovery
1 parent 92b6a91 commit ad4e2ef

5 files changed

Lines changed: 252 additions & 45 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,4 +4,4 @@
44

55
Keep new chat responses intact after Stop, including slow Stop acknowledgments and page reloads.
66
Sequence-free replies after Stop require a transcript reload before further messages.
7-
Loading a fresh transcript through `useLoadTranscript` restores blocked sessions without replaying old output.
7+
Loading a fresh transcript through `useLoadTranscript` restores blocked sessions only after its saved input cursor covers the stopped turn.

packages/trigger-sdk/src/v3/chat-react.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -65,8 +65,8 @@ export type UseLoadTranscriptOptions = {
6565
* the loaded transcript, so the live subscription opens just past the
6666
* persisted history instead of replaying it. Only applies once the
6767
* transport knows the session (from `sessions` or after `start`).
68-
* For a blocked session, a fresh load with a newer saved cursor also clears
69-
* the transcript reload requirement. A stale result leaves sends blocked.
68+
* A blocked session requires a newer output cursor and an input cursor that covers the stopped input.
69+
* A stale result leaves sends blocked.
7070
*/
7171
transport?: TriggerChatTransport;
7272
/** Page size passed to the action. */
@@ -85,11 +85,11 @@ export type UseLoadTranscriptOptions = {
8585
export function seedTranscriptCursor(
8686
transport: Pick<TriggerChatTransport, "seedResumeCursor">,
8787
chatId: string,
88-
cursors: { lastOutEventId?: string } | undefined,
89-
completeRecovery?: (lastEventId: string | undefined) => boolean
88+
cursors: { lastOutEventId?: string; lastInEventId?: string } | undefined,
89+
completeRecovery?: (lastEventId: string | undefined, lastInEventId?: string) => boolean
9090
): boolean {
9191
const lastEventId = cursors?.lastOutEventId;
92-
if (completeRecovery) return completeRecovery(lastEventId);
92+
if (completeRecovery) return completeRecovery(lastEventId, cursors?.lastInEventId);
9393
if (!lastEventId) return false;
9494
transport.seedResumeCursor(chatId, lastEventId);
9595
return true;

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

Lines changed: 105 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -501,6 +501,106 @@ describe("Stop with a successor response", () => {
501501
}
502502
);
503503

504+
it("rejects a newer snapshot from before the stopped input", async () => {
505+
await hydrateBlockedSession("constructor");
506+
const recover = transport.prepareTranscriptRecovery("chat");
507+
if (!recover) throw new Error("Expected transcript recovery");
508+
expect(recover("5", "9")).toBe(false);
509+
await expect(send()).rejects.toThrow("Stopped chat response cannot be matched");
510+
expect(inputSeq).toBe(13);
511+
});
512+
513+
it.each([
514+
["explicit", "constructor"],
515+
["explicit", "setSession"],
516+
["abort", "constructor"],
517+
["abort", "setSession"],
518+
] as const)(
519+
"rejects an older snapshot after cursor-free %s Stop and %s hydration",
520+
async (stopMode, hydrate) => {
521+
const abort = new AbortController();
522+
const resumed = await transport.reconnectToStream({
523+
chatId: "chat",
524+
abortSignal: abort.signal,
525+
stopOnAbort: true,
526+
});
527+
if (!resumed) throw new Error("Expected a resumed stream");
528+
await vi.waitFor(() => expect(outputs).toHaveLength(1));
529+
if (stopMode === "explicit") await transport.stopGeneration("chat");
530+
else abort.abort();
531+
await vi.waitFor(() =>
532+
expect(transport.getSession("chat")?.transcriptRecoveryInputSeq).toBe(10)
533+
);
534+
includeSequence = false;
535+
await expect(send()).rejects.toThrow("Stopped chat response cannot be matched");
536+
const session = transport.getSession("chat");
537+
if (!session) throw new Error("Expected persisted state");
538+
transport.dispose();
539+
transport = createTransport(
540+
hydrate === "constructor" ? session : { publicAccessToken: "test-token" }
541+
);
542+
if (hydrate === "setSession") transport.setSession("chat", session);
543+
const stale = transport.prepareTranscriptRecovery("chat");
544+
if (!stale) throw new Error("Expected transcript recovery");
545+
expect(stale("5", "9")).toBe(false);
546+
await expect(send()).rejects.toThrow("Stopped chat response cannot be matched");
547+
const fresh = transport.prepareTranscriptRecovery("chat");
548+
if (!fresh) throw new Error("Expected transcript recovery");
549+
expect(fresh("11", "10")).toBe(true);
550+
const next = await transport.reconnectToStream({ chatId: "chat" });
551+
if (!next) throw new Error("Expected a resumed stream");
552+
await vi.waitFor(() => expect(outputs).toHaveLength(2));
553+
emit([...reply(12), complete(17, 11)]);
554+
await expect(readText(next)).resolves.toBe("New response");
555+
expect(inputSeq).toBe(12);
556+
}
557+
);
558+
559+
it.each([false, true])(
560+
"retains the first Stop input through repeated Stop (delayed first acknowledgment: %s)",
561+
async (delayed) => {
562+
holdStop = delayed;
563+
const firstStop = transport.stopGeneration("chat");
564+
if (delayed) await vi.waitFor(() => expect(pendingStop).toBeDefined());
565+
else expect(await firstStop).toBe(true);
566+
includeSequence = false;
567+
await expect(send()).rejects.toThrow("Stopped chat response cannot be matched");
568+
includeSequence = true;
569+
holdStop = false;
570+
expect(await transport.stopGeneration("chat")).toBe(true);
571+
if (delayed) {
572+
const pending = pendingStop!;
573+
appendResponse(pending.response, pending.seq);
574+
expect(await firstStop).toBe(true);
575+
}
576+
const recover = transport.prepareTranscriptRecovery("chat");
577+
if (!recover) throw new Error("Expected transcript recovery");
578+
expect(recover("17", "11")).toBe(true);
579+
const next = await send();
580+
emit([...reply(18), complete(23, 13)]);
581+
await expect(readText(next)).resolves.toBe("New response");
582+
}
583+
);
584+
585+
it("does not install an old Stop sequence into a replacement session", async () => {
586+
holdStop = true;
587+
const stopped = transport.stopGeneration("chat");
588+
await vi.waitFor(() => expect(pendingStop).toBeDefined());
589+
transport.setSession("chat", {
590+
publicAccessToken: "replacement-token",
591+
skipToTurnComplete: true,
592+
requiresTranscriptReload: true,
593+
isStreaming: false,
594+
});
595+
const pending = pendingStop!;
596+
appendResponse(pending.response, pending.seq);
597+
expect(await stopped).toBe(true);
598+
const recover = transport.prepareTranscriptRecovery("chat");
599+
if (!recover) throw new Error("Expected transcript recovery");
600+
expect(recover("17", "11")).toBe(false);
601+
await expect(send()).rejects.toThrow("Stopped chat response cannot be matched");
602+
});
603+
504604
it("accepts a sequence-free response without a stopped boundary", async () => {
505605
includeSequence = false;
506606
const stream = await send();
@@ -532,7 +632,7 @@ describe("Stop with a successor response", () => {
532632
await hydrateBlockedSession(hydrate);
533633
const recover = transport.prepareTranscriptRecovery("chat");
534634
if (!recover) throw new Error("Expected transcript recovery");
535-
expect(recover("11")).toBe(true);
635+
expect(recover("11", "11")).toBe(true);
536636
transport.seedResumeCursor("chat", "11");
537637
expect(saved).toMatchObject({
538638
lastEventId: "11",
@@ -558,7 +658,7 @@ describe("Stop with a successor response", () => {
558658
await hydrateBlockedSession("constructor");
559659
const recover = transport.prepareTranscriptRecovery("chat");
560660
if (!recover) throw new Error("Expected transcript recovery");
561-
expect(recover(cursor)).toBe(false);
661+
expect(recover(cursor, "11")).toBe(false);
562662
await expect(send()).rejects.toThrow("Stopped chat response cannot be matched");
563663
await expect(transport.sendAction("chat", { type: "undo" })).rejects.toThrow(
564664
"Stopped chat response cannot be matched"
@@ -585,7 +685,7 @@ describe("Stop with a successor response", () => {
585685
await hydrateBlockedSession("constructor");
586686
const recover = transport.prepareTranscriptRecovery("chat");
587687
if (!recover) throw new Error("Expected transcript recovery");
588-
expect(recover("11")).toBe(true);
688+
expect(recover("11", "11")).toBe(true);
589689
if (hydrate !== "none") {
590690
const session = transport.getSession("chat");
591691
if (!session) throw new Error("Expected persisted state");
@@ -611,7 +711,7 @@ describe("Stop with a successor response", () => {
611711
await hydrateBlockedSession("constructor");
612712
const recover = transport.prepareTranscriptRecovery("chat");
613713
if (!recover) throw new Error("Expected transcript recovery");
614-
expect(recover("11")).toBe(true);
714+
expect(recover("11", "11")).toBe(true);
615715
resumeAfterStoppedCheckpoint = true;
616716
emptyRecoveredOutput = true;
617717
const resumed = await transport.reconnectToStream({ chatId: "chat" });
@@ -642,7 +742,7 @@ describe("Stop with a successor response", () => {
642742
holdStop = true;
643743
const stopped = transport.stopGeneration("chat");
644744
await vi.waitFor(() => expect(pendingStop).toBeDefined());
645-
expect(recover("11")).toBe(false);
745+
expect(recover("11", "11")).toBe(false);
646746
await expect(send()).rejects.toThrow("Stopped chat response cannot be matched");
647747
const pending = pendingStop!;
648748
appendResponse(pending.response, pending.seq);

0 commit comments

Comments
 (0)