Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398import 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 }); } },);