import assert from "node:assert/strict"; import { test } from "node:test"; import { build } from "esbuild"; import { Miniflare, convertV4MiniflareOptions } from "miniflare"; import { mkdtemp, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; async function until(read, predicate, description) { const deadline = Date.now() + 10_000; do { const value = await read(); if (predicate(value)) return value; await new Promise((resolve) => setTimeout(resolve, 10)); } while (Date.now() < deadline); assert.fail(`Timed out: ${description}`); } test( "Think MCP actions use native parent RPC, current owner authority, bounded results and cancellation", { timeout: 90_000 }, async () => { const bundle = await build({ entryPoints: ["tests/fixtures/mcp-invocation-worker.ts"], alias: { path: "node:path", crypto: "node:crypto", async_hooks: "node:async_hooks", }, target: "es2022", bundle: true, write: false, format: "esm", platform: "neutral", conditions: ["workerd", "worker", "browser"], mainFields: ["module", "main"], external: ["cloudflare:*", "node:*"], }); const persistence = await mkdtemp( join(tmpdir(), "flarebot-mcp-invocation-"), ); const calls = []; const cancellations = []; const held = new Map(); let changed = false; const secret = "fixture-provider-secret-do-not-expose"; const options = { ...convertV4MiniflareOptions({ modules: true, script: bundle.outputFiles[0].text, compatibilityDate: "2026-09-04", compatibilityFlags: ["nodejs_compat"], durableObjects: { MODEL: { className: "McpInvocationFixture", useSQLite: true }, }, async outboundService(request) { if (request.method !== "POST") return new Response(null, { status: 405 }); const body = await request.json(); if (body.method === "notifications/cancelled") { cancellations.push(body.params.requestId); held.get(body.params.requestId)?.(); } if (!("id" in body)) return new Response(null, { status: 202 }); const reply = (result) => Response.json({ jsonrpc: "2.0", id: body.id, result }); if (body.method === "initialize") return reply({ protocolVersion: "2025-06-18", capabilities: { tools: {} }, serverInfo: { name: "Notes", version: "1" }, }); if (body.method === "tools/list") return reply({ tools: [ { name: "send_note", description: "Send a note", inputSchema: { type: "object", properties: { note: { type: changed ? "number" : "string" }, }, required: ["note"], additionalProperties: false, }, }, { name: "disabled_tool", inputSchema: { type: "object" } }, ], }); if (body.method === "ping") return reply({}); if (body.method !== "tools/call") return Response.json({ jsonrpc: "2.0", id: body.id, error: { code: -32601, message: "Unknown method" }, }); calls.push(body); const note = body.params.arguments.note; if (note === "protocol-error") return Response.json({ jsonrpc: "2.0", id: body.id, error: { code: -32603, message: secret }, }); if (note === "tool-error") return reply({ isError: true, content: [{ type: "text", text: secret }], }); if (note === "large") return reply({ content: [{ type: "text", text: "x".repeat(70_000) }], }); if (note === "slow") await new Promise((resolve) => held.set(body.id, resolve)); return reply({ content: [{ type: "text", text: `Delivered ${note}` }], }); }, }), resourcePersistencePath: persistence, }; const worker = new Miniflare(options); const request = async (path, data) => { const response = await worker.dispatchFetch( `https://fixture.example/${path}`, { method: "POST", body: JSON.stringify(data) }, ); assert.equal(response.status, 200); const result = await response.json(); assert.ok(!result.fixtureError, JSON.stringify(result)); return result; }; const runtime = async (action, args = {}) => { const result = await request("runtime", { action, ...args }); assert.equal(result.ok, true, JSON.stringify(result)); return result.value; }; const call = (action, value, name = "conversation") => request("control", { action, value, name }); const connection = async () => (await runtime("list")).connections[0]; const ready = () => until(connection, (item) => item?.state === "ready", "native MCP ready"); const choose = async (enabled) => { const item = await connection(); const tool = item.capabilities.find((tool) => tool.name === "send_note"); return runtime("tools", { id: item.id, value: { settingsRevision: item.settingsRevision, enabled, tools: [{ id: tool.id, fingerprint: tool.fingerprint }], }, }); }; const reference = (definition) => ({ source: definition.metadata.source, id: definition.id, fingerprint: definition.fingerprint, }); const direct = (ref, input, more = {}) => call("direct", { reference: ref, input, ...more }, "direct"); const pause = async (name, note = name) => { const count = calls.length; await call("submit", { note }, name); const pending = await call("pending", undefined, name); assert.equal(pending.length, 1); assert.deepEqual(pending[0].arguments, { note }); assert.equal( calls.length, count, "Ask pauses before any outbound tools/call", ); return pending[0]; }; const decide = (name, pending, decision = "approve") => call("decide", { executionId: pending.executionId, decision }, name); try { await runtime("add", { value: { name: "Fixture notes", endpoint: "https://notes.example.com/mcp", authMode: "none", }, }); await ready(); assert.deepEqual(await call("definitions"), []); await choose(true); const definitions = await call("definitions"); assert.deepEqual( definitions.map((definition) => definition.metadata.name), ["send_note"], ); const ref = reference(definitions[0]); const approved = await pause("approved"); assert.match(approved.summary, /send_note.*Fixture notes/); assert.deepEqual(await decide("approved", approved), { resolved: true }); assert.equal(calls.length, 1); assert.equal( calls[0].params.name, "send_note", "native manager removes only the server namespace", ); assert.deepEqual(calls[0].params.arguments, { note: "approved" }); let snapshot = await call("snapshot", undefined, "approved"); assert.equal( snapshot.activities.find( (item) => item.toolCallId === approved.toolCallId, )?.status, "succeeded", ); assert.ok( JSON.stringify(snapshot.messages).includes("Delivered approved"), ); assert.deepEqual( snapshot.activities.find( (item) => item.toolCallId === approved.toolCallId, ).capability, { id: ref.id, fingerprint: ref.fingerprint, name: "send_note", source: { ...ref.source, name: "Fixture notes" }, }, ); assert.deepEqual(await decide("approved", approved), { resolved: false }); assert.equal(calls.length, 1); const denied = await pause("denied"); await decide("denied", denied, "deny"); assert.equal(calls.length, 1); snapshot = await call("snapshot", undefined, "denied"); assert.equal( snapshot.activities.find( (item) => item.toolCallId === denied.toolCallId, )?.reason, "approval-denied", ); const disabled = await pause("disabled"); await choose(false); await decide("disabled", disabled); assert.equal(calls.length, 1); assert.equal( (await direct(ref, { note: "disabled" }, { approvedOnce: true })).result .error.code, "unavailable", ); await choose(true); const never = await pause("never"); await call("policy", "never"); assert.deepEqual(await call("definitions"), []); await decide("never", never); assert.equal(calls.length, 1); assert.equal( (await direct(ref, { note: "never" }, { approvedOnce: true })).result .error.code, "never", ); await call("policy", "ask"); assert.equal( (await direct(ref, { note: "no-grant" })).result.error.code, "ask", ); assert.equal(calls.length, 1); assert.equal( ( await direct( ref, { note: "pre-cancel" }, { approvedOnce: true, mode: "pre-cancel" }, ) ).result.error.code, "cancelled", ); assert.equal(calls.length, 1); const once = await direct( ref, { note: "once" }, { approvedOnce: true, mode: "duplicate" }, ); assert.equal(once.duplicate.error.code, "already_started"); assert.equal(calls.length, 2); assert.equal( (await direct(ref, { note: 123 }, { approvedOnce: true })).result.error .code, "invalid_arguments", ); assert.equal(calls.length, 2); assert.equal( ( await direct( ref, { note: "x".repeat(65_536) }, { approvedOnce: true }, ) ).result.error.code, "arguments_too_large", ); assert.equal(calls.length, 2); for (const [note, code] of [ ["protocol-error", "mcp_failed"], ["tool-error", "tool_failed"], ["large", "result_too_large"], ]) { const result = await direct(ref, { note }, { approvedOnce: true }); assert.equal(result.result.error.code, code); assert.ok(!JSON.stringify(result).includes(secret)); assert.ok(JSON.stringify(result).length < 1000); } for (const [note, code] of [ ["protocol-error", "mcp_failed"], ["tool-error", "tool_failed"], ["large", "result_too_large"], ]) { const name = `think-${note}`; const pending = await pause(name, note); await decide(name, pending); const failed = await call("snapshot", undefined, name); assert.equal( failed.activities.find( (item) => item.toolCallId === pending.toolCallId, )?.status, "failed", ); assert.ok(JSON.stringify(failed.messages).includes(code)); assert.ok(!JSON.stringify(failed).includes(secret)); assert.ok(JSON.stringify(failed).length < 10_000); } await call("policy", "allow"); const beforeCancellation = cancellations.length; const active = call("submit", { note: "slow" }, "cancelled"); await until( async () => calls.at(-1), (item) => item?.params.arguments.note === "slow", "outbound MCP call started", ); await call("cancel", undefined, "cancelled"); await active; await until( async () => cancellations, (items) => items.length > beforeCancellation && items.includes(calls.at(-1).id), "native MCP cancellation notification", ); snapshot = await call("snapshot", undefined, "cancelled"); assert.ok( snapshot.activities.some((item) => item.status === "cancelled"), ); await call("policy", "ask"); const stale = await pause("stale"); const beforeChange = calls.length; changed = true; const item = await connection(); await runtime("connect", item); await ready(); await choose(true); assert.notEqual( (await call("definitions"))[0].fingerprint, ref.fingerprint, ); await decide("stale", stale); assert.equal(calls.length, beforeChange); assert.equal( (await direct(ref, { note: "stale" }, { approvedOnce: true })).result .error.code, "unavailable", ); const latestDefinition = (await call("definitions"))[0]; const liveConnection = await connection(); await runtime("enable", { id: liveConnection.id, revision: liveConnection.revision, value: false, }); assert.deepEqual(await call("definitions"), []); assert.equal( ( await direct( reference(latestDefinition), { note: 123 }, { approvedOnce: true }, ) ).result.error.code, "unavailable", ); assert.equal(calls.length, beforeChange); } finally { for (const release of held.values()) release(); await worker.dispose(); await rm(persistence, { recursive: true, force: true }); } }, );