From 1b174b1dce2e6a58d3b0adb9825fc165d112cc8c Mon Sep 17 00:00:00 2001 From: Schreezer <96696271+Schreezer@users.noreply.github.com> Date: Mon, 10 Aug 2026 23:57:23 +0530 Subject: [PATCH 1/2] feat(mcp): add 2026-07-28 dual-era support --- .changeset/modern-mcp-transport.md | 5 + apps/cloud/src/env-augment.d.ts | 2 + apps/cloud/src/mcp/agent-handler.ts | 43 +- apps/host-cloudflare/package.json | 1 + apps/host-cloudflare/src/config.ts | 2 + apps/host-cloudflare/src/mcp/agent-handler.ts | 42 +- .../src/worker.e2e.node.test.ts | 56 ++ apps/host-selfhost/package.json | 1 + apps/host-selfhost/src/mcp/mcp.test.ts | 35 ++ apps/host-selfhost/src/mcp/session-store.ts | 2 +- apps/local/src/main.ts | 1 + apps/local/src/mcp.ts | 44 +- bun.lock | 20 + .../src/mcp/agent-session-durable-object.ts | 133 ++++- .../hosts/cloudflare/src/mcp/session-stub.ts | 8 + packages/hosts/mcp/package.json | 6 + packages/hosts/mcp/src/envelope.test.ts | 67 +++ packages/hosts/mcp/src/envelope.ts | 66 ++- .../hosts/mcp/src/in-memory-session-store.ts | 24 +- packages/hosts/mcp/src/index.ts | 2 + .../hosts/mcp/src/modern-tool-server.test.ts | 198 +++++++ packages/hosts/mcp/src/modern-tool-server.ts | 487 ++++++++++++++++++ packages/hosts/mcp/src/seams.ts | 8 + packages/hosts/mcp/src/tool-server.ts | 6 +- packages/plugins/mcp/package.json | 2 + .../plugins/mcp/src/sdk/connection-pool.ts | 18 +- packages/plugins/mcp/src/sdk/connection.ts | 50 +- packages/plugins/mcp/src/sdk/http-status.ts | 20 +- packages/plugins/mcp/src/sdk/invoke.test.ts | 18 +- packages/plugins/mcp/src/sdk/invoke.ts | 54 +- packages/plugins/mcp/src/sdk/plugin.ts | 14 +- .../plugins/mcp/src/sdk/probe-shape.test.ts | 84 ++- packages/plugins/mcp/src/sdk/probe-shape.ts | 63 ++- .../plugins/mcp/src/sdk/stdio-connector.ts | 5 +- packages/plugins/mcp/src/testing/server.ts | 2 +- 35 files changed, 1507 insertions(+), 82 deletions(-) create mode 100644 .changeset/modern-mcp-transport.md create mode 100644 packages/hosts/mcp/src/modern-tool-server.test.ts create mode 100644 packages/hosts/mcp/src/modern-tool-server.ts diff --git a/.changeset/modern-mcp-transport.md b/.changeset/modern-mcp-transport.md new file mode 100644 index 0000000000..b7d2a48df2 --- /dev/null +++ b/.changeset/modern-mcp-transport.md @@ -0,0 +1,5 @@ +--- +"@executor-js/plugin-mcp": minor +--- + +Add automatic MCP 2026-07-28 negotiation for HTTP, SSE, and stdio connections while preserving legacy server compatibility. diff --git a/apps/cloud/src/env-augment.d.ts b/apps/cloud/src/env-augment.d.ts index 4ef130bfa6..1f50e62337 100644 --- a/apps/cloud/src/env-augment.d.ts +++ b/apps/cloud/src/env-augment.d.ts @@ -67,6 +67,8 @@ declare global { MCP_RESOURCE_ORIGIN?: string; MCP_SESSION_TIMEOUT_MS?: string; MCP_PAUSED_SESSION_IDLE_TIMEOUT_MS?: string; + /** Emergency rollback for inbound MCP 2026-07-28 traffic only. */ + MCP_2026_07_28_ENABLED?: string; NODE_ENV?: string; // Shared with frontend diff --git a/apps/cloud/src/mcp/agent-handler.ts b/apps/cloud/src/mcp/agent-handler.ts index 7738e79c6a..f7c5c4cffd 100644 --- a/apps/cloud/src/mcp/agent-handler.ts +++ b/apps/cloud/src/mcp/agent-handler.ts @@ -5,6 +5,9 @@ import { McpAuthProvider, jsonRpcErrorBody, defaultMcpResource, + isLegacyMcpRequest, + mcpResourceKey, + validateMcpRequestAuthority, UNAVAILABLE_RETRY_AFTER_SECONDS, type AuthOutcome, type McpResource, @@ -31,7 +34,7 @@ const corsPreflightResponse = (): Response => "access-control-allow-origin": "*", "access-control-allow-methods": "GET, POST, DELETE, OPTIONS", "access-control-allow-headers": - "content-type, authorization, mcp-session-id, accept, mcp-protocol-version", + "content-type, authorization, mcp-session-id, accept, mcp-protocol-version, mcp-method, mcp-name", "access-control-expose-headers": "mcp-session-id, WWW-Authenticate", }, }); @@ -46,6 +49,17 @@ const jsonRpcResponse = ( ? jsonRpcErrorBody(status, code, message) : jsonRpcErrorBody(status, code, message, { challenge }); +const withCors = (response: Response): Response => { + const headers = new Headers(response.headers); + headers.set("access-control-allow-origin", "*"); + headers.set("access-control-expose-headers", "mcp-session-id, WWW-Authenticate"); + return new Response(response.body, { + status: response.status, + statusText: response.statusText, + headers, + }); +}; + const renderAuthError = ( auth: McpAuthProvider["Service"], request: Request, @@ -167,6 +181,8 @@ export const makeCloudMcpAgentHandler = () => { if (!ALLOWED_METHODS.has(request.method)) { return jsonRpcResponse(405, -32001, "Method not allowed"); } + const authorityRejection = validateMcpRequestAuthority(request); + if (authorityRejection) return authorityRejection; const sessionId = request.headers.get("mcp-session-id"); const { auth, outcome } = await runTraced(request, authenticate(request)); @@ -188,6 +204,30 @@ export const makeCloudMcpAgentHandler = () => { return renderAuthError(auth, request, outcome); } + const resource = resourceFromPath(request); + if (!(await isLegacyMcpRequest(request))) { + if (env.MCP_2026_07_28_ENABLED === "false") { + return jsonRpcResponse(503, -32022, "MCP 2026-07-28 support is disabled"); + } + const props = await runTraced( + request, + propsForPrincipal(request, outcome.principal, resource), + ); + const flowId = JSON.stringify([ + "modern", + outcome.principal.accountId, + outcome.principal.organizationId, + mcpResourceKey(resource), + ]); + const response = await mcpSessionStub(env.MCP_SESSION, flowId).handleModernRequest( + request, + outcome.principal, + props.session, + props.propagation, + ); + return wrapMcpSseResponse(request, env, withCors(response)); + } + if (!sessionId && request.method === "DELETE") { // Matches the old envelope's contract (@modelcontextprotocol/sdk's // `WebStandardStreamableHTTPServerTransport.handleDeleteRequest`): 200, @@ -217,7 +257,6 @@ export const makeCloudMcpAgentHandler = () => { } } - const resource = resourceFromPath(request); const props = await runTraced(request, propsForPrincipal(request, outcome.principal, resource)); (ctx as ExecutionContext & { props?: McpSessionProps }).props = props; const forwarded = withVerifiedIdentityHeaders( diff --git a/apps/host-cloudflare/package.json b/apps/host-cloudflare/package.json index b693834278..cb26403b7e 100644 --- a/apps/host-cloudflare/package.json +++ b/apps/host-cloudflare/package.json @@ -48,6 +48,7 @@ "@cloudflare/workers-types": "^4.20250410.0", "@effect/vitest": "catalog:", "@executor-js/vite-plugin": "workspace:*", + "@modelcontextprotocol/client": "2.0.0", "@tailwindcss/vite": "catalog:", "@tanstack/router-plugin": "^1.167.12", "@tanstack/virtual-file-routes": "^1.162.0", diff --git a/apps/host-cloudflare/src/config.ts b/apps/host-cloudflare/src/config.ts index c397c4ef87..55ee21fb09 100644 --- a/apps/host-cloudflare/src/config.ts +++ b/apps/host-cloudflare/src/config.ts @@ -56,6 +56,8 @@ export interface CloudflareEnv { * behind Access, or the instance is wide open. */ readonly ENABLE_DEV_AUTH?: string; + /** Emergency rollback for inbound MCP 2026-07-28 traffic only. */ + readonly MCP_2026_07_28_ENABLED?: string; } export interface CloudflareConfig { diff --git a/apps/host-cloudflare/src/mcp/agent-handler.ts b/apps/host-cloudflare/src/mcp/agent-handler.ts index 5ec09fc1b3..a545ab3f87 100644 --- a/apps/host-cloudflare/src/mcp/agent-handler.ts +++ b/apps/host-cloudflare/src/mcp/agent-handler.ts @@ -4,6 +4,9 @@ import { McpAuthProvider, jsonRpcErrorBody, defaultMcpResource, + isLegacyMcpRequest, + mcpResourceKey, + validateMcpRequestAuthority, type AuthOutcome, type Principal, } from "@executor-js/host-mcp"; @@ -27,7 +30,7 @@ const corsPreflightResponse = (): Response => "access-control-allow-origin": "*", "access-control-allow-methods": "GET, POST, DELETE, OPTIONS", "access-control-allow-headers": - "content-type, authorization, mcp-session-id, accept, mcp-protocol-version", + "content-type, authorization, mcp-session-id, accept, mcp-protocol-version, mcp-method, mcp-name", "access-control-expose-headers": "mcp-session-id, WWW-Authenticate", }, }); @@ -42,6 +45,17 @@ const jsonRpcResponse = ( ? jsonRpcErrorBody(status, code, message) : jsonRpcErrorBody(status, code, message, { challenge }); +const withCors = (response: Response): Response => { + const headers = new Headers(response.headers); + headers.set("access-control-allow-origin", "*"); + headers.set("access-control-expose-headers", "mcp-session-id, WWW-Authenticate"); + return new Response(response.body, { + status: response.status, + statusText: response.statusText, + headers, + }); +}; + const renderAuthError = ( auth: McpAuthProvider["Service"], request: Request, @@ -95,9 +109,15 @@ export const makeCloudflareMcpAgentHandler = (config: CloudflareConfig) => { binding: "MCP_SESSION", transport: "streamable-http", }); + const ALLOWED_METHODS = new Set(["GET", "POST", "DELETE", "OPTIONS"]); return async (request: Request, env: CloudflareEnv, ctx: ExecutionContext): Promise => { if (request.method === "OPTIONS") return corsPreflightResponse(); + if (!ALLOWED_METHODS.has(request.method)) { + return jsonRpcResponse(405, -32001, "Method not allowed"); + } + const authorityRejection = validateMcpRequestAuthority(request); + if (authorityRejection) return authorityRejection; const sessionId = request.headers.get("mcp-session-id"); const { auth, outcome } = await Effect.runPromise(authenticate(request, config)); @@ -114,6 +134,26 @@ export const makeCloudflareMcpAgentHandler = (config: CloudflareConfig) => { return renderAuthError(auth, request, outcome); } + if (!(await isLegacyMcpRequest(request))) { + if (env.MCP_2026_07_28_ENABLED === "false") { + return jsonRpcResponse(503, -32022, "MCP 2026-07-28 support is disabled"); + } + const props = await Effect.runPromise(propsForPrincipal(request, outcome.principal)); + const flowId = JSON.stringify([ + "modern", + outcome.principal.accountId, + outcome.principal.organizationId, + mcpResourceKey(defaultMcpResource), + ]); + const response = await mcpSessionStub(env.MCP_SESSION, flowId).handleModernRequest( + request, + outcome.principal, + props.session, + props.propagation, + ); + return withCors(response); + } + if (!sessionId && request.method === "DELETE") { return new Response(null, { status: 204, headers: { "access-control-allow-origin": "*" } }); } diff --git a/apps/host-cloudflare/src/worker.e2e.node.test.ts b/apps/host-cloudflare/src/worker.e2e.node.test.ts index 738a358827..0fed0b27f5 100644 --- a/apps/host-cloudflare/src/worker.e2e.node.test.ts +++ b/apps/host-cloudflare/src/worker.e2e.node.test.ts @@ -9,6 +9,10 @@ import { unstable_dev, type Unstable_DevWorker } from "wrangler"; import { Client } from "@modelcontextprotocol/sdk/client/index.js"; import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js"; import { ElicitRequestSchema } from "@modelcontextprotocol/sdk/types.js"; +import { + Client as ModernClient, + StreamableHTTPClientTransport as ModernStreamableHTTPClientTransport, +} from "@modelcontextprotocol/client"; import { microsoftCatalog } from "@executor-js/plugin-openapi/providers/microsoft"; // --------------------------------------------------------------------------- @@ -389,6 +393,58 @@ describe("cloudflare host e2e (workerd/miniflare)", () => { expect(result.result?.structuredContent?.result).toBe(42); }, 60_000); + it("discovers, lists, and executes over stateless MCP 2026-07-28", async () => { + const transport = new ModernStreamableHTTPClientTransport( + new URL("/mcp", `http://${worker.address}:${worker.port}`), + ); + const client = new ModernClient( + { name: "cloudflare-modern-test", version: "1" }, + { + capabilities: { elicitation: { form: {}, url: {} } }, + versionNegotiation: { mode: "auto" }, + }, + ); + + await client.connect(transport); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test cleanup boundary: the network client must close when an assertion fails. + try { + expect(client.getProtocolEra()).toBe("modern"); + const tools = await client.listTools(); + expect(tools.tools.map((tool) => tool.name)).toEqual(["execute", "skills"]); + expect(transport.sessionId).toBeUndefined(); + + const result = await client.callTool({ + name: "execute", + arguments: { code: "export default 6 * 7" }, + }); + expect(result.structuredContent).toMatchObject({ status: "completed", result: 42 }); + expect(transport.sessionId).toBeUndefined(); + + let receivedElicitation = false; + client.setRequestHandler("elicitation/create", () => { + receivedElicitation = true; + return { action: "accept", content: {} }; + }); + const resumed = await client.callTool({ + name: "execute", + arguments: { + code: [ + "return await tools.executor.coreTools.policies.create({", + ' owner: "org",', + ` pattern: "modern-input-required-${runId}.*",`, + ' action: "require_approval"', + "});", + ].join("\n"), + }, + }); + expect(receivedElicitation).toBe(true); + expect(resumed.isError).toBeFalsy(); + expect(transport.sessionId).toBeUndefined(); + } finally { + await client.close(); + } + }, 60_000); + it("delivers native elicitation on the approval-gated tool call stream", async () => { const client = new Client( { name: "native-elicitation-test", version: "1.0.0" }, diff --git a/apps/host-selfhost/package.json b/apps/host-selfhost/package.json index 9a6f70ce13..224b5a3303 100644 --- a/apps/host-selfhost/package.json +++ b/apps/host-selfhost/package.json @@ -53,6 +53,7 @@ "devDependencies": { "@effect/vitest": "catalog:", "@executor-js/vite-plugin": "workspace:*", + "@modelcontextprotocol/client": "2.0.0", "@tailwindcss/vite": "catalog:", "@tanstack/router-plugin": "^1.167.12", "@tanstack/virtual-file-routes": "^1.162.0", diff --git a/apps/host-selfhost/src/mcp/mcp.test.ts b/apps/host-selfhost/src/mcp/mcp.test.ts index 42a9379744..1cb13d87d9 100644 --- a/apps/host-selfhost/src/mcp/mcp.test.ts +++ b/apps/host-selfhost/src/mcp/mcp.test.ts @@ -3,6 +3,10 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { afterAll, expect, test } from "@effect/vitest"; +import { + Client as ModernClient, + StreamableHTTPClientTransport, +} from "@modelcontextprotocol/client"; import { mintInviteCode } from "../testing/mint-invite"; @@ -87,6 +91,37 @@ test("an authenticated MCP client initializes, lists tools, and executes code", expect(JSON.stringify(await call.json())).toContain("42"); }); +test("an authenticated MCP 2026-07-28 client discovers, lists, and executes statelessly", async () => { + const token = await signUp("modern@mcp.test"); + const transport = new StreamableHTTPClientTransport(new URL(`${BASE}/mcp`), { + requestInit: { headers: { authorization: `Bearer ${token}` } }, + fetch: (url, init) => + handler(url instanceof Request ? new Request(url, init) : new Request(url.toString(), init)), + }); + const client = new ModernClient( + { name: "selfhost-modern-test", version: "1" }, + { versionNegotiation: { mode: "auto" } }, + ); + + await client.connect(transport); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test cleanup boundary: the client must close when an assertion fails. + try { + expect(client.getProtocolEra()).toBe("modern"); + const tools = await client.listTools(); + expect(tools.tools.map((tool) => tool.name)).toEqual(["execute", "skills"]); + expect(transport.sessionId).toBeUndefined(); + + const result = await client.callTool({ + name: "execute", + arguments: { code: "export default 6 * 7" }, + }); + expect(result.structuredContent).toMatchObject({ status: "completed", result: 42 }); + expect(transport.sessionId).toBeUndefined(); + } finally { + await client.close(); + } +}); + test("an MCP session cannot be reused by another user, and unauth is rejected", async () => { const alice = await signUp("alice2@mcp.test"); const bob = await signUp("bob2@mcp.test"); diff --git a/apps/host-selfhost/src/mcp/session-store.ts b/apps/host-selfhost/src/mcp/session-store.ts index c17a8fa4db..2189c4d28d 100644 --- a/apps/host-selfhost/src/mcp/session-store.ts +++ b/apps/host-selfhost/src/mcp/session-store.ts @@ -48,7 +48,7 @@ export const makeSelfHostMcpSessionStore = ( selfHostAnalytics.record(`artifact_${action}`, { via: "agent" }), }, ), - { webBaseUrl }, + { webBaseUrl, modernEnabled: process.env.MCP_2026_07_28_ENABLED !== "false" }, ); /** The `McpSessionStore` envelope seam over a freshly built in-process store. */ diff --git a/apps/local/src/main.ts b/apps/local/src/main.ts index 26378f58f1..14e272f37a 100644 --- a/apps/local/src/main.ts +++ b/apps/local/src/main.ts @@ -119,6 +119,7 @@ export const createServerHandlers = async (token: string): Promise; +class ModernMcpEngineUnavailable extends Data.TaggedError("ModernMcpEngineUnavailable")<{}> {} + // --------------------------------------------------------------------------- // Streamable HTTP handler // --------------------------------------------------------------------------- @@ -51,6 +55,8 @@ export interface LocalMcpServerConfig { export interface LocalMcpRequestHandlerConfig { readonly defaultConfig: ExecutorMcpServerConfig; + /** Emergency rollback for inbound MCP 2026-07-28 traffic only. */ + readonly modernEnabled?: boolean; readonly createConfigForResource?: ( resource: McpResource, ) => Promise | LocalMcpServerConfig; @@ -150,6 +156,35 @@ export const createMcpRequestHandler = ( return handlerConfig.createConfigForResource(resource); }; + const modern = makeModernMcpDispatcher((principal, options) => + Effect.promise(async () => { + const resourceConfig = await configForResource(options?.resource ?? defaultMcpResource); + const engine = engineFromConfig(resourceConfig.config); + return { resourceConfig, engine }; + }).pipe( + Effect.flatMap(({ resourceConfig, engine }) => { + if (!engine) return Effect.fail(new ModernMcpEngineUnavailable()); + return engine.getDescription.pipe( + Effect.map((description) => ({ + engine, + description, + ...(resourceConfig.close ? { close: resourceConfig.close } : {}), + })), + ); + }), + ), + ); + + const localPrincipal = { + accountId: "local", + organizationId: "local", + organizationName: "Local", + email: "local@executor.invalid", + name: "Local", + avatarUrl: null, + roles: [] as string[], + }; + const dispose = async (id: string, opts: { transport?: boolean; server?: boolean } = {}) => { const t = transports.get(id); const s = servers.get(id); @@ -168,6 +203,12 @@ export const createMcpRequestHandler = ( handleRequest: async (request) => { const resource = resourceFromRequest(request); if (!resource) return jsonError(404, -32001, "MCP resource not found"); + if (!(await isLegacyMcpRequest(request))) { + if (handlerConfig.modernEnabled === false) { + return jsonError(503, -32022, "MCP 2026-07-28 support is disabled"); + } + return Effect.runPromise(modern.dispatch(request, localPrincipal, resource)); + } const sessionId = request.headers.get("mcp-session-id"); if (sessionId) { @@ -284,6 +325,7 @@ export const createMcpRequestHandler = ( close: async () => { const ids = new Set([...transports.keys(), ...servers.keys()]); await Promise.all([...ids].map((id) => dispose(id, { transport: true, server: true }))); + await modern.close(); }, }; }; diff --git a/bun.lock b/bun.lock index adf616729e..92d90b70d6 100644 --- a/bun.lock +++ b/bun.lock @@ -204,6 +204,7 @@ "@cloudflare/workers-types": "^4.20250410.0", "@effect/vitest": "catalog:", "@executor-js/vite-plugin": "workspace:*", + "@modelcontextprotocol/client": "2.0.0", "@tailwindcss/vite": "catalog:", "@tanstack/router-plugin": "^1.167.12", "@tanstack/virtual-file-routes": "^1.162.0", @@ -257,6 +258,7 @@ "devDependencies": { "@effect/vitest": "catalog:", "@executor-js/vite-plugin": "workspace:*", + "@modelcontextprotocol/client": "2.0.0", "@tailwindcss/vite": "catalog:", "@tanstack/router-plugin": "^1.167.12", "@tanstack/virtual-file-routes": "^1.162.0", @@ -704,11 +706,13 @@ "@executor-js/sdk": "workspace:*", "@modelcontextprotocol/ext-apps": "^1.7.4", "@modelcontextprotocol/sdk": "^1.12.1", + "@modelcontextprotocol/server": "2.0.0", "effect": "catalog:", "zod": "4.3.6", }, "devDependencies": { "@effect/vitest": "catalog:", + "@modelcontextprotocol/client": "2.0.0", "@types/node": "catalog:", "bun-types": "catalog:", "vitest": "catalog:", @@ -994,6 +998,8 @@ "@effect/platform-node": "catalog:", "@executor-js/config": "workspace:*", "@executor-js/sdk": "workspace:*", + "@modelcontextprotocol/client": "2.0.0", + "@modelcontextprotocol/core": "2.0.0", "@modelcontextprotocol/sdk": "^1.29.0", "zod": "4.3.6", }, @@ -2126,10 +2132,16 @@ "@mishieck/ink-titled-box": ["@mishieck/ink-titled-box@0.3.0", "", { "peerDependencies": { "ink": "^6.0.0", "react": "^19.1.0", "typescript": "^5" } }, "sha512-ugzVH9hixp3hwKfQ8On/qnsrdAxS3y9rTu/aGOFed4zVUvtZyGZNIR4rxAwXult8HKI4vJEh0OM8wib9NPrwUg=="], + "@modelcontextprotocol/client": ["@modelcontextprotocol/client@2.0.0", "", { "dependencies": { "@modelcontextprotocol/core": "2.0.0", "cross-spawn": "^7.0.5", "eventsource": "^3.0.2", "eventsource-parser": "^3.0.0", "jose": "^6.1.3", "pkce-challenge": "^5.0.0", "zod": "^4.2.0" } }, "sha512-8f1OghQ2rjzIOfqgUCP+8GiUWqRs89njoWLNqAe8kWmDePv3s1fZXseej+QXemssEuuOvLLmLO/kqM3IQHtISw=="], + + "@modelcontextprotocol/core": ["@modelcontextprotocol/core@2.0.0", "", { "dependencies": { "zod": "^4.2.0" } }, "sha512-pJCEwGG7Lfr/+PQp9ZTwKXNeO5wzbfKL7H3MYpCorM4oFBoQrdjnBgEoqG+RjhsvS1FKrDbKux+M1HhlnGWqcA=="], + "@modelcontextprotocol/ext-apps": ["@modelcontextprotocol/ext-apps@1.7.5", "", { "dependencies": { "@standard-schema/spec": "^1.1.0" }, "peerDependencies": { "@modelcontextprotocol/sdk": "^1.29.0", "react": "^17.0.0 || ^18.0.0 || ^19.0.0", "react-dom": "^17.0.0 || ^18.0.0 || ^19.0.0", "zod": "^3.25.0 || ^4.0.0" }, "optionalPeers": ["react", "react-dom"] }, "sha512-TjPH2S2y5UEGKhmI6+XGFuqfqOV4ppe1x6DA3txnUaEWkgtA4G5vo14jGKFZmegdkZ1H4QMLyujLvoU1BEdnAg=="], "@modelcontextprotocol/sdk": ["@modelcontextprotocol/sdk@1.29.0", "", { "dependencies": { "@hono/node-server": "^1.19.9", "ajv": "^8.17.1", "ajv-formats": "^3.0.1", "content-type": "^1.0.5", "cors": "^2.8.5", "cross-spawn": "^7.0.5", "eventsource": "^3.0.2", "eventsource-parser": "^3.0.0", "express": "^5.2.1", "express-rate-limit": "^8.2.1", "hono": "^4.11.4", "jose": "^6.1.3", "json-schema-typed": "^8.0.2", "pkce-challenge": "^5.0.0", "raw-body": "^3.0.0", "zod": "^3.25 || ^4.0", "zod-to-json-schema": "^3.25.1" }, "peerDependencies": { "@cfworker/json-schema": "^4.1.1" }, "optionalPeers": ["@cfworker/json-schema"] }, "sha512-zo37mZA9hJWpULgkRpowewez1y6ML5GsXJPY8FI0tBBCd77HEvza4jDqRKOXgHNn867PVGCyTdzqpz0izu5ZjQ=="], + "@modelcontextprotocol/server": ["@modelcontextprotocol/server@2.0.0", "", { "dependencies": { "@modelcontextprotocol/core": "2.0.0", "zod": "^4.2.0" } }, "sha512-YhHWdHfpFMQfd0prsEnxKeS3Qz3ytIGmsS0sth4KDjnacIT7hxk6hXHkJ9KysxlkvTM+WZAtQbbcUhdoP4Hvtw=="], + "@msgpackr-extract/msgpackr-extract-darwin-arm64": ["@msgpackr-extract/msgpackr-extract-darwin-arm64@3.0.3", "", { "os": "darwin", "cpu": "arm64" }, "sha512-QZHtlVgbAdy2zAqNA9Gu1UpIuI8Xvsd1v8ic6B2pZmeFnFcMWiPLfWXh7TVw4eGEZ/C9TH281KwhVoeQUKbyjw=="], "@msgpackr-extract/msgpackr-extract-darwin-x64": ["@msgpackr-extract/msgpackr-extract-darwin-x64@3.0.3", "", { "os": "darwin", "cpu": "x64" }, "sha512-mdzd3AVzYKuUmiWOQ8GNhl64/IoFGol569zNRdkLReh6LRLHOXxU4U8eq0JwaD8iFHdVGqSy4IjFL4reoWCDFw=="], @@ -6118,8 +6130,16 @@ "@manypkg/get-packages/fs-extra": ["fs-extra@8.1.0", "", { "dependencies": { "graceful-fs": "^4.2.0", "jsonfile": "^4.0.0", "universalify": "^0.1.0" } }, "sha512-yhlQgA6mnOJUKOsRUFsgJdQCvkKhcz8tlZG5HBQfReYZy46OwLcY+Zia0mtdHsOo9y/hP+CxMN0TU9QxoOtG4g=="], + "@modelcontextprotocol/client/jose": ["jose@6.2.2", "", {}, "sha512-d7kPDd34KO/YnzaDOlikGpOurfF0ByC2sEV4cANCtdqLlTfBlw2p14O/5d/zv40gJPbIQxfES3nSx1/oYNyuZQ=="], + + "@modelcontextprotocol/client/zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], + + "@modelcontextprotocol/core/zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], + "@modelcontextprotocol/sdk/jose": ["jose@6.2.2", "", {}, "sha512-d7kPDd34KO/YnzaDOlikGpOurfF0ByC2sEV4cANCtdqLlTfBlw2p14O/5d/zv40gJPbIQxfES3nSx1/oYNyuZQ=="], + "@modelcontextprotocol/server/zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], + "@octokit/request/content-type": ["content-type@2.0.0", "", {}, "sha512-j/O/d7GcZCyNl7/hwZAb606rzqkyvaDctLmckbxLzHvFBzTJHuGEdodATcP3yIRoDrLHkIATJuvzbFlp/ki2cQ=="], "@opentelemetry/exporter-logs-otlp-proto/@opentelemetry/resources": ["@opentelemetry/resources@2.6.1", "", { "dependencies": { "@opentelemetry/core": "2.6.1", "@opentelemetry/semantic-conventions": "^1.29.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.3.0 <1.10.0" } }, "sha512-lID/vxSuKWXM55XhAKNoYXu9Cutoq5hFdkbTdI/zDKQktXzcWBVhNsOkiZFTMU9UtEWuGRNe0HUgmsFldIdxVA=="], 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 4ccc51965a..f10b45bae6 100644 --- a/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts +++ b/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts @@ -1,4 +1,4 @@ -import { Cause, Deferred, Effect, Exit, Option, Schema } from "effect"; +import { Cause, Data, Deferred, Effect, Exit, Option, Schema } from "effect"; import type * as Tracer from "effect/Tracer"; import type { Connection, ConnectionContext } from "agents"; import { McpAgent } from "agents/mcp"; @@ -18,7 +18,16 @@ import { type PausedExecutionHooks, type ResumeFallbackOutcome, } from "@executor-js/host-mcp/tool-server"; -import { defaultMcpResource, type McpResource } from "@executor-js/host-mcp"; +import { + defaultMcpResource, + mcpResourceKey, + type McpResource, + type Principal, +} from "@executor-js/host-mcp"; +import { + makeModernMcpDispatcher, + type ModernMcpDispatcher, +} from "@executor-js/host-mcp/modern-tool-server"; import type { IncomingPropagationHeaders, McpElicitationMode } from "./do-headers"; import type { @@ -137,12 +146,15 @@ const LAST_ACTIVITY_KEY = "last-activity-ms"; const PARTYSERVER_NAME_KEY = "__ps_name"; /** The agents SDK's durable "condemned" marker (`_cf_scheduleDestroy`). */ const AGENTS_DESTROY_PENDING_KEY = "cf_agents_destroy_pending"; +const MODERN_REQUEST_STATE_KEY = "modern-request-state-key"; const MCP_HTTP_METHOD_HEADER = "cf-mcp-method"; const MCP_MESSAGE_HEADER = "cf-mcp-message"; const MODEL_RESUME_FORWARD_TIMEOUT_MS = 10_000; const MCP_STREAM_REQS_KEY_PREFIX = "__mcp_stream_reqs__:"; const approvalResponseKey = (executionId: string) => `approval-response:${executionId}`; +class ModernMcpRuntimeUnavailable extends Data.TaggedError("ModernMcpRuntimeUnavailable")<{}> {} + type JsonRpcRequestId = string | number; const JsonRpcRequestWithId = Schema.Struct({ id: Schema.Union([Schema.String, Schema.Number]), @@ -235,6 +247,7 @@ export abstract class McpAgentSessionDOBase< private approvalResponses = new Map(); private approvalWaiters = new Map>(); private pendingApprovalLeases = new Map(); + private modernDispatcher: ModernMcpDispatcher | null = null; protected abstract openSessionDb(): TDbHandle | Promise; @@ -623,6 +636,11 @@ export abstract class McpAgentSessionDOBase< const self = this; return Effect.gen(function* () { yield* self.releaseAllPendingApprovalLeases(); + if (self.modernDispatcher) { + const modern = self.modernDispatcher; + self.modernDispatcher = null; + yield* Effect.promise(() => modern.close()).pipe(Effect.ignore); + } if (options.closeStreams ?? true) { yield* Effect.sync(() => self.closeActiveStreams()); } @@ -657,6 +675,117 @@ export abstract class McpAgentSessionDOBase< }).pipe(Effect.withSpan("McpSessionDO.ensure_runtime_for_approval")); } + private ensureModernOwner( + principal: Principal, + token: McpSessionInit, + ): Effect.Effect { + const self = this; + return Effect.gen(function* () { + let sessionMeta = yield* self.loadSessionMeta(); + if (sessionMeta) { + if ( + sessionMeta.userId !== principal.accountId || + sessionMeta.organizationId !== principal.organizationId || + mcpResourceKey(sessionMeta.resource) !== mcpResourceKey(token.resource) + ) { + return "forbidden" as const; + } + } else { + sessionMeta = yield* self.resolveAndStoreSessionMeta(token); + } + + yield* Effect.promise(() => self.markActivity()); + return sessionMeta; + }); + } + + private ensureModernRuntime( + principal: Principal, + token: McpSessionInit, + ): Effect.Effect<"ok" | "forbidden"> { + const self = this; + return Effect.gen(function* () { + const sessionMeta = yield* self.ensureModernOwner(principal, token); + if (sessionMeta === "forbidden") return "forbidden" as const; + if (!self.initialized || !self.engine) { + const dbHandle = yield* self.openSessionDbHandle(); + const built = yield* self.buildRuntime(sessionMeta, dbHandle); + self.dbHandle = dbHandle; + self.server = built.mcpServer; + self.engine = built.engine; + self.initialized = true; + } + return "ok" as const; + }); + } + + /** Serve one modern stateless HTTP exchange inside an owner/resource-keyed + * Durable Object. The live engine stays here across MRTR rounds; signed + * requestState is only the client-carried continuation pointer. */ + async handleModernRequest( + request: Request, + principal: Principal, + token: McpSessionInit, + incoming?: IncomingTraceHeaders, + ): Promise { + const self = this; + const program = Effect.gen(function* () { + yield* self.prepareErrorCaptureScope(); + const access = yield* self.ensureModernOwner(principal, token); + if (access === "forbidden") { + return new Response( + JSON.stringify({ + jsonrpc: "2.0", + id: null, + error: { code: -32003, message: "Modern MCP flow belongs to another principal" }, + }), + { status: 403, headers: { "content-type": "application/json" } }, + ); + } + + if (!self.modernDispatcher) { + let key = yield* Effect.promise(() => + self.ctx.storage.get(MODERN_REQUEST_STATE_KEY), + ); + if (!key) { + key = crypto.getRandomValues(new Uint8Array(32)); + yield* Effect.promise(() => self.ctx.storage.put(MODERN_REQUEST_STATE_KEY, key)); + } + self.modernDispatcher = makeModernMcpDispatcher( + () => + self.ensureModernRuntime(principal, token).pipe( + Effect.flatMap((runtimeAccess) => { + if (runtimeAccess === "forbidden" || !self.engine) { + return Effect.fail(new ModernMcpRuntimeUnavailable()); + } + return self.engine.getDescription.pipe( + Effect.map((description) => ({ engine: self.engine!, description })), + ); + }), + ), + { + requestStateKey: key, + onExecutionPaused: (executionId) => + Effect.runPromise( + self.startPendingApprovalLease(executionId, { + ttlMs: PAUSED_APPROVAL_TIMEOUT_MS, + expiresAt: new Date(Date.now() + PAUSED_APPROVAL_TIMEOUT_MS).toISOString(), + }), + ), + onResumeStarted: (executionId) => + Effect.runPromise(self.beginPendingApprovalResume(executionId)), + onResumeSettled: (executionId) => + Effect.runPromise(self.finishPendingApprovalResume(executionId)), + }, + ); + } + return yield* self.modernDispatcher.dispatch(request, principal, token.resource); + }).pipe(Effect.withSpan("McpSessionDO.handleModernRequest"), (effect) => + self.withSpanFlush(effect), + ); + return Effect.runPromise(self.withTelemetry(program, incoming)); + } + private startRuntimeFromOnStart(props?: McpSessionProps): Effect.Effect { const self = this; return Effect.gen(function* () { diff --git a/packages/hosts/cloudflare/src/mcp/session-stub.ts b/packages/hosts/cloudflare/src/mcp/session-stub.ts index 3a003ff0cc..b56b4e1072 100644 --- a/packages/hosts/cloudflare/src/mcp/session-stub.ts +++ b/packages/hosts/cloudflare/src/mcp/session-stub.ts @@ -6,7 +6,9 @@ import type { McpSessionApprovalResult, McpSessionModelResumeResult, McpSessionResumeApprovalResult, + McpSessionInit, } from "./agent-session-durable-object"; +import type { Principal } from "@executor-js/host-mcp"; import { mcpSessionDurableObjectName } from "./execution-owner-directory"; export interface McpSessionNamespace { @@ -15,6 +17,12 @@ export interface McpSessionNamespace { } export interface McpSessionStub { + readonly handleModernRequest: ( + request: Request, + principal: Principal, + token: McpSessionInit, + incoming?: IncomingTraceHeaders, + ) => Promise; readonly validateMcpSessionOwner: ( identity: McpApprovalOwner, ) => Promise<"ok" | "not_found" | "forbidden" | "terminated">; diff --git a/packages/hosts/mcp/package.json b/packages/hosts/mcp/package.json index 2c0345b466..5e6b376fc9 100644 --- a/packages/hosts/mcp/package.json +++ b/packages/hosts/mcp/package.json @@ -12,6 +12,10 @@ "types": "./src/tool-server.ts", "default": "./src/tool-server.ts" }, + "./modern-tool-server": { + "types": "./src/modern-tool-server.ts", + "default": "./src/modern-tool-server.ts" + }, "./create-artifact": { "types": "./src/create-artifact.ts", "default": "./src/create-artifact.ts" @@ -49,11 +53,13 @@ "@executor-js/sdk": "workspace:*", "@modelcontextprotocol/ext-apps": "^1.7.4", "@modelcontextprotocol/sdk": "^1.12.1", + "@modelcontextprotocol/server": "2.0.0", "effect": "catalog:", "zod": "4.3.6" }, "devDependencies": { "@effect/vitest": "catalog:", + "@modelcontextprotocol/client": "2.0.0", "@types/node": "catalog:", "bun-types": "catalog:", "vitest": "catalog:" diff --git a/packages/hosts/mcp/src/envelope.test.ts b/packages/hosts/mcp/src/envelope.test.ts index 523dc60a0b..8b22d1e3d1 100644 --- a/packages/hosts/mcp/src/envelope.test.ts +++ b/packages/hosts/mcp/src/envelope.test.ts @@ -176,6 +176,73 @@ describe("McpServingRoutes envelope", () => { expect(response.status).toBe(403); expect(await Effect.runPromise(Ref.get(disposed))).toEqual([]); }); + + it("classifies modern traffic before legacy session rules and dispatches statelessly", async () => { + const calls = await Effect.runPromise(Ref.make([])); + const ModernStoreLive = Layer.succeed(McpSessionStore)({ + dispatch: (): Effect.Effect => + Effect.die("legacy dispatch should not run"), + dispatchModern: ({ request }) => + Ref.update(calls, (all) => [...all, request.headers.get("mcp-session-id") ?? "none"]).pipe( + Effect.as( + new Response( + JSON.stringify({ + jsonrpc: "2.0", + id: 1, + result: { resultType: "complete", supportedVersions: ["2026-07-28"] }, + }), + { status: 200 }, + ), + ), + ), + dispose: () => Effect.void, + }); + const handler = buildHandler(ModernStoreLive, McpErrorReporterNoop); + const response = await handler( + new Request("https://host.test/mcp", { + method: "POST", + headers: { + authorization: "Bearer x", + "content-type": "application/json", + "mcp-protocol-version": "2026-07-28", + "mcp-method": "server/discover", + "mcp-name": "server", + }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: 1, + method: "server/discover", + params: { + _meta: { + "io.modelcontextprotocol/protocolVersion": "2026-07-28", + "io.modelcontextprotocol/clientInfo": { name: "test", version: "1" }, + "io.modelcontextprotocol/clientCapabilities": {}, + }, + }, + }), + }), + ); + + expect(response.status).toBe(200); + expect(response.headers.get("access-control-allow-origin")).toBe("*"); + expect(await Effect.runPromise(Ref.get(calls))).toEqual(["none"]); + }); + + it("rejects mismatched Host and cross-origin browser requests before dispatch", async () => { + const handler = buildHandler(OkStoreLive, McpErrorReporterNoop); + const badHost = await handler( + new Request("https://host.test/mcp", { method: "GET", headers: { host: "evil.test" } }), + ); + expect(badHost.status).toBe(421); + + const badOrigin = await handler( + new Request("https://host.test/mcp", { + method: "GET", + headers: { origin: "https://evil.test" }, + }), + ); + expect(badOrigin.status).toBe(403); + }); }); it("dispatches toolkit MCP routes with the parsed toolkit resource", async () => { diff --git a/packages/hosts/mcp/src/envelope.ts b/packages/hosts/mcp/src/envelope.ts index fe5483978e..43e8e11009 100644 --- a/packages/hosts/mcp/src/envelope.ts +++ b/packages/hosts/mcp/src/envelope.ts @@ -1,5 +1,6 @@ import { Effect, Match, Predicate } from "effect"; import { HttpRouter, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"; +import { isLegacyRequest } from "@modelcontextprotocol/server"; import { defaultMcpResource, @@ -75,7 +76,7 @@ const corsPreflightResponse = (): Response => "access-control-allow-origin": "*", "access-control-allow-methods": "GET, POST, DELETE, OPTIONS", "access-control-allow-headers": - "content-type, authorization, mcp-session-id, accept, mcp-protocol-version", + "content-type, authorization, mcp-session-id, accept, mcp-protocol-version, mcp-method, mcp-name", "access-control-expose-headers": "mcp-session-id, WWW-Authenticate", }, }); @@ -146,6 +147,42 @@ const jsonRpcResponse = ( ? jsonRpcErrorBody(status, code, message) : jsonRpcErrorBody(status, code, message, { challenge }); +const withMcpCors = (response: Response): Response => { + const headers = new Headers(response.headers); + headers.set("access-control-allow-origin", "*"); + headers.set("access-control-expose-headers", "mcp-session-id, WWW-Authenticate"); + return new Response(response.body, { + status: response.status, + statusText: response.statusText, + headers, + }); +}; + +/** Reject DNS-rebinding and cross-origin browser requests before credentials + * are evaluated. Non-browser clients normally send no Origin and are accepted. */ +export const validateMcpRequestAuthority = (request: Request): Response | null => { + const url = new URL(request.url); + const host = request.headers.get("host"); + if (host && host.toLowerCase() !== url.host.toLowerCase()) { + return jsonRpcResponse(421, -32600, "Host header does not match request URL"); + } + const origin = request.headers.get("origin"); + if (!origin) return null; + if (!URL.canParse(origin) || new URL(origin).origin !== url.origin) { + return jsonRpcResponse(403, -32600, "Origin is not allowed"); + } + return null; +}; + +/** Classify conservatively: an unreadable or malformed request never enters + * the modern stateless adapter by accident. The selected transport will render + * the protocol error on its own path. */ +export const isLegacyMcpRequest = (request: Request): Promise => + isLegacyRequest(request).then( + (legacy) => legacy, + () => true, + ); + /** * Reconstruct a WHATWG `Request` from the Effect HTTP request. Prefer the * underlying source `Request` (preserves the body stream the transport reads); @@ -233,6 +270,9 @@ const mcpDispatch = (resource: McpResource) => return fromWebResponse(jsonRpcResponse(405, -32001, "Method not allowed")); } + const authorityRejection = validateMcpRequestAuthority(request); + if (authorityRejection) return fromWebResponse(authorityRejection); + const sessionId = request.headers.get("mcp-session-id"); // Authenticate (and, for session-aware providers, authorize) on EVERY @@ -245,6 +285,30 @@ const mcpDispatch = (resource: McpResource) => } const principal = outcome.principal; + // Era classification MUST happen before any legacy session-id rules. A + // modern request is a complete exchange and neither requires nor creates + // an `mcp-session-id`; the existing session store remains untouched for + // claim-less 2025-era traffic. + const legacy = yield* Effect.promise(() => isLegacyMcpRequest(request)); + if (!legacy) { + if (!store.dispatchModern) { + return fromWebResponse( + jsonRpcErrorBody( + 400, + -32022, + "Unsupported protocol version: this host has not enabled MCP 2026-07-28", + ), + ); + } + const modern = yield* store.dispatchModern({ + request, + principal, + resource, + method: request.method, + }); + return fromWebResponse(withMcpCors(modern)); + } + // No session id: per the streamable-HTTP transport contract, only POST opens // a session. A GET needs an existing id (400); a DELETE on nothing is a // no-op (204). Both short-circuit BEFORE dispatch so the store never spins up diff --git a/packages/hosts/mcp/src/in-memory-session-store.ts b/packages/hosts/mcp/src/in-memory-session-store.ts index 2cd870fc1b..745de9fe31 100644 --- a/packages/hosts/mcp/src/in-memory-session-store.ts +++ b/packages/hosts/mcp/src/in-memory-session-store.ts @@ -16,6 +16,7 @@ import { type InProcessBrowserApprovalStore, } from "./browser-approval-store"; import { jsonRpcErrorBody } from "./envelope"; +import { makeModernMcpDispatcher } from "./modern-tool-server"; import { McpSessionStore, defaultMcpResource, @@ -155,13 +156,27 @@ export const makeInMemoryMcpSessionStore = ( // proxy) it is preferred over the request URL — whose host would be the // internal bind address (127.0.0.1:PORT), unreachable for the user. Omit it on // loopback hosts (local/desktop), where the request URL is already correct. - options: { readonly webBaseUrl?: string } = {}, + options: { readonly webBaseUrl?: string; readonly modernEnabled?: boolean } = {}, ): InMemoryMcpSessionStore => { const transports = new Map(); const servers = new Map(); const owners = new Map(); const engines = new Map>(); const approvals: InProcessBrowserApprovalStore = makeInProcessBrowserApprovalStore(); + const modern = makeModernMcpDispatcher((principal, modernOptions) => + buildServer(principal, { + resource: modernOptions?.resource, + elicitationMode: { mode: "model" }, + artifactsEnabled: false, + }).pipe( + Effect.flatMap(({ mcpServer, engine }) => + engine.getDescription.pipe( + Effect.map((description) => ({ engine, description })), + Effect.ensuring(Effect.promise(() => ignoreClose(() => mcpServer.close()))), + ), + ), + ), + ); const dispose = async (id: string, opts: { transport?: boolean; server?: boolean } = {}) => { const transport = transports.get(id); @@ -293,6 +308,12 @@ export const makeInMemoryMcpSessionStore = ( sessionId ? forward(sessionId, principal, resource, request) : create(principal, resource ?? defaultMcpResource, request), + ...(options.modernEnabled === false + ? {} + : { + dispatchModern: ({ request, principal, resource }) => + modern.dispatch(request, principal, resource), + }), dispose: (sessionId) => Effect.promise(() => dispose(sessionId, { transport: true, server: true })), }; @@ -379,6 +400,7 @@ export const makeInMemoryMcpSessionStore = ( close: async () => { const ids = new Set([...transports.keys(), ...servers.keys()]); await Promise.all([...ids].map((id) => dispose(id, { transport: true, server: true }))); + await modern.close(); }, }; }; diff --git a/packages/hosts/mcp/src/index.ts b/packages/hosts/mcp/src/index.ts index 2e536296d9..9445443fb4 100644 --- a/packages/hosts/mcp/src/index.ts +++ b/packages/hosts/mcp/src/index.ts @@ -40,5 +40,7 @@ export { McpServingRoutes, McpDiscoveryRoutes, jsonRpcErrorBody, + isLegacyMcpRequest, + validateMcpRequestAuthority, UNAVAILABLE_RETRY_AFTER_SECONDS, } from "./envelope"; diff --git a/packages/hosts/mcp/src/modern-tool-server.test.ts b/packages/hosts/mcp/src/modern-tool-server.test.ts new file mode 100644 index 0000000000..cdb18049e9 --- /dev/null +++ b/packages/hosts/mcp/src/modern-tool-server.test.ts @@ -0,0 +1,198 @@ +import { describe, expect, it } from "@effect/vitest"; +import { Effect } from "effect"; +import { Client, StreamableHTTPClientTransport } from "@modelcontextprotocol/client"; + +import { FormElicitation, ToolAddress } from "@executor-js/sdk"; +import type { ExecutionEngine, ExecutionResult } from "@executor-js/execution"; + +import { defaultMcpResource, type Principal } from "./seams"; +import { makeModernMcpDispatcher } from "./modern-tool-server"; + +type DispatcherRequest = Parameters["dispatch"]>[0]; + +const cloneDispatcherRequest = (request: DispatcherRequest): DispatcherRequest => + // oxlint-disable-next-line executor/no-double-cast -- test adapter boundary: Bun and the v2 client's undici types disagree about Headers extensions on the same runtime Request. + request.clone() as unknown as DispatcherRequest; + +const principal: Principal = { + accountId: "account-1", + organizationId: "org-1", + organizationName: "Test Org", + email: "test@example.com", + name: "Test", + avatarUrl: null, + roles: [], +}; + +const completed = (value: unknown): ExecutionResult => ({ + status: "completed", + result: { result: value }, +}); + +const stubEngine = ( + overrides: { + readonly executeWithPause?: ExecutionEngine["executeWithPause"]; + readonly resume?: ExecutionEngine["resume"]; + readonly getPausedExecution?: ExecutionEngine["getPausedExecution"]; + } = {}, +): ExecutionEngine => ({ + execute: () => Effect.succeed({ result: "unused" }), + executeWithPause: overrides.executeWithPause ?? (() => Effect.succeed(completed("ok"))), + resume: overrides.resume ?? (() => Effect.succeed(null)), + getPausedExecution: overrides.getPausedExecution ?? (() => Effect.succeed(null)), + pausedExecutionCount: () => Effect.succeed(0), + hasPausedExecutions: () => Effect.succeed(false), + getDescription: Effect.succeed("Live integrations:\n- test"), +}); + +const withClient = async ( + dispatcher: ReturnType, + use: (client: Client) => Promise, + onRequest?: (request: DispatcherRequest) => void, +): Promise => { + const transport = new StreamableHTTPClientTransport(new URL("https://executor.test/mcp"), { + fetch: (url, init) => { + // oxlint-disable-next-line executor/no-double-cast -- test adapter boundary: v2 client bundles undici fetch types, while the runtime objects implement the same web Request contract. + const request = (url instanceof Request + ? new Request(url, init) + : new Request(url.toString(), init)) as unknown as DispatcherRequest; + onRequest?.(cloneDispatcherRequest(request)); + return Effect.runPromise(dispatcher.dispatch(request, principal, defaultMcpResource)); + }, + }); + const client = new Client( + { name: "modern-test", version: "1" }, + { + capabilities: { elicitation: { form: {}, url: {} } }, + versionNegotiation: { mode: "auto" }, + }, + ); + await client.connect(transport); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test cleanup boundary: both client and dispatcher must close when an assertion fails. + try { + await use(client); + } finally { + await client.close(); + await dispatcher.close(); + } +}; + +describe("modern MCP dispatcher", () => { + it("answers discover and tools/list without building an execution engine", async () => { + let builds = 0; + const dispatcher = makeModernMcpDispatcher(() => { + builds += 1; + return Effect.succeed({ engine: stubEngine(), description: "test" }); + }); + + await withClient(dispatcher, async (client) => { + expect(client.getProtocolEra()).toBe("modern"); + const tools = await client.listTools(); + expect(tools.tools.map((tool) => tool.name)).toEqual(["execute", "skills"]); + expect(builds).toBe(0); + }); + }); + + it("executes a complete call and reuses the principal-scoped workspace", async () => { + let builds = 0; + const engine = stubEngine({ + executeWithPause: (code) => Effect.succeed(completed(`ran:${code}`)), + }); + const dispatcher = makeModernMcpDispatcher(() => { + builds += 1; + return Effect.succeed({ engine, description: "test" }); + }); + + await withClient(dispatcher, async (client) => { + const first = await client.callTool({ name: "execute", arguments: { code: "return 1" } }); + const second = await client.callTool({ name: "execute", arguments: { code: "return 2" } }); + expect(first.structuredContent).toMatchObject({ result: "ran:return 1" }); + expect(second.structuredContent).toMatchObject({ result: "ran:return 2" }); + expect(builds).toBe(1); + }); + }); + + it("round-trips elicitation through signed requestState and resumes once", async () => { + const request = FormElicitation.make({ + message: "Approve the action?", + requestedSchema: { + type: "object", + properties: { approved: { type: "boolean" } }, + required: ["approved"], + }, + }); + const paused: Extract = { + status: "paused", + execution: { + id: "exec-modern-1", + elicitationContext: { + address: ToolAddress.make("tools.test.org.main.action"), + args: {}, + request, + }, + }, + }; + let resumes = 0; + const engine = stubEngine({ + executeWithPause: () => Effect.succeed(paused), + getPausedExecution: () => Effect.succeed(paused.execution), + resume: (_id, response) => { + resumes += 1; + return Effect.succeed(completed(response)); + }, + }); + const dispatcher = makeModernMcpDispatcher(() => + Effect.succeed({ engine, description: "test" }), + ); + let retryRequest: DispatcherRequest | undefined; + + await withClient( + dispatcher, + async (client) => { + client.setRequestHandler("elicitation/create", () => ({ + action: "accept", + content: { approved: true }, + })); + const result = await client.callTool({ + name: "execute", + arguments: { code: "await tools.test()" }, + }); + expect(result.structuredContent).toMatchObject({ + status: "completed", + result: { action: "accept", content: { approved: true } }, + }); + expect(resumes).toBe(1); + + expect(retryRequest).toBeDefined(); + const retry = await Effect.runPromise( + dispatcher.dispatch(cloneDispatcherRequest(retryRequest!), principal, defaultMcpResource), + ); + expect(await retry.json()).toMatchObject({ + result: { + structuredContent: { + status: "completed", + result: { action: "accept", content: { approved: true } }, + }, + }, + }); + expect(resumes).toBe(1); + + const otherPrincipal = { ...principal, accountId: "account-2" }; + const crossPrincipal = await Effect.runPromise( + dispatcher.dispatch( + cloneDispatcherRequest(retryRequest!), + otherPrincipal, + defaultMcpResource, + ), + ); + expect(await crossPrincipal.json()).toMatchObject({ + error: { message: expect.stringMatching(/request.?state/i) }, + }); + expect(resumes).toBe(1); + }, + (request) => { + if (request.headers.get("mcp-method") === "tools/call") retryRequest = request; + }, + ); + }); +}); diff --git a/packages/hosts/mcp/src/modern-tool-server.ts b/packages/hosts/mcp/src/modern-tool-server.ts new file mode 100644 index 0000000000..7f7f69133a --- /dev/null +++ b/packages/hosts/mcp/src/modern-tool-server.ts @@ -0,0 +1,487 @@ +import { Cause, Effect, Predicate } from "effect"; +import { + createMcpHandler, + createRequestStateCodec, + inputRequired, + McpServer, + type InputRequiredResult, + type ServerContext, +} from "@modelcontextprotocol/server"; +import * as z from "zod/v4"; + +import { + EXECUTE_SKILL, + findSkill, + renderSkillsIndex, + skillCatalogFor, + type ExecutionEngine, + type ExecutionResult, +} from "@executor-js/execution"; + +import { mcpResourceKey, principalOwns, type McpResource, type Principal } from "./seams"; +import { + formatMcpExecutionFailure, + formatMcpExecutionOutcome, + type McpToolResult, +} from "./tool-server"; + +const MODERN_FLOW_TTL_MS = 10 * 60 * 1000; +const MAX_WORKSPACES = 100; +const MAX_TERMINAL_FLOWS = 1_000; + +interface ModernRequestState { + readonly version: 1; + readonly kind: "executor.execute"; + readonly tool: "execute"; + readonly executionId: string; + readonly argsDigest: string; + readonly phase: "awaiting_input"; +} + +interface ModernFlow { + readonly principal: Principal; + readonly resource: McpResource; + readonly engine: ExecutionEngine; + readonly argsDigest: string; + readonly workspaceKey: string; + touchedAt: number; +} + +interface ModernTerminalFlow { + readonly principal: Principal; + readonly resource: McpResource; + readonly argsDigest: string; + readonly result: McpToolResult; + touchedAt: number; +} + +interface ModernWorkspace { + readonly engine: ExecutionEngine; + readonly description: string; + readonly close?: () => Promise; + touchedAt: number; +} + +export interface ModernMcpBuild { + readonly engine: ExecutionEngine; + readonly description: string; + readonly close?: () => Promise; +} + +export type BuildModernMcp = ( + principal: Principal, + options?: { readonly resource?: McpResource }, +) => Effect.Effect; + +export interface ModernMcpDispatcher { + readonly dispatch: ( + request: Request, + principal: Principal, + resource: McpResource, + ) => Effect.Effect; + readonly close: () => Promise; +} + +const ownerKey = (principal: Principal, resource: McpResource): string => + `${principal.accountId}\0${principal.organizationId}\0${mcpResourceKey(resource)}`; + +const digest = async (value: string): Promise => { + const bytes = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(value)); + return [...new Uint8Array(bytes)].map((byte) => byte.toString(16).padStart(2, "0")).join(""); +}; + +const isModernRequestState = (value: unknown): value is ModernRequestState => { + if (typeof value !== "object" || value === null || Array.isArray(value)) return false; + const state = value as Partial; + return ( + state.version === 1 && + state.kind === "executor.execute" && + state.tool === "execute" && + typeof state.executionId === "string" && + typeof state.argsDigest === "string" && + state.phase === "awaiting_input" + ); +}; + +const unavailableFlowResult = (executionId: string): McpToolResult => ({ + content: [ + { + type: "text", + text: `The paused execution ${executionId} is no longer available. Run execute again to start a fresh flow.`, + }, + ], + structuredContent: { + status: "execution_expired", + executionId, + recovery: "re_execute", + }, + isError: true, +}); + +const invalidFlowResult = (message: string): McpToolResult => ({ + content: [{ type: "text", text: message }], + structuredContent: { status: "invalid_request_state" }, + isError: true, +}); + +const isMcpToolResult = (value: McpToolResult | ExecutionResult): value is McpToolResult => + "content" in value; + +const resumeResponse = ( + ctx: ServerContext, +): { + readonly action: "accept" | "decline" | "cancel"; + readonly content?: Record; +} | null => { + const raw = ctx.mcpReq.inputResponses?.elicitation; + if (typeof raw !== "object" || raw === null || Array.isArray(raw)) return null; + const response = raw as { action?: unknown; content?: unknown }; + if ( + response.action !== "accept" && + response.action !== "decline" && + response.action !== "cancel" + ) { + return null; + } + const content = + typeof response.content === "object" && + response.content !== null && + !Array.isArray(response.content) + ? (response.content as Record) + : undefined; + return content === undefined ? { action: response.action } : { action: response.action, content }; +}; + +const inputRequestFor = ( + execution: Extract["execution"], +) => { + const request = execution.elicitationContext.request; + return Predicate.isTagged(request, "UrlElicitation") + ? inputRequired.elicitUrl({ message: request.message, url: request.url }) + : inputRequired.elicit({ + message: request.message, + requestedSchema: request.requestedSchema as Parameters< + typeof inputRequired.elicit + >[0]["requestedSchema"], + }); +}; + +export const makeModernMcpDispatcher = ( + build: BuildModernMcp, + options?: { + readonly requestStateKey?: Uint8Array | string; + readonly onExecutionPaused?: (executionId: string) => Promise; + readonly onResumeStarted?: (executionId: string) => Promise; + readonly onResumeSettled?: (executionId: string) => Promise; + }, +): ModernMcpDispatcher => { + const requestStateKey = options?.requestStateKey ?? crypto.getRandomValues(new Uint8Array(32)); + const flows = new Map(); + const terminalFlows = new Map(); + const resumeRuns = new Map>(); + const workspaces = new Map(); + const workspaceBuilds = new Map>(); + + const sweep = (): void => { + const cutoff = Date.now() - MODERN_FLOW_TTL_MS; + for (const [executionId, flow] of flows) { + if (flow.touchedAt < cutoff) { + flows.delete(executionId); + void options?.onResumeSettled?.(executionId).then( + () => undefined, + () => undefined, + ); + } + } + for (const [executionId, flow] of terminalFlows) { + if (flow.touchedAt < cutoff) terminalFlows.delete(executionId); + } + while (terminalFlows.size > MAX_TERMINAL_FLOWS) { + const oldest = terminalFlows.keys().next().value; + if (oldest === undefined) break; + terminalFlows.delete(oldest); + } + const activeWorkspaceKeys = new Set([...flows.values()].map((flow) => flow.workspaceKey)); + for (const [key, workspace] of workspaces) { + if (workspace.touchedAt < cutoff && !activeWorkspaceKeys.has(key)) { + workspaces.delete(key); + void workspace.close?.().then( + () => undefined, + () => undefined, + ); + } + } + while (workspaces.size > MAX_WORKSPACES) { + const evictable = [...workspaces.entries()] + .filter(([key]) => !activeWorkspaceKeys.has(key)) + .sort(([, a], [, b]) => a.touchedAt - b.touchedAt)[0]; + if (!evictable) break; + const [oldest, workspace] = evictable; + workspaces.delete(oldest); + void workspace?.close?.().then( + () => undefined, + () => undefined, + ); + } + }; + + const workspaceFor = async ( + principal: Principal, + resource: McpResource, + ): Promise => { + const key = ownerKey(principal, resource); + const existing = workspaces.get(key); + if (existing) { + existing.touchedAt = Date.now(); + return existing; + } + const pending = workspaceBuilds.get(key); + if (pending) return pending; + const building = Effect.runPromise(build(principal, { resource })).then((built) => { + const workspace: ModernWorkspace = { ...built, touchedAt: Date.now() }; + workspaces.set(key, workspace); + return workspace; + }); + workspaceBuilds.set(key, building); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- adapter boundary: clear the in-flight build cache on both fulfillment and rejection without changing the original promise result. + try { + return await building; + } finally { + workspaceBuilds.delete(key); + } + }; + + const dispatch = ( + request: Request, + principal: Principal, + resource: McpResource, + ): Effect.Effect => + Effect.promise(async () => { + sweep(); + const binding = `${ownerKey(principal, resource)}\0tools/call\0execute`; + const codec = createRequestStateCodec({ + key: requestStateKey, + ttlSeconds: MODERN_FLOW_TTL_MS / 1000, + bind: () => binding, + }); + + const serverFactory = () => { + const server = new McpServer( + { name: "executor", version: "1.0.0" }, + { + capabilities: { tools: {} }, + requestState: { verify: codec.verify }, + }, + ); + + server.registerTool( + "execute", + { + description: [ + "Execute TypeScript in Executor's sandbox against the authenticated caller's integrations.", + 'Call `skills({ name: "execute" })` first for the workflow, live integration inventory, and result rules.', + ].join("\n"), + inputSchema: z.object({ code: z.string().trim().min(1) }), + }, + async ({ code }, ctx): Promise => { + const argsDigest = await digest(code); + const decoded = ctx.mcpReq.requestState(); + let outcome: ExecutionResult; + let activeEngine: ExecutionEngine; + + if (decoded !== undefined) { + if (!isModernRequestState(decoded) || decoded.argsDigest !== argsDigest) { + return invalidFlowResult( + "The execute retry does not match its signed request state.", + ); + } + const terminal = terminalFlows.get(decoded.executionId); + if ( + terminal && + principalOwns(terminal.principal, principal) && + mcpResourceKey(terminal.resource) === mcpResourceKey(resource) && + terminal.argsDigest === argsDigest + ) { + terminal.touchedAt = Date.now(); + return terminal.result; + } + const flow = flows.get(decoded.executionId); + if ( + !flow || + !principalOwns(flow.principal, principal) || + mcpResourceKey(flow.resource) !== mcpResourceKey(resource) || + flow.argsDigest !== argsDigest + ) { + return unavailableFlowResult(decoded.executionId); + } + flow.touchedAt = Date.now(); + const workspace = workspaces.get(flow.workspaceKey); + if (workspace) workspace.touchedAt = Date.now(); + activeEngine = flow.engine; + const response = resumeResponse(ctx); + if (!response) { + const paused = await Effect.runPromise( + flow.engine.getPausedExecution(decoded.executionId), + ); + if (!paused) return unavailableFlowResult(decoded.executionId); + return inputRequired({ + inputRequests: { + elicitation: inputRequestFor({ status: "paused", execution: paused }.execution), + }, + requestState: await codec.mint(decoded, ctx), + }); + } + let resumed: ExecutionResult | McpToolResult | null; + let run = resumeRuns.get(decoded.executionId); + if (!run) { + run = (async () => { + await options?.onResumeStarted?.(decoded.executionId); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- adapter boundary: the DO keepalive lease must settle even if an engine adapter rejects unexpectedly. + try { + return await Effect.runPromise( + flow.engine + .resume(decoded.executionId, response) + .pipe( + Effect.catchCause((cause) => + Effect.succeed(formatMcpExecutionFailure(cause)), + ), + ), + ); + } finally { + await options?.onResumeSettled?.(decoded.executionId); + } + })(); + resumeRuns.set(decoded.executionId, run); + void run.then( + () => resumeRuns.delete(decoded.executionId), + () => resumeRuns.delete(decoded.executionId), + ); + } + resumed = await run; + if (resumed === null) return unavailableFlowResult(decoded.executionId); + if (isMcpToolResult(resumed)) { + flows.delete(decoded.executionId); + terminalFlows.set(decoded.executionId, { + principal, + resource, + argsDigest, + result: resumed, + touchedAt: Date.now(), + }); + return resumed; + } + outcome = resumed; + } else { + const workspace = await workspaceFor(principal, resource); + activeEngine = workspace.engine; + const executed = await Effect.runPromise( + workspace.engine + .executeWithPause(code) + .pipe( + Effect.catchCause((cause) => Effect.succeed(formatMcpExecutionFailure(cause))), + ), + ); + if (isMcpToolResult(executed)) return executed; + outcome = executed; + } + + if (outcome.status === "completed") { + const result = formatMcpExecutionOutcome(outcome); + if (decoded !== undefined) { + flows.delete(decoded.executionId); + terminalFlows.set(decoded.executionId, { + principal, + resource, + argsDigest, + result, + touchedAt: Date.now(), + }); + } + return result; + } + + const state: ModernRequestState = { + version: 1, + kind: "executor.execute", + tool: "execute", + executionId: outcome.execution.id, + argsDigest, + phase: "awaiting_input", + }; + if (decoded !== undefined) flows.delete(decoded.executionId); + flows.set(outcome.execution.id, { + principal, + resource, + engine: activeEngine, + argsDigest, + workspaceKey: ownerKey(principal, resource), + touchedAt: Date.now(), + }); + await options?.onExecutionPaused?.(outcome.execution.id); + return inputRequired({ + inputRequests: { elicitation: inputRequestFor(outcome.execution) }, + requestState: await codec.mint(state, ctx), + }); + }, + ); + + server.registerTool( + "skills", + { + description: "Fetch Executor's execution guide and live integration inventory.", + inputSchema: z.object({ name: z.string().optional() }), + }, + async ({ name }): Promise => { + const catalog = skillCatalogFor({ artifacts: false }); + const trimmed = name?.trim(); + if (!trimmed) return { content: [{ type: "text", text: renderSkillsIndex(catalog) }] }; + const skill = findSkill(trimmed, catalog); + if (!skill) { + return { + content: [ + { + type: "text", + text: `No skill named "${trimmed}".\n\n${renderSkillsIndex(catalog)}`, + }, + ], + isError: true, + }; + } + if (skill.name !== EXECUTE_SKILL.name) { + return { content: [{ type: "text", text: skill.body }] }; + } + const workspace = await workspaceFor(principal, resource); + return { + content: [{ type: "text", text: `${skill.body}\n\n${workspace.description}` }], + }; + }, + ); + + return server; + }; + + const handler = createMcpHandler(serverFactory, { legacy: "reject" }); + return handler.fetch(request); + }); + + return { + dispatch, + close: async () => { + const flowIds = [...flows.keys()]; + flows.clear(); + terminalFlows.clear(); + resumeRuns.clear(); + workspaceBuilds.clear(); + const closes = [...workspaces.values()].flatMap((workspace) => + workspace.close ? [workspace.close()] : [], + ); + workspaces.clear(); + await Promise.allSettled([ + ...closes, + ...flowIds.flatMap((executionId) => + options?.onResumeSettled ? [options.onResumeSettled(executionId)] : [], + ), + ]); + }, + }; +}; diff --git a/packages/hosts/mcp/src/seams.ts b/packages/hosts/mcp/src/seams.ts index 12e713dd91..1821af81d6 100644 --- a/packages/hosts/mcp/src/seams.ts +++ b/packages/hosts/mcp/src/seams.ts @@ -255,6 +255,14 @@ export class McpSessionStore extends Context.Service< * `"not-found"` / `"forbidden"` discriminant for the envelope to render. */ readonly dispatch: (input: McpDispatchInput) => Effect.Effect; + /** + * Serve one stateless 2026-07-28 request. Optional during staged rollout: + * hosts that have not installed the modern adapter receive an explicit + * unsupported-protocol response while their legacy sessions continue to work. + */ + readonly dispatchModern?: ( + input: Omit, + ) => Effect.Effect; /** * Tear down a session by id (idempotent). * diff --git a/packages/hosts/mcp/src/tool-server.ts b/packages/hosts/mcp/src/tool-server.ts index 41c61bcd0a..575bff6b9b 100644 --- a/packages/hosts/mcp/src/tool-server.ts +++ b/packages/hosts/mcp/src/tool-server.ts @@ -671,7 +671,7 @@ const formatResumeApprovalRequired = (input: { }, }); -const toMcpFailureResult = (cause: Cause.Cause): McpToolResult => { +export const formatMcpExecutionFailure = (cause: Cause.Cause): McpToolResult => { const correlationId = newCorrelationId(); const defect = Cause.findDefect(cause); const nativeElicitationFailed = @@ -784,7 +784,7 @@ const fallbackOutcomeResult = ( // The catalog is per-session: a connection that opted out of artifacts never // sees the artifact skills, so the index cannot advertise a how-to for tools it // does not have, and fetching one by name misses like any unknown skill. -const skillsResult = ( +export const skillsResult = ( name: string | undefined, executeInventory: string, catalog: readonly Skill[], @@ -1164,7 +1164,7 @@ export const createExecutorMcpServer = ( const runToolEffect = (effect: Effect.Effect) => Effect.runPromiseWith(context)( anchor(effect).pipe( - Effect.catchCause((cause) => Effect.succeed(toMcpFailureResult(cause))), + Effect.catchCause((cause) => Effect.succeed(formatMcpExecutionFailure(cause))), ), ); diff --git a/packages/plugins/mcp/package.json b/packages/plugins/mcp/package.json index d899b77212..13aedbbfb2 100644 --- a/packages/plugins/mcp/package.json +++ b/packages/plugins/mcp/package.json @@ -65,6 +65,8 @@ "@effect/platform-node": "catalog:", "@executor-js/config": "workspace:*", "@executor-js/sdk": "workspace:*", + "@modelcontextprotocol/client": "2.0.0", + "@modelcontextprotocol/core": "2.0.0", "@modelcontextprotocol/sdk": "^1.29.0", "zod": "4.3.6" }, diff --git a/packages/plugins/mcp/src/sdk/connection-pool.ts b/packages/plugins/mcp/src/sdk/connection-pool.ts index 6b732caf78..292689f62b 100644 --- a/packages/plugins/mcp/src/sdk/connection-pool.ts +++ b/packages/plugins/mcp/src/sdk/connection-pool.ts @@ -21,7 +21,7 @@ const closeQuietly = (connection: McpConnection): Effect.Effect => const isMcpInvocationError = (error: unknown): error is McpInvocationError => Predicate.isTagged(error, "McpInvocationError"); -const isDeadConnectionFailure = (error: unknown): boolean => { +const isDeadConnectionFailure = (connection: McpConnection, error: unknown): boolean => { if (Predicate.isTagged(error, "McpConnectionError")) return true; if (!isMcpInvocationError(error)) return false; return ( @@ -34,16 +34,16 @@ const isDeadConnectionFailure = (error: unknown): boolean => { // freshly minted token can only reach the server over a NEW session. Drop // it so the retry dials with the new credential. error.status === 401 || - error.status === 404 || + (error.status === 404 && connection.protocolEra() !== "modern") || error.status === 408 ); }; -const shouldDropConnection = (exit: Exit.Exit): boolean => { +const shouldDropConnection = (connection: McpConnection, exit: Exit.Exit): boolean => { if (Exit.isSuccess(exit)) return false; const failures = exit.cause.reasons.filter(Cause.isFailReason); if (failures.length === 0) return true; - return failures.some((failure) => isDeadConnectionFailure(failure.error)); + return failures.some((failure) => isDeadConnectionFailure(connection, failure.error)); }; /** A per-plugin-instance pool that gives each invocation an exclusive MCP @@ -92,7 +92,7 @@ export const createMcpConnectionPool = (): McpConnectionPool => { const release = (key: string, lease: ConnectionLease, exit: Exit.Exit) => Effect.gen(function* () { - if (shouldDropConnection(exit)) { + if (shouldDropConnection(lease.connection, exit)) { yield* closeQuietly(lease.connection); return; } @@ -107,11 +107,15 @@ export const createMcpConnectionPool = (): McpConnectionPool => { const withConnection: McpConnectionPool["withConnection"] = (key, connector, use) => { let reused = false; + let reusedEra: ReturnType; const run = (forceFresh: boolean) => Effect.acquireUseRelease( acquire(key, connector, forceFresh), (lease) => { - if (!forceFresh) reused = lease.reused; + if (!forceFresh) { + reused = lease.reused; + reusedEra = lease.connection.protocolEra(); + } return use(lease.connection); }, (lease, exit) => release(key, lease, exit), @@ -119,7 +123,7 @@ export const createMcpConnectionPool = (): McpConnectionPool => { return run(false).pipe( Effect.catch((error) => - reused && isMcpInvocationError(error) && error.status === 404 + reused && reusedEra !== "modern" && isMcpInvocationError(error) && error.status === 404 ? run(true) : Effect.fail(error), ), diff --git a/packages/plugins/mcp/src/sdk/connection.ts b/packages/plugins/mcp/src/sdk/connection.ts index 82a98cd0d1..ca80e416e5 100644 --- a/packages/plugins/mcp/src/sdk/connection.ts +++ b/packages/plugins/mcp/src/sdk/connection.ts @@ -1,14 +1,15 @@ -import type { OAuthClientProvider } from "@modelcontextprotocol/sdk/client/auth.js"; -import { Client } from "@modelcontextprotocol/sdk/client/index.js"; -import { SSEClientTransport } from "@modelcontextprotocol/sdk/client/sse.js"; -import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js"; -import type { FetchLike } from "@modelcontextprotocol/sdk/shared/transport.js"; -import { CfWorkerJsonSchemaValidator } from "@modelcontextprotocol/sdk/validation/cfworker"; +import { CfWorkerJsonSchemaValidator } from "@modelcontextprotocol/client/validators/cf-worker"; +import { + Client, + SSEClientTransport, + StreamableHTTPClientTransport, +} from "@modelcontextprotocol/client"; +import type { FetchLike, OAuthClientProvider, ProtocolEra } from "@modelcontextprotocol/client"; import { Effect, Layer, Predicate, Stream } from "effect"; import { HttpClient, HttpClientRequest } from "effect/unstable/http"; // NOTE: `StdioClientTransport` is NOT imported eagerly. The upstream module -// (`@modelcontextprotocol/sdk/client/stdio.js`) touches `node:child_process` +// (`@modelcontextprotocol/client/stdio`) touches `node:child_process` // at evaluation time, which crashes workerd (incl. vitest-pool-workers) at // SIGSEGV on module instantiation. Cloud callers set // `dangerouslyAllowStdioMCP: false` and never reach the stdio branch below; @@ -30,6 +31,8 @@ import { detectInsufficientScope } from "@executor-js/sdk/core"; export type McpConnection = { readonly client: Client; + readonly protocolEra: () => ProtocolEra | undefined; + readonly setToolListChangedHandler: (handler: (() => void) | undefined) => void; readonly close: () => Promise; }; @@ -201,18 +204,36 @@ const fetchFromHttpClientLayer = ( // MCP plugin runs inside a Cloudflare Worker (executor.sh). The // cfworker validator does not use code generation and works in every // runtime we ship to. -const createClient = (): Client => - new Client( +const createClient = (): Pick => { + let onToolListChanged: (() => void) | undefined; + const client = new Client( { name: "executor-mcp", version: "0.1.0" }, { capabilities: { elicitation: { form: {}, url: {} } }, jsonSchemaValidator: new CfWorkerJsonSchemaValidator(), + versionNegotiation: { mode: "auto" }, + listChanged: { + tools: { + autoRefresh: false, + onChanged: () => onToolListChanged?.(), + }, + }, }, ); + return { + client, + setToolListChangedHandler: (handler) => { + onToolListChanged = handler; + }, + }; +}; -const connectionFromClient = (client: Client): McpConnection => ({ - client, - close: () => client.close(), +const connectionFromClient = ( + built: Pick, +): McpConnection => ({ + ...built, + protocolEra: () => built.client.getProtocolEra(), + close: () => built.client.close(), }); const connectionFailure = ( @@ -249,7 +270,8 @@ const connectClient = (input: { createTransport: () => Parameters[0]; }): Effect.Effect => Effect.gen(function* () { - const client = createClient(); + const built = createClient(); + const client = built.client; const transportInstance = input.createTransport(); yield* Effect.tryPromise({ @@ -262,7 +284,7 @@ const connectClient = (input: { }), ); - return connectionFromClient(client); + return connectionFromClient(built); }); // --------------------------------------------------------------------------- diff --git a/packages/plugins/mcp/src/sdk/http-status.ts b/packages/plugins/mcp/src/sdk/http-status.ts index a4442f8d23..31bdff3c4e 100644 --- a/packages/plugins/mcp/src/sdk/http-status.ts +++ b/packages/plugins/mcp/src/sdk/http-status.ts @@ -1,7 +1,9 @@ +import { SdkHttpError } from "@modelcontextprotocol/client"; + // --------------------------------------------------------------------------- -// Extract the HTTP status from an MCP SDK transport error. The SDK surfaces -// transport failures two ways: a `StreamableHTTPError` subclass carrying a -// numeric `code`, and an SSE POST failure whose message embeds `(HTTP nnn)`. +// Extract the HTTP status from an MCP SDK transport error. The v2 SDK exposes +// it on `SdkHttpError.status`; legacy SSE POST failures still encode the status +// in their message. // Shared by the invoke path (classifies tool-call failures) and the connect // path (so a 401/403 during the handshake reaches the liveness health check). // --------------------------------------------------------------------------- @@ -9,8 +11,6 @@ import { Option, Schema } from "effect"; import { insufficientScopeFromEmbeddedJson } from "@executor-js/sdk/core"; -import { StreamableHTTPError } from "@modelcontextprotocol/sdk/client/streamableHttp.js"; - const SsePostErrorCause = Schema.Struct({ message: Schema.String }); const decodeSsePostErrorCause = Schema.decodeUnknownOption(SsePostErrorCause); @@ -28,9 +28,9 @@ const statusFromSsePostError = (cause: unknown): number | undefined => const statusFromStreamableHttpError = (cause: unknown): number | undefined => { // oxlint-disable-next-line executor/no-instanceof-tagged-error -- boundary: MCP SDK exposes transport HTTP failures as this Error subclass; protocol errors can carry the same numeric code - if (!(cause instanceof StreamableHTTPError)) return undefined; - const code = cause.code; - return code !== undefined && code >= 100 && code <= 599 ? code : undefined; + if (!(cause instanceof SdkHttpError)) return undefined; + const status = cause.status; + return status >= 100 && status <= 599 ? status : undefined; }; export const httpStatusFromCause = (cause: unknown): number | undefined => @@ -42,8 +42,8 @@ export const httpStatusFromCause = (cause: unknown): number | undefined => // StreamableHTTP transport consumes the insufficient_scope challenge ITSELF: // it re-runs auth requesting the broader scope, and only when that upscoped // retry still 403s does it throw — with the fixed message matched below -// (verified against @modelcontextprotocol/sdk streamableHttp.js; re-verify on -// SDK bumps). Both paths mean the same thing: the grant does not cover the +// (verified against @modelcontextprotocol/client v2; re-verify on SDK bumps). +// Both paths mean the same thing: the grant does not cover the // operation, and re-running the identical flow cannot help. Strict matching // (exact serialized field forms via the shared core detector, or the SDK's // exact upscoping message) — a miss stays on the generic auth path. diff --git a/packages/plugins/mcp/src/sdk/invoke.test.ts b/packages/plugins/mcp/src/sdk/invoke.test.ts index 3b5aaae66f..5322f7bb47 100644 --- a/packages/plugins/mcp/src/sdk/invoke.test.ts +++ b/packages/plugins/mcp/src/sdk/invoke.test.ts @@ -2,9 +2,13 @@ import { describe, expect, it } from "@effect/vitest"; import { Effect, Predicate } from "effect"; import { HttpServerResponse } from "effect/unstable/http"; -import type { OAuthClientProvider } from "@modelcontextprotocol/sdk/client/auth.js"; -import { StreamableHTTPError } from "@modelcontextprotocol/sdk/client/streamableHttp.js"; -import { McpError } from "@modelcontextprotocol/sdk/types.js"; +import { + ProtocolError, + ProtocolErrorCode, + SdkErrorCode, + SdkHttpError, + type OAuthClientProvider, +} from "@modelcontextprotocol/client"; import { ElicitationResponse } from "@executor-js/sdk"; import { serveTestHttpApp } from "@executor-js/sdk/testing"; @@ -22,6 +26,8 @@ const rejectingConnector = (cause: unknown): McpConnector => // oxlint-disable-next-line executor/no-promise-reject -- boundary: fake MCP client rejects to exercise invocation error wrapping callTool: () => Promise.reject(cause), } as unknown as McpConnection["client"], + protocolEra: () => "legacy", + setToolListChangedHandler: () => undefined, close: () => Promise.resolve(), }); @@ -108,14 +114,16 @@ const invocationRejectionCases = [ name: "wraps callTool rejection with a stable message and status", toolId: "blocked", transport: "streamable-http", - cause: new StreamableHTTPError(401, "token=do-not-leak"), + cause: new SdkHttpError(SdkErrorCode.ClientHttpAuthentication, "token=do-not-leak", { + status: 401, + }), expectedStatus: 401 as number | undefined, }, { name: "does not treat MCP protocol error codes as HTTP statuses", toolId: "protocol_error", transport: "streamable-http", - cause: new McpError(401, "application-level do-not-leak"), + cause: new ProtocolError(401 as ProtocolErrorCode, "application-level do-not-leak"), expectedStatus: undefined, }, { diff --git a/packages/plugins/mcp/src/sdk/invoke.ts b/packages/plugins/mcp/src/sdk/invoke.ts index 3e5aa7695f..0b53c5de7e 100644 --- a/packages/plugins/mcp/src/sdk/invoke.ts +++ b/packages/plugins/mcp/src/sdk/invoke.ts @@ -1,4 +1,5 @@ // --------------------------------------------------------------------------- +import { ProtocolError, ProtocolErrorCode } from "@modelcontextprotocol/client"; // MCP tool invocation — shared helper called from plugin.invokeTool. // // Responsible for: @@ -15,14 +16,6 @@ // --------------------------------------------------------------------------- import { Cause, Effect, Exit, Option, Predicate, Schema } from "effect"; - -import { - ElicitRequestSchema, - ErrorCode, - McpError, - ToolListChangedNotificationSchema, -} from "@modelcontextprotocol/sdk/types.js"; - import { ElicitationId, FormElicitation, @@ -66,8 +59,9 @@ export const isUnknownToolMessage = (message: string, toolName: string): boolean const isUnknownToolCause = (cause: unknown, toolName: string): boolean => // oxlint-disable-next-line executor/no-instanceof-tagged-error -- boundary: MCP SDK surfaces JSON-RPC protocol errors as this Error subclass - cause instanceof McpError && - (cause.code === ErrorCode.InvalidParams || cause.code === ErrorCode.MethodNotFound) && + cause instanceof ProtocolError && + (cause.code === ProtocolErrorCode.InvalidParams || + cause.code === ProtocolErrorCode.MethodNotFound) && // oxlint-disable-next-line executor/no-unknown-error-message -- boundary: instanceof narrows to the SDK's McpError, whose message carries the only unknown-tool discriminator the protocol provides isUnknownToolMessage(cause.message, toolName); @@ -93,6 +87,24 @@ const McpElicitParams = Schema.Union([ type McpElicitParams = typeof McpElicitParams.Type; const decodeElicitParams = Schema.decodeUnknownSync(McpElicitParams); +const McpFormContent = Schema.Record( + Schema.String, + Schema.Union([Schema.String, Schema.Number, Schema.Boolean, Schema.Array(Schema.String)]), +); +const decodeMcpFormContent = Schema.decodeUnknownSync(McpFormContent); +const mutableMcpFormValue = ( + value: string | number | boolean | ReadonlyArray, +): string | number | boolean | string[] => + Array.isArray(value) ? Array.from(value) : (value as string | number | boolean); +const mutableMcpFormContent = ( + content: unknown, +): Record => + Object.fromEntries( + Object.entries(decodeMcpFormContent(content)).map(([key, value]) => [ + key, + mutableMcpFormValue(value), + ]), + ); const toElicitationRequest = (params: McpElicitParams): ElicitationRequest => params.mode === "url" @@ -107,7 +119,7 @@ const toElicitationRequest = (params: McpElicitParams): ElicitationRequest => }); const installElicitationHandler = (client: McpConnection["client"], elicit: Elicit): void => { - client.setRequestHandler(ElicitRequestSchema, async (request: { params: unknown }) => { + client.setRequestHandler("elicitation/create", async (request: { params: unknown }) => { const params = decodeElicitParams(request.params); const req = toElicitationRequest(params); // Use runPromiseExit so we can inspect typed failures — `elicit` @@ -119,7 +131,9 @@ const installElicitationHandler = (client: McpConnection["client"], elicit: Elic const response = exit.value; return { action: response.action, - ...(response.action === "accept" && response.content ? { content: response.content } : {}), + ...(response.action === "accept" && response.content + ? { content: mutableMcpFormContent(response.content) } + : {}), }; } const failure = exit.cause.reasons.find(Cause.isFailReason); @@ -145,13 +159,15 @@ const installElicitationHandler = (client: McpConnection["client"], elicit: Elic // --------------------------------------------------------------------------- const installToolListChangedHandler = ( - client: McpConnection["client"], + connection: McpConnection, onToolListChanged: (() => void) | undefined, ): void => { - if (!onToolListChanged) return; - client.setNotificationHandler(ToolListChangedNotificationSchema, () => { - onToolListChanged(); - }); + connection.setToolListChangedHandler(onToolListChanged); + if (connection.protocolEra() !== "modern" && onToolListChanged) { + connection.client.setNotificationHandler("notifications/tools/list_changed", () => { + onToolListChanged(); + }); + } }; // --------------------------------------------------------------------------- @@ -167,7 +183,7 @@ const useConnection = ( ): Effect.Effect => Effect.gen(function* () { installElicitationHandler(connection.client, elicit); - installToolListChangedHandler(connection.client, onToolListChanged); + installToolListChangedHandler(connection, onToolListChanged); return yield* Effect.tryPromise({ try: () => connection.client.callTool({ name: toolName, arguments: args }), catch: (cause) => { @@ -190,7 +206,7 @@ const useConnection = ( } const status = httpStatusFromCause(cause); // oxlint-disable-next-line executor/no-instanceof-tagged-error -- boundary: MCP SDK protocol failures are its McpError subclass; transport failures use other error shapes - const protocolFailure = cause instanceof McpError; + const protocolFailure = cause instanceof ProtocolError; return new McpInvocationError({ toolName, message: `MCP tool call failed for ${toolName}`, diff --git a/packages/plugins/mcp/src/sdk/plugin.ts b/packages/plugins/mcp/src/sdk/plugin.ts index e3b7a6857e..b50d67d667 100644 --- a/packages/plugins/mcp/src/sdk/plugin.ts +++ b/packages/plugins/mcp/src/sdk/plugin.ts @@ -1,9 +1,7 @@ import { Effect, Layer, Option, Result, Schema } from "effect"; import type { HttpClient } from "effect/unstable/http"; - -import type { OAuthClientProvider } from "@modelcontextprotocol/sdk/client/auth.js"; -import { CallToolResultSchema } from "@modelcontextprotocol/sdk/types.js"; -import * as z from "zod/v4"; +import { CallToolResultSchema } from "@modelcontextprotocol/core"; +import type { OAuthClientProvider } from "@modelcontextprotocol/client"; import { authToolFailure, @@ -392,11 +390,13 @@ type JsonSchemaObject = Record & { readonly properties?: Record; }; -const McpCallToolResultJsonSchema = z.toJSONSchema(CallToolResultSchema) as JsonSchemaObject; +const McpCallToolResultJsonSchema = CallToolResultSchema.toJSONSchema() as JsonSchemaObject; const mcpCallToolResultOutputSchema = (structuredContentSchema?: unknown): JsonSchemaObject => { - const defaultStructuredContentSchema = - McpCallToolResultJsonSchema.properties?.structuredContent ?? {}; + const defaultStructuredContentSchema = { + type: "object", + additionalProperties: true, + }; return { ...McpCallToolResultJsonSchema, diff --git a/packages/plugins/mcp/src/sdk/probe-shape.test.ts b/packages/plugins/mcp/src/sdk/probe-shape.test.ts index 7eb124ebaf..1eab783674 100644 --- a/packages/plugins/mcp/src/sdk/probe-shape.test.ts +++ b/packages/plugins/mcp/src/sdk/probe-shape.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it } from "@effect/vitest"; -import { Effect, Ref } from "effect"; +import { Effect, Ref, Schema } from "effect"; import { HttpServerResponse } from "effect/unstable/http"; import { serveTestHttpApp } from "@executor-js/sdk/testing"; @@ -14,6 +14,13 @@ interface CapturedProbeRequest { type ProbeHandler = (request: CapturedProbeRequest) => HttpServerResponse.HttpServerResponse; +const decodeJson = Schema.decodeUnknownSync(Schema.fromJsonString(Schema.Unknown)); +const decodeProbe = Schema.decodeUnknownSync( + Schema.fromJsonString( + Schema.Struct({ id: Schema.optional(Schema.Unknown), method: Schema.optional(Schema.String) }), + ), +); + const serveProbeEndpoint = (handler: ProbeHandler) => Effect.gen(function* () { const requests = yield* Ref.make([]); @@ -50,6 +57,81 @@ const withServer = (handler: ProbeHandler, use: (endpoint: string) => Effe ); describe("probeMcpEndpointShape", () => { + it.effect("probes server/discover with the modern per-request envelope first", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveProbeEndpoint((request) => { + const body = decodeProbe(request.body); + if (body.method !== "server/discover") { + return HttpServerResponse.empty({ status: 500 }); + } + return HttpServerResponse.jsonUnsafe({ + jsonrpc: "2.0", + id: 1, + result: { + resultType: "complete", + supportedVersions: ["2026-07-28"], + capabilities: { tools: {} }, + }, + }); + }); + + const result = yield* probeMcpEndpointShape(server.endpoint); + const requests = yield* server.requests; + + expect(result).toEqual({ kind: "mcp", requiresAuth: false }); + expect(requests).toHaveLength(1); + expect(requests[0]?.headers["mcp-protocol-version"]).toBe("2026-07-28"); + expect(requests[0]?.headers["mcp-method"]).toBe("server/discover"); + expect(decodeJson(requests[0]?.body ?? "{}")).toMatchObject({ + method: "server/discover", + params: { + _meta: { + "io.modelcontextprotocol/protocolVersion": "2026-07-28", + "io.modelcontextprotocol/clientInfo": { name: "executor-probe", version: "0" }, + "io.modelcontextprotocol/clientCapabilities": {}, + }, + }, + }); + }), + ), + ); + + it.effect("falls back to legacy initialize when discover is unsupported", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveProbeEndpoint((request) => { + const body = decodeProbe(request.body); + if (body.method === "server/discover") { + return HttpServerResponse.jsonUnsafe({ + jsonrpc: "2.0", + id: body.id, + error: { code: -32601, message: "Method not found" }, + }); + } + return HttpServerResponse.jsonUnsafe({ + jsonrpc: "2.0", + id: body.id, + result: { + protocolVersion: "2025-06-18", + capabilities: {}, + serverInfo: { name: "legacy", version: "0" }, + }, + }); + }); + + const result = yield* probeMcpEndpointShape(server.endpoint); + const requests = yield* server.requests; + + expect(result).toEqual({ kind: "mcp", requiresAuth: false }); + expect(requests.map((request) => decodeProbe(request.body).method)).toEqual([ + "server/discover", + "initialize", + ]); + }), + ), + ); + it.effect("classifies 2xx as unauth-OK MCP", () => withServer( () => diff --git a/packages/plugins/mcp/src/sdk/probe-shape.ts b/packages/plugins/mcp/src/sdk/probe-shape.ts index 3b3be94297..868b02b41e 100644 --- a/packages/plugins/mcp/src/sdk/probe-shape.ts +++ b/packages/plugins/mcp/src/sdk/probe-shape.ts @@ -36,6 +36,26 @@ import { Data, Duration, Effect, Layer, Option, Schema } from "effect"; import { FetchHttpClient, HttpClient, HttpClientRequest } from "effect/unstable/http"; +import { + CLIENT_CAPABILITIES_META_KEY, + CLIENT_INFO_META_KEY, + PROTOCOL_VERSION_META_KEY, +} from "@modelcontextprotocol/client"; + +const MODERN_PROTOCOL_VERSION = "2026-07-28"; + +const DISCOVER_BODY = JSON.stringify({ + jsonrpc: "2.0", + id: 1, + method: "server/discover", + params: { + _meta: { + [PROTOCOL_VERSION_META_KEY]: MODERN_PROTOCOL_VERSION, + [CLIENT_INFO_META_KEY]: { name: "executor-probe", version: "0" }, + [CLIENT_CAPABILITIES_META_KEY]: {}, + }, + }, +}); /** MCP initialize request body used as the shape probe. Any real MCP * server either answers it (unauth-OK server) or returns the spec- @@ -95,6 +115,20 @@ const isJsonRpcEnvelope = (body: string): boolean => { return "result" in obj || "error" in obj || "method" in obj; }; +const isDiscoverResultEnvelope = (body: string): boolean => { + const envelope = asObject(body); + if (!envelope || envelope.jsonrpc !== "2.0") return false; + const result = envelope.result; + if (typeof result !== "object" || result === null || Array.isArray(result)) return false; + const supportedVersions = (result as Record).supportedVersions; + return ( + Array.isArray(supportedVersions) && + supportedVersions.some( + (version) => typeof version === "string" && version >= MODERN_PROTOCOL_VERSION, + ) + ); +}; + /** Quick check that a body parses as an RFC 6750 OAuth Bearer error * envelope (`{error: "invalid_token", error_description?: ..., ...}`). * Real MCP servers like Atlassian return this shape on unauth requests @@ -267,6 +301,7 @@ export const probeMcpEndpointShape = ( readonly text: Effect.Effect; }, method: "GET" | "POST", + expected: "discover" | "initialize", ): Effect.Effect => Effect.gen(function* () { const contentType = readHeader(response.headers, "content-type") ?? ""; @@ -354,7 +389,10 @@ export const probeMcpEndpointShape = ( // JSON-RPC envelope so we don't accept HTML/REST 200 responses. if (isSse) return { kind: "mcp", requiresAuth: false } as const; const body = yield* readBody(response); - if (!isJsonRpcEnvelope(body)) { + if ( + expected === "discover" ? !isDiscoverResultEnvelope(body) : !isJsonRpcEnvelope(body) + ) { + if (expected === "discover") return null; return { kind: "not-mcp", category: "wrong-shape", @@ -372,6 +410,25 @@ export const probeMcpEndpointShape = ( url.searchParams.set(key, value); } + let discoverRequest = HttpClientRequest.post(url.toString()).pipe( + HttpClientRequest.setHeader("content-type", "application/json"), + HttpClientRequest.setHeader("accept", "application/json, text/event-stream"), + HttpClientRequest.bodyText(DISCOVER_BODY, "application/json"), + ); + for (const [name, value] of Object.entries(options.headers ?? {})) { + discoverRequest = HttpClientRequest.setHeader(discoverRequest, name, value); + } + discoverRequest = discoverRequest.pipe( + HttpClientRequest.setHeader("mcp-protocol-version", MODERN_PROTOCOL_VERSION), + HttpClientRequest.setHeader("mcp-method", "server/discover"), + ); + + const discoverResponse = yield* client + .execute(discoverRequest) + .pipe(Effect.timeout(Duration.millis(timeoutMs))); + const discoverResult = yield* classify(discoverResponse, "POST", "discover"); + if (discoverResult) return discoverResult; + let postRequest = HttpClientRequest.post(url.toString()).pipe( HttpClientRequest.setHeader("content-type", "application/json"), HttpClientRequest.setHeader("accept", "application/json, text/event-stream"), @@ -385,7 +442,7 @@ export const probeMcpEndpointShape = ( .execute(postRequest) .pipe(Effect.timeout(Duration.millis(timeoutMs))); - const postResult = yield* classify(postResponse, "POST"); + const postResult = yield* classify(postResponse, "POST", "initialize"); if (postResult) return postResult; if ([404, 405, 406, 415].includes(postResponse.status)) { @@ -398,7 +455,7 @@ export const probeMcpEndpointShape = ( const getResponse = yield* client .execute(getRequest) .pipe(Effect.timeout(Duration.millis(timeoutMs))); - const getResult = yield* classify(getResponse, "GET"); + const getResult = yield* classify(getResponse, "GET", "initialize"); if (getResult) return getResult; } diff --git a/packages/plugins/mcp/src/sdk/stdio-connector.ts b/packages/plugins/mcp/src/sdk/stdio-connector.ts index 99a0f72e37..501b2001f8 100644 --- a/packages/plugins/mcp/src/sdk/stdio-connector.ts +++ b/packages/plugins/mcp/src/sdk/stdio-connector.ts @@ -1,3 +1,5 @@ +import { StdioClientTransport } from "@modelcontextprotocol/client/stdio"; + // --------------------------------------------------------------------------- // Stdio transport factory — loaded only on demand // --------------------------------------------------------------------------- @@ -12,9 +14,6 @@ // in `connection.ts`. Remote-only consumers (cloud/marketing) never execute // the import and therefore never touch `node:child_process`. // --------------------------------------------------------------------------- - -import { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js"; - export type StdioTransportConfig = { readonly command: string; readonly args?: ReadonlyArray; diff --git a/packages/plugins/mcp/src/testing/server.ts b/packages/plugins/mcp/src/testing/server.ts index 0b6edf0f04..7558a0c6e4 100644 --- a/packages/plugins/mcp/src/testing/server.ts +++ b/packages/plugins/mcp/src/testing/server.ts @@ -160,11 +160,11 @@ export const serveMcpServer = (factory: () => McpServer, options: McpTestServerO const transport = new StreamableHTTPServerTransport({ sessionIdGenerator: () => crypto.randomUUID(), onsessioninitialized: (sid) => { + sessions += 1; transports.set(sid, transport); }, }); allTransports.add(transport); - sessions += 1; const mcpServer = factory(); yield* Effect.tryPromise({ From 927cea4b1f5c585a8ee685568fce88761e222284 Mon Sep 17 00:00:00 2001 From: Schreezer <96696271+Schreezer@users.noreply.github.com> Date: Tue, 11 Aug 2026 12:24:48 +0530 Subject: [PATCH 2/2] fix(mcp): harden rollback and cloud lifecycle --- apps/cloud/src/mcp/agent-handler.ts | 2 +- apps/host-cloudflare/src/mcp/agent-handler.ts | 2 +- .../src/worker.e2e.node.test.ts | 34 ++++ apps/local/src/mcp.ts | 2 +- .../mcp/agent-session-durable-object.test.ts | 191 +++++++++++++++++- .../src/mcp/agent-session-durable-object.ts | 179 ++++++++++++---- packages/hosts/mcp/src/envelope.test.ts | 19 +- packages/hosts/mcp/src/envelope.ts | 10 +- 8 files changed, 376 insertions(+), 63 deletions(-) diff --git a/apps/cloud/src/mcp/agent-handler.ts b/apps/cloud/src/mcp/agent-handler.ts index f7c5c4cffd..3d6612b520 100644 --- a/apps/cloud/src/mcp/agent-handler.ts +++ b/apps/cloud/src/mcp/agent-handler.ts @@ -207,7 +207,7 @@ export const makeCloudMcpAgentHandler = () => { const resource = resourceFromPath(request); if (!(await isLegacyMcpRequest(request))) { if (env.MCP_2026_07_28_ENABLED === "false") { - return jsonRpcResponse(503, -32022, "MCP 2026-07-28 support is disabled"); + return jsonRpcResponse(400, -32022, "MCP 2026-07-28 support is disabled"); } const props = await runTraced( request, diff --git a/apps/host-cloudflare/src/mcp/agent-handler.ts b/apps/host-cloudflare/src/mcp/agent-handler.ts index a545ab3f87..2fb2836f9b 100644 --- a/apps/host-cloudflare/src/mcp/agent-handler.ts +++ b/apps/host-cloudflare/src/mcp/agent-handler.ts @@ -136,7 +136,7 @@ export const makeCloudflareMcpAgentHandler = (config: CloudflareConfig) => { if (!(await isLegacyMcpRequest(request))) { if (env.MCP_2026_07_28_ENABLED === "false") { - return jsonRpcResponse(503, -32022, "MCP 2026-07-28 support is disabled"); + return jsonRpcResponse(400, -32022, "MCP 2026-07-28 support is disabled"); } const props = await Effect.runPromise(propsForPrincipal(request, outcome.principal)); const flowId = JSON.stringify([ diff --git a/apps/host-cloudflare/src/worker.e2e.node.test.ts b/apps/host-cloudflare/src/worker.e2e.node.test.ts index 0fed0b27f5..9a8b383258 100644 --- a/apps/host-cloudflare/src/worker.e2e.node.test.ts +++ b/apps/host-cloudflare/src/worker.e2e.node.test.ts @@ -445,6 +445,40 @@ describe("cloudflare host e2e (workerd/miniflare)", () => { } }, 60_000); + it("falls back to legacy MCP when modern support is rolled back", async () => { + const rollbackWorker = await unstable_dev(resolve(dir, "worker.ts"), { + config: resolve(dir, "../wrangler.jsonc"), + ip: "127.0.0.1", + local: true, + persist: false, + experimental: { disableExperimentalWarning: true }, + vars: { + EXECUTOR_SECRET_KEY: "test-secret-key-0123456789abcdef", + ENABLE_DEV_AUTH: "true", + MCP_2026_07_28_ENABLED: "false", + }, + }); + const transport = new ModernStreamableHTTPClientTransport( + new URL("/mcp", `http://${rollbackWorker.address}:${rollbackWorker.port}`), + ); + const client = new ModernClient( + { name: "cloudflare-rollback-test", version: "1" }, + { versionNegotiation: { mode: "auto" } }, + ); + + // oxlint-disable-next-line executor/no-try-catch-or-throw -- test cleanup boundary: both the network client and temporary worker must close on failure. + try { + await client.connect(transport); + expect(client.getProtocolEra()).toBe("legacy"); + const tools = await client.listTools(); + expect(tools.tools.map((tool) => tool.name)).toContain("execute"); + expect(transport.sessionId).toBeTruthy(); + } finally { + await client.close(); + await rollbackWorker.stop(); + } + }, 120_000); + it("delivers native elicitation on the approval-gated tool call stream", async () => { const client = new Client( { name: "native-elicitation-test", version: "1.0.0" }, diff --git a/apps/local/src/mcp.ts b/apps/local/src/mcp.ts index d87f0ce8b6..76096c4fb3 100644 --- a/apps/local/src/mcp.ts +++ b/apps/local/src/mcp.ts @@ -205,7 +205,7 @@ export const createMcpRequestHandler = ( if (!resource) return jsonError(404, -32001, "MCP resource not found"); if (!(await isLegacyMcpRequest(request))) { if (handlerConfig.modernEnabled === false) { - return jsonError(503, -32022, "MCP 2026-07-28 support is disabled"); + return jsonError(400, -32022, "MCP 2026-07-28 support is disabled"); } return Effect.runPromise(modern.dispatch(request, localPrincipal, resource)); } 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 8072df7e54..8cb52a5216 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 @@ -4,12 +4,14 @@ import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import type { Transport } from "@modelcontextprotocol/sdk/shared/transport.js"; import type { JSONRPCMessage, MessageExtraInfo } from "@modelcontextprotocol/sdk/types.js"; -import { defaultMcpResource } from "@executor-js/host-mcp"; +import { defaultMcpResource, type Principal } from "@executor-js/host-mcp"; +import type { ModernMcpDispatcher } from "@executor-js/host-mcp/modern-tool-server"; import type { ExecutionEngine, ExecutionResult, ResumeResponse } from "@executor-js/execution"; import { McpAgentSessionDOBase, type McpApprovalOwner, + type McpSessionInit, type McpSessionModelResumeResult, type SessionMeta, } from "./agent-session-durable-object"; @@ -80,21 +82,34 @@ class MemoryStorage { } type HarnessSession = { + activeModernRequestCount: number; alarm: () => Promise; + beginModernRequest: () => Promise<() => void>; + closeRuntime: () => Effect.Effect; ctx: MemoryStorage; dbHandle: { readonly end: () => void } | null; engine: ExecutionEngine | null; getConnections?: () => Iterable; getSessionId: () => string; + handleModernRequest: ( + request: Request, + principal: Principal, + token: McpSessionInit, + ) => Promise; initialized: boolean; lastActivityMs: number; maxPausedSessionIdleMs: () => number; + modernDispatcher: ModernMcpDispatcher | null; + modernDispatcherPromise: Promise | null; + modernRequestDrainWaiters: Set<() => void>; onStart: () => Promise; pendingApprovalLeases: Map; props: Record; runMcpAgentOnStart: () => Promise; + runningExecutionCount: () => Promise; + runtimeClosePromise: Promise | null; server?: McpServer; - sessionMeta: SessionMeta; + sessionMeta: SessionMeta | null; sessionTimeoutMs: () => number; resumeExecutionForModel: ( executionId: string, @@ -143,6 +158,34 @@ const makeDeferred = (): { readonly promise: Promise; readonly resolve: () return { promise, resolve }; }; +class GatedModernStorage extends MemoryStorage { + readonly readEntered = makeDeferred(); + readonly releaseRead = makeDeferred(); + modernKeyReads = 0; + + override async get(key: string): Promise { + if (key === "modern-request-state-key") { + this.modernKeyReads += 1; + this.readEntered.resolve(); + await this.releaseRead.promise; + } + return super.get(key); + } +} + +class GatedSessionMetaStorage extends MemoryStorage { + readonly readEntered = makeDeferred(); + readonly releaseRead = makeDeferred(); + + override async get(key: string): Promise { + if (key === "session-meta") { + this.readEntered.resolve(); + await this.releaseRead.promise; + } + return super.get(key); + } +} + type ResumeCall = { readonly executionId: string; readonly response: ResumeResponse; @@ -181,7 +224,9 @@ const approval = { content: { approved: true }, } satisfies ResumeResponse; -const makeHarnessSession = async (): Promise => { +const makeHarnessSession = async ( + storage: MemoryStorage = new MemoryStorage(), +): Promise => { const sessionId = "session-reconnect"; const sessionMeta: SessionMeta = { organizationId: "org-1", @@ -189,11 +234,11 @@ const makeHarnessSession = async (): Promise => { userId: "user-1", resource: defaultMcpResource, }; - const storage = new MemoryStorage(); const server = makeServer(); await server.connect(new StaleCloseTransport()); const session = Object.create(McpAgentSessionDOBase.prototype) as HarnessSession; + session.activeModernRequestCount = 0; session.ctx = storage; session.dbHandle = { end: () => undefined }; session.engine = makeEngine().engine; @@ -201,6 +246,9 @@ const makeHarnessSession = async (): Promise => { session.initialized = true; session.lastActivityMs = Date.now() - 10; session.maxPausedSessionIdleMs = () => 1_000; + session.modernDispatcher = null; + session.modernDispatcherPromise = null; + session.modernRequestDrainWaiters = new Set(); session.pendingApprovalLeases = new Map(); session.props = {}; session.server = server; @@ -213,10 +261,60 @@ const makeHarnessSession = async (): Promise => { session.engine = makeEngine().engine; session.initialized = true; }; + session.runtimeClosePromise = null; return session; }; +describe("McpAgentSessionDOBase modern dispatcher initialization", () => { + type ModernDispatcherHarness = { + ctx: MemoryStorage; + modernDispatcher: ModernMcpDispatcher | null; + modernDispatcherPromise: Promise | null; + runtimeClosePromise: Promise | null; + ensureModernDispatcher: ( + principal: Principal, + token: McpSessionInit, + ) => Promise; + }; + + it("single-flights concurrent first requests onto one dispatcher", async () => { + const storage = new GatedModernStorage(); + const session = Object.create(McpAgentSessionDOBase.prototype) as ModernDispatcherHarness; + session.ctx = storage; + session.modernDispatcher = null; + session.modernDispatcherPromise = null; + session.runtimeClosePromise = null; + const principal: Principal = { + accountId: "user-1", + organizationId: "org-1", + organizationName: "Org 1", + email: "user-1@example.com", + name: "User 1", + avatarUrl: null, + roles: [], + }; + const token: McpSessionInit = { + organizationId: "org-1", + userId: "user-1", + elicitationMode: "native", + resource: defaultMcpResource, + }; + + const first = session.ensureModernDispatcher(principal, token); + await storage.readEntered.promise; + const second = session.ensureModernDispatcher(principal, token); + expect(storage.modernKeyReads).toBe(1); + + storage.releaseRead.resolve(); + const [firstDispatcher, secondDispatcher] = await Promise.all([first, second]); + expect(firstDispatcher).toBe(secondDispatcher); + expect(session.modernDispatcher).toBe(firstDispatcher); + expect(session.modernDispatcherPromise).toBeNull(); + await firstDispatcher.close(); + }); +}); + // The negotiated MCP-Apps capability arrives once, at `initialize`, and lives // in the rebuilt server's memory. These pin the storage round-trip that lets a // cold-restored session rebuild with it instead of silently downgrading every @@ -308,6 +406,91 @@ describe("McpAgentSessionDOBase apps capability persistence", () => { }); describe("McpAgentSessionDOBase transport restore", () => { + it("re-checks a modern request lease before acting on an idle alarm snapshot", async () => { + const session = await makeHarnessSession(); + session.lastActivityMs = Date.now() - 2_000; + session.sessionTimeoutMs = () => 1_000; + const originalEngine = session.engine; + const countSnapshotEntered = makeDeferred(); + const releaseCountSnapshot = makeDeferred(); + session.runningExecutionCount = async () => { + countSnapshotEntered.resolve(); + await releaseCountSnapshot.promise; + return 0; + }; + + const alarm = session.alarm(); + await countSnapshotEntered.promise; + const releaseRequest = await session.beginModernRequest(); + releaseCountSnapshot.resolve(); + + await alarm; + expect(session.initialized).toBe(true); + expect(session.engine).toBe(originalEngine); + expect(session.ctx.alarm).toBeGreaterThan(Date.now()); + releaseRequest(); + }); + + it("lets a leased full modern request finish when shutdown starts before dispatcher init", async () => { + const storage = new GatedSessionMetaStorage(); + const session = await makeHarnessSession(storage); + const sessionMeta = session.sessionMeta; + expect(sessionMeta).not.toBeNull(); + await storage.put("session-meta", sessionMeta); + session.sessionMeta = null; + const principal: Principal = { + accountId: "user-1", + organizationId: "org-1", + organizationName: "Org 1", + email: "user-1@example.com", + name: "User 1", + avatarUrl: null, + roles: [], + }; + const token: McpSessionInit = { + organizationId: "org-1", + userId: "user-1", + elicitationMode: "native", + resource: defaultMcpResource, + }; + const request = session.handleModernRequest( + new Request("https://executor.test/mcp", { + method: "POST", + headers: { + "content-type": "application/json", + "mcp-method": "server/discover", + "mcp-name": "server", + "mcp-protocol-version": "2026-07-28", + }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: 1, + method: "server/discover", + params: { + _meta: { + "io.modelcontextprotocol/protocolVersion": "2026-07-28", + "io.modelcontextprotocol/clientInfo": { name: "test", version: "1" }, + "io.modelcontextprotocol/clientCapabilities": {}, + }, + }, + }), + }), + principal, + token, + ); + await storage.readEntered.promise; + const close = Effect.runPromise(session.closeRuntime()); + await Promise.resolve(); + expect(session.runtimeClosePromise).not.toBeNull(); + + storage.releaseRead.resolve(); + const [response] = await Promise.all([request, close]); + expect(response.status).toBe(200); + expect(session.runtimeClosePromise).toBeNull(); + expect(session.initialized).toBe(false); + expect(session.modernDispatcher).toBeNull(); + }); + it("preserves hibernated response streams when a cold isolate starts", async () => { const session = await makeHarnessSession(); let closeCalls = 0; 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 f10b45bae6..3050c5ee60 100644 --- a/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts +++ b/packages/hosts/cloudflare/src/mcp/agent-session-durable-object.ts @@ -248,6 +248,10 @@ export abstract class McpAgentSessionDOBase< private approvalWaiters = new Map>(); private pendingApprovalLeases = new Map(); private modernDispatcher: ModernMcpDispatcher | null = null; + private modernDispatcherPromise: Promise | null = null; + private activeModernRequestCount = 0; + private modernRequestDrainWaiters = new Set<() => void>(); + private runtimeClosePromise: Promise | null = null; protected abstract openSessionDb(): TDbHandle | Promise; @@ -441,7 +445,7 @@ export abstract class McpAgentSessionDOBase< for (const requestIds of rows.values()) { if (Array.isArray(requestIds)) count += requestIds.length; } - return count; + return count + this.activeModernRequestCount; } private closeActiveStreams(): void { @@ -632,14 +636,42 @@ export abstract class McpAgentSessionDOBase< : built; } - private closeRuntime(options: { readonly closeStreams?: boolean } = {}): Effect.Effect { + private waitForModernRequestsToDrain(): Promise { + if (this.activeModernRequestCount === 0) return Promise.resolve(); + return new Promise((resolve) => this.modernRequestDrainWaiters.add(resolve)); + } + + private async beginModernRequest(): Promise<() => void> { + while (this.runtimeClosePromise) await this.runtimeClosePromise; + this.activeModernRequestCount += 1; + let active = true; + return () => { + if (!active) return; + active = false; + this.activeModernRequestCount -= 1; + if (this.activeModernRequestCount !== 0) return; + const waiters = [...this.modernRequestDrainWaiters]; + this.modernRequestDrainWaiters.clear(); + for (const resolve of waiters) resolve(); + }; + } + + private closeRuntimeNow(options: { readonly closeStreams?: boolean } = {}): Effect.Effect { const self = this; return Effect.gen(function* () { + yield* Effect.promise(() => self.waitForModernRequestsToDrain()); yield* self.releaseAllPendingApprovalLeases(); - if (self.modernDispatcher) { - const modern = self.modernDispatcher; - self.modernDispatcher = null; - yield* Effect.promise(() => modern.close()).pipe(Effect.ignore); + const modern = self.modernDispatcher; + const pendingModern = self.modernDispatcherPromise; + let initializedModern = modern; + if (pendingModern) { + const completed = yield* Effect.exit(Effect.promise(() => pendingModern)); + if (Exit.isSuccess(completed)) initializedModern ??= completed.value; + } + self.modernDispatcher = null; + self.modernDispatcherPromise = null; + if (initializedModern) { + yield* Effect.promise(() => initializedModern.close()).pipe(Effect.ignore); } if (options.closeStreams ?? true) { yield* Effect.sync(() => self.closeActiveStreams()); @@ -660,6 +692,32 @@ export abstract class McpAgentSessionDOBase< }); } + private startRuntimeClose(options: { readonly closeStreams?: boolean }): Promise { + if (this.runtimeClosePromise) return this.runtimeClosePromise; + const self = this; + let resolveClose: () => void = () => undefined; + let rejectClose: (reason?: unknown) => void = () => undefined; + const closing = new Promise((resolve, reject) => { + resolveClose = resolve; + rejectClose = reject; + }); + self.runtimeClosePromise = closing; + void Effect.runPromise(self.closeRuntimeNow(options)).then(resolveClose, rejectClose); + void closing.then( + () => { + if (self.runtimeClosePromise === closing) self.runtimeClosePromise = null; + }, + () => { + if (self.runtimeClosePromise === closing) self.runtimeClosePromise = null; + }, + ); + return closing; + } + + private closeRuntime(options: { readonly closeStreams?: boolean } = {}): Effect.Effect { + return Effect.promise(() => this.startRuntimeClose(options)); + } + private ensureRuntimeForApproval(): Effect.Effect { const self = this; return Effect.gen(function* () { @@ -719,6 +777,65 @@ export abstract class McpAgentSessionDOBase< }); } + /** Single-flight dispatcher construction so concurrent first requests share + * the same continuation maps and the runtime build de-duplication they own. */ + private async ensureModernDispatcher( + principal: Principal, + token: McpSessionInit, + ): Promise { + if (this.modernDispatcher) return this.modernDispatcher; + if (this.modernDispatcherPromise) return this.modernDispatcherPromise; + + const self = this; + const starting = (async () => { + let key = await self.ctx.storage.get(MODERN_REQUEST_STATE_KEY); + if (!key) { + key = crypto.getRandomValues(new Uint8Array(32)); + await self.ctx.storage.put(MODERN_REQUEST_STATE_KEY, key); + } + return makeModernMcpDispatcher( + () => + self.ensureModernRuntime(principal, token).pipe( + Effect.flatMap((runtimeAccess) => { + if (runtimeAccess === "forbidden" || !self.engine) { + return Effect.fail(new ModernMcpRuntimeUnavailable()); + } + return self.engine.getDescription.pipe( + Effect.map((description) => ({ engine: self.engine!, description })), + ); + }), + ), + { + requestStateKey: key, + onExecutionPaused: (executionId) => + Effect.runPromise( + self.startPendingApprovalLease(executionId, { + ttlMs: PAUSED_APPROVAL_TIMEOUT_MS, + expiresAt: new Date(Date.now() + PAUSED_APPROVAL_TIMEOUT_MS).toISOString(), + }), + ), + onResumeStarted: (executionId) => + Effect.runPromise(self.beginPendingApprovalResume(executionId)), + onResumeSettled: (executionId) => + Effect.runPromise(self.finishPendingApprovalResume(executionId)), + }, + ); + })(); + self.modernDispatcherPromise = starting; + void starting.then( + (dispatcher) => { + if (self.modernDispatcherPromise === starting) { + self.modernDispatcher = dispatcher; + self.modernDispatcherPromise = null; + } + }, + () => { + if (self.modernDispatcherPromise === starting) self.modernDispatcherPromise = null; + }, + ); + return starting; + } + /** Serve one modern stateless HTTP exchange inside an owner/resource-keyed * Durable Object. The live engine stays here across MRTR rounds; signed * requestState is only the client-carried continuation pointer. */ @@ -729,6 +846,7 @@ export abstract class McpAgentSessionDOBase< incoming?: IncomingTraceHeaders, ): Promise { const self = this; + const releaseModernRequest = await self.beginModernRequest(); const program = Effect.gen(function* () { yield* self.prepareErrorCaptureScope(); const access = yield* self.ensureModernOwner(principal, token); @@ -743,47 +861,12 @@ export abstract class McpAgentSessionDOBase< ); } - if (!self.modernDispatcher) { - let key = yield* Effect.promise(() => - self.ctx.storage.get(MODERN_REQUEST_STATE_KEY), - ); - if (!key) { - key = crypto.getRandomValues(new Uint8Array(32)); - yield* Effect.promise(() => self.ctx.storage.put(MODERN_REQUEST_STATE_KEY, key)); - } - self.modernDispatcher = makeModernMcpDispatcher( - () => - self.ensureModernRuntime(principal, token).pipe( - Effect.flatMap((runtimeAccess) => { - if (runtimeAccess === "forbidden" || !self.engine) { - return Effect.fail(new ModernMcpRuntimeUnavailable()); - } - return self.engine.getDescription.pipe( - Effect.map((description) => ({ engine: self.engine!, description })), - ); - }), - ), - { - requestStateKey: key, - onExecutionPaused: (executionId) => - Effect.runPromise( - self.startPendingApprovalLease(executionId, { - ttlMs: PAUSED_APPROVAL_TIMEOUT_MS, - expiresAt: new Date(Date.now() + PAUSED_APPROVAL_TIMEOUT_MS).toISOString(), - }), - ), - onResumeStarted: (executionId) => - Effect.runPromise(self.beginPendingApprovalResume(executionId)), - onResumeSettled: (executionId) => - Effect.runPromise(self.finishPendingApprovalResume(executionId)), - }, - ); - } - return yield* self.modernDispatcher.dispatch(request, principal, token.resource); + const dispatcher = yield* Effect.promise(() => self.ensureModernDispatcher(principal, token)); + return yield* dispatcher.dispatch(request, principal, token.resource); }).pipe(Effect.withSpan("McpSessionDO.handleModernRequest"), (effect) => self.withSpanFlush(effect), ); - return Effect.runPromise(self.withTelemetry(program, incoming)); + return Effect.runPromise(self.withTelemetry(program, incoming)).finally(releaseModernRequest); } private startRuntimeFromOnStart(props?: McpSessionProps): Effect.Effect { @@ -1115,6 +1198,14 @@ export abstract class McpAgentSessionDOBase< return; } + // The alarm snapshot spans durable-storage awaits. A modern request may + // start or finish `markActivity` after an earlier count/read, so re-check + // the in-memory lease and activity immediately before destructive cleanup. + if (this.activeModernRequestCount > 0 || this.lastActivityMs > lastActivityMs) { + await this.ctx.storage.setAlarm(Date.now() + this.sessionTimeoutMs()); + return; + } + await this.disposeIdleRuntime({ idleMs, pausedExecutionCount }); } diff --git a/packages/hosts/mcp/src/envelope.test.ts b/packages/hosts/mcp/src/envelope.test.ts index 8b22d1e3d1..b49f1d27f1 100644 --- a/packages/hosts/mcp/src/envelope.test.ts +++ b/packages/hosts/mcp/src/envelope.test.ts @@ -207,6 +207,7 @@ describe("McpServingRoutes envelope", () => { "mcp-protocol-version": "2026-07-28", "mcp-method": "server/discover", "mcp-name": "server", + origin: "https://claude.ai", }, body: JSON.stringify({ jsonrpc: "2.0", @@ -228,20 +229,28 @@ describe("McpServingRoutes envelope", () => { expect(await Effect.runPromise(Ref.get(calls))).toEqual(["none"]); }); - it("rejects mismatched Host and cross-origin browser requests before dispatch", async () => { + it("rejects a mismatched Host before dispatch", async () => { const handler = buildHandler(OkStoreLive, McpErrorReporterNoop); const badHost = await handler( new Request("https://host.test/mcp", { method: "GET", headers: { host: "evil.test" } }), ); expect(badHost.status).toBe(421); + }); - const badOrigin = await handler( + it("permits cross-origin legacy browser requests to reach dispatch", async () => { + const handler = buildHandler(OkStoreLive, McpErrorReporterNoop); + const response = await handler( new Request("https://host.test/mcp", { - method: "GET", - headers: { origin: "https://evil.test" }, + method: "POST", + headers: { + authorization: "Bearer x", + "content-type": "application/json", + origin: "https://claude.ai", + }, + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "tools/list" }), }), ); - expect(badOrigin.status).toBe(403); + expect(response.status).toBe(200); }); }); diff --git a/packages/hosts/mcp/src/envelope.ts b/packages/hosts/mcp/src/envelope.ts index 43e8e11009..14d41ebf05 100644 --- a/packages/hosts/mcp/src/envelope.ts +++ b/packages/hosts/mcp/src/envelope.ts @@ -158,19 +158,15 @@ const withMcpCors = (response: Response): Response => { }); }; -/** Reject DNS-rebinding and cross-origin browser requests before credentials - * are evaluated. Non-browser clients normally send no Origin and are accepted. */ +/** Reject DNS-rebinding attempts before credentials are evaluated. Origin is + * intentionally not restricted here: this envelope advertises wildcard CORS, + * and hosts with a narrower browser policy enforce it at their own boundary. */ export const validateMcpRequestAuthority = (request: Request): Response | null => { const url = new URL(request.url); const host = request.headers.get("host"); if (host && host.toLowerCase() !== url.host.toLowerCase()) { return jsonRpcResponse(421, -32600, "Host header does not match request URL"); } - const origin = request.headers.get("origin"); - if (!origin) return null; - if (!URL.canParse(origin) || new URL(origin).origin !== url.origin) { - return jsonRpcResponse(403, -32600, "Origin is not allowed"); - } return null; };