diff --git a/.changeset/chat-custom-agent-capture.md b/.changeset/chat-custom-agent-capture.md new file mode 100644 index 00000000000..8729cda8c94 --- /dev/null +++ b/.changeset/chat-custom-agent-capture.md @@ -0,0 +1,22 @@ +--- +"@trigger.dev/sdk": patch +--- + +Custom chat agent loops get two ergonomic wins for owning the turn loop. + +`chat.writeTurnComplete()` now returns the turn boundary's resume cursors (`lastEventId` for the output stream and `sessionInEventId` for the input stream), so you can persist them straight from the task instead of round-tripping them back from the client. + +```ts +const { lastEventId, sessionInEventId } = await chat.writeTurnComplete(); +await db.chats.update(chatId, { lastEventId, sessionInEventId }); +``` + +`chat.pipeAndCapture()` no longer throws when a stream is stopped or fails. It now returns a `PipeAndCaptureResult` whose `message` holds any partial output captured before the stop or failure, alongside a typed `status` (`"complete" | "aborted" | "error"`) and, on failure, the `error`. Read the message off the result: + +```ts +const { message, status, error } = await chat.pipeAndCapture(result, { signal }); +if (message) conversation.addResponse(message); +if (status === "error") logger.error("turn failed", { error }); +``` + +Note: `pipeAndCapture` previously resolved to `UIMessage | undefined`. Update call sites to read `.message` from the returned result. diff --git a/docs/ai-chat/changelog.mdx b/docs/ai-chat/changelog.mdx index 4bd667a4148..22f496b91bf 100644 --- a/docs/ai-chat/changelog.mdx +++ b/docs/ai-chat/changelog.mdx @@ -55,10 +55,10 @@ if (isFinal) { const result = streamText({ model, messages: conversation.modelMessages, tools }); // Pass originalMessages so the handed-over tool round merges into the // step-1 assistant instead of starting a new message. - const response = await chat.pipeAndCapture(result, { + const { message } = await chat.pipeAndCapture(result, { originalMessages: conversation.uiMessages, }); - if (response) await conversation.addResponse(response); + if (message) await conversation.addResponse(message); } ``` diff --git a/docs/ai-chat/compaction.mdx b/docs/ai-chat/compaction.mdx index 02db7d7b13c..0a8ed48c7ee 100644 --- a/docs/ai-chat/compaction.mdx +++ b/docs/ai-chat/compaction.mdx @@ -365,8 +365,8 @@ for (let turn = 0; turn < 100; turn++) { stopWhen: stepCountIs(15), }); - const response = await chat.pipeAndCapture(result); - if (response) await conversation.addResponse(response); + const { message } = await chat.pipeAndCapture(result); + if (message) await conversation.addResponse(message); // Outer-loop compaction const usage = await result.totalUsage; diff --git a/docs/ai-chat/custom-agents.mdx b/docs/ai-chat/custom-agents.mdx index 1cac571ad42..197bff6b5e1 100644 --- a/docs/ai-chat/custom-agents.mdx +++ b/docs/ai-chat/custom-agents.mdx @@ -160,11 +160,11 @@ for await (const turn of session) { }); // Manual: pipe and capture separately - const response = await chat.pipeAndCapture(result, { signal: turn.signal }); + const { message } = await chat.pipeAndCapture(result, { signal: turn.signal }); - if (response) { + if (message) { // Custom processing before accumulating - await turn.addResponse(response); + await turn.addResponse(message); } // Custom persistence, analytics, etc. @@ -215,8 +215,8 @@ For full control, skip `createSession` and compose the primitives directly: | ------------------------------- | -------------------------------------------------------------------------------------------- | | `chat.messages` | Input stream for incoming messages — use `.waitWithIdleTimeout()` to wait for the next turn | | `chat.createStopSignal()` | Create a managed stop signal wired to the stop input stream | -| `chat.pipeAndCapture(result)` | Pipe a `StreamTextResult` to the chat stream and capture the response | -| `chat.writeTurnComplete()` | Signal the frontend that the current turn is complete | +| `chat.pipeAndCapture(result)` | Pipe a stream and capture the response; returns `{ message, status, error }` | +| `chat.writeTurnComplete()` | Signal turn complete; returns `{ lastEventId, sessionInEventId }` resume cursors | | `chat.MessageAccumulator` | Accumulates conversation messages across turns | | `chat.pipe(stream)` | Pipe a stream to the frontend (no response capture) | | `chat.cleanupAbortedParts(msg)` | Clean up incomplete parts from a stopped response | @@ -285,21 +285,16 @@ export const myChat = chat.customAgent({ stopWhen: stepCountIs(15), }); - let response; - try { - response = await chat.pipeAndCapture(result, { signal: combinedSignal }); - } catch (error) { - if (error instanceof Error && error.name === "AbortError") { - if (runSignal.aborted) break; - // Stop — fall through to accumulate partial - } else { - throw error; - } - } + const { message, status, error } = await chat.pipeAndCapture(result, { + signal: combinedSignal, + }); + // pipeAndCapture never throws: a user stop returns status "aborted" with + // the partial message, and a failure returns status "error" with `error`. + if (status === "error") throw error; - if (response) { + if (message) { const cleaned = - stop.signal.aborted && !runSignal.aborted ? chat.cleanupAbortedParts(response) : response; + stop.signal.aborted && !runSignal.aborted ? chat.cleanupAbortedParts(message) : message; await conversation.addResponse(cleaned); } @@ -346,8 +341,8 @@ const messages = await conversation.addIncoming( ); // After piping, add the response -const response = await chat.pipeAndCapture(result); -if (response) await conversation.addResponse(response); +const { message } = await chat.pipeAndCapture(result); +if (message) await conversation.addResponse(message); // Access accumulated messages for persistence conversation.uiMessages; // UIMessage[] diff --git a/docs/ai-chat/fast-starts.mdx b/docs/ai-chat/fast-starts.mdx index 40d7b62fd20..78bb7ed9f13 100644 --- a/docs/ai-chat/fast-starts.mdx +++ b/docs/ai-chat/fast-starts.mdx @@ -589,10 +589,10 @@ if (turn === 0 && payload.trigger === "handover-prepare") { messages: conversation.modelMessages, stopWhen: stepCountIs(10), }); - const response = await chat.pipeAndCapture(result, { + const { message } = await chat.pipeAndCapture(result, { originalMessages: conversation.uiMessages, }); - if (response) await conversation.addResponse(response); + if (message) await conversation.addResponse(message); } await chat.writeTurnComplete(); // on isFinal the warm partial is already the response return; diff --git a/docs/ai-chat/pending-messages.mdx b/docs/ai-chat/pending-messages.mdx index 5aa668099a8..c77600f470a 100644 --- a/docs/ai-chat/pending-messages.mdx +++ b/docs/ai-chat/pending-messages.mdx @@ -179,10 +179,14 @@ for (let turn = 0; turn < 100; turn++) { stopWhen: stepCountIs(15), }); - const response = await chat.pipeAndCapture(result); + const { message, status, error } = await chat.pipeAndCapture(result); sub.off(); - if (response) await conversation.addResponse(response); + // Keep any partial output, then branch on the outcome: pipeAndCapture no + // longer throws, so rethrow a real failure and don't complete a failed turn. + if (message) await conversation.addResponse(message); + if (status === "error") throw error; + await chat.writeTurnComplete(); } ``` diff --git a/docs/ai-chat/reference.mdx b/docs/ai-chat/reference.mdx index 340422050fb..da08f9a0473 100644 --- a/docs/ai-chat/reference.mdx +++ b/docs/ai-chat/reference.mdx @@ -503,8 +503,8 @@ All methods available on the `chat` object from `@trigger.dev/sdk/ai`. | `chat.agent(options)` | Create a chat agent | | `chat.createSession(payload, options)` | Create an async iterator for chat turns | | `chat.pipe(source, options?)` | Pipe a stream to the frontend (from anywhere inside a task) | -| `chat.pipeAndCapture(source, options?)` | Pipe and capture the response `UIMessage` | -| `chat.writeTurnComplete(options?)` | Signal the frontend that the current turn is complete | +| `chat.pipeAndCapture(source, options?)` | Pipe and capture the response; returns `{ message, status, error }` | +| `chat.writeTurnComplete(options?)` | Signal turn complete; returns `{ lastEventId, sessionInEventId }` resume cursors | | `chat.createStopSignal()` | Create a managed stop signal wired to the stop input stream | | `chat.messages` | Input stream for incoming messages — use `.waitWithIdleTimeout()` | | `chat.local({ id })` | Create a per-run typed local (see [`chat.local`](/ai-chat/chat-local)) | diff --git a/packages/trigger-sdk/src/v3/ai.ts b/packages/trigger-sdk/src/v3/ai.ts index ee70c79c755..4612302c4ae 100644 --- a/packages/trigger-sdk/src/v3/ai.ts +++ b/packages/trigger-sdk/src/v3/ai.ts @@ -8828,14 +8828,117 @@ function createStopSignal(): { * The `TriggerChatTransport` intercepts this to close the ReadableStream * for the current turn. Call after piping the response stream. * + * Returns two resume cursors for the turn boundary, both saveable from the + * task instead of round-tripping them back from the client: + * - `lastEventId` — the turn-complete control record's seq_num on + * `session.out`; where the next turn's output stream resumes. + * - `sessionInEventId` — the committed-consume cursor on `session.in` as of + * this turn-complete, letting a raw loop correlate the boundary with the + * exact input record it acknowledged. Trigger owns input-cursor recovery, + * so this is for correlation / out-of-sync detection, not required. + * + * Either is `undefined` when the corresponding cursor isn't available. + * * @example * ```ts * await chat.pipe(result); - * await chat.writeTurnComplete(); + * const { lastEventId, sessionInEventId } = await chat.writeTurnComplete(); + * await db.chats.update(chatId, { lastEventId, sessionInEventId }); * ``` */ -async function chatWriteTurnComplete(options?: { publicAccessToken?: string }): Promise { - await writeTurnCompleteChunk(undefined, options?.publicAccessToken); +async function chatWriteTurnComplete(options?: { + publicAccessToken?: string; +}): Promise<{ lastEventId?: string; sessionInEventId?: string }> { + const result = await writeTurnCompleteChunk(undefined, options?.publicAccessToken); + // Same cursor written to the `session-in-event-id` header inside + // `writeTurnCompleteChunk`; surfaced here so the caller can persist it. + const inCursor = getChatSession().in.lastDispatchedSeqNum(); + return { + lastEventId: result?.lastEventId, + ...(inCursor !== undefined ? { sessionInEventId: String(inCursor) } : {}), + }; +} + +/** + * The outcome of a turn's stream, reported by {@link pipeChatAndCapture}. + * + * - `complete` — the stream finished on its own. + * - `aborted` — the stream was stopped via the `signal` (user stop / cancel). + * - `error` — the stream threw; `error` carries what was thrown. + * + * `message` holds whatever the assistant produced and is present for every + * status — including `aborted` and `error` — as long as any output streamed + * before the stop or failure, so partial responses are never lost. + */ +export type PipeAndCaptureResult = { + /** The captured assistant message, or `undefined` if nothing streamed. */ + message: UIMessage | undefined; + /** Coarse outcome of the stream. */ + status: "complete" | "aborted" | "error"; + /** The AI SDK finish reason, when the stream reported one. */ + finishReason?: FinishReason; + /** What the stream threw, present only when `status === "error"`. */ + error?: unknown; +}; + +/** + * Pass every chunk through untouched while recording it in `buffer`. Handles + * both the `AsyncIterable` and `ReadableStream` shapes `toUIMessageStream()` + * can return, and propagates a source error to the consumer after buffering + * whatever streamed first. See {@link pipeChatAndCapture} for why. + */ +async function* tapUIMessageChunks( + source: AsyncIterable | ReadableStream, + buffer: UIMessageChunk[] +): AsyncGenerator { + if (isReadableStream(source)) { + const reader = source.getReader(); + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + buffer.push(value as UIMessageChunk); + yield value; + } + } finally { + reader.releaseLock(); + } + } else { + for await (const chunk of source) { + buffer.push(chunk as UIMessageChunk); + yield chunk; + } + } +} + +/** + * Reconstruct a partial assistant `UIMessage` from the raw chunks that + * streamed before a failure — the fallback for {@link pipeChatAndCapture} + * when a transport error abandons the stream before `onFinish` runs. Uses the + * same `readUIMessageStream` reducer as the boot-time replay path. Returns + * `undefined` if there's nothing to assemble or the reducer throws. + */ +async function assemblePartialFromChunks(chunks: UIMessageChunk[]): Promise { + const relevant = chunks.filter((c) => { + const type = (c as { type?: unknown }).type; + return typeof type === "string" && !type.startsWith("trigger:"); + }); + if (relevant.length === 0) return undefined; + try { + const stream = new ReadableStream({ + start(controller) { + for (const chunk of relevant) controller.enqueue(chunk); + controller.close(); + }, + }); + let last: UIMessage | undefined; + for await (const message of readUIMessageStream({ stream })) { + last = message; + } + return last; + } catch { + return undefined; + } } /** @@ -8843,20 +8946,26 @@ async function chatWriteTurnComplete(options?: { publicAccessToken?: string }): * the assistant's response message via `onFinish`. * * Combines `toUIMessageStream()` + `onFinish` callback + `chat.pipe()`. - * Returns the captured `UIMessage`, or `undefined` if capture failed. + * Never throws on a stopped or failed stream: it returns a + * {@link PipeAndCaptureResult} whose `message` holds any partial output + * captured before the stop/failure, alongside a typed `status` (and `error` + * on failure). Save the partial after a stop or error without separate + * capture logic. * * @example * ```ts * const result = streamText({ model, messages, abortSignal: signal }); - * const response = await chat.pipeAndCapture(result, { signal }); - * if (response) conversation.addResponse(response); + * const { message, status, error } = await chat.pipeAndCapture(result, { signal }); + * if (message) conversation.addResponse(message); + * if (status === "error") logger.error("turn failed", { error }); * ``` */ async function pipeChatAndCapture( source: UIMessageStreamable, options?: { signal?: AbortSignal; spanName?: string; originalMessages?: UIMessage[] } -): Promise { +): Promise { let captured: UIMessage | undefined; + let capturedFinishReason: FinishReason | undefined; let resolveOnFinish: () => void; const onFinishPromise = new Promise((r) => { resolveOnFinish = r; @@ -8875,19 +8984,66 @@ async function pipeChatAndCapture( // the frontend replaces the partial message — wiping the // pre-injection text from the UI and the captured response. generateMessageId: resolvedOptions.generateMessageId ?? generateMessageId, - onFinish: ({ responseMessage }: { responseMessage: UIMessage }) => { + onFinish: ({ + responseMessage, + finishReason, + }: { + responseMessage: UIMessage; + finishReason?: FinishReason; + }) => { captured = responseMessage; + capturedFinishReason = finishReason; resolveOnFinish!(); }, }); - await pipeChat(uiStream, { - signal: options?.signal, - spanName: options?.spanName ?? "stream response", - }); - await onFinishPromise; + // Buffer chunks as they flow to the pipe so a transport failure that + // abandons the UI stream before `onFinish` fires — the one termination path + // that skips `onFinish` — can still reconstruct the partial. This retains + // chunk references only; the reassembly runs solely in the failure fallback + // below, so the happy path pays nothing extra. + const bufferedChunks: UIMessageChunk[] = []; + const tappedStream = tapUIMessageChunks(uiStream, bufferedChunks); - return captured; + let status: PipeAndCaptureResult["status"] = "complete"; + let error: unknown; + try { + await pipeChat(tappedStream, { + signal: options?.signal, + spanName: options?.spanName ?? "stream response", + }); + // The pipe can drain cleanly on a stop — the source stream just ends + // early — so classify by the signal rather than relying on a throw. + if (options?.signal?.aborted) { + status = "aborted"; + } + } catch (err) { + if ((err instanceof Error && err.name === "AbortError") || options?.signal?.aborted) { + status = "aborted"; + } else { + status = "error"; + error = err; + } + } + + // `onFinish` fires even on abort, carrying the partial — but a hard stop can + // prevent it from firing at all, so race it against a timeout to avoid + // hanging the caller. Mirrors chat.agent's capture path. + await Promise.race([onFinishPromise, new Promise((r) => setTimeout(r, 2_000))]); + + // A transport failure can abandon the UI stream before `onFinish` fires, so + // reconstruct the partial from the buffered chunks rather than losing output + // that already streamed. Only runs when `onFinish` produced nothing. + if (!captured && bufferedChunks.length > 0) { + captured = await assemblePartialFromChunks(bufferedChunks); + } + + return { + message: captured, + status, + finishReason: capturedFinishReason, + ...(error !== undefined ? { error } : {}), + }; } /** @@ -8903,8 +9059,8 @@ async function pipeChatAndCapture( * for (let turn = 0; turn < 100; turn++) { * const messages = await conversation.addIncoming(payload.messages, payload.trigger, turn); * const result = streamText({ model, messages }); - * const response = await chat.pipeAndCapture(result); - * if (response) await conversation.addResponse(response); + * const { message } = await chat.pipeAndCapture(result); + * if (message) await conversation.addResponse(message); * } * ``` */ @@ -9570,7 +9726,7 @@ function createChatSession( } let response: UIMessage | undefined; try { - response = await pipeChatAndCapture(source, { + const captured = await pipeChatAndCapture(source, { signal: combinedSignal, // On a non-final handover turn, thread the spliced partial so a // resumed tool round's tool-output chunks merge into the @@ -9579,18 +9735,18 @@ function createChatSession( // fresh response into the prior assistant message). ...(handoverThisTurn ? { originalMessages: accumulator.uiMessages } : {}), }); - } catch (error) { - if (error instanceof Error && error.name === "AbortError") { - if (runSignal.aborted) { - // Full cancel — don't accumulate - sessionMsgSub.off(); - await chatWriteTurnComplete(); - return undefined; - } - // Stop — fall through to accumulate partial response - } else { - throw error; + if (runSignal.aborted) { + // Full cancel — don't accumulate + sessionMsgSub.off(); + await chatWriteTurnComplete(); + return undefined; + } + // Surface a genuine stream failure to the caller. A user stop + // (status "aborted") falls through so the partial is accumulated. + if (captured.status === "error") { + throw captured.error; } + response = captured.message; } finally { // Detach at stream end (like the agent loop): the steering queue // can't inject anymore, so later arrivals must buffer for the next turn. diff --git a/packages/trigger-sdk/src/v3/test/test-session-handle.ts b/packages/trigger-sdk/src/v3/test/test-session-handle.ts index f97cdbcc061..860b7694e6e 100644 --- a/packages/trigger-sdk/src/v3/test/test-session-handle.ts +++ b/packages/trigger-sdk/src/v3/test/test-session-handle.ts @@ -124,32 +124,37 @@ export class TestSessionOutputChannel extends SessionOutputChannel { ): PipeStreamResult { const state = this.state; const readChunks: T[] = []; + let pipeError: unknown; let resolveDone!: () => void; const done = new Promise((resolve) => { resolveDone = resolve; }); (async () => { - const readable = ensureReadableStream(value); - const reader = readable.getReader(); + let reader: ReadableStreamDefaultReader | undefined; try { + const readable = ensureReadableStream(value); + reader = readable.getReader(); while (true) { const { done: d, value: v } = await reader.read(); - if (d) return; + if (d) break; readChunks.push(v as T); notify(state, v); } + } catch (err) { + // Mirror production: a source-stream error (or a throw from stream + // setup) rejects waitUntilComplete instead of being silently + // swallowed, so callers (e.g. chat.pipeAndCapture) can observe it. + pipeError = err; } finally { try { - reader.releaseLock(); + reader?.releaseLock(); } catch { // ignore } resolveDone(); } - })().catch(() => { - resolveDone(); - }); + })(); const replayStream = new ReadableStream({ async start(controller) { @@ -167,6 +172,7 @@ export class TestSessionOutputChannel extends SessionOutputChannel { }, waitUntilComplete: async () => { await done; + if (pipeError) throw pipeError; return emptyResult; }, }; @@ -262,7 +268,11 @@ export class TestSessionOutputChannel extends SessionOutputChannel { } } notify(this.state, synthetic); - return {}; + // Project a synthetic monotonic seq_num as the ack's `lastEventId`, + // mirroring what S2 returns in production (there it's the control + // record's seq_num). Using the running `.out` record count lets + // `chat.writeTurnComplete()` surface a real resume cursor in tests. + return { lastEventId: String(this.state.chunks.length) }; } /** diff --git a/packages/trigger-sdk/test/chat-pipe-and-capture.test.ts b/packages/trigger-sdk/test/chat-pipe-and-capture.test.ts new file mode 100644 index 00000000000..a25cb2b21b6 --- /dev/null +++ b/packages/trigger-sdk/test/chat-pipe-and-capture.test.ts @@ -0,0 +1,230 @@ +// Import the test harness FIRST — this installs the resource catalog so +// `chat.customAgent()` calls below register their task functions correctly. +import { mockChatAgent } from "../src/v3/test/index.js"; + +import { describe, expect, it } from "vitest"; +import type { UIMessage } from "ai"; +import { simulateReadableStream, streamText } from "ai"; +import { MockLanguageModelV3 } from "ai/test"; +import type { LanguageModelV3StreamPart } from "@ai-sdk/provider"; +import { chat } from "../src/v3/ai.js"; +import type { PipeAndCaptureResult } from "../src/v3/ai.js"; + +// ── Helpers ──────────────────────────────────────────────────────────── + +function userMessage(text: string, id: string): UIMessage { + return { id, role: "user", parts: [{ type: "text", text }] }; +} + +function textChunks(text: string, opts?: { split?: boolean }): LanguageModelV3StreamPart[] { + const deltas = opts?.split ? text.split(" ").map((w, i) => (i === 0 ? w : ` ${w}`)) : [text]; + return [ + { type: "text-start", id: "t1" }, + ...deltas.map((delta) => ({ type: "text-delta" as const, id: "t1", delta })), + { type: "text-end", id: "t1" }, + { + type: "finish", + finishReason: { unified: "stop", raw: "stop" }, + usage: { + inputTokens: { total: 5, noCache: 5, cacheRead: undefined, cacheWrite: undefined }, + outputTokens: { total: 5, text: 5, reasoning: undefined }, + }, + }, + ]; +} + +/** Model that streams `text` in one fast pass with a `stop` finish. */ +function fastModel(text: string) { + return new MockLanguageModelV3({ + doStream: async () => ({ stream: simulateReadableStream({ chunks: textChunks(text) }) }), + }); +} + +/** Model that streams `text` word-by-word with a wide gap before the final + * chunk, leaving a window to abort mid-stream after the first delta. */ +function slowModel(text: string) { + return new MockLanguageModelV3({ + doStream: async () => ({ + stream: simulateReadableStream({ + chunks: textChunks(text, { split: true }), + initialDelayInMs: 0, + chunkDelayInMs: 500, + }), + }), + }); +} + +function extractText(message: UIMessage | undefined): string { + if (!message) return ""; + return (message.parts as Array<{ type: string; text?: string }>) + .filter((p) => p.type === "text") + .map((p) => p.text ?? "") + .join(""); +} + +async function waitFor(check: () => boolean, timeoutMs = 5_000) { + const start = Date.now(); + while (Date.now() - start < timeoutMs) { + if (check()) return; + await new Promise((r) => setTimeout(r, 20)); + } + throw new Error("waitFor timed out"); +} + +function deltaCount(harness: { allChunks: unknown[] }): number { + return (harness.allChunks as { type?: string }[]).filter((c) => c.type === "text-delta").length; +} + +// ── Tests ────────────────────────────────────────────────────────────── + +describe("chat.pipeAndCapture", () => { + it("returns status 'complete' with the message and finish reason on a normal turn", async () => { + const captures: PipeAndCaptureResult[] = []; + const turnCompletes: Array<{ lastEventId?: string; sessionInEventId?: string }> = []; + + const agent = chat.customAgent({ + id: "pipe-capture.complete", + run: async () => { + const conversation = new chat.MessageAccumulator(); + const next = await chat.messages.waitWithIdleTimeout({ + idleTimeoutInSeconds: 60, + timeout: "1h", + }); + if (!next.ok) return; + const wire = next.output as { message?: UIMessage; trigger: string }; + const incoming = wire.message ? [wire.message] : []; + const messages = await conversation.addIncoming(incoming, wire.trigger, 0); + const result = streamText({ model: fastModel("hello world"), messages }); + const captured = await chat.pipeAndCapture(result); + captures.push(captured); + if (captured.message) await conversation.addResponse(captured.message); + turnCompletes.push(await chat.writeTurnComplete()); + }, + }); + + const harness = mockChatAgent(agent, { chatId: "pc-complete" }); + try { + await harness.sendMessage(userMessage("hi", "u-1")); + await waitFor(() => captures.length >= 1 && turnCompletes.length >= 1); + + expect(captures[0]!.status).toBe("complete"); + expect(extractText(captures[0]!.message)).toBe("hello world"); + expect(captures[0]!.finishReason).toBe("stop"); + expect(captures[0]!.error).toBeUndefined(); + // chat.writeTurnComplete() surfaces the .out resume cursor for the next + // turn. (sessionInEventId is a passthrough of the same value written to + // the session-in-event-id header; the in-memory harness doesn't track + // the .in dispatch cursor, so its value isn't asserted here.) + expect(typeof turnCompletes[0]!.lastEventId).toBe("string"); + expect(turnCompletes[0]!.lastEventId!.length).toBeGreaterThan(0); + } finally { + await harness.close(); + } + }); + + it("returns status 'error' with the thrown error and does not throw when the stream fails", async () => { + const captures: PipeAndCaptureResult[] = []; + let runThrew = false; + + // Synthetic source whose UI stream errors after emitting a partial. This + // deterministically drives the pipe-failure path without depending on the + // AI SDK's model-error handling. Chunks are delivered one-per-pull before + // the error so they aren't discarded — calling controller.error() in the + // same tick as enqueue() would reset the queue and drop them. + const partialChunks = [ + { type: "start", messageId: "a-err" }, + { type: "text-start", id: "t1" }, + { type: "text-delta", id: "t1", delta: "partial" }, + ]; + const erroringSource = { + toUIMessageStream() { + let i = 0; + return new ReadableStream({ + pull(controller) { + if (i < partialChunks.length) { + controller.enqueue(partialChunks[i++]); + } else { + controller.error(new Error("boom")); + } + }, + }); + }, + }; + + const agent = chat.customAgent({ + id: "pipe-capture.error", + run: async () => { + const next = await chat.messages.waitWithIdleTimeout({ + idleTimeoutInSeconds: 60, + timeout: "1h", + }); + if (!next.ok) return; + try { + captures.push(await chat.pipeAndCapture(erroringSource as never)); + } catch { + runThrew = true; + } + await chat.writeTurnComplete(); + }, + }); + + const harness = mockChatAgent(agent, { chatId: "pc-error" }); + try { + await harness.sendMessage(userMessage("go", "u-1")); + await waitFor(() => captures.length >= 1); + + expect(runThrew).toBe(false); + expect(captures[0]!.status).toBe("error"); + expect(captures[0]!.error).toBeInstanceOf(Error); + expect((captures[0]!.error as Error).message).toBe("boom"); + // The partial that streamed before the failure is reconstructed from the + // buffered chunks even though onFinish never fired on this hard-error path. + expect(extractText(captures[0]!.message)).toBe("partial"); + } finally { + await harness.close(); + } + }); + + it("returns status 'aborted' and preserves the partial message on a mid-stream stop", async () => { + const captures: PipeAndCaptureResult[] = []; + + const agent = chat.customAgent({ + id: "pipe-capture.aborted", + run: async () => { + const stop = chat.createStopSignal(); + const conversation = new chat.MessageAccumulator(); + const next = await chat.messages.waitWithIdleTimeout({ + idleTimeoutInSeconds: 60, + timeout: "1h", + }); + if (!next.ok) return; + const wire = next.output as { message?: UIMessage; trigger: string }; + const incoming = wire.message ? [wire.message] : []; + const messages = await conversation.addIncoming(incoming, wire.trigger, 0); + const result = streamText({ + model: slowModel("one two three four"), + messages, + abortSignal: stop.signal, + }); + captures.push(await chat.pipeAndCapture(result, { signal: stop.signal })); + await chat.writeTurnComplete(); + stop.cleanup(); + }, + }); + + const harness = mockChatAgent(agent, { chatId: "pc-aborted" }); + try { + void harness.sendMessage(userMessage("hi", "u-1")); + // Stop once the first delta has streamed but before the turn finishes. + await waitFor(() => deltaCount(harness) >= 1); + await harness.sendStop(); + await waitFor(() => captures.length >= 1); + + expect(captures[0]!.status).toBe("aborted"); + // The partial that streamed before the stop is preserved. + expect(extractText(captures[0]!.message).startsWith("one")).toBe(true); + } finally { + await harness.close(); + } + }); +});