Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/686-stream-generated-flight.md
Original file line number Diff line number Diff line change
@@ -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)
67 changes: 34 additions & 33 deletions packages/agent-bundle/src/adapters/hook-contract.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {',
Expand Down
13 changes: 9 additions & 4 deletions packages/agent-bundle/src/build/entry-shell.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 },',
Expand Down Expand Up @@ -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 {',
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand All @@ -379,11 +378,24 @@ const streamFromWorker = (
}>): Promise<ReadableStream<Uint8Array>> => {
const id = ++sequence;
let controller!: ReadableStreamDefaultController<Uint8Array>;
const stream = new ReadableStream<Uint8Array>({ 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<Uint8Array>({
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({
Expand Down
96 changes: 56 additions & 40 deletions packages/agent-bundle/src/mcp-server-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -515,9 +515,8 @@ export const createFlightWorkerHost = (
): WarmFlightHost => {
interface PendingRender {
readonly abort: () => void;
readonly controller: ReadableStreamDefaultController<Uint8Array>;
readonly progress?: AgentProgressReporter;
readonly reject: (error: Error) => void;
readonly resolve: (stream: ReadableStream<Uint8Array>) => void;
readonly signal: AbortSignal;
}
// Generated route modules may write to stdout; stdout is the stdio
Expand All @@ -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 => {
Expand Down Expand Up @@ -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<Uint8Array>({
start(controller) {
controller.enqueue(message.bytes!);
controller.close();
},
}));
settle(message.id, request, message.type === 'error' ? workerError(message) : undefined);
});
const warmHost = createWarmFlightHost({
artifactEpoch,
Expand All @@ -597,33 +601,45 @@ export const createFlightWorkerHost = (
}
const context = await agent();
const id = ++sequence;
return new Promise<ReadableStream<Uint8Array>>((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<Uint8Array>;
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<Uint8Array>({
cancel: () => {
cancelRender();
signal.removeEventListener('abort', abort);
},
start: (opened) => { controller = opened; },
});
Comment thread
ScriptedAlchemy marked this conversation as resolved.
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;
},
},
});
Expand Down
11 changes: 10 additions & 1 deletion packages/agent-bundle/tests/entry-shell.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down
53 changes: 53 additions & 0 deletions packages/agent-bundle/tests/generated-route-server.test.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -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;
Expand Down
Loading
Loading