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/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); 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 }))); }, };