From bebd8d74f7ea718f8c6f158b858ae8ab04b42c3f Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Mon, 14 Sep 2026 18:11:39 -0700 Subject: [PATCH 1/2] Add search and invoke MCP mode --- .changeset/mcp-passthrough-mode.md | 7 + apps/cli/src/main.ts | 33 +- apps/cloud/src/mcp/agent-handler.ts | 2 + apps/cloud/src/mcp/session-durable-object.ts | 4 + apps/cloud/src/mcp/session-meta.ts | 1 + apps/host-cloudflare/src/mcp/agent-handler.ts | 2 + .../src/mcp/session-durable-object.ts | 3 + apps/local/src/main.ts | 3 + apps/local/src/mcp.ts | 2 + .../passthrough-opencode-codemode.test.ts | 327 +++++++++ e2e/cloud/passthrough-scale.test.ts | 75 ++ e2e/scenarios/mcp-passthrough.test.ts | 300 ++++++++ e2e/scenarios/support/search-invoke.ts | 25 + e2e/src/surfaces/mcp.ts | 48 +- packages/core/api/src/server/mcp-build.ts | 3 + packages/core/execution/src/tool-invoker.ts | 2 +- packages/core/sdk/src/elicitation.ts | 10 + packages/core/sdk/src/executor.ts | 3 +- packages/core/sdk/src/host-internal.ts | 3 + packages/core/sdk/src/index.ts | 1 + .../src/mcp/agent-session-durable-object.ts | 9 + .../hosts/cloudflare/src/mcp/do-headers.ts | 2 + packages/hosts/mcp/src/browser-approval.ts | 15 + .../hosts/mcp/src/in-memory-session-store.ts | 28 +- .../hosts/mcp/src/passthrough-tools.test.ts | 668 ++++++++++++++++++ packages/hosts/mcp/src/passthrough-tools.ts | 29 + packages/hosts/mcp/src/tool-server.ts | 506 ++++++++++--- packages/react/src/api/analytics.tsx | 1 + .../react/src/components/mcp-install-card.tsx | 64 +- 29 files changed, 2065 insertions(+), 111 deletions(-) create mode 100644 .changeset/mcp-passthrough-mode.md create mode 100644 e2e/cloud/passthrough-opencode-codemode.test.ts create mode 100644 e2e/cloud/passthrough-scale.test.ts create mode 100644 e2e/scenarios/mcp-passthrough.test.ts create mode 100644 e2e/scenarios/support/search-invoke.ts create mode 100644 packages/hosts/mcp/src/passthrough-tools.test.ts create mode 100644 packages/hosts/mcp/src/passthrough-tools.ts diff --git a/.changeset/mcp-passthrough-mode.md b/.changeset/mcp-passthrough-mode.md new file mode 100644 index 0000000000..9008960f31 --- /dev/null +++ b/.changeset/mcp-passthrough-mode.md @@ -0,0 +1,7 @@ +--- +"@executor-js/sdk": minor +"@executor-js/execution": minor +"executor": minor +--- + +Add a search and invoke MCP mode (`?mode=passthrough`, `executor mcp --mode passthrough`). Search returns bounded pages of matching tool IDs and input schemas. Invoke validates arguments and runs the selected tool, with native client approval and workspace blocks enforced. The MCP catalog stays at two tools regardless of integration count. diff --git a/apps/cli/src/main.ts b/apps/cli/src/main.ts index fc2a9bf391..e41a723b89 100644 --- a/apps/cli/src/main.ts +++ b/apps/cli/src/main.ts @@ -1363,6 +1363,7 @@ const mcpUrlForActiveLocalServer = (input: { readonly elicitationMode: "browser" | "model"; readonly artifacts: boolean; readonly searchTools: boolean; + readonly toolMode: "codemode" | "passthrough"; }): URL => { const url = new URL("/mcp", input.connection.origin); if (input.elicitationMode === "browser") { @@ -1378,6 +1379,10 @@ const mcpUrlForActiveLocalServer = (input: { if (input.searchTools) { url.searchParams.set("search_tools", "true"); } + // Passthrough is the non-default surface; only it is spelled out. + if (input.toolMode === "passthrough") { + url.searchParams.set("mode", "passthrough"); + } return url; }; @@ -1394,6 +1399,7 @@ const runMcpHttpBridge = async (input: { readonly elicitationMode: "browser" | "model"; readonly artifacts: boolean; readonly searchTools: boolean; + readonly toolMode: "codemode" | "passthrough"; }): Promise => { const stdio = new StdioServerTransport(); const authorization = getExecutorServerAuthorizationHeader(input.manifest.connection); @@ -1403,6 +1409,7 @@ const runMcpHttpBridge = async (input: { elicitationMode: input.elicitationMode, artifacts: input.artifacts, searchTools: input.searchTools, + toolMode: input.toolMode, }), authorization ? { requestInit: { headers: { Authorization: authorization } } } : undefined, ); @@ -1482,6 +1489,7 @@ const runStdioMcpSession = (input: { readonly elicitationMode: "browser" | "model"; readonly artifacts: boolean; readonly searchTools: boolean; + readonly toolMode: "codemode" | "passthrough"; }) => Effect.gen(function* () { // `executor mcp` never owns the local database. If a local server is already @@ -1499,6 +1507,7 @@ const runStdioMcpSession = (input: { elicitationMode: input.elicitationMode, artifacts: input.artifacts, searchTools: input.searchTools, + toolMode: input.toolMode, }), ); return; @@ -1526,6 +1535,7 @@ const runStdioMcpSession = (input: { elicitationMode: input.elicitationMode, artifacts: input.artifacts, searchTools: input.searchTools, + toolMode: input.toolMode, }), ); }); @@ -2898,11 +2908,30 @@ const mcpCommand = Command.make( "Serve one search_ tool per connected integration. Off by default; each routes through the same flow as tools.search inside execute.", ), ), + toolMode: Options.choice("mode", ["codemode", "passthrough"] as const) + .pipe(Options.withDefault("codemode")) + .pipe( + Options.withDescription( + "codemode (default) serves the execute tool; passthrough serves search and invoke, with input schemas in search results and client approval for invoke.", + ), + ), }, - ({ scope, elicitationMode, noArtifacts, searchTools }) => + ({ scope, elicitationMode, noArtifacts, searchTools, toolMode }) => Effect.gen(function* () { applyScope(scope); - yield* runStdioMcpSession({ elicitationMode, artifacts: !noArtifacts, searchTools }); + if (toolMode === "passthrough" && searchTools) { + return yield* Effect.fail( + new Error( + "--search-tools is a codemode option; passthrough already provides search. Drop --search-tools or --mode passthrough.", + ), + ); + } + yield* runStdioMcpSession({ + elicitationMode, + artifacts: !noArtifacts, + searchTools, + toolMode, + }); }), ).pipe(Command.withDescription("Start an MCP server over stdio")); diff --git a/apps/cloud/src/mcp/agent-handler.ts b/apps/cloud/src/mcp/agent-handler.ts index 05c2107ed0..dc858cee3f 100644 --- a/apps/cloud/src/mcp/agent-handler.ts +++ b/apps/cloud/src/mcp/agent-handler.ts @@ -16,6 +16,7 @@ import { readArtifactsEnabled, readElicitationMode, readSearchToolsEnabled, + readToolMode, withVerifiedIdentityHeaders, } from "@executor-js/cloudflare/mcp/do-headers"; import type { McpSessionProps } from "@executor-js/cloudflare/mcp/agent-durable-object"; @@ -189,6 +190,7 @@ const propsForPrincipal = ( elicitationMode: readElicitationMode(request), artifactsEnabled: readArtifactsEnabled(request), searchToolsEnabled: readSearchToolsEnabled(request), + toolMode: readToolMode(request), resource, webOrigin: new URL(request.url).origin, }, diff --git a/apps/cloud/src/mcp/session-durable-object.ts b/apps/cloud/src/mcp/session-durable-object.ts index 2646a785a0..3f6be9ef28 100644 --- a/apps/cloud/src/mcp/session-durable-object.ts +++ b/apps/cloud/src/mcp/session-durable-object.ts @@ -384,6 +384,7 @@ export class McpSessionDOSqlite extends McpAgentSessionDOBase `${prefix}_${randomBytes(4).toString("hex")}`; + +/** The one integration that is actually callable: a read and a write against + * the recording upstream. Everything else in the catalog is discovery mass. */ +const notesSpec = (baseUrl: string): string => + JSON.stringify({ + openapi: "3.0.3", + info: { title: "Notes API", version: "1.0.0" }, + servers: [{ url: baseUrl }], + paths: { + "/notes": { + get: { + operationId: "listNotes", + summary: "List every note in the notebook", + responses: { "200": { description: "ok" } }, + }, + post: { + operationId: "createNote", + summary: "Create a note in the notebook", + requestBody: { + required: true, + content: { + "application/json": { + schema: { + type: "object", + properties: { text: { type: "string" } }, + required: ["text"], + }, + }, + }, + }, + responses: { "200": { description: "ok" } }, + }, + }, + }, + }); + +interface RecordedRequest { + readonly method: string; + readonly authorization: string | undefined; +} + +const serveRecordingUpstream = Effect.acquireRelease( + Effect.callback<{ + readonly url: string; + readonly requests: RecordedRequest[]; + close: () => void; + }>((resume) => { + const requests: RecordedRequest[] = []; + const server = createServer((request, response) => { + request.on("data", () => undefined); + request.on("end", () => { + requests.push({ + method: request.method ?? "", + authorization: request.headers.authorization, + }); + response.writeHead(200, { "content-type": "application/json" }); + response.end(JSON.stringify({ notes: [{ id: "note_0", text: "existing note" }] })); + }); + }); + server.listen(0, "127.0.0.1", () => { + const address = server.address(); + const port = typeof address === "object" && address ? address.port : 0; + resume( + Effect.succeed({ + url: `http://127.0.0.1:${port}`, + requests, + close: () => { + server.close(); + server.closeAllConnections(); + }, + }), + ); + }); + }), + (upstream) => Effect.sync(() => upstream.close()), +); + +scenario( + "Passthrough · the real OpenCode binary runs its own codemode over 10,200 Executor tools", + { timeout: 900_000 }, + Effect.scoped( + Effect.gen(function* () { + const target = yield* Target; + const opencode = yield* OpenCode; + const runDir = yield* RunDir; + const cli = yield* Cli; + const { client: makeClient } = yield* Api; + + const identity = yield* target.newIdentity(); + const email = identity.credentials?.email ?? identity.label; + const client = yield* makeClient(catalogApi, identity); + const upstream = yield* serveRecordingUpstream; + const notesSlug = unique("notes"); + + // --- Seed: 10,198 discovery-mass tools + 2 callable ones. --- + const seeded = yield* seedLargeCatalog(client, { + syntheticIntegrations: SYNTHETIC_INTEGRATIONS, + opsPerIntegration: OPS_PER_INTEGRATION, + }); + const cleanup = Effect.gen(function* () { + yield* client.connections + .remove({ + params: { + owner: "org", + integration: IntegrationSlug.make(notesSlug), + name: ConnectionName.make("main"), + }, + }) + .pipe(Effect.ignore); + yield* client.openapi.removeSpec({ params: { slug: notesSlug } }).pipe(Effect.ignore); + yield* seeded.cleanup; + }); + + yield* Effect.ensuring( + Effect.gen(function* () { + yield* client.openapi.addSpec({ + payload: { + spec: { kind: "blob", value: notesSpec(upstream.url) }, + slug: notesSlug, + baseUrl: upstream.url, + authenticationTemplate: [ + { + slug: "apiKey", + type: "apiKey", + headers: { authorization: ["Bearer ", { type: "variable", name: "token" }] }, + }, + ], + }, + }); + yield* client.connections.create({ + payload: { + owner: "org", + name: ConnectionName.make("main"), + integration: IntegrationSlug.make(notesSlug), + template: AuthTemplateSlug.make("apiKey"), + value: "tok_notes", + }, + }); + const visible = (yield* client.tools.list({ query: {} })).filter( + (tool) => tool.static !== true, + ); + expect(visible.length, "the catalog is exactly the advertised size").toBe( + EXPECTED_TOOL_COUNT, + ); + + // --- The replay brain: two scripted turns of OpenCode codemode. --- + // Turn 0: discover the notes tool through Executor search. + // Turn 1: call the path search returned. Turn 2: summarize, stop. + let discoveredId: string | undefined; + const brain = yield* serveReplayBrain((ctx) => { + // OpenCode also asks the model for a session title, with no tools + // offered. Answer it with text and keep it out of the script. + if (ctx.toolNames.length === 0) return { text: "Notebook" }; + if (ctx.lastRole === "user") { + return { + text: "Searching the connected tools.", + tool: { + name: "execute", + args: { + code: `return await tools.${SERVER_NAME}.search({ query: "${notesSlug} listNotes", limit: 5 });`, + }, + }, + }; + } + if (discoveredId === undefined) { + const result = ctx.lastToolResult ?? ""; + // Search returns an opaque tool ID for the next invoke call. + const match = new RegExp(`"id":\\s*"([^"\\n]*${notesSlug}[^"]*listNotes)"`).exec( + result, + ); + if (!match) { + throw new Error(`search did not surface the notes tool: ${result.slice(0, 600)}`); + } + discoveredId = match[1]!; + return { + text: "Found it. Listing the notes.", + tool: { + name: "execute", + args: { + code: `return await tools.${SERVER_NAME}.invoke({tool: ${JSON.stringify(discoveredId)}, arguments: {}});`, + }, + }, + }; + } + return { text: "The notebook has one existing note." }; + }); + + const passthroughUrl = new URL("/mcp?mode=passthrough", target.baseUrl).toString(); + const home = opencode.makeHome(SERVER_NAME, passthroughUrl, { + chatBrainUrl: brain.baseUrl, + }); + const env = { + ...home.env, + OPENCODE_EXPERIMENTAL_CODE_MODE: "true", + PS1: "$ ", + BASH_SILENCE_DEPRECATION_WARNING: "1", + }; + // First-run database migration happens off camera. + yield* Effect.sync(() => opencode.warmUp(home)); + + let connectMs = -1; + yield* cli.session( + ["bash", "--norc"], + async (term) => { + await term.screen.waitForText("$", { timeoutMs: 10_000 }); + + const outputAfter = (text: string, line: string): string | null => { + const echoed = text.lastIndexOf(line); + if (echoed === -1) return null; + const after = text.slice(echoed + line.length); + return after.trimEnd().endsWith("\n$") ? after : null; + }; + const sh = async (line: string, timeoutMs: number) => { + await term.keyboard.type(line); + await term.keyboard.press("Enter"); + const snapshot = await term.screen.waitUntil( + (current) => outputAfter(current.text, line) !== null, + { timeoutMs }, + ); + return outputAfter(snapshot.text, line) ?? ""; + }; + + // OpenCode's own OAuth against the target: discovery, DCR, PKCE. + const consent = opencode.completeOAuthConsent(home, email, home.openedUrls().length); + const auth = await sh(`opencode mcp auth ${SERVER_NAME}`, 90_000); + await consent; + expect(auth, "opencode mcp auth completes").not.toContain("failed"); + + // The load-bearing connect: catalog loading and two MCP definitions + // inside OpenCode's own 30s connect timeout. + const startedAt = Date.now(); + const listed = await sh("opencode mcp list", 120_000); + connectMs = Date.now() - startedAt; + expect( + listed, + `OpenCode connects to the 10,200-tool passthrough endpoint (took ${connectMs}ms)`, + ).toContain("connected"); + + // A real agent turn: OpenCode's codemode over our catalog. + const ran = await sh(`opencode run "List the notes in my notebook"`, 300_000); + expect(ran, "the run did not error").not.toContain("UnknownError"); + }, + { + cwd: home.projectDir, + env, + record: join(runDir, "terminal.cast"), + viewport: { cols: 100, rows: 40 }, + }, + ); + + // --- What OpenCode showed the model, and what came back. --- + expect(brain.errors(), "the scripted brain hit no surprises").toEqual([]); + // Only the turns that carried tools are the agent loop; the + // title request is OpenCode housekeeping. + const requests = brain.requests().filter((request) => request.toolNames.length > 0); + expect(requests.length, "three model turns: search, call, summary").toBe(3); + + // Codemode: the model sees ONE execute tool, not 10,200 functions. + const offered = requests[0]!.toolNames; + expect(offered, "OpenCode offers its codemode execute tool").toContain("execute"); + expect( + offered.filter((name) => name.startsWith(`${SERVER_NAME}_`)), + "no MCP tool is flattened into the model's tool list", + ).toEqual([]); + expect( + offered.length, + "the model's tool list stays small in front of a 10,200-tool server", + ).toBeLessThan(40); + + // Discovery returned the full underlying tool ID. + expect(discoveredId, "search returned a tool ID").toBeDefined(); + expect(discoveredId, "the ID names the integration tool").toContain(notesSlug); + + // The call went over the wire with the connection's credential and + // the payload came back through OpenCode's interpreter. + const lastToolResult = [...requests[2]!.messages] + .reverse() + .find((message) => message.role === "tool")?.content; + expect(lastToolResult, "the executed program returned the upstream payload").toContain( + "existing note", + ); + const upstreamGet = upstream.requests.find((request) => request.method === "GET"); + expect(upstreamGet, "the GET reached the upstream").toBeDefined(); + expect(upstreamGet?.authorization, "the connection's credential was applied").toBe( + "Bearer tok_notes", + ); + + // Recorded for the run report; the hard bound is OpenCode's own + // connect timeout, which `mcp list` above already proved. + expect(connectMs, "connect stays inside OpenCode's timeout").toBeLessThan( + OPENCODE_CONNECT_TIMEOUT_MS * 4, + ); + }), + cleanup, + ); + }), + ), +); diff --git a/e2e/cloud/passthrough-scale.test.ts b/e2e/cloud/passthrough-scale.test.ts new file mode 100644 index 0000000000..c9f92fff59 --- /dev/null +++ b/e2e/cloud/passthrough-scale.test.ts @@ -0,0 +1,75 @@ +// Large catalogs stay server-side and remain searchable within the connect timeout. +import { expect } from "@effect/vitest"; +import { Effect } from "effect"; + +import { decodeToolSearch } from "../scenarios/support/search-invoke"; + +import { scenario } from "../src/scenario"; +import { Api, Mcp, Target } from "../src/services"; +import { catalogApi, seedLargeCatalog } from "../scenarios/support/large-catalog"; + +// Catalog loading and the two-tool handshake must complete within 20 seconds. +const MAX_PASSTHROUGH_CONNECT_MS = 20_000; + +scenario( + "Passthrough · a production-shaped catalog is served completely, in bounded time", + { timeout: 300_000 }, + Effect.scoped( + Effect.gen(function* () { + const target = yield* Target; + const mcp = yield* Mcp; + const { client: makeClient } = yield* Api; + + const identity = yield* target.newIdentity(); + const client = yield* makeClient(catalogApi, identity); + const seeded = yield* seedLargeCatalog(client); + + yield* Effect.ensuring( + Effect.gen(function* () { + // What the caller can see through the typed API is the ground truth + // search results must cover completely — minus the plugins' + // static configuration tools (`executor.*`, `openapi.addSpec`, …), + // which are codemode affordances and deliberately not served here. + const visible = (yield* client.tools.list({ query: {} })).filter( + (tool) => tool.static !== true, + ); + expect(visible.length, "the seeded catalog is large").toBeGreaterThan(3000); + + const session = mcp.session(identity, { mode: "passthrough" }); + const startedAt = Date.now(); + const served = yield* session.describeTools(); + const elapsedMs = Date.now() - startedAt; + + expect( + elapsedMs, + `a ${visible.length}-tool passthrough connect stays bounded (took ${elapsedMs}ms)`, + ).toBeLessThan(MAX_PASSTHROUGH_CONNECT_MS); + + expect(served.map((tool) => tool.name).sort()).toEqual(["invoke", "search"]); + const first = decodeToolSearch( + (yield* session.call("search", { query: "org", limit: 20 })).raw, + ).structuredContent; + expect(first.total).toBe(visible.length); + expect(first.items).toHaveLength(20); + expect(first.nextOffset).toBe(20); + const last = decodeToolSearch( + (yield* session.call("search", { query: "org", offset: visible.length - 1 })).raw, + ).structuredContent; + expect(last.items).toHaveLength(1); + expect(last.hasMore).toBe(false); + expect(last.nextOffset).toBeNull(); + for (const slug of seeded.integrationSlugs) { + const found = decodeToolSearch( + (yield* session.call("search", { query: slug })).raw, + ).structuredContent; + expect( + found.items.some((tool) => tool.integration === slug), + `${slug} is discoverable`, + ).toBe(true); + } + }), + seeded.cleanup, + ); + }), + ), +); diff --git a/e2e/scenarios/mcp-passthrough.test.ts b/e2e/scenarios/mcp-passthrough.test.ts new file mode 100644 index 0000000000..34aaf54f4f --- /dev/null +++ b/e2e/scenarios/mcp-passthrough.test.ts @@ -0,0 +1,300 @@ +// Search discovers schemas; invoke reaches the upstream and enforces workspace blocks. +import { randomBytes } from "node:crypto"; +import { createServer, type IncomingMessage } from "node:http"; + +import { expect } from "@effect/vitest"; +import { Effect } from "effect"; +import { composePluginApi } from "@executor-js/api/server"; +import { openApiHttpPlugin } from "@executor-js/plugin-openapi/api"; +import { + AuthTemplateSlug, + ConnectionName, + IntegrationSlug, + ToolAddress, +} from "@executor-js/sdk/shared"; + +import { decodeToolSearch } from "./support/search-invoke"; + +import { scenario } from "../src/scenario"; +import { Api, Mcp, Target } from "../src/services"; + +const api = composePluginApi([openApiHttpPlugin()] as const); + +const unique = (prefix: string) => `${prefix}_${randomBytes(4).toString("hex")}`; + +/** A two-operation API: a read and a write, so the surface carries one tool + * per policy outcome. The write takes a JSON body with a shared `$ref`, so + * the advertised schema must be self-contained to be usable. */ +const spec = (baseUrl: string): string => + JSON.stringify({ + openapi: "3.0.3", + info: { title: "Passthrough API", version: "1.0.0" }, + servers: [{ url: baseUrl }], + components: { + schemas: { + NewNote: { + type: "object", + properties: { text: { type: "string" } }, + required: ["text"], + }, + }, + }, + paths: { + "/notes": { + get: { + operationId: "listNotes", + summary: "List notes", + responses: { "200": { description: "ok" } }, + }, + post: { + operationId: "createNote", + summary: "Create a note", + requestBody: { + required: true, + content: { + "application/json": { schema: { $ref: "#/components/schemas/NewNote" } }, + }, + }, + responses: { "200": { description: "ok" } }, + }, + }, + }, + }); + +interface RecordedRequest { + readonly method: string; + readonly path: string; + readonly authorization: string | undefined; + readonly body: string; +} + +/** A real upstream that records what reached it, so a passthrough call can be + * proven to have gone over the wire with the connection's credential. */ +const serveRecordingUpstream = Effect.acquireRelease( + Effect.callback<{ + readonly url: string; + readonly requests: RecordedRequest[]; + close: () => void; + }>((resume) => { + const requests: RecordedRequest[] = []; + const readBody = (request: IncomingMessage) => + new Promise((resolve) => { + const chunks: Buffer[] = []; + request.on("data", (chunk: Buffer) => chunks.push(chunk)); + request.on("end", () => resolve(Buffer.concat(chunks).toString("utf8"))); + }); + const server = createServer((request, response) => { + void readBody(request).then((body) => { + requests.push({ + method: request.method ?? "", + path: request.url ?? "", + authorization: request.headers.authorization, + body, + }); + response.writeHead(200, { "content-type": "application/json" }); + response.end( + JSON.stringify( + request.method === "POST" + ? { id: "note_1", ...(body ? (JSON.parse(body) as object) : {}) } + : { notes: [{ id: "note_0", text: "existing" }] }, + ), + ); + }); + }); + server.listen(0, "127.0.0.1", () => { + const address = server.address(); + const port = typeof address === "object" && address ? address.port : 0; + resume( + Effect.succeed({ + url: `http://127.0.0.1:${port}`, + requests, + close: () => { + server.close(); + server.closeAllConnections(); + }, + }), + ); + }); + }), + (upstream) => Effect.sync(() => upstream.close()), +); + +const rawResultOf = (result: { readonly raw: unknown }) => + result.raw as { + content?: ReadonlyArray<{ type: string; text?: string }>; + structuredContent?: Record; + isError?: boolean; + }; + +scenario( + "Passthrough · a session connected with mode=passthrough serves search and invoke with schemas and enforced blocks", + { timeout: 180_000 }, + Effect.scoped( + Effect.gen(function* () { + const target = yield* Target; + const mcp = yield* Mcp; + const { client: makeClient } = yield* Api; + + const identity = yield* target.newIdentity(); + const client = yield* makeClient(api, identity); + const upstream = yield* serveRecordingUpstream; + const slug = unique("ptapi"); + const otherSlug = unique("ptother"); + + yield* Effect.ensuring( + Effect.gen(function* () { + // Two integrations, one connection each: both must appear on the surface. + for (const s of [slug, otherSlug]) { + yield* client.openapi.addSpec({ + payload: { + spec: { kind: "blob", value: spec(upstream.url) }, + slug: s, + baseUrl: upstream.url, + authenticationTemplate: [ + { + slug: "apiKey", + type: "apiKey", + headers: { authorization: ["Bearer ", { type: "variable", name: "token" }] }, + }, + ], + }, + }); + yield* client.connections.create({ + payload: { + owner: "org", + name: ConnectionName.make("main"), + integration: IntegrationSlug.make(s), + template: AuthTemplateSlug.make("apiKey"), + value: `tok_${s}`, + }, + }); + } + + const codemode = mcp.session(identity); + expect(yield* codemode.listTools()).toContain("execute"); + + const passthrough = mcp.session(identity, { mode: "passthrough" }); + const described = yield* passthrough.describeTools(); + expect(described.map((tool) => tool.name).sort()).toEqual(["invoke", "search"]); + expect(described.find((tool) => tool.name === "search")?.annotations).toMatchObject({ + readOnlyHint: true, + destructiveHint: false, + }); + expect(described.find((tool) => tool.name === "invoke")?.annotations).toMatchObject({ + readOnlyHint: false, + destructiveHint: true, + }); + const found = decodeToolSearch( + (yield* passthrough.call("search", { query: slug })).raw, + ).structuredContent; + const listDef = found.items.find((tool) => tool.id.endsWith(".listNotes")); + const createDef = found.items.find((tool) => tool.id.endsWith(".createNote")); + if (!listDef || !createDef) return yield* Effect.die("Search omitted notes operations"); + const listTool = listDef.id; + const createTool = createDef.id; + expect(listDef).toMatchObject({ + integration: slug, + owner: "org", + connection: "main", + }); + const schemaView = yield* client.tools.schema({ + query: { address: ToolAddress.make(createTool) }, + }); + expect(createDef.annotations).toEqual(schemaView.annotations); + expect(JSON.stringify(createDef.inputSchema)).toContain("text"); + const other = decodeToolSearch( + (yield* passthrough.call("search", { query: otherSlug })).raw, + ).structuredContent; + expect(other.items.some((tool) => tool.integration === otherSlug)).toBe(true); + + // --- A read call reaches the upstream with the connection's credential. --- + const listed = yield* passthrough.call("invoke", { tool: listTool, arguments: {} }); + expect(listed.ok, `the read call completes: ${listed.text}`).toBe(true); + expect(listed.text, "the upstream payload comes back").toContain("existing"); + const listReq = upstream.requests.find((r) => r.method === "GET"); + expect(listReq, "the GET reached the upstream").toBeDefined(); + expect(listReq?.authorization, "the connection's credential was applied").toBe( + `Bearer tok_${slug}`, + ); + + // --- An approval-gated call runs to completion: no pause, no resume. --- + const created = yield* passthrough.call("invoke", { + tool: createTool, + arguments: { body: { text: "hello" } }, + }); + expect(created.ok, `the gated call completes without a pause: ${created.text}`).toBe( + true, + ); + expect(created.text, "the call did not pause").not.toContain("Execution paused"); + expect(created.text, "the call did not ask for a resume").not.toContain("executionId"); + const createReq = upstream.requests.find((r) => r.method === "POST"); + expect(createReq, "the POST reached the upstream").toBeDefined(); + expect(createReq?.body, "the JSON body went over the wire").toContain('"text":"hello"'); + expect( + rawResultOf(created).structuredContent?.status, + "the result is a completed execution", + ).toBe("completed"); + + // --- Arguments are validated against the advertised schema. --- + // An MCP error result must not reach the upstream. + const invalid = yield* passthrough + .call("invoke", { tool: createTool, arguments: { body: { wrong: 1 } } }) + .pipe( + Effect.map((r) => r.ok), + Effect.catchCause(() => Effect.succeed(false)), + ); + expect(invalid, "a body missing its required field is refused").toBe(false); + expect( + upstream.requests.filter((r) => r.method === "POST").length, + "the invalid call never reached the upstream", + ).toBe(1); + + // --- `block` is enforced on the list AND the call. --- + const blockRule = yield* client.policies.create({ + payload: { owner: "org", pattern: `${slug}.*.*.*.createNote`, action: "block" }, + }); + yield* Effect.ensuring( + Effect.gen(function* () { + const afterBlock = mcp.session(identity, { mode: "passthrough" }); + const afterNames = decodeToolSearch( + (yield* afterBlock.call("search", { query: slug })).raw, + ).structuredContent.items.map((tool) => tool.id); + expect(afterNames, "a blocked tool is not listed").not.toContain(createTool); + expect(afterNames, "the unblocked sibling still is").toContain(listTool); + // A client that cached the old name cannot call it either: the + // executor refuses the call at invoke time, which passthrough + // renders as an MCP error result. + const stale = yield* passthrough.call("invoke", { + tool: createTool, + arguments: { body: { text: "again" } }, + }); + expect(stale.ok, "a blocked tool fails when called").toBe(false); + expect(stale.text, "the failure names the policy").toContain("blocked"); + expect( + upstream.requests.filter((r) => r.method === "POST").length, + "the blocked call never reached the upstream", + ).toBe(1); + }), + client.policies + .remove({ params: { policyId: blockRule.id }, payload: { owner: "org" } }) + .pipe(Effect.ignore), + ); + }), + Effect.gen(function* () { + for (const s of [slug, otherSlug]) { + yield* client.connections + .remove({ + params: { + owner: "org", + integration: IntegrationSlug.make(s), + name: ConnectionName.make("main"), + }, + }) + .pipe(Effect.ignore); + yield* client.openapi.removeSpec({ params: { slug: s } }).pipe(Effect.ignore); + } + }), + ); + }), + ), +); diff --git a/e2e/scenarios/support/search-invoke.ts b/e2e/scenarios/support/search-invoke.ts new file mode 100644 index 0000000000..4c956eadd8 --- /dev/null +++ b/e2e/scenarios/support/search-invoke.ts @@ -0,0 +1,25 @@ +import { Schema } from "effect"; + +/** Parse the public search result, including schemas and account identity. */ +export const decodeToolSearch = Schema.decodeUnknownSync( + Schema.Struct({ + structuredContent: Schema.Struct({ + items: Schema.Array( + Schema.Struct({ + id: Schema.String, + name: Schema.String, + integration: Schema.String, + owner: Schema.String, + connection: Schema.String, + inputSchema: Schema.Record(Schema.String, Schema.Unknown), + annotations: Schema.optional( + Schema.Struct({ requiresApproval: Schema.optional(Schema.Boolean) }), + ), + }), + ), + total: Schema.Number, + hasMore: Schema.Boolean, + nextOffset: Schema.NullOr(Schema.Number), + }), + }), +); diff --git a/e2e/src/surfaces/mcp.ts b/e2e/src/surfaces/mcp.ts index e1fa12c5f1..8e01f29258 100644 --- a/e2e/src/surfaces/mcp.ts +++ b/e2e/src/surfaces/mcp.ts @@ -100,6 +100,11 @@ export interface McpCallResult { export interface McpToolDef { readonly name: string; readonly description: string; + /** The advertised MCP annotations (`readOnlyHint`, `destructiveHint`, …), + * for scenarios that assert on what a harness's native approval reads. */ + readonly annotations?: Record; + /** The advertised input JSON Schema, verbatim. */ + readonly inputSchema?: unknown; } /** How a connection surfaces a paused (approval-gated) execution. `browser` is @@ -172,6 +177,9 @@ export interface McpSurface { * `search_` tools (`?search_tools=true`). Omitted means * the product default: none. */ readonly searchTools?: boolean; + /** `passthrough` serves search and invoke for the visible catalog + * (`?mode=passthrough`). Omitted means the product default: codemode. */ + readonly mode?: "codemode" | "passthrough"; readonly url?: string; }, ) => McpSession; @@ -315,6 +323,7 @@ export const makeMcpSurface = (target: Target, runDir?: string): McpSurface => ( ...(options?.elicitationMode ? [`elicitation_mode=${options.elicitationMode}`] : []), ...(options?.artifacts === false ? ["artifacts=false"] : []), ...(options?.searchTools === true ? ["search_tools=true"] : []), + ...(options?.mode === "passthrough" ? ["mode=passthrough"] : []), ].join("&"); const sessionUrl = sessionQuery ? `${mcpUrl}?${sessionQuery}` : mcpUrl; @@ -358,6 +367,10 @@ export const makeMcpSurface = (target: Target, runDir?: string): McpSurface => ( return listed.tools.map((tool) => ({ name: tool.name, description: tool.description ?? "", + ...(tool.annotations + ? { annotations: tool.annotations as Record } + : {}), + inputSchema: tool.inputSchema, })); }), call, @@ -408,14 +421,37 @@ export const makeMcpSurface = (target: Target, runDir?: string): McpSurface => ( return defs.map((tool: { name: string }) => tool.name); }); + // mcporter's `listTools` projects annotations away, so read the raw + // client it holds: the same connection (and cached OAuth), the full tool + // definition. const describeTools = () => Effect.promise(async (): Promise> => { - const defs = await (await runtime()).listTools(serverName, callOptions); - connected = true; - return defs.map((tool: { name: string; description?: string }) => ({ - name: tool.name, - description: tool.description ?? "", - })); + const rt = await runtime(); + if (!connected) { + await rt.listTools(serverName, callOptions); + connected = true; + } + const context = await rt.connect(serverName, { + allowCachedAuth: true, + oauthSessionOptions: callOptions.oauthSessionOptions, + }); + const out: McpToolDef[] = []; + let cursor: string | undefined; + do { + const page = await context.client.listTools(cursor ? { cursor } : undefined); + for (const tool of page.tools) { + out.push({ + name: tool.name, + description: tool.description ?? "", + ...(tool.annotations + ? { annotations: tool.annotations as Record } + : {}), + inputSchema: tool.inputSchema, + }); + } + cursor = page.nextCursor ?? undefined; + } while (cursor); + return out; }); const call = (name: string, args: Record = {}) => diff --git a/packages/core/api/src/server/mcp-build.ts b/packages/core/api/src/server/mcp-build.ts index 299aa853ec..e2d57acb54 100644 --- a/packages/core/api/src/server/mcp-build.ts +++ b/packages/core/api/src/server/mcp-build.ts @@ -71,6 +71,7 @@ export const makeMcpBuildServer = engine, artifacts: executor.artifacts, connections: executor.connections, + tools: executor.tools, ...(hostOptions?.loadAppShellHtml ? { loadAppShellHtml: hostOptions.loadAppShellHtml } : {}), @@ -88,6 +89,8 @@ export const makeMcpBuildServer = ...(options ?? {}), }).pipe( Effect.withSpan("mcp.server.create"), + // Catalog failures use the same retryable build envelope. + Effect.mapError((cause) => new McpEngineBuildError({ cause })), Effect.map((mcpServer) => ({ mcpServer, engine, executor })), ), ), diff --git a/packages/core/execution/src/tool-invoker.ts b/packages/core/execution/src/tool-invoker.ts index b8458929a4..a0c59a2cb8 100644 --- a/packages/core/execution/src/tool-invoker.ts +++ b/packages/core/execution/src/tool-invoker.ts @@ -672,7 +672,7 @@ const scoreToolMatch = (tool: SearchableTool, query: string): ToolDiscoveryResul /** What `tools.search()` calls inside the sandbox. */ export const searchTools = Effect.fn("executor.tools.search")(function* ( - executor: Executor, + executor: { readonly tools: Pick }, query: string, limit = 12, options?: { readonly namespace?: string; readonly offset?: number }, diff --git a/packages/core/sdk/src/elicitation.ts b/packages/core/sdk/src/elicitation.ts index 290349e3cd..5f7370649c 100644 --- a/packages/core/sdk/src/elicitation.ts +++ b/packages/core/sdk/src/elicitation.ts @@ -48,11 +48,21 @@ export const ElicitationResponse = Schema.Struct({ }); export type ElicitationResponse = typeof ElicitationResponse.Type; +/** Who raised an elicitation. `"policy"` is the executor's own approval gate + * (`enforceApproval`): a consent-only form whose terms are exactly "run this + * tool with these arguments". `"tool"` is anything the tool itself asked for + * mid-call, which may carry its own terms (a permanent site grant, a scope + * choice) even when the schema is empty. A host that has already obtained + * consent for the tool call may auto-accept the former and must never + * auto-accept the latter. Absent means unknown, which reads as `"tool"`. */ +export type ElicitationSource = "policy" | "tool"; + /** Handler input — the tool address being invoked, its args, and the request. */ export interface ElicitationContext { readonly address: ToolAddress; readonly args: unknown; readonly request: ElicitationRequest; + readonly source?: ElicitationSource; } /** Host-provided handler the SDK calls when a tool suspends for input. */ diff --git a/packages/core/sdk/src/executor.ts b/packages/core/sdk/src/executor.ts index 87208a3cc2..b01c62dbc7 100644 --- a/packages/core/sdk/src/executor.ts +++ b/packages/core/sdk/src/executor.ts @@ -6246,6 +6246,7 @@ export const createExecutor = ` toolkit), so the tool catalog is scoped to it. */ readonly resource: McpResource; @@ -151,6 +155,11 @@ interface SessionMetaBase { * {@link McpSessionInit}). Absent — including for sessions persisted before * the flag existed — means the default (disabled). */ readonly searchToolsEnabled?: boolean; + /** The tool surface (carried from {@link McpSessionInit}). Absent — + * including for sessions persisted before the field existed — means + * codemode. A cold restore MUST rebuild the same surface the client first + * saw, or its cached tool names stop resolving mid-conversation. */ + readonly toolMode?: McpToolMode; /** The MCP resource the session serves (carried from {@link McpSessionInit}); * `buildMcpServer` scopes the tool catalog to it. */ readonly resource: McpResource; diff --git a/packages/hosts/cloudflare/src/mcp/do-headers.ts b/packages/hosts/cloudflare/src/mcp/do-headers.ts index 02c4b196d0..6690331b1b 100644 --- a/packages/hosts/cloudflare/src/mcp/do-headers.ts +++ b/packages/hosts/cloudflare/src/mcp/do-headers.ts @@ -133,5 +133,7 @@ export { readArtifactsEnabled, readElicitationMode, readSearchToolsEnabled, + readToolMode, type McpElicitationMode, + type McpToolMode, } from "@executor-js/host-mcp/browser-approval"; diff --git a/packages/hosts/mcp/src/browser-approval.ts b/packages/hosts/mcp/src/browser-approval.ts index a524b89fbb..aa6ed5efab 100644 --- a/packages/hosts/mcp/src/browser-approval.ts +++ b/packages/hosts/mcp/src/browser-approval.ts @@ -79,6 +79,21 @@ export const readSearchToolsEnabled = (request: Request): boolean => { return TRUE_QUERY_VALUES.has(value.toLowerCase()); }; +export type McpToolMode = "codemode" | "passthrough"; + +/** + * Read the tool surface mode off an MCP request's `?mode=` query. The default, + * `codemode`, serves `execute` (the model writes sandboxed TypeScript against + * `tools.*`). `?mode=passthrough` instead serves search and invoke, with no + * `execute`, `skills`, or `resume`. Invoke is marked destructive for native + * client approval. Any other value + * reads as the default. + */ +export const readToolMode = (request: Request): McpToolMode => { + const value = new URL(request.url).searchParams.get("mode"); + return value === "passthrough" ? "passthrough" : "codemode"; +}; + /** * Build the console approval URL for a paused execution: * `//resume/?mcp_session_id=` diff --git a/packages/hosts/mcp/src/in-memory-session-store.ts b/packages/hosts/mcp/src/in-memory-session-store.ts index 47747e1a24..952f9bd83f 100644 --- a/packages/hosts/mcp/src/in-memory-session-store.ts +++ b/packages/hosts/mcp/src/in-memory-session-store.ts @@ -12,6 +12,8 @@ import { readArtifactsEnabled, readElicitationMode, readSearchToolsEnabled, + readToolMode, + type McpToolMode, } from "./browser-approval"; import { makeInProcessBrowserApprovalStore, @@ -31,7 +33,7 @@ import { type Principal, type McpResource, } from "./seams"; -import type { BrowserApprovalStore } from "./tool-server"; +import type { BrowserApprovalStore, McpPassthroughUnavailableError } from "./tool-server"; // --------------------------------------------------------------------------- // In-process McpSessionStore — the single-node serving store, shared by every @@ -117,13 +119,15 @@ export interface McpBuildServerOptions { /** Whether this session serves the per-integration `search_` * tools. False unless the client connected with `?search_tools=true`. */ readonly searchToolsEnabled?: boolean; + /** The tool surface (`?mode=`): codemode (default) or passthrough. */ + readonly mode?: McpToolMode; } /** Build the per-session `McpServer` + engine for a principal (the host's engine + tools). */ export type McpBuildServer = ( principal: Principal, options?: McpBuildServerOptions, -) => Effect.Effect; +) => Effect.Effect; export interface InMemoryMcpSessionStore { /** The `McpSessionStore` seam value to hand to `inMemoryMcpSessionsLayer`. */ @@ -414,13 +418,18 @@ export const makeInMemoryMcpSessionStore = ( ): McpBuildServerOptions => { const artifactsEnabled = readArtifactsEnabled(request); const searchToolsEnabled = readSearchToolsEnabled(request); + const toolMode = readToolMode(request); + const surface = { + artifactsEnabled, + searchToolsEnabled, + mode: toolMode, + }; const mode = readElicitationMode(request); if (mode !== "browser") { - return { artifactsEnabled, searchToolsEnabled, elicitationMode: { mode } }; + return { ...surface, elicitationMode: { mode } }; } return { - artifactsEnabled, - searchToolsEnabled, + ...surface, elicitationMode: { mode: "browser", // Prefer the pinned public origin; fall back to the request URL (correct @@ -489,9 +498,12 @@ export const makeInMemoryMcpSessionStore = ( }), ), // A build failure has nowhere typed to go in the envelope; render a 500. - Effect.catchTag("McpEngineBuildError", () => - Effect.succeed(jsonRpcError(500, -32603, "Internal server error")), - ), + Effect.catchTags({ + McpEngineBuildError: () => + Effect.succeed(jsonRpcError(500, -32603, "Internal server error")), + McpPassthroughUnavailableError: () => + Effect.succeed(jsonRpcError(500, -32603, "Internal server error")), + }), ); }; diff --git a/packages/hosts/mcp/src/passthrough-tools.test.ts b/packages/hosts/mcp/src/passthrough-tools.test.ts new file mode 100644 index 0000000000..fe6764be79 --- /dev/null +++ b/packages/hosts/mcp/src/passthrough-tools.test.ts @@ -0,0 +1,668 @@ +import { describe, expect, it } from "@effect/vitest"; +import { Effect, Schema } from "effect"; +import { Client } from "@modelcontextprotocol/sdk/client/index.js"; +import { InMemoryTransport } from "@modelcontextprotocol/sdk/inMemory.js"; +import type * as Cause from "effect/Cause"; + +import type { ExecutionEngine } from "@executor-js/execution"; +import { + ToolAddress, + IntegrationSlug, + ConnectionName, + ToolName, + type Tool, + type ToolSchemaView, +} from "@executor-js/sdk"; + +import { readToolMode } from "./browser-approval"; +import { passthroughCallCode } from "./passthrough-tools"; +import { + createExecutorMcpServer, + McpPassthroughUnavailableError, + type ExecutorMcpServerConfig, + type McpToolsPort, +} from "./tool-server"; + +// --------------------------------------------------------------------------- +// Fixtures +// --------------------------------------------------------------------------- + +const projection = (input: { + readonly integration: string; + readonly name: string; + readonly owner?: "org" | "user"; + readonly connection?: string; + readonly description?: string; + readonly inputSchema?: unknown; + readonly requiresApproval?: boolean; + readonly static?: boolean; +}): Tool => { + const owner = input.owner ?? "org"; + const connection = input.connection ?? "main"; + return { + address: ToolAddress.make(`tools.${input.integration}.${owner}.${connection}.${input.name}`), + integration: IntegrationSlug.make(input.integration), + owner, + connection: ConnectionName.make(connection), + name: ToolName.make(input.name), + pluginId: "test", + description: input.description ?? `${input.integration} ${input.name}`, + ...(input.inputSchema === undefined ? {} : { inputSchema: input.inputSchema }), + ...(input.requiresApproval === undefined + ? {} + : { annotations: { requiresApproval: input.requiresApproval } }), + ...(input.static === undefined ? {} : { static: input.static }), + }; +}; + +/** Exercise the existing list/schema seam and record which schemas were requested. */ +const toolPort = ( + catalog: readonly Tool[], + schemaReads: string[] = [], + lists: string[] = [], +): McpToolsPort => ({ + list: () => + Effect.sync(() => { + lists.push("list"); + return catalog; + }), + schema: (address) => + Effect.sync(() => { + schemaReads.push(String(address)); + const tool = catalog.find((item) => item.address === address); + if (!tool) return null; + return { + address, + name: tool.name, + description: tool.description, + inputSchema: tool.inputSchema, + annotations: tool.annotations, + } satisfies ToolSchemaView; + }), +}); + +/** A stub engine that records every executed code string and answers with a + * fixed value, so a test can prove a passthrough call became the expected + * single-call code and took `execute` (never `executeWithPause`). */ +const makeRecordingEngine = (result: unknown = { ok: true, data: { hello: "world" } }) => { + const executed: string[] = []; + let pausedCalls = 0; + const engine: ExecutionEngine = { + execute: (code) => + Effect.sync(() => { + executed.push(code); + return { result }; + }), + executeWithPause: () => + Effect.sync(() => { + pausedCalls += 1; + return { status: "completed" as const, result: { result } }; + }), + resume: () => Effect.succeed(null), + isExecutionSettled: undefined, + getPausedExecution: () => Effect.succeed(null), + pausedExecutionCount: () => Effect.succeed(0), + hasPausedExecutions: () => Effect.succeed(false), + getDescription: Effect.succeed("test executor"), + shutdown: Effect.void, + }; + return { engine, executed, pausedCalls: () => pausedCalls }; +}; + +const withClient = async ( + config: ExecutorMcpServerConfig, + fn: (client: Client) => Promise, +) => { + const mcpServer = await Effect.runPromise(createExecutorMcpServer(config)); + const [clientTransport, serverTransport] = InMemoryTransport.createLinkedPair(); + const client = new Client({ name: "test-client", version: "1.0.0" }, { capabilities: {} }); + await mcpServer.connect(serverTransport); + await client.connect(clientTransport); + // oxlint-disable-next-line executor/no-try-catch-or-throw -- boundary: test helper must close MCP transports after async client assertions + try { + await fn(client); + } finally { + await clientTransport.close(); + await serverTransport.close(); + } +}; + +const decodeJsonString = Schema.decodeUnknownSync(Schema.fromJsonString(Schema.String)); +const decodeJsonRecord = Schema.decodeUnknownSync( + Schema.fromJsonString(Schema.Record(Schema.String, Schema.Unknown)), +); + +const decodeSearchItems = Schema.decodeUnknownSync( + Schema.Struct({ items: Schema.Array(Schema.Struct({ id: Schema.String })) }), +); + +const CATALOG: readonly Tool[] = [ + projection({ + integration: "github", + name: "issues.create", + requiresApproval: true, + inputSchema: { + type: "object", + properties: { title: { type: "string" }, body: { $ref: "#/$defs/Body" } }, + required: ["title"], + $defs: { Body: { type: "string" } }, + }, + }), + projection({ integration: "github", name: "issues.list" }), + projection({ integration: "linear", name: "issueCreate", requiresApproval: true }), +]; + +describe("passthrough catalog", () => { + it("emits exactly one awaited tool call with the whole address as one string literal", () => { + expect(passthroughCallCode("tools.github.org.main.issues.create", { title: "hi" })).toBe( + 'return await tools["github.org.main.issues.create"]({"title":"hi"});', + ); + expect(passthroughCallCode("linear.org.main.issueCreate", undefined)).toBe( + 'return await tools["linear.org.main.issueCreate"]({});', + ); + // `then` is reserved by every sandbox proxy; as part of one key it is + // just text, so such a tool stays callable. + expect(passthroughCallCode("tools.svc.org.main.items.then", {})).toBe( + 'return await tools["svc.org.main.items.then"]({});', + ); + }); + + it("keeps a hostile tool segment as data, never as code", () => { + // An OpenAPI spec controls its tool paths (`x-executor-toolPath`), so a + // segment can contain anything. It must land inside a JSON string. + const hostile = 'x"](await tools.victim.org.main.destroy({}))["'; + const code = passthroughCallCode(`tools.evil.org.main.${hostile}`, {}); + // Structural proof the payload never escapes the string literal: the + // source is exactly `return await tools[]();`, and + // that one string decodes back to the raw address. + const shape = /^return await tools\[("(?:[^"\\]|\\.)*")\]\((\{.*\})\);$/s.exec(code); + expect(shape).not.toBeNull(); + expect(decodeJsonString(shape![1]!)).toBe(`evil.org.main.${hostile}`); + // And the call's argument is the JSON we passed, untouched by the address. + expect(decodeJsonRecord(shape![2]!)).toEqual({}); + }); +}); + +// --------------------------------------------------------------------------- +// Wire flags +// --------------------------------------------------------------------------- + +describe("readToolMode", () => { + const request = (query: string) => new Request(`https://example.test/mcp${query}`); + + it("defaults to codemode and only accepts the exact passthrough spelling", () => { + expect(readToolMode(request(""))).toBe("codemode"); + expect(readToolMode(request("?mode=passthrough"))).toBe("passthrough"); + expect(readToolMode(request("?mode=Passthrough"))).toBe("codemode"); + expect(readToolMode(request("?mode=direct"))).toBe("codemode"); + }); +}); + +// --------------------------------------------------------------------------- +// Server: the served surface +// --------------------------------------------------------------------------- + +describe("passthrough mode server", () => { + it("serves exactly search and invoke even for 10000 tools", async () => { + const { engine, executed } = makeRecordingEngine(); + const schemaReads: string[] = []; + const lists: string[] = []; + const catalog = Array.from({ length: 10000 }, (_, i) => + projection({ + integration: "bench", + name: `record${i}`, + description: i === 9999 ? "cobalt orchard sentinel" : `benchmark record ${i}`, + }), + ); + await withClient( + { + engine, + mode: "passthrough", + searchToolsEnabled: true, + tools: toolPort(catalog, schemaReads, lists), + }, + async (client) => { + const listed = await client.listTools(); + expect(listed.tools.map((tool) => tool.name).sort()).toEqual(["invoke", "search"]); + expect(JSON.stringify(listed).length).toBeLessThan(4000); + expect(lists).toEqual([]); + expect(schemaReads).toEqual([]); + const result = await client.callTool({ + name: "search", + arguments: { query: "cobalt orchard sentinel" }, + }); + expect(result.structuredContent).toMatchObject({ + total: 1, + hasMore: false, + nextOffset: null, + items: [{ id: "tools.bench.org.main.record9999" }], + }); + expect(schemaReads).toEqual(["tools.bench.org.main.record9999"]); + const found = decodeSearchItems(result.structuredContent).items[0]; + expect(found).toBeDefined(); + const invoked = await client.callTool({ + name: "invoke", + arguments: { tool: found?.id, arguments: {} }, + }); + expect(invoked.isError ?? false).toBe(false); + expect(executed).toEqual(['return await tools["bench.org.main.record9999"]({});']); + const page = await client.callTool({ + name: "search", + arguments: { query: "bench", limit: 2 }, + }); + expect(page.structuredContent).toMatchObject({ + total: 10000, + hasMore: true, + nextOffset: 2, + }); + expect(decodeSearchItems(page.structuredContent).items).toHaveLength(2); + expect(schemaReads).toHaveLength(4); + const next = await client.callTool({ + name: "search", + arguments: { query: "bench", limit: 2, offset: 2 }, + }); + expect(next.structuredContent).toMatchObject({ total: 10000, nextOffset: 4 }); + expect(next.structuredContent).not.toEqual(page.structuredContent); + const missing = await client.callTool({ + name: "search", + arguments: { query: "nonexistent quasar" }, + }); + expect(missing.structuredContent).toEqual({ + items: [], + total: 0, + hasMore: false, + nextOffset: null, + }); + }, + ); + }); + + it("uses current schemas and excludes static configuration tools", async () => { + const recording = makeRecordingEngine(); + const dynamic = projection({ + integration: "notes", + name: "create", + inputSchema: { type: "object" }, + }); + const catalog = [dynamic, projection({ integration: "settings", name: "erase", static: true })]; + let current: ToolSchemaView = { address: dynamic.address, inputSchema: dynamic.inputSchema }; + const tools: McpToolsPort = { ...toolPort(catalog), schema: () => Effect.succeed(current) }; + await withClient({ engine: recording.engine, mode: "passthrough", tools }, async (client) => { + const hidden = await client.callTool({ name: "search", arguments: { query: "settings" } }); + expect(decodeSearchItems(hidden.structuredContent).items).toEqual([]); + const staticCall = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.settings.org.main.erase", arguments: {} }, + }); + expect(staticCall.isError).toBe(true); + await client.callTool({ name: "search", arguments: { query: "notes" } }); + current = { + address: dynamic.address, + inputSchema: { + type: "object", + properties: { title: { type: "string" } }, + required: ["title"], + }, + }; + const invalid = await client.callTool({ + name: "invoke", + arguments: { tool: String(dynamic.address), arguments: {} }, + }); + expect(invalid.isError).toBe(true); + expect(recording.executed).toEqual([]); + const refreshed = await client.callTool({ name: "search", arguments: { query: "notes" } }); + expect(refreshed.structuredContent).toMatchObject({ + items: [{ inputSchema: { required: ["title"] } }], + }); + }); + }); + + it("returns schemas and account details from search and marks invoke destructive", async () => { + const { engine } = makeRecordingEngine(); + await withClient({ engine, mode: "passthrough", tools: toolPort(CATALOG) }, async (client) => { + const listed = await client.listTools(); + expect(listed.tools.find((tool) => tool.name === "search")?.annotations).toMatchObject({ + readOnlyHint: true, + destructiveHint: false, + }); + expect(listed.tools.find((tool) => tool.name === "invoke")?.annotations).toMatchObject({ + readOnlyHint: false, + destructiveHint: true, + }); + const result = await client.callTool({ + name: "search", + arguments: { query: "github issues create", limit: 1 }, + }); + expect(result.structuredContent).toMatchObject({ + items: [ + { + id: "tools.github.org.main.issues.create", + owner: "org", + connection: "main", + annotations: { requiresApproval: true }, + inputSchema: { + type: "object", + properties: { title: { type: "string" }, body: { $ref: "#/$defs/Body" } }, + required: ["title"], + $defs: { Body: { type: "string" } }, + }, + }, + ], + }); + }); + }); + + it("runs a call as one execute of single-call code, never a pause", async () => { + const recording = makeRecordingEngine({ ok: true, data: { number: 7 } }); + await withClient( + { + engine: recording.engine, + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const result = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.issues.create", arguments: { title: "hello" } }, + }); + expect(recording.executed).toEqual([ + 'return await tools["github.org.main.issues.create"]({"title":"hello"});', + ]); + expect(recording.pausedCalls()).toBe(0); + // The tool's `data` is the result: nothing sits between the tool and + // the client to unwrap the `{ ok, data }` envelope for it. + expect(result.isError ?? false).toBe(false); + expect(result.structuredContent).toEqual({ + status: "completed", + result: { number: 7 }, + logs: [], + }); + }, + ); + }); + + it("surfaces an expected tool failure as an MCP error result", async () => { + const recording = makeRecordingEngine({ + ok: false, + error: { code: "tool_blocked", message: "Tool blocked by policy: github.org.main.x" }, + }); + await withClient( + { + engine: recording.engine, + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const result = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.issues.list", arguments: {} }, + }); + expect(result.isError).toBe(true); + expect(result.structuredContent).toEqual({ + status: "error", + error: { code: "tool_blocked", message: "Tool blocked by policy: github.org.main.x" }, + logs: [], + }); + const text = (result.content as Array<{ type: string; text?: string }>)[0]?.text ?? ""; + expect(text).toContain("tool_blocked"); + }, + ); + }); + + /** An engine whose tool raises the given elicitations in order and records + * each answer. `source` is what the executor stamps: `policy` for its own + * approval gate, `tool` for anything the tool asked for itself. */ + const elicitingEngine = ( + requests: ReadonlyArray<{ readonly source: "policy" | "tool"; readonly request: any }>, + seen: string[], + ): ExecutionEngine => ({ + ...makeRecordingEngine().engine, + execute: (_code, options) => + Effect.gen(function* () { + for (const { source, request } of requests) { + const answer = yield* options.onElicitation({ + address: CATALOG[0]!.address, + args: {}, + request, + source, + }); + seen.push(`${source}:${answer.action}`); + if (answer.action !== "accept") { + return { result: { ok: false, error: { code: "declined", message: "declined" } } }; + } + } + return { result: { ok: true, data: null } }; + }), + }); + + const approvalGate = { + _tag: "FormElicitation" as const, + message: "Approve github.org.main.issues.create?", + requestedSchema: { type: "object", properties: {} }, + }; + + it("accepts the executor's own approval gate inline", async () => { + const seen: string[] = []; + await withClient( + { + engine: elicitingEngine([{ source: "policy", request: approvalGate }], seen), + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const result = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.issues.create", arguments: { title: "x" } }, + }); + expect(seen).toEqual(["policy:accept"]); + expect(result.isError ?? false).toBe(false); + }, + ); + }); + + it("keeps every invoke marked destructive even when the selected tool has no approval annotation", async () => { + const seen: string[] = []; + await withClient( + { + engine: elicitingEngine([{ source: "policy", request: approvalGate }], seen), + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const invoke = (await client.listTools()).tools.find((tool) => tool.name === "invoke"); + expect(invoke?.annotations).toMatchObject({ readOnlyHint: false, destructiveHint: true }); + const result = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.issues.list", arguments: {} }, + }); + expect(seen).toEqual(["policy:accept"]); + expect(result.isError ?? false).toBe(false); + }, + ); + }); + + it("never auto-accepts a tool-raised prompt, even one with an empty schema", async () => { + // Same wire shape as the approval gate, but raised by the TOOL: a + // per-site grant whose terms live in `meta`. Provenance, not shape, + // decides. With no elicitation capability on the client, the call fails + // and says so — it is not silently granted. + const seen: string[] = []; + const siteGrant = { + _tag: "FormElicitation" as const, + message: "Allow Browser use to access example.com?", + requestedSchema: {}, + meta: { persist: "always", origin: "https://example.com" }, + }; + await withClient( + { + engine: elicitingEngine([{ source: "tool", request: siteGrant }], seen), + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const result = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.issues.create", arguments: { title: "x" } }, + }); + expect(seen).toEqual(["tool:decline"]); + expect(result.isError).toBe(true); + expect(result.structuredContent).toMatchObject({ + status: "error", + error: { code: "elicitation_unsupported", request: siteGrant.message }, + }); + }, + ); + }); + + it("reports an unanswerable URL request with the URL, not as a user decline", async () => { + const seen: string[] = []; + const reconnect = { + _tag: "UrlElicitation" as const, + message: "Reconnect GitHub", + url: "https://example.test/oauth/start", + elicitationId: "elic_1", + }; + await withClient( + { + engine: elicitingEngine([{ source: "tool", request: reconnect }], seen), + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const result = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.issues.create", arguments: { title: "x" } }, + }); + expect(seen).toEqual(["tool:decline"]); + expect(result.isError).toBe(true); + const text = (result.content as Array<{ text?: string }>)[0]?.text ?? ""; + expect(text).toContain("does not support elicitation"); + expect(text).toContain("https://example.test/oauth/start"); + expect(text).not.toContain("declined by the user"); + expect(result.structuredContent).toMatchObject({ + error: { code: "elicitation_unsupported", url: reconnect.url }, + }); + }, + ); + }); + + it("serves no artifact tools in passthrough even when artifacts are requested", async () => { + const { engine } = makeRecordingEngine(); + await withClient( + { + engine, + mode: "passthrough", + artifactsEnabled: true, + loadAppShellHtml: async () => "", + artifacts: { + list: () => Effect.succeed([]), + get: () => Effect.die("unused"), + save: () => Effect.die("unused"), + }, + tools: toolPort(CATALOG), + }, + async (client) => { + const names = (await client.listTools()).tools.map((tool) => tool.name); + expect(names.sort()).toEqual(["invoke", "search"]); + }, + ); + }); + + it("rejects arguments that fail the advertised schema before running anything", async () => { + const recording = makeRecordingEngine(); + await withClient( + { + engine: recording.engine, + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const result = await client.callTool({ + name: "invoke", + arguments: { + tool: "tools.github.org.main.issues.create", + arguments: { body: "no title" }, + }, + }); + expect(result.isError).toBe(true); + expect(JSON.stringify(result.content)).toContain("Invalid arguments"); + expect(recording.executed).toEqual([]); + }, + ); + }); + + it("answers an unknown tool name with a not-found error", async () => { + const { engine } = makeRecordingEngine(); + await withClient( + { + engine, + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + const result = await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.nope", arguments: {} }, + }); + expect(result.isError).toBe(true); + expect(JSON.stringify(result.content)).toContain("not found"); + }, + ); + }); + + it("rejects empty searches, oversized pages and malformed invoke inputs", async () => { + const recording = makeRecordingEngine(); + await withClient( + { + engine: recording.engine, + mode: "passthrough", + tools: toolPort(CATALOG), + }, + async (client) => { + for (const arguments_ of [ + { query: " " }, + { query: "github", limit: 21 }, + { query: "github", offset: -1 }, + ]) { + expect((await client.callTool({ name: "search", arguments: arguments_ })).isError).toBe( + true, + ); + } + expect( + ( + await client.callTool({ + name: "invoke", + arguments: { tool: "tools.github.org.main.issues.create", arguments: "{}" }, + }) + ).isError, + ).toBe(true); + expect(recording.executed).toEqual([]); + }, + ); + }); + + it("leaves codemode untouched when the mode is absent", async () => { + const { engine } = makeRecordingEngine(); + await withClient( + { + engine, + description: "Execute TypeScript in a sandboxed runtime.", + tools: toolPort(CATALOG), + }, + async (client) => { + const names = (await client.listTools()).tools.map((tool) => tool.name); + expect(names).toContain("execute"); + expect(names).toContain("skills"); + expect(names).not.toContain("github__issues_create"); + }, + ); + }); + + it("fails the build, not the session, when the host provides no catalog", async () => { + const { engine } = makeRecordingEngine(); + const outcome = await Effect.runPromise( + createExecutorMcpServer({ engine, mode: "passthrough" }).pipe(Effect.flip), + ); + expect(outcome).toBeInstanceOf(McpPassthroughUnavailableError); + }); +}); diff --git a/packages/hosts/mcp/src/passthrough-tools.ts b/packages/hosts/mcp/src/passthrough-tools.ts new file mode 100644 index 0000000000..c79d59861f --- /dev/null +++ b/packages/hosts/mcp/src/passthrough-tools.ts @@ -0,0 +1,29 @@ +/** + * The sandbox code a passthrough call runs. Built HERE from the session's + * resolved address and a JSON-encoded argument — never concatenated from raw + * model input — and shaped exactly like the artifact `execute-action` grammar + * (`return await tools.()`), so it takes the same engine path as + * every other execution: billing, rate limits, shape memory and analytics all + * see it as one execution. + */ +export const passthroughCallCode = (address: string, args: unknown): string => { + // The whole dotted address is ONE JSON string literal in bracket notation: + // `tools["github.org.main.items.then"](...)`. Two reasons it is not a chain + // of property accesses. The tool segment is customer-controlled (an OpenAPI + // spec may set `x-executor-toolPath`), so it must be data in the generated + // source, never syntax. And every sandbox proxy reserves the property name + // `then` (a thenable check would otherwise await the proxy itself), so a + // per-segment chain could never reach a tool whose path contains `then`. + // Each proxy joins the accessed keys with `.` to form the dispatch path, so + // a single key holding the dotted address reassembles to exactly the same + // path the chain would have. + const bare = address.startsWith("tools.") ? address.slice("tools.".length) : address; + return `return await tools[${JSON.stringify(bare)}](${JSON.stringify(args ?? {})});`; +}; + +/** Describe the fixed search/invoke surface without listing the underlying catalog. */ +export const passthroughInstructions = (): string => + "Find connected integration tools with search, then call invoke with the returned tool ID and JSON arguments. " + + "Search returns input schemas and account details. Use its nextOffset to get more matches. " + + "Invoke can change external state; your client handles approval for each call. Workspace block policies remain enforced. " + + "No JavaScript, execute, resume, or artifact tools are exposed in this mode."; diff --git a/packages/hosts/mcp/src/tool-server.ts b/packages/hosts/mcp/src/tool-server.ts index 730c5bfd89..f611cf9d41 100644 --- a/packages/hosts/mcp/src/tool-server.ts +++ b/packages/hosts/mcp/src/tool-server.ts @@ -1,3 +1,4 @@ +import { reattachDefs } from "@executor-js/sdk/host-internal"; import { Data, Duration, Effect, Match, Option, Predicate, Result, Schema } from "effect"; import * as Cause from "effect/Cause"; import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; @@ -22,7 +23,10 @@ import * as z from "zod/v4"; import { CurrentOrgWriteAccess, + ToolAddress, + parseToolAddress, isToolFile, + isToolResult, makeOrgWriteAccessState, sanitizeArtifactPreviewMarkup, type OrgWriteAccess, @@ -37,10 +41,13 @@ import type { ElicitationRequest, SaveArtifactInput, ToolFileValue, + Executor, + ToolSchemaView, } from "@executor-js/sdk"; import type * as Tracer from "effect/Tracer"; import { createExecutionEngine, + searchTools, formatExecuteResult, formatPausedExecution, formatTtlDuration, @@ -74,6 +81,8 @@ import { type BindableConnection, } from "./artifact-bindings"; import { MCP_ORG_WRITE_ACCESS_HEADER } from "./seams"; +import { passthroughCallCode, passthroughInstructions } from "./passthrough-tools"; +import type { McpToolMode } from "./browser-approval"; // --------------------------------------------------------------------------- // Workers-compatible JSON Schema validator (replaces Ajv which uses new Function()) @@ -180,9 +189,26 @@ type SharedMcpServerConfig = { * the `execute` description lists). The tools exist to carry the namespaces * into the model's context as tool names; each call routes through the same * execution flow as `tools.search({ namespace })` inside `execute`, so the - * results match what code-side search returns. + * results match what code-side search returns. Codemode only: passthrough + * ignores it (it has its own search tool, and codemode search + * results point at an `execute` tool passthrough does not serve). */ readonly searchToolsEnabled?: boolean; + /** + * The tool surface this connection serves. `codemode` (the default) is the + * `execute` tool plus `skills`/`resume` and the artifact surface. + * `passthrough` (`?mode=passthrough`) serves search and invoke for the + * visible catalog, with no execute, skills, resume, or artifact tools. + * Invoke is marked destructive for client approval. Requires `tools`. + */ + readonly mode?: McpToolMode; + /** + * The scoped executor's tool catalog, for passthrough mode. Structurally + * satisfied by `executor.tools`. Hosts that never serve passthrough may + * leave it unset; a passthrough session without it fails at build time + * rather than silently serving an empty surface. + */ + readonly tools?: McpToolsPort; /** * Renders an artifact once, server-side, before it is saved — so a component * that throws on its first render is refused at create time with the real @@ -272,6 +298,15 @@ export type McpConnectionsPort = { readonly list: () => Effect.Effect; }; +/** The same list and schema APIs used by codemode discovery. */ +export type McpToolsPort = Pick; + +/** A passthrough session was requested but the host gave the factory no + * catalog to serve. A configuration defect, not a runtime condition. */ +export class McpPassthroughUnavailableError extends Data.TaggedError( + "McpPassthroughUnavailableError", +)<{ readonly reason: string }> {} + export type ExecutorMcpServerConfig = | (ExecutionEngineConfig & SharedMcpServerConfig) | ({ readonly engine: ExecutionEngine } & SharedMcpServerConfig) @@ -638,6 +673,64 @@ const toMcpResult = (result: FormattedExecuteInput): McpToolResult => { }; }; +/** + * A passthrough call's result IS the tool's `ToolResult`. Inside `execute` + * the model reads `{ ok, data | error }` and branches; here nothing runs + * between the tool and the client, so an expected failure (`ok: false` — a + * 4xx wall, a blocked policy, a validation miss) has to be an MCP error + * result, and a success unwraps to the tool's `data`. Everything else + * (sandbox error, emitted output) keeps the codemode rendering. + */ +const toPassthroughResult = (outcome: FormattedExecuteInput): McpToolResult => { + const value = outcome.result; + if (outcome.error || !isToolResult(value)) return toMcpResult(outcome); + if (value.ok) { + return toMcpResult({ ...outcome, result: value.data }); + } + const message = `${value.error.code}: ${value.error.message}`; + return { + content: [{ type: "text", text: `Error: ${message}` }], + structuredContent: { + status: "error", + error: value.error, + logs: outcome.logs ?? [], + }, + isError: true, + }; +}; + +/** + * A passthrough tool asked the user for something and the connected client + * advertises no elicitation capability, so nobody could answer. Say exactly + * that, and carry the request — a reconnect/OAuth URL is the usual content — + * so the model can relay it and the user can act outside the client. + */ +const elicitationUnsupportedResult = ( + toolName: string, + request: ElicitationRequest, +): McpToolResult => { + const url = elicitationRequestUrl(request); + const lines = [ + `Tool ${toolName} needs input from the user, but this MCP client does not support elicitation, so the call could not complete.`, + `Request: ${request.message}`, + ...(url ? [`Open this URL to continue, then retry the call: ${url}`] : []), + ]; + return { + content: [{ type: "text", text: `Error: ${lines.join("\n")}` }], + structuredContent: { + status: "error", + error: { + code: "elicitation_unsupported", + message: lines[0]!, + request: request.message, + ...(url ? { url } : {}), + }, + logs: [], + }, + isError: true, + }; +}; + const toMcpPausedResult = (formatted: ReturnType): McpToolResult => ({ content: [{ type: "text", text: formatted.text }], structuredContent: formatted.structured, @@ -1104,13 +1197,174 @@ const parseJsonContent = (raw: string): Record | undefined => { return Option.isSome(parsed) ? parsed.value : undefined; }; +// --------------------------------------------------------------------------- +// Passthrough surface +// --------------------------------------------------------------------------- + +/** Serialize one existing schema view as a self-contained MCP input schema. */ +const passthroughInputSchema = (view: ToolSchemaView): unknown => + reattachDefs( + view.inputSchema ?? { type: "object", properties: {} }, + new Map(Object.entries(view.schemaDefinitions ?? {})), + ); + +/** Register discovery over the existing APIs, with no catalog work at connection time. */ +const registerPassthroughTools = ( + server: McpServer, + tools: McpToolsPort, + run: ( + address: ToolAddress, + args: unknown, + extra: McpRequestJoinKeys, + ) => Effect.Effect, +): Effect.Effect => + Effect.gen(function* () { + const context = yield* Effect.context(); + const validator = new CfWorkerJsonSchemaValidator(); + const discovery = { + tools: { + list: (filter?: Parameters[0]) => + tools + .list(filter) + .pipe(Effect.map((items) => items.filter((tool) => tool.static !== true))), + }, + }; + const boundary = ( + effect: Effect.Effect, + extra: McpRequestJoinKeys, + ) => + Effect.runPromiseWith(context)( + effect.pipe( + Effect.provideService( + CurrentOrgWriteAccess, + makeOrgWriteAccessState(requestOrgWriteAccess(extra)), + ), + Effect.catchCause((cause) => Effect.succeed(toMcpFailureResult(cause))), + ), + ); + yield* Effect.sync(() => { + server.registerTool( + "search", + { + description: + "Search connected integration tools by action, integration, or account. Returns matching tool IDs, account details, and full JSON input schemas. Pass the returned ID and arguments to invoke. Use nextOffset to page through matches.", + inputSchema: { + query: z + .string() + .trim() + .min(1) + .max(500) + .describe("Keywords describing the tool or task, such as github create issue."), + limit: z.number().int().min(1).max(20).default(10), + offset: z.number().int().min(0).default(0), + }, + annotations: { readOnlyHint: true, destructiveHint: false, openWorldHint: false }, + }, + ({ query, limit, offset }, extra) => + boundary( + Effect.gen(function* () { + const page = yield* searchTools(discovery, query, limit, { offset }); + const candidates = yield* Effect.forEach( + page.items, + (match) => + Effect.gen(function* () { + const address = ToolAddress.make(`tools.${match.path}`); + const identity = parseToolAddress(String(address)); + if (!identity) return null; + const schema = yield* tools.schema(address); + // Visibility can change between listing and schema lookup. + if (!schema) return null; + return { + id: String(address), + name: match.name, + integration: identity.integration, + owner: identity.owner, + connection: identity.connection, + description: match.description, + inputSchema: passthroughInputSchema(schema), + ...(schema.annotations ? { annotations: schema.annotations } : {}), + }; + }), + { concurrency: 4 }, + ); + const result = { ...page, items: candidates.filter(Predicate.isNotNull) }; + return { + content: [{ type: "text" as const, text: JSON.stringify(result) }], + structuredContent: result, + }; + }), + extra, + ), + ); + server.registerTool( + "invoke", + { + description: + "Call one connected integration tool using the exact ID and JSON input schema returned by search. May read or change external state. Your client handles approval for this call; workspace blocks remain enforced.", + inputSchema: { + tool: z.string().min(1).describe("Exact tool ID returned by search."), + arguments: z + .record(z.string(), z.unknown()) + .describe("Tool arguments matching the inputSchema returned by search."), + }, + annotations: { readOnlyHint: false, destructiveHint: true, openWorldHint: true }, + }, + ({ tool: id, arguments: args }, extra) => + boundary( + Effect.gen(function* () { + const identity = parseToolAddress(id); + const unavailable = { + isError: true, + content: [ + { + type: "text" as const, + text: "Tool not found or blocked by policy. Search for an available tool.", + }, + ], + }; + if (!identity) return unavailable; + const address = ToolAddress.make(id); + // Use the existing visibility filter and exclude static configuration tools. + const visible = yield* discovery.tools.list({ + integration: identity.integration, + owner: identity.owner, + connection: identity.connection, + query: String(identity.tool), + includeAnnotations: false, + }); + if (!visible.some((tool) => tool.address === address)) return unavailable; + const schema = yield* tools.schema(address); + if (!schema) return unavailable; + // The SDK validator checks this dynamic JSON schema at the MCP boundary. + const validate = validator.getValidator( + passthroughInputSchema(schema) as JsonSchemaType, + ); + const checked = validate(args); + if (!checked.valid) + return { + isError: true, + content: [ + { + type: "text" as const, + text: `Invalid arguments for tool ${id}: ${checked.errorMessage ?? "invalid"}`, + }, + ], + }; + return yield* run(address, checked.data, extra); + }), + extra, + ), + ); + }); + }).pipe(Effect.withSpan("mcp.host.register_search_invoke")); + // --------------------------------------------------------------------------- // Server factory // --------------------------------------------------------------------------- export const createExecutorMcpServer = ( config: ExecutorMcpServerConfig, -): Effect.Effect => +): Effect.Effect => Effect.gen(function* () { const engine = "engine" in config ? config.engine : createExecutionEngine(config); const description = @@ -1122,11 +1376,24 @@ export const createExecutorMcpServer = ( // Artifacts are on unless this connection opted out (`?artifacts=false`). // One flag decides the whole surface: the tools, the shell resource, and // the skills catalog below. - const artifactsEnabled = config.artifactsEnabled ?? true; + // Search/invoke serves no artifact tools: artifacts run sandboxed code. + const artifactsEnabled = + config.mode === "passthrough" ? false : (config.artifactsEnabled ?? true); const skillCatalog: readonly Skill[] = skillCatalogFor({ artifacts: artifactsEnabled }); // Per-integration search tools are off unless this connection opted in // (`?search_tools=true`). const searchToolsEnabled = config.searchToolsEnabled ?? false; + // Passthrough (`?mode=passthrough`) replaces the codemode surface + // wholesale. The flag is read once here and every codemode-only + // registration below is gated on it, so the two surfaces cannot leak into + // each other. + const mode: McpToolMode = config.mode ?? "codemode"; + const passthrough = mode === "passthrough"; + if (passthrough && !config.tools) { + return yield* new McpPassthroughUnavailableError({ + reason: "passthrough mode requires tool list and schema APIs", + }); + } // Captured at construction time. SDK callbacks fire later (often // deferred past the outer Effect's await), so we use the runtime to @@ -1227,6 +1494,11 @@ export const createExecutorMcpServer = ( // per host. capabilities: { resources: {}, tools: {} }, jsonSchemaValidator: new CfWorkerJsonSchemaValidator(), + ...(passthrough + ? { + instructions: passthroughInstructions(), + } + : {}), }, ), ).pipe(Effect.withSpan("mcp.host.create_server")); @@ -1545,103 +1817,171 @@ export const createExecutorMcpServer = ( Effect.annotateSpans(joinKeyAttributes(extra)), ); + // --- passthrough call path --- + // + // Invoke runs one generated call through the existing execution engine. + // The client approves the generic destructive tool; upstream prompts use + // native elicitation, or fail with an actionable result when unsupported. + const executePassthroughCall = ( + address: ToolAddress, + args: unknown, + extra: McpRequestJoinKeys, + ): Effect.Effect => + Effect.gen(function* () { + yield* startMarker("mcp.host.tool.execute.start", { + "mcp.tool.id": String(address), + "mcp.tool.mode": "passthrough", + "executor.tool.address": address, + }); + const { url: supportsUrl } = getElicitationSupport(server); + const native = makeMcpElicitationHandler(server, extra.requestId, debugLog); + const { form: supportsForm } = getElicitationSupport(server); + // Set when the tool asked the user for something this client cannot + // relay. The handler has no error channel (a non-accept is a decline + // to the executor), so the request is kept here and the whole call is + // reported as unanswerable below — with what was asked, URL included — + // instead of as "declined by the user", which nobody did. + let unanswerable: ElicitationRequest | undefined; + const onElicitation: ElicitationHandler = (ctx) => { + // Every invoke is advertised as destructive, so the client's native + // approval covers the selected ID and arguments, even if policy changed. + // Tool-raised prompts still require their own response below. + if (ctx.source === "policy") { + return Effect.succeed({ action: "accept" as const, content: {} }); + } + // Anything the tool itself asked for goes to the client natively + // when it can take it; the native bridge already turns a URL + // request into a form for form-only clients. + if (supportsForm || (supportsUrl && Predicate.isTagged(ctx.request, "UrlElicitation"))) { + return native(ctx); + } + unanswerable = ctx.request; + return Effect.succeed({ action: "decline" as const }); + }; + const outcome = yield* engine.execute(passthroughCallCode(address, args), { + onElicitation, + }); + if (unanswerable) return elicitationUnsupportedResult(String(address), unanswerable); + return toPassthroughResult(outcome); + }).pipe( + Effect.withSpan("mcp.host.tool.execute", { + attributes: { + "mcp.tool.id": String(address), + "mcp.tool.mode": "passthrough", + "executor.integration": parseToolAddress(String(address))?.integration, + }, + }), + Effect.annotateSpans(joinKeyAttributes(extra)), + ); + // --- tools --- - yield* Effect.sync(() => - server.registerTool( - "execute", - { - description, - inputSchema: { code: z.string().trim().min(1) }, - }, - ({ code }, extra) => runToolEffect(executeCode(code, extra), extra), - ), - ).pipe( - Effect.withSpan("mcp.host.register_tool", { - attributes: { "mcp.tool.name": "execute" }, - }), - ); + // Passthrough serves search and invoke in place of the codemode tools. + if (passthrough && config.tools) { + yield* registerPassthroughTools(server, config.tools, executePassthroughCall); + } - yield* Effect.sync(() => - server.registerTool( - "skills", - { - description: [ - "Documentation for THIS server's own tools. Not a general skill reader: it serves a short, fixed set of how-to docs about using `execute` and artifacts here, and it cannot reach your harness's skills, a SKILL.md on disk, or any user- or project-authored skill. The argument is a name from its own catalog, never a path or an outside skill's id.", - "These docs hold the long-form guidance that would otherwise bloat another tool's always-loaded description.", - 'Call `skills({ name: "execute" })` for the full guide to writing code for the `execute` tool (search the catalog, call tools, emit results, resume paused runs).', - "Call with no name to list the few docs available.", - ].join("\n"), - inputSchema: { - name: z - .string() - .optional() - .describe( - 'A doc from this server\'s own catalog, e.g. "execute" — not a path or an outside skill name. Omit to list the catalog.', - ), + if (!passthrough) + yield* Effect.sync(() => + server.registerTool( + "execute", + { + description, + inputSchema: { code: z.string().trim().min(1) }, }, - }, - ({ name }, extra) => - runToolEffect(Effect.succeed(skillsResult(name, executeInventory, skillCatalog)), extra), - ), - ).pipe( - Effect.withSpan("mcp.host.register_tool", { - attributes: { "mcp.tool.name": "skills" }, - }), - ); - - yield* Effect.sync(() => { - if (elicitationMode.mode === "native") { - return undefined; - } + ({ code }, extra) => runToolEffect(executeCode(code, extra), extra), + ), + ).pipe( + Effect.withSpan("mcp.host.register_tool", { + attributes: { "mcp.tool.name": "execute" }, + }), + ); - if (elicitationMode.mode === "model") { - return server.registerTool( - "resume", + if (!passthrough) + yield* Effect.sync(() => + server.registerTool( + "skills", { description: [ - "Resume a paused execution using the executionId returned by execute.", - "This connection explicitly allows model-side resume via elicitation_mode=model.", + "Documentation for THIS server's own tools. Not a general skill reader: it serves a short, fixed set of how-to docs about using `execute` and artifacts here, and it cannot reach your harness's skills, a SKILL.md on disk, or any user- or project-authored skill. The argument is a name from its own catalog, never a path or an outside skill's id.", + "These docs hold the long-form guidance that would otherwise bloat another tool's always-loaded description.", + 'Call `skills({ name: "execute" })` for the full guide to writing code for the `execute` tool (search the catalog, call tools, emit results, resume paused runs).', + "Call with no name to list the few docs available.", ].join("\n"), inputSchema: { - executionId: z.string().describe("The execution ID from the paused result"), - action: z - .enum(["accept", "decline", "cancel"]) - .describe("How to respond to the interaction"), - content: z + name: z .string() - .describe("Optional JSON-encoded response content for form elicitations") - .default("{}"), + .optional() + .describe( + 'A doc from this server\'s own catalog, e.g. "execute" — not a path or an outside skill name. Omit to list the catalog.', + ), }, }, - ({ executionId, action, content: rawContent }, extra) => + ({ name }, extra) => runToolEffect( - resumeExecution(executionId, action, parseJsonContent(rawContent), extra), + Effect.succeed(skillsResult(name, executeInventory, skillCatalog)), extra, ), - ); - } + ), + ).pipe( + Effect.withSpan("mcp.host.register_tool", { + attributes: { "mcp.tool.name": "skills" }, + }), + ); - return server.registerTool( - "resume", - { - description: [ - "Request user approval to resume a paused execution.", - "Call this with the executionId returned by execute. If the user has not approved in the browser yet, tell them to open the returned approval URL. If they have approved, this returns the resumed execution result.", - "This connection does not allow the model to choose accept, decline, cancel, or content.", - ].join("\n"), - inputSchema: { - executionId: z.string().describe("The execution ID from the paused result"), + if (!passthrough) + yield* Effect.sync(() => { + if (elicitationMode.mode === "native") { + return undefined; + } + + if (elicitationMode.mode === "model") { + return server.registerTool( + "resume", + { + description: [ + "Resume a paused execution using the executionId returned by execute.", + "This connection explicitly allows model-side resume via elicitation_mode=model.", + ].join("\n"), + inputSchema: { + executionId: z.string().describe("The execution ID from the paused result"), + action: z + .enum(["accept", "decline", "cancel"]) + .describe("How to respond to the interaction"), + content: z + .string() + .describe("Optional JSON-encoded response content for form elicitations") + .default("{}"), + }, + }, + ({ executionId, action, content: rawContent }, extra) => + runToolEffect( + resumeExecution(executionId, action, parseJsonContent(rawContent), extra), + extra, + ), + ); + } + + return server.registerTool( + "resume", + { + description: [ + "Request user approval to resume a paused execution.", + "Call this with the executionId returned by execute. If the user has not approved in the browser yet, tell them to open the returned approval URL. If they have approved, this returns the resumed execution result.", + "This connection does not allow the model to choose accept, decline, cancel, or content.", + ].join("\n"), + inputSchema: { + executionId: z.string().describe("The execution ID from the paused result"), + }, }, - }, - ({ executionId }, extra) => - runToolEffect(resumeAfterBrowserApproval(executionId, extra), extra), + ({ executionId }, extra) => + runToolEffect(resumeAfterBrowserApproval(executionId, extra), extra), + ); + }).pipe( + Effect.withSpan("mcp.host.register_tool", { + attributes: { "mcp.tool.name": "resume" }, + }), ); - }).pipe( - Effect.withSpan("mcp.host.register_tool", { - attributes: { "mcp.tool.name": "resume" }, - }), - ); // --- per-integration search tools (opt-in, `?search_tools=true`) --- // @@ -1659,7 +1999,7 @@ export const createExecutorMcpServer = ( // would only repeat the name) and a single bare `query` parameter — no // paging knobs, because anything past the first page belongs in `execute`. // `namespace-search-tools.test.ts` pins the serialized size. - if (searchToolsEnabled) { + if (searchToolsEnabled && !passthrough) { // The MCP tool-name grammar ([A-Za-z0-9_-]). Integration slugs already // conform (they are `tools.` property names in sandbox code); one // that somehow doesn't is skipped rather than failing the whole session. diff --git a/packages/react/src/api/analytics.tsx b/packages/react/src/api/analytics.tsx index 465e7b88c3..c637b59bd6 100644 --- a/packages/react/src/api/analytics.tsx +++ b/packages/react/src/api/analytics.tsx @@ -169,6 +169,7 @@ export interface AnalyticsEvents { mcp_install_elicitation_mode_changed: { elicitation_mode: string }; mcp_install_artifacts_toggled: { artifacts: boolean }; mcp_install_search_tools_toggled: { search_tools: boolean }; + mcp_install_tool_mode_changed: { tool_mode: "codemode" | "passthrough" }; // ── Command palette ────────────────────────────────────────────────────── command_palette_navigated: { diff --git a/packages/react/src/components/mcp-install-card.tsx b/packages/react/src/components/mcp-install-card.tsx index cbefa82644..e77a3bef73 100644 --- a/packages/react/src/components/mcp-install-card.tsx +++ b/packages/react/src/components/mcp-install-card.tsx @@ -26,6 +26,9 @@ const McpInstallPreferencesSchema = Schema.Struct({ httpElicitationMode: Schema.Literals(["browser", "model", "native"]), artifacts: Schema.Boolean, searchTools: Schema.Boolean, + // Added after v1 shipped; optional so a stored preference from before it + // existed still decodes and takes the default (codemode). + toolMode: Schema.optional(Schema.Literals(["codemode", "passthrough"])), }); type McpInstallPreferences = typeof McpInstallPreferencesSchema.Type; @@ -52,6 +55,7 @@ const DEFAULT_MCP_INSTALL_PREFERENCES: McpInstallPreferences = { httpElicitationMode: "model", artifacts: true, searchTools: false, + toolMode: "codemode", }; const decodeMcpInstallPreferences = Schema.decodeUnknownOption( Schema.fromJsonString(McpInstallPreferencesSchema), @@ -116,6 +120,9 @@ export const buildMcpHttpEndpoint = (input: { /** Per-integration search tools are off by default, so only the opt-in is * spelled out on the URL (`&search_tools=true`). */ readonly searchTools?: boolean; + /** Codemode is the default, so only passthrough is spelled out on the URL + * (`&mode=passthrough`). */ + readonly toolMode?: "codemode" | "passthrough"; // Cloud only: pins the URL to `//mcp` (the server also accepts the // legacy `//mcp` form). Desktop/local pass nothing and get the bare // `/mcp` path. @@ -137,7 +144,12 @@ export const buildMcpHttpEndpoint = (input: { params.push(["elicitation_mode", input.elicitationMode]); } if (input.artifacts === false) params.push(["artifacts", "false"]); - if (input.searchTools === true) params.push(["search_tools", "true"]); + // Per-integration search tools and artifacts are codemode affordances; passthrough serves + // neither, so the URL never claims them alongside it. + if (input.searchTools === true && input.toolMode !== "passthrough") { + params.push(["search_tools", "true"]); + } + if (input.toolMode === "passthrough") params.push(["mode", "passthrough"]); if (params.length === 0) return endpoint; const query = params.map(([key, value]) => `${key}=${value}`).join("&"); @@ -160,6 +172,7 @@ export const buildMcpInstallCommand = (input: { readonly elicitationMode?: McpElicitationMode; readonly artifacts?: boolean; readonly searchTools?: boolean; + readonly toolMode?: "codemode" | "passthrough"; readonly devCliCwd?: string; readonly organizationSlug?: string | null; }): string => { @@ -170,6 +183,7 @@ export const buildMcpInstallCommand = (input: { elicitationMode: input.elicitationMode, artifacts: input.artifacts, searchTools: input.searchTools, + toolMode: input.toolMode, organizationSlug: input.organizationSlug, }); const headerFlags: string[] = []; @@ -197,9 +211,12 @@ export const buildMcpInstallCommand = (input: { if (input.artifacts === false) { innerArgs.push("--no-artifacts"); } - if (input.searchTools === true) { + if (input.searchTools === true && input.toolMode !== "passthrough") { innerArgs.push("--search-tools"); } + if (input.toolMode === "passthrough") { + innerArgs.push("--mode", "passthrough"); + } return `npx add-mcp ${shellQuoteWord(innerArgs.map(shellQuoteWord).join(" "))} --name executor`; }; @@ -227,6 +244,7 @@ export function McpInstallCard(props: { className?: string }) { } const [advancedOpen, setAdvancedOpen] = useState(false); const { mode, httpElicitationMode, artifacts, searchTools } = preferences; + const toolMode = preferences.toolMode ?? "codemode"; useEffect(() => { writeMcpInstallPreferences(storageKey, preferences); @@ -268,6 +286,7 @@ export function McpInstallCard(props: { className?: string }) { elicitationMode, artifacts, searchTools, + toolMode, devCliCwd, organizationSlug, }); @@ -290,16 +309,38 @@ export function McpInstallCard(props: { className?: string }) {
+
+
Search and invoke
+
+ {toolMode === "passthrough" + ? "Find connected tools with search, then call them with invoke. Your client handles approval for each call." + : "Disabled: agents write code against your tools through one execute tool."} +
+
+ { + const nextMode = next ? "passthrough" : "codemode"; + setPreferences((current) => ({ ...current, toolMode: nextMode })); + trackEvent("mcp_install_tool_mode_changed", { tool_mode: nextMode }); + }} + aria-label="Search and invoke" + /> +
+
Artifacts
- {artifacts - ? "Generated UI components are saved to your workspace." - : "Disabled: this connection serves no artifact tools."} + {toolMode === "passthrough" + ? "Not available in search and invoke mode." + : artifacts + ? "Generated UI components are saved to your workspace." + : "Disabled: this connection serves no artifact tools."}
{ setPreferences((current) => ({ ...current, artifacts: next })); trackEvent("mcp_install_artifacts_toggled", { artifacts: next }); @@ -311,13 +352,16 @@ export function McpInstallCard(props: { className?: string }) {
Integration search tools
- {searchTools - ? "One search tool per connected integration, so agents see your integrations as tool names." - : "Disabled: agents discover tools through search inside execute."} + {toolMode === "passthrough" + ? "This mode already provides one search tool for all integrations." + : searchTools + ? "One search tool per connected integration, so agents see your integrations as tool names." + : "Disabled: agents discover tools through search inside execute."}
{ setPreferences((current) => ({ ...current, searchTools: next })); trackEvent("mcp_install_search_tools_toggled", { search_tools: next }); From 9876c2232cf22a59de40c81497c5c822f9582d60 Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Mon, 14 Sep 2026 18:26:07 -0700 Subject: [PATCH 2/2] Preserve full annotations in MCP scenario decoding --- e2e/scenarios/support/search-invoke.ts | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/e2e/scenarios/support/search-invoke.ts b/e2e/scenarios/support/search-invoke.ts index 4c956eadd8..0de674d34f 100644 --- a/e2e/scenarios/support/search-invoke.ts +++ b/e2e/scenarios/support/search-invoke.ts @@ -1,4 +1,5 @@ import { Schema } from "effect"; +import { ToolAnnotationsView } from "@executor-js/sdk"; /** Parse the public search result, including schemas and account identity. */ export const decodeToolSearch = Schema.decodeUnknownSync( @@ -12,9 +13,7 @@ export const decodeToolSearch = Schema.decodeUnknownSync( owner: Schema.String, connection: Schema.String, inputSchema: Schema.Record(Schema.String, Schema.Unknown), - annotations: Schema.optional( - Schema.Struct({ requiresApproval: Schema.optional(Schema.Boolean) }), - ), + annotations: Schema.optional(ToolAnnotationsView), }), ), total: Schema.Number,