feat: expose Fusion custom tools to ACP agents via loopback MCP bridge (#3476)
## Summary - ACP runtimes can expose Fusion custom tools (`fn_*`) to external agents such as Hermes ACP and Prime. - When the engine passes `customTools`, `AcpRuntimeAdapter` starts a per-session loopback tool bridge and registers it as a stdio MCP server in `session/new.mcpServers`. - The bridge uses a per-session bearer token, binds only to loopback, preserves the MCP request ID as `toolCallId`, propagates tool `isError` results, and exposes only runnable non-built-in tools. - The MCP shim supports `initialize`, `ping`, `tools/list`, and `tools/call`, reporting transport, authentication, malformed-request, and unknown-tool failures correctly. - Bridge startup failures degrade gracefully: the ACP session remains usable without custom tools and exposes a fixed `fusionToolBridgeError` reason code for engine auditing. - Startup, request, and teardown paths clean temporary schemas and listeners; disposal aborts cooperative tools, bounds non-cooperative drains, prevents post-dispose execution, and preserves the existing synchronous session contract through an awaitable `disposePromise`. - The packaged CLI stages `mcp-schema-server.cjs` beside the ACP runtime distribution. ## Verification Validated at source commit `2fa6fd2e26923ac6241d1b9379b3819f1832c079`: - `pnpm --filter @fusion-plugin-examples/acp-runtime test` — 21 files / 251 tests passed - `pnpm --filter @fusion-plugin-examples/acp-runtime typecheck` — passed - Plugin build with `.cjs` asset copy — passed - `pnpm check:changesets` — passed - `git diff --check` — passed - MCP shim smoke test — passed - Greptile Review — passed - CodeRabbit — passed - Devin review — 0 bugs, 6 analyses, `errored_tasks=[]`, quality `4/10` ## Scope ACP runtime plugin only. The Hermes CLI runtime remains untouched; Hermes ACP is enabled by selecting the generic ACP runtime with `acpBinaryPath: hermes` and `acpArgs: ["acp"]`. --------- Co-authored-by: gsxdsm <gsxdsm@users.noreply.github.com>
This commit is contained in:
7
.changeset/hermes-acp-custom-tools-bridge.md
Normal file
7
.changeset/hermes-acp-custom-tools-bridge.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": minor
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: ACP runtimes can now expose Fusion custom tools (fn_*) to external agents such as Hermes ACP and Prime.
|
||||||
|
category: feature
|
||||||
|
dev: AcpRuntimeAdapter starts a per-session loopback tool bridge and registers it as a stdio MCP server in session/new.mcpServers when the engine passes customTools; the bridge authenticates requests with a per-session bearer token, threads the real MCP request id as the toolCallId, and is disposed on session/new failure and session teardown. Build copies mcp-schema-server.cjs beside dist (tsc does not copy .cjs assets).
|
||||||
@@ -33,6 +33,9 @@ const RUNTIME_PLUGINS_WITH_MCP_SCHEMA_SERVER = new Set([
|
|||||||
"fusion-plugin-claude-runtime",
|
"fusion-plugin-claude-runtime",
|
||||||
// FNXC:OmpAcp 2026-07-14-00:05: OMP ACP ships the same bridge asset for fn_* tools.
|
// FNXC:OmpAcp 2026-07-14-00:05: OMP ACP ships the same bridge asset for fn_* tools.
|
||||||
"fusion-plugin-omp-runtime",
|
"fusion-plugin-omp-runtime",
|
||||||
|
// FNXC:AcpCustomTools 2026-08-16-00:30: generic ACP runtime ships the same
|
||||||
|
// bridge asset for Hermes ACP / Prime fn_* tool forwarding.
|
||||||
|
"fusion-plugin-acp-runtime",
|
||||||
]);
|
]);
|
||||||
|
|
||||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||||
|
|||||||
@@ -21,7 +21,7 @@
|
|||||||
},
|
},
|
||||||
"private": true,
|
"private": true,
|
||||||
"scripts": {
|
"scripts": {
|
||||||
"build": "tsc",
|
"build": "tsc && node -e \"require('node:fs').copyFileSync('src/mcp-schema-server.cjs','dist/mcp-schema-server.cjs')\"",
|
||||||
"test": "vitest run --silent=passed-only --reporter=dot",
|
"test": "vitest run --silent=passed-only --reporter=dot",
|
||||||
"typecheck": "tsc --noEmit"
|
"typecheck": "tsc --noEmit"
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -1,8 +1,10 @@
|
|||||||
import { describe, it, expect, afterEach } from "vitest";
|
import { describe, it, expect, afterEach, vi } from "vitest";
|
||||||
import os from "node:os";
|
import os from "node:os";
|
||||||
import { fileURLToPath } from "node:url";
|
import { fileURLToPath } from "node:url";
|
||||||
import { AcpRuntimeAdapter } from "../runtime-adapter.js";
|
import { AcpRuntimeAdapter } from "../runtime-adapter.js";
|
||||||
import { killAllProcesses, activeProcessCount } from "../process-manager.js";
|
import { killAllProcesses, activeProcessCount } from "../process-manager.js";
|
||||||
|
import * as provider from "../provider.js";
|
||||||
|
import * as toolBridge from "../tool-bridge.js";
|
||||||
import type { AcpSession, AgentRuntimeOptions } from "../types.js";
|
import type { AcpSession, AgentRuntimeOptions } from "../types.js";
|
||||||
|
|
||||||
const FIXTURE = fileURLToPath(new URL("./fixtures/echo-agent.mjs", import.meta.url));
|
const FIXTURE = fileURLToPath(new URL("./fixtures/echo-agent.mjs", import.meta.url));
|
||||||
@@ -122,3 +124,142 @@ describe("AcpRuntimeAdapter (U3)", () => {
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("AcpRuntimeAdapter custom-tools bridge (FNXC:AcpCustomTools)", () => {
|
||||||
|
it("keeps the ACP session alive when the custom-tools bridge cannot start", async () => {
|
||||||
|
let captured: { mcpServers?: unknown[] } | undefined;
|
||||||
|
const onText = vi.fn();
|
||||||
|
const adapter = makeAdapter();
|
||||||
|
const bridgeSpy = vi.spyOn(toolBridge, "startFusionToolBridge").mockRejectedValue(new Error("bind failed"));
|
||||||
|
const providerSpy = vi.spyOn(provider, "newAcpSession").mockImplementation(async (_connection, opts) => {
|
||||||
|
captured = opts;
|
||||||
|
return { sessionId: "degraded-session" };
|
||||||
|
});
|
||||||
|
try {
|
||||||
|
const { session } = await adapter.createSession(makeOptions({
|
||||||
|
onText,
|
||||||
|
customTools: [{ name: "fn_task_list", execute: async () => "ok" }],
|
||||||
|
}));
|
||||||
|
expect(session.sessionId).toBe("degraded-session");
|
||||||
|
expect(session.fusionToolBridgeError).toEqual({ reasonCode: "bridge-start-failed" });
|
||||||
|
expect(captured?.mcpServers ?? []).toHaveLength(0);
|
||||||
|
expect(onText).toHaveBeenCalledWith("FUSION_TOOL_BRIDGE_FAILED: bridge-start-failed");
|
||||||
|
await adapter.dispose(session);
|
||||||
|
} finally {
|
||||||
|
providerSpy.mockRestore();
|
||||||
|
bridgeSpy.mockRestore();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("registers the tool bridge in session/new mcpServers when customTools are provided", async () => {
|
||||||
|
let captured: { mcpServers?: unknown[] } | undefined;
|
||||||
|
const spy = vi
|
||||||
|
.spyOn(provider, "newAcpSession")
|
||||||
|
.mockImplementation(async (_connection, opts) => {
|
||||||
|
captured = opts;
|
||||||
|
return { sessionId: "bridge-test-session" };
|
||||||
|
});
|
||||||
|
const adapter = makeAdapter();
|
||||||
|
let session: AcpSession | undefined;
|
||||||
|
try {
|
||||||
|
const created = await adapter.createSession(
|
||||||
|
makeOptions({
|
||||||
|
customTools: [
|
||||||
|
{
|
||||||
|
name: "fn_heartbeat_done",
|
||||||
|
description: "Finish heartbeat",
|
||||||
|
parameters: { type: "object", properties: {} },
|
||||||
|
execute: async () => ({ text: "ok" }),
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
session = created.session;
|
||||||
|
expect(captured?.mcpServers).toHaveLength(1);
|
||||||
|
expect(captured?.mcpServers?.[0]).toMatchObject({
|
||||||
|
name: "fusion-custom-tools",
|
||||||
|
command: process.execPath,
|
||||||
|
});
|
||||||
|
const bridgeUrl = (captured?.mcpServers?.[0] as { env?: Array<{ name: string; value: string }> }).env?.find(
|
||||||
|
(entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL",
|
||||||
|
)?.value;
|
||||||
|
expect(bridgeUrl).toMatch(/^http:\/\/127\.0\.0\.1:\d+$/);
|
||||||
|
await adapter.dispose(session);
|
||||||
|
await expect(fetch(`${bridgeUrl}/tool-call`)).rejects.toThrow();
|
||||||
|
} finally {
|
||||||
|
// Dispose in finally so a failed assertion cannot leak the bridge/socket.
|
||||||
|
if (session) await adapter.dispose(session);
|
||||||
|
spy.mockRestore();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("direct session.dispose exposes completion of bridge shutdown", async () => {
|
||||||
|
let captured: { mcpServers?: unknown[] } | undefined;
|
||||||
|
const spy = vi.spyOn(provider, "newAcpSession").mockImplementation(async (_connection, opts) => {
|
||||||
|
captured = opts;
|
||||||
|
return { sessionId: "direct-dispose-session" };
|
||||||
|
});
|
||||||
|
let session: AcpSession | undefined;
|
||||||
|
try {
|
||||||
|
session = (await makeAdapter().createSession(
|
||||||
|
makeOptions({ customTools: [{ name: "fn_direct_dispose", execute: async () => "ok" }] }),
|
||||||
|
)).session as AcpSession;
|
||||||
|
const bridgeUrl = (captured?.mcpServers?.[0] as { env: Array<{ name: string; value: string }> })
|
||||||
|
.env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value;
|
||||||
|
session.dispose();
|
||||||
|
await expect(session.disposePromise).resolves.toBeUndefined();
|
||||||
|
await expect(fetch(`${bridgeUrl}/tool-call`)).rejects.toThrow();
|
||||||
|
} finally {
|
||||||
|
session?.dispose();
|
||||||
|
await session?.disposePromise;
|
||||||
|
spy.mockRestore();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not add a bridge when no customTools are supplied", async () => {
|
||||||
|
let captured: { mcpServers?: unknown[] } | undefined;
|
||||||
|
const spy = vi
|
||||||
|
.spyOn(provider, "newAcpSession")
|
||||||
|
.mockImplementation(async (_connection, opts) => {
|
||||||
|
captured = opts;
|
||||||
|
return { sessionId: "plain-session" };
|
||||||
|
});
|
||||||
|
const adapter = makeAdapter();
|
||||||
|
try {
|
||||||
|
const { session } = await adapter.createSession(makeOptions());
|
||||||
|
expect(captured?.mcpServers ?? []).toHaveLength(0);
|
||||||
|
await adapter.dispose(session);
|
||||||
|
} finally {
|
||||||
|
spy.mockRestore();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("disposes the bridge when session/new fails", async () => {
|
||||||
|
const spy = vi
|
||||||
|
.spyOn(provider, "newAcpSession")
|
||||||
|
.mockImplementation(async () => {
|
||||||
|
throw new Error("session/new failed");
|
||||||
|
});
|
||||||
|
const adapter = makeAdapter();
|
||||||
|
try {
|
||||||
|
await expect(
|
||||||
|
adapter.createSession(
|
||||||
|
makeOptions({
|
||||||
|
customTools: [
|
||||||
|
{
|
||||||
|
name: "fn_task_list",
|
||||||
|
description: "List",
|
||||||
|
parameters: {},
|
||||||
|
execute: async () => ({ text: "ok" }),
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
),
|
||||||
|
).rejects.toThrow(/session\/new failed/);
|
||||||
|
// The subprocess must be cleaned up even though session/new failed.
|
||||||
|
expect(activeProcessCount()).toBe(0);
|
||||||
|
} finally {
|
||||||
|
spy.mockRestore();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -0,0 +1,442 @@
|
|||||||
|
import { spawn } from "node:child_process";
|
||||||
|
import { mkdtemp, rm, writeFile } from "node:fs/promises";
|
||||||
|
import { existsSync, readdirSync } from "node:fs";
|
||||||
|
import { Server, request as httpRequest } from "node:http";
|
||||||
|
import { tmpdir } from "node:os";
|
||||||
|
import { join } from "node:path";
|
||||||
|
import { describe, expect, it, vi } from "vitest";
|
||||||
|
import { fusionToolsMcpServerPath, startFusionToolBridge, toolsToMcpToolDefs } from "../tool-bridge.js";
|
||||||
|
|
||||||
|
async function mcpRequest(child: ReturnType<typeof spawn>, payload: string): Promise<Record<string, unknown>> {
|
||||||
|
const stdout = child.stdout;
|
||||||
|
const stdin = child.stdin;
|
||||||
|
if (!stdout || !stdin) {
|
||||||
|
throw new Error("MCP child stdio is unavailable");
|
||||||
|
}
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
let output = "";
|
||||||
|
const onData = (chunk: Buffer) => {
|
||||||
|
output += chunk.toString("utf8");
|
||||||
|
const lines = output.split("\n");
|
||||||
|
const line = lines.find((candidate) => {
|
||||||
|
try {
|
||||||
|
JSON.parse(candidate);
|
||||||
|
return true;
|
||||||
|
} catch {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
if (line) {
|
||||||
|
stdout.off("data", onData);
|
||||||
|
resolve(JSON.parse(line) as Record<string, unknown>);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
stdout.setEncoding("utf8");
|
||||||
|
stdout.on("data", onData);
|
||||||
|
child.once("error", reject);
|
||||||
|
stdin.write(`${payload}\n`);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("tool-bridge", () => {
|
||||||
|
it("cleans up the schema and server when startup fails", async () => {
|
||||||
|
const listen = vi.spyOn(Server.prototype, "listen").mockImplementation(function (this: Server, ...args: unknown[]) {
|
||||||
|
const callback = args.at(-1);
|
||||||
|
if (typeof callback === "function") queueMicrotask(callback as () => void);
|
||||||
|
return this;
|
||||||
|
});
|
||||||
|
const address = vi.spyOn(Server.prototype, "address").mockReturnValue(null);
|
||||||
|
const close = vi.spyOn(Server.prototype, "close").mockImplementation(function (this: Server, callback?: (error?: Error) => void) {
|
||||||
|
callback?.();
|
||||||
|
return this;
|
||||||
|
});
|
||||||
|
const before = new Set(readdirSync(tmpdir()).filter((name) => name.startsWith("fusion-acp-mcp-schemas-")));
|
||||||
|
try {
|
||||||
|
await expect(startFusionToolBridge([{ name: "fn_startup", execute: async () => "ok" }])).rejects.toThrow("failed to bind");
|
||||||
|
const after = readdirSync(tmpdir()).filter((name) => name.startsWith("fusion-acp-mcp-schemas-"));
|
||||||
|
expect(after.filter((name) => !before.has(name))).toEqual([]);
|
||||||
|
expect(close).toHaveBeenCalled();
|
||||||
|
} finally {
|
||||||
|
listen.mockRestore();
|
||||||
|
address.mockRestore();
|
||||||
|
close.mockRestore();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("filters built-ins and maps tool schemas", () => {
|
||||||
|
expect(
|
||||||
|
toolsToMcpToolDefs([
|
||||||
|
{ name: "read", description: "builtin", parameters: {} },
|
||||||
|
{ name: "fn_not_runnable", description: "missing execute", parameters: {} },
|
||||||
|
{ name: "fn_task_list", description: "List tasks", parameters: { type: "object", properties: {} }, execute: async () => ({}) },
|
||||||
|
]),
|
||||||
|
).toEqual([
|
||||||
|
{
|
||||||
|
name: "fn_task_list",
|
||||||
|
description: "List tasks",
|
||||||
|
inputSchema: { type: "object", properties: {} },
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("preserves an isError result and its text", async () => {
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{
|
||||||
|
name: "fn_failed",
|
||||||
|
execute: async () => ({ isError: true, text: "tool failed" }),
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const bridgeUrl = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value;
|
||||||
|
const token = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")!.value;
|
||||||
|
|
||||||
|
const response = await fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_failed", arguments: {} }),
|
||||||
|
});
|
||||||
|
expect(await response.json()).toEqual({
|
||||||
|
isError: true,
|
||||||
|
content: [{ type: "text", text: "tool failed" }],
|
||||||
|
});
|
||||||
|
await bridge!.dispose();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("returns null when there are no custom tools", async () => {
|
||||||
|
expect(await startFusionToolBridge([])).toBeNull();
|
||||||
|
expect(await startFusionToolBridge(undefined)).toBeNull();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("starts a bridge that executes Fusion custom tools over authenticated HTTP", async () => {
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{
|
||||||
|
name: "fn_heartbeat_done",
|
||||||
|
description: "Finish heartbeat",
|
||||||
|
parameters: { type: "object", properties: {} },
|
||||||
|
execute: async (toolCallId: string) => ({ text: `done:${typeof toolCallId}:${toolCallId}` }),
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
expect(bridge!.toolCount).toBe(1);
|
||||||
|
expect(bridge!.mcpServer.name).toBe("fusion-custom-tools");
|
||||||
|
expect(bridge!.mcpServer.command).toBe(process.execPath);
|
||||||
|
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const bridgeUrl = env.find((e) => e.name === "FUSION_ACP_TOOL_BRIDGE_URL")?.value;
|
||||||
|
const token = env.find((e) => e.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")?.value;
|
||||||
|
expect(bridgeUrl).toMatch(/^http:\/\/127\.0\.0\.1:\d+$/);
|
||||||
|
expect(token).toBeTruthy();
|
||||||
|
|
||||||
|
// No token → 401 (a same-host probe must not invoke Fusion closures).
|
||||||
|
const unauthorized = await fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { "content-type": "application/json" },
|
||||||
|
body: JSON.stringify({ name: "fn_heartbeat_done", arguments: {} }),
|
||||||
|
});
|
||||||
|
expect(unauthorized.status).toBe(401);
|
||||||
|
|
||||||
|
// Authenticated call threads the real toolCallId.
|
||||||
|
const res = await fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { "content-type": "application/json", authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_heartbeat_done", toolCallId: "call-42", arguments: {} }),
|
||||||
|
});
|
||||||
|
const body = (await res.json()) as { isError?: boolean; content?: Array<{ text?: string }> };
|
||||||
|
expect(body.isError).toBe(false);
|
||||||
|
expect(body.content?.[0]?.text).toBe("done:string:call-42");
|
||||||
|
|
||||||
|
const numericId = await fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { "content-type": "application/json", authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_heartbeat_done", toolCallId: 42, arguments: {} }),
|
||||||
|
});
|
||||||
|
const numericBody = (await numericId.json()) as { content?: Array<{ text?: string }> };
|
||||||
|
// The numeric JSON-RPC id must be threaded as the real toolCallId (coerced to
|
||||||
|
// string), never replaced by a fabricated UUID fallback.
|
||||||
|
expect(numericBody.content?.[0]?.text).toBe("done:string:42");
|
||||||
|
|
||||||
|
// Unknown tool → 404.
|
||||||
|
const unknown = await fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { "content-type": "application/json", authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_nope", arguments: {} }),
|
||||||
|
});
|
||||||
|
expect(unknown.status).toBe(404);
|
||||||
|
|
||||||
|
await bridge!.dispose();
|
||||||
|
// Port closed after dispose: a follow-up request must fail.
|
||||||
|
await expect(
|
||||||
|
fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { "content-type": "application/json", authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_heartbeat_done", arguments: {} }),
|
||||||
|
}),
|
||||||
|
).rejects.toThrow();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not execute built-in tools posted directly to the bridge", async () => {
|
||||||
|
let executions = 0;
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{ name: "read", parameters: {}, execute: async () => { executions += 1; return "should not run"; } },
|
||||||
|
{ name: "fn_allowed", parameters: {}, execute: async () => "ok" },
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const bridgeUrl = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value;
|
||||||
|
const token = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")!.value;
|
||||||
|
const response = await fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "read", arguments: {} }),
|
||||||
|
});
|
||||||
|
expect(response.status).toBe(404);
|
||||||
|
expect(executions).toBe(0);
|
||||||
|
await bridge!.dispose();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not leak an unhandled rejection when a request aborts", async () => {
|
||||||
|
let unhandled: unknown;
|
||||||
|
let started!: () => void;
|
||||||
|
let rejectHandler!: (reason: Error) => void;
|
||||||
|
const handlerStarted = new Promise<void>((resolve) => { started = resolve; });
|
||||||
|
const handler = new Promise<never>((_resolve, reject) => { rejectHandler = reject; });
|
||||||
|
const onUnhandled = (reason: unknown) => { unhandled = reason; };
|
||||||
|
process.on("unhandledRejection", onUnhandled);
|
||||||
|
let bridge: Awaited<ReturnType<typeof startFusionToolBridge>> | undefined;
|
||||||
|
try {
|
||||||
|
bridge = await startFusionToolBridge([
|
||||||
|
{ name: "fn_reject", parameters: {}, execute: async () => { started(); return handler; } },
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const bridgeUrl = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value;
|
||||||
|
const token = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")!.value;
|
||||||
|
const controller = new AbortController();
|
||||||
|
const request = fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_reject", arguments: {} }),
|
||||||
|
signal: controller.signal,
|
||||||
|
}).catch(() => undefined);
|
||||||
|
await handlerStarted;
|
||||||
|
controller.abort();
|
||||||
|
rejectHandler(new Error("handler failed"));
|
||||||
|
await request;
|
||||||
|
expect(unhandled).toBeUndefined();
|
||||||
|
} finally {
|
||||||
|
await bridge?.dispose();
|
||||||
|
process.off("unhandledRejection", onUnhandled);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("dispose is idempotent and removes the temporary schema", async () => {
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{ name: "fn_task_list", description: "List", parameters: {}, execute: async () => ({ text: "ok" }) },
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
const schemaPath = bridge!.mcpServer.args[1] as string;
|
||||||
|
expect(existsSync(schemaPath)).toBe(true);
|
||||||
|
await bridge!.dispose();
|
||||||
|
await bridge!.dispose();
|
||||||
|
expect(existsSync(schemaPath)).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not resume a request into tool execution after disposal", async () => {
|
||||||
|
const close = vi.spyOn(Server.prototype, "close").mockImplementation(function (this: Server) {
|
||||||
|
return this;
|
||||||
|
});
|
||||||
|
const closeAllConnections = vi.spyOn(Server.prototype, "closeAllConnections").mockImplementation(() => undefined);
|
||||||
|
let executions = 0;
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{ name: "fn_after_dispose", parameters: {}, execute: async () => { executions += 1; return "must not run"; } },
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const bridgeUrl = new URL(env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value);
|
||||||
|
const token = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")!.value;
|
||||||
|
const pendingResponse = new Promise<{ status: number }>((resolve, reject) => {
|
||||||
|
const req = httpRequest({
|
||||||
|
hostname: bridgeUrl.hostname,
|
||||||
|
port: Number(bridgeUrl.port),
|
||||||
|
path: "/tool-call",
|
||||||
|
method: "POST",
|
||||||
|
headers: { authorization: `Bearer ${token}`, "content-type": "application/json" },
|
||||||
|
}, (response) => resolve({ status: response.statusCode ?? 0 }));
|
||||||
|
req.once("error", reject);
|
||||||
|
req.flushHeaders();
|
||||||
|
void bridge!.dispose().then(() => req.end(JSON.stringify({ name: "fn_after_dispose", arguments: {} })));
|
||||||
|
});
|
||||||
|
try {
|
||||||
|
await expect(pendingResponse).resolves.toMatchObject({ status: 503 });
|
||||||
|
expect(executions).toBe(0);
|
||||||
|
} finally {
|
||||||
|
close.mockRestore();
|
||||||
|
closeAllConnections.mockRestore();
|
||||||
|
await bridge?.dispose();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("aborts and awaits an in-flight tool before disposal completes", async () => {
|
||||||
|
let started!: () => void;
|
||||||
|
const executionStarted = new Promise<void>((resolve) => {
|
||||||
|
started = resolve;
|
||||||
|
});
|
||||||
|
let finished = false;
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{
|
||||||
|
name: "fn_slow",
|
||||||
|
parameters: {},
|
||||||
|
execute: async (_id, _args, signal) => {
|
||||||
|
started();
|
||||||
|
await new Promise<void>((resolve) => {
|
||||||
|
if (signal?.aborted) return resolve();
|
||||||
|
signal?.addEventListener("abort", () => resolve(), { once: true });
|
||||||
|
});
|
||||||
|
finished = true;
|
||||||
|
return { text: "done" };
|
||||||
|
},
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const bridgeUrl = env.find((e) => e.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value;
|
||||||
|
const token = env.find((e) => e.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")!.value;
|
||||||
|
const request = fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_slow", arguments: {} }),
|
||||||
|
});
|
||||||
|
await executionStarted;
|
||||||
|
let disposed = false;
|
||||||
|
const disposing = bridge!.dispose().then(() => {
|
||||||
|
disposed = true;
|
||||||
|
});
|
||||||
|
expect(disposed).toBe(false);
|
||||||
|
await disposing;
|
||||||
|
expect(finished).toBe(true);
|
||||||
|
await expect(request).rejects.toThrow();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not wait for the server close callback before draining", async () => {
|
||||||
|
const close = vi.spyOn(Server.prototype, "close").mockImplementation(function (this: Server) {
|
||||||
|
return this;
|
||||||
|
});
|
||||||
|
const closeAllConnections = vi.spyOn(Server.prototype, "closeAllConnections").mockImplementation(() => undefined);
|
||||||
|
let started!: () => void;
|
||||||
|
const executionStarted = new Promise<void>((resolve) => {
|
||||||
|
started = resolve;
|
||||||
|
});
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{
|
||||||
|
name: "fn_cooperative",
|
||||||
|
parameters: {},
|
||||||
|
execute: async (_id, _args, signal) => {
|
||||||
|
started();
|
||||||
|
await new Promise<void>((resolve) => signal?.addEventListener("abort", () => resolve(), { once: true }));
|
||||||
|
return "done";
|
||||||
|
},
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
try {
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const request = fetch(`${env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { authorization: `Bearer ${env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")!.value}` },
|
||||||
|
body: JSON.stringify({ name: "fn_cooperative", arguments: {} }),
|
||||||
|
});
|
||||||
|
await executionStarted;
|
||||||
|
await expect(Promise.race([
|
||||||
|
bridge!.dispose(),
|
||||||
|
new Promise((_, reject) => setTimeout(() => reject(new Error("dispose timed out")), 500)),
|
||||||
|
])).resolves.toBeUndefined();
|
||||||
|
await request;
|
||||||
|
expect(close).toHaveBeenCalled();
|
||||||
|
expect(closeAllConnections).toHaveBeenCalled();
|
||||||
|
} finally {
|
||||||
|
close.mockRestore();
|
||||||
|
closeAllConnections.mockRestore();
|
||||||
|
await bridge?.dispose();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("completes disposal and removes the schema when a tool ignores abort", async () => {
|
||||||
|
let started!: () => void;
|
||||||
|
const executionStarted = new Promise<void>((resolve) => {
|
||||||
|
started = resolve;
|
||||||
|
});
|
||||||
|
const bridge = await startFusionToolBridge([
|
||||||
|
{
|
||||||
|
name: "fn_never_settles",
|
||||||
|
parameters: {},
|
||||||
|
execute: async () => {
|
||||||
|
started();
|
||||||
|
return new Promise(() => undefined);
|
||||||
|
},
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
expect(bridge).not.toBeNull();
|
||||||
|
const schemaPath = bridge!.mcpServer.args[1] as string;
|
||||||
|
const env = bridge!.mcpServer.env;
|
||||||
|
const bridgeUrl = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_URL")!.value;
|
||||||
|
const token = env.find((entry) => entry.name === "FUSION_ACP_TOOL_BRIDGE_TOKEN")!.value;
|
||||||
|
const request = fetch(`${bridgeUrl}/tool-call`, {
|
||||||
|
method: "POST",
|
||||||
|
headers: { authorization: `Bearer ${token}` },
|
||||||
|
body: JSON.stringify({ name: "fn_never_settles", arguments: {} }),
|
||||||
|
});
|
||||||
|
// Handle the eventual socket-close rejection up front so it can never surface
|
||||||
|
// as an unhandled rejection while dispose drains for TOOL_DRAIN_TIMEOUT_MS.
|
||||||
|
const requestOutcome = request.then(
|
||||||
|
() => "resolved",
|
||||||
|
() => "rejected",
|
||||||
|
);
|
||||||
|
// Wait until the tool body is running so the execution is tracked before dispose.
|
||||||
|
await executionStarted;
|
||||||
|
|
||||||
|
await expect(bridge!.dispose()).resolves.toBeUndefined();
|
||||||
|
expect(existsSync(schemaPath)).toBe(false);
|
||||||
|
expect(await requestOutcome).toBe("rejected");
|
||||||
|
}, 10_000);
|
||||||
|
|
||||||
|
it("serves the full MCP protocol from the co-located schema server", async () => {
|
||||||
|
const directory = await mkdtemp(join(tmpdir(), "fusion-acp-mcp-smoke-"));
|
||||||
|
const schemaPath = join(directory, "schemas.json");
|
||||||
|
await writeFile(
|
||||||
|
schemaPath,
|
||||||
|
JSON.stringify([
|
||||||
|
{ name: "fn_heartbeat_done", description: "Finish heartbeat", inputSchema: { type: "object", properties: {} } },
|
||||||
|
]),
|
||||||
|
);
|
||||||
|
const child = spawn(process.execPath, [fusionToolsMcpServerPath(), schemaPath], {
|
||||||
|
env: { ...process.env, FUSION_ACP_TOOL_BRIDGE_URL: "http://127.0.0.1:1" },
|
||||||
|
stdio: ["pipe", "pipe", "pipe"],
|
||||||
|
});
|
||||||
|
try {
|
||||||
|
const initialized = await mcpRequest(child, '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{}}');
|
||||||
|
expect((initialized.result as { serverInfo?: { name?: string } }).serverInfo?.name).toBe("fusion-custom-tools");
|
||||||
|
|
||||||
|
const pinged = await mcpRequest(child, '{"jsonrpc":"2.0","id":2,"method":"ping","params":{}}');
|
||||||
|
expect(pinged.result).toEqual({});
|
||||||
|
|
||||||
|
const listed = await mcpRequest(child, '{"jsonrpc":"2.0","id":3,"method":"tools/list","params":{}}');
|
||||||
|
const tools = (listed.result as { tools?: Array<{ name: string }> }).tools ?? [];
|
||||||
|
expect(tools.map((tool) => tool.name)).toEqual(["fn_heartbeat_done"]);
|
||||||
|
|
||||||
|
// tools/call against a dead bridge (port 1) must surface an error result,
|
||||||
|
// never an empty "success".
|
||||||
|
const called = await mcpRequest(
|
||||||
|
child,
|
||||||
|
'{"jsonrpc":"2.0","id":4,"method":"tools/call","params":{"name":"fn_heartbeat_done","arguments":{}}}',
|
||||||
|
);
|
||||||
|
expect((called.result as { isError?: boolean }).isError).toBe(true);
|
||||||
|
|
||||||
|
// Unknown method → -32601.
|
||||||
|
const missing = await mcpRequest(child, '{"jsonrpc":"2.0","id":5,"method":"bogus","params":{}}');
|
||||||
|
expect((missing.error as { code?: number }).code).toBe(-32601);
|
||||||
|
} finally {
|
||||||
|
child.kill();
|
||||||
|
await rm(directory, { recursive: true, force: true });
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
171
plugins/fusion-plugin-acp-runtime/src/mcp-schema-server.cjs
Normal file
171
plugins/fusion-plugin-acp-runtime/src/mcp-schema-server.cjs
Normal file
@@ -0,0 +1,171 @@
|
|||||||
|
#!/usr/bin/env node
|
||||||
|
/*
|
||||||
|
FNXC:AcpCustomTools 2026-08-16-00:30:
|
||||||
|
Executable MCP bridge for Fusion custom tools (fn_*) on generic ACP paths
|
||||||
|
(Hermes ACP, Prime, ...). tools/list is served from a schema file; tools/call
|
||||||
|
POSTs to a localhost bridge owned by AcpRuntimeAdapter so ToolDefinition.execute
|
||||||
|
runs in-process with the engine's closures. The bridge authenticates the POST
|
||||||
|
with a per-session bearer token passed via FUSION_ACP_TOOL_BRIDGE_TOKEN.
|
||||||
|
*/
|
||||||
|
"use strict";
|
||||||
|
|
||||||
|
const fs = require("fs");
|
||||||
|
const http = require("http");
|
||||||
|
const readline = require("readline");
|
||||||
|
// FNXC:AcpCustomTools 2026-08-16-00:30: CJS has no global URL under eslint no-undef; use node:url.
|
||||||
|
const { URL } = require("node:url");
|
||||||
|
|
||||||
|
const schemaPath = process.argv[2];
|
||||||
|
const bridgeUrl = process.env.FUSION_ACP_TOOL_BRIDGE_URL;
|
||||||
|
const bridgeToken = process.env.FUSION_ACP_TOOL_BRIDGE_TOKEN;
|
||||||
|
if (!schemaPath || !bridgeUrl) {
|
||||||
|
process.stderr.write("fusion-custom-tools-mcp: missing schema path or FUSION_ACP_TOOL_BRIDGE_URL\n");
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
|
||||||
|
let tools = [];
|
||||||
|
try {
|
||||||
|
tools = JSON.parse(fs.readFileSync(schemaPath, "utf-8"));
|
||||||
|
if (!Array.isArray(tools)) tools = [];
|
||||||
|
} catch {
|
||||||
|
process.exit(1);
|
||||||
|
}
|
||||||
|
|
||||||
|
function write(msg) {
|
||||||
|
process.stdout.write(JSON.stringify(msg) + "\n");
|
||||||
|
}
|
||||||
|
|
||||||
|
function callBridge(toolName, toolCallId, args) {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
const body = JSON.stringify({ name: toolName, toolCallId, arguments: args ?? {} });
|
||||||
|
const url = new URL("/tool-call", bridgeUrl);
|
||||||
|
const req = http.request(
|
||||||
|
{
|
||||||
|
hostname: url.hostname,
|
||||||
|
port: url.port,
|
||||||
|
path: url.pathname,
|
||||||
|
method: "POST",
|
||||||
|
headers: {
|
||||||
|
"content-type": "application/json",
|
||||||
|
"content-length": Buffer.byteLength(body),
|
||||||
|
...(bridgeToken ? { authorization: `Bearer ${bridgeToken}` } : {}),
|
||||||
|
},
|
||||||
|
timeout: 120_000,
|
||||||
|
},
|
||||||
|
(res) => {
|
||||||
|
let data = "";
|
||||||
|
res.on("data", (chunk) => {
|
||||||
|
data += chunk;
|
||||||
|
});
|
||||||
|
res.on("end", () => {
|
||||||
|
// A non-2xx bridge response is a transport failure, not a tool result:
|
||||||
|
// surface it as an error instead of fabricating an empty "success".
|
||||||
|
if (res.statusCode < 200 || res.statusCode >= 300) {
|
||||||
|
reject(new Error(`tool bridge responded ${res.statusCode}: ${data.slice(0, 200)}`));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
resolve(JSON.parse(data || "{}"));
|
||||||
|
} catch (err) {
|
||||||
|
reject(err);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
},
|
||||||
|
);
|
||||||
|
req.on("error", reject);
|
||||||
|
req.on("timeout", () => {
|
||||||
|
req.destroy(new Error("tool bridge timeout"));
|
||||||
|
});
|
||||||
|
req.write(body);
|
||||||
|
req.end();
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
const rl = readline.createInterface({ input: process.stdin });
|
||||||
|
rl.on("line", (line) => {
|
||||||
|
let msg;
|
||||||
|
try {
|
||||||
|
msg = JSON.parse(line);
|
||||||
|
} catch {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (msg.method === "initialize") {
|
||||||
|
write({
|
||||||
|
jsonrpc: "2.0",
|
||||||
|
id: msg.id,
|
||||||
|
result: {
|
||||||
|
protocolVersion: "2024-11-05",
|
||||||
|
capabilities: { tools: {} },
|
||||||
|
serverInfo: { name: "fusion-custom-tools", version: "1.0.0" },
|
||||||
|
},
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Ping must answer: some clients (Hermes ACP) treat a missing handler as a
|
||||||
|
// method-not-found and back off.
|
||||||
|
if (msg.method === "ping") {
|
||||||
|
write({ jsonrpc: "2.0", id: msg.id, result: {} });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (msg.method === "notifications/initialized" || msg.method === "initialized") {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (msg.method === "tools/list") {
|
||||||
|
write({
|
||||||
|
jsonrpc: "2.0",
|
||||||
|
id: msg.id,
|
||||||
|
result: {
|
||||||
|
tools: tools.map((tool) => ({
|
||||||
|
name: tool.name,
|
||||||
|
description: tool.description ?? "",
|
||||||
|
inputSchema: tool.inputSchema ?? { type: "object", properties: {} },
|
||||||
|
})),
|
||||||
|
},
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (msg.method === "tools/call") {
|
||||||
|
const toolName = msg.params?.name;
|
||||||
|
const args = msg.params?.arguments ?? {};
|
||||||
|
// JSON-RPC ids are commonly numbers; coerce so the bridge threads the real
|
||||||
|
// request id as the toolCallId instead of falling back to a fabricated id.
|
||||||
|
const toolCallId = typeof msg.id === "string" ? msg.id : String(msg.id);
|
||||||
|
callBridge(toolName, toolCallId, args)
|
||||||
|
.then((result) => {
|
||||||
|
write({
|
||||||
|
jsonrpc: "2.0",
|
||||||
|
id: msg.id,
|
||||||
|
result: {
|
||||||
|
content: Array.isArray(result.content)
|
||||||
|
? result.content
|
||||||
|
: [{ type: "text", text: typeof result.text === "string" ? result.text : JSON.stringify(result) }],
|
||||||
|
isError: result.isError === true,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
})
|
||||||
|
.catch((err) => {
|
||||||
|
write({
|
||||||
|
jsonrpc: "2.0",
|
||||||
|
id: msg.id,
|
||||||
|
result: {
|
||||||
|
content: [{ type: "text", text: err instanceof Error ? err.message : String(err) }],
|
||||||
|
isError: true,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (msg.id !== undefined) {
|
||||||
|
write({
|
||||||
|
jsonrpc: "2.0",
|
||||||
|
id: msg.id,
|
||||||
|
error: { code: -32601, message: `Method not found: ${msg.method}` },
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
@@ -18,6 +18,7 @@ import {
|
|||||||
} from "./provider.js";
|
} from "./provider.js";
|
||||||
import { buildSpawnEnv } from "./process-manager.js";
|
import { buildSpawnEnv } from "./process-manager.js";
|
||||||
import { buildPromptBlocks, extractPromptImagesFromOptions } from "./prompt-builder.js";
|
import { buildPromptBlocks, extractPromptImagesFromOptions } from "./prompt-builder.js";
|
||||||
|
import { startFusionToolBridge, type FusionToolBridge, type ToolLike } from "./tool-bridge.js";
|
||||||
import type {
|
import type {
|
||||||
AgentRuntime,
|
AgentRuntime,
|
||||||
AgentRuntimeOptions,
|
AgentRuntimeOptions,
|
||||||
@@ -88,29 +89,56 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
// caller supplied them (U10 — Route A); absent/empty keeps the Route B
|
// caller supplied them (U10 — Route A); absent/empty keeps the Route B
|
||||||
// read-only ask posture. Tool calls still route through the U5 permission floor.
|
// read-only ask posture. Tool calls still route through the U5 permission floor.
|
||||||
//
|
//
|
||||||
|
// FNXC:AcpCustomTools 2026-08-16-00:30:
|
||||||
|
// Engine customTools (fn_*) ride the same mcpServers channel: a loopback tool
|
||||||
|
// bridge is started in-process and registered as a stdio MCP server so any
|
||||||
|
// ACP agent (Hermes ACP, Prime, ...) can invoke Fusion closures. The bridge
|
||||||
|
// is disposed on session/new failure and on session teardown.
|
||||||
|
//
|
||||||
// FNXC:GrokAcp 2026-07-11-14:00:
|
// FNXC:GrokAcp 2026-07-11-14:00:
|
||||||
// Callers (Grok runtime) may also pass `_meta` (pluginDirs / rules /
|
// Callers (Grok runtime) may also pass `_meta` (pluginDirs / rules /
|
||||||
// systemPromptOverride) via options.sessionMeta so agent-specific skill and
|
// systemPromptOverride) via options.sessionMeta so agent-specific skill and
|
||||||
// prompt setup rides on session/new without a second protocol hop.
|
// prompt setup rides on session/new without a second protocol hop.
|
||||||
|
let toolBridge: FusionToolBridge | null = null;
|
||||||
|
let toolBridgeFailure: "mcp-schema-server-missing" | "bridge-start-failed" | undefined;
|
||||||
let sessionId: string;
|
let sessionId: string;
|
||||||
|
const customTools = Array.isArray(options.customTools) ? (options.customTools as ToolLike[]) : [];
|
||||||
try {
|
try {
|
||||||
|
if (customTools.length > 0) {
|
||||||
|
try {
|
||||||
|
toolBridge = await startFusionToolBridge(customTools);
|
||||||
|
} catch (error) {
|
||||||
|
toolBridgeFailure = (error as { code?: string }).code === "mcp-schema-server-missing"
|
||||||
|
? "mcp-schema-server-missing"
|
||||||
|
: "bridge-start-failed";
|
||||||
|
options.onText?.(`FUSION_TOOL_BRIDGE_FAILED: ${toolBridgeFailure}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
const sessionMeta =
|
const sessionMeta =
|
||||||
options && typeof options === "object" && "sessionMeta" in options
|
options && typeof options === "object" && "sessionMeta" in options
|
||||||
? (options as { sessionMeta?: Record<string, unknown> }).sessionMeta
|
? (options as { sessionMeta?: Record<string, unknown> }).sessionMeta
|
||||||
: undefined;
|
: undefined;
|
||||||
const opened = await newAcpSession(connection, {
|
const opened = await newAcpSession(connection, {
|
||||||
cwd: options.cwd,
|
cwd: options.cwd,
|
||||||
mcpServers: options.mcpServers,
|
mcpServers: [...(options.mcpServers ?? []), ...(toolBridge ? [toolBridge.mcpServer] : [])],
|
||||||
meta: sessionMeta,
|
meta: sessionMeta,
|
||||||
});
|
});
|
||||||
sessionId = opened.sessionId;
|
sessionId = opened.sessionId;
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
// Don't leak the subprocess if session/new fails after a good handshake.
|
// Don't leak the subprocess or the bridge if session/new fails after a
|
||||||
|
// good handshake.
|
||||||
|
await toolBridge?.dispose();
|
||||||
connection.dispose();
|
connection.dispose();
|
||||||
throw err;
|
throw err;
|
||||||
}
|
}
|
||||||
|
|
||||||
let disposed = false;
|
let disposed = false;
|
||||||
|
let bridgeDisposePromise: Promise<void> | undefined;
|
||||||
|
let disposePromise = Promise.resolve();
|
||||||
|
const disposeBridge = (): Promise<void> => {
|
||||||
|
bridgeDisposePromise ??= toolBridge?.dispose() ?? Promise.resolve();
|
||||||
|
return bridgeDisposePromise;
|
||||||
|
};
|
||||||
const session: AcpSession = {
|
const session: AcpSession = {
|
||||||
model,
|
model,
|
||||||
systemPrompt: options.systemPrompt,
|
systemPrompt: options.systemPrompt,
|
||||||
@@ -118,6 +146,7 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
cwd: options.cwd,
|
cwd: options.cwd,
|
||||||
lastModelDescription: `acp/${model}`,
|
lastModelDescription: `acp/${model}`,
|
||||||
callbacks,
|
callbacks,
|
||||||
|
fusionToolBridgeError: toolBridgeFailure ? { reasonCode: toolBridgeFailure } : undefined,
|
||||||
// Persist the per-run gate (KTD3) so U5/U7 can reach the live action gate.
|
// Persist the per-run gate (KTD3) so U5/U7 can reach the live action gate.
|
||||||
gate: options.actionGateContext,
|
gate: options.actionGateContext,
|
||||||
connection,
|
connection,
|
||||||
@@ -125,12 +154,19 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
// turn that trips the per-turn output cap can't latch and suppress every
|
// turn that trips the per-turn output cap can't latch and suppress every
|
||||||
// subsequent turn (FIX 1).
|
// subsequent turn (FIX 1).
|
||||||
resetTurn,
|
resetTurn,
|
||||||
|
disposeBridge,
|
||||||
|
get disposePromise() {
|
||||||
|
return disposePromise;
|
||||||
|
},
|
||||||
dispose: () => {
|
dispose: () => {
|
||||||
if (disposed) return;
|
if (disposed) return;
|
||||||
disposed = true;
|
disposed = true;
|
||||||
// Drain in-flight permission requests BEFORE the registry kill so a
|
// Drain in-flight permission requests BEFORE the registry kill so a
|
||||||
// blocked agent is released (KTD4a — the SIGKILL is still authoritative).
|
// blocked agent is released (KTD4a — the SIGKILL is still authoritative).
|
||||||
cancelPending();
|
cancelPending();
|
||||||
|
// Close the loopback tool bridge so no port or schema outlives the
|
||||||
|
// session (idempotent).
|
||||||
|
disposePromise = disposeBridge();
|
||||||
connection.dispose();
|
connection.dispose();
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
@@ -184,6 +220,7 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
if (acp.connection && acp.sessionId) {
|
if (acp.connection && acp.sessionId) {
|
||||||
await cancelAcpSession(acp.connection, acp.sessionId);
|
await cancelAcpSession(acp.connection, acp.sessionId);
|
||||||
}
|
}
|
||||||
|
await acp.disposeBridge?.();
|
||||||
session.dispose();
|
session.dispose();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
273
plugins/fusion-plugin-acp-runtime/src/tool-bridge.ts
Normal file
273
plugins/fusion-plugin-acp-runtime/src/tool-bridge.ts
Normal file
@@ -0,0 +1,273 @@
|
|||||||
|
/*
|
||||||
|
FNXC:AcpCustomTools 2026-08-16-00:30:
|
||||||
|
Host Fusion custom tools (fn_*) for any ACP agent (Hermes ACP, Prime, ...).
|
||||||
|
ToolDefinition.execute closures only run in-process, so AcpRuntimeAdapter starts
|
||||||
|
a loopback HTTP bridge and pairs it with mcp-schema-server.cjs (stdio MCP) that
|
||||||
|
the agent connects to via session/new.mcpServers. Dispose closes the bridge and
|
||||||
|
removes the temporary schema so no port or file outlives the session.
|
||||||
|
*/
|
||||||
|
import { createServer, type Server } from "node:http";
|
||||||
|
import { existsSync, unlinkSync, writeFileSync } from "node:fs";
|
||||||
|
import { tmpdir } from "node:os";
|
||||||
|
import { dirname, join } from "node:path";
|
||||||
|
import { fileURLToPath } from "node:url";
|
||||||
|
import { randomBytes, randomUUID } from "node:crypto";
|
||||||
|
import type { AcpMcpServerStdio } from "./types.js";
|
||||||
|
|
||||||
|
const BUILT_IN_TOOL_NAMES = new Set(["read", "write", "edit", "bash", "grep", "find"]);
|
||||||
|
// Custom tools may ignore AbortSignal; bound the drain so dispose never hangs
|
||||||
|
// on a tool that never settles (abort already unblocks cooperative tools).
|
||||||
|
const TOOL_DRAIN_TIMEOUT_MS = 5_000;
|
||||||
|
|
||||||
|
export interface ToolLike {
|
||||||
|
name: string;
|
||||||
|
description?: string;
|
||||||
|
parameters?: Record<string, unknown>;
|
||||||
|
execute?: (
|
||||||
|
toolCallId: string,
|
||||||
|
params: unknown,
|
||||||
|
signal?: AbortSignal,
|
||||||
|
onUpdate?: unknown,
|
||||||
|
ctx?: unknown,
|
||||||
|
) => Promise<unknown> | unknown;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface McpToolDef {
|
||||||
|
name: string;
|
||||||
|
description: string;
|
||||||
|
inputSchema: Record<string, unknown>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface FusionToolBridge {
|
||||||
|
mcpServer: AcpMcpServerStdio;
|
||||||
|
dispose: () => Promise<void>;
|
||||||
|
toolCount: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function toolsToMcpToolDefs(tools: ReadonlyArray<ToolLike> | undefined): McpToolDef[] {
|
||||||
|
if (!Array.isArray(tools)) return [];
|
||||||
|
return tools
|
||||||
|
.filter((tool) => tool && typeof tool.name === "string" && tool.name.trim().length > 0 && !BUILT_IN_TOOL_NAMES.has(tool.name) && typeof tool.execute === "function")
|
||||||
|
.map((tool) => ({
|
||||||
|
name: tool.name,
|
||||||
|
description: typeof tool.description === "string" ? tool.description : "",
|
||||||
|
inputSchema: tool.parameters ?? { type: "object", properties: {} },
|
||||||
|
}));
|
||||||
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AcpCustomToolsPackaging 2026-08-16-00:30:
|
||||||
|
The stdio MCP child resolves this asset beside the loaded bridge module. Keep the
|
||||||
|
source asset co-located for source-loaded plugins and copy it beside dist output
|
||||||
|
on builds (the Grok postbuild pattern); otherwise the host reports
|
||||||
|
`handshake failed: connection closed: initialize response` for fusion-custom-tools.
|
||||||
|
*/
|
||||||
|
export function fusionToolsMcpServerPath(): string {
|
||||||
|
return join(dirname(fileURLToPath(import.meta.url)), "mcp-schema-server.cjs");
|
||||||
|
}
|
||||||
|
|
||||||
|
function missingMcpSchemaServerError(serverPath: string): Error {
|
||||||
|
const error = new Error(`Fusion MCP schema server is missing: ${serverPath}`) as Error & { code?: string };
|
||||||
|
error.code = "mcp-schema-server-missing";
|
||||||
|
return error;
|
||||||
|
}
|
||||||
|
|
||||||
|
function resultToText(result: unknown): string {
|
||||||
|
if (result == null) return "";
|
||||||
|
if (typeof result === "string") return result;
|
||||||
|
if (typeof result === "object") {
|
||||||
|
const obj = result as { content?: unknown; text?: unknown };
|
||||||
|
if (typeof obj.text === "string") return obj.text;
|
||||||
|
if (Array.isArray(obj.content)) {
|
||||||
|
return obj.content
|
||||||
|
.map((block) => {
|
||||||
|
if (block && typeof block === "object" && "text" in block && typeof (block as { text: unknown }).text === "string") {
|
||||||
|
return (block as { text: string }).text;
|
||||||
|
}
|
||||||
|
return JSON.stringify(block);
|
||||||
|
})
|
||||||
|
.join("\n");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
return JSON.stringify(result);
|
||||||
|
} catch {
|
||||||
|
return String(result);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Start a loopback tool bridge and return the ACP mcpServers stdio entry the
|
||||||
|
* agent should connect to for Fusion custom tools. Returns null when there are
|
||||||
|
* no tools. The bridge authenticates every request with a per-session bearer
|
||||||
|
* token so a same-host probe cannot invoke Fusion closures.
|
||||||
|
*/
|
||||||
|
export async function startFusionToolBridge(tools: ReadonlyArray<ToolLike> | undefined): Promise<FusionToolBridge | null> {
|
||||||
|
const defs = toolsToMcpToolDefs(tools);
|
||||||
|
if (defs.length === 0) return null;
|
||||||
|
|
||||||
|
const serverPath = fusionToolsMcpServerPath();
|
||||||
|
if (!existsSync(serverPath)) {
|
||||||
|
throw missingMcpSchemaServerError(serverPath);
|
||||||
|
}
|
||||||
|
|
||||||
|
const byName = new Map<string, ToolLike>();
|
||||||
|
for (const tool of tools ?? []) {
|
||||||
|
if (tool && typeof tool.name === "string" && !BUILT_IN_TOOL_NAMES.has(tool.name) && typeof tool.execute === "function") {
|
||||||
|
byName.set(tool.name, tool);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
const token = randomBytes(24).toString("hex");
|
||||||
|
const schemaPath = join(tmpdir(), `fusion-acp-mcp-schemas-${process.pid}-${randomUUID()}.json`);
|
||||||
|
writeFileSync(schemaPath, JSON.stringify(defs));
|
||||||
|
const activeExecutions = new Set<Promise<void>>();
|
||||||
|
const activeControllers = new Set<AbortController>();
|
||||||
|
|
||||||
|
const server: Server = createServer((req, res) => {
|
||||||
|
void (async () => {
|
||||||
|
if (req.method !== "POST" || req.url !== "/tool-call") {
|
||||||
|
res.statusCode = 404;
|
||||||
|
res.end(JSON.stringify({ isError: true, text: "not found" }));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
// Per-session bearer token: the MCP shim carries it; anything else is rejected.
|
||||||
|
const auth = req.headers.authorization;
|
||||||
|
if (auth !== `Bearer ${token}`) {
|
||||||
|
res.statusCode = 401;
|
||||||
|
res.end(JSON.stringify({ isError: true, text: "unauthorized" }));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let body = "";
|
||||||
|
for await (const chunk of req) body += chunk;
|
||||||
|
if (disposed) {
|
||||||
|
res.statusCode = 503;
|
||||||
|
res.end(JSON.stringify({ isError: true, text: "tool bridge disposed" }));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let parsed: { name?: string; toolCallId?: string | number; arguments?: unknown };
|
||||||
|
try {
|
||||||
|
parsed = JSON.parse(body || "{}") as { name?: string; toolCallId?: string | number; arguments?: unknown };
|
||||||
|
} catch {
|
||||||
|
res.statusCode = 400;
|
||||||
|
res.end(JSON.stringify({ isError: true, text: "invalid JSON body" }));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const name = typeof parsed.name === "string" ? parsed.name : "";
|
||||||
|
const tool = byName.get(name);
|
||||||
|
if (!tool?.execute) {
|
||||||
|
res.statusCode = 404;
|
||||||
|
res.end(JSON.stringify({ isError: true, text: `Unknown Fusion tool: ${name}` }));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
const execute = tool.execute;
|
||||||
|
const controller = new AbortController();
|
||||||
|
activeControllers.add(controller);
|
||||||
|
const execution = (async () => {
|
||||||
|
try {
|
||||||
|
// Thread the real MCP request id as the toolCallId so correlation,
|
||||||
|
// cancellation, and dedupe keep working (never a fabricated id).
|
||||||
|
const result = await execute(
|
||||||
|
(typeof parsed.toolCallId === "string" || typeof parsed.toolCallId === "number") && String(parsed.toolCallId)
|
||||||
|
? String(parsed.toolCallId)
|
||||||
|
: `acp-mcp-${randomUUID()}`,
|
||||||
|
parsed.arguments ?? {},
|
||||||
|
controller.signal,
|
||||||
|
undefined,
|
||||||
|
undefined,
|
||||||
|
);
|
||||||
|
res.statusCode = 200;
|
||||||
|
res.setHeader("content-type", "application/json");
|
||||||
|
res.end(
|
||||||
|
JSON.stringify({
|
||||||
|
isError: typeof result === "object" && result !== null && "isError" in result && (result as { isError?: unknown }).isError === true,
|
||||||
|
content: [{ type: "text", text: resultToText(result) }],
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
} catch (err) {
|
||||||
|
res.statusCode = 200;
|
||||||
|
res.setHeader("content-type", "application/json");
|
||||||
|
res.end(
|
||||||
|
JSON.stringify({
|
||||||
|
isError: true,
|
||||||
|
content: [{ type: "text", text: err instanceof Error ? err.message : String(err) }],
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
})();
|
||||||
|
activeExecutions.add(execution);
|
||||||
|
try {
|
||||||
|
await execution;
|
||||||
|
} finally {
|
||||||
|
activeExecutions.delete(execution);
|
||||||
|
activeControllers.delete(controller);
|
||||||
|
}
|
||||||
|
})().catch(() => {
|
||||||
|
// Request streams and response sockets can abort independently of the handler.
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
let address: { port: number };
|
||||||
|
try {
|
||||||
|
address = await new Promise<{ port: number }>((resolve, reject) => {
|
||||||
|
server.once("error", reject);
|
||||||
|
// Bind loopback only — never expose Fusion tools on a public interface.
|
||||||
|
server.listen(0, "127.0.0.1", () => {
|
||||||
|
const addr = server.address();
|
||||||
|
if (!addr || typeof addr === "string") {
|
||||||
|
reject(new Error("tool bridge failed to bind"));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
resolve({ port: addr.port });
|
||||||
|
});
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
server.close();
|
||||||
|
try {
|
||||||
|
unlinkSync(schemaPath);
|
||||||
|
} catch {
|
||||||
|
// Schema cleanup is best effort after startup failure.
|
||||||
|
}
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
|
||||||
|
const bridgeUrl = `http://127.0.0.1:${address.port}`;
|
||||||
|
let disposed = false;
|
||||||
|
return {
|
||||||
|
toolCount: defs.length,
|
||||||
|
mcpServer: {
|
||||||
|
name: "fusion-custom-tools",
|
||||||
|
command: process.execPath,
|
||||||
|
args: [serverPath, schemaPath],
|
||||||
|
env: [
|
||||||
|
{ name: "FUSION_ACP_TOOL_BRIDGE_URL", value: bridgeUrl },
|
||||||
|
{ name: "FUSION_ACP_TOOL_BRIDGE_TOKEN", value: token },
|
||||||
|
],
|
||||||
|
},
|
||||||
|
dispose: async () => {
|
||||||
|
if (disposed) return;
|
||||||
|
disposed = true;
|
||||||
|
for (const controller of activeControllers) controller.abort();
|
||||||
|
server.close(() => undefined);
|
||||||
|
if (typeof server.closeAllConnections === "function") {
|
||||||
|
server.closeAllConnections();
|
||||||
|
}
|
||||||
|
await Promise.allSettled(
|
||||||
|
[...activeExecutions].map((execution) =>
|
||||||
|
Promise.race([
|
||||||
|
execution,
|
||||||
|
new Promise<void>((resolve) => {
|
||||||
|
const timer = setTimeout(resolve, TOOL_DRAIN_TIMEOUT_MS);
|
||||||
|
timer.unref?.();
|
||||||
|
}),
|
||||||
|
]),
|
||||||
|
),
|
||||||
|
);
|
||||||
|
try {
|
||||||
|
unlinkSync(schemaPath);
|
||||||
|
} catch {
|
||||||
|
// Schema may already be gone; disposal stays idempotent.
|
||||||
|
}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
@@ -116,6 +116,13 @@ export interface AgentRuntimeOptions {
|
|||||||
cwd: string;
|
cwd: string;
|
||||||
systemPrompt: string;
|
systemPrompt: string;
|
||||||
tools?: "coding" | "readonly";
|
tools?: "coding" | "readonly";
|
||||||
|
/**
|
||||||
|
* Engine-assembled Fusion custom tools (fn_*). ToolDefinition.execute closures
|
||||||
|
* only run in-process, so the ACP runtime exposes them to the agent through a
|
||||||
|
* loopback tool bridge registered in `session/new.mcpServers` (same pattern as
|
||||||
|
* the Grok runtime). Absent/empty keeps Route B's read-only ask posture.
|
||||||
|
*/
|
||||||
|
customTools?: unknown;
|
||||||
onText?: (text: string) => void;
|
onText?: (text: string) => void;
|
||||||
onThinking?: (text: string) => void;
|
onThinking?: (text: string) => void;
|
||||||
onToolStart?: (toolName: string, args?: unknown) => void;
|
onToolStart?: (toolName: string, args?: unknown) => void;
|
||||||
@@ -149,6 +156,7 @@ export interface AcpSession {
|
|||||||
/** Working directory the agent operates over (the task worktree). */
|
/** Working directory the agent operates over (the task worktree). */
|
||||||
cwd: string;
|
cwd: string;
|
||||||
lastModelDescription: string;
|
lastModelDescription: string;
|
||||||
|
fusionToolBridgeError?: { reasonCode: "mcp-schema-server-missing" | "bridge-start-failed" };
|
||||||
callbacks: AcpCallbacks;
|
callbacks: AcpCallbacks;
|
||||||
/** Per-run permission gate captured at createSession (U5/U7 read this). */
|
/** Per-run permission gate captured at createSession (U5/U7 read this). */
|
||||||
gate?: PermissionGate;
|
gate?: PermissionGate;
|
||||||
@@ -163,6 +171,10 @@ export interface AcpSession {
|
|||||||
* of each turn (FIX 1). Undefined for the bare session shell used in tests.
|
* of each turn (FIX 1). Undefined for the bare session shell used in tests.
|
||||||
*/
|
*/
|
||||||
resetTurn?: () => void;
|
resetTurn?: () => void;
|
||||||
|
/** Awaitable bridge cleanup used by AgentRuntime.dispose; absent for bare sessions. */
|
||||||
|
disposeBridge?: () => Promise<void>;
|
||||||
|
/** Completion of the most recent direct dispose call. */
|
||||||
|
disposePromise?: Promise<void>;
|
||||||
dispose(): void;
|
dispose(): void;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user