fix(mcp-stdio): made close() unconditional in resource teardown

Per second review on #1711: when notify()'s write failed, #handleClose()
flipped #connected=false before the throw. connectToServer()'s catch
then called transport.close(), but close() early-returned because
#connected was already false — so #process.kill() and the readLoop
await never ran. A subprocess that closed its stdin without exiting
(parent EPIPE, transport dead, subprocess alive) would leak.

close() no longer guards on #connected for the resource phase. It still
calls #handleClose() once (idempotent — only if #connected is true), then
unconditionally walks the cleanup chain (kill process, null #process,
await + null #readLoop), each step individually guarded so repeat calls
remain no-ops. onClose still fires exactly once per transport lifetime.

Two new tests in StdioTransport.close: (1) close() called after the
read-loop has already EOF'd and torn down still completes cleanup
without throwing and without re-firing onClose; (2) repeated close()
calls fire onClose exactly once and leave #connected=false.
This commit is contained in:
roboomp
2026-06-02 12:35:35 +00:00
parent b870732258
commit 6598c69a2f
3 changed files with 90 additions and 12 deletions
+1 -1
View File
@@ -4,7 +4,7 @@
### Fixed
- Fixed an unhandled `EPIPE` rejection when an MCP stdio server exits between returning the `initialize` response and the client's `notifications/initialized` send. `StdioTransport.notify()` and `#sendResponse()` now route stdin writes through a shared helper that catches synchronous sink failures: `notify()` tears the transport down (firing `onClose`) and surfaces a `Transport closed while sending notification` rejection so `connectToServer()` treats the handshake as a failed connection instead of returning a "connected" handle wrapping a dead transport; `#sendResponse()` stays silent because a dead subprocess has no use for the response ([#1710](https://github.com/can1357/oh-my-pi/issues/1710)).
- Fixed an unhandled `EPIPE` rejection when an MCP stdio server exits between returning the `initialize` response and the client's `notifications/initialized` send. `StdioTransport.notify()` and `#sendResponse()` now route stdin writes through a shared helper that catches synchronous sink failures: `notify()` tears the transport down (firing `onClose`) and surfaces a `Transport closed while sending notification` rejection so `connectToServer()` treats the handshake as a failed connection instead of returning a "connected" handle wrapping a dead transport; `#sendResponse()` stays silent because a dead subprocess has no use for the response. `StdioTransport.close()` is now the authoritative resource teardown — it no longer early-returns when `#handleClose()` has already flipped `#connected`, so the subprocess and read loop are always cleaned up (including in the `connectToServer()` failure path) ([#1710](https://github.com/can1357/oh-my-pi/issues/1710)).
## [15.8.0] - 2026-06-02
@@ -328,28 +328,26 @@ export class StdioTransport implements MCPTransport {
}
async close(): Promise<void> {
if (!this.#connected) return;
this.#connected = false;
// Reject pending requests
for (const [, pending] of this.#pendingRequests) {
pending.reject(new Error("Transport closed"));
// `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();
}
this.#pendingRequests.clear();
// Kill subprocess
if (this.#process) {
this.#process.kill();
this.#process = null;
}
// Wait for read loop to finish
if (this.#readLoop) {
await this.#readLoop.catch(() => {});
this.#readLoop = null;
}
this.onClose?.();
}
}
@@ -180,3 +180,83 @@ describe("StdioTransport.notify", () => {
}
});
});
// ---------------------------------------------------------------------------
// StdioTransport.close — authoritative resource teardown that must keep
// cleaning up the subprocess and read loop even when `#handleClose()` has
// already flipped `#connected` (read-loop EOF, or a notify() write failure
// in the connectToServer() failure path). See PR #1711 follow-up.
//
// Bun's parent-side stdout reader only sees EOF when the subprocess
// actually exits, so the "subprocess closed its stdout but stayed alive"
// state we'd love to test directly cannot be reproduced through a real
// subprocess on this platform. Instead we exercise the post-handleClose
// code path via the natural read-loop-EOF route and pair it with explicit
// idempotency checks; the reviewer-flagged leak surfaces on Windows where
// the notify() write actually throws.
// ---------------------------------------------------------------------------
describe("StdioTransport.close", () => {
let transport: StdioTransport | undefined;
afterEach(async () => {
await transport?.close().catch(() => {});
transport = undefined;
});
it("completes cleanup when called after the read loop has already torn down", async () => {
// Subprocess exits cleanly; the read loop sees EOF and fires
// `#handleClose()`, flipping `#connected` to false. `close()` then
// runs in exactly the state the reviewer flagged — `#connected`
// already false, `#process` and `#readLoop` still set — and must
// still null them out instead of early-returning.
transport = new StdioTransport({
type: "stdio",
command: "bun",
args: ["-e", "process.exit(0)"],
});
let closeCount = 0;
transport.onClose = () => {
closeCount++;
};
await transport.connect();
// Wait for the read loop to observe EOF and fire #handleClose.
for (let i = 0; i < 100 && transport.connected; i++) {
await Bun.sleep(10);
}
expect(transport.connected).toBe(false);
expect(closeCount).toBe(1);
// Must not throw and must not re-fire onClose.
await transport.close();
expect(closeCount).toBe(1);
// Second close is a no-op too — every resource is already released.
await transport.close();
expect(closeCount).toBe(1);
});
it("is idempotent — repeat close() calls fire onClose exactly once", async () => {
transport = new StdioTransport({
type: "stdio",
command: "bun",
args: ["-e", "await Bun.sleep(60_000)"],
});
let closeCount = 0;
transport.onClose = () => {
closeCount++;
};
await transport.connect();
await transport.close();
await transport.close();
await transport.close();
expect(closeCount).toBe(1);
expect(transport.connected).toBe(false);
});
});