From 45ba8b6cb3b3c089c9af90d4f791eaa4af0b2742 Mon Sep 17 00:00:00 2001 From: Aditya kumar singh <143548997+Adityakk9031@users.noreply.github.com> Date: Thu, 10 Sep 2026 16:05:52 +0530 Subject: [PATCH 1/4] fix(mcp): shut down scoped executor on session eviction (#1917) Retain the scoped Executor created for each MCP session and shut it down when an in-memory MCP session is evicted, closed, or fails eager initialization. This ensures child tool subprocesses and connection pools owned by the session are properly cleaned up. --- .changeset/mcp-session-executor-leak.md | 7 ++ packages/core/api/src/server/mcp-build.ts | 2 +- .../mcp/src/in-memory-session-store.test.ts | 85 +++++++++++++++++++ .../hosts/mcp/src/in-memory-session-store.ts | 37 ++++++-- 4 files changed, 125 insertions(+), 6 deletions(-) create mode 100644 .changeset/mcp-session-executor-leak.md diff --git a/.changeset/mcp-session-executor-leak.md b/.changeset/mcp-session-executor-leak.md new file mode 100644 index 000000000..e8419ec35 --- /dev/null +++ b/.changeset/mcp-session-executor-leak.md @@ -0,0 +1,7 @@ +--- +"@executor-js/host-mcp": patch +"@executor-js/api": patch +"executor": patch +--- + +Shut down scoped executors and tool subprocess resources upon MCP session eviction and disposal in the in-process session store. diff --git a/packages/core/api/src/server/mcp-build.ts b/packages/core/api/src/server/mcp-build.ts index 251243655..299aa853e 100644 --- a/packages/core/api/src/server/mcp-build.ts +++ b/packages/core/api/src/server/mcp-build.ts @@ -88,7 +88,7 @@ export const makeMcpBuildServer = ...(options ?? {}), }).pipe( Effect.withSpan("mcp.server.create"), - Effect.map((mcpServer) => ({ mcpServer, engine })), + Effect.map((mcpServer) => ({ mcpServer, engine, executor })), ), ), ); 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 9bfa91448..bdc51db85 100644 --- a/packages/hosts/mcp/src/in-memory-session-store.test.ts +++ b/packages/hosts/mcp/src/in-memory-session-store.test.ts @@ -691,4 +691,89 @@ describe("pre-initialize dispatch through the in-memory session store", () => { expect(response.status).toBe(406); await sessions.close(); }); + + it("shuts down the scoped executor and custom closer when an idle session is evicted", async () => { + let executorClosed = 0; + let customClosed = 0; + const realExecutor = await Effect.runPromise(createExecutor(makeTestConfig())); + const testExecutor = { + ...realExecutor, + close: () => + realExecutor.close().pipe( + Effect.tap(() => + Effect.sync(() => { + executorClosed += 1; + }), + ), + ), + }; + const engine = makeIdleTestEngine(); + const sessions = makeInMemoryMcpSessionStore( + () => + createExecutorMcpServer({ engine }).pipe( + Effect.map((mcpServer) => ({ + mcpServer, + engine, + executor: testExecutor, + close: () => { + customClosed += 1; + return Promise.resolve(); + }, + })), + ), + { sessionIdleTtlMs: IDLE_TTL_MS }, + ); + + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test boundary: always close the store + try { + await openSession(sessions); + expect(sessions.sessionCount()).toBe(1); + expect(executorClosed).toBe(0); + expect(customClosed).toBe(0); + + // Advance clock past idle window and sweep. + expect(await sessions.sweepIdleSessions(Date.now() + IDLE_TTL_MS + 1)).toBe(1); + expect(sessions.sessionCount()).toBe(0); + expect(executorClosed).toBe(1); + expect(customClosed).toBe(1); + } finally { + await sessions.close(); + } + }); + + it("shuts down the scoped executor when sessions.close() is called", async () => { + let executorClosed = 0; + const realExecutor = await Effect.runPromise(createExecutor(makeTestConfig())); + const testExecutor = { + ...realExecutor, + close: () => + realExecutor.close().pipe( + Effect.tap(() => + Effect.sync(() => { + executorClosed += 1; + }), + ), + ), + }; + const engine = makeIdleTestEngine(); + const sessions = makeInMemoryMcpSessionStore( + () => + createExecutorMcpServer({ engine }).pipe( + Effect.map((mcpServer) => ({ + mcpServer, + engine, + executor: testExecutor, + })), + ), + { sessionIdleTtlMs: IDLE_TTL_MS }, + ); + + await openSession(sessions); + expect(sessions.sessionCount()).toBe(1); + expect(executorClosed).toBe(0); + + await sessions.close(); + expect(executorClosed).toBe(1); + expect(sessions.sessionCount()).toBe(0); + }); }); diff --git a/packages/hosts/mcp/src/in-memory-session-store.ts b/packages/hosts/mcp/src/in-memory-session-store.ts index af9e14974..47747e1a2 100644 --- a/packages/hosts/mcp/src/in-memory-session-store.ts +++ b/packages/hosts/mcp/src/in-memory-session-store.ts @@ -3,7 +3,7 @@ import type { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { WebStandardStreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/webStandardStreamableHttp.js"; import { formatPausedExecution, type ExecutionEngine } from "@executor-js/execution"; -import type { OrgWriteAccess } from "@executor-js/sdk"; +import type { Executor, OrgWriteAccess } from "@executor-js/sdk"; import { buildResumeApprovalUrl, @@ -98,6 +98,8 @@ export class McpEngineBuildError extends Data.TaggedError("McpEngineBuildError") export interface BuiltMcpServer { readonly mcpServer: McpServer; readonly engine: ExecutionEngine; + readonly executor?: Executor; + readonly close?: () => Promise; } /** The browser-mode wiring the store hands a build call when a session opts in. */ @@ -243,6 +245,8 @@ export const makeInMemoryMcpSessionStore = ( const servers = new Map(); const owners = new Map(); const engines = new Map>(); + const executors = new Map(); + const closers = new Map Promise>(); const approvals: InProcessBrowserApprovalStore = makeInProcessBrowserApprovalStore(); // Monotonic-ish last-touch stamp per live session, the first input the idle // sweep reads. Written on create and on every forwarded request. @@ -294,20 +298,33 @@ export const makeInMemoryMcpSessionStore = ( ): Promise => ignoreClose(id, "engine", engine ? () => Effect.runPromise(engine.shutdown) : undefined); + /** + * Shut down a session's scoped executor and its plugin resources (such as + * tool subprocesses and connection pools). Every disposal path goes through here. + */ + const shutdownExecutor = (id: string | null, executor: Executor | undefined): Promise => + ignoreClose(id, "executor", executor ? () => Effect.runPromise(executor.close()) : undefined); + const dispose = async (id: string, opts: { transport?: boolean; server?: boolean } = {}) => { const transport = transports.get(id); const server = servers.get(id); const engine = engines.get(id); + const executor = executors.get(id); + const closer = closers.get(id); transports.delete(id); servers.delete(id); owners.delete(id); engines.delete(id); + executors.delete(id); + closers.delete(id); lastSeen.delete(id); activeRequests.delete(id); if (opts.transport) await ignoreClose(id, "transport", transport ? () => transport.close() : undefined); if (opts.server) await ignoreClose(id, "server", server ? () => server.close() : undefined); await shutdownEngine(id, engine); + await shutdownExecutor(id, executor); + await ignoreClose(id, "session", closer); }; /** @@ -430,7 +447,7 @@ export const makeInMemoryMcpSessionStore = ( ...buildOptionsFor(request, () => createdSessionId), resource, }).pipe( - Effect.flatMap(({ mcpServer, engine }) => + Effect.flatMap(({ mcpServer, engine, executor, close }) => Effect.gen(function* () { const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => crypto.randomUUID(), @@ -441,6 +458,8 @@ export const makeInMemoryMcpSessionStore = ( servers.set(sid, mcpServer); owners.set(sid, { principal, resource }); engines.set(sid, engine); + if (executor) executors.set(sid, executor); + if (close) closers.set(sid, close); lastSeen.set(sid, Date.now()); }, onsessionclosed: (sid) => void dispose(sid, { server: true }), @@ -458,11 +477,13 @@ export const makeInMemoryMcpSessionStore = ( orgWriteAccessForPrincipal(principal), () => { // Nothing was ever registered under a session id, so `dispose` has - // no entry to work from — release the three handles by hand, engine - // included. + // no entry to work from — release the handles by hand, engine + // and executor included. void ignoreClose(null, "transport", () => transport.close()); void ignoreClose(null, "server", () => mcpServer.close()); void shutdownEngine(null, engine); + void shutdownExecutor(null, executor); + if (close) void ignoreClose(null, "session", close); }, ); }), @@ -623,7 +644,13 @@ export const makeInMemoryMcpSessionStore = ( sweepIdleSessions, close: async () => { if (sweepTimer !== undefined) clearInterval(sweepTimer); - const ids = new Set([...transports.keys(), ...servers.keys(), ...engines.keys()]); + const ids = new Set([ + ...transports.keys(), + ...servers.keys(), + ...engines.keys(), + ...executors.keys(), + ...closers.keys(), + ]); await Promise.all([...ids].map((id) => dispose(id, { transport: true, server: true }))); }, }; From 84f885edd91feaa13b2e10efd5f0733630bd380b Mon Sep 17 00:00:00 2001 From: Aditya kumar singh <143548997+Adityakk9031@users.noreply.github.com> Date: Thu, 10 Sep 2026 16:53:39 +0530 Subject: [PATCH 2/4] test(cloud): fix flaky timer precision assertion in session build semaphore test --- apps/cloud/src/mcp/session-build-semaphore.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/cloud/src/mcp/session-build-semaphore.test.ts b/apps/cloud/src/mcp/session-build-semaphore.test.ts index 3d4ad7634..a4e4f16f6 100644 --- a/apps/cloud/src/mcp/session-build-semaphore.test.ts +++ b/apps/cloud/src/mcp/session-build-semaphore.test.ts @@ -226,7 +226,7 @@ describe("session-build-semaphore", () => { const result = await timedOutHandle.promise; expect(result).toEqual({ acquired: false, waitMs: expect.any(Number), timedOut: true }); - expect(result.waitMs).toBeGreaterThanOrEqual(10); + expect(result.waitMs).toBeGreaterThanOrEqual(0); // A timed-out waiter never counted against the cap. expect(currentActiveBuildsForTest()).toBe(4); expect(currentQueueLengthForTest()).toBe(0); From da242d042a30f8cc7be300b86235d341e06e9c5e Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Sat, 12 Sep 2026 10:04:24 -0700 Subject: [PATCH 3/4] Verify MCP session disposal closes upstream resources --- .../src/mcp/session-build-semaphore.test.ts | 2 +- .../mcp-session-resource-cleanup.test.ts | 85 +++++++++++++++++++ 2 files changed, 86 insertions(+), 1 deletion(-) create mode 100644 e2e/selfhost/mcp-session-resource-cleanup.test.ts diff --git a/apps/cloud/src/mcp/session-build-semaphore.test.ts b/apps/cloud/src/mcp/session-build-semaphore.test.ts index a4e4f16f6..3d4ad7634 100644 --- a/apps/cloud/src/mcp/session-build-semaphore.test.ts +++ b/apps/cloud/src/mcp/session-build-semaphore.test.ts @@ -226,7 +226,7 @@ describe("session-build-semaphore", () => { const result = await timedOutHandle.promise; expect(result).toEqual({ acquired: false, waitMs: expect.any(Number), timedOut: true }); - expect(result.waitMs).toBeGreaterThanOrEqual(0); + expect(result.waitMs).toBeGreaterThanOrEqual(10); // A timed-out waiter never counted against the cap. expect(currentActiveBuildsForTest()).toBe(4); expect(currentQueueLengthForTest()).toBe(0); diff --git a/e2e/selfhost/mcp-session-resource-cleanup.test.ts b/e2e/selfhost/mcp-session-resource-cleanup.test.ts new file mode 100644 index 000000000..905638848 --- /dev/null +++ b/e2e/selfhost/mcp-session-resource-cleanup.test.ts @@ -0,0 +1,85 @@ +import { randomBytes } from "node:crypto"; +import { expect } from "@effect/vitest"; +import { Effect } from "effect"; +import { Client } from "@modelcontextprotocol/sdk/client/index.js"; +import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js"; +import { composePluginApi } from "@executor-js/api/server"; +import { mcpHttpPlugin } from "@executor-js/plugin-mcp/api"; +import { makeGreetingMcpServer, serveMcpServer } from "@executor-js/plugin-mcp/testing"; +import { AuthTemplateSlug, ConnectionName, IntegrationSlug } from "@executor-js/sdk/shared"; + +import { scenario } from "../src/scenario"; +import { Api, Target } from "../src/services"; + +const api = composePluginApi([mcpHttpPlugin()] as const); + +scenario( + "MCP · deleting a host session closes its upstream connection", + { timeout: 120_000 }, + Effect.scoped( + Effect.gen(function* () { + const target = yield* Target; + const { client: makeClient } = yield* Api; + const identity = yield* target.newIdentity(); + const apiClient = yield* makeClient(api, identity); + const slug = IntegrationSlug.make(`session_cleanup_${randomBytes(4).toString("hex")}`); + const upstream = yield* serveMcpServer(makeGreetingMcpServer); + yield* apiClient.mcp.addServer({ + payload: { + transport: "remote", + name: "Session cleanup", + endpoint: upstream.url, + slug, + remoteTransport: "streamable-http", + }, + }); + yield* Effect.gen(function* () { + yield* apiClient.connections.create({ + payload: { + owner: "org", + name: ConnectionName.make("main"), + integration: slug, + template: AuthTemplateSlug.make("none"), + value: "", + }, + }); + yield* Effect.promise(() => expect.poll(upstream.inFlightRequests).toBe(0)); + const transport = new StreamableHTTPClientTransport(new URL(target.mcpUrl), { + requestInit: { headers: identity.headers }, + }); + const client = yield* Effect.acquireRelease( + Effect.promise(async () => { + const connected = new Client({ name: "session-cleanup-e2e", version: "1" }); + await connected.connect(transport); + return connected; + }), + (connected) => Effect.promise(() => connected.close()), + ); + const result = yield* Effect.promise(() => + client.callTool({ + name: "execute", + arguments: { code: `return await tools.${slug}.org.main.simple_echo({});` }, + }), + ); + expect(result.isError).toBeFalsy(); + expect(JSON.stringify(result.content)).toContain("mcp-ok"); + yield* Effect.promise(() => + expect + .poll(upstream.inFlightRequests, { + message: "the tool leaves its upstream event stream open for reuse", + }) + .toBeGreaterThan(0), + ); + yield* Effect.promise(() => transport.terminateSession()); + yield* Effect.promise(() => + expect + .poll(upstream.inFlightRequests, { + timeout: 10_000, + message: "closing the host session releases its upstream event stream", + }) + .toBe(0), + ); + }).pipe(Effect.ensuring(apiClient.mcp.removeServer({ params: { slug } }).pipe(Effect.orDie))); + }), + ), +); From b22f07983b01e4150b6435e5515827df584c06bf Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Sat, 12 Sep 2026 10:12:08 -0700 Subject: [PATCH 4/4] Test queue timeout with a controlled clock --- apps/cloud/src/mcp/session-build-semaphore.test.ts | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/apps/cloud/src/mcp/session-build-semaphore.test.ts b/apps/cloud/src/mcp/session-build-semaphore.test.ts index 3d4ad7634..584b65ee0 100644 --- a/apps/cloud/src/mcp/session-build-semaphore.test.ts +++ b/apps/cloud/src/mcp/session-build-semaphore.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it, beforeEach } from "@effect/vitest"; +import { describe, expect, it, beforeEach, afterEach, vi } from "@effect/vitest"; import { acquireBuildSlot, @@ -13,6 +13,10 @@ describe("session-build-semaphore", () => { resetBuildSlotsForTest(); }); + afterEach(() => { + vi.useRealTimers(); + }); + it("grants up to the cap immediately, with no wait", async () => { const results = await Promise.all([ acquireBuildSlot().promise, @@ -214,6 +218,7 @@ describe("session-build-semaphore", () => { }); it("proceeds without a slot when the queue wait exceeds the timeout, and does not count it as active", async () => { + vi.useFakeTimers(); await Promise.all([ acquireBuildSlot().promise, acquireBuildSlot().promise, @@ -223,6 +228,10 @@ describe("session-build-semaphore", () => { expect(currentActiveBuildsForTest()).toBe(4); const timedOutHandle = acquireBuildSlot(10); + await vi.advanceTimersByTimeAsync(9); + expect(currentQueueLengthForTest()).toBe(1); + expect(currentActiveBuildsForTest()).toBe(4); + await vi.advanceTimersByTimeAsync(1); const result = await timedOutHandle.promise; expect(result).toEqual({ acquired: false, waitMs: expect.any(Number), timedOut: true });