From 6ac545a095556db3e7cfba8dcc0e2ec63c22adbd Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Thu, 23 Apr 2026 10:54:10 +0100 Subject: [PATCH] =?UTF-8?q?feat(sdk):=20chat.agent=20=E2=86=92=20Sessions?= =?UTF-8?q?=20migration=20(phases=20B=20+=20C=20+=20min=20E)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Rewires chat.agent's internal I/O and TriggerChatTransport's send + subscribe paths onto the Session primitive. Minimum token-scope work included so the transport's session endpoints actually authenticate. Phase B — chat.agent internals (ai.ts) - New ChatInputChunk tagged union (`kind: "message" | "stop"`) — replaces the two-stream split (chat-messages + chat-stop) with a single Session `.in` channel. - New chatSessionHandleKey locals slot populated at run start from `payload.sessionId ?? payload.chatId`. Every module-level helper now resolves to the per-run session handle. - Module-level `chatStream`, `messagesInput`, `stopInput` become thin facades over the session. `chatStream` mirrors `RealtimeDefinedStream` and delegates to `handle.out`. `messagesInput` / `stopInput` mirror `RealtimeDefinedInputStream<…>` and filter `.in` by kind — the two internal `.on()`/`.waitWithIdleTimeout()` callers and the `chat.messages` / `chat.createStopSignal` public exposures keep their existing shapes. - Every `streams.writer(CHAT_STREAM_KEY, …)` callsite swaps to `chatStream.writer(…)` so all chat output flows through `session.out` → `SessionStreamInstance` → direct-to-S2. - Threaded `sessionId` through `ChatTaskWirePayload` / `ChatTaskPayload` / `ChatTaskRunPayload` so advanced users can `sessions.open(sessionId)` directly from `run()`. Phase C — TriggerChatTransport (chat.ts) - `ChatSessionState` keys durable identity on `sessionId` (friendlyId); `runId` becomes an optional hint about whether a run is live. - `ensureSession(chatId)` lazily upserts the Session via `apiClient.createSession({type: "chat.agent", externalId: chatId})` on the direct `accessToken` path. Idempotent — two tabs on the same chat converge. - `sendMessages`, `sendPendingMessage`, `stopGeneration`, `sendAction` all go through `appendToSessionStream(sessionId, "in", serializeInputChunk({kind: …}))` — one endpoint, one tag per record. - SSE subscribe URL moves from `/realtime/v1/streams/{runId}/chat` to `/realtime/v1/sessions/{sessionId}/out`. The old run-scoped `subscribeToStream` is replaced by `subscribeToSessionStream`. Incoming chunks come back as JSON strings on the session channel (server wraps records as `{data, id}` on S2), so the subscribe loop parses them back into objects to keep the rest of the control flow (turn-complete / upgrade-required / skipToTurnComplete) unchanged. - Upgrade-required re-trigger keeps the same Session and swaps only the runId + token. - `getSession` / `setSession` / `setOnSessionChange` / persistence shape all grow a `sessionId` field (runId now optional). Phase E — minimum token scopes - `chat.createTriggerAction` (server side) now creates the Session before triggering so it can (a) thread `sessionId` into the run payload and (b) mint a token with both run and session scopes. Returns `sessionId` in its result so the transport can skip its own `sessions.create` call on the server-side-trigger path. - `TriggerChatTaskResult` gains optional `sessionId`. - The two in-run PAT refresh sites (preloadAccessToken, turnAccessToken) add `read:sessions:{sessionId}` + `write:sessions:{sessionId}` alongside the existing run scopes. Known follow-ups (deferred to later passes) - Phase D: `AgentChat` / `ChatStream` in chat-client.ts still uses the old `/realtime/v1/streams/{runId}/chat` path. Used by server- side task-to-task compositions, not the browser transport. - Phase F: delete CHAT_STREAM_KEY, CHAT_MESSAGES_STREAM_ID, CHAT_STOP_STREAM_ID from chat-constants.ts + ai-chat smoke verify. --- packages/trigger-sdk/src/v3/ai.ts | 435 ++++++++++++++- packages/trigger-sdk/src/v3/chat.test.ts | 2 + packages/trigger-sdk/src/v3/chat.ts | 640 +++++++++++++++-------- 3 files changed, 831 insertions(+), 246 deletions(-) diff --git a/packages/trigger-sdk/src/v3/ai.ts b/packages/trigger-sdk/src/v3/ai.ts index f28d45bb7..cd33c5c82 100644 --- a/packages/trigger-sdk/src/v3/ai.ts +++ b/packages/trigger-sdk/src/v3/ai.ts @@ -2,11 +2,23 @@ import { accessoryAttributes, AnyTask, getSchemaParseFn, + InputStreamOncePromise, + type InputStreamOnceOptions, + type InputStreamWaitOptions, + type InputStreamWaitWithIdleTimeoutOptions, isSchemaZodEsque, logger, + ManualWaitpointPromise, + type PipeStreamResult, + type RealtimeDefinedInputStream, + type RealtimeDefinedStream, + type ReadStreamOptions, SemanticInternalAttributes, + type SendInputStreamOptions, Task, taskContext, + type AppendStreamOptions, + type InputStreamOnceResult, type inferSchemaIn, type inferSchemaOut, type PipeStreamOptions, @@ -15,6 +27,7 @@ import { type TaskSchema, type TaskRunContext, type TaskWithSchema, + type WriterStreamOptions, } from "@trigger.dev/core/v3"; import type { FinishReason, @@ -48,6 +61,14 @@ import { spawn } from "node:child_process"; import * as fs from "node:fs/promises"; import * as nodePath from "node:path"; import { streams } from "./streams.js"; +import { + sessions, + type SessionHandle, + type SessionInputChannel, + type SessionOutputChannel, + type SessionPipeStreamOptions, + type SessionSubscribeOptions, +} from "./sessions.js"; import { createTask, trigger as triggerTaskInternal } from "./shared.js"; import { resourceCatalog } from "@trigger.dev/core/v3"; import type { TriggerChatTaskParams, TriggerChatTaskResult } from "./chat.js"; @@ -100,6 +121,34 @@ type ChatTurnContext = { }; const chatTurnContextKey = locals.create("chat.turnContext"); +/** + * Per-run slot holding the Session handle that backs this chat's `.in` / + * `.out` channels. Populated at the top of `chatAgent`'s run function from + * `payload.sessionId`; read by every module-level helper (`chatStream`, + * `messagesInput`, `stopInput`, `streams.writer(CHAT_STREAM_KEY, …)` + * callers) so the chat.agent internals can remain the same module-level + * shape they were when the I/O was run-scoped. + * @internal + */ +const chatSessionHandleKey = locals.create("chat.sessionHandle"); + +/** + * Resolve the Session handle for the current chat.agent run. Throws if + * called outside of a chat.agent `run()` — every internal consumer is + * inside the run, and every external consumer goes through the public + * `sessions.open(id)` entry point. + * @internal + */ +function getChatSession(): SessionHandle { + const handle = locals.get(chatSessionHandleKey); + if (!handle) { + throw new Error( + "chat.agent session handle is not initialized. This indicates a chat.agent helper was used outside of a chat.agent run, or the transport did not send a sessionId." + ); + } + return handle; +} + type ToolResultContent = Array< | { type: "text"; @@ -442,8 +491,9 @@ export const CHAT_STREAM_KEY = _CHAT_STREAM_KEY; export { CHAT_MESSAGES_STREAM_ID, CHAT_STOP_STREAM_ID }; /** - * Typed chat output stream. Provides `.writer()`, `.pipe()`, `.append()`, - * and `.read()` methods pre-bound to the chat stream key and typed to `UIMessageChunk`. + * Typed chat output stream — `.writer()`, `.pipe()`, `.append()`, and + * `.read()` methods pre-bound to this run's Session `.out` channel and + * typed to `UIMessageChunk`. * * Use from within a `chat.agent` run to write custom chunks: * ```ts @@ -457,12 +507,36 @@ export { CHAT_MESSAGES_STREAM_ID, CHAT_STOP_STREAM_ID }; * await waitUntilComplete(); * ``` * - * Use from a subtask to stream back to the parent chat: - * ```ts - * chat.stream.pipe(myStream, { target: "root" }); - * ``` + * Backed by the Session primitive so a chat's output outlives any single + * run — subscribers (browser transport, server-side `ChatStream`) read + * the session's `.out`, not a per-run stream. Run-scoped `target` + * options on `.pipe()` are honoured as no-ops; the session is the target. */ -const chatStream = streams.define({ id: _CHAT_STREAM_KEY }); +const chatStream: RealtimeDefinedStream = { + id: _CHAT_STREAM_KEY, + pipe(value, options) { + const { target: _target, ...sessionOptions } = (options ?? {}) as PipeStreamOptions; + return getChatSession().out.pipe( + value, + sessionOptions as SessionPipeStreamOptions + ); + }, + async read(_runId, options) { + // Session channels don't need a runId — the session is the address. + // Keep the signature for backward compatibility with the run-scoped + // RealtimeDefinedStream shape, but ignore the argument. + return getChatSession().out.read( + options as SessionSubscribeOptions | undefined + ); + }, + async append(value, options) { + const { target: _target, ...sessionOptions } = (options ?? {}) as AppendStreamOptions; + return getChatSession().out.append(value, sessionOptions as SessionPipeStreamOptions); + }, + writer(options) { + return getChatSession().out.writer(options); + }, +}; // --------------------------------------------------------------------------- // chat.response — write data parts that persist to the response message @@ -492,7 +566,7 @@ const chatResponse = { */ write(part: UIMessageChunk): void { queueResponsePart(part); - const { waitUntilComplete } = streams.writer(CHAT_STREAM_KEY, { + const { waitUntilComplete } = chatStream.writer({ spanName: "chat.response.write", collapsed: true, execute: ({ write }) => { @@ -532,7 +606,7 @@ const chatStoreListenersKey = locals.create>( /** @internal — write a store chunk onto the chat output stream. */ function writeStoreChunk(chunk: ChatStoreChunk): void { - const { waitUntilComplete } = streams.writer(CHAT_STREAM_KEY, { + const { waitUntilComplete } = chatStream.writer({ spanName: chunk.type === "store-snapshot" ? "chat.store.set" : "chat.store.patch", collapsed: true, execute: ({ write }) => { @@ -736,6 +810,16 @@ export type ChatTaskWirePayload = + | { + kind: "message"; + /** + * Full wire payload for a new user message or regeneration. Mirrors + * what the legacy `chat-messages` input stream carried. + */ + payload: ChatTaskWirePayload; + } + | { + kind: "stop"; + /** Optional human-readable reason. Maps to the legacy `chat-stop` record. */ + message?: string; + }; + /** * The payload shape passed to the `chatAgent` run function. * @@ -789,6 +897,14 @@ export type ChatTaskPayload = { previousRunId?: string; /** Whether this run was preloaded before the first message. */ preloaded: boolean; + /** + * The friendlyId of the Session primitive backing this chat. Use with + * `sessions.open(sessionId)` when you need direct access to the session's + * `.in` / `.out` channels outside the hooks the agent already wires for + * you. Undefined only for legacy transports that predate the sessions + * migration. + */ + sessionId?: string; }; /** @@ -821,8 +937,213 @@ export type ChatTaskRunPayload = ChatTaskPayload({ id: CHAT_MESSAGES_STREAM_ID }); -const stopInput = streams.input<{ stop: true; message?: string }>({ id: CHAT_STOP_STREAM_ID }); +// +// Both `messagesInput` and `stopInput` are thin facades over the current +// run's Session `.in` channel. The Session carries a single tagged stream +// (`ChatInputChunk`); these facades filter by `kind` so existing call +// sites (both internal and exposed via `chat.messages` / `chat.createStopSignal`) +// keep their original shape. Each accessor resolves the session handle +// lazily via `getChatSession()` so the module-level references stay +// compatible with the pre-migration wiring. +const messagesInput: RealtimeDefinedInputStream = { + id: CHAT_MESSAGES_STREAM_ID, + on(handler) { + return getChatSession().in.on((chunk) => { + if (chunk.kind === "message") { + return handler(chunk.payload); + } + }); + }, + once(options) { + const ctx = taskContext.ctx; + const runId = ctx?.run.id; + + return new InputStreamOncePromise((resolve, reject) => { + tracer + .startActiveSpan( + options?.spanName ?? `chat.messages.once()`, + async () => { + while (true) { + const result = await getChatSession().in.once(options); + if (!result.ok) { + resolve(result as InputStreamOnceResult); + return; + } + if (result.output.kind === "message") { + resolve({ ok: true, output: result.output.payload }); + return; + } + // Non-message chunks (stops) are handled by the stopInput + // facade's persistent listener; loop and wait for the next. + } + }, + { + attributes: { + [SemanticInternalAttributes.STYLE_ICON]: "streams", + [SemanticInternalAttributes.ENTITY_TYPE]: "input-stream", + ...(runId + ? { + [SemanticInternalAttributes.ENTITY_ID]: `${runId}:${CHAT_MESSAGES_STREAM_ID}`, + } + : {}), + streamId: CHAT_MESSAGES_STREAM_ID, + ...accessoryAttributes({ + items: [{ text: CHAT_MESSAGES_STREAM_ID, variant: "normal" }], + style: "codepath", + }), + }, + } + ) + .catch(reject); + }); + }, + peek() { + const chunk = getChatSession().in.peek(); + if (chunk && chunk.kind === "message") return chunk.payload; + return undefined; + }, + wait(options) { + return new ManualWaitpointPromise(async (resolve, reject) => { + try { + while (true) { + const result = await getChatSession().in.wait(options); + if (!result.ok) { + resolve(result); + return; + } + if (result.output.kind === "message") { + resolve({ ok: true, output: result.output.payload }); + return; + } + // Stop chunks are handled by the stopInput facade's persistent + // listener; loop back into the suspending wait. + } + } catch (error) { + reject(error); + } + }); + }, + async waitWithIdleTimeout(options) { + while (true) { + const result = await getChatSession().in.waitWithIdleTimeout(options); + if (!result.ok) return result; + if (result.output.kind === "message") { + return { ok: true, output: result.output.payload }; + } + // Swallow stop-kind chunks — persistent stop listener already handled + // the abort; we just loop for the next message. + } + }, + async send(_runId, data, options) { + // The `runId` argument is kept for signature parity with + // `RealtimeDefinedInputStream` but ignored — sessions are addressed + // by sessionId, not runId. Callers producing messages from outside + // the run should prefer the transport's `session.in.send(...)` path. + await getChatSession().in.send( + { kind: "message", payload: data } satisfies ChatInputChunk, + options?.requestOptions + ); + }, +}; + +const stopInput: RealtimeDefinedInputStream<{ stop: true; message?: string }> = { + id: CHAT_STOP_STREAM_ID, + on(handler) { + return getChatSession().in.on((chunk) => { + if (chunk.kind === "stop") { + return handler({ stop: true, message: chunk.message }); + } + }); + }, + once(options) { + const ctx = taskContext.ctx; + const runId = ctx?.run.id; + + return new InputStreamOncePromise<{ stop: true; message?: string }>((resolve, reject) => { + tracer + .startActiveSpan( + options?.spanName ?? `chat.stop.once()`, + async () => { + while (true) { + const result = await getChatSession().in.once(options); + if (!result.ok) { + resolve(result as InputStreamOnceResult<{ stop: true; message?: string }>); + return; + } + if (result.output.kind === "stop") { + resolve({ + ok: true, + output: { stop: true, message: result.output.message }, + }); + return; + } + } + }, + { + attributes: { + [SemanticInternalAttributes.STYLE_ICON]: "streams", + [SemanticInternalAttributes.ENTITY_TYPE]: "input-stream", + ...(runId + ? { + [SemanticInternalAttributes.ENTITY_ID]: `${runId}:${CHAT_STOP_STREAM_ID}`, + } + : {}), + streamId: CHAT_STOP_STREAM_ID, + ...accessoryAttributes({ + items: [{ text: CHAT_STOP_STREAM_ID, variant: "normal" }], + style: "codepath", + }), + }, + } + ) + .catch(reject); + }); + }, + peek() { + const chunk = getChatSession().in.peek(); + if (chunk && chunk.kind === "stop") { + return { stop: true, message: chunk.message }; + } + return undefined; + }, + wait(options) { + return new ManualWaitpointPromise<{ stop: true; message?: string }>(async (resolve, reject) => { + try { + while (true) { + const result = await getChatSession().in.wait(options); + if (!result.ok) { + resolve(result); + return; + } + if (result.output.kind === "stop") { + resolve({ + ok: true, + output: { stop: true, message: result.output.message }, + }); + return; + } + } + } catch (error) { + reject(error); + } + }); + }, + async waitWithIdleTimeout(options) { + while (true) { + const result = await getChatSession().in.waitWithIdleTimeout(options); + if (!result.ok) return result; + if (result.output.kind === "stop") { + return { ok: true, output: { stop: true, message: result.output.message } }; + } + } + }, + async send(_runId, data, options) { + await getChatSession().in.send( + { kind: "stop", message: data?.message } satisfies ChatInputChunk, + options?.requestOptions + ); + }, +}; /** * Per-turn deferred promises. Registered via `chat.defer()`, awaited @@ -1551,11 +1872,14 @@ async function chatCompact( const compactionId = generateMessageId(); let summary!: string; - const { waitUntilComplete } = streams.writer(CHAT_STREAM_KEY, { + const { waitUntilComplete } = chatStream.writer({ spanName: "stream compaction chunks", collapsed: true, execute: async ({ write, merge }) => { - write({ type: "step-start" }); + // Control chunks aren't part of UIMessageChunk's discriminated + // union but flow on the same session.out so subscribers can + // intercept them — cast on the way out. + write({ type: "step-start" } as unknown as UIMessageChunk); write({ type: "data-compaction", id: compactionId, @@ -1753,7 +2077,7 @@ async function drainSteeringQueue( // knows which messages were injected and where in the response. if (injected.length > 0) { try { - const { waitUntilComplete } = streams.writer(CHAT_STREAM_KEY, { + const { waitUntilComplete } = chatStream.writer({ collapsed: true, execute: ({ write }) => { write({ @@ -3386,6 +3710,17 @@ function chatAgent< ) => { locals.set(chatAgentRunContextKey, ctx); + // Bind the run to its backing Session so every module-level helper + // (chat.stream, chat.messages, streams.writer(CHAT_STREAM_KEY, …)) + // resolves to this chat's `.in` / `.out` channels. + // + // The transport opens/creates the session with `externalId = chatId` + // and threads its friendlyId through `payload.sessionId`. For legacy + // transports that predate the migration the field is missing; fall + // back to opening by chatId (which the transport uses as externalId). + const sessionIdForHandle = payload.sessionId ?? payload.chatId; + locals.set(chatSessionHandleKey, sessions.open(sessionIdForHandle)); + // Set gen_ai.conversation.id on the run-level span for dashboard context const activeSpan = trace.getActiveSpan(); if (activeSpan) { @@ -3457,8 +3792,14 @@ function chatAgent< try { preloadAccessToken = await auth.createPublicToken({ scopes: { - read: { runs: currentRunId }, - write: { inputStreams: currentRunId }, + read: { + runs: currentRunId, + ...(sessionIdForHandle ? { sessions: sessionIdForHandle } : {}), + }, + write: { + inputStreams: currentRunId, + ...(sessionIdForHandle ? { sessions: sessionIdForHandle } : {}), + }, }, expirationTime: chatAccessTokenTTL, }); @@ -4031,8 +4372,14 @@ function chatAgent< try { turnAccessToken = await auth.createPublicToken({ scopes: { - read: { runs: currentRunId }, - write: { inputStreams: currentRunId }, + read: { + runs: currentRunId, + ...(sessionIdForHandle ? { sessions: sessionIdForHandle } : {}), + }, + write: { + inputStreams: currentRunId, + ...(sessionIdForHandle ? { sessions: sessionIdForHandle } : {}), + }, }, expirationTime: chatAccessTokenTTL, }); @@ -4443,7 +4790,7 @@ function chatAgent< async (compactionSpan) => { const compactionId = generateMessageId(); - const { waitUntilComplete } = streams.writer(CHAT_STREAM_KEY, { + const { waitUntilComplete } = chatStream.writer({ spanName: "stream compaction chunks", collapsed: true, execute: async ({ write, merge }) => { @@ -6635,7 +6982,33 @@ function createChatTriggerAction( options?: CreateChatTriggerActionOptions ): (params: TriggerChatTaskParams) => Promise { return async (params: TriggerChatTaskParams): Promise => { - const handle = await triggerTaskInternal(taskId, params.payload, { + // Create (or upsert) the backing Session first so we can thread its + // friendlyId through the run's payload and the PAT's scopes. Using + // the server's secret key here means the browser's `accessToken` + // doesn't need `write:sessions` — the session lifecycle lives on + // the server action side, which is what the `triggerTask` pattern + // is for. + const chatId = (params.payload as { chatId?: string }).chatId; + let sessionId = (params.payload as { sessionId?: string }).sessionId; + if (!sessionId) { + if (!chatId) { + throw new Error( + "chat.createTriggerAction: payload.chatId is required so the backing Session can be keyed on externalId." + ); + } + const session = await sessions.create({ + type: "chat.agent", + externalId: chatId, + }); + sessionId = session.id; + } + + const payloadWithSession = { + ...params.payload, + sessionId, + }; + + const handle = await triggerTaskInternal(taskId, payloadWithSession, { tags: params.options.tags, queue: params.options.queue, maxAttempts: params.options.maxAttempts, @@ -6645,13 +7018,19 @@ function createChatTriggerAction( const publicAccessToken = await auth.createPublicToken({ scopes: { - read: { runs: handle.id }, - write: { inputStreams: handle.id }, + read: { + runs: handle.id, + sessions: sessionId, + }, + write: { + inputStreams: handle.id, + sessions: sessionId, + }, }, expirationTime: options?.tokenTTL ?? "1h", }); - return { runId: handle.id, publicAccessToken }; + return { runId: handle.id, publicAccessToken, sessionId }; }; } @@ -6792,14 +7171,16 @@ async function writeTurnCompleteChunk( chatId?: string, publicAccessToken?: string ): Promise { - const { waitUntilComplete } = streams.writer(CHAT_STREAM_KEY, { + const { waitUntilComplete } = chatStream.writer({ spanName: "turn complete", collapsed: true, execute: ({ write }) => { + // Transport-intercepted control chunk — not a valid UIMessageChunk + // type but travels on the same session.out stream. write({ type: "trigger:turn-complete", ...(publicAccessToken ? { publicAccessToken } : {}), - }); + } as unknown as UIMessageChunk); }, }); return await waitUntilComplete(); @@ -6811,13 +7192,13 @@ async function writeTurnCompleteChunk( * @internal */ async function writeUpgradeRequiredChunk(): Promise { - const { waitUntilComplete } = streams.writer(CHAT_STREAM_KEY, { + const { waitUntilComplete } = chatStream.writer({ spanName: "upgrade required", collapsed: true, execute: ({ write }) => { write({ type: "trigger:upgrade-required", - }); + } as unknown as UIMessageChunk); }, }); return await waitUntilComplete(); diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index 25ed39e62..69d4b14f2 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -510,6 +510,7 @@ describe("TriggerChatTransport", () => { accessToken: "token", sessions: { "chat-completed": { + sessionId: "session_completed", runId: "run_completed", publicAccessToken: "pub_token", lastEventId: "42", @@ -554,6 +555,7 @@ describe("TriggerChatTransport", () => { baseURL: "https://api.test.trigger.dev", sessions: { "chat-streaming": { + sessionId: "session_streaming", runId: "run_streaming", publicAccessToken: "pub_token", lastEventId: "10", diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 8cddc02cb..91aa9587d 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -24,6 +24,7 @@ import type { ChatTransport, UIMessage, UIMessageChunk, ChatRequestOptions } from "ai"; import { ApiClient, SSEStreamSubscription } from "@trigger.dev/core/v3"; +import type { ChatInputChunk } from "./ai.js"; /** * Detect 401/403 from realtime/input-stream calls without relying on `instanceof` @@ -36,7 +37,11 @@ function isRunPatAuthError(error: unknown): boolean { const e = error as { name?: string; status?: number }; return e.name === "TriggerApiError" && (e.status === 401 || e.status === 403); } -import { CHAT_MESSAGES_STREAM_ID, CHAT_STOP_STREAM_ID } from "./chat-constants.js"; +// Stream-ID constants are no longer used — the transport writes a tagged +// ChatInputChunk to the session's `.in` channel (append route) and reads +// UIMessageChunks from `.out` (SSE subscribe). Kept imported from +// chat-constants.js only for callers that still import them directly; the +// legacy constants will be deleted in a follow-up cleanup. import { ChatTabCoordinator } from "./chat-tab-coordinator.js"; const DEFAULT_STREAM_KEY = "chat"; @@ -97,8 +102,21 @@ export type TriggerChatTaskParams = { export type TriggerChatTaskResult = { /** The run ID from the triggered task. */ runId: string; - /** A run-scoped public access token for stream subscription and input stream writes. */ + /** + * Public access token scoped to the run and the backing session. Must + * include `read:runs:{runId}`, `write:inputStreams:{runId}`, + * `read:sessions:{sessionId}`, and `write:sessions:{sessionId}` so the + * transport can both send on `session.in` and subscribe to + * `session.out`. + */ publicAccessToken: string; + /** + * Session friendlyId backing this chat. Optional for callers that still + * create the session from the transport side; when present, the + * transport skips its own `sessions.create` call and threads the + * provided id through subsequent sends + subscribes. + */ + sessionId?: string; }; /** Common options shared by all TriggerChatTransport configurations. */ @@ -169,12 +187,21 @@ type TriggerChatTransportOptionsBase = { * task: "my-chat", * accessToken, * sessions: { - * "chat-abc": { runId: "run_123", publicAccessToken: "...", lastEventId: "42" }, + * "chat-abc": { sessionId: "session_abc", runId: "run_123", publicAccessToken: "...", lastEventId: "42" }, * }, * }); * ``` */ - sessions?: Record; + sessions?: Record< + string, + { + sessionId: string; + runId?: string; + publicAccessToken: string; + lastEventId?: string; + isStreaming?: boolean; + } + >; /** * Called whenever a chat session's state changes. @@ -204,7 +231,15 @@ type TriggerChatTransportOptionsBase = { */ onSessionChange?: ( chatId: string, - session: { runId: string; publicAccessToken: string; lastEventId?: string; isStreaming?: boolean } | null + session: + | { + sessionId: string; + runId?: string; + publicAccessToken: string; + lastEventId?: string; + isStreaming?: boolean; + } + | null ) => void; /** @@ -342,10 +377,29 @@ export type TriggerChatTransportOptions = /** * Internal state for tracking active chat sessions. + * + * After the Sessions migration, the durable identity of a chat is the + * Session primitive (keyed on `sessionId`, externalId `chatId`). The + * transport produces messages by appending to `.in` on that session and + * consumes responses by subscribing to `.out`. The `runId` is kept only + * as a hint for "is there a live run right now?" — e.g. to decide + * whether the next message needs to trigger a fresh run. * @internal */ type ChatSessionState = { - runId: string; + /** Session friendlyId (`session_*`). Durable across runs. */ + sessionId: string; + /** + * The live run's friendlyId, if any. Cleared on turn-complete / + * run-end signals. Used to know whether the next sendMessages should + * trigger a new run or rely on the existing run reading from + * `session.in`. + */ + runId?: string; + /** + * Public access token. Scoped to both the session and (for legacy + * run-scoped endpoints still in the transition) the active run. + */ publicAccessToken: string; /** Last SSE event ID — used to resume the stream without replaying old events. */ lastEventId?: string; @@ -403,7 +457,15 @@ export class TriggerChatTransport implements ChatTransport { private _onSessionChange: | (( chatId: string, - session: { runId: string; publicAccessToken: string; lastEventId?: string; isStreaming?: boolean } | null + session: + | { + sessionId: string; + runId?: string; + publicAccessToken: string; + lastEventId?: string; + isStreaming?: boolean; + } + | null ) => void) | undefined; @@ -455,7 +517,15 @@ export class TriggerChatTransport implements ChatTransport { // Restore sessions from external storage if (options.sessions) { for (const [chatId, session] of Object.entries(options.sessions)) { + // Accept either the new `sessionId` field or the legacy + // shape (runId only). If only runId is present there's no way + // to recover the sessionId synchronously — we leave it missing + // and `ensureSession` will upsert on next send (externalId = + // chatId, so the server returns the existing friendlyId). + const sessionId = (session as { sessionId?: string }).sessionId; + if (!sessionId) continue; this.sessions.set(chatId, { + sessionId, runId: session.runId, publicAccessToken: session.publicAccessToken, lastEventId: session.lastEventId, @@ -465,6 +535,46 @@ export class TriggerChatTransport implements ChatTransport { } } + /** + * Resolve the Session backing this chat, creating it on the server + * if the transport doesn't already have a `sessionId` cached. The + * create call is idempotent — it upserts on `externalId = chatId`, + * so two tabs of the same chat converge to one row. + */ + private async ensureSession(chatId: string): Promise { + const existing = this.sessions.get(chatId); + if (existing?.sessionId) return existing; + + const token = await this.resolveAccessToken({ chatId, purpose: "trigger" }); + const apiClient = new ApiClient(this.baseURL, token); + const created = await apiClient.createSession({ + type: "chat.agent", + externalId: chatId, + }); + + const state: ChatSessionState = { + sessionId: created.id, + runId: existing?.runId, + publicAccessToken: existing?.publicAccessToken ?? token, + lastEventId: existing?.lastEventId, + isStreaming: existing?.isStreaming, + skipToTurnComplete: existing?.skipToTurnComplete, + }; + this.sessions.set(chatId, state); + this.notifySessionChange(chatId, state); + return state; + } + + /** + * Serialize a ChatInputChunk body for `POST …/sessions/:session/:io/append`. + * Session channel records are raw JSON strings — the server wraps them + * in `{ data: , id }` for S2 storage and the subscribe side + * parses the string back for consumers. + */ + private serializeInputChunk(chunk: ChatInputChunk): string { + return JSON.stringify(chunk); + } + sendMessages = async ( options: { trigger: "submit-message" | "regenerate-message"; @@ -498,12 +608,30 @@ export class TriggerChatTransport implements ChatTransport { metadata: mergedMetadata, }; - const session = this.sessions.get(chatId); + // Resolve (or lazily create) the backing Session. For the + // `triggerTask` callback path we defer session creation to the + // server action (see `createChatTriggerAction`) — it has the secret + // key and returns the sessionId in its result. For the direct + // `accessToken` path we create the session from here (requires the + // token to have `write:sessions` scope). + let state: ChatSessionState | undefined = this.sessions.get(chatId); + if (!state?.sessionId) { + if (this.triggerTaskFn) { + // Leave state as-is (undefined or runless); first-message path + // below triggers and adopts the sessionId from the result. + state = state; + } else { + state = await this.ensureSession(chatId); + } + } let isContinuation = false; let previousRunId: string | undefined; - // If we have an existing run, send the message via input stream - // to resume the conversation in the same run. - if (session?.runId) { + + // If a run is already live for this chat, deliver the new message + // via `session.in.send({kind: "message", payload})`. The agent's + // `chat.messages.waitWithIdleTimeout` (now backed by session.in) + // picks it up without tearing down the run. + if (state?.runId) { const slicedMessages = trigger === "submit-message" ? messages.slice(-1) : messages; const minimalPayload = { ...payload, @@ -512,17 +640,21 @@ export class TriggerChatTransport implements ChatTransport { const sendChatMessages = async (token: string) => { const apiClient = new ApiClient(this.baseURL, token); - await apiClient.sendInputStream(session.runId, CHAT_MESSAGES_STREAM_ID, minimalPayload); + await apiClient.appendToSessionStream( + state.sessionId, + "in", + this.serializeInputChunk({ kind: "message", payload: minimalPayload }) + ); }; let inputSendOk = false; try { - await sendChatMessages(session.publicAccessToken); + await sendChatMessages(state.publicAccessToken); inputSendOk = true; } catch (err) { if (isRunPatAuthError(err) && this.renewRunAccessToken) { - const newToken = await this.renewRunPatForSession(chatId, session.runId); + const newToken = await this.renewRunPatForSession(chatId, state.runId); if (newToken) { try { await sendChatMessages(newToken); @@ -536,32 +668,28 @@ export class TriggerChatTransport implements ChatTransport { } else if (isRunPatAuthError(err)) { throw err; } else { - previousRunId = session.runId; - this.sessions.delete(chatId); + // Unknown non-auth error: assume the run has ended and fall + // back to triggering a new one on the same session. + previousRunId = state.runId; + state.runId = undefined; this.coordinator?.release(chatId); - this.notifySessionChange(chatId, null); + this.notifySessionChange(chatId, state); isContinuation = true; } } if (inputSendOk) { - const currentSession = this.sessions.get(chatId); - if (!currentSession?.runId) { - throw new Error("TriggerChatTransport: session missing after input stream send"); - } - const activeStream = this.activeStreams.get(chatId); if (activeStream) { activeStream.abort(); this.activeStreams.delete(chatId); } - currentSession.isStreaming = true; - this.notifySessionChange(chatId, currentSession); + state.isStreaming = true; + this.notifySessionChange(chatId, state); - return this.subscribeToStream( - currentSession.runId, - currentSession.publicAccessToken, + return this.subscribeToSessionStream( + state, abortSignal, chatId, { upgradeRetry: { payload, messages } } @@ -569,19 +697,44 @@ export class TriggerChatTransport implements ChatTransport { } } - // First message or run has ended — trigger a new run - const triggerPayload = { + // First message or previous run ended — trigger a new run on the + // same session. The sessionId is threaded through the payload so + // the agent's runtime can `sessions.open(sessionId)` and attach to + // `.in` / `.out`. For the `triggerTask` callback path the server + // action creates the session and returns its id back to us. + const triggerPayload: Record = { ...payload, continuation: isContinuation, ...(previousRunId ? { previousRunId } : {}), }; + if (state?.sessionId) { + triggerPayload.sessionId = state.sessionId; + } - const { runId, publicAccessToken } = await this.triggerNewRun(chatId, triggerPayload, "trigger"); + const { + runId, + publicAccessToken, + sessionId: triggeredSessionId, + } = await this.triggerNewRun(chatId, triggerPayload, "trigger"); - const newSession: ChatSessionState = { runId, publicAccessToken, isStreaming: true }; - this.sessions.set(chatId, newSession); - this.notifySessionChange(chatId, newSession); - return this.subscribeToStream(runId, publicAccessToken, abortSignal, chatId, { + const resolvedSessionId = state?.sessionId ?? triggeredSessionId; + if (!resolvedSessionId) { + throw new Error( + "TriggerChatTransport: triggerTask did not return a sessionId. Use chat.createTriggerAction on the server to ensure the Session is created before the run is triggered." + ); + } + + const nextState: ChatSessionState = { + sessionId: resolvedSessionId, + runId, + publicAccessToken, + lastEventId: state?.lastEventId, + isStreaming: true, + skipToTurnComplete: state?.skipToTurnComplete, + }; + this.sessions.set(chatId, nextState); + this.notifySessionChange(chatId, nextState); + return this.subscribeToSessionStream(nextState, abortSignal, chatId, { upgradeRetry: { payload, messages }, }); }; @@ -607,8 +760,8 @@ export class TriggerChatTransport implements ChatTransport { message: UIMessage, metadata?: Record ): Promise => { - const session = this.sessions.get(chatId); - if (!session?.runId) return false; + const state = this.sessions.get(chatId); + if (!state?.runId) return false; const mergedMetadata = this.defaultMetadata || metadata @@ -624,15 +777,19 @@ export class TriggerChatTransport implements ChatTransport { const sendPending = async (token: string) => { const apiClient = new ApiClient(this.baseURL, token); - await apiClient.sendInputStream(session.runId, CHAT_MESSAGES_STREAM_ID, payload); + await apiClient.appendToSessionStream( + state.sessionId, + "in", + this.serializeInputChunk({ kind: "message", payload }) + ); }; try { - await sendPending(session.publicAccessToken); + await sendPending(state.publicAccessToken); return true; } catch (err) { if (isRunPatAuthError(err) && this.renewRunAccessToken) { - const newToken = await this.renewRunPatForSession(chatId, session.runId); + const newToken = await this.renewRunPatForSession(chatId, state.runId); if (newToken) { try { await sendPending(newToken); @@ -656,41 +813,32 @@ export class TriggerChatTransport implements ChatTransport { abortSignal?: AbortSignal | undefined; } & ChatRequestOptions ): Promise | null> => { - const session = this.sessions.get(options.chatId); - if (!session) { - return null; - } + const state = this.sessions.get(options.chatId); + if (!state) return null; // No active stream — the last turn completed before the page refreshed. // Return null so useChat settles into "ready" state instead of hanging. - if (!session.isStreaming) { - return null; - } + if (!state.isStreaming) return null; // Deduplicate: if there's already an active stream for this chatId, // return null so the second caller no-ops. - if (this.activeStreams.has(options.chatId)) { - return null; - } + if (this.activeStreams.has(options.chatId)) return null; const abortController = new AbortController(); this.activeStreams.set(options.chatId, abortController); // When the AI SDK (or caller) provides an abortSignal (e.g. from - // useChat's stop()), use it as the stream signal so stop sends - // the stop input stream signal to the backend. Fall back to the - // internal controller for stream lifecycle management. + // useChat's stop()), use it as the stream signal so stop sends a + // `{kind: "stop"}` chunk to the session. Fall back to the internal + // controller for stream lifecycle management. const abortSignal = options.abortSignal ? AbortSignal.any([options.abortSignal, abortController.signal]) : abortController.signal; - return this.subscribeToStream( - session.runId, - session.publicAccessToken, + return this.subscribeToSessionStream( + state, abortSignal, options.chatId, - // Send stop when the caller's signal fires (user-initiated stop). - // The internal abortController is only for stream management. { sendStopOnAbort: !!options.abortSignal } ); }; @@ -718,19 +866,23 @@ export class TriggerChatTransport implements ChatTransport { * ``` */ stopGeneration = async (chatId: string): Promise => { - const session = this.sessions.get(chatId); - if (!session?.runId) return false; + const state = this.sessions.get(chatId); + if (!state?.sessionId) return false; const sendStop = async (token: string) => { const api = new ApiClient(this.baseURL, token); - await api.sendInputStream(session.runId, CHAT_STOP_STREAM_ID, { stop: true }); + await api.appendToSessionStream( + state.sessionId, + "in", + this.serializeInputChunk({ kind: "stop" }) + ); }; try { - await sendStop(session.publicAccessToken); + await sendStop(state.publicAccessToken); } catch (err) { - if (isRunPatAuthError(err) && this.renewRunAccessToken) { - const newToken = await this.renewRunPatForSession(chatId, session.runId); + if (isRunPatAuthError(err) && this.renewRunAccessToken && state.runId) { + const newToken = await this.renewRunPatForSession(chatId, state.runId); if (newToken) { try { await sendStop(newToken); @@ -745,7 +897,7 @@ export class TriggerChatTransport implements ChatTransport { } } - session.skipToTurnComplete = true; + state.skipToTurnComplete = true; // Abort the active stream (if any) to close the SSE connection // and end the ReadableStream, causing useChat to finalize. @@ -785,9 +937,9 @@ export class TriggerChatTransport implements ChatTransport { this.coordinator.claim(chatId); } - const session = this.sessions.get(chatId); + const state = this.sessions.get(chatId); - if (session?.runId) { + if (state?.runId) { const mergedMetadata = this.defaultMetadata ?? undefined; const payload = { @@ -798,20 +950,18 @@ export class TriggerChatTransport implements ChatTransport { metadata: mergedMetadata, }; - const apiClient = new ApiClient(this.baseURL, session.publicAccessToken); + const apiClient = new ApiClient(this.baseURL, state.publicAccessToken); + + const body = this.serializeInputChunk({ kind: "message", payload }); try { - await apiClient.sendInputStream(session.runId, CHAT_MESSAGES_STREAM_ID, payload); + await apiClient.appendToSessionStream(state.sessionId, "in", body); } catch (err) { if (isRunPatAuthError(err) && this.renewRunAccessToken) { - const newToken = await this.renewRunPatForSession(chatId, session.runId); + const newToken = await this.renewRunPatForSession(chatId, state.runId); if (newToken) { const renewedClient = new ApiClient(this.baseURL, newToken); - await renewedClient.sendInputStream( - session.runId, - CHAT_MESSAGES_STREAM_ID, - payload - ); + await renewedClient.appendToSessionStream(state.sessionId, "in", body); } else { throw err; } @@ -820,12 +970,7 @@ export class TriggerChatTransport implements ChatTransport { } } - return this.subscribeToStream( - session.runId, - session.publicAccessToken, - undefined, - chatId - ); + return this.subscribeToSessionStream(state, undefined, chatId); } throw new Error(`No active session for chatId "${chatId}". Cannot send action.`); @@ -848,14 +993,23 @@ export class TriggerChatTransport implements ChatTransport { */ getSession = ( chatId: string - ): { runId: string; publicAccessToken: string; lastEventId?: string; isStreaming?: boolean } | undefined => { - const session = this.sessions.get(chatId); - if (!session) return undefined; + ): + | { + sessionId: string; + runId?: string; + publicAccessToken: string; + lastEventId?: string; + isStreaming?: boolean; + } + | undefined => { + const state = this.sessions.get(chatId); + if (!state) return undefined; return { - runId: session.runId, - publicAccessToken: session.publicAccessToken, - lastEventId: session.lastEventId, - isStreaming: session.isStreaming, + sessionId: state.sessionId, + runId: state.runId, + publicAccessToken: state.publicAccessToken, + lastEventId: state.lastEventId, + isStreaming: state.isStreaming, }; }; @@ -867,7 +1021,15 @@ export class TriggerChatTransport implements ChatTransport { callback: | (( chatId: string, - session: { runId: string; publicAccessToken: string; lastEventId?: string; isStreaming?: boolean } | null + session: + | { + sessionId: string; + runId?: string; + publicAccessToken: string; + lastEventId?: string; + isStreaming?: boolean; + } + | null ) => void) | undefined ): void { @@ -949,9 +1111,15 @@ export class TriggerChatTransport implements ChatTransport { */ setSession( chatId: string, - session: { runId: string; publicAccessToken: string; lastEventId?: string } + session: { + sessionId: string; + runId?: string; + publicAccessToken: string; + lastEventId?: string; + } ): void { this.sessions.set(chatId, { + sessionId: session.sessionId, runId: session.runId, publicAccessToken: session.publicAccessToken, lastEventId: session.lastEventId, @@ -975,7 +1143,7 @@ export class TriggerChatTransport implements ChatTransport { chatId: string, options?: { idleTimeoutInSeconds?: number; metadata?: Record } ): Promise { - // Don't preload if session already exists + // Don't preload if a run is already live for this chat. if (this.sessions.get(chatId)?.runId) return; // Deduplicate concurrent preload calls (e.g. React strict mode double-firing effects) @@ -983,6 +1151,8 @@ export class TriggerChatTransport implements ChatTransport { if (pending) return pending; const doPreload = async () => { + const state = await this.ensureSession(chatId); + const mergedMetadata = this.defaultMetadata || options?.metadata ? { ...(this.defaultMetadata ?? {}), ...(options?.metadata ?? {}) } @@ -991,6 +1161,7 @@ export class TriggerChatTransport implements ChatTransport { const payload = { messages: [] as never[], chatId, + sessionId: state.sessionId, trigger: "preload" as const, metadata: mergedMetadata, ...(options?.idleTimeoutInSeconds !== undefined @@ -1000,9 +1171,10 @@ export class TriggerChatTransport implements ChatTransport { const { runId, publicAccessToken } = await this.triggerNewRun(chatId, payload, "preload"); - const newSession: ChatSessionState = { runId, publicAccessToken }; - this.sessions.set(chatId, newSession); - this.notifySessionChange(chatId, newSession); + state.runId = runId; + state.publicAccessToken = publicAccessToken; + this.sessions.set(chatId, state); + this.notifySessionChange(chatId, state); }; const promise = doPreload().finally(() => { @@ -1028,13 +1200,17 @@ export class TriggerChatTransport implements ChatTransport { chatId: string, payload: Record, purpose: "trigger" | "preload" - ): Promise<{ runId: string; publicAccessToken: string }> { + ): Promise<{ runId: string; publicAccessToken: string; sessionId?: string }> { const autoTags = purpose === "preload" ? [`chat:${chatId}`, "preload:true"] : [`chat:${chatId}`]; const userTags = this.triggerOptions?.tags ?? []; const tags = [...autoTags, ...userTags].slice(0, 5); if (this.triggerTaskFn) { + // The server-side callback is expected to (a) create the Session + // if needed and (b) mint a token scoped to both the run and the + // session. It returns `sessionId` in its result when it created + // the session itself. return await this.triggerTaskFn({ payload: payload as TriggerChatTaskParams["payload"], options: { @@ -1068,6 +1244,9 @@ export class TriggerChatTransport implements ChatTransport { ? (triggerResponse as { publicAccessToken?: string }).publicAccessToken : undefined; + // In the direct `accessToken` path the session was created by + // `ensureSession` before this call, so the caller already has the + // sessionId — no need to surface it back. return { runId, publicAccessToken: publicAccessToken ?? currentToken }; } @@ -1075,6 +1254,7 @@ export class TriggerChatTransport implements ChatTransport { if (!this._onSessionChange) return; if (session) { this._onSessionChange(chatId, { + sessionId: session.sessionId, runId: session.runId, publicAccessToken: session.publicAccessToken, lastEventId: session.lastEventId, @@ -1110,9 +1290,8 @@ export class TriggerChatTransport implements ChatTransport { } } - private subscribeToStream( - runId: string, - accessToken: string, + private subscribeToSessionStream( + state: ChatSessionState, abortSignal: AbortSignal | undefined, chatId?: string, options?: { @@ -1124,9 +1303,8 @@ export class TriggerChatTransport implements ChatTransport { }; } ): ReadableStream { - // When resuming a run, skip past previously-seen events - // so we only receive the new turn's response. - const session = chatId ? this.sessions.get(chatId) : undefined; + const sessionId = state.sessionId; + const accessToken = state.publicAccessToken; // Create an internal AbortController so we can terminate the underlying // fetch connection when we're done reading (e.g. after intercepting the @@ -1137,16 +1315,22 @@ export class TriggerChatTransport implements ChatTransport { : internalAbort.signal; // When the caller aborts (user calls stop()), close the SSE connection. - // Only send a stop signal to the task if this is a user-initiated stop - // (sendStopOnAbort), not an internal stream management abort. + // Only send a stop chunk to the session if this is a user-initiated + // stop (sendStopOnAbort), not an internal stream management abort. if (abortSignal) { abortSignal.addEventListener( "abort", () => { - if (options?.sendStopOnAbort !== false && session) { - session.skipToTurnComplete = true; - const api = new ApiClient(this.baseURL, session.publicAccessToken); - api.sendInputStream(session.runId, CHAT_STOP_STREAM_ID, { stop: true }).catch(() => {}); // Best-effort + if (options?.sendStopOnAbort !== false) { + state.skipToTurnComplete = true; + const api = new ApiClient(this.baseURL, state.publicAccessToken); + api + .appendToSessionStream( + sessionId, + "in", + this.serializeInputChunk({ kind: "stop" }) + ) + .catch(() => {}); // Best-effort } internalAbort.abort(); }, @@ -1154,7 +1338,7 @@ export class TriggerChatTransport implements ChatTransport { ); } - const streamUrl = `${this.baseURL}/realtime/v1/streams/${runId}/${this.streamKey}`; + const streamUrl = `${this.baseURL}/realtime/v1/sessions/${encodeURIComponent(sessionId)}/out`; return new ReadableStream({ start: async (controller) => { @@ -1166,7 +1350,7 @@ export class TriggerChatTransport implements ChatTransport { }, signal: combinedSignal, timeoutInSeconds: this.streamTimeoutSeconds, - lastEventId: session?.lastEventId, + lastEventId: state.lastEventId, }); const sseStream = await subscription.subscribe(); const reader = sseStream.getReader(); @@ -1200,8 +1384,13 @@ export class TriggerChatTransport implements ChatTransport { reader = opened.reader; primed = opened.primed; } catch (e) { - if (isRunPatAuthError(e) && chatId && this.renewRunAccessToken) { - const newToken = await this.renewRunPatForSession(chatId, runId); + if ( + isRunPatAuthError(e) && + chatId && + this.renewRunAccessToken && + state.runId + ) { + const newToken = await this.renewRunPatForSession(chatId, state.runId); if (newToken) { const opened = await connectSseOnce(newToken); if (opened === null) { @@ -1247,123 +1436,136 @@ export class TriggerChatTransport implements ChatTransport { } // Track the last event ID so we can resume from here - if (value.id && session) { - session.lastEventId = value.id; + if (value.id) { + state.lastEventId = value.id; } - // Guard against heartbeat or malformed SSE events - if (value.chunk != null && typeof value.chunk === "object") { - const chunk = value.chunk as Record; - - // After a stop, skip leftover chunks from the stopped turn - // until we see the trigger:turn-complete marker. - if (session?.skipToTurnComplete) { - if (chunk.type === "trigger:turn-complete") { - session.skipToTurnComplete = false; - chunkCount = 0; + // Session SSE subscribers receive the raw record body as a + // string (the server wraps `{data, id}` for S2). Parse here + // so the rest of this loop can treat chunks as objects like + // the old run-scoped stream did. + let chunkObj: Record | null = null; + if (value.chunk != null) { + if (typeof value.chunk === "string") { + try { + chunkObj = JSON.parse(value.chunk) as Record; + } catch { + chunkObj = null; } + } else if (typeof value.chunk === "object") { + chunkObj = value.chunk as Record; + } + } + if (!chunkObj) continue; + + const chunk = chunkObj; + + // After a stop, skip leftover chunks from the stopped turn + // until we see the trigger:turn-complete marker. + if (state.skipToTurnComplete) { + if (chunk.type === "trigger:turn-complete") { + state.skipToTurnComplete = false; + chunkCount = 0; + } + continue; + } + + if (chunk.type === "trigger:upgrade-required" && chatId && options?.upgradeRetry) { + // Agent requested a version upgrade — re-trigger with the same + // message on the latest version and pipe the new stream through. + internalAbort.abort(); + const retryInfo = options.upgradeRetry; + const previousRunId = state.runId; + + // Keep the sessionId (sessions outlive runs) but clear the + // runId so triggerNewRun kicks a fresh run on the same session. + state.runId = undefined; + this.coordinator?.release(chatId); + this.notifySessionChange(chatId, state); + + try { + const triggerPayload = { + ...retryInfo.payload, + messages: retryInfo.messages, + sessionId, + continuation: true, + ...(previousRunId ? { previousRunId } : {}), + }; + + const { runId: newRunId, publicAccessToken: newToken } = + await this.triggerNewRun(chatId, triggerPayload, "trigger"); + + state.runId = newRunId; + state.publicAccessToken = newToken; + this.sessions.set(chatId, state); + this.notifySessionChange(chatId, state); + + // Subscribe to the new run's session stream and pipe through + const newStream = this.subscribeToSessionStream( + state, + abortSignal, + chatId + ); + const newReader = newStream.getReader(); + try { + while (true) { + const next = await newReader.read(); + if (next.done) break; + controller.enqueue(next.value); + } + } finally { + newReader.releaseLock(); + } + } catch (retryError) { + controller.error(retryError); + return; + } + try { + controller.close(); + } catch { + // Controller may already be closed + } + return; + } + + if (chunk.type === "trigger:turn-complete" && chatId) { + // Update token if a refreshed one was provided in the chunk + if (typeof chunk.publicAccessToken === "string") { + state.publicAccessToken = chunk.publicAccessToken; + } + // Mark streaming as complete so reconnectToStream doesn't + // hang on page refresh when no turn is in-flight. + state.isStreaming = false; + this.notifySessionChange(chatId, state); + + // Release multi-tab claim so other tabs can send + this.coordinator?.release(chatId); + + // Broadcast session to other tabs so they have the latest + // lastEventId (prevents replaying old SSE events on next send) + this.coordinator?.broadcastSession(chatId, { + lastEventId: state.lastEventId, + }); + + // Watch mode: keep the subscription open across turn + // boundaries so the consumer sees turn 2, 3, etc. through + // a single long-lived ReadableStream. Filter the control + // chunk and continue the read loop instead of closing. + if (this.watchMode) { continue; } - if (chunk.type === "trigger:upgrade-required" && chatId && options?.upgradeRetry) { - // Agent requested a version upgrade — re-trigger with the same - // message on the latest version and pipe the new stream through. - internalAbort.abort(); - const retryInfo = options.upgradeRetry; - const previousRunId = session?.runId; - - // Clear session so triggerNewRun creates a fresh one - this.sessions.delete(chatId); - this.coordinator?.release(chatId); - this.notifySessionChange(chatId, null); - - try { - const triggerPayload = { - ...retryInfo.payload, - messages: retryInfo.messages, - continuation: true, - ...(previousRunId ? { previousRunId } : {}), - }; - - const { runId: newRunId, publicAccessToken: newToken } = - await this.triggerNewRun(chatId, triggerPayload, "trigger"); - - const newSession: ChatSessionState = { runId: newRunId, publicAccessToken: newToken }; - this.sessions.set(chatId, newSession); - this.notifySessionChange(chatId, newSession); - - // Subscribe to the new run's stream and pipe through - const newStream = this.subscribeToStream( - newRunId, - newToken, - abortSignal, - chatId - ); - const newReader = newStream.getReader(); - try { - while (true) { - const next = await newReader.read(); - if (next.done) break; - controller.enqueue(next.value); - } - } finally { - newReader.releaseLock(); - } - } catch (retryError) { - controller.error(retryError); - return; - } - try { - controller.close(); - } catch { - // Controller may already be closed - } - return; + internalAbort.abort(); + try { + controller.close(); + } catch { + // Controller may already be closed } - - if (chunk.type === "trigger:turn-complete" && chatId) { - // Update token if a refreshed one was provided in the chunk - if (session && typeof chunk.publicAccessToken === "string") { - session.publicAccessToken = chunk.publicAccessToken; - } - // Mark streaming as complete so reconnectToStream doesn't - // hang on page refresh when no turn is in-flight. - if (session) { - session.isStreaming = false; - this.notifySessionChange(chatId, session); - } - - // Release multi-tab claim so other tabs can send - this.coordinator?.release(chatId); - - // Broadcast session to other tabs so they have the latest - // lastEventId (prevents replaying old SSE events on next send) - if (session) { - this.coordinator?.broadcastSession(chatId, { - lastEventId: session.lastEventId, - }); - } - - // Watch mode: keep the subscription open across turn - // boundaries so the consumer sees turn 2, 3, etc. through - // a single long-lived ReadableStream. Filter the control - // chunk and continue the read loop instead of closing. - if (this.watchMode) { - continue; - } - - internalAbort.abort(); - try { - controller.close(); - } catch { - // Controller may already be closed - } - return; - } - - chunkCount++; - controller.enqueue(chunk as unknown as UIMessageChunk); + return; } + + chunkCount++; + controller.enqueue(chunk as unknown as UIMessageChunk); } } catch (readError) { reader.releaseLock();