Skip to content
Open
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
7 changes: 7 additions & 0 deletions .changeset/mcp-abort-releases-response-body.md
Original file line number Diff line number Diff line change
@@ -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.
64 changes: 64 additions & 0 deletions packages/plugins/mcp/src/sdk/connection-socket-release.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>((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);
}),
),
);
});
21 changes: 20 additions & 1 deletion packages/plugins/mcp/src/sdk/connection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> =>
Effect.callback<void>((resume) => {
if (signal.aborted) {
resume(Effect.void);
return;
}
signal.addEventListener("abort", () => resume(Effect.void), { once: true });
});

const fetchFromHttpClientLayer = (
httpClientLayer: Layer.Layer<HttpClient.HttpClient>,
): FetchLike => {
Expand All @@ -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,
Expand Down
12 changes: 12 additions & 0 deletions packages/plugins/mcp/src/testing/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 McpTestRequest[]>;
readonly clearRequests: Effect.Effect<void>;
/** Drops all server-side session registrations without notifying clients. */
Expand Down Expand Up @@ -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));
});

Expand All @@ -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()),
Expand Down