diff --git a/.changeset/mcp-abort-releases-response-body.md b/.changeset/mcp-abort-releases-response-body.md new file mode 100644 index 0000000000..00cdfe5d80 --- /dev/null +++ b/.changeset/mcp-abort-releases-response-body.md @@ -0,0 +1,7 @@ +--- +"@executor-js/plugin-mcp": patch +--- + +**Closing a remote MCP connection now ends its streamable-http SSE request** + +On a supplied `httpClientLayer`, the fetch adapter wired the caller's `AbortSignal` only to the pending response promise, never to the response body, so closing a connection left the long-lived `GET` channel in flight forever — one abandoned request per dial. Under Bun each holds one of the 256 concurrent-request slots, so a long-running process eventually exhausts the pool and every connection starts failing with `MCP discovery timed out after 15000ms`. The response stream is now interrupted when the signal aborts. diff --git a/packages/plugins/mcp/src/sdk/connection-socket-release.test.ts b/packages/plugins/mcp/src/sdk/connection-socket-release.test.ts new file mode 100644 index 0000000000..fa9b0b06fd --- /dev/null +++ b/packages/plugins/mcp/src/sdk/connection-socket-release.test.ts @@ -0,0 +1,64 @@ +import { describe, expect, it } from "@effect/vitest"; +import { Effect } from "effect"; +import { FetchHttpClient } from "effect/unstable/http"; + +import { createMcpConnector } from "./connection"; +import { makeEchoMcpServer, serveMcpServer } from "../testing"; + +// Raw timers, not Effect.sleep: `it.effect`'s TestClock never advances one. +const settle = (attempts: number, done: () => boolean) => + Effect.promise(async () => { + for (let i = 0; i < attempts; i++) { + if (done()) return; + await new Promise((resolve) => setTimeout(resolve, 25)); + } + }); + +const connect = (endpoint: string) => + createMcpConnector({ + transport: "remote", + endpoint, + remoteTransport: "streamable-http", + // Without a layer the SDK uses global fetch and skips the adapter. + httpClientLayer: FetchHttpClient.layer, + }); + +// `close()` has to end the request, not just drop the session: an SSE `GET` +// left in flight per dial exhausts Bun's 256-request pool, after which every +// outbound fetch queues forever. +describe("MCP connection request release", () => { + it.effect("closing a connection ends the SSE request", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveMcpServer(() => makeEchoMcpServer()); + + const connection = yield* connect(server.endpoint); + // The SSE `GET` opens asynchronously, after the handshake POSTs. + yield* settle(40, () => server.inFlightRequests() > 0); + expect(server.inFlightRequests()).toBeGreaterThan(0); + + yield* Effect.promise(() => connection.close()); + yield* settle(40, () => server.inFlightRequests() === 0); + + expect(server.inFlightRequests()).toBe(0); + }), + ), + ); + + it.effect("repeated connect/close does not accumulate open requests", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveMcpServer(() => makeEchoMcpServer()); + + for (let i = 0; i < 5; i++) { + const connection = yield* connect(server.endpoint); + yield* Effect.promise(() => connection.close()); + } + yield* settle(40, () => server.inFlightRequests() === 0); + + expect(server.sessionCount()).toBe(5); + expect(server.inFlightRequests()).toBe(0); + }), + ), + ); +}); diff --git a/packages/plugins/mcp/src/sdk/connection.ts b/packages/plugins/mcp/src/sdk/connection.ts index 5ac15669c6..a8fbb91101 100644 --- a/packages/plugins/mcp/src/sdk/connection.ts +++ b/packages/plugins/mcp/src/sdk/connection.ts @@ -123,6 +123,16 @@ const abortError = (signal: AbortSignal): unknown => { return error; }; +/** An effect that completes when `signal` aborts (already-aborted = now). */ +const awaitAbort = (signal: AbortSignal): Effect.Effect => + Effect.callback((resume) => { + if (signal.aborted) { + resume(Effect.void); + return; + } + signal.addEventListener("abort", () => resume(Effect.void), { once: true }); + }); + const fetchFromHttpClientLayer = ( httpClientLayer: Layer.Layer, ): FetchLike => { @@ -139,10 +149,19 @@ const fetchFromHttpClientLayer = ( for (const [key, value] of Object.entries(response.headers)) { if (value !== undefined) responseHeaders.set(key, value); } + // Abort must reach the body, not just the pending request: this stream + // fiber outlives the `runPromise` below, so without this streamable + // http's SSE `GET` stays in flight after `close()`. Interrupted at the + // source because the SDK holds a locked reader on that same stream, + // which rules out cancelling the ReadableStream. + const stream = + init?.signal == null + ? response.stream + : Stream.interruptWhen(response.stream, awaitAbort(init.signal)); const body = response.status === 204 || response.status === 205 || response.status === 304 ? null - : Stream.toReadableStream(response.stream); + : Stream.toReadableStream(stream); return new Response(body, { status: response.status, headers: responseHeaders, diff --git a/packages/plugins/mcp/src/testing/server.ts b/packages/plugins/mcp/src/testing/server.ts index 0b6edf0f04..2836aa719c 100644 --- a/packages/plugins/mcp/src/testing/server.ts +++ b/packages/plugins/mcp/src/testing/server.ts @@ -10,6 +10,8 @@ export type McpTestServer = { readonly endpoint: string; /** Number of MCP sessions created (each connect = 1 session) */ readonly sessionCount: () => number; + /** Requests the server has accepted and not yet finished answering. */ + readonly inFlightRequests: () => number; readonly requests: Effect.Effect; readonly clearRequests: Effect.Effect; /** Drops all server-side session registrations without notifying clients. */ @@ -187,7 +189,16 @@ export const serveMcpServer = (factory: () => McpServer, options: McpTestServerO ), ); + // An abandoned SSE `GET` leaves the session gone but the request open, + // which `sessionCount` cannot see and socket counting cannot either + // (keep-alive holds idle sockets open regardless). + let inFlight = 0; + const nodeServer = http.createServer((request, response) => { + inFlight += 1; + response.once("close", () => { + inFlight -= 1; + }); void Effect.runPromise(handleMcpRequest(request, response)); }); @@ -214,6 +225,7 @@ export const serveMcpServer = (factory: () => McpServer, options: McpTestServerO url: endpoint, endpoint, sessionCount: () => sessions, + inFlightRequests: () => inFlight, requests: Ref.get(requests), clearRequests: Ref.set(requests, []), forgetSessions: Effect.sync(() => transports.clear()),