From 7a6cac1407f07576d0dcf45d6cc69749a2fc1a81 Mon Sep 17 00:00:00 2001 From: Jarell Cheong Date: Wed, 5 Aug 2026 21:55:03 +0800 Subject: [PATCH] Preserve hibernated MCP streams during cold restore PartyServer can restore a Durable Object's WebSockets before onStart rebuilds Executor's in-memory MCP runtime. The existing cleanup path closed those restored response streams even though there was no prior runtime to replace, which could leave the startup block waiting on its own stream teardown until Cloudflare reset the object. Only close active streams when onStart is replacing live in-memory runtime state. Keep the existing cleanup behavior for warm restarts and failed starts, and cover both lifecycle paths with regression tests. --- .../mcp/agent-session-durable-object.test.ts | 51 +++++++++++++++++++ .../src/mcp/agent-session-durable-object.ts | 16 ++++-- 2 files changed, 64 insertions(+), 3 deletions(-) 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();