diff --git a/docs/contributing/architecture/authentication.md b/docs/contributing/architecture/authentication.md index ad9fba30a8..41118cf291 100644 --- a/docs/contributing/architecture/authentication.md +++ b/docs/contributing/architecture/authentication.md @@ -470,7 +470,17 @@ routed from `packages/worker/src/index.ts`. - Authorization endpoint: `/oauth/authorize` - Token endpoint: `/oauth/token` (via provider) -- Client registration: `/oauth/register` (via provider) +- Client registration: `/oauth/register` (via provider), plus Client ID Metadata + Documents (`clientIdMetadataDocumentEnabled` in + `packages/worker/src/index.ts`): a client may present an HTTPS URL as its + `client_id` with no registration step. MCP `2026-07-28` deprecates RFC 7591 + dynamic registration in favor of CIMD, so both stay enabled: clients that do + not use CIMD register via `/oauth/register`, and a failed CIMD metadata fetch + returns `invalid_client` (any DCR retry after that is the client's own + recovery, not a server-side fallback). CIMD metadata fetches rely on the + `global_fetch_strictly_public` compatibility flag in + `packages/worker/wrangler.jsonc` for SSRF safety; the provider only advertises + `client_id_metadata_document_supported` when both are set. - Supported scopes: `profile`, `email` - On `/oauth/authorize`, unauthenticated users can log in inline or via top-nav auth links; those links preserve the full authorize URL in `redirectTo` so diff --git a/docs/contributing/architecture/data-storage.md b/docs/contributing/architecture/data-storage.md index 020e4e3122..c9427d39d1 100644 --- a/docs/contributing/architecture/data-storage.md +++ b/docs/contributing/architecture/data-storage.md @@ -939,6 +939,10 @@ Bindings are configured per environment in `packages/worker/wrangler.jsonc` [Usage metering](./usage-metering.md)) - `EMAIL_EVENTS` (Analytics Engine dataset, production/preview only; indexed by stable user id and read only through role-gated platform aggregates) +- `MCP_PROTOCOL_EVENTS` (Analytics Engine dataset, production/preview only; one + point per authenticated `/mcp` request recording which protocol lane served it + — legacy sessionful vs stateless 2026-07-28 — for legacy-lane retirement; see + `packages/worker/src/mcp/protocol-metrics.ts`) `packages/worker/wrangler.jsonc` also configures the `EMAIL` send binding, dispatch queues, worker loaders (`LOADER` / `APP_LOADER`), the `AI` binding, and diff --git a/docs/contributing/architecture/request-lifecycle.md b/docs/contributing/architecture/request-lifecycle.md index 7d24d74c92..f41f25cee1 100644 --- a/docs/contributing/architecture/request-lifecycle.md +++ b/docs/contributing/architecture/request-lifecycle.md @@ -51,7 +51,16 @@ Requests are handled in this order: the `/mcp` suffix path only): - `/.well-known/oauth-protected-resource/mcp` 5. MCP endpoint: - - `/mcp` (requires OAuth bearer token) + - `/mcp` (requires OAuth bearer token). After authentication, + `packages/worker/src/mcp-auth.ts` routes by protocol era: 2025-era requests + go to the sessionful `MCP` Durable Object (`McpAgent`, MCP SDK v1), and + `2026-07-28` envelope requests are served statelessly per request by + `packages/worker/src/mcp/stateless-lane.ts` (MCP SDK v2, no Durable + Object). Both lanes share one tool registration; every authenticated + request records a lane data point to the `MCP_PROTOCOL_EVENTS` Analytics + Engine dataset so the legacy lane can be retired once its traffic stops + (see + [decision 0005](../decisions/0005-mcp-dual-lane-stateless-migration.md)). 6. Public `@username` ingress handled in `packages/worker/src/index.ts` before the OAuth provider / app router (needs `ExecutionContext` for background work): diff --git a/docs/contributing/decisions/0005-mcp-dual-lane-stateless-migration.md b/docs/contributing/decisions/0005-mcp-dual-lane-stateless-migration.md new file mode 100644 index 0000000000..c1c743195f --- /dev/null +++ b/docs/contributing/decisions/0005-mcp-dual-lane-stateless-migration.md @@ -0,0 +1,44 @@ +# 0005: MCP dual-lane serving with metrics-driven legacy retirement + +- **Status:** accepted +- **Date:** 2026-08-05 + +## Context + +MCP protocol revision `2026-07-28` made the protocol stateless: the `initialize` +handshake and `Mcp-Session-Id` header are removed and every request carries its +own `_meta` envelope. The Cloudflare Agents SDK deprecated and feature-froze +`McpAgent`, which hosts kody's `/mcp` as a sessionful Durable Object on MCP SDK +v1. Nearly all installed MCP clients still speak 2025-era revisions, so dropping +the sessionful path outright would break real traffic, while staying on +`McpAgent` alone pins kody to a frozen stack. + +## Decision + +Serve `/mcp` as two lanes behind one route and one shared tool registration +(`packages/worker/src/mcp/register-tools.ts`): 2025-era requests keep the +`McpAgent` Durable Object lane unchanged, and `2026-07-28` envelope requests are +served by a per-request stateless SDK v2 server +(`packages/worker/src/mcp/stateless-lane.ts`). Routing uses the SDK's own +`isLegacyRequest` predicate, and every authenticated request records a lane data +point to the `MCP_PROTOCOL_EVENTS` Analytics Engine dataset — deliberately not +the primary D1 database (aggregate-only readout, off the hot path; consistent +with [0002](./0002-data-placement.md)). The legacy lane is removed when the +metrics show its traffic has gone (a sustained window of zero or negligible +legacy-lane requests from real clients), not on a calendar date. + +## Consequences + +Modern clients get stateless serving with no MCP session Durable Object on the +request path (the account write lease taken at the auth boundary is unchanged +and applies to both lanes — it is the deletion-safety guard every kody surface +takes, not MCP session state), while every existing client keeps byte-identical +behavior. Tool definitions cannot drift between lanes, but the two SDK +generations meet at a typed seam (`asMcpToolServer` in +`packages/worker/src/mcp/mcp-registration-agent.ts`) that a future SDK bump must +revisit. Retiring the legacy lane later also deletes the `mcp_agent_sessions` +registry, the `MCP_OBJECT` Durable Object, and the session purge path. Tasks +(the `io.modelcontextprotocol/tasks` extension) are deliberately not implemented +yet: the SDK v2 ships the vocabulary without a runtime, no major client supports +it, and execute's idempotency-key + `run_get` flow already covers the need; +revisit when a major host ships task support. diff --git a/docs/contributing/decisions/index.md b/docs/contributing/decisions/index.md index a1a0316281..1352453a6f 100644 --- a/docs/contributing/decisions/index.md +++ b/docs/contributing/decisions/index.md @@ -23,3 +23,4 @@ behavior (see [documentation principles](../documentation.md)). - [0002 — Data placement: D1, per-user Durable Objects, Analytics Engine](./0002-data-placement.md) - [0003 — Repos are the base primitive; packages are an explicit extension](./0003-repos-as-base-primitive.md) - [0004 — Status page as a separate worker with its own storage](./0004-status-page-separate-worker.md) +- [0005 — MCP dual-lane serving with metrics-driven legacy retirement](./0005-mcp-dual-lane-stateless-migration.md) diff --git a/packages/worker/src/env-schema.ts b/packages/worker/src/env-schema.ts index cd89924a2a..6234b1574d 100644 --- a/packages/worker/src/env-schema.ts +++ b/packages/worker/src/env-schema.ts @@ -205,6 +205,7 @@ export const EnvSchema = object({ EMAIL_EVENTS: optionalAnalyticsEngineDatasetSchema, USAGE_EVENTS: optionalAnalyticsEngineDatasetSchema, FLAG_EXPOSURES: optionalAnalyticsEngineDatasetSchema, + MCP_PROTOCOL_EVENTS: optionalAnalyticsEngineDatasetSchema, SENTRY_DSN: optionalUrlStringSchema, SENTRY_ENVIRONMENT: optionalNonEmptyStringSchema, SENTRY_TRACES_SAMPLE_RATE: optionalSentryTracesSampleRateSchema, diff --git a/packages/worker/src/index.ts b/packages/worker/src/index.ts index 6eda72c433..1b036c9764 100644 --- a/packages/worker/src/index.ts +++ b/packages/worker/src/index.ts @@ -425,6 +425,17 @@ const oauthProvider = new OAuthProvider({ tokenEndpoint: oauthPaths.token, clientRegistrationEndpoint: oauthPaths.register, scopesSupported: oauthScopes, + // Client ID Metadata Documents (MCP 2025-11-25 SEP-991): clients may use + // an HTTPS URL as their client_id instead of registering via DCR. The + // 2026-07-28 revision deprecates RFC 7591 DCR in favor of CIMD, so both + // stay enabled: CIMD clients present their URL client_id with no + // registration step, and clients that do not use CIMD register via + // /oauth/register. A failed CIMD metadata fetch returns invalid_client; + // whether a client then registers via DCR is the client's own recovery. + // Requires the global_fetch_strictly_public compatibility flag (set in + // wrangler.jsonc) so metadata fetches are SSRF-safe; the provider only + // advertises CIMD support when both are on. + clientIdMetadataDocumentEnabled: true, // Provider default onError logs every structured OAuth error via console.warn. // Keep those responses on the wire without duplicating them into worker logs / // test console guards; unexpected throws still reach our fetch catch + Sentry. diff --git a/packages/worker/src/mcp-auth.ts b/packages/worker/src/mcp-auth.ts index 04fd5fc2a2..ef7f67b5d9 100644 --- a/packages/worker/src/mcp-auth.ts +++ b/packages/worker/src/mcp-auth.ts @@ -11,6 +11,11 @@ import { } from './mcp/auth-audit.ts' import { withAccountWriteLease } from '#worker/account/deletion-state.ts' import { createMcpCallerContext, type McpServerProps } from './mcp/context.ts' +import { + classifyMcpProtocolRequest, + recordMcpProtocolEvent, +} from './mcp/protocol-metrics.ts' +import { handleStatelessMcpRequest } from './mcp/stateless-lane.ts' import { oauthScopes } from './oauth-handlers.ts' export const mcpResourcePath = '/mcp' @@ -250,16 +255,41 @@ export async function handleMcpRequest({ }) context.props = props + // Lane classification: 2025-era ("legacy") requests keep the sessionful + // Durable Object McpAgent lane; 2026-07-28 envelope requests are served + // by the stateless SDK v2 lane. Every authenticated request records a + // lane data point so legacy-lane retirement is a metrics decision — see + // ./mcp/protocol-metrics.ts for the readout query. + const classification = await classifyMcpProtocolRequest(request) + recordMcpProtocolEvent(env, { + lane: classification.lane, + method: classification.method, + protocolVersion: classification.protocolVersion, + clientName: classification.clientName, + clientVersion: classification.clientVersion, + userId: mcpUser.userId, + }) + return await withAccountWriteLease({ db: env.APP_DB, stableUserId: mcpUser.userId, holder: `mcp:${request.method} ${url.pathname}`, env, write: async () => - await fetchMcp( - request, - env, - context as ExecutionContext, - ), + classification.lane === 'legacy' + ? await fetchMcp( + request, + env, + context as ExecutionContext, + ) + : await handleStatelessMcpRequest({ + request, + env, + ctx, + callerContext: props, + ...(classification.parsedBody === undefined + ? {} + : { parsedBody: classification.parsedBody }), + }), }) } diff --git a/packages/worker/src/mcp-auth.workers.test.ts b/packages/worker/src/mcp-auth.workers.test.ts index 2577682ee8..1597ed87ce 100644 --- a/packages/worker/src/mcp-auth.workers.test.ts +++ b/packages/worker/src/mcp-auth.workers.test.ts @@ -620,6 +620,145 @@ test('mcp request enforces token audience and forwards caller props', async () = ) }, 15_000) +test('mcp requests route by protocol era and record lane metrics', async () => { + const origin = 'https://example.com' + const validToken: TokenSummary = { + id: 'token', + grantId: 'grant', + userId: 'user', + createdAt: 0, + expiresAt: 999999, + audience: `${origin}${mcpResourcePath}`, + grant: { + clientId: 'client', + scope: oauthScopes, + props: { userId: 'user', email: 'user@example.com' }, + }, + } + const dataPoints: Array = [] + const env = createEnv( + createHelpers({ unwrapToken: async () => validToken }), + { + MCP_PROTOCOL_EVENTS: { + writeDataPoint: (point: AnalyticsEngineDataPoint) => { + dataPoints.push(point) + }, + } as AnalyticsEngineDataset, + }, + { emailVerifiedAt: new Date(0).toISOString() }, + ) + let legacyLaneCalls = 0 + const fetchMcp = () => { + legacyLaneCalls += 1 + return new Response('legacy-lane') + } + + // 2025-era handshake stays on the sessionful Durable Object lane. + const legacyResponse = await handleMcpRequestAndDrain({ + request: new Request(`${origin}${mcpResourcePath}`, { + method: 'POST', + headers: { + Authorization: 'Bearer token', + 'Content-Type': 'application/json', + Accept: 'application/json, text/event-stream', + }, + body: JSON.stringify({ + jsonrpc: '2.0', + id: 1, + method: 'initialize', + params: { + protocolVersion: '2025-06-18', + capabilities: {}, + clientInfo: { name: 'legacy-client', version: '1.0.0' }, + }, + }), + }), + env, + ctx: createContext(), + fetchMcp, + }) + expect(await legacyResponse.text()).toBe('legacy-lane') + expect(legacyLaneCalls).toBe(1) + expect(dataPoints).toHaveLength(1) + expect(dataPoints[0]).toEqual({ + indexes: ['legacy'], + blobs: [ + 'legacy', + 'initialize', + '2025-06-18', + 'legacy-client', + '1.0.0', + 'user', + ], + doubles: [1], + }) + + // 2026-07-28 envelope requests are served by the stateless lane and + // never reach the Durable Object; the advertised tools carry the shared + // definitions including output schemas and icons. + const modernResponse = await handleMcpRequestAndDrain({ + request: new Request(`${origin}${mcpResourcePath}`, { + method: 'POST', + headers: { + Authorization: 'Bearer token', + 'Content-Type': 'application/json', + Accept: 'application/json, text/event-stream', + 'MCP-Protocol-Version': '2026-07-28', + 'Mcp-Method': 'tools/list', + }, + body: JSON.stringify({ + jsonrpc: '2.0', + id: 2, + method: 'tools/list', + params: { + _meta: { + 'io.modelcontextprotocol/protocolVersion': '2026-07-28', + 'io.modelcontextprotocol/clientCapabilities': {}, + 'io.modelcontextprotocol/clientInfo': { + name: 'modern-client', + version: '2.0.0', + }, + }, + }, + }), + }), + env, + ctx: createContext(), + fetchMcp, + }) + expect(legacyLaneCalls).toBe(1) + expect(modernResponse.status).toBe(200) + const modernBody = (await modernResponse.json()) as { + result: { + resultType?: string + tools: Array<{ + name: string + outputSchema?: Record + icons?: Array<{ src: string }> + }> + } + } + const toolNames = modernBody.result.tools.map((tool) => tool.name).sort() + expect(toolNames).toEqual(['execute', 'search']) + for (const tool of modernBody.result.tools) { + expect(tool.outputSchema).toMatchObject({ type: 'object' }) + expect(tool.icons?.[0]?.src).toBe(`${origin}/android-chrome-192x192.png`) + } + expect(dataPoints).toHaveLength(2) + expect(dataPoints[1]).toEqual({ + indexes: ['modern'], + blobs: [ + 'modern', + 'tools/list', + '2026-07-28', + 'modern-client', + '2.0.0', + 'user', + ], + doubles: [1], + }) +}) + test('mcp request rejects unverified and unidentifiable accounts fail-closed', async () => { const request = new Request(`https://example.com${mcpResourcePath}`, { headers: { Authorization: 'Bearer token' }, diff --git a/packages/worker/src/mcp/index.ts b/packages/worker/src/mcp/index.ts index ff22dfc455..6a860a88bc 100644 --- a/packages/worker/src/mcp/index.ts +++ b/packages/worker/src/mcp/index.ts @@ -8,6 +8,10 @@ import { buildSentryOptions } from '../sentry-options.ts' import { parseMcpCallerContext, type McpServerProps } from './context.ts' import { buildMcpServerInstructions } from './server-instructions.ts' import { registerTools } from './register-tools.ts' +import { + asMcpToolServer, + type McpRegistrationAgent, +} from './mcp-registration-agent.ts' import { createKodyMcpServer } from './sentry-mcp-server.ts' import { getMcpUserServerInstructions } from './user-server-instructions-repo.ts' import { getCapabilityRegistryForContext } from './capabilities/registry.ts' @@ -63,7 +67,32 @@ class MCPBase extends McpAgent { }), jsonSchemaValidator: new CfWorkerJsonSchemaValidator(), }) - await registerTools(this) + await registerTools(this.getRegistrationAgent()) + } + /** + * Registration surface shared with the stateless lane (see + * `asMcpToolServer` for the SDK v1/v2 seam). `state`/`setState` are + * forwarded live so tool runners keep their per-session behavior + * (search preamble dedup, raw-fetch host nudges) on this lane. + */ + getRegistrationAgent() { + const self = this + const agent: McpRegistrationAgent & { + state?: State + setState?: (state: State) => void + } = { + server: asMcpToolServer(this.server), + getEnv: () => self.getEnv(), + getCallerContext: () => self.getCallerContext(), + requireDomain: () => self.requireDomain(), + getLoopbackExports: () => self.getLoopbackExports(), + waitUntil: (promise) => self.waitUntil(promise), + get state() { + return self.state + }, + setState: (state) => self.setState(state), + } + return agent } getCallerContext() { return parseMcpCallerContext(this.props) diff --git a/packages/worker/src/mcp/mcp-registration-agent.ts b/packages/worker/src/mcp/mcp-registration-agent.ts index d47fa94bbd..cda7b84c3e 100644 --- a/packages/worker/src/mcp/mcp-registration-agent.ts +++ b/packages/worker/src/mcp/mcp-registration-agent.ts @@ -1,8 +1,65 @@ -import { type McpServer } from '@modelcontextprotocol/sdk/server/mcp.js' +import { type McpServer as McpServerV1 } from '@modelcontextprotocol/sdk/server/mcp.js' +import { type McpServer as McpServerV2 } from '@modelcontextprotocol/server' +import { + type CallToolResult, + type ToolAnnotations, +} from '@modelcontextprotocol/sdk/types.js' +import { type z } from 'zod' import { type McpCallerContext } from '@kody-internal/shared/chat.ts' +/** + * Tool icon metadata (protocol revision 2025-11-25, SEP-973). Advertised by + * the stateless SDK v2 lane; the SDK v1 `McpAgent` lane has no icon support + * and ignores the config key at runtime (its `registerTool` destructures only + * the keys it knows), so one shared config serves both lanes. + */ +export type McpToolIcon = { + src: string + mimeType?: string + sizes?: Array +} + +export type McpToolConfig = { + title?: string + description?: string + inputSchema?: Record + outputSchema?: Record + annotations?: ToolAnnotations + icons?: Array +} + +/** + * The registration surface kody's tools use, satisfied by both MCP server + * generations: `@modelcontextprotocol/sdk` v1 (`McpAgent` Durable Object + * lane, 2025-era protocol revisions) and `@modelcontextprotocol/server` v2 + * (stateless lane, 2026-07-28). Both accept the raw-shape config form with + * a callback returning content + structured content, so one registration + * definition backs both lanes and they can never drift apart. + * + * The callback parameter is typed `never` so any concrete argument type is + * accepted; call sites annotate their args explicitly and the concrete + * server validates inputs against `inputSchema` at runtime. + */ +export type McpToolServer = { + registerTool( + name: string, + config: McpToolConfig, + cb: (args: never) => CallToolResult | Promise, + ): unknown +} + +/** + * Bridge one concrete server generation onto the shared registration + * surface. Runtime-compatible by construction (both generations implement + * the raw-shape `registerTool` form); the cast only bridges the two SDKs' + * incompatible generic signatures, which TypeScript cannot relate directly. + */ +export function asMcpToolServer(server: McpServerV1 | McpServerV2) { + return server as unknown as McpToolServer +} + export type McpRegistrationAgent = { - server: McpServer + server: McpToolServer getEnv(): Env getCallerContext(): McpCallerContext requireDomain(): string diff --git a/packages/worker/src/mcp/protocol-metrics.node.test.ts b/packages/worker/src/mcp/protocol-metrics.node.test.ts new file mode 100644 index 0000000000..28ae183f8d --- /dev/null +++ b/packages/worker/src/mcp/protocol-metrics.node.test.ts @@ -0,0 +1,207 @@ +import { expect, test, vi } from 'vitest' +import { + classifyMcpProtocolRequest, + recordMcpProtocolEvent, + type McpProtocolEventEnv, +} from './protocol-metrics.ts' + +const mcpUrl = 'https://example.com/mcp' + +function jsonRequest(body: unknown, headers: Record = {}) { + return new Request(mcpUrl, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + Accept: 'application/json, text/event-stream', + ...headers, + }, + body: JSON.stringify(body), + }) +} + +test('legacy initialize handshake classifies legacy with client info', async () => { + const request = jsonRequest({ + jsonrpc: '2.0', + id: 1, + method: 'initialize', + params: { + protocolVersion: '2025-06-18', + capabilities: {}, + clientInfo: { name: 'claude-ai', version: '0.1.0' }, + }, + }) + const classification = await classifyMcpProtocolRequest(request) + expect(classification).toMatchObject({ + lane: 'legacy', + method: 'initialize', + protocolVersion: '2025-06-18', + clientName: 'claude-ai', + clientVersion: '0.1.0', + }) + // The request body stays readable for the lane that serves it. + expect(await request.text()).toContain('initialize') +}) + +test('legacy post-handshake call reports the negotiated header version', async () => { + const classification = await classifyMcpProtocolRequest( + jsonRequest( + { + jsonrpc: '2.0', + id: 2, + method: 'tools/call', + params: { name: 'search', arguments: { query: 'email' } }, + }, + { 'mcp-protocol-version': '2025-03-26' }, + ), + ) + expect(classification).toMatchObject({ + lane: 'legacy', + method: 'tools/call', + protocolVersion: '2025-03-26', + clientName: '', + clientVersion: '', + }) +}) + +test('modern envelope request classifies modern', async () => { + const classification = await classifyMcpProtocolRequest( + jsonRequest( + { + jsonrpc: '2.0', + id: 3, + method: 'tools/call', + params: { + name: 'search', + arguments: { query: 'email' }, + _meta: { + 'io.modelcontextprotocol/protocolVersion': '2026-07-28', + 'io.modelcontextprotocol/clientCapabilities': {}, + 'io.modelcontextprotocol/clientInfo': { + name: 'modern-client', + version: '2.0.0', + }, + }, + }, + }, + { + 'MCP-Protocol-Version': '2026-07-28', + 'Mcp-Method': 'tools/call', + 'Mcp-Name': 'search', + }, + ), + ) + expect(classification).toMatchObject({ + lane: 'modern', + method: 'tools/call', + protocolVersion: '2026-07-28', + clientName: 'modern-client', + clientVersion: '2.0.0', + }) + expect(classification.parsedBody).toMatchObject({ method: 'tools/call' }) +}) + +test('body-less session operations classify legacy http methods', async () => { + const get = await classifyMcpProtocolRequest( + new Request(mcpUrl, { + headers: { + Accept: 'text/event-stream', + 'mcp-protocol-version': '2025-06-18', + }, + }), + ) + expect(get).toMatchObject({ + lane: 'legacy', + method: 'http:GET', + protocolVersion: '2025-06-18', + }) + + const del = await classifyMcpProtocolRequest( + new Request(mcpUrl, { method: 'DELETE' }), + ) + expect(del).toMatchObject({ lane: 'legacy', method: 'http:DELETE' }) +}) + +test('invalid JSON body classifies legacy without throwing', async () => { + const classification = await classifyMcpProtocolRequest( + new Request(mcpUrl, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: 'not json', + }), + ) + expect(classification.lane).toBe('legacy') + expect(classification.method).toBe('unknown') + expect(classification.parsedBody).toBeUndefined() +}) + +test('recordMcpProtocolEvent writes one lane data point', () => { + const writeDataPoint = vi.fn<(point: AnalyticsEngineDataPoint) => void>() + const env = { + MCP_PROTOCOL_EVENTS: { + writeDataPoint, + } as unknown as AnalyticsEngineDataset, + } satisfies McpProtocolEventEnv + recordMcpProtocolEvent(env, { + lane: 'legacy', + method: 'tools/call', + protocolVersion: '2025-06-18', + clientName: 'claude-ai', + clientVersion: '0.1.0', + userId: 'user-1', + }) + expect(writeDataPoint).toHaveBeenCalledExactlyOnceWith({ + indexes: ['legacy'], + blobs: [ + 'legacy', + 'tools/call', + '2025-06-18', + 'claude-ai', + '0.1.0', + 'user-1', + ], + doubles: [1], + }) +}) + +test('recordMcpProtocolEvent is a no-op without the binding and never throws', () => { + expect(() => + recordMcpProtocolEvent( + {}, + { + lane: 'modern', + method: 'tools/list', + protocolVersion: '2026-07-28', + clientName: '', + clientVersion: '', + }, + ), + ).not.toThrow() + + const consoleWarn = vi.spyOn(console, 'warn').mockImplementation(() => {}) + try { + expect(() => + recordMcpProtocolEvent( + { + MCP_PROTOCOL_EVENTS: { + writeDataPoint: () => { + throw new Error('sink offline') + }, + } as unknown as AnalyticsEngineDataset, + }, + { + lane: 'modern', + method: 'tools/list', + protocolVersion: '2026-07-28', + clientName: '', + clientVersion: '', + }, + ), + ).not.toThrow() + expect(consoleWarn).toHaveBeenCalledWith( + 'mcp-protocol-event-failed', + expect.any(Error), + ) + } finally { + consoleWarn.mockRestore() + } +}) diff --git a/packages/worker/src/mcp/protocol-metrics.ts b/packages/worker/src/mcp/protocol-metrics.ts new file mode 100644 index 0000000000..621cb7971a --- /dev/null +++ b/packages/worker/src/mcp/protocol-metrics.ts @@ -0,0 +1,168 @@ +/** + * MCP protocol lane metrics. + * + * `/mcp` serves two lanes: the sessionful Durable Object `McpAgent` lane for + * 2025-era protocol revisions ("legacy") and the stateless SDK v2 lane for + * 2026-07-28+ revisions ("modern"). The legacy lane can only be retired once + * real traffic stops using it, so every authenticated `/mcp` request records + * one data point describing which lane served it, the protocol revision, the + * JSON-RPC method, and the connecting client. + * + * Data points go to the `MCP_PROTOCOL_EVENTS` Analytics Engine dataset — + * deliberately not the primary D1 database: the retirement readout is a pure + * aggregate (no joins), and `writeDataPoint` is non-blocking so recording + * stays off the request hot path. Without the binding (local dev, tests) + * recording is a no-op. Recording never throws. + * + * Retirement readout (Analytics Engine SQL API; counts must be weighted by + * `_sample_interval` because Analytics Engine samples at high volume): + * + * ```sql + * SELECT blob1 AS lane, blob4 AS client_name, SUM(_sample_interval) AS requests + * FROM kody_mcp_protocol_events + * WHERE timestamp > NOW() - INTERVAL '30' DAY + * GROUP BY lane, client_name + * ORDER BY requests DESC + * ``` + */ + +import { + CLIENT_INFO_META_KEY, + PROTOCOL_VERSION_META_KEY, + isLegacyRequest, +} from '@modelcontextprotocol/server' + +export type McpProtocolLane = 'legacy' | 'modern' + +export type McpProtocolEventEnv = { + MCP_PROTOCOL_EVENTS?: AnalyticsEngineDataset +} + +export type McpProtocolClassification = { + lane: McpProtocolLane + /** JSON-RPC method for POSTs, `http:GET` / `http:DELETE` for session operations. */ + method: string + /** Protocol revision claimed by the request (`_meta` envelope or negotiated header). */ + protocolVersion: string + clientName: string + clientVersion: string + /** + * Parsed JSON body for POST requests, when parseable. Callers forward it + * to the stateless handler so the body is parsed at most once per request. + */ + parsedBody?: unknown +} + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +function readString(value: unknown): string { + return typeof value === 'string' ? value : '' +} + +/** First JSON-RPC message of a body (batch arrays classify by their first entry). */ +function firstJsonRpcMessage(body: unknown): Record | null { + if (Array.isArray(body)) { + const first = body.find(isRecord) + return first ?? null + } + return isRecord(body) ? body : null +} + +/** + * Classify one authenticated `/mcp` request into its serving lane and extract + * the metric fields. Uses the SDK's own {@link isLegacyRequest} predicate — + * the exact classifier `createMcpHandler` runs — so routing and metrics can + * never disagree. The request body is read from a clone; the request stays + * fully readable for whichever lane serves it. + */ +export async function classifyMcpProtocolRequest( + request: Request, +): Promise { + let parsedBody: unknown + let hasParsedBody = false + if (request.method === 'POST') { + try { + parsedBody = await request.clone().json() + hasParsedBody = true + } catch { + // Empty or invalid JSON classifies legacy; the lane handler owns + // the error response. + } + } + const legacy = hasParsedBody + ? await isLegacyRequest(request, parsedBody) + : await isLegacyRequest(request) + const lane: McpProtocolLane = legacy ? 'legacy' : 'modern' + + if (request.method !== 'POST') { + return { + lane, + method: `http:${request.method}`, + protocolVersion: readString(request.headers.get('mcp-protocol-version')), + clientName: '', + clientVersion: '', + } + } + + const message = firstJsonRpcMessage(parsedBody) + const method = readString(message?.['method']) || 'unknown' + const meta = isRecord(message?.['params']) + ? message['params']['_meta'] + : undefined + const params = isRecord(message?.['params']) ? message['params'] : undefined + + // Modern requests carry the revision in the `_meta` envelope; legacy + // requests send the negotiated `mcp-protocol-version` header (or name the + // offered revision in the `initialize` params). + const protocolVersion = + (isRecord(meta) ? readString(meta[PROTOCOL_VERSION_META_KEY]) : '') || + readString(request.headers.get('mcp-protocol-version')) || + (method === 'initialize' ? readString(params?.['protocolVersion']) : '') + + const clientInfo = isRecord(meta) + ? meta[CLIENT_INFO_META_KEY] + : method === 'initialize' + ? params?.['clientInfo'] + : undefined + + return { + lane, + method, + protocolVersion, + clientName: isRecord(clientInfo) ? readString(clientInfo['name']) : '', + clientVersion: isRecord(clientInfo) + ? readString(clientInfo['version']) + : '', + ...(hasParsedBody ? { parsedBody } : {}), + } +} + +/** + * Record one lane data point. Non-blocking, never throws, no-op without the + * Analytics Engine binding. + */ +export function recordMcpProtocolEvent( + env: McpProtocolEventEnv, + input: Omit & { userId?: string }, +): void { + try { + env.MCP_PROTOCOL_EVENTS?.writeDataPoint({ + // Indexed by lane so Analytics Engine sampling stays independent + // per lane and low-volume modern traffic is never drowned out. + indexes: [input.lane], + blobs: [ + input.lane, + input.method, + input.protocolVersion, + input.clientName, + input.clientVersion, + input.userId ?? '', + ], + doubles: [1], + }) + } catch (error) { + console.warn('mcp-protocol-event-failed', error) + } +} diff --git a/packages/worker/src/mcp/stateless-lane.mcp-e2e.test.ts b/packages/worker/src/mcp/stateless-lane.mcp-e2e.test.ts new file mode 100644 index 0000000000..16fca037b3 --- /dev/null +++ b/packages/worker/src/mcp/stateless-lane.mcp-e2e.test.ts @@ -0,0 +1,42 @@ +import { expect, test } from 'vitest' +import { + createModernMcpClient, + createTestDatabase, + startDevServer, +} from '../../../../tools/mcp-test-support.ts' + +/** + * One smoke journey for the stateless 2026-07-28 lane: a real SDK v2 client + * pinned to the modern revision negotiates over HTTP with OAuth and calls a + * tool. Everything else about lane behavior is covered by faster node and + * workers tests beside the implementation (`mcp-auth.workers.test.ts`, + * `protocol-metrics.node.test.ts`). + */ + +test('pinned 2026-07-28 client negotiates the stateless lane and calls search', async () => { + await using database = await createTestDatabase() + await using server = await startDevServer(database.persistDir) + await using modern = await createModernMcpClient( + server.origin, + database.user, + { persistDir: database.persistDir }, + ) + + const tools = await modern.client.listTools() + const toolNames = tools.tools.map((tool) => tool.name).sort() + expect(toolNames).toEqual(['execute', 'search']) + const searchTool = tools.tools.find((tool) => tool.name === 'search') + expect(searchTool?.outputSchema).toMatchObject({ type: 'object' }) + expect(searchTool?.icons?.[0]?.src).toBe( + new URL('/android-chrome-192x192.png', server.origin).toString(), + ) + + const searchResult = await modern.client.callTool({ + name: 'search', + arguments: { query: 'jobs', limit: 5 }, + }) + expect(searchResult.isError ?? false).toBe(false) + expect(searchResult.structuredContent).toMatchObject({ + conversationId: expect.any(String), + }) +}) diff --git a/packages/worker/src/mcp/stateless-lane.ts b/packages/worker/src/mcp/stateless-lane.ts new file mode 100644 index 0000000000..37d375ac3e --- /dev/null +++ b/packages/worker/src/mcp/stateless-lane.ts @@ -0,0 +1,138 @@ +/** + * Stateless `/mcp` lane for protocol revision 2026-07-28. + * + * 2025-era ("legacy") traffic keeps its sessionful Durable Object `McpAgent` + * lane unchanged; requests carrying the 2026-07-28 per-request `_meta` + * envelope are served here by a per-request SDK v2 server with no session, + * no Durable Object, and no cross-request state. Routing between the lanes + * happens in `handleMcpRequest` via the SDK's own `isLegacyRequest` + * predicate, and this handler keeps `legacy: 'reject'` so the two lanes can + * never both answer the same era. + * + * Tool definitions are shared with the legacy lane through `registerTools`, + * so the lanes cannot drift apart. Lane usage is recorded to Analytics + * Engine (see `protocol-metrics.ts`); the legacy lane is retired once its + * traffic stops. + */ + +import { type exports as workerExports } from 'cloudflare:workers' +import { invariant } from '@epic-web/invariant' +import { McpServer, createMcpHandler } from '@modelcontextprotocol/server' +import { CfWorkerJsonSchemaValidator } from '@modelcontextprotocol/server/validators/cf-worker' +import { type McpCallerContext } from '@kody-internal/shared/chat.ts' +import { buildMcpServerInstructions } from './server-instructions.ts' +import { registerTools } from './register-tools.ts' +import { + asMcpToolServer, + type McpRegistrationAgent, +} from './mcp-registration-agent.ts' +import { getMcpUserServerInstructions } from './user-server-instructions-repo.ts' +import { getCapabilityRegistryForContext } from './capabilities/registry.ts' +import { listPopularAgentPackagesForUser } from '#worker/usage/agent-package-conversation-uses.ts' + +const kodyMcpServerInfo = { + name: 'kody-mcp', + version: '1.0.0', +} as const + +/** + * Serve one modern-era `/mcp` request. Called after `handleMcpRequest` has + * authenticated the bearer token and built the caller context; this handler + * performs no token verification of its own. + */ +export async function handleStatelessMcpRequest(input: { + request: Request + env: Env + ctx: ExecutionContext + callerContext: McpCallerContext + /** Body already parsed by lane classification, so it is read only once. */ + parsedBody?: unknown +}): Promise { + const { request, env, ctx, callerContext } = input + const handler = createMcpHandler( + async () => { + // `server/discover` is the only modern method that returns server + // metadata, so instruction assembly (three D1 reads) stays off the + // tools/call hot path. Modern Streamable HTTP requires the + // Mcp-Method header, and the SDK rejects header/body mismatches. + const instructions = + request.headers.get('Mcp-Method') === 'server/discover' + ? await buildStatelessServerInstructions({ env, callerContext }) + : undefined + const server = new McpServer(kodyMcpServerInfo, { + ...(instructions ? { instructions } : {}), + jsonSchemaValidator: new CfWorkerJsonSchemaValidator(), + }) + await registerTools( + createStatelessRegistrationAgent({ server, env, ctx, callerContext }), + ) + return server + }, + { + // The Durable Object McpAgent lane owns 2025-era traffic (sessions, + // SSE streams, resumability); rejecting legacy here keeps exactly + // one lane authoritative per era. + legacy: 'reject', + onerror: (error) => console.warn('mcp-stateless-lane-error', error), + }, + ) + return handler.fetch( + request, + input.parsedBody === undefined + ? undefined + : { parsedBody: input.parsedBody }, + ) +} + +async function buildStatelessServerInstructions(input: { + env: Env + callerContext: McpCallerContext +}) { + const userId = input.callerContext.user?.userId ?? null + const [overlay, registry, popularPackages] = await Promise.all([ + userId !== null + ? getMcpUserServerInstructions(input.env.APP_DB, userId) + : Promise.resolve(null), + getCapabilityRegistryForContext({ + env: input.env, + callerContext: input.callerContext, + }), + userId !== null + ? listPopularAgentPackagesForUser(input.env.APP_DB, { userId }) + : Promise.resolve([]), + ]) + return buildMcpServerInstructions({ + userOverlay: overlay, + domains: registry.capabilityDomains, + popularPackages, + }) +} + +/** + * Per-request registration agent for the stateless lane. Unlike the Durable + * Object lane there is no agent state, so stateful niceties (search preamble + * dedup, raw-fetch host nudge accumulation) degrade to per-call behavior — + * the tool runners already handle a stateless agent defensively. + */ +function createStatelessRegistrationAgent(input: { + server: McpServer + env: Env + ctx: ExecutionContext + callerContext: McpCallerContext +}): McpRegistrationAgent { + return { + server: asMcpToolServer(input.server), + getEnv: () => input.env, + getCallerContext: () => input.callerContext, + requireDomain: () => { + const { baseUrl } = input.callerContext + invariant( + baseUrl, + 'This should never happen, but somehow we did not get the baseUrl from the request handler', + ) + return baseUrl + }, + getLoopbackExports: () => input.ctx.exports as typeof workerExports, + waitUntil: (promise) => input.ctx.waitUntil(promise), + } +} diff --git a/packages/worker/src/mcp/tools/execute.ts b/packages/worker/src/mcp/tools/execute.ts index 22e16ae25c..2909d880c3 100644 --- a/packages/worker/src/mcp/tools/execute.ts +++ b/packages/worker/src/mcp/tools/execute.ts @@ -40,6 +40,7 @@ import { } from './memory-tool-context.ts' import { finishToolTiming, startToolTiming } from './tool-timing.ts' import { prependToolMetadataContent } from './tool-response-content.ts' +import { buildKodyToolIcons } from './tool-icons.ts' import { applyRawFetchHostCounts, codeUsesIntegrationAuthHelpers, @@ -79,12 +80,94 @@ const executeTool = { } satisfies ToolAnnotations, } as const +/** + * Advertised MCP output schema for the execute tool's `structuredContent` + * envelope. Deliberately loose: every field is optional and compound values + * are `z.unknown()`, so server-side output validation (which runs on every + * successful call once a schema is advertised) can never reject a real + * response. The module's own return value is arbitrary caller JSON and stays + * `unknown` by construction. + */ +export const executeToolOutputSchema = { + conversationId: z + .string() + .optional() + .describe( + 'Tool conversation id; pass it back on subsequent search/execute calls.', + ), + timing: z + .unknown() + .optional() + .describe( + 'Server-side timing: startedAt, endedAt, durationMs, optional serverTiming phases.', + ), + storage: z + .unknown() + .optional() + .describe('Bound durable storage descriptor ({ id }) when active.'), + runId: z + .string() + .optional() + .describe('Run record id for keyed or persisted runs; poll via run_get.'), + inProgress: z + .boolean() + .optional() + .describe( + 'True when a keyed retry found the original run still executing.', + ), + status: z + .string() + .optional() + .describe('Run status accompanying inProgress lookups.'), + replayed: z + .boolean() + .optional() + .describe('True when a keyed retry returned a retained earlier result.'), + returnedBytes: z + .number() + .optional() + .describe('Serialized size of the returned value.'), + truncated: z + .boolean() + .optional() + .describe('True when the result was truncated to fit responseLimit.'), + note: z.string().optional().describe('Truncation note when truncated.'), + result: z + .unknown() + .optional() + .describe("The module default export's return value (arbitrary JSON)."), + error: z + .string() + .optional() + .describe('Error summary when execution failed (isError is set).'), + errorName: z.string().optional().describe('Error name for replayed errors.'), + errorDetails: z + .unknown() + .optional() + .describe('Structured details for sandbox errors.'), + logs: z + .array(z.unknown()) + .optional() + .describe('Console output captured from the sandboxed module.'), + warnings: z + .array(z.string()) + .optional() + .describe('Server guidance, e.g. raw-fetch host nudges.'), + memories: z + .unknown() + .optional() + .describe('Relevant stored memories surfaced for this call.'), +} + export async function registerExecuteTool(agent: McpRegistrationAgent) { + const icons = buildKodyToolIcons(agent.getCallerContext().baseUrl) agent.server.registerTool( executeTool.name, { title: executeTool.title, description: executeTool.description, + outputSchema: executeToolOutputSchema, + ...(icons ? { icons } : {}), inputSchema: { code: z .string() diff --git a/packages/worker/src/mcp/tools/output-schemas.node.test.ts b/packages/worker/src/mcp/tools/output-schemas.node.test.ts new file mode 100644 index 0000000000..63ded0eabb --- /dev/null +++ b/packages/worker/src/mcp/tools/output-schemas.node.test.ts @@ -0,0 +1,137 @@ +/** + * Advertising an MCP `outputSchema` makes the SDK validate every successful + * result's `structuredContent` against it server-side (v1 `McpServer` runs + * `safeParseAsync` on the zod object; v2 validates the converted JSON + * schema). A schema that rejects a real response would turn working calls + * into protocol errors, so these tests run representative structured + * responses from every return path through the advertised schemas exactly + * the way the SDK does. + */ + +import { expect, test } from 'vitest' +import { z } from 'zod' +import { searchToolOutputSchema } from './search-tool-definition.ts' +import { executeToolOutputSchema } from './execute.ts' + +const timing = { + startedAt: '2026-08-05T00:00:00.000Z', + endedAt: '2026-08-05T00:00:01.000Z', + durationMs: 1000, +} + +const searchSchema = z.object(searchToolOutputSchema) +const executeSchema = z.object(executeToolOutputSchema) + +test.each([ + [ + 'ranked list result', + { + conversationId: 'c1', + timing, + result: { + offline: false, + warnings: [], + telemetry: { topResultTypes: ['capability'] }, + phaseTimings: { formattingMs: 2 }, + matches: [{ type: 'capability', name: 'email_send' }], + }, + }, + ], + [ + 'entity detail result', + { conversationId: 'c1', timing, result: { entityRef: 'x:capability' } }, + ], + [ + 'entity batch with partial failures', + { + conversationId: 'c1', + timing, + result: [{ entityRef: 'a:capability', error: 'not found' }], + error: 'All entity lookups failed.', + }, + ], + [ + 'validation error', + { conversationId: 'c1', timing, error: 'Provide "query".' }, + ], +])( + 'search structured content passes advertised schema: %s', + async (_, payload) => { + const parsed = await searchSchema.safeParseAsync(payload) + expect(parsed.success).toBe(true) + }, +) + +test.each([ + [ + 'success with storage and warnings', + { + conversationId: 'c1', + timing: { ...timing, serverTiming: [{ name: 'registry', ms: 5 }] }, + storage: { id: 'bucket-1' }, + returnedBytes: 42, + result: { anything: ['goes', 1, null] }, + logs: ['log line', { level: 'warn' }], + warnings: ['Consider integration auth helpers for api.example.com.'], + memories: { surfaced: [], suppressedCount: 0 }, + }, + ], + [ + 'truncated result', + { + conversationId: 'c1', + timing, + returnedBytes: 100000, + truncated: true, + note: 'Result truncated to fit responseLimit.', + result: 'partial…', + logs: [], + }, + ], + [ + 'sandbox error', + { + conversationId: 'c1', + timing, + returnedBytes: 0, + error: 'ReferenceError: foo is not defined', + errorDetails: { phase: 'sandbox' }, + logs: [], + }, + ], + [ + 'keyed replay', + { + conversationId: 'c1', + timing, + runId: 'run-1', + replayed: true, + returnedBytes: 0, + result: { ok: true }, + logs: [], + }, + ], + [ + 'keyed in-progress lookup', + { + conversationId: 'c1', + timing, + runId: 'run-1', + inProgress: true, + status: 'running', + }, + ], +])( + 'execute structured content passes advertised schema: %s', + async (_, payload) => { + const parsed = await executeSchema.safeParseAsync(payload) + expect(parsed.success).toBe(true) + }, +) + +test('schemas convert to JSON Schema without throwing', () => { + // The SDK advertises the zod schemas as JSON Schema on tools/list; a + // conversion failure would drop the advertisement or break listing. + expect(() => z.toJSONSchema(searchSchema, { io: 'output' })).not.toThrow() + expect(() => z.toJSONSchema(executeSchema, { io: 'output' })).not.toThrow() +}) diff --git a/packages/worker/src/mcp/tools/search-register.ts b/packages/worker/src/mcp/tools/search-register.ts index c78cf367e3..1482d97294 100644 --- a/packages/worker/src/mcp/tools/search-register.ts +++ b/packages/worker/src/mcp/tools/search-register.ts @@ -3,18 +3,23 @@ import { type McpRegistrationAgent } from '#mcp/mcp-registration-agent.ts' import { searchTool, searchToolInputSchema, + searchToolOutputSchema, type SearchToolArgs, } from './search-tool-definition.ts' import { runSearchTool } from './search-tool-runner.ts' +import { buildKodyToolIcons } from './tool-icons.ts' export async function registerSearchTool(agent: McpRegistrationAgent) { + const icons = buildKodyToolIcons(agent.getCallerContext().baseUrl) agent.server.registerTool( searchTool.name, { title: searchTool.title, description: searchTool.description, inputSchema: searchToolInputSchema, + outputSchema: searchToolOutputSchema, annotations: searchTool.annotations, + ...(icons ? { icons } : {}), }, async (args: SearchToolArgs) => runSearchTool({ agent, args }), ) diff --git a/packages/worker/src/mcp/tools/search-tool-definition.ts b/packages/worker/src/mcp/tools/search-tool-definition.ts index d16f5a816e..112e8f0df6 100644 --- a/packages/worker/src/mcp/tools/search-tool-definition.ts +++ b/packages/worker/src/mcp/tools/search-tool-definition.ts @@ -116,6 +116,39 @@ export const searchToolInputSchema = { ), } +/** + * Advertised MCP output schema for the search tool's `structuredContent` + * envelope. Deliberately loose: every field is optional and compound values + * are `z.unknown()`, so server-side output validation (which runs on every + * successful call once a schema is advertised) can never reject a real + * response. The schema documents the envelope for clients; mode-specific + * payload shapes stay in the tool description and docs. + */ +export const searchToolOutputSchema = { + conversationId: z + .string() + .optional() + .describe( + 'Tool conversation id; pass it back on subsequent search/execute calls.', + ), + timing: z + .unknown() + .optional() + .describe( + 'Server-side timing: startedAt, endedAt, durationMs, optional serverTiming phases.', + ), + result: z + .unknown() + .optional() + .describe( + 'Mode-specific structured payload: ranked matches with telemetry, domain browse listing, entity detail, or entity-batch results.', + ), + error: z + .string() + .optional() + .describe('Error summary when the search failed (isError is set).'), +} + export type SearchToolArgs = { query?: string entity?: string | Array diff --git a/packages/worker/src/mcp/tools/tool-icons.ts b/packages/worker/src/mcp/tools/tool-icons.ts new file mode 100644 index 0000000000..c8a4a2414f --- /dev/null +++ b/packages/worker/src/mcp/tools/tool-icons.ts @@ -0,0 +1,26 @@ +import { type McpToolIcon } from '#mcp/mcp-registration-agent.ts' + +/** + * Kody tool icons (protocol revision 2025-11-25, SEP-973). Servers must use + * absolute URIs, so icons are derived from the caller's base URL at + * registration time. Advertised by the stateless SDK v2 lane; the SDK v1 + * `McpAgent` lane ignores the config key. Returns `undefined` when the base + * URL is unknown (some unit-test agents), which omits the field entirely. + */ +export function buildKodyToolIcons( + baseUrl: string | null | undefined, +): Array | undefined { + if (!baseUrl) return undefined + return [ + { + src: new URL('/android-chrome-192x192.png', baseUrl).toString(), + mimeType: 'image/png', + sizes: ['192x192'], + }, + { + src: new URL('/android-chrome-512x512.png', baseUrl).toString(), + mimeType: 'image/png', + sizes: ['512x512'], + }, + ] +} diff --git a/packages/worker/worker-configuration.d.ts b/packages/worker/worker-configuration.d.ts index f435ad2417..3af361f79c 100644 --- a/packages/worker/worker-configuration.d.ts +++ b/packages/worker/worker-configuration.d.ts @@ -1,5 +1,5 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types --config=packages/worker/wrangler.jsonc --env=production ./packages/worker/worker-configuration.d.ts` (hash: ff721fcd9d867202d94fb7f9ffbd1993) +// Generated by Wrangler by running `wrangler types --config=packages/worker/wrangler.jsonc --env=production ./packages/worker/worker-configuration.d.ts` (hash: 8be293561e7f31f9f513d88c67c3b61a) // Runtime types generated with workerd@1.20260730.1 2026-04-16 global_fetch_strictly_public,nodejs_compat interface __BaseEnv_Env { OAUTH_KV: KVNamespace; @@ -13,6 +13,7 @@ interface __BaseEnv_Env { USAGE_EVENTS: AnalyticsEngineDataset; FLAG_EXPOSURES: AnalyticsEngineDataset; EMAIL_EVENTS: AnalyticsEngineDataset; + MCP_PROTOCOL_EVENTS: AnalyticsEngineDataset; PLATFORM_FEEDBACK_DISPATCH_QUEUE: Queue; COMMUNITY_ACTIVITY_DISPATCH_QUEUE: Queue; SCHEDULED_DISPATCH_QUEUE: Queue; diff --git a/packages/worker/wrangler.jsonc b/packages/worker/wrangler.jsonc index 59fd6cac7a..ff755b13bc 100644 --- a/packages/worker/wrangler.jsonc +++ b/packages/worker/wrangler.jsonc @@ -365,6 +365,10 @@ "binding": "EMAIL_EVENTS", "dataset": "kody_email_events", }, + { + "binding": "MCP_PROTOCOL_EVENTS", + "dataset": "kody_mcp_protocol_events", + }, ], "ai": { "binding": "AI", @@ -565,6 +569,10 @@ "binding": "EMAIL_EVENTS", "dataset": "kody_email_events_preview", }, + { + "binding": "MCP_PROTOCOL_EVENTS", + "dataset": "kody_mcp_protocol_events_preview", + }, ], "ai": { "binding": "AI", diff --git a/tools/mcp-test-support.ts b/tools/mcp-test-support.ts index 770e1487ec..a1b2c48108 100644 --- a/tools/mcp-test-support.ts +++ b/tools/mcp-test-support.ts @@ -8,6 +8,10 @@ import { stripVTControlCharacters } from 'node:util' import { Client } from '@modelcontextprotocol/sdk/client/index.js' import { StreamableHTTPClientTransport } from '@modelcontextprotocol/sdk/client/streamableHttp.js' import { type CallToolRequest } from '@modelcontextprotocol/sdk/types.js' +import { + Client as ModernClient, + StreamableHTTPClientTransport as ModernStreamableHTTPClientTransport, +} from '@modelcontextprotocol/client' import getPort from 'get-port' import { captureOutput, @@ -513,6 +517,57 @@ export async function createMcpClient( } } +/** + * Modern-era MCP client (protocol revision 2026-07-28) from the SDK v2 + * client package, pinned so the connection fails loudly unless the server + * serves the stateless lane. Reuses the same signup + OAuth plumbing as + * `createMcpClient`. + */ +export async function createModernMcpClient( + origin: string, + user: TestUser, + options: { + persistDir: string + }, +) { + const cookieHeader = await loginToApp(origin, user) + await markEmailVerifiedInMcpTestDatabase({ + persistDir: options.persistDir, + email: user.email, + }) + const clientRegistration = await registerOAuthClient(origin) + const code = await authorizeOAuthClient( + origin, + clientRegistration, + cookieHeader, + ) + const accessToken = await exchangeAuthorizationCode( + origin, + clientRegistration, + code, + ) + const client = new ModernClient( + { name: 'kody-mcp-e2e-modern-client', version: '1.0.0' }, + { versionNegotiation: { mode: { pin: '2026-07-28' } } }, + ) + const transport = new ModernStreamableHTTPClientTransport( + new URL('/mcp', origin), + { + requestInit: { + headers: { Authorization: `Bearer ${accessToken}` }, + }, + }, + ) + await client.connect(transport) + return { + client, + async [Symbol.asyncDispose]() { + await client.close().catch(() => undefined) + await transport.close().catch(() => undefined) + }, + } +} + async function fetchJson>( origin: string, pathname: string,