diff --git a/.changeset/686-stream-generated-flight.md b/.changeset/686-stream-generated-flight.md new file mode 100644 index 000000000..547e7fb65 --- /dev/null +++ b/.changeset/686-stream-generated-flight.md @@ -0,0 +1,5 @@ +--- +'agent-bundle': patch +--- + +Stream Flight bytes from generated MCP and event-route workers as they render instead of buffering the whole document: a `Suspense` fallback authored in a compiled tool or event route now reaches the stdio MCP client (`notifications/progress`), standalone hook wrappers, and the Workbench's production route surface before the suspended child resolves, and a consumer that stops reading releases the worker render instead of crashing the host on a late chunk. (#718) diff --git a/packages/agent-bundle/src/adapters/hook-contract.ts b/packages/agent-bundle/src/adapters/hook-contract.ts index f09158a66..7ce358fcb 100644 --- a/packages/agent-bundle/src/adapters/hook-contract.ts +++ b/packages/agent-bundle/src/adapters/hook-contract.ts @@ -734,40 +734,41 @@ const eventRouteHookWrapperSource = ( ' execute: async (dispatch) => {', ' const context = await agent();', ' const id = ++sequence;', - ' return new Promise((resolvePromise, rejectPromise) => {', - ' const abort = () => { worker.postMessage({ id, type: "cancel" }); rejectPromise(new DOMException("Agent render was aborted", "AbortError")); };', - ' const cleanup = () => { dispatch.signal.removeEventListener("abort", abort); worker.off("error", onError); worker.off("exit", onExit); worker.off("message", onMessage); };', - ' const reject = (error) => { cleanup(); rejectPromise(error); };', - ' const onError = (error) => { reject(error); };', - ' const onExit = (code) => { reject(new Error(`Generated hook Flight worker exited with code ${String(code)}.`)); };', - ' const onMessage = (message) => {', - ' if (message.id !== id) return;', - ' if (message.type === "progress") { Promise.resolve(dispatch.progress?.report(message.update)).catch(reject); return; }', - ' if (message.type === "error") { reject(new Error(message.message)); return; }', - ' if (message.type !== "complete") return;', - ' cleanup();', - ' resolvePromise(new ReadableStream({ start(controller) { controller.enqueue(message.bytes); controller.close(); } }));', - ' };', - ' worker.on("error", onError);', - ' worker.on("exit", onExit);', - ' worker.on("message", onMessage);', - ' dispatch.signal.addEventListener("abort", abort, { once: true });', - ' if (dispatch.signal.aborted) { abort(); return; }', - ' worker.postMessage({', - ' actor: context.actor,', - ' artifactEpoch: flightArtifactEpoch,', - ' host: context.host,', - ' id,', - ' invocation: dispatch.invocation,', - ' lineage: context.lineage,', - ' plugin: context.plugin,', - ' requestInvocation: context.invocation,', - ' session: context.session,', - ' terminal: context.terminal,', - ' type: "render",', - ' workspace: context.workspace,', - ' });', + ' let controller;', + ' const cleanup = () => { dispatch.signal.removeEventListener("abort", abort); worker.off("error", fail); worker.off("exit", onExit); worker.off("message", onMessage); };', + ' const stream = new ReadableStream({ cancel() { worker.postMessage({ id, type: "cancel" }); cleanup(); }, start(opened) { controller = opened; } });', + ' const fail = (error) => { cleanup(); controller.error(error); };', + ' const abort = () => { worker.postMessage({ id, type: "cancel" }); fail(new DOMException("Agent render was aborted", "AbortError")); };', + ' const onExit = (code) => { fail(new Error(`Generated hook Flight worker exited with code ${String(code)}.`)); };', + ' const onMessage = (message) => {', + ' if (message.id !== id) return;', + ' if (message.type === "progress") { Promise.resolve(dispatch.progress?.report(message.update)).catch(fail); return; }', + ' if (message.type === "chunk") { controller.enqueue(message.bytes); return; }', + ' if (message.type === "error") { fail(new Error(message.message)); return; }', + ' if (message.type !== "end") return;', + ' cleanup();', + ' controller.close();', + ' };', + ' if (dispatch.signal.aborted) { controller.error(new DOMException("Agent render was aborted", "AbortError")); return stream; }', + ' worker.on("error", fail);', + ' worker.on("exit", onExit);', + ' worker.on("message", onMessage);', + ' dispatch.signal.addEventListener("abort", abort, { once: true });', + ' worker.postMessage({', + ' actor: context.actor,', + ' artifactEpoch: flightArtifactEpoch,', + ' host: context.host,', + ' id,', + ' invocation: dispatch.invocation,', + ' lineage: context.lineage,', + ' plugin: context.plugin,', + ' requestInvocation: context.invocation,', + ' session: context.session,', + ' terminal: context.terminal,', + ' type: "render",', + ' workspace: context.workspace,', ' });', + ' return stream;', ' },', ' });', ' try {', diff --git a/packages/agent-bundle/src/build/entry-shell.ts b/packages/agent-bundle/src/build/entry-shell.ts index b0ae6a780..00f34d58d 100644 --- a/packages/agent-bundle/src/build/entry-shell.ts +++ b/packages/agent-bundle/src/build/entry-shell.ts @@ -1143,7 +1143,7 @@ export const generatedRouteFlightWorkerSource = (options: GeneratedRouteFlightWo ? [] : [' const bindings = await runtimeState.requestBindings({ signal: controller.signal });', ' try {']), ' const plugin = message.plugin ?? pluginRoot.identity;', - ' const bytes = await runAgentRequest({', + ' await runAgentRequest({', ' ...(message.actor === undefined ? {} : { actor: message.actor }),', ' ...(message.host === undefined ? {} : { host: message.host }),', ' invocation: { ...message.requestInvocation, artifactEpoch: ARTIFACT_EPOCH, kind: message.invocation.kind, operationId: route.id, surface: route.name },', @@ -1178,14 +1178,19 @@ export const generatedRouteFlightWorkerSource = (options: GeneratedRouteFlightWo ' const renderStartedAt = performance.now();', " const element = validationError === undefined ? composeLayouts(observedRoute, props, controller.signal) : createElement(Agent.Result, null, createElement(Agent.Error, { code: 'invalid-input' }, `Input validation error: ${validationError instanceof Error ? validationError.message : String(validationError)}`));", ' const flight = renderAgentFlight(element, { signal: controller.signal });', - ' const renderedBytes = new Uint8Array(await new Response(flight).arrayBuffer());', + ' const reader = flight.getReader();', + ' while (true) {', + ' const next = await reader.read();', + ' if (next.done) break;', + ' const bytes = next.value;', + " parentPort.postMessage({ bytes, id: message.id, type: 'chunk' }, [bytes.buffer]);", + ' }', " if (message.observe === true) parentPort.postMessage({ durationMs: performance.now() - renderStartedAt, id: message.id, type: 'observed-render-finish' });", - ' return renderedBytes;', ' });', - ' parentPort.postMessage({ bytes, id: message.id, type: \'complete\' }, [bytes.buffer]);', ...(options.state === undefined ? [] : [' } finally {', ' await bindings.close();', ' }']), + " parentPort.postMessage({ id: message.id, type: 'end' });", ' } catch (error) {', " parentPort.postMessage({ id: message.id, message: error instanceof Error ? error.message : String(error), type: 'error' });", ' } finally {', diff --git a/packages/agent-bundle/src/dev/routes/route-invocation-production.ts b/packages/agent-bundle/src/dev/routes/route-invocation-production.ts index 7263bcdc7..bcc1b5a0a 100644 --- a/packages/agent-bundle/src/dev/routes/route-invocation-production.ts +++ b/packages/agent-bundle/src/dev/routes/route-invocation-production.ts @@ -361,12 +361,11 @@ const streamFromWorker = ( } pending.delete(message.id); entry.dispatchSignal.removeEventListener('abort', entry.abort); - if (message.type === 'complete' && message.bytes !== undefined) { - entry.controller.enqueue(message.bytes); - entry.controller.close(); - return; - } - if (message.type === 'end') { + // `complete` is the whole render in one message from a Flight worker + // compiled before #718; the epoch store restores such artifacts across + // dev-server restarts until the project rebuilds. + if (message.type === 'complete' && message.bytes !== undefined) entry.controller.enqueue(message.bytes); + if (message.type === 'end' || message.type === 'complete') { entry.controller.close(); return; } @@ -379,11 +378,24 @@ const streamFromWorker = ( }>): Promise> => { const id = ++sequence; let controller!: ReadableStreamDefaultController; - const stream = new ReadableStream({ start: (opened) => { controller = opened; } }); - const abort = (): void => { + const cancelRender = (): void => { worker.postMessage({ id, type: 'cancel' }); + pending.delete(id); + }; + const abort = (): void => { + cancelRender(); controller.error(new DOMException('Agent render was aborted.', 'AbortError')); }; + // The dispatcher cancels the stream when its session closes, which can + // precede the worker's `end`; dropping the entry keeps later chunks off + // the closed controller. + const stream = new ReadableStream({ + cancel: () => { + cancelRender(); + dispatch.signal.removeEventListener('abort', abort); + }, + start: (opened) => { controller = opened; }, + }); pending.set(id, { abort, controller, dispatchSignal: dispatch.signal }); dispatch.signal.addEventListener('abort', abort, { once: true }); worker.postMessage({ diff --git a/packages/agent-bundle/src/mcp-server-runtime.ts b/packages/agent-bundle/src/mcp-server-runtime.ts index e5943e508..5a1300de3 100644 --- a/packages/agent-bundle/src/mcp-server-runtime.ts +++ b/packages/agent-bundle/src/mcp-server-runtime.ts @@ -515,9 +515,8 @@ export const createFlightWorkerHost = ( ): WarmFlightHost => { interface PendingRender { readonly abort: () => void; + readonly controller: ReadableStreamDefaultController; readonly progress?: AgentProgressReporter; - readonly reject: (error: Error) => void; - readonly resolve: (stream: ReadableStream) => void; readonly signal: AbortSignal; } // Generated route modules may write to stdout; stdout is the stdio @@ -529,7 +528,7 @@ export const createFlightWorkerHost = ( let sequence = 0; let exited = false; const failPending = (error: Error): void => { - for (const request of pending.values()) request.reject(error); + for (const request of pending.values()) request.controller.error(error); pending.clear(); }; const workerError = (message: FlightWorkerMessage): Error => { @@ -560,25 +559,30 @@ export const createFlightWorkerHost = ( : `The MCP render runtime restarted; worker exited with code ${String(code)}.`, )); }); + const settle = (id: number, request: PendingRender, error?: Error): void => { + pending.delete(id); + request.signal.removeEventListener('abort', request.abort); + if (error === undefined) request.controller.close(); + else request.controller.error(error); + }; worker.on('message', (message: FlightWorkerMessage) => { const request = pending.get(message.id); if (request === undefined) return; if (message.type === 'progress') { - void request.progress?.report(message.update as never); + // A report the session no longer accepts (it completed or failed first) + // settles a still-pending render once instead of rejecting unhandled. + Promise.resolve(request.progress?.report(message.update as never)).catch((error: unknown) => { + if (pending.get(message.id) !== request) return; + worker.postMessage({ id: message.id, type: 'cancel' }); + settle(message.id, request, error instanceof Error ? error : new Error(String(error))); + }); return; } - pending.delete(message.id); - request.signal.removeEventListener('abort', request.abort); - if (message.type === 'error') { - request.reject(workerError(message)); + if (message.type === 'chunk') { + request.controller.enqueue(message.bytes!); return; } - request.resolve(new ReadableStream({ - start(controller) { - controller.enqueue(message.bytes!); - controller.close(); - }, - })); + settle(message.id, request, message.type === 'error' ? workerError(message) : undefined); }); const warmHost = createWarmFlightHost({ artifactEpoch, @@ -597,33 +601,45 @@ export const createFlightWorkerHost = ( } const context = await agent(); const id = ++sequence; - return new Promise>((resolve, reject) => { - const abort = (): void => { - worker.postMessage({ id, type: 'cancel' }); - pending.delete(id); - reject(new DOMException('Agent render was aborted', 'AbortError')); - }; - pending.set(id, { abort, ...(progress === undefined ? {} : { progress }), reject, resolve, signal }); - signal.addEventListener('abort', abort, { once: true }); - if (signal.aborted) { - abort(); - return; - } - worker.postMessage({ - actor: context.actor, - artifactEpoch: requestEpoch ?? artifactEpoch, - host: context.host, - id, - invocation, - lineage: context.lineage, - plugin: context.plugin, - requestInvocation: context.invocation, - session: context.session, - terminal: context.terminal, - type: 'render', - workspace: context.workspace, - }); + let controller!: ReadableStreamDefaultController; + const cancelRender = (): void => { + worker.postMessage({ id, type: 'cancel' }); + pending.delete(id); + }; + const abort = (): void => { + cancelRender(); + controller.error(new DOMException('Agent render was aborted', 'AbortError')); + }; + // A consumer that cancels the stream closes its controller; dropping the + // pending entry keeps later worker chunks from reaching a closed stream. + const stream = new ReadableStream({ + cancel: () => { + cancelRender(); + signal.removeEventListener('abort', abort); + }, + start: (opened) => { controller = opened; }, + }); + if (signal.aborted) { + controller.error(new DOMException('Agent render was aborted', 'AbortError')); + return stream; + } + pending.set(id, { abort, controller, ...(progress === undefined ? {} : { progress }), signal }); + signal.addEventListener('abort', abort, { once: true }); + worker.postMessage({ + actor: context.actor, + artifactEpoch: requestEpoch ?? artifactEpoch, + host: context.host, + id, + invocation, + lineage: context.lineage, + plugin: context.plugin, + requestInvocation: context.invocation, + session: context.session, + terminal: context.terminal, + type: 'render', + workspace: context.workspace, }); + return stream; }, }, }); diff --git a/packages/agent-bundle/tests/entry-shell.test.ts b/packages/agent-bundle/tests/entry-shell.test.ts index 8b420c907..566d0fd6a 100644 --- a/packages/agent-bundle/tests/entry-shell.test.ts +++ b/packages/agent-bundle/tests/entry-shell.test.ts @@ -688,8 +688,17 @@ it('generates the warm react-server Flight worker separately from the MCP dispat expect(source).toContain('route.module.inputSchema.parse(message.invocation.props.input)'); expect(source).toContain('message.validateInput !== true ? { input: message.invocation.props.input'); expect(source).toContain("createElement(Agent.Error, { code: 'invalid-input' }"); + // Flight bytes leave the worker chunk by chunk inside the request scope — + // the same `chunk`/`end`/`error` transport the rendered CLI worker speaks — + // so a Suspense fallback reaches the consumer before the render completes (#686). + expect(source).toContain("parentPort.postMessage({ bytes, id: message.id, type: 'chunk' }, [bytes.buffer]);"); + expect(source).toContain("parentPort.postMessage({ id: message.id, type: 'end' });"); + expect(source).not.toContain('arrayBuffer'); + expect(source).not.toContain("type: 'complete'"); + expect(source.indexOf("type: 'chunk'")).toBeLessThan(source.indexOf("type: 'observed-render-finish'")); + expect(source.indexOf("type: 'end'")).toBeGreaterThan(source.indexOf("type: 'observed-render-finish'")); expect(createHash('sha256').update(source).digest('hex')).toBe( - '68b593d21fdf4aaa5c51d99cffb1a106773e50ef1837072bfb0774699719f98c', + 'f363a6abcf9002e413be930cd00e1514da975ddfb8ca58a70be07f392601a4f4', ); expect(generate({ artifactEpoch: 'route-fixture@1.2.3', diff --git a/packages/agent-bundle/tests/generated-route-server.test.ts b/packages/agent-bundle/tests/generated-route-server.test.ts index 34d0c378a..2ee67534f 100644 --- a/packages/agent-bundle/tests/generated-route-server.test.ts +++ b/packages/agent-bundle/tests/generated-route-server.test.ts @@ -1,6 +1,7 @@ import { Client } from '@modelcontextprotocol/client'; import { StdioClientTransport } from '@modelcontextprotocol/client/stdio'; import { spawn } from 'node:child_process'; +import { existsSync } from 'node:fs'; import { mkdir, mkdtemp, readFile, rm, stat, symlink, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { dirname, join, resolve } from 'node:path'; @@ -876,6 +877,58 @@ it('emits MCP progress notifications only when a progress token is supplied', { } }); +it('streams an authored Suspense fallback to the stdio client before the render completes (#686)', { retry: 2, timeout: 60_000 }, async () => { + const root = await mkdtemp(join(tmpdir(), 'agent-bundle-generated-streaming-')); + roots.push(root); + await writeGeneratedProject(root, { + 'src/mcp/curator/tools/gated.tsx': [ + "import { existsSync } from 'node:fs';", + "import { Agent } from '@agent-bundle/runtime';", + "import { createElement, Suspense } from 'react';", + "import { z } from 'zod';", + "export const config = { annotations: { readOnlyHint: true }, description: 'Wait for a gate.' };", + 'export const inputSchema = z.object({ gate: z.string() }).strict();', + 'export const resultSchema = z.object({ ok: z.literal(true) }).strict();', + 'const Gated = async ({ gate, signal }) => {', + ' while (!existsSync(gate)) {', + " if (signal.aborted) throw new DOMException('gate abandoned', 'AbortError');", + ' await new Promise((resolve) => setTimeout(resolve, 25));', + ' }', + " return createElement(Agent.Text, null, 'released');", + '};', + 'export default async function GatedTool({ input, signal }) {', + " return createElement(Agent.Result, { value: { ok: true } }, createElement(Suspense, { fallback: createElement(Agent.Progress, { completed: 0, message: 'waiting', total: 1 }) }, createElement(Gated, { gate: input.gate, signal })));", + '}', + '', + ].join('\n'), + }); + const session = await connectGeneratedServer(root); + const notifications: Array<{ readonly message?: string; readonly progress: number }> = []; + session.client.setNotificationHandler('notifications/progress', (notification) => { + notifications.push(notification.params); + }); + try { + const gate = join(root, 'gate.marker'); + const pending = session.client.callTool({ + arguments: { gate }, + name: 'gated', + _meta: { progressToken: 'tok-gate' }, + }, { signal: AbortSignal.timeout(20_000) }); + // The fallback's progress node reaches the wire while the child is still + // blocked: the worker streamed the shell instead of buffering the render. + await expect.poll(() => notifications, { timeout: 15_000 }).toEqual([{ message: 'waiting', progress: 0, progressToken: 'tok-gate', total: 1 }]); + expect(existsSync(gate)).toBe(false); + await writeFile(gate, 'open\n'); + await expect(pending).resolves.toMatchObject({ + content: [{ text: 'released', type: 'text' }], + structuredContent: { ok: true }, + }); + expect(notifications).toHaveLength(1); + } finally { + await session.close(); + } +}); + /** * #492: what a thrown (not represented) route error is on each MCP surface of * a real generated stdio server. A tool throw is the SDK's default tool error; diff --git a/packages/agent-bundle/tests/helpers/gated-routes.ts b/packages/agent-bundle/tests/helpers/gated-routes.ts new file mode 100644 index 000000000..b39804163 --- /dev/null +++ b/packages/agent-bundle/tests/helpers/gated-routes.ts @@ -0,0 +1,38 @@ +/** + * Fixture routes whose render suspends until a gate file appears under the + * project's `.agent-bundle` directory, so a test can observe the authored + * Suspense fallback on the wire before releasing the child (#686). + */ +export const gatedRouteFiles: Readonly> = Object.freeze({ + 'src/gate.ts': [ + "import { existsSync } from 'node:fs';", + "import { join } from 'node:path';", + '', + 'export const awaitGate = async (gate, signal) => {', + " const path = join(process.cwd(), '.agent-bundle', gate);", + ' while (!existsSync(path)) {', + " if (signal.aborted) throw new DOMException('gate abandoned', 'AbortError');", + ' await new Promise((resolve) => setTimeout(resolve, 25));', + ' }', + '};', + '', + ].join('\n'), + 'src/mcp/status/tools/live.tsx': [ + "import { Agent } from '@agent-bundle/runtime';", + "import { createElement, Suspense } from 'react';", + "import { z } from 'zod';", + "import { awaitGate } from '../../../gate.js';", + '', + 'export const inputSchema = z.object({ gate: z.string().optional() }).strict();', + 'export const resultSchema = z.object({ done: z.boolean() }).strict();', + 'const Slow = async ({ gate, signal }) => {', + ' if (gate !== undefined) await awaitGate(gate, signal);', + " return createElement(Agent.Text, null, 'stream complete');", + '};', + '', + 'export default async function Live({ input, signal }) {', + " return createElement(Agent.Result, { value: { done: true } }, createElement(Suspense, { fallback: createElement(Agent.Progress, { completed: 0, message: 'streaming', total: 1 }) }, createElement(Slow, { gate: input.gate, signal })));", + '}', + '', + ].join('\n'), +}); diff --git a/packages/agent-bundle/tests/route-invocation-dev-server.test.ts b/packages/agent-bundle/tests/route-invocation-dev-server.test.ts index 92f6601e9..862f1a0e2 100644 --- a/packages/agent-bundle/tests/route-invocation-dev-server.test.ts +++ b/packages/agent-bundle/tests/route-invocation-dev-server.test.ts @@ -2,6 +2,7 @@ import { existsSync } from 'node:fs'; import { mkdir, readFile, readdir, rm, symlink, writeFile } from 'node:fs/promises'; import { join } from 'node:path'; +import type { AgentRenderEvent } from '@agent-bundle/runtime'; import { expect, it } from '@rstest/core'; import { Client } from '@modelcontextprotocol/client'; import { StdioClientTransport } from '@modelcontextprotocol/client/stdio'; @@ -15,6 +16,7 @@ import { stableJson } from '../src/core/digest.ts'; import { pluginRootEnvAnchor, pluginStateRootEnvAnchor } from '../src/core/types.ts'; import { createWorkbenchAssetSource } from '../src/dev/workbench-assets.ts'; import { startDevServer } from '../src/dev/workbench-server.ts'; +import { gatedRouteFiles } from './helpers/gated-routes.ts'; import { createProjectFixture } from './helpers/project-fixture.ts'; import { agentBundleNodeModules } from './helpers/workspace-paths.ts'; import { replaceWatchedSourceAndAwaitRebuild } from './support/watched-files.ts'; @@ -44,26 +46,52 @@ const readEvent = async ( } }; -const readInvocationStream = async (response: Response): Promise[]> => { +type StreamMessage = Record; + +/** Reads one invocation SSE stream message at a time, so the test can act between messages. */ +const invocationMessages = (response: Response): Readonly<{ + readonly next: (matches: (message: StreamMessage) => boolean) => Promise; + readonly seen: readonly StreamMessage[]; +}> => { const reader = response.body!.pipeThrough(new TextDecoderStream()).getReader(); - const messages: Record[] = []; + const seen: StreamMessage[] = []; + const queue: StreamMessage[] = []; let buffered = ''; - for (;;) { - const next = await reader.read(); - if (next.done) throw new Error('Invocation stream ended before final.'); - buffered += next.value; - const frames = buffered.split('\n\n'); - buffered = frames.pop() ?? ''; - for (const frame of frames) { - const data = frame.split('\n').find((line) => line.startsWith('data: ')); - if (data === undefined) continue; - const message = JSON.parse(data.slice('data: '.length)) as Record; - messages.push(message); - if (message.type === 'final') return messages; + const pull = async (): Promise => { + while (queue.length === 0) { + const next = await reader.read(); + if (next.done) throw new Error('Invocation stream ended.'); + buffered += next.value; + const frames = buffered.split('\n\n'); + buffered = frames.pop() ?? ''; + for (const frame of frames) { + const data = frame.split('\n').find((line) => line.startsWith('data: ')); + if (data !== undefined) queue.push(JSON.parse(data.slice('data: '.length)) as StreamMessage); + } } - } + const message = queue.shift()!; + seen.push(message); + return message; + }; + return { + next: async (matches) => { + for (;;) { + const message = await pull(); + if (matches(message)) return message; + } + }, + seen, + }; }; +const renderEvents = (messages: readonly StreamMessage[]): readonly AgentRenderEvent[] => + messages.filter((message) => message.type === 'render').map((message) => message.event as AgentRenderEvent); + +const isFallbackRender = (fallback: string) => (message: StreamMessage): boolean => + message.type === 'render' && JSON.stringify(message.event).includes(JSON.stringify(fallback)); + +const isFinal = (message: StreamMessage): boolean => message.type === 'final'; + it('invokes compiled tool and event routes through the foreground server', { timeout: 180_000 }, async () => { const project = await createProjectFixture({ config: [ @@ -125,19 +153,26 @@ it('invokes compiled tool and event routes through the foreground server', { tim "import { appendFileSync } from 'node:fs';", "import { join } from 'node:path';", "import { Agent, agent } from '@agent-bundle/runtime';", - "import { createElement } from 'react';", + "import { createElement, Suspense } from 'react';", + "import { awaitGate } from '../../gate.js';", "export { default as preflight } from './after.preflight.js';", '', "export const config = { providers: ['clock'], runtime: 'standalone' };", '', - 'export default async function AfterTool({ canonical, preflight }) {', + 'const Observed = async ({ gate, signal, toolName }) => {', + ' if (gate !== undefined) await awaitGate(gate, signal);', + ' return createElement(Agent.Context, null, `Observed ${toolName}.`);', + '};', + '', + 'export default async function AfterTool({ canonical, preflight, signal }) {', ' const context = await agent();', " appendFileSync(join(process.cwd(), '.agent-bundle', 'defer-handler.marker'), 'run\\n');", " const value = { outcome: 'defer', providers: Object.keys(context.providers).sort(), ticket: preflight.ticket };", - " return createElement(Agent.Result, { value }, createElement(Agent.Context, null, `Observed ${canonical.payload.toolName}.`));", + " return createElement(Agent.Result, { value }, createElement(Suspense, { fallback: createElement(Agent.Progress, { completed: 0, message: 'event streaming', total: 1 }) }, createElement(Observed, { gate: canonical.payload.toolInput?.value?.gate, signal, toolName: canonical.payload.toolName?.value })));", '}', '', ].join('\n'), + ...gatedRouteFiles, 'src/events/prompt/submit.preflight.ts': "export default () => ({ outcome: 'continue' });\n", 'src/events/prompt/submit.tsx': [ "import { writeFileSync } from 'node:fs';", @@ -222,20 +257,17 @@ it('invokes compiled tool and event routes through the foreground server', { tim '}', '', ].join('\n'), - 'src/mcp/status/tools/live.tsx': [ + 'src/mcp/status/tools/crash.tsx': [ "import { Agent } from '@agent-bundle/runtime';", "import { createElement, Suspense } from 'react';", "import { z } from 'zod';", '', 'export const inputSchema = z.object({}).strict();', - 'export const resultSchema = z.object({ done: z.boolean() }).strict();', - 'const Slow = async () => {', - ' await new Promise((resolve) => setTimeout(resolve, 1_000));', - " return createElement(Agent.Text, null, 'stream complete');", - '};', + 'export const resultSchema = z.object({ ok: z.boolean() }).strict();', + "const Boom = async () => { throw new Error('render exploded'); };", '', - 'export default async function Live() {', - " return createElement(Agent.Result, { value: { done: true } }, createElement(Suspense, { fallback: createElement(Agent.Progress, { completed: 0, message: 'streaming', total: 1 }) }, createElement(Slow)));", + 'export default async function Crash() {', + " return createElement(Agent.Result, { value: { ok: true } }, createElement(Suspense, { fallback: createElement(Agent.Progress, { completed: 0, message: 'crashing', total: 1 }) }, createElement(Boom)));", '}', '', ].join('\n'), @@ -310,30 +342,69 @@ it('invokes compiled tool and event routes through the foreground server', { tim headers: { cookie, origin: server.url }, }); const stateRoot = join(project.root, '.agent-bundle', 'state'); - const startLive = async () => { + const startStreaming = async (body: Record) => { const response = await fetch(`${server!.url}/api/routes/invocations`, { - body: JSON.stringify({ routeId: 'tool:status/live', stream: true }), + body: JSON.stringify({ ...body, stream: true }), headers, method: 'POST', }); expect(response.status).toBe(202); - return response.json() as Promise<{ readonly invocation: { readonly id: string; readonly status: string } }>; + const started = await response.json() as RouteInvocationResponse; + expect(started.invocation.status).toBe('running'); + const stream = await fetch(`${server!.url}/api/routes/invocations/${started.invocation.id}/stream`, { headers }); + expect(stream.status).toBe(200); + return { id: started.invocation.id, stream: invocationMessages(stream) }; }; - const live = await startLive(); - expect(live.invocation.status).toBe('running'); - const liveStreamResponse = await fetch(`${server.url}/api/routes/invocations/${live.invocation.id}/stream`, { headers }); - expect(liveStreamResponse.status).toBe(200); - const liveMessages = await readInvocationStream(liveStreamResponse); - expect(liveMessages.findIndex((message) => message.type === 'render')).toBeGreaterThanOrEqual(0); - expect(liveMessages.at(-1)).toMatchObject({ - invocation: { status: 'succeeded' }, - type: 'final', + const gatePath = (gate: string): string => join(project.root, '.agent-bundle', gate); + // Releasing returns the release instant so a final envelope can prove the + // render completed after it — the fallback was streamed from a blocked + // render, not replayed from a finished one. + const releaseGate = async (gate: string): Promise => { + expect(existsSync(gatePath(gate))).toBe(false); + const releasedAt = Date.now(); + await writeFile(gatePath(gate), 'open\n'); + return releasedAt; + }; + const completedAfter = (message: StreamMessage, releasedAt: number): void => { + expect(Date.parse((message.invocation as RouteInvocationResponse['invocation']).completedAt)).toBeGreaterThanOrEqual(releasedAt); + }; + + // The compiled MCP tool streams its authored Suspense fallback to the + // Workbench consumer while the child is still blocked on the gate; no + // terminal event exists yet (#686). + const live = await startStreaming({ input: { gate: 'live-gate' }, routeId: 'tool:status/live' }); + const liveFallback = await live.stream.next(isFallbackRender('streaming')); + expect(liveFallback.event).toMatchObject({ type: 'shell' }); + expect(renderEvents(live.stream.seen).map((event) => event.type)).toEqual(['shell']); + const liveReleasedAt = await releaseGate('live-gate'); + const liveFinal = await live.stream.next(isFinal); + expect(liveFinal).toMatchObject({ invocation: { status: 'succeeded' }, type: 'final' }); + completedAfter(liveFinal, liveReleasedAt); + const liveEvents = renderEvents(live.stream.seen); + expect(liveEvents.map((event) => event.type)).toEqual(['shell', 'replace', 'complete']); + const liveComplete = liveEvents.at(-1)!; + if (liveComplete.type !== 'complete') throw new Error('Expected a complete render event.'); + expect(JSON.stringify(liveComplete.document)).toContain('stream complete'); + expect(JSON.stringify(liveComplete.document)).not.toContain('"streaming"'); + // Parity control: the same route completing without a gate renders the + // same final document the streamed run replaced its fallback with. + const controlResponse = await fetch(`${server.url}/api/routes/invocations`, { + body: JSON.stringify({ input: {}, routeId: 'tool:status/live' }), + headers, + method: 'POST', }); + expect(controlResponse.status).toBe(200); + const control = await controlResponse.json() as RouteInvocationResponse; + expect(control.invocation.status).toBe('succeeded'); + expect(stableJson(control.invocation.document)).toBe(stableJson(liveComplete.document)); + expect(stableJson((liveFinal.invocation as RouteInvocationResponse['invocation']).document)).toBe(stableJson(control.invocation.document)); - const cancelling = await startLive(); - const cancellingStream = await fetch(`${server.url}/api/routes/invocations/${cancelling.invocation.id}/stream`, { headers }); - const cancellingMessages = readInvocationStream(cancellingStream); - const cancelResponse = await fetch(`${server.url}/api/routes/invocations/${cancelling.invocation.id}/cancel`, { + // Cancelling while blocked: the fallback is the only render published, + // the terminal message says cancelled, and the released gate publishes + // nothing afterwards because the stream has already closed. + const cancelling = await startStreaming({ input: { gate: 'cancel-gate' }, routeId: 'tool:status/live' }); + await cancelling.stream.next(isFallbackRender('streaming')); + const cancelResponse = await fetch(`${server.url}/api/routes/invocations/${cancelling.id}/cancel`, { headers, method: 'POST', }); @@ -341,11 +412,35 @@ it('invokes compiled tool and event routes through the foreground server', { tim const cancelled = await cancelResponse.json() as RouteInvocationResponse; expect(cancelled.invocation).toMatchObject({ status: 'cancelled' }); expect(cancelled.invocation).not.toHaveProperty('outcome'); - expect((await cancellingMessages).at(-1)).toMatchObject({ + await expect(cancelling.stream.next(isFinal)).resolves.toMatchObject({ invocation: { status: 'cancelled' }, type: 'final', }); - const finalCancelResponse = await fetch(`${server.url}/api/routes/invocations/${cancelling.invocation.id}/cancel`, { + expect(renderEvents(cancelling.stream.seen).map((event) => event.type)).toEqual(['shell']); + await releaseGate('cancel-gate'); + await expect(cancelling.stream.next(() => true)).rejects.toThrow('Invocation stream ended.'); + + // A render that fails behind its fallback settles once: one terminal + // render event, one final message, and the same outcome the completed + // (non-streamed) run reports. + const crashing = await startStreaming({ input: {}, routeId: 'tool:status/crash' }); + const crashFinal = await crashing.stream.next(isFinal); + const crashEvents = renderEvents(crashing.stream.seen); + expect(crashEvents.filter((entry) => entry.type === 'complete')).toHaveLength(1); + expect(crashEvents.at(-1)).toMatchObject({ type: 'complete' }); + expect(JSON.stringify(crashEvents.at(-1))).toContain('render exploded'); + await expect(crashing.stream.next(() => true)).rejects.toThrow('Invocation stream ended.'); + const crashControlResponse = await fetch(`${server.url}/api/routes/invocations`, { + body: JSON.stringify({ input: {}, routeId: 'tool:status/crash' }), + headers, + method: 'POST', + }); + const crashControl = await crashControlResponse.json() as RouteInvocationResponse; + const crashInvocation = crashFinal.invocation as RouteInvocationResponse['invocation']; + expect(crashInvocation.status).toBe(crashControl.invocation.status); + expect(crashInvocation.outcome).toEqual(crashControl.invocation.outcome); + expect(stableJson(crashInvocation.document)).toBe(stableJson(crashControl.invocation.document)); + const finalCancelResponse = await fetch(`${server.url}/api/routes/invocations/${cancelling.id}/cancel`, { headers, method: 'POST', }); @@ -490,6 +585,36 @@ it('invokes compiled tool and event routes through the foreground server', { tim expect(await readFile(join(project.root, '.agent-bundle', 'defer-gate.marker'), 'utf8')).toBe('gate\n'); expect(await readFile(join(project.root, '.agent-bundle', 'defer-handler.marker'), 'utf8')).toBe('run\n'); expect(event.invocation.projection.hosts?.[0]).toMatchObject({ host: 'claude' }); + // The same compiled event route, streamed: the manifest-selected host's + // preflight runs, the authored fallback arrives while the child is gated, + // and the released render matches the completed run's document (#686). + const gatedEvent = await startStreaming({ + input: { + cwd: project.root, + hook_event_name: 'PostToolUse', + session_id: 'session-1', + tool_input: { gate: 'event-gate' }, + tool_name: 'Write', + tool_response: { ok: true }, + tool_use_id: 'use-2', + transcript_path: join(project.root, 'transcript.json'), + }, + routeId: 'event:tool/after', + surface: { host: 'claude', kind: 'event' }, + }); + const eventFallback = await gatedEvent.stream.next(isFallbackRender('event streaming')); + expect(eventFallback.event).toMatchObject({ type: 'shell' }); + expect(renderEvents(gatedEvent.stream.seen).map((entry) => entry.type)).toEqual(['shell']); + const eventReleasedAt = await releaseGate('event-gate'); + const gatedEventFinal = await gatedEvent.stream.next(isFinal); + expect(gatedEventFinal).toMatchObject({ invocation: { status: 'succeeded' }, type: 'final' }); + completedAfter(gatedEventFinal, eventReleasedAt); + expect(renderEvents(gatedEvent.stream.seen).map((entry) => entry.type)).toEqual(['shell', 'replace', 'complete']); + const gatedEventInvocation = gatedEventFinal.invocation as RouteInvocationResponse['invocation']; + expect(gatedEventInvocation.result).toEqual(event.invocation.result); + expect(stableJson(gatedEventInvocation.document)).toBe(stableJson(event.invocation.document)); + expect(gatedEventInvocation.projection.hosts?.[0]).toMatchObject({ host: 'claude' }); + expect(gatedEventInvocation.trace?.map((trace) => trace.kind)).toEqual(event.invocation.trace?.map((trace) => trace.kind)); expect(event.invocation.context.session).toEqual({ source: 'receipt', state: 'available', diff --git a/packages/workbench/tests/streaming-render.e2e.test.ts b/packages/workbench/tests/streaming-render.e2e.test.ts new file mode 100644 index 000000000..91c94f71f --- /dev/null +++ b/packages/workbench/tests/streaming-render.e2e.test.ts @@ -0,0 +1,86 @@ +import { existsSync } from 'node:fs'; +import { mkdir, symlink, writeFile } from 'node:fs/promises'; +import { join } from 'node:path'; + +import { expect } from '@rstest/playwright'; + +import { inspectWorkbenchSurface } from '../../agent-bundle/src/test/index.ts'; +import { gatedRouteFiles } from '../../agent-bundle/tests/helpers/gated-routes.ts'; +import { createProjectFixture, removeProjectFixture } from '../../agent-bundle/tests/helpers/project-fixture.ts'; +import { agentBundleNodeModules } from '../../agent-bundle/tests/helpers/workspace-paths.ts'; +import { timeScale } from '../../agent-bundle/tests/support/time-scale.ts'; +import { applicationLeafForRouteId } from '../src/application/application-tree-model.ts'; +import { + expectRenderedDocument, + fillRouteInput, + runSelectedRoute, + selectApplicationLeaf, + workbenchTestId, +} from './support/workbench-acceptance.ts'; +import { buildWorkbench, e2e, startWorkbenchDevServer, withWorkbenchServer } from './support/workbench-e2e.ts'; + +const browserTimeout = 15_000 * timeScale; +const runTimeout = 60_000 * timeScale; + +/** + * #686: the production MCP surface streams the compiled tool's authored + * Suspense fallback into the Rendered pane while the child is still blocked, + * and the document that replaces it is the one the unblocked run renders. + */ +e2e('renders a compiled MCP tool\'s Suspense fallback before its gated child releases', { timeout: 180_000 * timeScale }, async ({ page }) => { + await buildWorkbench(); + await withWorkbenchServer({ + createProject: () => createProjectFixture({ + config: [ + 'export default {', + " plugin: { name: 'streaming-render-e2e', version: '1.0.0' },", + " targets: ['portable'],", + '};', + '', + ].join('\n'), + files: { + ...gatedRouteFiles, + 'package.json': '{"dependencies":{"@agent-bundle/runtime":"workspace:*","react":"19.2.8","zod":"4.5.4"},"type":"module"}\n', + }, + prefix: 'agent-bundle-streaming-render-e2e-', + }), + dispose: (project) => removeProjectFixture(project.root), + setup: async (project) => { + await mkdir(join(project.root, '.agent-bundle'), { recursive: true }); + await symlink(agentBundleNodeModules, join(project.root, 'node_modules'), 'dir'); + }, + start: (project) => startWorkbenchDevServer(project), + }, async (server, project) => { + const pageErrors: Error[] = []; + page.on('pageerror', (error) => pageErrors.push(error)); + const surface = await inspectWorkbenchSurface({ root: project.root }); + const leaf = applicationLeafForRouteId(surface.application, 'tool:status/live'); + if (leaf?.ref.kind !== 'tool') throw new Error('inspectWorkbenchSurface did not project tool:status/live as a tool leaf.'); + await selectApplicationLeaf(page, server.url, leaf); + + const gate = join(project.root, '.agent-bundle', 'browser-gate'); + await fillRouteInput(page, { gate: 'browser-gate' }); + await workbenchTestId(page, 'routeRun').click(); + await expect(workbenchTestId(page, 'routeRunningStatus')).toBeVisible({ timeout: browserTimeout }); + await workbenchTestId(page, 'resultTabRendered').click(); + const document = workbenchTestId(page, 'renderedDocument'); + const body = document.locator('.rendered-document-body'); + await expect(body.locator('.agent-document-progress')).toContainText('streaming', { timeout: browserTimeout }); + await expect(document).toHaveAttribute('aria-busy', 'true'); + expect(existsSync(gate)).toBe(false); + + await writeFile(gate, 'open\n'); + await expect(workbenchTestId(page, 'routeStatus')).toHaveClass(/route-status--succeeded/u, { timeout: runTimeout }); + const streamed = await expectRenderedDocument(page, runTimeout); + const streamedText = await streamed.locator('.rendered-document-body').innerText(); + expect(streamedText).toContain('stream complete'); + expect(streamedText).not.toContain('streaming'); + + // Parity control: the same route with an already-open gate never suspends + // and renders the same final document. + await runSelectedRoute(page, runTimeout); + const control = await expectRenderedDocument(page, runTimeout); + expect(await control.locator('.rendered-document-body').innerText()).toBe(streamedText); + expect(pageErrors).toEqual([]); + }); +}); diff --git a/rstest.integration-tests.ts b/rstest.integration-tests.ts index 92cc0a00b..ee34df486 100644 --- a/rstest.integration-tests.ts +++ b/rstest.integration-tests.ts @@ -119,6 +119,7 @@ export const integrationTestFiles: readonly string[] = [ 'packages/workbench/tests/rsbuild-workbench.test.ts', 'packages/workbench/tests/route-editor-atoms-disposal.test.ts', 'packages/workbench/tests/sessions.e2e.test.ts', + 'packages/workbench/tests/streaming-render.e2e.test.ts', 'packages/workbench/tests/web-command.e2e.test.ts', 'packages/workbench/tests/workbench-dev-command.test.ts', ];