Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ type HarnessSession = {
ctx: MemoryStorage;
dbHandle: { readonly end: () => void } | null;
engine: ExecutionEngine<Cause.YieldableError> | null;
getConnections?: () => Iterable<unknown>;
getSessionId: () => string;
initialized: boolean;
lastActivityMs: number;
Expand Down Expand Up @@ -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();

Expand Down
16 changes: 13 additions & 3 deletions packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts
Original file line number Diff line number Diff line change
Expand Up @@ -619,11 +619,13 @@ export abstract class McpAgentSessionDOBase<
: built;
}

private closeRuntime(): Effect.Effect<void> {
private closeRuntime(options: { readonly closeStreams?: boolean } = {}): Effect.Effect<void> {
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;
Expand Down Expand Up @@ -658,7 +660,15 @@ export abstract class McpAgentSessionDOBase<
private startRuntimeFromOnStart(props?: McpSessionProps): Effect.Effect<void> {
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();
Expand Down
Loading