Skip to content
Draft
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
11 changes: 11 additions & 0 deletions .changeset/stdio-server-stdout-epipe.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
---
'@modelcontextprotocol/sdk': patch
---

Handle `stdout` errors in `StdioServerTransport` instead of crashing the process. When a stdio client disconnects, the server's next write fails with `EPIPE`, and because nothing listened for `'error'` on `stdout` that became an unhandled `'error'` event, which takes the whole
Node process down. The transport now listens for it, reports it through `onerror`, and closes itself, so a client that goes away ends the session rather than killing the server.

`send()` no longer waits for a `'drain'` event that a destroyed stream can never emit. It settles from an error listener armed before the write and removes both listeners on every exit path, so a write that fails rejects with the underlying error instead of leaving the promise
pending, and `send()` after `close()` rejects rather than writing to a stream the transport has let go of. `close()` is idempotent too, so an error-triggered close followed by an explicit one fires `onclose` once.

This backports the fix that shipped for the v2 line, so both lines now behave the same way on a disconnecting client.
45 changes: 42 additions & 3 deletions src/server/stdio.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import { Transport } from '../shared/transport.js';
export class StdioServerTransport implements Transport {
private _readBuffer: ReadBuffer;
private _started = false;
private _closed = false;

constructor(
private _stdin: Readable = process.stdin,
Expand Down Expand Up @@ -46,6 +47,12 @@ export class StdioServerTransport implements Transport {
_onerror = (error: Error) => {
this.onerror?.(error);
};
_onstdouterror = (error: Error) => {
this.onerror?.(error);
this.close().catch(() => {
// Ignore errors during close — we're already in an error path
});
};

/**
* Starts listening for messages on stdin.
Expand All @@ -60,6 +67,7 @@ export class StdioServerTransport implements Transport {
this._started = true;
this._stdin.on('data', this._ondata);
this._stdin.on('error', this._onerror);
this._stdout.on('error', this._onstdouterror);
}

private processReadBuffer() {
Expand All @@ -78,9 +86,15 @@ export class StdioServerTransport implements Transport {
}

async close(): Promise<void> {
if (this._closed) {
return;
}
this._closed = true;

// Remove our event listeners first
this._stdin.off('data', this._ondata);
this._stdin.off('error', this._onerror);
this._stdout.off('error', this._onstdouterror);

// Check if we were the only data listener
const remainingDataListeners = this._stdin.listenerCount('data');
Expand All @@ -96,12 +110,37 @@ export class StdioServerTransport implements Transport {
}

send(message: JSONRPCMessage): Promise<void> {
return new Promise(resolve => {
if (this._closed) {
return Promise.reject(new Error('StdioServerTransport is closed'));
}
return new Promise((resolve, reject) => {
const json = serializeMessage(message);

let settled = false;
const onError = (error: Error) => {
if (settled) return;
settled = true;
this._stdout.off('error', onError);
this._stdout.off('drain', onDrain);
reject(error);
};
const onDrain = () => {
if (settled) return;
settled = true;
this._stdout.off('error', onError);
this._stdout.off('drain', onDrain);
resolve();
};

this._stdout.once('error', onError);

if (this._stdout.write(json)) {
if (settled) return;
settled = true;
this._stdout.off('error', onError);
resolve();
} else {
this._stdout.once('drain', resolve);
} else if (!settled) {
this._stdout.once('drain', onDrain);
}
});
}
Expand Down
77 changes: 77 additions & 0 deletions test/server/stdio.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -148,3 +148,80 @@ test('should fire onerror and close when ReadBuffer overflows', async () => {
expect(receivedError?.message).toMatch(/ReadBuffer exceeded maximum size/);
expect(closeCount).toBe(1);
});

test('should close and fire onerror when stdout errors', async () => {
const server = new StdioServerTransport(input, output);

let receivedError: Error | undefined;
server.onerror = err => {
receivedError = err;
};
let closeCount = 0;
server.onclose = () => {
closeCount++;
};

await server.start();
output.emit('error', new Error('EPIPE'));

expect(receivedError?.message).toBe('EPIPE');
expect(closeCount).toBe(1);
});

test('should fire onerror before onclose on stdout error', async () => {
const server = new StdioServerTransport(input, output);

const events: string[] = [];
server.onerror = () => events.push('error');
server.onclose = () => events.push('close');

await server.start();
output.emit('error', new Error('EPIPE'));

expect(events).toEqual(['error', 'close']);
});

test('should not fire onclose twice when close() is called after stdout error', async () => {
const server = new StdioServerTransport(input, output);
server.onerror = () => {};

let closeCount = 0;
server.onclose = () => {
closeCount++;
};

await server.start();
output.emit('error', new Error('EPIPE'));
await server.close();

expect(closeCount).toBe(1);
});

test('should reject send() when stdout errors before drain', async () => {
let completeWrite: ((error?: Error | null) => void) | undefined;
const slowOutput = new Writable({
highWaterMark: 0,
write(_chunk, _encoding, callback) {
completeWrite = callback;
}
});

const server = new StdioServerTransport(input, slowOutput);
server.onerror = () => {};
await server.start();

const sendPromise = server.send({ jsonrpc: '2.0', id: 1, method: 'ping' });
completeWrite!(new Error('write EPIPE'));

await expect(sendPromise).rejects.toThrow('write EPIPE');
expect(slowOutput.listenerCount('drain')).toBe(0);
expect(slowOutput.listenerCount('error')).toBe(0);
});

test('should reject send() after transport is closed', async () => {
const server = new StdioServerTransport(input, output);
await server.start();
await server.close();

await expect(server.send({ jsonrpc: '2.0', id: 1, method: 'ping' })).rejects.toThrow('closed');
});
Loading