Skip to content

Commit a057fdf

Browse files
fix: stream generated MCP/event Flight bytes before render completion (#718)
Closes #686
1 parent b805eec commit a057fdf

11 files changed

Lines changed: 481 additions & 130 deletions
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'agent-bundle': patch
3+
---
4+
5+
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)

‎packages/agent-bundle/src/adapters/hook-contract.ts‎

Lines changed: 34 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -734,40 +734,41 @@ const eventRouteHookWrapperSource = (
734734
' execute: async (dispatch) => {',
735735
' const context = await agent();',
736736
' const id = ++sequence;',
737-
' return new Promise((resolvePromise, rejectPromise) => {',
738-
' const abort = () => { worker.postMessage({ id, type: "cancel" }); rejectPromise(new DOMException("Agent render was aborted", "AbortError")); };',
739-
' const cleanup = () => { dispatch.signal.removeEventListener("abort", abort); worker.off("error", onError); worker.off("exit", onExit); worker.off("message", onMessage); };',
740-
' const reject = (error) => { cleanup(); rejectPromise(error); };',
741-
' const onError = (error) => { reject(error); };',
742-
' const onExit = (code) => { reject(new Error(`Generated hook Flight worker exited with code ${String(code)}.`)); };',
743-
' const onMessage = (message) => {',
744-
' if (message.id !== id) return;',
745-
' if (message.type === "progress") { Promise.resolve(dispatch.progress?.report(message.update)).catch(reject); return; }',
746-
' if (message.type === "error") { reject(new Error(message.message)); return; }',
747-
' if (message.type !== "complete") return;',
748-
' cleanup();',
749-
' resolvePromise(new ReadableStream({ start(controller) { controller.enqueue(message.bytes); controller.close(); } }));',
750-
' };',
751-
' worker.on("error", onError);',
752-
' worker.on("exit", onExit);',
753-
' worker.on("message", onMessage);',
754-
' dispatch.signal.addEventListener("abort", abort, { once: true });',
755-
' if (dispatch.signal.aborted) { abort(); return; }',
756-
' worker.postMessage({',
757-
' actor: context.actor,',
758-
' artifactEpoch: flightArtifactEpoch,',
759-
' host: context.host,',
760-
' id,',
761-
' invocation: dispatch.invocation,',
762-
' lineage: context.lineage,',
763-
' plugin: context.plugin,',
764-
' requestInvocation: context.invocation,',
765-
' session: context.session,',
766-
' terminal: context.terminal,',
767-
' type: "render",',
768-
' workspace: context.workspace,',
769-
' });',
737+
' let controller;',
738+
' const cleanup = () => { dispatch.signal.removeEventListener("abort", abort); worker.off("error", fail); worker.off("exit", onExit); worker.off("message", onMessage); };',
739+
' const stream = new ReadableStream({ cancel() { worker.postMessage({ id, type: "cancel" }); cleanup(); }, start(opened) { controller = opened; } });',
740+
' const fail = (error) => { cleanup(); controller.error(error); };',
741+
' const abort = () => { worker.postMessage({ id, type: "cancel" }); fail(new DOMException("Agent render was aborted", "AbortError")); };',
742+
' const onExit = (code) => { fail(new Error(`Generated hook Flight worker exited with code ${String(code)}.`)); };',
743+
' const onMessage = (message) => {',
744+
' if (message.id !== id) return;',
745+
' if (message.type === "progress") { Promise.resolve(dispatch.progress?.report(message.update)).catch(fail); return; }',
746+
' if (message.type === "chunk") { controller.enqueue(message.bytes); return; }',
747+
' if (message.type === "error") { fail(new Error(message.message)); return; }',
748+
' if (message.type !== "end") return;',
749+
' cleanup();',
750+
' controller.close();',
751+
' };',
752+
' if (dispatch.signal.aborted) { controller.error(new DOMException("Agent render was aborted", "AbortError")); return stream; }',
753+
' worker.on("error", fail);',
754+
' worker.on("exit", onExit);',
755+
' worker.on("message", onMessage);',
756+
' dispatch.signal.addEventListener("abort", abort, { once: true });',
757+
' worker.postMessage({',
758+
' actor: context.actor,',
759+
' artifactEpoch: flightArtifactEpoch,',
760+
' host: context.host,',
761+
' id,',
762+
' invocation: dispatch.invocation,',
763+
' lineage: context.lineage,',
764+
' plugin: context.plugin,',
765+
' requestInvocation: context.invocation,',
766+
' session: context.session,',
767+
' terminal: context.terminal,',
768+
' type: "render",',
769+
' workspace: context.workspace,',
770770
' });',
771+
' return stream;',
771772
' },',
772773
' });',
773774
' try {',

‎packages/agent-bundle/src/build/entry-shell.ts‎

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1143,7 +1143,7 @@ export const generatedRouteFlightWorkerSource = (options: GeneratedRouteFlightWo
11431143
? []
11441144
: [' const bindings = await runtimeState.requestBindings({ signal: controller.signal });', ' try {']),
11451145
' const plugin = message.plugin ?? pluginRoot.identity;',
1146-
' const bytes = await runAgentRequest({',
1146+
' await runAgentRequest({',
11471147
' ...(message.actor === undefined ? {} : { actor: message.actor }),',
11481148
' ...(message.host === undefined ? {} : { host: message.host }),',
11491149
' 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
11781178
' const renderStartedAt = performance.now();',
11791179
" 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)}`));",
11801180
' const flight = renderAgentFlight(element, { signal: controller.signal });',
1181-
' const renderedBytes = new Uint8Array(await new Response(flight).arrayBuffer());',
1181+
' const reader = flight.getReader();',
1182+
' while (true) {',
1183+
' const next = await reader.read();',
1184+
' if (next.done) break;',
1185+
' const bytes = next.value;',
1186+
" parentPort.postMessage({ bytes, id: message.id, type: 'chunk' }, [bytes.buffer]);",
1187+
' }',
11821188
" if (message.observe === true) parentPort.postMessage({ durationMs: performance.now() - renderStartedAt, id: message.id, type: 'observed-render-finish' });",
1183-
' return renderedBytes;',
11841189
' });',
1185-
' parentPort.postMessage({ bytes, id: message.id, type: \'complete\' }, [bytes.buffer]);',
11861190
...(options.state === undefined
11871191
? []
11881192
: [' } finally {', ' await bindings.close();', ' }']),
1193+
" parentPort.postMessage({ id: message.id, type: 'end' });",
11891194
' } catch (error) {',
11901195
" parentPort.postMessage({ id: message.id, message: error instanceof Error ? error.message : String(error), type: 'error' });",
11911196
' } finally {',

‎packages/agent-bundle/src/dev/routes/route-invocation-production.ts‎

Lines changed: 20 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -361,12 +361,11 @@ const streamFromWorker = (
361361
}
362362
pending.delete(message.id);
363363
entry.dispatchSignal.removeEventListener('abort', entry.abort);
364-
if (message.type === 'complete' && message.bytes !== undefined) {
365-
entry.controller.enqueue(message.bytes);
366-
entry.controller.close();
367-
return;
368-
}
369-
if (message.type === 'end') {
364+
// `complete` is the whole render in one message from a Flight worker
365+
// compiled before #718; the epoch store restores such artifacts across
366+
// dev-server restarts until the project rebuilds.
367+
if (message.type === 'complete' && message.bytes !== undefined) entry.controller.enqueue(message.bytes);
368+
if (message.type === 'end' || message.type === 'complete') {
370369
entry.controller.close();
371370
return;
372371
}
@@ -379,11 +378,24 @@ const streamFromWorker = (
379378
}>): Promise<ReadableStream<Uint8Array>> => {
380379
const id = ++sequence;
381380
let controller!: ReadableStreamDefaultController<Uint8Array>;
382-
const stream = new ReadableStream<Uint8Array>({ start: (opened) => { controller = opened; } });
383-
const abort = (): void => {
381+
const cancelRender = (): void => {
384382
worker.postMessage({ id, type: 'cancel' });
383+
pending.delete(id);
384+
};
385+
const abort = (): void => {
386+
cancelRender();
385387
controller.error(new DOMException('Agent render was aborted.', 'AbortError'));
386388
};
389+
// The dispatcher cancels the stream when its session closes, which can
390+
// precede the worker's `end`; dropping the entry keeps later chunks off
391+
// the closed controller.
392+
const stream = new ReadableStream<Uint8Array>({
393+
cancel: () => {
394+
cancelRender();
395+
dispatch.signal.removeEventListener('abort', abort);
396+
},
397+
start: (opened) => { controller = opened; },
398+
});
387399
pending.set(id, { abort, controller, dispatchSignal: dispatch.signal });
388400
dispatch.signal.addEventListener('abort', abort, { once: true });
389401
worker.postMessage({

‎packages/agent-bundle/src/mcp-server-runtime.ts‎

Lines changed: 56 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -515,9 +515,8 @@ export const createFlightWorkerHost = (
515515
): WarmFlightHost => {
516516
interface PendingRender {
517517
readonly abort: () => void;
518+
readonly controller: ReadableStreamDefaultController<Uint8Array>;
518519
readonly progress?: AgentProgressReporter;
519-
readonly reject: (error: Error) => void;
520-
readonly resolve: (stream: ReadableStream<Uint8Array>) => void;
521520
readonly signal: AbortSignal;
522521
}
523522
// Generated route modules may write to stdout; stdout is the stdio
@@ -529,7 +528,7 @@ export const createFlightWorkerHost = (
529528
let sequence = 0;
530529
let exited = false;
531530
const failPending = (error: Error): void => {
532-
for (const request of pending.values()) request.reject(error);
531+
for (const request of pending.values()) request.controller.error(error);
533532
pending.clear();
534533
};
535534
const workerError = (message: FlightWorkerMessage): Error => {
@@ -560,25 +559,30 @@ export const createFlightWorkerHost = (
560559
: `The MCP render runtime restarted; worker exited with code ${String(code)}.`,
561560
));
562561
});
562+
const settle = (id: number, request: PendingRender, error?: Error): void => {
563+
pending.delete(id);
564+
request.signal.removeEventListener('abort', request.abort);
565+
if (error === undefined) request.controller.close();
566+
else request.controller.error(error);
567+
};
563568
worker.on('message', (message: FlightWorkerMessage) => {
564569
const request = pending.get(message.id);
565570
if (request === undefined) return;
566571
if (message.type === 'progress') {
567-
void request.progress?.report(message.update as never);
572+
// A report the session no longer accepts (it completed or failed first)
573+
// settles a still-pending render once instead of rejecting unhandled.
574+
Promise.resolve(request.progress?.report(message.update as never)).catch((error: unknown) => {
575+
if (pending.get(message.id) !== request) return;
576+
worker.postMessage({ id: message.id, type: 'cancel' });
577+
settle(message.id, request, error instanceof Error ? error : new Error(String(error)));
578+
});
568579
return;
569580
}
570-
pending.delete(message.id);
571-
request.signal.removeEventListener('abort', request.abort);
572-
if (message.type === 'error') {
573-
request.reject(workerError(message));
581+
if (message.type === 'chunk') {
582+
request.controller.enqueue(message.bytes!);
574583
return;
575584
}
576-
request.resolve(new ReadableStream<Uint8Array>({
577-
start(controller) {
578-
controller.enqueue(message.bytes!);
579-
controller.close();
580-
},
581-
}));
585+
settle(message.id, request, message.type === 'error' ? workerError(message) : undefined);
582586
});
583587
const warmHost = createWarmFlightHost({
584588
artifactEpoch,
@@ -597,33 +601,45 @@ export const createFlightWorkerHost = (
597601
}
598602
const context = await agent();
599603
const id = ++sequence;
600-
return new Promise<ReadableStream<Uint8Array>>((resolve, reject) => {
601-
const abort = (): void => {
602-
worker.postMessage({ id, type: 'cancel' });
603-
pending.delete(id);
604-
reject(new DOMException('Agent render was aborted', 'AbortError'));
605-
};
606-
pending.set(id, { abort, ...(progress === undefined ? {} : { progress }), reject, resolve, signal });
607-
signal.addEventListener('abort', abort, { once: true });
608-
if (signal.aborted) {
609-
abort();
610-
return;
611-
}
612-
worker.postMessage({
613-
actor: context.actor,
614-
artifactEpoch: requestEpoch ?? artifactEpoch,
615-
host: context.host,
616-
id,
617-
invocation,
618-
lineage: context.lineage,
619-
plugin: context.plugin,
620-
requestInvocation: context.invocation,
621-
session: context.session,
622-
terminal: context.terminal,
623-
type: 'render',
624-
workspace: context.workspace,
625-
});
604+
let controller!: ReadableStreamDefaultController<Uint8Array>;
605+
const cancelRender = (): void => {
606+
worker.postMessage({ id, type: 'cancel' });
607+
pending.delete(id);
608+
};
609+
const abort = (): void => {
610+
cancelRender();
611+
controller.error(new DOMException('Agent render was aborted', 'AbortError'));
612+
};
613+
// A consumer that cancels the stream closes its controller; dropping the
614+
// pending entry keeps later worker chunks from reaching a closed stream.
615+
const stream = new ReadableStream<Uint8Array>({
616+
cancel: () => {
617+
cancelRender();
618+
signal.removeEventListener('abort', abort);
619+
},
620+
start: (opened) => { controller = opened; },
621+
});
622+
if (signal.aborted) {
623+
controller.error(new DOMException('Agent render was aborted', 'AbortError'));
624+
return stream;
625+
}
626+
pending.set(id, { abort, controller, ...(progress === undefined ? {} : { progress }), signal });
627+
signal.addEventListener('abort', abort, { once: true });
628+
worker.postMessage({
629+
actor: context.actor,
630+
artifactEpoch: requestEpoch ?? artifactEpoch,
631+
host: context.host,
632+
id,
633+
invocation,
634+
lineage: context.lineage,
635+
plugin: context.plugin,
636+
requestInvocation: context.invocation,
637+
session: context.session,
638+
terminal: context.terminal,
639+
type: 'render',
640+
workspace: context.workspace,
626641
});
642+
return stream;
627643
},
628644
},
629645
});

‎packages/agent-bundle/tests/entry-shell.test.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -688,8 +688,17 @@ it('generates the warm react-server Flight worker separately from the MCP dispat
688688
expect(source).toContain('route.module.inputSchema.parse(message.invocation.props.input)');
689689
expect(source).toContain('message.validateInput !== true ? { input: message.invocation.props.input');
690690
expect(source).toContain("createElement(Agent.Error, { code: 'invalid-input' }");
691+
// Flight bytes leave the worker chunk by chunk inside the request scope —
692+
// the same `chunk`/`end`/`error` transport the rendered CLI worker speaks —
693+
// so a Suspense fallback reaches the consumer before the render completes (#686).
694+
expect(source).toContain("parentPort.postMessage({ bytes, id: message.id, type: 'chunk' }, [bytes.buffer]);");
695+
expect(source).toContain("parentPort.postMessage({ id: message.id, type: 'end' });");
696+
expect(source).not.toContain('arrayBuffer');
697+
expect(source).not.toContain("type: 'complete'");
698+
expect(source.indexOf("type: 'chunk'")).toBeLessThan(source.indexOf("type: 'observed-render-finish'"));
699+
expect(source.indexOf("type: 'end'")).toBeGreaterThan(source.indexOf("type: 'observed-render-finish'"));
691700
expect(createHash('sha256').update(source).digest('hex')).toBe(
692-
'68b593d21fdf4aaa5c51d99cffb1a106773e50ef1837072bfb0774699719f98c',
701+
'f363a6abcf9002e413be930cd00e1514da975ddfb8ca58a70be07f392601a4f4',
693702
);
694703
expect(generate({
695704
artifactEpoch: 'route-fixture@1.2.3',

‎packages/agent-bundle/tests/generated-route-server.test.ts‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { Client } from '@modelcontextprotocol/client';
22
import { StdioClientTransport } from '@modelcontextprotocol/client/stdio';
33
import { spawn } from 'node:child_process';
4+
import { existsSync } from 'node:fs';
45
import { mkdir, mkdtemp, readFile, rm, stat, symlink, writeFile } from 'node:fs/promises';
56
import { tmpdir } from 'node:os';
67
import { dirname, join, resolve } from 'node:path';
@@ -876,6 +877,58 @@ it('emits MCP progress notifications only when a progress token is supplied', {
876877
}
877878
});
878879

880+
it('streams an authored Suspense fallback to the stdio client before the render completes (#686)', { retry: 2, timeout: 60_000 }, async () => {
881+
const root = await mkdtemp(join(tmpdir(), 'agent-bundle-generated-streaming-'));
882+
roots.push(root);
883+
await writeGeneratedProject(root, {
884+
'src/mcp/curator/tools/gated.tsx': [
885+
"import { existsSync } from 'node:fs';",
886+
"import { Agent } from '@agent-bundle/runtime';",
887+
"import { createElement, Suspense } from 'react';",
888+
"import { z } from 'zod';",
889+
"export const config = { annotations: { readOnlyHint: true }, description: 'Wait for a gate.' };",
890+
'export const inputSchema = z.object({ gate: z.string() }).strict();',
891+
'export const resultSchema = z.object({ ok: z.literal(true) }).strict();',
892+
'const Gated = async ({ gate, signal }) => {',
893+
' while (!existsSync(gate)) {',
894+
" if (signal.aborted) throw new DOMException('gate abandoned', 'AbortError');",
895+
' await new Promise((resolve) => setTimeout(resolve, 25));',
896+
' }',
897+
" return createElement(Agent.Text, null, 'released');",
898+
'};',
899+
'export default async function GatedTool({ input, signal }) {',
900+
" 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 })));",
901+
'}',
902+
'',
903+
].join('\n'),
904+
});
905+
const session = await connectGeneratedServer(root);
906+
const notifications: Array<{ readonly message?: string; readonly progress: number }> = [];
907+
session.client.setNotificationHandler('notifications/progress', (notification) => {
908+
notifications.push(notification.params);
909+
});
910+
try {
911+
const gate = join(root, 'gate.marker');
912+
const pending = session.client.callTool({
913+
arguments: { gate },
914+
name: 'gated',
915+
_meta: { progressToken: 'tok-gate' },
916+
}, { signal: AbortSignal.timeout(20_000) });
917+
// The fallback's progress node reaches the wire while the child is still
918+
// blocked: the worker streamed the shell instead of buffering the render.
919+
await expect.poll(() => notifications, { timeout: 15_000 }).toEqual([{ message: 'waiting', progress: 0, progressToken: 'tok-gate', total: 1 }]);
920+
expect(existsSync(gate)).toBe(false);
921+
await writeFile(gate, 'open\n');
922+
await expect(pending).resolves.toMatchObject({
923+
content: [{ text: 'released', type: 'text' }],
924+
structuredContent: { ok: true },
925+
});
926+
expect(notifications).toHaveLength(1);
927+
} finally {
928+
await session.close();
929+
}
930+
});
931+
879932
/**
880933
* #492: what a thrown (not represented) route error is on each MCP surface of
881934
* a real generated stdio server. A tool throw is the SDK's default tool error;

0 commit comments

Comments
 (0)