Skip to content

Commit 93c63c6

Browse files
ericallamTrigger.dev RepoOps
authored andcommitted
fix(sdk): stop chat.agent replay blocking the event loop on long turns
Recovering a long unfinished `chat.agent` response on a continuation run could block the worker's event loop for minutes, long enough to miss heartbeats and fail the run. Capturing the response at the end of a long turn had the same cost. Both paths rebuilt the message by draining `readUIMessageStream` and keeping its last snapshot. That reducer emits a `structuredClone` of the whole message after almost every chunk, and every chunk was enqueued into one web stream up front, where each dequeue is linear in the queue length on Node. Both costs are quadratic in turn length. The reductions now go through one helper: ```ts const message = await reduceUIMessageChunks(chunks, { message: original }); ``` - It runs the same AI SDK reducer through `createUIMessageStream`'s `onFinish`, which mutates one state and hands back only the final message. That message is cloned once so it shares nothing with the recorded chunks. - Chunks are fed one at a time, each after the previous one comes out, and the loop yields to the event loop every 10ms, so heartbeats keep running. - The result matches the last `readUIMessageStream` snapshot, including `undefined` when no chunk would emit one. A malformed chunk is identified exactly, and the chunks before it are reduced again on the same path, so content before it is kept without a slow fallback. A 320,000-chunk text response now reduces in under a second. Tests cover equivalence with `readUIMessageStream` across finished, unfinished, resumed and continued messages, a malformed chunk after a long prefix, detachment from the input chunks, and bounded timer gaps during large reductions. Mono-RevId: 4ce17e79a4750a7c949bbc3448881ff0a98dc705
1 parent 7607372 commit 93c63c6

8 files changed

Lines changed: 373 additions & 62 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@trigger.dev/sdk": patch
3+
---
4+
5+
Recovering a long unfinished `chat.agent` response on a continuation run no longer blocks the worker. Rebuilding the response from its streamed chunks now takes time linear in its length and yields to the event loop as it goes, so heartbeats keep firing and the run is not killed mid-replay. Capturing the response at the end of a long turn gets the same speedup.

‎packages/trigger-sdk/src/imports/ai-runtime-cjs.cts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,8 @@ const ai = require("ai");
77
// @ts-ignore
88
module.exports.convertToModelMessages = ai.convertToModelMessages;
99
// @ts-ignore
10+
module.exports.createUIMessageStream = ai.createUIMessageStream;
11+
// @ts-ignore
1012
module.exports.dynamicTool = ai.dynamicTool;
1113
// @ts-ignore
1214
module.exports.generateId = ai.generateId;

‎packages/trigger-sdk/src/imports/ai-runtime.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@
1313
// @ts-ignore
1414
import {
1515
convertToModelMessages,
16+
createUIMessageStream,
1617
dynamicTool,
1718
generateId,
1819
getToolName,
@@ -28,6 +29,7 @@ import {
2829
// @ts-ignore
2930
export {
3031
convertToModelMessages,
32+
createUIMessageStream,
3133
dynamicTool,
3234
generateId,
3335
getToolName,

‎packages/trigger-sdk/src/v3/ai.ts‎

Lines changed: 13 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -75,7 +75,6 @@ import {
7575
getToolName,
7676
isToolUIPart,
7777
jsonSchema,
78-
readUIMessageStream,
7978
streamText as aiStreamText,
8079
zodSchema,
8180
} from "../imports/ai-runtime.js";
@@ -97,6 +96,7 @@ import {
9796
import { responseAfterCompaction } from "./compactionResponse.js";
9897
import { withToolResultsInCallOrder } from "./toolResultOrder.js";
9998
import { ManagedChatResponse, createOrderedChatWriter } from "./managedChatResponse.js";
99+
import { reduceUIMessageChunks } from "./uiMessageChunks.js";
100100
import {
101101
convertSteeredMessages,
102102
retainStepMessages,
@@ -439,9 +439,9 @@ export function __setReplaySessionOutTailImplForTests(
439439
* 2. Filter out the agent's control chunks (`type: "trigger:*"`) — they
440440
* ride on the same stream as the user-visible UIMessageChunks.
441441
* 3. Split chunks at `start`/`finish` boundaries so each segment is a
442-
* single message, then feed each segment through the AI SDK's
443-
* `readUIMessageStream` reducer (the same one `useChat` uses on the
444-
* browser side) and grab the final emitted snapshot.
442+
* single message, then reduce each segment with the AI SDK's reducer
443+
* (the same one `useChat` uses on the browser side) via
444+
* {@link reduceUIMessageChunks}, which builds only the final state.
445445
* 4. The trailing message — if it never received a `finish` chunk —
446446
* goes through `cleanupAbortedParts` so partial in-flight parts
447447
* don't leak into the next turn's accumulator. Drop it entirely
@@ -497,8 +497,8 @@ async function replaySessionOutTail<TUIMessage extends UIMessage>(
497497
}
498498
if (!current) {
499499
// Chunk arrived before any `start`. Synthesize a segment so the reducer
500-
// has something to work with — `readUIMessageStream` tolerates a missing
501-
// `start` because we pass `message: undefined`.
500+
// has something to work with. The reducer tolerates a missing `start`
501+
// and leaves the message ID empty.
502502
current = { chunks: [], closed: false };
503503
segments.push(current);
504504
}
@@ -515,17 +515,9 @@ async function replaySessionOutTail<TUIMessage extends UIMessage>(
515515
for (let i = 0; i < segments.length; i++) {
516516
const seg = segments[i]!;
517517
const isTrailing = i === segments.length - 1 && !seg.closed;
518-
const segmentStream = new ReadableStream<UIMessageChunk>({
519-
start(controller) {
520-
for (const c of seg.chunks) controller.enqueue(c);
521-
controller.close();
522-
},
523-
});
524518
let last: UIMessage | undefined;
525519
try {
526-
for await (const snapshot of readUIMessageStream({ stream: segmentStream })) {
527-
last = snapshot;
528-
}
520+
last = await reduceUIMessageChunks(seg.chunks);
529521
} catch (error) {
530522
// Reducer error — the segment is malformed. Skip it and keep going so a
531523
// single corrupt chunk doesn't sink the entire replay.
@@ -555,15 +547,15 @@ async function replaySessionOutTail<TUIMessage extends UIMessage>(
555547

556548
/**
557549
* Test-only entry point that bypasses `__setReplaySessionOutTailImplForTests`
558-
* and reaches the real `apiClient.subscribeToSessionStream` + chunk-segment
559-
* splitter + `readUIMessageStream` reducer. Pairs with the snapshot
550+
* and reaches the real `apiClient.readSessionStreamRecords` + chunk-segment
551+
* splitter + {@link reduceUIMessageChunks} reducer. Pairs with the snapshot
560552
* production-path wrappers above. Lets `replay-session-out.test.ts` drive
561553
* synthetic chunk sequences through the real reducer to lock down chunk-
562554
* stream → `UIMessage[]` correctness — if the AI SDK's chunk semantics
563555
* shift in a future version, the test catches it before customers do.
564556
*
565-
* Tests should mock `apiClient.subscribeToSessionStream` (e.g. via
566-
* `vi.spyOn(apiClient, ...)`) to feed a `ReadableStream<UIMessageChunk>`.
557+
* Tests should stub `apiClient.readSessionStreamRecords` to return the
558+
* recorded chunks as session stream records.
567559
*
568560
* Not part of the public API.
569561
* @internal
@@ -11689,7 +11681,7 @@ function tapUIMessageChunks(
1168911681
* Reconstruct a partial assistant `UIMessage` from the raw chunks that
1169011682
* streamed before a failure — the fallback for {@link pipeChatAndCapture}
1169111683
* when a transport error abandons the stream before `onFinish` runs. Uses the
11692-
* same `readUIMessageStream` reducer as the boot-time replay path. Returns
11684+
* same {@link reduceUIMessageChunks} reducer as the boot-time replay path. Returns
1169311685
* `undefined` if there's nothing to assemble or the reducer throws.
1169411686
*/
1169511687
async function assemblePartialFromChunks(chunks: UIMessageChunk[]): Promise<UIMessage | undefined> {
@@ -11699,17 +11691,7 @@ async function assemblePartialFromChunks(chunks: UIMessageChunk[]): Promise<UIMe
1169911691
});
1170011692
if (relevant.length === 0) return undefined;
1170111693
try {
11702-
const stream = new ReadableStream<UIMessageChunk>({
11703-
start(controller) {
11704-
for (const chunk of relevant) controller.enqueue(chunk);
11705-
controller.close();
11706-
},
11707-
});
11708-
let last: UIMessage | undefined;
11709-
for await (const message of readUIMessageStream({ stream })) {
11710-
last = message;
11711-
}
11712-
return last;
11694+
return await reduceUIMessageChunks(relevant);
1171311695
} catch {
1171411696
return undefined;
1171511697
}

‎packages/trigger-sdk/src/v3/managedChatResponse.ts‎

Lines changed: 5 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import type { UIMessage, UIMessageChunk } from "ai";
2-
import { generateId, readUIMessageStream } from "../imports/ai-runtime.js";
2+
import { generateId } from "../imports/ai-runtime.js";
3+
import { reduceUIMessageChunks } from "./uiMessageChunks.js";
34

45
/** One ordered output channel for a turn's managed model and data chunks. */
56
export class ManagedChatResponse {
@@ -147,21 +148,10 @@ export class ManagedChatResponse {
147148
async snapshot(options?: { message: UIMessage; from: number }): Promise<UIMessage | undefined> {
148149
if (!this.chunks.length) return options?.message;
149150
const chunks = this.chunks.slice(options?.from ?? 0);
150-
let message: UIMessage | undefined = options?.message;
151151
const original = options?.message ?? this.original;
152-
const stream = new ReadableStream<UIMessageChunk>({
153-
start(controller) {
154-
for (const chunk of chunks) controller.enqueue(chunk);
155-
controller.close();
156-
},
157-
});
158-
for await (const update of readUIMessageStream({
159-
stream,
160-
...(original ? { message: structuredClone(original) } : {}),
161-
// An error chunk must leave previously streamed content recoverable.
162-
terminateOnError: false,
163-
}))
164-
message = update;
152+
const message =
153+
(await reduceUIMessageChunks(chunks, original ? { message: original } : undefined)) ??
154+
options?.message;
165155
return message ? { ...message, id: message.id || this.fallbackId } : undefined;
166156
}
167157

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

Lines changed: 4 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import {
1616
type ChatSnapshotV1,
1717
} from "../ai.js";
1818
import { createTestSessionHandle, type TestSessionOutState } from "./test-session-handle.js";
19+
import { reduceUIMessageChunks } from "../uiMessageChunks.js";
1920

2021
/** Pre-seed locals before the agent's `run()` starts. */
2122
type SetupLocals = (locals: { set<T>(key: LocalsKey<T>, value: T): void }) => void | Promise<void>;
@@ -898,22 +899,17 @@ export function mockChatAgent(
898899
/**
899900
* Reduce a synthetic UIMessageChunk[] sequence into the UIMessage[] that
900901
* the runtime's `replaySessionOutTail` would produce. Splits chunks at
901-
* `start` boundaries and feeds each segment through AI SDK's
902-
* `readUIMessageStream`. The trailing un-finished segment goes through
902+
* `start` boundaries and reduces each segment with
903+
* `reduceUIMessageChunks`. The trailing un-finished segment goes through
903904
* `cleanupAbortedParts`. Mirrors the production reducer used in
904905
* `ai.ts:replaySessionOutTail`.
905906
*/
906907
async function reduceChunksToMessages(chunks: UIMessageChunk[]): Promise<UIMessage[]> {
907908
if (chunks.length === 0) return [];
908909
const aiModule = (await import("ai")) as {
909-
readUIMessageStream?: (args: {
910-
stream: ReadableStream<UIMessageChunk>;
911-
}) => AsyncIterable<UIMessage>;
912910
cleanupAbortedParts?: (msg: UIMessage) => UIMessage;
913911
};
914-
const readUIMessageStream = aiModule.readUIMessageStream;
915912
const cleanupAbortedParts = aiModule.cleanupAbortedParts;
916-
if (!readUIMessageStream) return [];
917913

918914
type Segment = { chunks: UIMessageChunk[]; closed: boolean };
919915
const segments: Segment[] = [];
@@ -939,17 +935,9 @@ async function reduceChunksToMessages(chunks: UIMessageChunk[]): Promise<UIMessa
939935
for (let i = 0; i < segments.length; i++) {
940936
const seg = segments[i]!;
941937
const isTrailing = i === segments.length - 1 && !seg.closed;
942-
const segmentStream = new ReadableStream<UIMessageChunk>({
943-
start(controller) {
944-
for (const c of seg.chunks) controller.enqueue(c);
945-
controller.close();
946-
},
947-
});
948938
let last: UIMessage | undefined;
949939
try {
950-
for await (const snapshot of readUIMessageStream({ stream: segmentStream })) {
951-
last = snapshot;
952-
}
940+
last = await reduceUIMessageChunks(seg.chunks);
953941
} catch {
954942
// Skip malformed segment — tests can assert by inspecting what makes it through.
955943
continue;

0 commit comments

Comments
 (0)