From 7ba0a66a382f5bbd08db111e6553631bc34f18f4 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Thu, 23 Apr 2026 13:46:47 +0100 Subject: [PATCH] fix(sdk,chat): route pipeChat through session.out + chat-agent smoke test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit pipeChat (the internal that auto-pipes a chat.agent's returned streamText result to the chat output) was still calling streams.pipe(CHAT_STREAM_KEY, stream) — a run-scoped run-streams path. After the session migration, the module-level facades (chatStream, messagesInput, stopInput) routed correctly, but pipeChat bypassed the facade and went straight to the old run-scoped pipe. Result: the turn-complete control chunk reached the session.out subscriber (written via chatStream.writer in writeTurnCompleteChunk) but every streamed UIMessageChunk from the LLM's turn was written to the dead run-scoped stream and never surfaced on session.out. Swap the pipe target to chatStream.pipe (the session-routed facade). The target / streamKey options still type-check for API parity but are no longer meaningful — sessions are the address, and sub-agents that need to write to a parent chat open the parent's Session explicitly. Smoke now catches all 14 chunks (start / start-step / text-start / 7x text-delta / text-end / finish-step / finish / trigger:turn-complete) with ids 0 through 13 from session.out, match: true. Also adds references/hello-world/src/trigger/chatAgentSmoke.ts — end-to-end validation: - sessions.create with externalId = chatId - trigger test-agent with {chatId, sessionId, messages, …} - handle.out.read(...) SSE subscribe, capture chunks by id+type - sessions.close on completion Triggered from the dashboard or MCP as chat-agent-smoke. Requires OPENAI_API_KEY in the dev env (the test-agent uses openai:gpt-4o-mini). --- packages/trigger-sdk/src/v3/ai.ts | 20 +++- .../hello-world/src/trigger/chatAgentSmoke.ts | 111 ++++++++++++++++++ 2 files changed, 125 insertions(+), 6 deletions(-) create mode 100644 references/hello-world/src/trigger/chatAgentSmoke.ts diff --git a/packages/trigger-sdk/src/v3/ai.ts b/packages/trigger-sdk/src/v3/ai.ts index cd33c5c82..ddf5d091a 100644 --- a/packages/trigger-sdk/src/v3/ai.ts +++ b/packages/trigger-sdk/src/v3/ai.ts @@ -2687,7 +2687,6 @@ async function pipeChat( options?: PipeChatOptions ): Promise { locals.set(chatPipeCountKey, (locals.get(chatPipeCountKey) ?? 0) + 1); - const streamKey = options?.streamKey ?? CHAT_STREAM_KEY; let stream: AsyncIterable | ReadableStream; @@ -2702,18 +2701,27 @@ async function pipeChat( ); } - const pipeOptions: PipeStreamOptions = {}; + const pipeOptions: SessionPipeStreamOptions = {}; if (options?.signal) { pipeOptions.signal = options.signal; } - if (options?.target) { - pipeOptions.target = options.target; - } if (options?.spanName) { pipeOptions.spanName = options.spanName; } + // `options.target` / `options.streamKey` are accepted for API parity + // with the pre-migration run-scoped pipe but no longer have meaning — + // sessions are the address (single stream per session, no sub-run + // targeting). Sub-agents that need to write into a parent's chat now + // open that session explicitly via `sessions.open(parentSessionId).out.pipe`. - const { waitUntilComplete } = streams.pipe(streamKey, stream, pipeOptions); + // The generic is typed for `UIMessageChunk`, but `pipeChat` also + // accepts opaque UIMessageStreamable / raw iterables whose element + // type we don't know at compile time. Cast — runtime behaviour is + // identical (bytes go to session.out either way). + const { waitUntilComplete } = chatStream.pipe( + stream as ReadableStream | AsyncIterable, + pipeOptions + ); await waitUntilComplete(); } diff --git a/references/hello-world/src/trigger/chatAgentSmoke.ts b/references/hello-world/src/trigger/chatAgentSmoke.ts new file mode 100644 index 000000000..84b043e13 --- /dev/null +++ b/references/hello-world/src/trigger/chatAgentSmoke.ts @@ -0,0 +1,111 @@ +import { logger, sessions, task, tasks } from "@trigger.dev/sdk"; + +/** + * End-to-end smoke test for the chat.agent -> Sessions migration. + * + * Flow: + * 1. Create a Session with a deterministic externalId so the + * `test-agent` run can `sessions.open(...)` it on startup. + * 2. Trigger `test-agent` with `{chatId, sessionId, messages, trigger, + * metadata}` — mirrors what TriggerChatTransport would send for a + * first message, minus the browser-triggered access token layer. + * 3. `session.out.read({...})` — consume the agent's UIMessageChunks + * as they stream out. Bail after the first text-delta (good + * enough to prove output flow + SSE subscription). + * 4. `sessions.close(...)` — tidy up. + * + * Trigger from the dashboard or MCP: + * + * mcp__trigger__trigger_task(taskId: "chat-agent-smoke", payload: {}) + * + * Expects OPENAI_API_KEY set in the env (the test-agent uses + * `openai:gpt-4o-mini`). If the key is missing the smoke reports an + * error payload without crashing. + */ +export const chatAgentSmoke = task({ + id: "chat-agent-smoke", + run: async () => { + const stamp = Date.now(); + const chatId = `chat-agent-smoke-${stamp}`; + + logger.info("creating chat.agent backing session", { externalId: chatId }); + const session = await sessions.create({ + type: "chat.agent", + externalId: chatId, + tags: ["chat-agent-smoke"], + }); + + logger.info("triggering test-agent run", { + chatId, + sessionId: session.id, + }); + await tasks.trigger("test-agent", { + chatId, + sessionId: session.id, + trigger: "submit-message", + messages: [ + { + id: `m-${stamp}`, + role: "user", + parts: [{ type: "text", text: "Say hello in five words." }], + }, + ], + metadata: { userId: "smoke", model: "openai:gpt-4o-mini" }, + }); + + logger.info("subscribing to session.out, waiting for first chunks"); + const handle = sessions.open(session.id); + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), 60_000); + + const received: Array<{ type?: string; id?: string }> = []; + let firstTextDelta: string | undefined; + let turnCompleteSeen = false; + + try { + const stream = await handle.out.read>({ + signal: controller.signal, + timeoutInSeconds: 30, + // Start from seq 0 so we don't race the agent's early writes. + lastEventId: "-1", + onPart: (part) => { + // Record the event id alongside the chunk so we can see the + // full sequence that came down the wire. + received.push({ + id: part.id, + type: (part.chunk as { type?: string } | null)?.type, + }); + }, + }); + + for await (const chunk of stream) { + if (chunk.type === "text-delta" && typeof chunk.delta === "string") { + firstTextDelta ??= chunk.delta; + } + if (chunk.type === "trigger:turn-complete") { + turnCompleteSeen = true; + break; + } + if (received.length > 500) break; + } + } catch (err) { + if ((err as Error).name !== "AbortError") throw err; + } finally { + clearTimeout(timeout); + } + + await sessions.close(session.id, { reason: "chat-agent-smoke-done" }); + + return { + ok: received.length > 0, + chatId, + sessionId: session.id, + chunkCount: received.length, + firstTextDelta, + turnCompleteSeen, + types: [...new Set(received.map((c) => c.type ?? ""))], + firstFiveIds: received.slice(0, 5).map((c) => `${c.id}:${c.type ?? ""}`), + lastFiveIds: received.slice(-5).map((c) => `${c.id}:${c.type ?? ""}`), + }; + }, +});