From 6dd658c7d558d84d745a15639f7095636b5d72cc Mon Sep 17 00:00:00 2001 From: utpal singh Date: Fri, 18 Sep 2026 23:29:05 +0530 Subject: [PATCH 1/2] fix: keep quiet MCP streams alive --- .changeset/quiet-mcp-sse-keepalive.md | 5 + apps/local/src/mcp-sse-heartbeat.test.ts | 141 +++++++++++++ apps/local/src/mcp.ts | 5 +- e2e/local/mcp-standalone-get.test.ts | 95 +++++++++ packages/hosts/mcp/package.json | 4 + packages/hosts/mcp/src/envelope.test.ts | 28 +++ .../mcp/src/in-memory-session-store.test.ts | 114 +++++++++++ .../hosts/mcp/src/in-memory-session-store.ts | 2 + packages/hosts/mcp/src/sse-heartbeat.test.ts | 189 ++++++++++++++++++ packages/hosts/mcp/src/sse-heartbeat.ts | 101 ++++++++++ 10 files changed, 682 insertions(+), 2 deletions(-) create mode 100644 .changeset/quiet-mcp-sse-keepalive.md create mode 100644 apps/local/src/mcp-sse-heartbeat.test.ts create mode 100644 e2e/local/mcp-standalone-get.test.ts create mode 100644 packages/hosts/mcp/src/sse-heartbeat.test.ts create mode 100644 packages/hosts/mcp/src/sse-heartbeat.ts diff --git a/.changeset/quiet-mcp-sse-keepalive.md b/.changeset/quiet-mcp-sse-keepalive.md new file mode 100644 index 0000000000..c4739c07f2 --- /dev/null +++ b/.changeset/quiet-mcp-sse-keepalive.md @@ -0,0 +1,5 @@ +--- +"executor": patch +--- + +Keep quiet MCP GET streams alive with SSE comment heartbeats so Bun clients do not time out. diff --git a/apps/local/src/mcp-sse-heartbeat.test.ts b/apps/local/src/mcp-sse-heartbeat.test.ts new file mode 100644 index 0000000000..00488d0d13 --- /dev/null +++ b/apps/local/src/mcp-sse-heartbeat.test.ts @@ -0,0 +1,141 @@ +// --------------------------------------------------------------------------- +// Quiet standalone GET `/mcp` under Bun — regression for #1983 +// --------------------------------------------------------------------------- +// +// After initialize, the SDK GET stream is a 200 `text/event-stream` with no +// body bytes until a server-initiated message. Bun fetch does not settle a +// silent stream even when `Bun.serve` has `idleTimeout: 0`. The handler must +// emit a legal SSE comment so the GET resolves promptly, then drop the +// upstream stream on cancel so a reconnect is not 409. +// --------------------------------------------------------------------------- + +import { describe, expect, it } from "@effect/vitest"; +import { Effect } from "effect"; + +import type { ExecutionEngine } from "@executor-js/execution"; + +import { createMcpRequestHandler } from "./mcp"; + +const MCP_POST_HEADERS = { + "content-type": "application/json", + accept: "application/json, text/event-stream", +} as const; + +const stubEngine: ExecutionEngine = { + execute: () => Effect.succeed({ result: "unused" }), + executeWithPause: () => Effect.succeed({ status: "completed", result: { result: "unused" } }), + resume: () => Effect.succeed(null), + getPausedExecution: () => Effect.succeed(null), + pausedExecutionCount: () => Effect.succeed(0), + hasPausedExecutions: () => Effect.succeed(false), + getDescription: Effect.succeed("test executor"), + shutdown: Effect.void, +}; + +const initializeBody = { + jsonrpc: "2.0", + id: 1, + method: "initialize", + params: { + protocolVersion: "2025-06-18", + capabilities: {}, + clientInfo: { name: "bun-sse-test", version: "1.0.0" }, + }, +}; + +const openLiveSession = async ( + origin: string, + path = "/mcp", +): Promise<{ readonly sessionId: string }> => { + const init = await fetch(`${origin}${path}`, { + method: "POST", + headers: MCP_POST_HEADERS, + body: JSON.stringify(initializeBody), + }); + expect(init.status).toBe(200); + const sessionId = init.headers.get("mcp-session-id"); + expect(sessionId).toBeTruthy(); + await init.body?.cancel(); + + const initialized = await fetch(`${origin}${path}`, { + method: "POST", + headers: { ...MCP_POST_HEADERS, "mcp-session-id": sessionId! }, + body: JSON.stringify({ jsonrpc: "2.0", method: "notifications/initialized" }), + }); + expect(initialized.status).toBe(202); + await initialized.body?.cancel(); + return { sessionId: sessionId! }; +}; + +describe("local MCP handler, quiet GET stream", () => { + it("resolves a live-session GET under Bun with a keepalive comment and reconnects after cancel", async () => { + const handler = createMcpRequestHandler({ engine: stubEngine }); + const server = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch: (request) => handler.handleRequest(request), + }); + + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: always stop the server + try { + const origin = `http://127.0.0.1:${server.port}`; + const { sessionId } = await openLiveSession(origin); + + const get = await fetch(`${origin}/mcp`, { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": sessionId }, + signal: AbortSignal.timeout(1_000), + }); + expect(get.status).toBe(200); + expect(get.headers.get("content-type")).toContain("text/event-stream"); + expect(get.headers.get("mcp-session-id")).toBe(sessionId); + + const reader = get.body!.getReader(); + const first = await reader.read(); + expect(new TextDecoder().decode(first.value)).toBe(": keepalive\n\n"); + await reader.cancel(); + + const reconnect = await fetch(`${origin}/mcp`, { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": sessionId }, + signal: AbortSignal.timeout(1_000), + }); + expect(reconnect.status).toBe(200); + expect(reconnect.headers.get("content-type")).toContain("text/event-stream"); + await reconnect.body?.cancel(); + } finally { + server.stop(true); + await handler.close(); + } + }); + + it("keeps the same GET contract on a toolkit MCP path", async () => { + const handler = createMcpRequestHandler({ engine: stubEngine }); + const server = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch: (request) => handler.handleRequest(request), + }); + + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: always stop the server + try { + const origin = `http://127.0.0.1:${server.port}`; + const path = "/mcp/toolkits/deploy"; + const { sessionId } = await openLiveSession(origin, path); + + const get = await fetch(`${origin}${path}`, { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": sessionId }, + signal: AbortSignal.timeout(1_000), + }); + expect(get.status).toBe(200); + expect(get.headers.get("content-type")).toContain("text/event-stream"); + const reader = get.body!.getReader(); + expect(new TextDecoder().decode((await reader.read()).value)).toBe(": keepalive\n\n"); + await reader.cancel(); + } finally { + server.stop(true); + await handler.close(); + } + }); +}); diff --git a/apps/local/src/mcp.ts b/apps/local/src/mcp.ts index 8f6181bbb9..908deb9f4e 100644 --- a/apps/local/src/mcp.ts +++ b/apps/local/src/mcp.ts @@ -10,6 +10,7 @@ import { preInitializeMethodNotFound, type McpResource, } from "@executor-js/host-mcp"; +import { withMcpSseHeartbeat } from "@executor-js/host-mcp/sse-heartbeat"; import { createExecutorMcpServer, type ExecutorMcpServerConfig, @@ -196,7 +197,7 @@ export const createMcpRequestHandler = ( if (!sessionResource || mcpResourceKey(sessionResource) !== mcpResourceKey(resource)) { return jsonError(403, -32003, "Session belongs to a different MCP resource"); } - return transport.handleRequest(request); + return withMcpSseHeartbeat(request, await transport.handleRequest(request)); } // Pre-initialize dispatch: only `initialize` opens a session here, so a @@ -263,7 +264,7 @@ export const createMcpRequestHandler = ( }), ); await created.connect(transport); - const response = await transport.handleRequest(request); + const response = withMcpSseHeartbeat(request, await transport.handleRequest(request)); if (!transport.sessionId) { await ignoreClose(() => transport.close()); diff --git a/e2e/local/mcp-standalone-get.test.ts b/e2e/local/mcp-standalone-get.test.ts new file mode 100644 index 0000000000..51671f8600 --- /dev/null +++ b/e2e/local/mcp-standalone-get.test.ts @@ -0,0 +1,95 @@ +// Black-box regression for #1983: a quiet Streamable HTTP GET `/mcp` on the +// local daemon must resolve under Bun with a legal SSE comment, then release +// the standalone stream on cancel so a reconnect is not 409. +import { expect } from "@effect/vitest"; +import { Effect } from "effect"; + +import { scenario } from "../src/scenario"; +import { Cli, RunDir } from "../src/services"; +import { withLocalServer } from "./local-server"; + +const MCP_POST_HEADERS = { + "content-type": "application/json", + accept: "application/json, text/event-stream", +} as const; + +scenario( + "Local · a quiet MCP GET stream stays alive with an SSE keepalive comment", + { timeout: 300_000 }, + Effect.gen(function* () { + const cli = yield* Cli; + const runDir = yield* RunDir; + + yield* withLocalServer(cli, runDir, (server) => + Effect.gen(function* () { + const auth = { authorization: `Bearer ${server.token}` }; + + const init = yield* Effect.promise(() => + fetch(`${server.origin}/mcp`, { + method: "POST", + headers: { ...MCP_POST_HEADERS, ...auth }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: 1, + method: "initialize", + params: { + protocolVersion: "2025-06-18", + capabilities: {}, + clientInfo: { name: "e2e-local-sse-keepalive", version: "1.0.0" }, + }, + }), + }), + ); + expect(init.status).toBe(200); + const sessionId = init.headers.get("mcp-session-id"); + expect(sessionId).toBeTruthy(); + yield* Effect.promise(() => init.body?.cancel() ?? Promise.resolve()); + + const initialized = yield* Effect.promise(() => + fetch(`${server.origin}/mcp`, { + method: "POST", + headers: { ...MCP_POST_HEADERS, ...auth, "mcp-session-id": sessionId! }, + body: JSON.stringify({ jsonrpc: "2.0", method: "notifications/initialized" }), + }), + ); + expect(initialized.status).toBe(202); + yield* Effect.promise(() => initialized.body?.cancel() ?? Promise.resolve()); + + const get = yield* Effect.promise(() => + fetch(`${server.origin}/mcp`, { + method: "GET", + headers: { + ...auth, + accept: "text/event-stream", + "mcp-session-id": sessionId!, + }, + signal: AbortSignal.timeout(5_000), + }), + ); + expect(get.status).toBe(200); + expect(get.headers.get("content-type")).toContain("text/event-stream"); + expect(get.headers.get("mcp-session-id")).toBe(sessionId); + + const reader = get.body!.getReader(); + const first = yield* Effect.promise(() => reader.read()); + expect(new TextDecoder().decode(first.value)).toBe(": keepalive\n\n"); + yield* Effect.promise(() => reader.cancel()); + + const reconnect = yield* Effect.promise(() => + fetch(`${server.origin}/mcp`, { + method: "GET", + headers: { + ...auth, + accept: "text/event-stream", + "mcp-session-id": sessionId!, + }, + signal: AbortSignal.timeout(5_000), + }), + ); + expect(reconnect.status).toBe(200); + expect(reconnect.headers.get("content-type")).toContain("text/event-stream"); + yield* Effect.promise(() => reconnect.body?.cancel() ?? Promise.resolve()); + }), + ); + }), +); diff --git a/packages/hosts/mcp/package.json b/packages/hosts/mcp/package.json index 1e8dd7c4fc..a82fd9983f 100644 --- a/packages/hosts/mcp/package.json +++ b/packages/hosts/mcp/package.json @@ -28,6 +28,10 @@ "types": "./src/in-memory-session-store.ts", "default": "./src/in-memory-session-store.ts" }, + "./sse-heartbeat": { + "types": "./src/sse-heartbeat.ts", + "default": "./src/sse-heartbeat.ts" + }, "./browser-approval": { "types": "./src/browser-approval.ts", "default": "./src/browser-approval.ts" diff --git a/packages/hosts/mcp/src/envelope.test.ts b/packages/hosts/mcp/src/envelope.test.ts index 072f98fe11..0d61b7604c 100644 --- a/packages/hosts/mcp/src/envelope.test.ts +++ b/packages/hosts/mcp/src/envelope.test.ts @@ -225,6 +225,34 @@ it("dispatches toolkit MCP routes with the parsed toolkit resource", async () => }); }); +it("forwards a live GET SSE body without injecting keepalive comments", async () => { + const StoreLive = Layer.succeed(McpSessionStore)({ + dispatch: (): Effect.Effect => + Effect.succeed( + new Response("data: already-kept-alive\n\n", { + status: 200, + headers: { "content-type": "text/event-stream", "mcp-session-id": "s1" }, + }), + ), + dispose: () => Effect.void, + }); + + const handler = buildHandler(StoreLive, McpErrorReporterNoop); + const response = await handler( + new Request("https://host.test/mcp", { + method: "GET", + headers: { + authorization: "Bearer x", + accept: "text/event-stream", + "mcp-session-id": "s1", + }, + }), + ); + + expect(response.status).toBe(200); + expect(await response.text()).toBe("data: already-kept-alive\n\n"); +}); + // --------------------------------------------------------------------------- // The pre-initialize dispatch guard. Session-less, only `initialize` is servable, // and the transport's answer for everything else is a connection-killing HTTP diff --git a/packages/hosts/mcp/src/in-memory-session-store.test.ts b/packages/hosts/mcp/src/in-memory-session-store.test.ts index bdc51db85d..2e8fc47141 100644 --- a/packages/hosts/mcp/src/in-memory-session-store.test.ts +++ b/packages/hosts/mcp/src/in-memory-session-store.test.ts @@ -777,3 +777,117 @@ describe("pre-initialize dispatch through the in-memory session store", () => { expect(sessions.sessionCount()).toBe(0); }); }); + +describe("in-memory store, quiet GET stream", () => { + it("emits a keepalive comment on a live-session GET and reconnects after cancel", async () => { + const { sessions } = makeServingStore(); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: always close the store + try { + const sessionId = await openSession(sessions); + const initialized = await Effect.runPromise( + sessions.store.dispatch({ + request: new Request("https://executor.test/mcp", { + method: "POST", + headers: { ...MCP_POST_HEADERS, "mcp-session-id": sessionId }, + body: JSON.stringify({ jsonrpc: "2.0", method: "notifications/initialized" }), + }), + principal: TEST_PRINCIPAL, + resource: defaultMcpResource, + sessionId, + method: "POST", + }), + ); + expect(initialized).toBeInstanceOf(Response); + expect((initialized as Response).status).toBe(202); + + const get = (await Effect.runPromise( + sessions.store.dispatch({ + request: new Request("https://executor.test/mcp", { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": sessionId }, + }), + principal: TEST_PRINCIPAL, + resource: defaultMcpResource, + sessionId, + method: "GET", + }), + )) as Response; + expect(get.status).toBe(200); + expect(get.headers.get("content-type")).toContain("text/event-stream"); + const reader = get.body!.getReader(); + expect(new TextDecoder().decode((await reader.read()).value)).toBe(": keepalive\n\n"); + await reader.cancel(); + + const reconnect = (await Effect.runPromise( + sessions.store.dispatch({ + request: new Request("https://executor.test/mcp", { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": sessionId }, + }), + principal: TEST_PRINCIPAL, + resource: defaultMcpResource, + sessionId, + method: "GET", + }), + )) as Response; + expect(reconnect.status).toBe(200); + expect(reconnect.headers.get("content-type")).toContain("text/event-stream"); + await reconnect.body?.cancel(); + } finally { + await sessions.close(); + } + }); + + it("keeps the GET keepalive on a toolkit resource session", async () => { + const { sessions } = makeServingStore(); + const toolkit = { kind: "toolkit" as const, slug: "deploy" }; + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: always close the store + try { + const init = (await Effect.runPromise( + sessions.store.dispatch({ + request: new Request("https://executor.test/mcp/toolkits/deploy", { + method: "POST", + headers: MCP_POST_HEADERS, + body: JSON.stringify({ + jsonrpc: "2.0", + id: 1, + method: "initialize", + params: { + protocolVersion: "2025-06-18", + capabilities: {}, + clientInfo: { name: "idle-test", version: "1.0.0" }, + }, + }), + }), + principal: TEST_PRINCIPAL, + resource: toolkit, + sessionId: null, + method: "POST", + }), + )) as Response; + expect(init.status).toBe(200); + const sessionId = init.headers.get("mcp-session-id") ?? ""; + expect(sessionId).not.toBe(""); + + const get = (await Effect.runPromise( + sessions.store.dispatch({ + request: new Request("https://executor.test/mcp/toolkits/deploy", { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": sessionId }, + }), + principal: TEST_PRINCIPAL, + resource: toolkit, + sessionId, + method: "GET", + }), + )) as Response; + expect(get.status).toBe(200); + expect(get.headers.get("content-type")).toContain("text/event-stream"); + const reader = get.body!.getReader(); + expect(new TextDecoder().decode((await reader.read()).value)).toBe(": keepalive\n\n"); + await reader.cancel(); + } finally { + await sessions.close(); + } + }); +}); diff --git a/packages/hosts/mcp/src/in-memory-session-store.ts b/packages/hosts/mcp/src/in-memory-session-store.ts index 952f9bd83f..777d3bdc23 100644 --- a/packages/hosts/mcp/src/in-memory-session-store.ts +++ b/packages/hosts/mcp/src/in-memory-session-store.ts @@ -20,6 +20,7 @@ import { type InProcessBrowserApprovalStore, } from "./browser-approval-store"; import { jsonRpcErrorBody, preInitializeMethodNotFound } from "./envelope"; +import { withMcpSseHeartbeat } from "./sse-heartbeat"; import { McpSessionStore, MCP_ORG_WRITE_ACCESS_HEADER, @@ -372,6 +373,7 @@ export const makeInMemoryMcpSessionStore = ( transport.handleRequest(withOrgWriteAccess(request, orgWriteAccess)), ); return handle.pipe( + Effect.map((response) => withMcpSseHeartbeat(request, response)), Effect.tap(() => Effect.sync(finish)), Effect.catchCause((cause) => Effect.sync(() => { diff --git a/packages/hosts/mcp/src/sse-heartbeat.test.ts b/packages/hosts/mcp/src/sse-heartbeat.test.ts new file mode 100644 index 0000000000..b8b7c4c6a7 --- /dev/null +++ b/packages/hosts/mcp/src/sse-heartbeat.test.ts @@ -0,0 +1,189 @@ +import { afterEach, describe, expect, it, vi } from "@effect/vitest"; + +import { + MCP_SSE_KEEPALIVE_FRAME, + MCP_SSE_KEEPALIVE_INTERVAL_MS, + withMcpSseHeartbeat, +} from "./sse-heartbeat"; + +const encoder = new TextEncoder(); +const decoder = new TextDecoder(); + +const sseGet = (): Request => + new Request("https://executor.test/mcp", { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": "s1" }, + }); + +const sseResponse = ( + body: ReadableStream | string, + headers: Record = { "content-type": "text/event-stream", "mcp-session-id": "s1" }, +): Response => new Response(body, { status: 200, statusText: "OK", headers }); + +const pendingStream = (): { + readonly stream: ReadableStream; + readonly push: (chunk: string) => void; + readonly close: () => void; + readonly error: (cause: unknown) => void; + readonly cancelled: () => boolean; + readonly cancelReason: () => unknown; +} => { + let cancelled = false; + let cancelReason: unknown; + let controller: ReadableStreamDefaultController; + const stream = new ReadableStream({ + start(c) { + controller = c; + }, + cancel(reason) { + cancelled = true; + cancelReason = reason; + }, + }); + return { + stream, + push: (chunk) => controller.enqueue(encoder.encode(chunk)), + close: () => controller.close(), + error: (cause) => controller.error(cause), + cancelled: () => cancelled, + cancelReason: () => cancelReason, + }; +}; + +const readChunk = async (reader: { + read: () => Promise<{ done: boolean; value?: Uint8Array }>; +}): Promise => { + const { value, done } = await reader.read(); + expect(done).toBe(false); + return decoder.decode(value); +}; + +afterEach(() => { + vi.useRealTimers(); +}); + +describe("withMcpSseHeartbeat", () => { + it("returns non-GET responses unchanged", () => { + const request = new Request("https://executor.test/mcp", { method: "POST" }); + const response = sseResponse("data: x\n\n"); + expect(withMcpSseHeartbeat(request, response)).toBe(response); + }); + + it("returns non-SSE GET responses unchanged", () => { + const response = new Response("{}", { + status: 200, + headers: { "content-type": "application/json" }, + }); + expect(withMcpSseHeartbeat(sseGet(), response)).toBe(response); + }); + + it("returns error GET SSE responses unchanged", () => { + const response = new Response("nope", { + status: 409, + headers: { "content-type": "text/event-stream" }, + }); + expect(withMcpSseHeartbeat(sseGet(), response)).toBe(response); + }); + + it("emits the keepalive comment as the first body bytes", async () => { + const upstream = pendingStream(); + const wrapped = withMcpSseHeartbeat(sseGet(), sseResponse(upstream.stream)); + const reader = wrapped.body!.getReader(); + expect(await readChunk(reader)).toBe(MCP_SSE_KEEPALIVE_FRAME); + await reader.cancel(); + }); + + it("repeats the keepalive comment on the established interval", async () => { + vi.useFakeTimers(); + const upstream = pendingStream(); + const wrapped = withMcpSseHeartbeat(sseGet(), sseResponse(upstream.stream)); + const reader = wrapped.body!.getReader(); + expect(await readChunk(reader)).toBe(MCP_SSE_KEEPALIVE_FRAME); + const next = readChunk(reader); + await vi.advanceTimersByTimeAsync(MCP_SSE_KEEPALIVE_INTERVAL_MS); + expect(await next).toBe(MCP_SSE_KEEPALIVE_FRAME); + await reader.cancel(); + }); + + it("forwards original SSE bytes unchanged and in order after the comment", async () => { + const upstream = pendingStream(); + const wrapped = withMcpSseHeartbeat(sseGet(), sseResponse(upstream.stream)); + const reader = wrapped.body!.getReader(); + expect(await readChunk(reader)).toBe(MCP_SSE_KEEPALIVE_FRAME); + upstream.push('id: 1\ndata: {"jsonrpc":"2.0"}\n\n'); + expect(await readChunk(reader)).toBe('id: 1\ndata: {"jsonrpc":"2.0"}\n\n'); + upstream.push("data: two\n\n"); + expect(await readChunk(reader)).toBe("data: two\n\n"); + await reader.cancel(); + }); + + it("preserves status, status text, and every response header", () => { + const headers = { + "content-type": "text/event-stream", + "mcp-session-id": "sess-9", + "cache-control": "no-cache, no-transform", + connection: "keep-alive", + "access-control-allow-origin": "*", + "access-control-expose-headers": "mcp-session-id", + }; + const wrapped = withMcpSseHeartbeat(sseGet(), sseResponse(pendingStream().stream, headers)); + expect(wrapped.status).toBe(200); + expect(wrapped.statusText).toBe("OK"); + for (const [key, value] of Object.entries(headers)) { + expect(wrapped.headers.get(key)).toBe(value); + } + }); + + it("does not stamp an event id on the keepalive comment", async () => { + const wrapped = withMcpSseHeartbeat(sseGet(), sseResponse(pendingStream().stream)); + const reader = wrapped.body!.getReader(); + const first = await readChunk(reader); + expect(first.startsWith(": ")).toBe(true); + expect(first).not.toContain("id:"); + expect(first).not.toContain("data:"); + await reader.cancel(); + }); + + it("cancels the upstream body when the wrapped body is cancelled", async () => { + const upstream = pendingStream(); + const wrapped = withMcpSseHeartbeat(sseGet(), sseResponse(upstream.stream)); + const reader = wrapped.body!.getReader(); + await readChunk(reader); + await reader.cancel("client-gone"); + expect(upstream.cancelled()).toBe(true); + expect(upstream.cancelReason()).toBe("client-gone"); + }); + + it("clears the timer after cancel, completion, and error", async () => { + vi.useFakeTimers(); + + const cancelled = pendingStream(); + const cancelledWrap = withMcpSseHeartbeat(sseGet(), sseResponse(cancelled.stream)); + const cancelledReader = cancelledWrap.body!.getReader(); + await readChunk(cancelledReader); + expect(vi.getTimerCount()).toBeGreaterThan(0); + await cancelledReader.cancel(); + expect(vi.getTimerCount()).toBe(0); + + const completed = pendingStream(); + const completedWrap = withMcpSseHeartbeat(sseGet(), sseResponse(completed.stream)); + const completedReader = completedWrap.body!.getReader(); + await readChunk(completedReader); + completed.close(); + await completedReader.read(); + expect(vi.getTimerCount()).toBe(0); + + const failed = pendingStream(); + const failedWrap = withMcpSseHeartbeat(sseGet(), sseResponse(failed.stream)); + const failedReader = failedWrap.body!.getReader(); + await readChunk(failedReader); + failed.error("upstream closed"); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: drain the errored reader + try { + await failedReader.read(); + } catch { + // expected: upstream error surfaces on the wrapped reader + } + expect(vi.getTimerCount()).toBe(0); + }); +}); diff --git a/packages/hosts/mcp/src/sse-heartbeat.ts b/packages/hosts/mcp/src/sse-heartbeat.ts new file mode 100644 index 0000000000..0f813a96af --- /dev/null +++ b/packages/hosts/mcp/src/sse-heartbeat.ts @@ -0,0 +1,101 @@ +// SSE comment heartbeats for quiet standalone GET `/mcp` streams. +// +// The SDK's GET handler returns `200 text/event-stream` and stores the +// controller without writing bytes until a server-initiated message exists. +// Bun fetch (and some intermediaries) wait for the first body byte and then +// time out a silent stream. `idleTimeout: 0` on Bun.serve does not emit those +// bytes. A comment frame is legal SSE, ignored by the MCP parser, and does not +// carry an event id. + +/** SSE comment frame the MCP parser drops before any event dispatch. */ +export const MCP_SSE_KEEPALIVE_FRAME = ": keepalive\n\n"; + +/** Same interval the Cloudflare agents bridge already uses. */ +export const MCP_SSE_KEEPALIVE_INTERVAL_MS = 25_000; + +const encoder = new TextEncoder(); +const keepaliveBytes = encoder.encode(MCP_SSE_KEEPALIVE_FRAME); + +const isSuccessfulSseGet = (request: Request, response: Response): boolean => + request.method === "GET" && + response.status === 200 && + (response.headers.get("content-type") ?? "").includes("text/event-stream") && + response.body !== null; + +/** + * Wrap a successful GET `text/event-stream` response with an immediate SSE + * comment and a repeating comment on {@link MCP_SSE_KEEPALIVE_INTERVAL_MS}. + * Non-GET, non-SSE, and error responses are returned unchanged. + */ +export const withMcpSseHeartbeat = (request: Request, response: Response): Response => { + if (!isSuccessfulSseGet(request, response) || response.body === null) return response; + + const upstream = response.body; + let timer: ReturnType | undefined; + let reader: + | { + read: () => Promise<{ done: boolean; value?: Uint8Array }>; + cancel: (reason?: unknown) => Promise; + } + | undefined; + let cancelled = false; + + const stopTimer = (): void => { + if (timer === undefined) return; + clearInterval(timer); + timer = undefined; + }; + + const body = new ReadableStream({ + async start(controller) { + const writeKeepalive = (): void => { + // oxlint-disable-next-line executor/no-try-catch-or-throw -- boundary: enqueue throws after cancel/close + try { + controller.enqueue(keepaliveBytes); + } catch { + stopTimer(); + } + }; + + writeKeepalive(); + timer = setInterval(writeKeepalive, MCP_SSE_KEEPALIVE_INTERVAL_MS); + (timer as { unref?: () => void }).unref?.(); + + if (cancelled) { + stopTimer(); + await upstream.cancel(); + return; + } + + reader = upstream.getReader(); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- boundary: upstream read/close must not leak the interval + try { + for (;;) { + const { done, value } = await reader.read(); + if (done || cancelled) break; + if (value !== undefined) controller.enqueue(value); + } + if (!cancelled) controller.close(); + } catch (error) { + if (!cancelled) controller.error(error); + } finally { + stopTimer(); + } + }, + async cancel(reason) { + cancelled = true; + stopTimer(); + if (reader !== undefined) { + await reader.cancel(reason); + return; + } + await upstream.cancel(reason); + }, + }); + + return new Response(body, { + status: response.status, + statusText: response.statusText, + headers: new Headers(response.headers), + }); +}; From 5f7340ae208f296944d70fdd5f0510ab0e6b1b62 Mon Sep 17 00:00:00 2001 From: utpal singh Date: Sat, 19 Sep 2026 00:16:53 +0530 Subject: [PATCH 2/2] fix: drop quiet MCP GET stream on client abort --- apps/local/src/mcp-sse-heartbeat.test.ts | 6 +++- packages/hosts/mcp/src/sse-heartbeat.test.ts | 21 ++++++++++++ packages/hosts/mcp/src/sse-heartbeat.ts | 36 ++++++++++++++------ 3 files changed, 52 insertions(+), 11 deletions(-) diff --git a/apps/local/src/mcp-sse-heartbeat.test.ts b/apps/local/src/mcp-sse-heartbeat.test.ts index 00488d0d13..ee91e385ea 100644 --- a/apps/local/src/mcp-sse-heartbeat.test.ts +++ b/apps/local/src/mcp-sse-heartbeat.test.ts @@ -81,10 +81,12 @@ describe("local MCP handler, quiet GET stream", () => { const origin = `http://127.0.0.1:${server.port}`; const { sessionId } = await openLiveSession(origin); + const getAbort = new AbortController(); + const getTimeout = setTimeout(() => getAbort.abort(), 1_000); const get = await fetch(`${origin}/mcp`, { method: "GET", headers: { accept: "text/event-stream", "mcp-session-id": sessionId }, - signal: AbortSignal.timeout(1_000), + signal: getAbort.signal, }); expect(get.status).toBe(200); expect(get.headers.get("content-type")).toContain("text/event-stream"); @@ -94,6 +96,8 @@ describe("local MCP handler, quiet GET stream", () => { const first = await reader.read(); expect(new TextDecoder().decode(first.value)).toBe(": keepalive\n\n"); await reader.cancel(); + getAbort.abort(); + clearTimeout(getTimeout); const reconnect = await fetch(`${origin}/mcp`, { method: "GET", diff --git a/packages/hosts/mcp/src/sse-heartbeat.test.ts b/packages/hosts/mcp/src/sse-heartbeat.test.ts index b8b7c4c6a7..3fcca05f63 100644 --- a/packages/hosts/mcp/src/sse-heartbeat.test.ts +++ b/packages/hosts/mcp/src/sse-heartbeat.test.ts @@ -186,4 +186,25 @@ describe("withMcpSseHeartbeat", () => { } expect(vi.getTimerCount()).toBe(0); }); + + it("cancels the upstream body when the request aborts", async () => { + const abort = new AbortController(); + const request = new Request("https://executor.test/mcp", { + method: "GET", + headers: { accept: "text/event-stream", "mcp-session-id": "s1" }, + signal: abort.signal, + }); + const upstream = pendingStream(); + const wrapped = withMcpSseHeartbeat(request, sseResponse(upstream.stream)); + const reader = wrapped.body!.getReader(); + await readChunk(reader); + abort.abort("client-abort"); + expect(upstream.cancelled()).toBe(true); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: abort may already close the reader + try { + await reader.cancel(); + } catch { + // already cancelled by abort + } + }); }); diff --git a/packages/hosts/mcp/src/sse-heartbeat.ts b/packages/hosts/mcp/src/sse-heartbeat.ts index 0f813a96af..f140202f41 100644 --- a/packages/hosts/mcp/src/sse-heartbeat.ts +++ b/packages/hosts/mcp/src/sse-heartbeat.ts @@ -39,6 +39,7 @@ export const withMcpSseHeartbeat = (request: Request, response: Response): Respo } | undefined; let cancelled = false; + let released = false; const stopTimer = (): void => { if (timer === undefined) return; @@ -46,6 +47,26 @@ export const withMcpSseHeartbeat = (request: Request, response: Response): Respo timer = undefined; }; + const releaseUpstream = async (reason?: unknown): Promise => { + if (released) return; + released = true; + cancelled = true; + stopTimer(); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- boundary: cancel after close is expected + try { + if (reader !== undefined) await reader.cancel(reason); + else await upstream.cancel(reason); + } catch { + // already closed + } + }; + + const onAbort = (): void => { + void releaseUpstream(request.signal.reason); + }; + if (request.signal.aborted) onAbort(); + else request.signal.addEventListener("abort", onAbort, { once: true }); + const body = new ReadableStream({ async start(controller) { const writeKeepalive = (): void => { @@ -62,8 +83,7 @@ export const withMcpSseHeartbeat = (request: Request, response: Response): Respo (timer as { unref?: () => void }).unref?.(); if (cancelled) { - stopTimer(); - await upstream.cancel(); + await releaseUpstream(request.signal.reason); return; } @@ -79,17 +99,13 @@ export const withMcpSseHeartbeat = (request: Request, response: Response): Respo } catch (error) { if (!cancelled) controller.error(error); } finally { - stopTimer(); + request.signal.removeEventListener("abort", onAbort); + await releaseUpstream(); } }, async cancel(reason) { - cancelled = true; - stopTimer(); - if (reader !== undefined) { - await reader.cancel(reason); - return; - } - await upstream.cancel(reason); + request.signal.removeEventListener("abort", onAbort); + await releaseUpstream(reason); }, });