diff --git a/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.test.ts b/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.test.ts index fe70e4cb2..8072df7e5 100644 --- a/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.test.ts +++ b/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.test.ts @@ -84,6 +84,7 @@ type HarnessSession = { ctx: MemoryStorage; dbHandle: { readonly end: () => void } | null; engine: ExecutionEngine | null; + getConnections?: () => Iterable; getSessionId: () => string; initialized: boolean; lastActivityMs: number; @@ -307,6 +308,56 @@ describe("McpAgentSessionDOBase apps capability persistence", () => { }); describe("McpAgentSessionDOBase transport restore", () => { + it("preserves hibernated response streams when a cold isolate starts", async () => { + const session = await makeHarnessSession(); + let closeCalls = 0; + + session.initialized = false; + session.engine = null; + session.dbHandle = null; + delete session.server; + session.getConnections = () => [ + { + close: () => { + closeCalls += 1; + }, + }, + ]; + session.runMcpAgentOnStart = async () => { + session.server = makeServer(); + session.engine = makeEngine().engine; + session.initialized = true; + }; + + await session.onStart(); + + expect(closeCalls).toBe(0); + expect(session.initialized).toBe(true); + }); + + it("closes response streams when an in-memory runtime restarts", async () => { + const session = await makeHarnessSession(); + let closeCalls = 0; + + session.getConnections = () => [ + { + close: () => { + closeCalls += 1; + }, + }, + ]; + session.runMcpAgentOnStart = async () => { + session.server = makeServer(); + session.engine = makeEngine().engine; + session.initialized = true; + }; + + await session.onStart(); + + expect(closeCalls).toBe(1); + expect(session.initialized).toBe(true); + }); + it("restores a same-session request after idle disposal leaves a stale server transport", async () => { const session = await makeHarnessSession(); diff --git a/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts b/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts index 09734bdc3..4ccc51965 100644 --- a/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts +++ b/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts @@ -619,11 +619,13 @@ export abstract class McpAgentSessionDOBase< : built; } - private closeRuntime(): Effect.Effect { + private closeRuntime(options: { readonly closeStreams?: boolean } = {}): Effect.Effect { const self = this; return Effect.gen(function* () { yield* self.releaseAllPendingApprovalLeases(); - yield* Effect.sync(() => self.closeActiveStreams()); + if (options.closeStreams ?? true) { + yield* Effect.sync(() => self.closeActiveStreams()); + } if (self.server) { const server = self.server; delete (self as { server?: McpServer }).server; @@ -658,7 +660,15 @@ export abstract class McpAgentSessionDOBase< private startRuntimeFromOnStart(props?: McpSessionProps): Effect.Effect { const self = this; return Effect.gen(function* () { - yield* self.closeRuntime(); + // PartyServer can rehydrate WebSockets before onStart runs in a + // cold-restored isolate. With no in-memory runtime to replace, those + // sockets are the live MCP response streams that triggered the restore. + const hasInMemoryRuntime = + self.initialized || + self.engine !== null || + self.dbHandle !== null || + self.server !== undefined; + yield* self.closeRuntime({ closeStreams: hasInMemoryRuntime }); const started = yield* Effect.exit(Effect.promise(() => self.runMcpAgentOnStart(props))); if (Exit.isFailure(started)) { yield* self.closeRuntime();