mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-24 04:31:02 +02:00
feat: MCP tool client infrastructure for agent extensibility
Add the full MCP tool pipeline enabling agents to invoke external tools (like Brave Search) via MCP servers: - Add ToolRequest/ToolResponse types and mcp-tool topics to @trustgraph/base - Create McpToolService (FlowProcessor) that connects to external MCP servers via @modelcontextprotocol/sdk StreamableHTTP transport - Add createMcpTool() to wire MCP tools into the agent's ReAct loop - Implement config-driven tool registration in AgentService with backward- compatible fallback to hardcoded tools - Add tool filtering by group and state (port of Python tool_filter.py) - Register mcp-tool in gateway dispatcher and export from @trustgraph/flow - Fix flow restart race condition: skip restart when flow definitions unchanged - Update seed config with MCP server config and tool definitions - Add run scripts for MCP tool service and Brave Search MCP server Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
f2b376abef
commit
b854b56558
17 changed files with 600 additions and 17 deletions
1
ts/packages/flow/src/agent/mcp-tool/index.ts
Normal file
1
ts/packages/flow/src/agent/mcp-tool/index.ts
Normal file
|
|
@ -0,0 +1 @@
|
|||
export { McpToolService } from "./service.js";
|
||||
160
ts/packages/flow/src/agent/mcp-tool/service.ts
Normal file
160
ts/packages/flow/src/agent/mcp-tool/service.ts
Normal file
|
|
@ -0,0 +1,160 @@
|
|||
/**
|
||||
* MCP tool-calling service — calls external MCP tool servers.
|
||||
*
|
||||
* Receives ToolRequest (name + JSON-encoded parameters) over NATS,
|
||||
* connects to the appropriate MCP server via the MCP SDK,
|
||||
* invokes the tool, and returns the result as a ToolResponse.
|
||||
*
|
||||
* MCP service configs are pushed via the config system under the "mcp" key.
|
||||
*
|
||||
* Python reference: trustgraph-flow/trustgraph/agent/mcp_tool/service.py
|
||||
*/
|
||||
|
||||
import { Client } from "@modelcontextprotocol/sdk/client/index.js";
|
||||
import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js";
|
||||
|
||||
import {
|
||||
FlowProcessor,
|
||||
ConsumerSpec,
|
||||
ProducerSpec,
|
||||
type ProcessorConfig,
|
||||
type FlowContext,
|
||||
type ToolRequest,
|
||||
type ToolResponse,
|
||||
} from "@trustgraph/base";
|
||||
|
||||
interface McpServiceConfig {
|
||||
url: string;
|
||||
"remote-name"?: string;
|
||||
"auth-token"?: string;
|
||||
}
|
||||
|
||||
export class McpToolService extends FlowProcessor {
|
||||
private mcpServices: Record<string, McpServiceConfig> = {};
|
||||
|
||||
constructor(config: ProcessorConfig) {
|
||||
super(config);
|
||||
|
||||
this.registerSpecification(
|
||||
new ConsumerSpec<ToolRequest>("request", this.onRequest.bind(this)),
|
||||
);
|
||||
this.registerSpecification(new ProducerSpec<ToolResponse>("response"));
|
||||
|
||||
this.registerConfigHandler(this.onMcpConfig.bind(this));
|
||||
}
|
||||
|
||||
private async onMcpConfig(
|
||||
config: Record<string, unknown>,
|
||||
version: number,
|
||||
): Promise<void> {
|
||||
console.log(`[McpToolService] Got config version ${version}`);
|
||||
|
||||
if (!("mcp" in config) || typeof config.mcp !== "object" || config.mcp === null) {
|
||||
this.mcpServices = {};
|
||||
return;
|
||||
}
|
||||
|
||||
const mcpConfig = config.mcp as Record<string, string>;
|
||||
this.mcpServices = {};
|
||||
|
||||
for (const [name, value] of Object.entries(mcpConfig)) {
|
||||
try {
|
||||
this.mcpServices[name] = JSON.parse(value) as McpServiceConfig;
|
||||
console.log(`[McpToolService] Registered MCP service: ${name}`);
|
||||
} catch (err) {
|
||||
console.error(`[McpToolService] Failed to parse MCP config for ${name}:`, err);
|
||||
}
|
||||
}
|
||||
|
||||
console.log(
|
||||
`[McpToolService] ${Object.keys(this.mcpServices).length} MCP services configured`,
|
||||
);
|
||||
}
|
||||
|
||||
private async onRequest(
|
||||
msg: ToolRequest,
|
||||
properties: Record<string, string>,
|
||||
flowCtx: FlowContext,
|
||||
): Promise<void> {
|
||||
const requestId = properties.id;
|
||||
if (!requestId) return;
|
||||
|
||||
const responseProducer = flowCtx.flow.producer<ToolResponse>("response");
|
||||
|
||||
try {
|
||||
const result = await this.invokeTool(
|
||||
msg.name,
|
||||
msg.parameters ? JSON.parse(msg.parameters) : {},
|
||||
);
|
||||
|
||||
if (typeof result === "string") {
|
||||
await responseProducer.send(requestId, { text: result });
|
||||
} else {
|
||||
await responseProducer.send(requestId, { object: JSON.stringify(result) });
|
||||
}
|
||||
} catch (err) {
|
||||
console.error(`[McpToolService] Error invoking tool ${msg.name}:`, err);
|
||||
const message = err instanceof Error ? err.message : String(err);
|
||||
await responseProducer.send(requestId, {
|
||||
error: { type: "tool-error", message },
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
private async invokeTool(
|
||||
name: string,
|
||||
parameters: Record<string, unknown>,
|
||||
): Promise<string | unknown> {
|
||||
if (!(name in this.mcpServices)) {
|
||||
throw new Error(`MCP service "${name}" not known`);
|
||||
}
|
||||
|
||||
const svcConfig = this.mcpServices[name];
|
||||
if (!svcConfig.url) {
|
||||
throw new Error(`MCP service "${name}" URL not defined`);
|
||||
}
|
||||
|
||||
const remoteName = svcConfig["remote-name"] ?? name;
|
||||
|
||||
// Build headers with optional bearer token
|
||||
const headers: Record<string, string> = {};
|
||||
if (svcConfig["auth-token"]) {
|
||||
headers["Authorization"] = `Bearer ${svcConfig["auth-token"]}`;
|
||||
}
|
||||
|
||||
console.log(`[McpToolService] Invoking ${remoteName} at ${svcConfig.url}`);
|
||||
|
||||
// Connect to streamable HTTP MCP server
|
||||
const transport = new StreamableHTTPClientTransport(
|
||||
new URL(svcConfig.url),
|
||||
{ requestInit: { headers } },
|
||||
);
|
||||
|
||||
const client = new Client({ name: "trustgraph-mcp-client", version: "1.0.0" });
|
||||
|
||||
try {
|
||||
await client.connect(transport);
|
||||
|
||||
const result = await client.callTool({
|
||||
name: remoteName,
|
||||
arguments: parameters,
|
||||
});
|
||||
|
||||
// Extract response — prefer structured content, fall back to text
|
||||
if (result.structuredContent) {
|
||||
return result.structuredContent;
|
||||
}
|
||||
|
||||
if (result.content && Array.isArray(result.content)) {
|
||||
return result.content
|
||||
.filter((c): c is { type: "text"; text: string } => c.type === "text")
|
||||
.map((c) => c.text)
|
||||
.join("");
|
||||
}
|
||||
|
||||
return "No content";
|
||||
} finally {
|
||||
await transport.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue