diff --git a/.changeset/stdio-server-stdout-epipe.md b/.changeset/stdio-server-stdout-epipe.md new file mode 100644 index 0000000000..a2ae63c37e --- /dev/null +++ b/.changeset/stdio-server-stdout-epipe.md @@ -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. diff --git a/src/server/stdio.ts b/src/server/stdio.ts index ced465bb81..d43f66d6c7 100644 --- a/src/server/stdio.ts +++ b/src/server/stdio.ts @@ -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, @@ -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. @@ -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() { @@ -78,9 +86,15 @@ export class StdioServerTransport implements Transport { } async close(): Promise { + 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'); @@ -96,12 +110,37 @@ export class StdioServerTransport implements Transport { } send(message: JSONRPCMessage): Promise { - 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); } }); } diff --git a/test/server/stdio.test.ts b/test/server/stdio.test.ts index 9613e0ec48..897f17358b 100644 --- a/test/server/stdio.test.ts +++ b/test/server/stdio.test.ts @@ -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'); +});