diff --git a/.server-changes/agent-view-sessions.md b/.server-changes/agent-view-sessions.md new file mode 100644 index 000000000..757dcdc2f --- /dev/null +++ b/.server-changes/agent-view-sessions.md @@ -0,0 +1,12 @@ +--- +area: webapp +type: improvement +--- + +Migrate the dashboard Agent tab (span inspector) to subscribe to the backing Session's `.out` and `.in` channels instead of the run-scoped chat output + chat-messages input streams. Pairs with the SDK + MCP migrations on the ai-chat branch. + +- `SpanPresenter.server.ts` extracts `agentSession` from the run payload (prefers `sessionId`, falls back to `chatId` for pre-Sessions agent runs — matches `resolveSessionByIdOrExternalId`). +- Span route threads `agentSession` through `AgentViewAuth` and gates `agentView` creation on having one. +- New dashboard resource route `resources.orgs.../runs.$runParam/realtime/v1/sessions/$sessionId/$io` proxies `S2RealtimeStreams.streamResponseFromSessionStream` under dashboard session auth. The run param binds resource hierarchy; the session identity is verified against the environment. +- `AgentView.tsx` subscribes to `/out` and `/in` URLs, drops local `CHAT_STREAM_KEY`/`CHAT_MESSAGES_STREAM_ID` constants, and parses the `.in` stream as `ChatInputChunk` (`{kind: "message", payload}` for user turns; `{kind: "stop"}` ignored). Output-stream parsing is unchanged — session v2 SSE already delivers UIMessageChunk objects from `record.body.data`. +- Smoke: opened a prior `test-agent` run in the dashboard, Agent tab rendered user + assistant messages end-to-end with zero console errors. Both SSE endpoints (`/out`, `/in`) returned 200. diff --git a/apps/webapp/app/components/runs/v3/agent/AgentView.tsx b/apps/webapp/app/components/runs/v3/agent/AgentView.tsx index fd3dfb5cb..c54904d5b 100644 --- a/apps/webapp/app/components/runs/v3/agent/AgentView.tsx +++ b/apps/webapp/app/components/runs/v3/agent/AgentView.tsx @@ -12,11 +12,19 @@ import { useProject } from "~/hooks/useProject"; export type AgentViewAuth = { publicAccessToken: string; apiOrigin: string; + /** + * Session identifier the AgentView uses to address the backing + * {@link Session} when subscribing to `.in` / `.out`. Accepts either + * a `session_*` friendlyId or the transport-supplied externalId + * (typically the browser's `chatId`) — the dashboard resource route + * resolves either form via `resolveSessionByIdOrExternalId`. + */ + sessionId: string; /** * User messages extracted from the run's task payload at load time. - * Empty array for runs started with `trigger: "preload"` — in that case - * the first user message will arrive over the chat-messages input stream - * and get merged in by the AgentView subscription. + * Empty array for runs started with `trigger: "preload"` — in that + * case the first user message arrives over the session's `.in` + * channel and is merged in by the AgentView subscription. */ initialMessages: UIMessage[]; }; @@ -26,12 +34,6 @@ type AgentViewRun = { taskIdentifier: string; }; -// Default stream IDs for Trigger.dev chat tasks — kept as literals so we -// don't pull server-only constants from `@trigger.dev/core/v3/chat-client` -// into a browser bundle. -const CHAT_STREAM_KEY = "chat"; -const CHAT_MESSAGES_STREAM_ID = "chat-messages"; - /** * Max state-update interval while assistant chunks are streaming. Matches * the `experimental_throttle: 100` we previously passed to `useChat`. @@ -51,20 +53,23 @@ const INITIAL_PAYLOAD_TIMESTAMP = 0; /** * Renders an agent run's chat conversation as it unfolds. * - * Subscribes to two separate realtime streams for the run: - * - The **chat output stream** delivers assistant `UIMessageChunk`s (text - * deltas, tool calls, reasoning, etc.) produced by `pipeChat` on the - * task side. - * - The **chat-messages input stream** delivers user messages sent to the - * task via `sendInputStream` — each chunk carries a `ChatTaskWirePayload` - * with the most recent `messages` array. + * Subscribes to both channels of the run's backing {@link Session}: + * - **`.out`** delivers assistant `UIMessageChunk`s (text deltas, tool + * calls, reasoning, etc.) produced by the agent's + * `chatStream.writer(...)` calls — objects, already parsed by the S2 + * SSE reader. + * - **`.in`** delivers {@link ChatInputChunk}s sent by + * {@link TriggerChatTransport} (or any other session writer). Each + * chunk is a tagged union (`{kind: "message", payload}` for user + * turns, `{kind: "stop"}` for stop signals) — the AgentView only + * cares about `kind: "message"` and pulls `.payload.messages`. * * Both streams are read directly via `SSEStreamSubscription` through the * dashboard's session-authed resource routes — not through `useChat` or * `TriggerChatTransport`. This gives us per-chunk server-side timestamps - * (Redis stream IDs) from both streams, which we use to produce a + * (S2 sequence numbers) from both streams, which we use to produce a * chronologically correct merged message list that works for replays, - * multi-message turns, and steering messages. + * multi-message turns, cross-run session resumes, and steering messages. * * Intended to be mounted inside a scrollable container — the component * does not own its own scrollbar. @@ -82,6 +87,7 @@ export function AgentView({ const messages = useAgentRunMessages({ runFriendlyId: run.friendlyId, + sessionId: agentView.sessionId, apiOrigin: agentView.apiOrigin, orgSlug: organization.slug, projectSlug: project.slug, @@ -119,14 +125,24 @@ export function AgentView({ // --------------------------------------------------------------------------- /** - * Shape of each chunk on the chat-messages input stream. Each chunk is a - * `ChatTaskWirePayload` whose `messages` field holds either the latest user - * message (for `submit-message`) or the full history (for - * `regenerate-message`). We dedupe by ID either way. + * Shape of each chunk on the session's `.in` channel. Mirrors the + * `ChatInputChunk` tagged union produced by {@link TriggerChatTransport}: + * - `kind: "message"` carries a `ChatTaskWirePayload` in `.payload` + * (user-submitted messages or regenerate calls); we dedupe by id. + * - `kind: "stop"` is a stop signal — no messages, nothing to render + * here, so it's filtered. + * + * The server wraps records in `{data, id}` and writes `data` as a JSON + * string; SSE v2 delivers the parsed string back. {@link parseChunkPayload} + * re-parses to recover the object. */ type InputStreamChunk = { - messages?: Array<{ id?: string; role?: string; parts?: unknown[] }>; - trigger?: string; + kind?: "message" | "stop"; + payload?: { + messages?: Array<{ id?: string; role?: string; parts?: unknown[] }>; + trigger?: string; + }; + message?: string; }; /** @@ -173,17 +189,17 @@ type MessageOrchestrationState = { /** * `SSEStreamSubscription`'s v2 batch path delivers `parsedBody.data` as-is - * — but for streams written via `sendInputStream` (which stores the user - * payload as a JSON string in the record body), `data` is itself a string - * that needs a second `JSON.parse` to recover the actual object. This - * happens for the chat-messages input stream because the action handler - * does `JSON.stringify(body.data.data)` before storing. + * — but session channels diverge by direction: * - * Output streams from `pipeChat` write objects directly, so the v2 path - * delivers them already-parsed. Either way this helper accepts both shapes - * defensively: a string is parsed; an object is returned as-is. + * - `.in`: {@link TriggerChatTransport.serializeInputChunk} writes the + * `ChatInputChunk` as a JSON **string**, so `data` is a string that + * needs a second `JSON.parse` to recover the tagged union. + * - `.out`: the agent's `chatStream.writer(...)` writes + * {@link UIMessageChunk} **objects** directly; `data` arrives + * already-parsed. * - * Returns `null` for unparseable / unexpected payloads. + * This helper accepts both shapes defensively: a string is parsed; an + * object is returned as-is. Returns `null` for unparseable payloads. */ function parseChunkPayload(raw: unknown): Record | null { if (raw == null) return null; @@ -208,6 +224,7 @@ function createOrchestrationState(): MessageOrchestrationState { function useAgentRunMessages({ runFriendlyId, + sessionId, apiOrigin, orgSlug, projectSlug, @@ -215,6 +232,7 @@ function useAgentRunMessages({ initialMessages, }: { runFriendlyId: string; + sessionId: string; apiOrigin: string; orgSlug: string; projectSlug: string; @@ -268,13 +286,13 @@ function useAgentRunMessages({ useEffect(() => { const abort = new AbortController(); - const outputUrl = + const encodedSession = encodeURIComponent(sessionId); + const sessionBase = `${apiOrigin}/resources/orgs/${orgSlug}/projects/${projectSlug}/env/${envSlug}` + - `/runs/${runFriendlyId}/realtime/v1/streams/${runFriendlyId}/${CHAT_STREAM_KEY}`; + `/runs/${runFriendlyId}/realtime/v1/sessions/${encodedSession}`; - const inputUrl = - `${apiOrigin}/resources/orgs/${orgSlug}/projects/${projectSlug}/env/${envSlug}` + - `/runs/${runFriendlyId}/realtime/v1/streams/${runFriendlyId}/input/${CHAT_MESSAGES_STREAM_ID}`; + const outputUrl = `${sessionBase}/out`; + const inputUrl = `${sessionBase}/in`; const commonSubOptions = { signal: abort.signal, @@ -385,7 +403,12 @@ function useAgentRunMessages({ } }; - // ---- Input stream: user messages --------------------------------------- + // ---- Input channel: user messages (`ChatInputChunk`) ------------------- + // + // The transport appends a `{kind: "message", payload}` ChatInputChunk + // for every user turn (and `{kind: "stop"}` for stop signals). We pull + // user messages out of `payload.messages` for `kind: "message"` chunks + // and ignore the rest. const runInput = async () => { try { const sub = new SSEStreamSubscription(inputUrl, commonSubOptions); @@ -397,9 +420,11 @@ function useAgentRunMessages({ if (done) return; const chunk = parseChunkPayload(value.chunk) as InputStreamChunk | null; - if (!chunk || !Array.isArray(chunk.messages)) continue; + if (!chunk || chunk.kind !== "message") continue; + const payload = chunk.payload; + if (!payload || !Array.isArray(payload.messages)) continue; - const incomingUsers = chunk.messages.filter( + const incomingUsers = payload.messages.filter( (m): m is UIMessage => m != null && (m as { role?: string }).role === "user" && typeof m.id === "string" ); @@ -438,7 +463,7 @@ function useAgentRunMessages({ pendingTimerRef.current = null; } }; - }, [runFriendlyId, apiOrigin, orgSlug, projectSlug, envSlug]); + }, [runFriendlyId, sessionId, apiOrigin, orgSlug, projectSlug, envSlug]); return useMemo(() => { const timestamps = timestampsRef.current; diff --git a/apps/webapp/app/presenters/v3/SpanPresenter.server.ts b/apps/webapp/app/presenters/v3/SpanPresenter.server.ts index 6ba1ead0e..b186848f2 100644 --- a/apps/webapp/app/presenters/v3/SpanPresenter.server.ts +++ b/apps/webapp/app/presenters/v3/SpanPresenter.server.ts @@ -260,24 +260,43 @@ export class SpanPresenter extends BasePresenter { const taskKind = RunAnnotations.safeParse(run.annotations).data?.taskKind; const isAgentRun = taskKind === "AGENT"; - // For agent runs, extract the initial user messages that were supplied - // via the task payload (from the original `triggerTask({ payload: { messages: [...] } })` - // call). When the run was started with `trigger: "preload"`, this array - // will be empty — in that case the first user message arrives later via - // the chat-messages input stream and is picked up by the AgentView. + // For agent runs, extract the initial user messages + the backing + // Session handle from the task payload (from the original + // `triggerTask({ payload: { messages, sessionId, chatId, ... } })` + // call). When the run was started with `trigger: "preload"`, + // `messages` is empty — the first user message arrives later over + // the session `.in` channel and is merged in by the AgentView. + // + // `agentSession` is the identifier the dashboard uses to address the + // backing Session when subscribing to `.out` / `.in`. Prefer the + // explicit `sessionId` threaded by `TriggerChatTransport` / + // `chat.createTriggerAction`; fall back to `chatId` for pre-migration + // agent runs (the session resource route accepts either, matching + // `resolveSessionByIdOrExternalId`). let agentInitialMessages: AgentInitialMessage[] = []; + let agentSession: string | null = null; if (isAgentRun && run.payload && run.payloadType !== "application/store") { try { const parsed = await parsePacket({ data: typeof run.payload === "string" ? run.payload : JSON.stringify(run.payload), dataType: run.payloadType ?? "application/json", }); - if (parsed && typeof parsed === "object" && Array.isArray((parsed as any).messages)) { - agentInitialMessages = (parsed as any).messages as AgentInitialMessage[]; + if (parsed && typeof parsed === "object") { + if (Array.isArray((parsed as any).messages)) { + agentInitialMessages = (parsed as any).messages as AgentInitialMessage[]; + } + const sessionId = (parsed as any).sessionId; + const chatId = (parsed as any).chatId; + if (typeof sessionId === "string" && sessionId.length > 0) { + agentSession = sessionId; + } else if (typeof chatId === "string" && chatId.length > 0) { + agentSession = chatId; + } } } catch { - // Fall back to an empty initial message list — the AgentView will - // render whatever arrives over the input/output streams. + // Fall back to empty initial messages + null session — the + // AgentView will show a loading spinner and surface any stream + // subscription errors to the console. } } @@ -340,6 +359,7 @@ export class SpanPresenter extends BasePresenter { isError: isFailedRunStatus(run.status), isAgentRun, agentInitialMessages, + agentSession, payload, payloadType: run.payloadType, output, diff --git a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.env.$envParam.playground.$agentParam/route.tsx b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.env.$envParam.playground.$agentParam/route.tsx index aa6e30be5..1fe4dc657 100644 --- a/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.env.$envParam.playground.$agentParam/route.tsx +++ b/apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.env.$envParam.playground.$agentParam/route.tsx @@ -720,7 +720,15 @@ function PlaygroundSidebar({ onRegionChange: (val: string | undefined) => void; regions: Array<{ id: string; name: string; description?: string; isDefault: boolean }>; isDev: boolean; - session: { runId: string; publicAccessToken: string; lastEventId?: string } | undefined; + session: + | { + sessionId: string; + runId?: string; + publicAccessToken: string; + lastEventId?: string; + isStreaming?: boolean; + } + | undefined; messageCount: number; isStreaming: boolean; status: string; diff --git a/apps/webapp/app/routes/resources.orgs.$organizationSlug.projects.$projectParam.env.$envParam.runs.$runParam.realtime.v1.sessions.$sessionId.$io.ts b/apps/webapp/app/routes/resources.orgs.$organizationSlug.projects.$projectParam.env.$envParam.runs.$runParam.realtime.v1.sessions.$sessionId.$io.ts new file mode 100644 index 000000000..93047d220 --- /dev/null +++ b/apps/webapp/app/routes/resources.orgs.$organizationSlug.projects.$projectParam.env.$envParam.runs.$runParam.realtime.v1.sessions.$sessionId.$io.ts @@ -0,0 +1,98 @@ +import { type LoaderFunctionArgs } from "@remix-run/server-runtime"; +import { z } from "zod"; +import { $replica } from "~/db.server"; +import { findProjectBySlug } from "~/models/project.server"; +import { findEnvironmentBySlug } from "~/models/runtimeEnvironment.server"; +import { getRequestAbortSignal } from "~/services/httpAsyncStorage.server"; +import { S2RealtimeStreams } from "~/services/realtime/s2realtimeStreams.server"; +import { resolveSessionByIdOrExternalId } from "~/services/realtime/sessions.server"; +import { getRealtimeStreamInstance } from "~/services/realtime/v1StreamsGlobal.server"; +import { requireUserId } from "~/services/session.server"; +import { EnvironmentParamSchema } from "~/utils/pathBuilder"; + +const ParamsSchema = z.object({ + runParam: z.string(), + sessionId: z.string(), + io: z.enum(["out", "in"]), +}); + +// GET: SSE stream subscription for a backing Session's `.out` / `.in` +// channel. Dashboard-auth counterpart to the public API's +// `/realtime/v1/sessions/:sessionId/:io` endpoint. Used by the Agent tab +// in the span inspector to observe assistant chunks (`.out`) and +// user-side ChatInputChunk payloads (`.in`) for a chat.agent run. +// +// The `:sessionId` segment accepts either the `session_*` friendlyId or +// the externalId the transport registered for the chat (typically the +// browser's `chatId`). Runs pre-dating the Sessions migration that have +// `chatId` but no `sessionId` in the payload take the externalId path. +// +// Authenticated by the dashboard session — the user must have access to +// the project, environment, and run. The run binds this resource +// hierarchy; the session identity is verified against the environment. +export async function loader({ request, params }: LoaderFunctionArgs) { + const userId = await requireUserId(request); + const { organizationSlug, projectParam, envParam } = EnvironmentParamSchema.parse(params); + const { runParam, sessionId, io } = ParamsSchema.parse(params); + + const project = await findProjectBySlug(organizationSlug, projectParam, userId); + if (!project) { + return new Response("Project not found", { status: 404 }); + } + + const environment = await findEnvironmentBySlug(project.id, envParam, userId); + if (!environment) { + return new Response("Environment not found", { status: 404 }); + } + + // Verify the run lives in this environment — keeps callers from + // subscribing to arbitrary sessions via `/runs/$runParam/...`. + const run = await $replica.taskRun.findFirst({ + where: { + friendlyId: runParam, + runtimeEnvironmentId: environment.id, + }, + select: { id: true, friendlyId: true }, + }); + + if (!run) { + return new Response("Run not found", { status: 404 }); + } + + const session = await resolveSessionByIdOrExternalId( + $replica, + environment.id, + sessionId + ); + + if (!session) { + return new Response("Session not found", { status: 404 }); + } + + const realtimeStream = getRealtimeStreamInstance(environment, "v2"); + + if (!(realtimeStream instanceof S2RealtimeStreams)) { + return new Response("Session channels require the S2 realtime backend", { + status: 501, + }); + } + + const lastEventId = request.headers.get("Last-Event-ID") || undefined; + const timeoutInSecondsRaw = request.headers.get("Timeout-Seconds") ?? undefined; + const timeoutInSeconds = timeoutInSecondsRaw ? parseInt(timeoutInSecondsRaw) : undefined; + + if ( + timeoutInSeconds && + (isNaN(timeoutInSeconds) || timeoutInSeconds < 1 || timeoutInSeconds > 600) + ) { + return new Response("Invalid timeout", { status: 400 }); + } + + return realtimeStream.streamResponseFromSessionStream( + request, + session.friendlyId, + io, + getRequestAbortSignal(), + { lastEventId, timeoutInSeconds } + ); +} diff --git a/apps/webapp/app/routes/resources.orgs.$organizationSlug.projects.$projectParam.env.$envParam.runs.$runParam.spans.$spanParam/route.tsx b/apps/webapp/app/routes/resources.orgs.$organizationSlug.projects.$projectParam.env.$envParam.runs.$runParam.spans.$spanParam/route.tsx index e495ac8c3..263921733 100644 --- a/apps/webapp/app/routes/resources.orgs.$organizationSlug.projects.$projectParam.env.$envParam.runs.$runParam.spans.$spanParam/route.tsx +++ b/apps/webapp/app/routes/resources.orgs.$organizationSlug.projects.$projectParam.env.$envParam.runs.$runParam.spans.$spanParam/route.tsx @@ -143,12 +143,16 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => { } // For agent runs, mint a read-only run-scoped token so the Agent tab - // can subscribe to the run's chat output stream from the browser. We - // also forward the initial user messages extracted from the task - // payload so the AgentView can seed its merged message list (empty for - // runs started via `trigger: "preload"`). + // can subscribe to the run's backing Session from the browser. We + // also forward the initial user messages + the session identifier + // extracted from the task payload — the AgentView uses the session + // to subscribe to `.in` / `.out` (replaces the old run-scoped + // chat-messages + chat streams). Runs without an identifiable + // session (misformed payload, legacy pre-chat-agent runs) get + // `agentSession: null`; the AgentView renders a loading spinner + // without subscribing. let agentView: AgentViewAuth | null = null; - if (result.type === "run" && result.run.isAgentRun) { + if (result.type === "run" && result.run.isAgentRun && result.run.agentSession) { const project = await findProjectBySlug(organizationSlug, projectParam, userId); const environment = project ? await findEnvironmentBySlug(project.id, envParam, userId) @@ -158,6 +162,7 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => { agentView = { publicAccessToken, apiOrigin: env.API_ORIGIN || env.LOGIN_ORIGIN, + sessionId: result.run.agentSession, initialMessages: (result.run.agentInitialMessages ?? []) as AgentViewAuth["initialMessages"], }; } diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index b054d66bb..4c1dde06e 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -655,7 +655,7 @@ describe("TriggerChatTransport", () => { it("should set isStreaming to false via onSessionChange when turn completes", async () => { const sessionChanges: Array<{ chatId: string; - session: { isStreaming?: boolean } | null; + session: { isStreaming?: boolean; runId?: string; sessionId?: string } | null; }> = []; global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 91aa9587d..ef85d71e6 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -195,7 +195,14 @@ type TriggerChatTransportOptionsBase = { sessions?: Record< string, { - sessionId: string; + /** + * Optional. If omitted, the transport upserts the backing + * Session on first use via `POST /api/v1/sessions` (keyed on + * `chatId` as `externalId`). Pre-Sessions persisted state + * won't have a sessionId — this lets old localStorage records + * hydrate without migration. + */ + sessionId?: string; runId?: string; publicAccessToken: string; lastEventId?: string;