Skip to content

Commit 68da8c2

Browse files
authored
Fix MCP session transport restore after idle teardown (#1302)
* Fix MCP session transport restore after idle teardown * Serialize MCP session transport restores
1 parent 4e61a69 commit 68da8c2

4 files changed

Lines changed: 306 additions & 4 deletions

File tree

e2e/cloud/mcp-client-sessions.test.ts

Lines changed: 43 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,10 @@ import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/
1818
import { scenario } from "../src/scenario";
1919
import { Api, Mcp, Target } from "../src/services";
2020
import type { Identity } from "../src/target";
21-
import { configuredMcpPausedSessionIdleTimeoutMs } from "../setup/mcp-session-timeouts";
21+
import {
22+
configuredMcpPausedSessionIdleTimeoutMs,
23+
configuredMcpSessionTimeoutMs,
24+
} from "../setup/mcp-session-timeouts";
2225

2326
const coreApi = composePluginApi([] as const);
2427

@@ -134,6 +137,45 @@ scenario(
134137
}),
135138
);
136139

140+
const IDLE_RUNTIME_DISPOSE_BUFFER_MS = 2_000;
141+
const IDLE_RUNTIME_DISPOSE_GAP_MS =
142+
configuredMcpSessionTimeoutMs() + IDLE_RUNTIME_DISPOSE_BUFFER_MS;
143+
const IDLE_RUNTIME_DISPOSE_SCENARIO_TIMEOUT_MS = IDLE_RUNTIME_DISPOSE_GAP_MS + 120_000;
144+
145+
scenario(
146+
"MCP sessions · a resumed session after idle runtime disposal restores cleanly",
147+
{ timeout: IDLE_RUNTIME_DISPOSE_SCENARIO_TIMEOUT_MS },
148+
Effect.gen(function* () {
149+
const target = yield* Target;
150+
const mcp = yield* Mcp;
151+
const identity = yield* target.newIdentity();
152+
const bearer = yield* mcp.mintBearer(emailOf(identity));
153+
154+
const first = yield* Effect.promise(() => connectClient(target.mcpUrl, bearer));
155+
const sessionId = first.transport.sessionId;
156+
expect(sessionId, "the client got a session id").toEqual(expect.any(String));
157+
if (sessionId === undefined) return yield* Effect.die("missing session id");
158+
159+
const before = yield* Effect.promise(() =>
160+
first.client.callTool({ name: "execute", arguments: { code: 'return "before-idle";' } }),
161+
).pipe(Effect.ensuring(closeQuietly(first)));
162+
expect(before.isError, "the pre-idle call succeeds").not.toBe(true);
163+
expect(textOf(before), "the pre-idle call returns results").toContain("before-idle");
164+
165+
yield* Effect.sleep(IDLE_RUNTIME_DISPOSE_GAP_MS);
166+
167+
const second = yield* Effect.promise(() => connectClient(target.mcpUrl, bearer, sessionId));
168+
yield* Effect.gen(function* () {
169+
expect(second.transport.sessionId, "the session id is preserved").toBe(sessionId);
170+
const after = yield* Effect.promise(() =>
171+
second.client.callTool({ name: "execute", arguments: { code: 'return "after-idle";' } }),
172+
);
173+
expect(after.isError, "the post-idle call succeeds").not.toBe(true);
174+
expect(textOf(after), "the post-idle call returns results").toContain("after-idle");
175+
}).pipe(Effect.ensuring(closeQuietly(second)));
176+
}),
177+
);
178+
137179
scenario(
138180
"MCP sessions · an unknown session id fails fast with a clean error, not a hang",
139181
{},

e2e/setup/mcp-session-timeouts.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
const DEFAULT_E2E_MCP_SESSION_TIMEOUT_MS = 3_000;
22
const DEFAULT_E2E_MCP_PAUSED_SESSION_IDLE_TIMEOUT_MS = 6_000;
3+
const PRODUCTION_MCP_SESSION_TIMEOUT_MS = 5 * 60 * 1000;
34
const PRODUCTION_MCP_PAUSED_SESSION_IDLE_TIMEOUT_MS = 9 * 60 * 1000;
45

56
export const MCP_SESSION_TIMEOUT_ENV = "MCP_SESSION_TIMEOUT_MS";
@@ -32,3 +33,6 @@ export const ensureE2eMcpSessionTimeoutEnv = (): {
3233
export const configuredMcpPausedSessionIdleTimeoutMs = (): number =>
3334
positiveMilliseconds(process.env[MCP_PAUSED_SESSION_IDLE_TIMEOUT_ENV]) ??
3435
PRODUCTION_MCP_PAUSED_SESSION_IDLE_TIMEOUT_MS;
36+
37+
export const configuredMcpSessionTimeoutMs = (): number =>
38+
positiveMilliseconds(process.env[MCP_SESSION_TIMEOUT_ENV]) ?? PRODUCTION_MCP_SESSION_TIMEOUT_MS;
Lines changed: 216 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,216 @@
1+
import { describe, expect, it } from "@effect/vitest";
2+
import { Effect } from "effect";
3+
import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js";
4+
import type { Transport } from "@modelcontextprotocol/sdk/shared/transport.js";
5+
import type { JSONRPCMessage, MessageExtraInfo } from "@modelcontextprotocol/sdk/types.js";
6+
7+
import { defaultMcpResource } from "@executor-js/host-mcp";
8+
9+
import { McpAgentSessionDOBase, type SessionMeta } from "./agent-session-durable-object";
10+
11+
class MemoryStorage {
12+
private readonly data = new Map<string, unknown>();
13+
alarm: number | undefined;
14+
15+
readonly sql = {
16+
exec: () => [],
17+
};
18+
19+
async get<T>(key: string): Promise<T | undefined> {
20+
return this.data.get(key) as T | undefined;
21+
}
22+
23+
async put(key: string, value: unknown): Promise<void> {
24+
this.data.set(key, value);
25+
}
26+
27+
async setAlarm(time: number | Date): Promise<void> {
28+
this.alarm = typeof time === "number" ? time : time.getTime();
29+
}
30+
31+
async deleteAlarm(): Promise<void> {
32+
this.alarm = undefined;
33+
}
34+
35+
async delete(key: string | readonly string[]): Promise<void> {
36+
if (typeof key === "string") {
37+
this.data.delete(key);
38+
return;
39+
}
40+
for (const entry of key) {
41+
this.data.delete(entry);
42+
}
43+
}
44+
45+
async deleteAll(): Promise<void> {
46+
this.data.clear();
47+
}
48+
49+
async list<T>(
50+
options: { readonly prefix?: string; readonly limit?: number } = {},
51+
): Promise<Map<string, T>> {
52+
const rows = new Map<string, T>();
53+
for (const [key, value] of this.data) {
54+
if (options.prefix && !key.startsWith(options.prefix)) continue;
55+
rows.set(key, value as T);
56+
if (options.limit && rows.size >= options.limit) break;
57+
}
58+
return rows;
59+
}
60+
61+
async blockConcurrencyWhile<T>(callback: () => T | Promise<T>): Promise<T> {
62+
return callback();
63+
}
64+
65+
get id(): { readonly name: string } {
66+
return { name: "streamable-http:session-reconnect" };
67+
}
68+
69+
get storage(): MemoryStorage {
70+
return this;
71+
}
72+
73+
waitUntil(_promise: Promise<unknown>): void {}
74+
}
75+
76+
type HarnessSession = {
77+
alarm: () => Promise<void>;
78+
ctx: MemoryStorage;
79+
dbHandle: { readonly end: () => void } | null;
80+
engine: { readonly pausedExecutionCount: () => Effect.Effect<number> } | null;
81+
getSessionId: () => string;
82+
initialized: boolean;
83+
lastActivityMs: number;
84+
maxPausedSessionIdleMs: () => number;
85+
onStart: () => Promise<void>;
86+
pendingApprovalLeases: Map<string, never>;
87+
props: Record<string, unknown>;
88+
server?: McpServer;
89+
sessionMeta: SessionMeta;
90+
sessionTimeoutMs: () => number;
91+
validateMcpSessionOwner: (identity: {
92+
readonly accountId: string;
93+
readonly organizationId: string;
94+
}) => Promise<"ok" | "not_found" | "forbidden">;
95+
};
96+
97+
class StaleCloseTransport implements Transport {
98+
onclose?: () => void;
99+
onerror?: (error: Error) => void;
100+
onmessage?: (message: JSONRPCMessage, extra?: MessageExtraInfo) => void;
101+
102+
async start(): Promise<void> {}
103+
104+
async close(): Promise<void> {}
105+
106+
async send(_message: JSONRPCMessage): Promise<void> {}
107+
}
108+
109+
class RestoredTransport implements Transport {
110+
onclose?: () => void;
111+
onerror?: (error: Error) => void;
112+
onmessage?: (message: JSONRPCMessage, extra?: MessageExtraInfo) => void;
113+
114+
async start(): Promise<void> {}
115+
116+
async close(): Promise<void> {
117+
this.onclose?.();
118+
}
119+
120+
async send(_message: JSONRPCMessage): Promise<void> {}
121+
}
122+
123+
const makeServer = () => new McpServer({ name: "executor-test", version: "1.0.0" });
124+
125+
const makeDeferred = (): { readonly promise: Promise<void>; readonly resolve: () => void } => {
126+
let resolve: () => void = () => undefined;
127+
const promise = new Promise<void>((settle) => {
128+
resolve = settle;
129+
});
130+
return { promise, resolve };
131+
};
132+
133+
const makeHarnessSession = async (): Promise<HarnessSession> => {
134+
const sessionId = "session-reconnect";
135+
const sessionMeta: SessionMeta = {
136+
organizationId: "org-1",
137+
organizationName: "Org 1",
138+
userId: "user-1",
139+
resource: defaultMcpResource,
140+
};
141+
const storage = new MemoryStorage();
142+
const server = makeServer();
143+
await server.connect(new StaleCloseTransport());
144+
145+
const session = Object.create(McpAgentSessionDOBase.prototype) as HarnessSession;
146+
session.ctx = storage;
147+
session.dbHandle = { end: () => undefined };
148+
session.engine = { pausedExecutionCount: () => Effect.succeed(0) };
149+
session.getSessionId = () => sessionId;
150+
session.initialized = true;
151+
session.lastActivityMs = Date.now() - 10;
152+
session.maxPausedSessionIdleMs = () => 1_000;
153+
session.pendingApprovalLeases = new Map<string, never>();
154+
session.props = {};
155+
session.server = server;
156+
session.sessionMeta = sessionMeta;
157+
session.sessionTimeoutMs = () => 1;
158+
session.onStart = async () => {
159+
const restored = session.server ?? makeServer();
160+
session.server = restored;
161+
await restored.connect(new RestoredTransport());
162+
session.initialized = true;
163+
};
164+
165+
return session;
166+
};
167+
168+
describe("McpAgentSessionDOBase transport restore", () => {
169+
it("restores a same-session request after idle disposal leaves a stale server transport", async () => {
170+
const session = await makeHarnessSession();
171+
172+
await session.alarm();
173+
174+
await expect(
175+
session.validateMcpSessionOwner({ accountId: "user-1", organizationId: "org-1" }),
176+
).resolves.toBe("ok");
177+
});
178+
179+
it("single-flights concurrent same-session restore after idle disposal", async () => {
180+
const session = await makeHarnessSession();
181+
const firstRestoreEntered = makeDeferred();
182+
const finishRestore = makeDeferred();
183+
let onStartCalls = 0;
184+
let restoredServer: McpServer | undefined;
185+
186+
session.onStart = async () => {
187+
onStartCalls += 1;
188+
const restored = session.server ?? makeServer();
189+
restoredServer ??= restored;
190+
session.server = restored;
191+
firstRestoreEntered.resolve();
192+
await finishRestore.promise;
193+
await restored.connect(new RestoredTransport());
194+
session.initialized = true;
195+
};
196+
197+
await session.alarm();
198+
199+
const first = session.validateMcpSessionOwner({
200+
accountId: "user-1",
201+
organizationId: "org-1",
202+
});
203+
const second = session.validateMcpSessionOwner({
204+
accountId: "user-1",
205+
organizationId: "org-1",
206+
});
207+
208+
await firstRestoreEntered.promise;
209+
await Promise.resolve();
210+
finishRestore.resolve();
211+
212+
await expect(Promise.all([first, second])).resolves.toEqual(["ok", "ok"]);
213+
expect(onStartCalls).toBe(1);
214+
expect(session.server).toBe(restoredServer);
215+
});
216+
});

packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts

Lines changed: 43 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -200,6 +200,7 @@ export abstract class McpAgentSessionDOBase<
200200
private dbHandle: TDbHandle | null = null;
201201
private sessionMeta: SessionMeta | null = null;
202202
private initialized = false;
203+
private restoreTransportRuntimePromise: Promise<void> | null = null;
203204
private lastActivityMs = 0;
204205
private approvalResponses = new Map<string, ResumeResponse>();
205206
private approvalWaiters = new Map<string, Deferred.Deferred<ResumeResponse>>();
@@ -415,8 +416,10 @@ export abstract class McpAgentSessionDOBase<
415416
yield* self.releaseAllPendingApprovalLeases();
416417
if (self.server) {
417418
const server = self.server;
419+
delete (self as { server?: McpServer }).server;
418420
yield* Effect.promise(() => server.close()).pipe(Effect.ignore);
419421
}
422+
Reflect.set(self, "_transport", undefined);
420423
self.engine = null;
421424
if (self.dbHandle) {
422425
const dbHandle = self.dbHandle;
@@ -449,6 +452,43 @@ export abstract class McpAgentSessionDOBase<
449452
}).pipe(Effect.withSpan("McpSessionDO.ensure_runtime_for_approval"));
450453
}
451454

455+
private restoreTransportRuntimeOnce(): Effect.Effect<void> {
456+
const self = this;
457+
return Effect.gen(function* () {
458+
yield* self.closeRuntime();
459+
const restored = yield* Effect.exit(Effect.promise(() => self.onStart()));
460+
if (Exit.isFailure(restored)) {
461+
yield* self.closeRuntime();
462+
return yield* Effect.failCause(restored.cause);
463+
}
464+
});
465+
}
466+
467+
private restoreTransportRuntime(): Effect.Effect<void> {
468+
const self = this;
469+
return Effect.promise(() => {
470+
if (self.restoreTransportRuntimePromise) return self.restoreTransportRuntimePromise;
471+
472+
const restoring = Promise.resolve().then(() =>
473+
Effect.runPromise(self.restoreTransportRuntimeOnce()),
474+
);
475+
self.restoreTransportRuntimePromise = restoring;
476+
restoring.then(
477+
() => {
478+
if (self.restoreTransportRuntimePromise === restoring) {
479+
self.restoreTransportRuntimePromise = null;
480+
}
481+
},
482+
() => {
483+
if (self.restoreTransportRuntimePromise === restoring) {
484+
self.restoreTransportRuntimePromise = null;
485+
}
486+
},
487+
);
488+
return restoring;
489+
});
490+
}
491+
452492
async init(): Promise<void> {
453493
if (this.initialized) return;
454494
const props = isSessionProps(this.props) ? this.props : null;
@@ -511,9 +551,9 @@ export abstract class McpAgentSessionDOBase<
511551
Effect.withSpan("McpSessionDO.markActivity"),
512552
);
513553
} else {
514-
yield* Effect.promise(() => self.onStart()).pipe(
515-
Effect.withSpan("McpSessionDO.restore_transport_runtime"),
516-
);
554+
yield* self
555+
.restoreTransportRuntime()
556+
.pipe(Effect.withSpan("McpSessionDO.restore_transport_runtime"));
517557
}
518558
return identity.accountId === sessionMeta.userId &&
519559
identity.organizationId === sessionMeta.organizationId

0 commit comments

Comments
 (0)