Files
oh-my-pi/packages/coding-agent/src/mcp/transports/stdio.ts
T
VoidChecksum bab53df38a fix(coding-agent/mcp): handle async broken-pipe rejections in stdio transport
On Windows, Bun's FileSink surfaces a broken pipe as a *rejected Promise*
(the EPIPE arrives via a processTicksAndRejections tick), not a synchronous
throw. The sync-only `writeFrame()` from #1711 returns `true` for such a
write and lets the rejected promise float; `request()` likewise wrote stdin
without awaiting. Either path escapes as an unhandled rejection, which the
postmortem handler turns into a fatal process.exit(1) when an MCP server
exits mid-handshake.

- writeFrame: also neutralize a rejected Promise returned by write()/flush()
  (covers the notify + #sendResponse paths).
- request(): await write()/flush() so the rejection lands in the existing
  try/catch and rejects the request instead of floating.
- tests: cover the async-rejection path for writeFrame.

Follow-up to #1710 / #1711, which addressed only the synchronous throw. The
existing notify mid-handshake test fails on Windows without this change
because the real FileSink rejects asynchronously.
2026-06-03 10:51:42 +02:00

387 lines
11 KiB
TypeScript

/**
* MCP stdio transport.
*
* Implements JSON-RPC 2.0 over subprocess stdin/stdout.
* Messages are newline-delimited JSON.
*/
import { getProjectDir, readJsonl, Snowflake } from "@oh-my-pi/pi-utils";
import { type Subprocess, spawn } from "bun";
import type {
JsonRpcError,
JsonRpcMessage,
JsonRpcRequest,
JsonRpcResponse,
MCPRequestOptions,
MCPStdioServerConfig,
MCPTransport,
} from "../../mcp/types";
import { toJsonRpcError } from "../../mcp/types";
import { isMCPTimeoutEnabled, resolveMCPTimeoutMs } from "../timeout";
/** Minimal write surface of `Subprocess.stdin` we need for framed sends. */
interface FrameSink {
write(chunk: string): unknown;
flush(): unknown;
}
/** Narrow a value to a thenable so a rejection handler can be attached. */
function isThenable(value: unknown): value is PromiseLike<unknown> {
return (
value != null &&
(typeof value === "object" || typeof value === "function") &&
typeof (value as { then?: unknown }).then === "function"
);
}
/**
* Write a newline-delimited JSON-RPC frame to the subprocess's stdin sink,
* swallowing both synchronous throws and asynchronous rejections so the caller
* can decide how to react.
*
* Bun's `FileSink.write()`/`flush()` can fail two ways once the read end of the
* pipe has been closed by a subprocess that exited between read-loop ticks:
* - a synchronous throw (most reliably observed on Windows), and
* - a *rejected Promise* returned from `write()`/`flush()`, i.e. the EPIPE is
* surfaced asynchronously (note the `processTicksAndRejections` frame in the
* stack traces on #1710 and the follow-up report).
*
* A sibling `async` method's `try/catch` only catches the synchronous case; an
* un-awaited rejected Promise escapes as a fatal unhandled rejection. So we both
* catch the throw and neutralize any returned promise's rejection.
*
* Returns `true` when the frame was accepted synchronously, `false` when the
* sink threw — callers signal transport closure on `false`. An asynchronous
* failure cannot be reflected in the return value; it is neutralized here and
* the dead transport is detected by the read loop / request timeout instead.
*/
export function writeFrame(stdin: FrameSink, frame: string): boolean {
try {
const wrote = stdin.write(frame);
const flushed = stdin.flush();
if (isThenable(wrote)) wrote.then(undefined, () => {});
if (isThenable(flushed)) flushed.then(undefined, () => {});
return true;
} catch {
return false;
}
}
/**
* Stdio transport for MCP servers.
* Spawns a subprocess and communicates via stdin/stdout.
*/
export class StdioTransport implements MCPTransport {
#process: Subprocess<"pipe", "pipe", "pipe"> | null = null;
#pendingRequests = new Map<
string | number,
{
resolve: (value: unknown) => void;
reject: (error: Error) => void;
}
>();
#connected = false;
#readLoop: Promise<void> | null = null;
onClose?: () => void;
onError?: (error: Error) => void;
onNotification?: (method: string, params: unknown) => void;
onRequest?: (method: string, params: unknown) => Promise<unknown>;
constructor(private config: MCPStdioServerConfig) {}
get connected(): boolean {
return this.#connected;
}
/**
* Start the subprocess and begin reading.
*/
async connect(): Promise<void> {
if (this.#connected) return;
const args = this.config.args ?? [];
const env = {
...Bun.env,
...this.config.env,
};
this.#process = spawn({
cmd: [this.config.command, ...args],
cwd: this.config.cwd ?? getProjectDir(),
env,
stdin: "pipe",
stdout: "pipe",
stderr: "pipe",
});
this.#connected = true;
// Start reading stdout
this.#readLoop = this.#startReadLoop();
// Log stderr for debugging
this.#startStderrLoop();
}
async #startReadLoop(): Promise<void> {
if (!this.#process?.stdout) return;
try {
for await (const line of readJsonl(this.#process.stdout)) {
if (!this.#connected) break;
try {
this.#handleMessage(line as JsonRpcMessage);
} catch {
// Skip malformed lines
}
}
} catch (error) {
if (this.#connected) {
this.onError?.(error instanceof Error ? error : new Error(String(error)));
}
} finally {
this.#handleClose();
}
}
async #startStderrLoop(): Promise<void> {
if (!this.#process?.stderr) return;
const reader = this.#process.stderr.getReader();
const decoder = new TextDecoder();
try {
while (this.#connected) {
const { done, value } = await reader.read();
if (done) break;
// Log stderr but don't treat as error - servers use it for logging
const text = decoder.decode(value, { stream: true });
if (text.trim()) {
// Could expose via onStderr callback if needed
// For now, silent - MCP spec says clients MAY capture/ignore
}
}
} catch {
// Ignore stderr read errors
} finally {
reader.releaseLock();
}
}
#handleMessage(message: JsonRpcMessage | JsonRpcMessage[]): void {
if (Array.isArray(message)) {
for (const m of message) this.#handleMessage(m);
return;
}
// Server-to-client request: has both method and id
if ("method" in message && "id" in message && message.id != null) {
void this.#handleServerRequest(message as JsonRpcRequest);
return;
}
// Response to our request: has id
if ("id" in message && message.id != null) {
const response = message as JsonRpcResponse;
const pending = this.#pendingRequests.get(response.id);
if (pending) {
this.#pendingRequests.delete(response.id);
if (response.error) {
pending.reject(new Error(`MCP error ${response.error.code}: ${response.error.message}`));
} else {
pending.resolve(response.result);
}
}
return;
}
// Notification: has method but no id
if ("method" in message) {
const notification = message as { method: string; params?: unknown };
this.onNotification?.(notification.method, notification.params);
}
}
async #handleServerRequest(request: JsonRpcRequest): Promise<void> {
try {
if (!this.onRequest) {
this.#sendResponse(request.id, undefined, { code: -32601, message: "Method not found" });
return;
}
const result = await this.onRequest(request.method, request.params);
this.#sendResponse(request.id, result);
} catch (error) {
this.#sendResponse(request.id, undefined, toJsonRpcError(error));
}
}
#sendResponse(id: string | number, result?: unknown, error?: JsonRpcError): void {
if (!this.#connected || !this.#process?.stdin) return;
const response = error
? { jsonrpc: "2.0" as const, id, error }
: { jsonrpc: "2.0" as const, id, result: result ?? {} };
// Silent on failure — a dead subprocess has no use for the response,
// and the read loop will close the transport on EOF.
writeFrame(this.#process.stdin, `${JSON.stringify(response)}\n`);
}
#handleClose(): void {
if (!this.#connected) return;
this.#connected = false;
// Reject all pending requests
for (const [, pending] of this.#pendingRequests) {
pending.reject(new Error("Transport closed"));
}
this.#pendingRequests.clear();
this.onClose?.();
}
async request<T = unknown>(
method: string,
params?: Record<string, unknown>,
options?: MCPRequestOptions,
): Promise<T> {
if (!this.#connected || !this.#process?.stdin) {
throw new Error("Transport not connected");
}
const id = Snowflake.next();
const request = {
jsonrpc: "2.0" as const,
id,
method,
params: params ?? {},
};
const timeout = resolveMCPTimeoutMs(this.config.timeout);
const signal = options?.signal;
if (signal?.aborted) {
const reason = signal.reason instanceof Error ? signal.reason : new Error("Aborted");
return Promise.reject(reason);
}
const { promise, resolve, reject } = Promise.withResolvers<T>();
let timer: NodeJS.Timeout | undefined;
let settled = false;
const cleanup = () => {
if (settled) return;
settled = true;
if (timer) {
clearTimeout(timer);
timer = undefined;
}
if (signal) {
signal.removeEventListener("abort", onAbort);
}
this.#pendingRequests.delete(id);
};
const onAbort = () => {
cleanup();
const reason = signal?.reason instanceof Error ? signal.reason : new Error("Aborted");
reject(reason);
};
if (signal) {
signal.addEventListener("abort", onAbort, { once: true });
}
this.#pendingRequests.set(id, {
resolve: (value: unknown) => {
cleanup();
resolve(value as T);
},
reject: (error: Error) => {
cleanup();
reject(error);
},
});
if (isMCPTimeoutEnabled(timeout)) {
timer = setTimeout(() => {
cleanup();
reject(new Error(`Request timeout after ${timeout}ms`));
}, timeout);
}
const stdin = this.#process.stdin;
const message = `${JSON.stringify(request)}\n`;
try {
// Await both: Bun's FileSink can surface a broken pipe either as a
// synchronous throw or as a rejected Promise (the EPIPE arrives on a
// processTicksAndRejections tick). Awaiting funnels both into this catch
// so the request rejects cleanly instead of leaving a floating rejected
// promise that crashes the process via the unhandledRejection handler.
await stdin.write(message);
await stdin.flush();
} catch (error: unknown) {
cleanup();
reject(error instanceof Error ? error : new Error(String(error)));
}
return promise;
}
async notify(method: string, params?: Record<string, unknown>): Promise<void> {
if (!this.#connected || !this.#process?.stdin) {
throw new Error("Transport not connected");
}
const notification = {
jsonrpc: "2.0" as const,
method,
params: params ?? {},
};
// Bun's FileSink can throw EPIPE synchronously on Windows when the
// subprocess has exited between the last read-loop tick and this
// write (e.g. an MCP server that dies after returning `initialize`
// but before `notifications/initialized` is delivered). Tear the
// transport down so any wired `onClose` (and reconnect machinery)
// engages, then surface the failure to the caller so a write that
// dropped on the floor is never silently treated as delivered —
// `initializeConnection()` runs before the manager installs its
// `onClose` handler, so a swallowed failure there would yield a
// "connected" handle wrapping a dead transport. See #1710.
if (!writeFrame(this.#process.stdin, `${JSON.stringify(notification)}\n`)) {
this.#handleClose();
throw new Error(`Transport closed while sending notification "${method}"`);
}
}
async close(): Promise<void> {
// `close()` is the authoritative resource teardown. `#handleClose()`
// may have already run (read-loop EOF, or a notify() write failure
// that surfaces the dead transport to the caller) and flipped
// `#connected` to false — but the subprocess and read loop are still
// alive in that path, so we MUST keep cleaning up regardless. Each
// step is individually guarded so this remains idempotent across
// repeat calls.
if (this.#connected) {
this.#handleClose();
}
if (this.#process) {
this.#process.kill();
this.#process = null;
}
if (this.#readLoop) {
await this.#readLoop.catch(() => {});
this.#readLoop = null;
}
}
}
/**
* Create and connect a stdio transport.
*/
export async function createStdioTransport(config: MCPStdioServerConfig): Promise<StdioTransport> {
const transport = new StdioTransport(config);
await transport.connect();
return transport;
}