From adbf87be0deb33b099be74dbdf3c0a5595c7b0b2 Mon Sep 17 00:00:00 2001 From: Devin AI Date: Fri, 21 Aug 2026 03:33:02 +0000 Subject: [PATCH 1/2] fix(sdk): reset skipToTurnComplete when a new chat turn starts --- .changeset/fluffy-pans-argue.md | 5 +++++ packages/trigger-sdk/src/v3/chat.ts | 8 ++++++++ 2 files changed, 13 insertions(+) create mode 100644 .changeset/fluffy-pans-argue.md diff --git a/.changeset/fluffy-pans-argue.md b/.changeset/fluffy-pans-argue.md new file mode 100644 index 00000000000..5bbe81a0fa7 --- /dev/null +++ b/.changeset/fluffy-pans-argue.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/sdk": patch +--- + +Fix chat transport discarding the next turn after stopping generation. `skipToTurnComplete` is now reset when a new message or action is sent, so a message sent after `stopGeneration` streams normally instead of leaving the chat stuck in a streaming state. diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index a7c7125575e..55889bd58e5 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -870,6 +870,10 @@ export class TriggerChatTransport implements ChatTransport { this.activeStreams.delete(chatId); } + // A stop that never saw its TURN_COMPLETE leaves the flag set, and the new + // turn would be skipped record by record. + state.skipToTurnComplete = false; + state.isStreaming = true; this.notifySessionChange(chatId, state); @@ -1281,6 +1285,10 @@ export class TriggerChatTransport implements ChatTransport { this.activeStreams.delete(chatId); } + // A stop that never saw its TURN_COMPLETE leaves the flag set, and the new + // turn would be skipped record by record. + state.skipToTurnComplete = false; + // Mark streaming + persist so a reload mid-action resumes (reconnectToStream // no-ops when the persisted session says isStreaming: false). state.isStreaming = true; From 9902a08dc5b9e50e524cadfd1620b53252a2b391 Mon Sep 17 00:00:00 2001 From: wei-wei Date: Fri, 21 Aug 2026 19:19:07 +0000 Subject: [PATCH 2/2] test: cover new turn after a stop without turn-complete Co-authored-by: Wei-Wei Wu --- .../test/chat-transport-events.test.ts | 64 +++++++++++++++++++ 1 file changed, 64 insertions(+) diff --git a/packages/trigger-sdk/test/chat-transport-events.test.ts b/packages/trigger-sdk/test/chat-transport-events.test.ts index 39f4e53d722..a35e78ad519 100644 --- a/packages/trigger-sdk/test/chat-transport-events.test.ts +++ b/packages/trigger-sdk/test/chat-transport-events.test.ts @@ -174,6 +174,70 @@ describe("transport send events", () => { }); }); +describe("stopped turn followed by a new turn", () => { + /** + * `.out` stub that honours the `Last-Event-ID` cursor like the server does, so + * a resubscribe cannot replay records the reader already consumed. A stop that + * never saw its turn-complete is therefore unrecoverable unless the new send + * clears the skip state. + */ + function cursoredOneTurnTransport() { + const frames = [ + { id: "1", data: `{"type":"text-delta","id":"t1","delta":"hello"}` }, + { id: "2", data: `{"type":"trigger:turn-complete"}` }, + ]; + + return makeTransport({ + fetch: async (_url, init, ctx) => { + if (ctx.endpoint === "in") return jsonOk(); + + const cursor = new Headers(init.headers).get("Last-Event-ID"); + const from = cursor ? frames.findIndex((f) => f.id === cursor) + 1 : 0; + const remaining = frames.slice(from); + const response = sseResponse( + remaining.map((f) => `id: ${f.id}\ndata: ${f.data}\n\n`).join("") + ); + // Nothing left to send: the session is settled, so the reader stops + // instead of resubscribing. + if (remaining.length === 0) response.headers.set("X-Session-Settled", "true"); + return response; + }, + }); + } + + it("streams a sendMessages turn after a stop that never saw turn-complete", async () => { + const { transport, events } = cursoredOneTurnTransport(); + + expect(await transport.stopGeneration("c1")).toBe(true); + events.length = 0; + + const stream = await transport.sendMessages({ + trigger: "submit-message", + chatId: "c1", + messageId: undefined, + messages: [user("after stop", "u-2")], + abortSignal: undefined, + }); + const chunks = await readAll(stream); + + expect(chunks).toEqual([{ type: "text-delta", id: "t1", delta: "hello" }]); + expect(events.some((e) => e.type === "turn-completed")).toBe(true); + }); + + it("streams a sendAction turn after a stop that never saw turn-complete", async () => { + const { transport, events } = cursoredOneTurnTransport(); + + expect(await transport.stopGeneration("c1")).toBe(true); + events.length = 0; + + const stream = await transport.sendAction("c1", { type: "undo" }); + const chunks = await readAll(stream); + + expect(chunks).toEqual([{ type: "text-delta", id: "t1", delta: "hello" }]); + expect(events.some((e) => e.type === "turn-completed")).toBe(true); + }); +}); + describe("transport stream events", () => { it("marks reconnectToStream subscriptions as resumed", async () => { const { transport, events } = makeTransport({