From 11295eb72eb3cbe08ed774b79343af64b17b497e Mon Sep 17 00:00:00 2001 From: Rhys Sullivan Date: Fri, 14 Aug 2026 06:44:31 -0700 Subject: [PATCH] Migrate outbound MCP client to SDK v2 (spec 2026-07-28) --- bun.lock | 12 +++ packages/plugins/mcp/package.json | 2 + .../plugins/mcp/src/sdk/catalog-sync.test.ts | 8 ++ .../mcp/src/sdk/connection-pool.test.ts | 16 ++-- .../plugins/mcp/src/sdk/connection-pool.ts | 4 + packages/plugins/mcp/src/sdk/connection.ts | 30 +++--- .../plugins/mcp/src/sdk/elicitation.test.ts | 6 +- .../plugins/mcp/src/sdk/http-status.test.ts | 20 +++- packages/plugins/mcp/src/sdk/http-status.ts | 45 +++++---- packages/plugins/mcp/src/sdk/invoke.test.ts | 15 ++- packages/plugins/mcp/src/sdk/invoke.ts | 39 +++++--- packages/plugins/mcp/src/sdk/plugin.test.ts | 8 ++ packages/plugins/mcp/src/sdk/plugin.ts | 11 ++- .../plugins/mcp/src/sdk/probe-shape.test.ts | 95 ++++++++++++++++++- packages/plugins/mcp/src/sdk/probe-shape.ts | 61 ++++++++++-- .../plugins/mcp/src/sdk/stdio-connector.ts | 7 +- .../mcp/src/sdk/testing-fixtures.test.ts | 12 ++- 17 files changed, 304 insertions(+), 87 deletions(-) diff --git a/bun.lock b/bun.lock index c004c4bcc2..43c3519f48 100644 --- a/bun.lock +++ b/bun.lock @@ -994,6 +994,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,6 +2128,10 @@ "@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=="], @@ -6118,6 +6124,12 @@ "@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=="], "@octokit/request/content-type": ["content-type@2.0.0", "", {}, "sha512-j/O/d7GcZCyNl7/hwZAb606rzqkyvaDctLmckbxLzHvFBzTJHuGEdodATcP3yIRoDrLHkIATJuvzbFlp/ki2cQ=="], 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/catalog-sync.test.ts b/packages/plugins/mcp/src/sdk/catalog-sync.test.ts index 5948485714..853719ca04 100644 --- a/packages/plugins/mcp/src/sdk/catalog-sync.test.ts +++ b/packages/plugins/mcp/src/sdk/catalog-sync.test.ts @@ -182,6 +182,13 @@ const decodeJsonRpcRequest = Schema.decodeUnknownOption(Schema.fromJsonString(Js const jsonRpcResult = (request: JsonRpcRequest, result: unknown) => HttpServerResponse.jsonUnsafe({ jsonrpc: "2.0", id: request.id ?? null, result }); +const jsonRpcMethodNotFound = (request: JsonRpcRequest) => + HttpServerResponse.jsonUnsafe({ + jsonrpc: "2.0", + id: request.id ?? null, + error: { code: -32601, message: "Method not found" }, + }); + const pageTool = (name: string) => ({ name, description: `Tool ${name}`, @@ -199,6 +206,7 @@ const servePaginatedListServer = () => return Option.match(decodeJsonRpcRequest(body), { onNone: () => HttpServerResponse.text("Invalid JSON-RPC fixture request", { status: 400 }), onSome: (rpc) => { + if (rpc.method === "server/discover") return jsonRpcMethodNotFound(rpc); if (rpc.method === "initialize") { return jsonRpcResult(rpc, { protocolVersion: "2025-06-18", diff --git a/packages/plugins/mcp/src/sdk/connection-pool.test.ts b/packages/plugins/mcp/src/sdk/connection-pool.test.ts index 5efb84b083..dfb366276d 100644 --- a/packages/plugins/mcp/src/sdk/connection-pool.test.ts +++ b/packages/plugins/mcp/src/sdk/connection-pool.test.ts @@ -7,6 +7,10 @@ import { createMcpConnectionPool } from "./connection-pool"; import { invokeMcpTool } from "./invoke"; import { makeEchoMcpServer, serveMcpServer } from "../testing"; +// The public v1 fixture creates one throwaway transport for v2's rejected +// server/discover probe, then one initialized legacy session per real dial. +const V1_FIXTURE_TRANSPORTS_PER_DIAL = 2; + const acceptAll: Elicit = () => Effect.succeed(ElicitationResponse.make({ action: "accept", content: { approved: true } })); @@ -54,7 +58,7 @@ describe("MCP connection pool", () => { expect(first).toMatchObject({ content: [{ type: "text", text: "first" }] }); expect(second).toMatchObject({ content: [{ type: "text", text: "second" }] }); - expect(server.sessionCount()).toBe(1); + expect(server.sessionCount()).toBe(V1_FIXTURE_TRANSPORTS_PER_DIAL); yield* pool.close(); }), ), @@ -88,7 +92,7 @@ describe("MCP connection pool", () => { expect.objectContaining({ content: [{ type: "text", text: "left" }] }), expect.objectContaining({ content: [{ type: "text", text: "right" }] }), ]); - expect(server.sessionCount()).toBe(2); + expect(server.sessionCount()).toBe(2 * V1_FIXTURE_TRANSPORTS_PER_DIAL); yield* pool.close(); }), ), @@ -115,7 +119,7 @@ describe("MCP connection pool", () => { }); expect(after).toMatchObject({ content: [{ type: "text", text: "after" }] }); - expect(server.sessionCount()).toBe(2); + expect(server.sessionCount()).toBe(2 * V1_FIXTURE_TRANSPORTS_PER_DIAL); yield* pool.close(); }), ), @@ -148,7 +152,7 @@ describe("MCP connection pool", () => { }); expect(after).toMatchObject({ content: [{ type: "text", text: "after" }] }); - expect(server.sessionCount()).toBe(2); + expect(server.sessionCount()).toBe(2 * V1_FIXTURE_TRANSPORTS_PER_DIAL); yield* pool.close(); }), ), @@ -188,8 +192,8 @@ describe("MCP connection pool", () => { }); expect(after).toMatchObject({ content: [{ type: "text", text: "after" }] }); - // A second session was dialled rather than the dead one being reused. - expect(server.sessionCount()).toBe(2); + // A second connection was dialled rather than the dead one being reused. + expect(server.sessionCount()).toBe(2 * V1_FIXTURE_TRANSPORTS_PER_DIAL); yield* pool.close(); }), ), diff --git a/packages/plugins/mcp/src/sdk/connection-pool.ts b/packages/plugins/mcp/src/sdk/connection-pool.ts index 6b732caf78..bf25c866a2 100644 --- a/packages/plugins/mcp/src/sdk/connection-pool.ts +++ b/packages/plugins/mcp/src/sdk/connection-pool.ts @@ -3,6 +3,10 @@ import { Cause, Effect, Exit, Predicate } from "effect"; import type { McpConnection, McpConnector } from "./connection"; import type { McpInvocationError } from "./errors"; +// The pool preserves sessions for sessionful legacy servers. Stateless +// 2026-07-28 servers do not need it, but retaining a cheap idle client is +// harmless and keeps one lifecycle for both protocol eras. + const IDLE_TTL_MS = 5 * 60 * 1_000; type IdleConnection = { diff --git a/packages/plugins/mcp/src/sdk/connection.ts b/packages/plugins/mcp/src/sdk/connection.ts index 82a98cd0d1..22fb1647cd 100644 --- a/packages/plugins/mcp/src/sdk/connection.ts +++ b/packages/plugins/mcp/src/sdk/connection.ts @@ -1,16 +1,18 @@ -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 { + Client, + SSEClientTransport, + StreamableHTTPClientTransport, + type FetchLike, + type OAuthClientProvider, +} from "@modelcontextprotocol/client"; +import { CfWorkerJsonSchemaValidator } from "@modelcontextprotocol/client/validators/cf-worker"; 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` -// at evaluation time, which crashes workerd (incl. vitest-pool-workers) at -// SIGSEGV on module instantiation. Cloud callers set +// (`@modelcontextprotocol/client/stdio`) still imports Node process/stream and +// `cross-spawn` eagerly at evaluation time, which crashes workerd (including +// vitest-pool-workers) with SIGSEGV on module instantiation. Cloud callers set // `dangerouslyAllowStdioMCP: false` and never reach the stdio branch below; // prod bundles that DO use stdio load it via a dynamic import inside the // stdio branch of `createMcpConnector`. @@ -201,12 +203,13 @@ 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 => +const createClient = (versionNegotiation?: { readonly mode: "auto" }): Client => new Client( { name: "executor-mcp", version: "0.1.0" }, { capabilities: { elicitation: { form: {}, url: {} } }, jsonSchemaValidator: new CfWorkerJsonSchemaValidator(), + ...(versionNegotiation === undefined ? {} : { versionNegotiation }), }, ); @@ -247,9 +250,10 @@ const connectionFailure = ( const connectClient = (input: { transport: string; createTransport: () => Parameters[0]; + versionNegotiation?: { readonly mode: "auto" }; }): Effect.Effect => Effect.gen(function* () { - const client = createClient(); + const client = createClient(input.versionNegotiation); const transportInstance = input.createTransport(); yield* Effect.tryPromise({ @@ -314,8 +318,12 @@ export const createMcpConnector = (input: ConnectorInput): McpConnector => { const endpoint = buildEndpointUrl(input.endpoint, input.queryParams ?? {}); + // Auto-negotiate the 2026-07-28 era only on Streamable HTTP. SSE is a + // legacy-only transport, and stdio servers are spawned per call where the + // SDK recommends retaining its legacy-default handshake. const connectStreamableHttp = connectClient({ transport: "streamable-http", + versionNegotiation: { mode: "auto" }, createTransport: () => new StreamableHTTPClientTransport(endpoint, { requestInit, diff --git a/packages/plugins/mcp/src/sdk/elicitation.test.ts b/packages/plugins/mcp/src/sdk/elicitation.test.ts index 4ca64f32a2..cc2bed354f 100644 --- a/packages/plugins/mcp/src/sdk/elicitation.test.ts +++ b/packages/plugins/mcp/src/sdk/elicitation.test.ts @@ -1,7 +1,7 @@ import { describe, expect, it } from "@effect/vitest"; import { Effect, Predicate, Schema, Semaphore } from "effect"; -import { CfWorkerJsonSchemaValidator } from "@modelcontextprotocol/sdk/validation/cfworker"; -import type { JsonSchemaType } from "@modelcontextprotocol/sdk/validation/types"; +import type { JsonSchemaType } from "@modelcontextprotocol/client"; +import { CfWorkerJsonSchemaValidator } from "@modelcontextprotocol/client/validators/cf-worker"; import { AuthTemplateSlug, @@ -226,7 +226,7 @@ describe("MCP elicitation (end-to-end)", () => { ]), ); expect(schema?.outputTypeScript).toContain('type: "text"'); - expect(schema?.outputTypeScript).toContain("structuredContent?: { [k: string]: unknown; }"); + expect(schema?.outputTypeScript).toContain("structuredContent?: unknown;"); const result = yield* executor.execute( simpleEcho.address, diff --git a/packages/plugins/mcp/src/sdk/http-status.test.ts b/packages/plugins/mcp/src/sdk/http-status.test.ts index 2276509f7e..0344a71986 100644 --- a/packages/plugins/mcp/src/sdk/http-status.test.ts +++ b/packages/plugins/mcp/src/sdk/http-status.test.ts @@ -1,4 +1,5 @@ import { describe, expect, it } from "@effect/vitest"; +import { InsufficientScopeError, SdkErrorCode, SdkHttpError } from "@modelcontextprotocol/client"; // oxlint-disable executor/no-error-constructor -- boundary: these tests reproduce the MCP SDK's own transport rejections, which are built-in Errors import { insufficientScopeFromCause } from "./http-status"; @@ -9,7 +10,8 @@ import { insufficientScopeFromCause } from "./http-status"; // - with an authProvider (the production OAuth path): the StreamableHTTP // transport consumes the insufficient_scope challenge itself, retries // with the broader scope, and only when THAT fails throws the fixed -// "Server returned 403 after trying upscoping" message. +// typed `InsufficientScopeError`, or after retry exhaustion the fixed +// `SdkHttpError` step-up message. describe("insufficientScopeFromCause", () => { it("detects the OAuth error body embedded in a transport message", () => { expect( @@ -31,9 +33,21 @@ describe("insufficientScopeFromCause", () => { ).toBe(true); }); - it("detects the SDK's exhausted-upscoping failure (the authProvider path)", () => { + it("detects the SDK's typed insufficient-scope failure", () => { expect( - insufficientScopeFromCause(new Error("Server returned 403 after trying upscoping")), + insufficientScopeFromCause(new InsufficientScopeError({ requiredScope: "files.read" })), + ).toBe(true); + }); + + it("detects the SDK's exhausted step-up failure (the authProvider path)", () => { + expect( + insufficientScopeFromCause( + new SdkHttpError( + SdkErrorCode.ClientHttpForbidden, + "Server returned 403 insufficient_scope after step-up re-authorization (retry limit 2 reached)", + { status: 403 }, + ), + ), ).toBe(true); }); diff --git a/packages/plugins/mcp/src/sdk/http-status.ts b/packages/plugins/mcp/src/sdk/http-status.ts index a4442f8d23..541c631f34 100644 --- a/packages/plugins/mcp/src/sdk/http-status.ts +++ b/packages/plugins/mcp/src/sdk/http-status.ts @@ -1,7 +1,8 @@ // --------------------------------------------------------------------------- // 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)`. +// transport failures two ways: an `SdkHttpError` carrying a numeric `status`, +// and an `SseError` carrying a numeric `code`. The SSE transport also retains +// its historic POST-failure message for errors created below EventSource. // 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,13 +10,13 @@ import { Option, Schema } from "effect"; import { insufficientScopeFromEmbeddedJson } from "@executor-js/sdk/core"; -import { StreamableHTTPError } from "@modelcontextprotocol/sdk/client/streamableHttp.js"; +import { InsufficientScopeError, SdkHttpError, SseError } from "@modelcontextprotocol/client"; const SsePostErrorCause = Schema.Struct({ message: Schema.String }); const decodeSsePostErrorCause = Schema.decodeUnknownOption(SsePostErrorCause); -// Matches the SDK's SSEClientTransport POST-failure message (sse.js); re-verify -// on SDK bumps. A format drift just yields undefined (generic error, no crash). +// V2 still constructs this exact message in SSEClientTransport._send. A format +// drift just yields undefined (generic error, no crash). const statusFromSsePostError = (cause: unknown): number | undefined => Option.match(decodeSsePostErrorCause(cause), { onNone: () => undefined, @@ -26,32 +27,36 @@ 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; +const statusFromTypedTransportError = (cause: unknown): number | undefined => { + if (SdkHttpError.isInstance(cause)) return cause.status; + if (SseError.isInstance(cause)) { + const code = cause.code; + return code !== undefined && code >= 100 && code <= 599 ? code : undefined; + } + return undefined; }; export const httpStatusFromCause = (cause: unknown): number | undefined => - statusFromStreamableHttpError(cause) ?? statusFromSsePostError(cause); + statusFromTypedTransportError(cause) ?? statusFromSsePostError(cause); // The SDK embeds the upstream response text in the transport error message // ("Error POSTing to endpoint: "), which is the only place a 403's body // survives for connections without an authProvider. For OAuth connections the -// 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 -// operation, and re-running the identical flow cannot help. Strict matching +// StreamableHTTP transport consumes the insufficient_scope challenge itself. +// V2 throws `InsufficientScopeError` when configured not to reauthorize; after +// exhausting its step-up retries it throws `SdkHttpError` with the exact fixed +// message matched below (verified against the installed v2 transport source). +// 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. -const SDK_UPSCOPING_EXHAUSTED_RE = /Server returned 403 after trying upscoping/; +// exact step-up message) — a miss stays on the generic auth path. +const SDK_STEP_UP_EXHAUSTED_RE = + /^Server returned 403 insufficient_scope after step-up re-authorization \(retry limit \d+ reached\)$/; export const insufficientScopeFromCause = (cause: unknown): boolean => + InsufficientScopeError.isInstance(cause) || Option.match(decodeSsePostErrorCause(cause), { onNone: () => false, onSome: ({ message }) => - insufficientScopeFromEmbeddedJson(message) || SDK_UPSCOPING_EXHAUSTED_RE.test(message), + insufficientScopeFromEmbeddedJson(message) || SDK_STEP_UP_EXHAUSTED_RE.test(message), }); diff --git a/packages/plugins/mcp/src/sdk/invoke.test.ts b/packages/plugins/mcp/src/sdk/invoke.test.ts index 3b5aaae66f..ce133bfd5f 100644 --- a/packages/plugins/mcp/src/sdk/invoke.test.ts +++ b/packages/plugins/mcp/src/sdk/invoke.test.ts @@ -2,9 +2,12 @@ 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, + SdkErrorCode, + SdkHttpError, + type OAuthClientProvider, +} from "@modelcontextprotocol/client"; import { ElicitationResponse } from "@executor-js/sdk"; import { serveTestHttpApp } from "@executor-js/sdk/testing"; @@ -108,14 +111,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, "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..46a08bc663 100644 --- a/packages/plugins/mcp/src/sdk/invoke.ts +++ b/packages/plugins/mcp/src/sdk/invoke.ts @@ -8,7 +8,7 @@ // legitimately retain state in that session. The pool keeps one idle // connection per resolved identity and leases it exclusively per invoke; // stdio and callers without a pool remain strictly per-call. -// 2. Installing a per-invocation `ElicitRequestSchema` handler that bridges +// 2. Installing a per-invocation `elicitation/create` handler that bridges // MCP's elicit capability into the host's elicit function threaded via // `InvokeToolInput.elicit`. // 3. Calling `client.callTool({ name, arguments })`. @@ -16,12 +16,7 @@ import { Cause, Effect, Exit, Option, Predicate, Schema } from "effect"; -import { - ElicitRequestSchema, - ErrorCode, - McpError, - ToolListChangedNotificationSchema, -} from "@modelcontextprotocol/sdk/types.js"; +import { ProtocolError, ProtocolErrorCode } from "@modelcontextprotocol/client"; import { ElicitationId, @@ -66,9 +61,10 @@ 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) && - // 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 + 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 ProtocolError, whose message carries the only unknown-tool discriminator the protocol provides isUnknownToolMessage(cause.message, toolName); // --------------------------------------------------------------------------- @@ -93,6 +89,17 @@ const McpElicitParams = Schema.Union([ type McpElicitParams = typeof McpElicitParams.Type; const decodeElicitParams = Schema.decodeUnknownSync(McpElicitParams); +const decodeElicitContent = Schema.decodeUnknownSync( + Schema.Record( + Schema.String, + Schema.Union([ + Schema.String, + Schema.Number, + Schema.Boolean, + Schema.mutable(Schema.Array(Schema.String)), + ]), + ), +); const toElicitationRequest = (params: McpElicitParams): ElicitationRequest => params.mode === "url" @@ -107,7 +114,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 +126,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: decodeElicitContent(response.content) } + : {}), }; } const failure = exit.cause.reasons.find(Cause.isFailReason); @@ -149,7 +158,7 @@ const installToolListChangedHandler = ( onToolListChanged: (() => void) | undefined, ): void => { if (!onToolListChanged) return; - client.setNotificationHandler(ToolListChangedNotificationSchema, () => { + client.setNotificationHandler("notifications/tools/list_changed", () => { onToolListChanged(); }); }; @@ -189,8 +198,8 @@ 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; + // oxlint-disable-next-line executor/no-instanceof-tagged-error -- boundary: MCP SDK protocol failures are its ProtocolError subclass; transport failures use other error shapes + 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.test.ts b/packages/plugins/mcp/src/sdk/plugin.test.ts index 0b8338684b..6ff1415cab 100644 --- a/packages/plugins/mcp/src/sdk/plugin.test.ts +++ b/packages/plugins/mcp/src/sdk/plugin.test.ts @@ -55,12 +55,20 @@ const jsonRpcResult = (request: JsonRpcRequest, result: unknown) => result, }); +const jsonRpcMethodNotFound = (request: JsonRpcRequest) => + HttpServerResponse.jsonUnsafe({ + jsonrpc: "2.0", + id: request.id ?? null, + error: { code: -32601, message: "Method not found" }, + }); + // The call-tool fixtures share one JSON-RPC scaffold (handshake, tool listing, // unknown-method rejection); only the `tools/call` response varies. Each // scenario supplies that branch via a `CallToolResponder`. type CallToolResponder = (rpc: JsonRpcRequest) => ReturnType; const callToolFixtureResponse = (rpc: JsonRpcRequest, callTool: CallToolResponder) => { + if (rpc.method === "server/discover") return jsonRpcMethodNotFound(rpc); if (rpc.method === "initialize") { return jsonRpcResult(rpc, { protocolVersion: "2025-06-18", diff --git a/packages/plugins/mcp/src/sdk/plugin.ts b/packages/plugins/mcp/src/sdk/plugin.ts index e3b7a6857e..70ec0e69e6 100644 --- a/packages/plugins/mcp/src/sdk/plugin.ts +++ b/packages/plugins/mcp/src/sdk/plugin.ts @@ -1,9 +1,8 @@ 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 type { OAuthClientProvider } from "@modelcontextprotocol/client"; +import { CallToolResultSchema } from "@modelcontextprotocol/core"; import { authToolFailure, @@ -392,7 +391,7 @@ type JsonSchemaObject = Record & { readonly properties?: Record; }; -const McpCallToolResultJsonSchema = z.toJSONSchema(CallToolResultSchema) as JsonSchemaObject; +const McpCallToolResultJsonSchema: JsonSchemaObject = CallToolResultSchema.toJSONSchema(); const mcpCallToolResultOutputSchema = (structuredContentSchema?: unknown): JsonSchemaObject => { const defaultStructuredContentSchema = @@ -493,7 +492,9 @@ export const userFacingProbeMessage = ( // MCP-SDK OAuth provider adapter — wraps a pre-resolved access token so the // transport sends it as a Bearer header. Refresh is core's responsibility // (the connection row carries the OAuth grant); this adapter never initiates -// a new flow and fails loudly if the SDK tries to. +// a new flow and fails loudly if the SDK tries to. V2 stamps stored credentials +// with the authorization-server issuer and offers scoped invalidation; this +// single-token boundary intentionally persists neither. // --------------------------------------------------------------------------- const makeOAuthProvider = (accessToken: string): OAuthClientProvider => ({ diff --git a/packages/plugins/mcp/src/sdk/probe-shape.test.ts b/packages/plugins/mcp/src/sdk/probe-shape.test.ts index 7eb124ebaf..8051125386 100644 --- a/packages/plugins/mcp/src/sdk/probe-shape.test.ts +++ b/packages/plugins/mcp/src/sdk/probe-shape.test.ts @@ -1,6 +1,12 @@ import { describe, expect, it } from "@effect/vitest"; -import { Effect, Ref } from "effect"; -import { HttpServerResponse } from "effect/unstable/http"; +import { Effect, Layer, Predicate, Ref } from "effect"; +import { + HttpClient, + HttpClientError, + HttpClientRequest, + HttpClientResponse, + HttpServerResponse, +} from "effect/unstable/http"; import { serveTestHttpApp } from "@executor-js/sdk/testing"; import { probeMcpEndpointShape } from "./probe-shape"; @@ -320,6 +326,91 @@ describe("probeMcpEndpointShape", () => { ), ); + it.effect("falls through a wrong-shape legacy GET retry to modern server discovery", () => + Effect.scoped( + Effect.gen(function* () { + const server = yield* serveProbeEndpoint((request) => { + if (request.body.includes('"method":"server/discover"')) { + return HttpServerResponse.jsonUnsafe({ + jsonrpc: "2.0", + id: 2, + error: { code: -32601, message: "Method not found" }, + }); + } + if (request.method === "GET") { + return HttpServerResponse.jsonUnsafe({ error: "legacy SSE is unsupported" }); + } + return HttpServerResponse.empty({ status: 405 }); + }); + + const result = yield* probeMcpEndpointShape(server.endpoint); + expect(result).toEqual({ kind: "mcp", requiresAuth: false }); + + const requests = yield* server.requests; + expect(requests).toHaveLength(3); + expect(requests[0]?.body).toContain('"protocolVersion":"2025-11-25"'); + expect(requests[1]?.method).toBe("GET"); + expect(requests[2]?.body).toBe( + JSON.stringify({ + jsonrpc: "2.0", + id: 2, + method: "server/discover", + params: { + _meta: { "io.modelcontextprotocol/protocolVersion": "2026-07-28" }, + }, + }), + ); + expect(requests[2]?.headers["mcp-protocol-version"]).toBe("2026-07-28"); + }), + ), + ); + + it.effect("keeps the initialize verdict when the discover fallback fails at the transport", () => + Effect.gen(function* () { + const bodyText = (request: HttpClientRequest.HttpClientRequest): string => + Predicate.isTagged(request.body, "Uint8Array") + ? new TextDecoder().decode(request.body.body) + : ""; + const httpClientLayer = Layer.succeed(HttpClient.HttpClient)( + HttpClient.make( + ( + request: HttpClientRequest.HttpClientRequest, + ): Effect.Effect< + HttpClientResponse.HttpClientResponse, + HttpClientError.HttpClientError + > => + bodyText(request).includes('"method":"server/discover"') + ? Effect.fail( + new HttpClientError.HttpClientError({ + reason: new HttpClientError.TransportError({ + request, + description: "connection reset by peer", + }), + }), + ) + : Effect.succeed( + HttpClientResponse.fromWeb( + request, + new Response("not mcp", { + status: 200, + headers: { "content-type": "text/html" }, + }), + ), + ), + ), + ); + + const result = yield* probeMcpEndpointShape("https://internal.example/mcp", { + httpClientLayer, + }); + expect(result).toEqual({ + kind: "not-mcp", + category: "wrong-shape", + reason: "2xx POST body is not a JSON-RPC envelope", + }); + }), + ); + it.effect("rejects 2xx with HTML body as wrong-shape", () => withServer( () => diff --git a/packages/plugins/mcp/src/sdk/probe-shape.ts b/packages/plugins/mcp/src/sdk/probe-shape.ts index 3b3be94297..1049fb9f22 100644 --- a/packages/plugins/mcp/src/sdk/probe-shape.ts +++ b/packages/plugins/mcp/src/sdk/probe-shape.ts @@ -13,7 +13,7 @@ // and (b) plenty of real MCP servers authenticate with static API // keys and publish no OAuth metadata at all (e.g. cubic.dev). // -// The probe issues an unauth JSON-RPC `initialize` POST and accepts +// The primary probe issues an unauth JSON-RPC `initialize` POST and accepts // only the wire shapes a real MCP server can return: // // - 2xx with `Content-Type: text/event-stream` — streamable HTTP @@ -30,8 +30,14 @@ // only accepts 2xx with `text/event-stream` or the same 401+Bearer // shape. // -// One `fetch` (occasionally two), no MCP-SDK session state, no OAuth -// round-trip, no DCR — every non-MCP endpoint exits here. +// If initialize ultimately has the wrong shape, a second JSON-RPC POST probes +// `server/discover` using the 2026-07-28 envelope and protocol header. This +// catches modern-only servers that reject initialize with a non-JSON-RPC +// response. Authentication outcomes remain terminal because they do not vary +// by transport era. +// +// One primary request (occasionally plus legacy GET and modern discover), no +// MCP-SDK session state, no OAuth round-trip, no DCR. // --------------------------------------------------------------------------- import { Data, Duration, Effect, Layer, Option, Schema } from "effect"; @@ -48,12 +54,21 @@ const INITIALIZE_BODY = JSON.stringify({ id: 1, method: "initialize", params: { - protocolVersion: "2025-06-18", + protocolVersion: "2025-11-25", capabilities: {}, clientInfo: { name: "executor-probe", version: "0" }, }, }); +const DISCOVER_BODY = JSON.stringify({ + jsonrpc: "2.0", + id: 2, + method: "server/discover", + params: { + _meta: { "io.modelcontextprotocol/protocolVersion": "2026-07-28" }, + }, +}); + /** Header-name lookup is case-insensitive per RFC 7230. `fetch`'s * `Response.headers` already lower-cases, but we normalise explicitly * to stay robust against test mocks that construct `Headers` loosely. */ @@ -385,10 +400,9 @@ export const probeMcpEndpointShape = ( .execute(postRequest) .pipe(Effect.timeout(Duration.millis(timeoutMs))); - const postResult = yield* classify(postResponse, "POST"); - if (postResult) return postResult; + let initializeResult = yield* classify(postResponse, "POST"); - if ([404, 405, 406, 415].includes(postResponse.status)) { + if (initializeResult === null && [404, 405, 406, 415].includes(postResponse.status)) { let getRequest = HttpClientRequest.get(url.toString()).pipe( HttpClientRequest.setHeader("accept", "text/event-stream"), ); @@ -398,15 +412,42 @@ export const probeMcpEndpointShape = ( const getResponse = yield* client .execute(getRequest) .pipe(Effect.timeout(Duration.millis(timeoutMs))); - const getResult = yield* classify(getResponse, "GET"); - if (getResult) return getResult; + initializeResult = yield* classify(getResponse, "GET"); } - return { + initializeResult ??= { kind: "not-mcp", category: "wrong-shape", reason: `unexpected status ${postResponse.status} for initialize`, } as const; + + if (initializeResult.kind !== "not-mcp" || initializeResult.category !== "wrong-shape") { + return initializeResult; + } + + 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 = HttpClientRequest.setHeader( + discoverRequest, + "MCP-Protocol-Version", + "2026-07-28", + ); + + // The endpoint already answered the primary probe, so a transport + // failure on this secondary request must not overwrite that verdict + // with "unreachable" — keep the initialize classification instead. + const discoverResult = yield* client.execute(discoverRequest).pipe( + Effect.timeout(Duration.millis(timeoutMs)), + Effect.flatMap((discoverResponse) => classify(discoverResponse, "POST")), + Effect.catch(() => Effect.succeed(null)), + ); + return discoverResult ?? initializeResult; }).pipe( Effect.provide(options.httpClientLayer ?? FetchHttpClient.layer), Effect.mapError( diff --git a/packages/plugins/mcp/src/sdk/stdio-connector.ts b/packages/plugins/mcp/src/sdk/stdio-connector.ts index 99a0f72e37..6fec6f0617 100644 --- a/packages/plugins/mcp/src/sdk/stdio-connector.ts +++ b/packages/plugins/mcp/src/sdk/stdio-connector.ts @@ -3,8 +3,9 @@ // --------------------------------------------------------------------------- // // Kept in its own module so `connection.ts` never imports it eagerly at -// module load. `@modelcontextprotocol/sdk/client/stdio.js` pulls in -// `node:child_process` at evaluation time; under `@cloudflare/vitest-pool-workers` +// module load. The v2 `@modelcontextprotocol/client/stdio` entry still eagerly +// evaluates Node-only process/stream imports and `cross-spawn` (which loads +// `node:child_process`); under `@cloudflare/vitest-pool-workers` // that crashes workerd at module instantiation with SIGSEGV (prod bundles // tree-shake it away when `dangerouslyAllowStdioMCP: false`, tests do not). // @@ -13,7 +14,7 @@ // the import and therefore never touch `node:child_process`. // --------------------------------------------------------------------------- -import { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js"; +import { StdioClientTransport } from "@modelcontextprotocol/client/stdio"; export type StdioTransportConfig = { readonly command: string; diff --git a/packages/plugins/mcp/src/sdk/testing-fixtures.test.ts b/packages/plugins/mcp/src/sdk/testing-fixtures.test.ts index b804335399..f29ef707a1 100644 --- a/packages/plugins/mcp/src/sdk/testing-fixtures.test.ts +++ b/packages/plugins/mcp/src/sdk/testing-fixtures.test.ts @@ -1,7 +1,6 @@ import { expect, layer } from "@effect/vitest"; import { Effect } from "effect"; -import { Client } from "@modelcontextprotocol/sdk/client/index.js"; -import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js"; +import { Client, StreamableHTTPClientTransport } from "@modelcontextprotocol/client"; import { OAuthTestServer } from "@executor-js/sdk/testing"; import { makeEchoMcpServer, serveMcpServerWithOAuth } from "../testing"; @@ -16,7 +15,10 @@ const createGreetingMcpServer = () => }); const makeClient = (endpoint: string, accessToken: string) => { - const client = new Client({ name: "executor-test-client", version: "1.0.0" }); + const client = new Client( + { name: "executor-test-client", version: "1.0.0" }, + { versionNegotiation: { mode: "auto" } }, + ); const transport = new StreamableHTTPClientTransport(new URL(endpoint), { requestInit: { headers: { authorization: `Bearer ${accessToken}` }, @@ -47,7 +49,9 @@ layer(OAuthTestServer.layer(), { timeout: "15 seconds" })("MCP testing fixtures" expect(result).toMatchObject({ content: [{ type: "text", text: "Hello Ada" }], }); - expect(server.sessionCount()).toBe(1); + // The v1 fixture creates a throwaway transport for the rejected modern + // discovery probe before the initialized legacy session. + expect(server.sessionCount()).toBe(2); const requests = yield* server.requests; expect(