Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346import assert from "node:assert/strict";import { test } from "node:test";import { mkdtemp, rm } from "node:fs/promises";import { tmpdir } from "node:os";import { join } from "node:path";import { build } from "esbuild";import { Miniflare, convertV4MiniflareOptions } from "miniflare";
test( "Effect MCP operations use native discovery and remove disabled connections", { timeout: 60000 }, async () => { const native = await build({ entryPoints: ["tests/fixtures/mcp-runtime-worker.ts"], alias: { path: "node:path" }, bundle: true, write: false, format: "esm", platform: "neutral", conditions: ["workerd", "worker", "browser"], mainFields: ["module", "main"], external: ["cloudflare:*", "node:*"], }); const requests = []; let rejectDiscovery = false; let changedCatalog = false; let ambiguousCatalog = false; const persistence = await mkdtemp(join(tmpdir(), "flarebot-mcp-runtime-")); let releaseDiscovery; let enteredDiscovery; let entered = new Promise((resolve) => { enteredDiscovery = resolve; }); let gate = new Promise((resolve) => { releaseDiscovery = resolve; }); const options = { ...convertV4MiniflareOptions({ modules: true, script: native.outputFiles[0].text, compatibilityDate: "2026-09-04", compatibilityFlags: ["nodejs_compat"], durableObjects: { MODEL: { className: "McpRuntimeFixture", useSQLite: true }, }, async outboundService(request) { if (request.method !== "POST") return new Response(null, { status: 405 }); const body = await request.json(); requests.push(body.method); if ( new URL(request.url).hostname === "slow.example.com" && body.method === "initialize" ) { enteredDiscovery(); await gate; } if (rejectDiscovery && body.method === "tools/list") return Response.json({ jsonrpc: "2.0", id: body.id, error: { code: -32603, message: "Discovery unavailable" }, }); const results = { initialize: { protocolVersion: "2025-06-18", capabilities: { tools: {}, resources: {}, prompts: {} }, serverInfo: { name: "Fixture MCP", version: "1" }, }, "tools/list": { tools: [ ...(ambiguousCatalog ? [ { name: "search", description: "x".repeat(4001), inputSchema: { type: "object" }, }, ] : []), { name: "search", description: "Search records", inputSchema: { type: "object", properties: { query: { type: changedCatalog ? "number" : "string" }, }, }, }, ...(changedCatalog ? [ { name: "added", inputSchema: { type: "object" } }, { name: "malformed", inputSchema: { type: "object", properties: { value: { type: "nonsense" } }, }, }, ] : []), ], }, "resources/list": { resources: [ { name: "Guide", uri: "docs://guide", mimeType: "text/plain" }, ], }, "resources/templates/list": { resourceTemplates: [] }, "prompts/list": { prompts: [ { name: "summarize", description: "Summarize records" }, ], }, }; if (!("id" in body)) return new Response(null, { status: 202 }); return Response.json({ jsonrpc: "2.0", id: body.id, ...(results[body.method] ? { result: results[body.method] } : { error: { code: -32601, message: "Method not found" } }), }); }, }), resourcePersistencePath: persistence, }; let worker = new Miniflare(options); const call = async (action, args = {}) => { const response = await worker.dispatchFetch("http://fixture", { method: "POST", body: JSON.stringify({ action, ...args }), }); assert.equal(response.status, 200); const result = await response.json(); assert.equal(result.ok, true, JSON.stringify(result)); return result.value; }; const waitFor = async (predicate) => { for (let i = 0; i < 100; i++) { const snapshot = await call("list"); if (predicate(snapshot)) return snapshot; await new Promise((resolve) => setTimeout(resolve, 50)); } assert.fail("MCP state did not settle"); }; const settled = async (item) => ( await waitFor((snapshot) => snapshot.connections.some( (row) => row.id === item.id && row.state !== "connecting", ), ) ).connections.find((row) => row.id === item.id); try { let item = await call("add", { value: { name: "Fixture", endpoint: "https://mcp.example.com/mcp", authMode: "none", }, }); assert.equal(item.state, "connecting"); item = await settled(item); assert.equal(item.state, "ready", JSON.stringify(item)); assert.deepEqual( item.capabilities.map((capability) => capability.kind).sort(), ["prompt", "resource", "tool"], ); assert.ok(requests.includes("tools/list")); const originalTool = item.capabilities.find( (capability) => capability.kind === "tool", ); assert.equal(originalTool.serverId, item.id); assert.equal(originalTool.inputSchema.properties.query.type, "string"); changedCatalog = true; item = await call("connect", item); item = await settled(item); assert.equal( item.state, "ready", "malformed nested schema does not break sibling discovery", ); assert.deepEqual(item.catalogRefresh.changed, [originalTool.id]); assert.equal(item.catalogRefresh.added.length, 2); assert.equal(item.catalogRefresh.invalid, 1); assert.equal( item.capabilities.find((capability) => capability.name === "malformed") .issue, "invalid_schema", ); changedCatalog = false; item = await call("connect", item); item = await settled(item); assert.equal(item.catalogRefresh.removed.length, 2); assert.deepEqual(item.catalogRefresh.changed, [originalTool.id]); ambiguousCatalog = true; item = await call("connect", item); item = await settled(item); assert.equal(item.state, "ready"); assert.ok( !item.capabilities.some((capability) => capability.name === "search"), ); assert.equal(item.catalogRefresh.omitted, 2); ambiguousCatalog = false; item = await call("connect", item); item = await settled(item); assert.deepEqual((await call("list")).native, [item.id]); item = await call("enable", { ...item, value: false }); assert.equal(item.state, "disabled"); await waitFor((snapshot) => snapshot.native.length === 0); item = await call("enable", { ...item, value: true }); item = await settled(item); assert.equal(item.state, "ready"); await worker.dispose(); worker = new Miniflare(options); item = await settled(item); assert.equal( item.state, "ready", "startup reconciles native restoration into settings", ); await call("remove", item); await waitFor((snapshot) => snapshot.native.length === 0); assert.deepEqual(await call("list"), { connections: [], native: [] }); const beforeSlowAdd = Date.now(); const slow = await call("add", { value: { name: "Slow", endpoint: "https://slow.example.com/mcp", authMode: "none", }, }); assert.equal(slow.state, "connecting"); assert.ok( Date.now() - beforeSlowAdd < 2000, "acknowledgement does not wait for discovery", ); await entered; for (const authMode of ["oauth", "headers"]) { const count = requests.length; let pending = await call("add", { value: { name: authMode, endpoint: "https://auth.example.com/mcp", authMode, }, }); pending = await settled(pending); assert.equal(pending.state, "authenticating"); assert.equal(pending.lastError, "authentication_required"); assert.equal(requests.length, count); } const disabledSlow = await call("enable", { ...slow, value: false }); assert.equal( disabledSlow.state, "disabled", "disable acknowledges while discovery is held", ); releaseDiscovery(); await waitFor((snapshot) => !snapshot.native.includes(slow.id)); assert.equal( (await call("list")).connections.find((row) => row.id === slow.id) .state, "disabled", "late discovery cannot re-enable the row", ); let restoredSlow = await call("enable", { ...disabledSlow, value: true }); restoredSlow = await settled(restoredSlow); let fast = await call("add", { value: { name: "Fast", endpoint: "https://fast.example.com/mcp", authMode: "none", }, }); fast = await settled(fast); await worker.dispose(); entered = new Promise((resolve) => { enteredDiscovery = resolve; }); gate = new Promise((resolve) => { releaseDiscovery = resolve; }); worker = new Miniflare(options); await call("list"); await entered; fast = await settled(fast); assert.equal( fast.state, "ready", "healthy restoration never waits for another endpoint", ); const restored = await call("list"); assert.equal( restored.connections.find((row) => row.id === restoredSlow.id).state, "connecting", ); assert.ok( restored.connections .filter((row) => row.authMode !== "none") .every((row) => row.state === "authenticating"), ); fast = await call("enable", { ...fast, value: false }); await waitFor((snapshot) => !snapshot.native.includes(fast.id)); releaseDiscovery(); assert.equal((await settled(restoredSlow)).state, "ready"); await worker.dispose(); rejectDiscovery = true; worker = new Miniflare(options); const failedDiscovery = await settled(restoredSlow); assert.equal(failedDiscovery.state, "error"); // Restored manager errors lose their cause. Retry once through the // instrumented application transport to establish an actionable category. assert.equal(failedDiscovery.lastError, "connection_failed"); assert.equal(failedDiscovery.health.failures, 1); await ( await worker.dispatchFetch( `https://fixture.example/__age?id=${failedDiscovery.id}`, ) ).arrayBuffer(); await ( await worker.dispatchFetch("https://fixture.example/__tick") ).arrayBuffer(); const classified = await waitFor((snapshot) => snapshot.connections.some( (row) => row.id === failedDiscovery.id && row.lastError === "discovery_failed", ), ); assert.equal( classified.connections.find((row) => row.id === failedDiscovery.id) .health.retryAt, null, ); assert.deepEqual(failedDiscovery.capabilities, restoredSlow.capabilities); } finally { releaseDiscovery(); await worker.dispose(); await rm(persistence, { recursive: true, force: true }); } },);