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-session-executor-leak.md
Original file line number Diff line number Diff line change
@@ -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.
2 changes: 1 addition & 1 deletion apps/cloud/src/mcp/session-build-semaphore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
2 changes: 1 addition & 1 deletion packages/core/api/src/server/mcp-build.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 })),
),
),
);
Expand Down
85 changes: 85 additions & 0 deletions packages/hosts/mcp/src/in-memory-session-store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});
});
37 changes: 32 additions & 5 deletions packages/hosts/mcp/src/in-memory-session-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -98,6 +98,8 @@ export class McpEngineBuildError extends Data.TaggedError("McpEngineBuildError")
export interface BuiltMcpServer {
readonly mcpServer: McpServer;
readonly engine: ExecutionEngine<Cause.YieldableError>;
readonly executor?: Executor;
readonly close?: () => Promise<void>;
}

/** The browser-mode wiring the store hands a build call when a session opts in. */
Expand Down Expand Up @@ -243,6 +245,8 @@ export const makeInMemoryMcpSessionStore = (
const servers = new Map<string, McpServer>();
const owners = new Map<string, SessionOwner>();
const engines = new Map<string, ExecutionEngine<Cause.YieldableError>>();
const executors = new Map<string, Executor>();
const closers = new Map<string, () => Promise<void>>();
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.
Expand Down Expand Up @@ -294,20 +298,33 @@ export const makeInMemoryMcpSessionStore = (
): Promise<void> =>
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<void> =>
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);
};

/**
Expand Down Expand Up @@ -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(),
Expand All @@ -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 }),
Expand All @@ -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);
},
);
}),
Expand Down Expand Up @@ -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 })));
},
};
Expand Down
Loading