1039 lines
37 KiB
TypeScript
1039 lines
37 KiB
TypeScript
/**
|
||
* MCP Manager Service — MCP Server 生命周期管理
|
||
*
|
||
* 负责:
|
||
* 1. MCP Server 的连接/断开/重连
|
||
* 2. 工具发现与动态注册到 ToolRegistry
|
||
* 3. Server 状态管理与健康检查
|
||
*
|
||
* 使用 @modelcontextprotocol/sdk 官方库。
|
||
*
|
||
* @see docs/生产级通用 AI Agent 智能体桌面应用:完整设计与构建指南.html — 第七章
|
||
* @see standard/开发规范.md — 优先使用第三方成熟库
|
||
*/
|
||
|
||
import { Client } from '@modelcontextprotocol/sdk/client/index.js';
|
||
import { StdioClientTransport } from '@modelcontextprotocol/sdk/client/stdio.js';
|
||
import { SSEClientTransport } from '@modelcontextprotocol/sdk/client/sse.js';
|
||
// v0.4.1: streamable HTTP 传输(MCP 当前主流远程传输方式)
|
||
import { StreamableHTTPClientTransport } from '@modelcontextprotocol/sdk/client/streamableHttp.js';
|
||
import type { Tool } from '@modelcontextprotocol/sdk/types.js';
|
||
import { nanoid } from 'nanoid';
|
||
import type Database from 'better-sqlite3';
|
||
import log from 'electron-log';
|
||
import type { ToolRegistry } from '../harness/tools/registry';
|
||
import type { IMetonaTool, ToolExecutionContext } from '../harness/types/metona-tool';
|
||
import type { MetonaToolDef, MetonaParamField } from '../harness/types';
|
||
import { MetonaToolCategory, MetonaRiskLevel } from '../harness/types';
|
||
// v0.7.3 P3-2: 子进程环境净化收敛到 utils/safe-env.ts 单源(与 run_command 共用)
|
||
import { buildSafeChildEnv } from '../utils/safe-env';
|
||
import { encryptConfigValue, decryptConfigValue, isEncryptedValue } from '../utils/secure-config';
|
||
|
||
// v0.3.0 修复: 安全解析 JSON args,防止数据库中存储了非法 JSON 导致初始化崩溃
|
||
/** @visibleForTesting 纯函数,供安全表测直接断言 */
|
||
export function safeParseArgs(raw: string): string[] {
|
||
try {
|
||
const parsed = JSON.parse(raw);
|
||
return Array.isArray(parsed) ? parsed : [];
|
||
} catch {
|
||
log.warn(`MCP args parse failed, using empty array: ${raw.slice(0, 100)}`);
|
||
return [];
|
||
}
|
||
}
|
||
|
||
// #6 修复: MCP stdio 命令安全校验
|
||
/**
|
||
* 允许的 MCP Server 启动命令白名单。
|
||
* 仅允许常见的 MCP Server 运行时,防止任意命令执行。
|
||
*/
|
||
const ALLOWED_MCP_COMMANDS = new Set([
|
||
'npx',
|
||
'node',
|
||
'npm',
|
||
'python',
|
||
'python3',
|
||
'uv',
|
||
'uvx',
|
||
'bun',
|
||
'deno',
|
||
]);
|
||
|
||
/**
|
||
* #6 修复: 校验 MCP Server 的 command 和 args,防止命令注入
|
||
*
|
||
* 1. 命令白名单:只允许已知的运行时命令
|
||
* 2. 参数注入检测:拒绝包含 shell 元字符的参数
|
||
*
|
||
* @param command MCP Server 启动命令
|
||
* @param args MCP Server 启动参数
|
||
* @throws 如果命令不在白名单或参数包含 shell 元字符
|
||
*/
|
||
/** @visibleForTesting 纯函数,供安全表测直接断言 */
|
||
export function validateMcpCommand(command: string, args: string[]): void {
|
||
// 提取命令 basename(处理 /usr/bin/node、C:\node\node.exe 等路径)
|
||
const baseCmd =
|
||
command
|
||
.split(/[\\/]/)
|
||
.pop()
|
||
?.replace(/\.exe$/i, '') ?? command;
|
||
|
||
if (!ALLOWED_MCP_COMMANDS.has(baseCmd)) {
|
||
throw new Error(
|
||
`MCP command "${baseCmd}" is not in the allowed list: ${[...ALLOWED_MCP_COMMANDS].join(', ')}. ` +
|
||
`For security reasons, only standard MCP runtimes are permitted.`,
|
||
);
|
||
}
|
||
|
||
// 审查修复: 移除 {}() 字符 — StdioClientTransport 用 spawn(不经 shell),
|
||
// 这些字符无注入风险,但 MCP Server 的 args 常含 JSON 配置(如 --config {"port":3000})会被误拒。
|
||
// 保留 ; & | ` $ < > 换行 等高危字符。
|
||
const shellMetacharPattern = /[;&|`$<>\n\r]/;
|
||
for (const arg of args) {
|
||
if (shellMetacharPattern.test(arg)) {
|
||
throw new Error(
|
||
`MCP command argument contains shell metacharacters and was rejected: ${arg.slice(0, 100)}`,
|
||
);
|
||
}
|
||
}
|
||
}
|
||
|
||
/**
|
||
* #6 修复 + 审查修复: 构建安全的子进程环境变量
|
||
*
|
||
* v0.7.3 P3-2: 实现收敛到 utils/safe-env.ts(buildSafeChildEnv)——与
|
||
* run_command 共用同一黑名单(历史双实现已漂移)。MCP 侧无运行时差异注入。
|
||
* 注意: GITHUB_TOKEN / SLACK_BOT_TOKEN 等含 _TOKEN 后缀的变量会被过滤;
|
||
* 如果 MCP Server 需要这些凭证,应通过 MCP Server 配置文件传递,而非环境变量。
|
||
*/
|
||
/** @visibleForTesting 纯函数,供安全表测直接断言 */
|
||
export function buildSafeEnv(): Record<string, string> {
|
||
return buildSafeChildEnv();
|
||
}
|
||
|
||
/**
|
||
* v0.7.4 P2-5 根治: MCP headers 中敏感键(鉴权类)的值加密落库。
|
||
*
|
||
* 背景:旧实现直接把 config.headers JSON.stringify 写入 mcp_servers.headers 列,
|
||
* 绕过 secure-config —— Authorization: Bearer xxx 等远程 MCP 鉴权头明文存于本地
|
||
* DB 文件,任何人拿到 DB 即可读取所有远程 MCP 凭据。
|
||
*
|
||
* 规则:仅对敏感键(authorization / proxy-authorization / x-api-key / api-key /
|
||
* 含 token/secret/password 的键)的值走 encryptConfigValue(safeStorage 加密 +
|
||
* metona-enc:v1: 前缀);普通键(如 Accept / Content-Type)明文存储,保持
|
||
* 可读性。读取侧(parseStoredHeaders)对加密值解密。
|
||
*/
|
||
const SENSITIVE_HEADER_KEY_RE =
|
||
/^(authorization|proxy-authorization|x-api-key|api-key)$|(token|secret|password|credential)/i;
|
||
|
||
/** @visibleForTesting 纯函数:判断 header 键是否敏感(需加密值) */
|
||
export function isSensitiveHeaderKey(key: string): boolean {
|
||
return SENSITIVE_HEADER_KEY_RE.test(key);
|
||
}
|
||
|
||
/**
|
||
* v0.7.4 P2-5: 落库前序列化 headers —— 敏感键值加密。
|
||
* @param headers 运行时 headers(明文)
|
||
* @returns 可写入 DB 的 JSON 字符串(敏感值带 metona-enc:v1: 前缀)
|
||
*/
|
||
export function serializeHeadersForStorage(headers: Record<string, string>): string {
|
||
const out: Record<string, string> = {};
|
||
for (const [k, v] of Object.entries(headers)) {
|
||
out[k] = isSensitiveHeaderKey(k) ? (encryptConfigValue(v) as string) : v;
|
||
}
|
||
return JSON.stringify(out);
|
||
}
|
||
|
||
/**
|
||
* v0.7.4 P2-5: 读取 DB 后反序列化 headers —— 加密值解密为明文。
|
||
* 兼容历史明文(无前缀值原样返回)与加密格式(metona-enc:v1:)。
|
||
* @param raw DB headers 列的原始 JSON 字符串
|
||
* @returns 明文 headers(供 requestInit 注入)
|
||
*/
|
||
export function parseStoredHeaders(
|
||
raw: string | null | undefined,
|
||
): Record<string, string> | undefined {
|
||
if (!raw) return undefined;
|
||
try {
|
||
const parsed: unknown = JSON.parse(raw);
|
||
if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) return undefined;
|
||
const out: Record<string, string> = {};
|
||
for (const [k, v] of Object.entries(parsed as Record<string, unknown>)) {
|
||
if (typeof v !== 'string') {
|
||
log.warn(`MCP headers entry "${k}" dropped: value is not a string`);
|
||
continue;
|
||
}
|
||
out[k] = isEncryptedValue(v) ? (decryptConfigValue(v) as string) : v;
|
||
}
|
||
return Object.keys(out).length > 0 ? out : undefined;
|
||
} catch {
|
||
log.warn(`MCP headers parse failed, ignoring: ${raw.slice(0, 100)}`);
|
||
return undefined;
|
||
}
|
||
}
|
||
|
||
/**
|
||
* v0.7.2 P2-8: 安全解析 MCP headers 列(JSON 对象,键值均为字符串)。
|
||
*
|
||
* 数据库中存储为 JSON 字符串(mcp_servers.headers 列自建表起即存在,
|
||
* 此前从未被读写)。解析规则:
|
||
* - 空/null → undefined(匿名连接)
|
||
* - 非对象 / 数组 / 含非字符串键值 → 丢弃该条目并 WARN(保持可用性优先,
|
||
* 不因单列损坏阻断整个 server 连接)
|
||
*
|
||
* v0.7.4 P2-5: 本函数保留为"解析不含加密值的 headers"(历史明文路径/测试兼容);
|
||
* 运行时读取统一走 parseStoredHeaders(含解密)。
|
||
*/
|
||
/** @visibleForTesting 纯函数,供安全表测直接断言 */
|
||
export function safeParseHeaders(
|
||
raw: string | null | undefined,
|
||
): Record<string, string> | undefined {
|
||
return parseStoredHeaders(raw);
|
||
}
|
||
|
||
// ===== 类型定义 =====
|
||
|
||
export type MCPServerStatus =
|
||
| 'connecting'
|
||
| 'connected'
|
||
| 'disconnected'
|
||
| 'error'
|
||
| 'reconnecting';
|
||
|
||
// ===== v0.7.3 P4-2: 自动重连策略常量 =====
|
||
|
||
/** 最大自动重连次数(超过后停留 error 态,等待用户手动 toggle) */
|
||
export const MAX_RECONNECT_ATTEMPTS = 3;
|
||
|
||
/**
|
||
* 重连退避间隔(毫秒):5s / 15s / 60s。
|
||
* 纯函数 nextRetryDelayMs 消费,表测锁定(vitest fake timers 场景)。
|
||
*/
|
||
export const RECONNECT_DELAYS_MS = [5_000, 15_000, 60_000] as const;
|
||
|
||
/** @visibleForTesting 纯函数 —— 第 attempt 次(1-based)重试前的等待毫秒数 */
|
||
export function nextRetryDelayMs(attempt: number): number {
|
||
const idx = Math.min(Math.max(attempt, 1), RECONNECT_DELAYS_MS.length) - 1;
|
||
return RECONNECT_DELAYS_MS[idx];
|
||
}
|
||
|
||
export interface MCPServerConfig {
|
||
id: string;
|
||
name: string;
|
||
/** v0.4.1: 新增 'streamable-http'(MCP 当前主流远程传输);'sse' 保留向后兼容 */
|
||
transport: 'stdio' | 'sse' | 'streamable-http';
|
||
command?: string;
|
||
args?: string[];
|
||
url?: string;
|
||
/**
|
||
* v0.7.2 P2-8: 远程传输(sse / streamable-http)的自定义请求头。
|
||
* 用于 Bearer/Basic 鉴权或网关路由头;仅 HTTP 传输生效(stdio 走 buildSafeEnv)。
|
||
* 通过 requestInit 注入 —— MCP SDK 的 SSE/StreamableHTTP 客户端在
|
||
* SSE GET 流与 JSON-RPC POST 中统一合并 _commonHeaders(含 requestInit.headers)。
|
||
*/
|
||
headers?: Record<string, string>;
|
||
enabled: boolean;
|
||
}
|
||
|
||
/** v0.8.0 P2-5: MCP Resource / Prompt 发现条目(渲染端展示用宽松形态) */
|
||
export interface MCPResourceInfo {
|
||
uri: string;
|
||
name: string;
|
||
description?: string;
|
||
mimeType?: string;
|
||
}
|
||
|
||
export interface MCPPromptInfo {
|
||
name: string;
|
||
description?: string;
|
||
arguments?: Array<{ name: string; description?: string; required?: boolean }>;
|
||
}
|
||
|
||
export interface MCPServerState {
|
||
config: MCPServerConfig;
|
||
status: MCPServerStatus;
|
||
client: Client | null;
|
||
tools: Tool[];
|
||
error?: string;
|
||
connectedAt?: number;
|
||
/** v0.8.0 P2-5: 发现的 Resources / Prompts(可选能力,未声明/失败时为空数组) */
|
||
resources: MCPResourceInfo[];
|
||
prompts: MCPPromptInfo[];
|
||
}
|
||
|
||
// ===== MCP Tool Adapter =====
|
||
|
||
/**
|
||
* 将 MCP Tool 适配为 IMetonaTool 接口
|
||
*/
|
||
class MCPToolAdapter implements IMetonaTool {
|
||
readonly definition: MetonaToolDef;
|
||
|
||
constructor(
|
||
private mcpTool: Tool,
|
||
private client: Client,
|
||
private serverName: string,
|
||
) {
|
||
this.definition = {
|
||
name: `mcp_${serverName}_${mcpTool.name}`,
|
||
description: mcpTool.description ?? `MCP tool from ${serverName}`,
|
||
parameters: this.convertSchema(mcpTool.inputSchema),
|
||
category: MetonaToolCategory.MCP,
|
||
riskLevel: MetonaRiskLevel.MEDIUM,
|
||
requiresPermission: false,
|
||
timeoutMs: 30_000,
|
||
};
|
||
}
|
||
|
||
async execute(args: Record<string, unknown>, _context: ToolExecutionContext): Promise<unknown> {
|
||
const result = await this.client.callTool({
|
||
name: this.mcpTool.name,
|
||
arguments: args,
|
||
});
|
||
const content = result.content as Array<Record<string, unknown>> | undefined;
|
||
|
||
// v0.8.2 P2-7 根治: MCP 返回的 image block 此前被直接透传 —— registry 的
|
||
// 内联图片白名单只识别顶层 dataUrl/image 字段,嵌套在 content 数组中的图片
|
||
// 块不匹配白名单,被 50KB 截断为破损 base64。现将 text/image block 归并:
|
||
// text 拼接为顶层文本,首个 image block 提升为顶层 `image` 字段(data URI,
|
||
// 命中 registry 白名单整段放行 → 渲染端可内联预览)。
|
||
if (Array.isArray(content)) {
|
||
const texts: string[] = [];
|
||
let imageDataUri: string | null = null;
|
||
for (const item of content) {
|
||
const type = item?.type;
|
||
if (type === 'text' && typeof item.text === 'string') {
|
||
texts.push(item.text);
|
||
} else if (type === 'image' && !imageDataUri) {
|
||
const data = typeof item.data === 'string' ? item.data : '';
|
||
const mimeType =
|
||
typeof item.mimeType === 'string' && item.mimeType ? item.mimeType : 'image/png';
|
||
if (data) {
|
||
imageDataUri = `data:${mimeType};base64,${data}`;
|
||
}
|
||
}
|
||
}
|
||
if (imageDataUri || texts.length > 0) {
|
||
return {
|
||
...(texts.length > 0 ? { text: texts.join('\n\n') } : {}),
|
||
...(imageDataUri ? { image: imageDataUri } : {}),
|
||
};
|
||
}
|
||
}
|
||
return content;
|
||
}
|
||
|
||
/**
|
||
* 将 MCP JSON Schema 转换为 MetonaToolParams
|
||
*
|
||
* v0.8.2 P2-7 根治: 旧实现只保留顶层 properties 的 type/description ——
|
||
* 丢弃 enum/anyOf/oneOf/嵌套对象/items/default,复杂 MCP 工具的参数约束
|
||
* 对 LLM 不可见,易产生非法参数。现递归保留 IR 支持的全部结构(见
|
||
* MetonaParamField 的 P2-7 扩展字段)。
|
||
*/
|
||
private convertSchema(schema: Record<string, unknown>): MetonaToolDef['parameters'] {
|
||
const properties: Record<string, MetonaToolDef['parameters']['properties'][string]> = {};
|
||
const schemaProps = (schema.properties ?? {}) as Record<string, Record<string, unknown>>;
|
||
|
||
for (const [key, prop] of Object.entries(schemaProps)) {
|
||
properties[key] = this.convertSchemaField(prop);
|
||
}
|
||
|
||
return {
|
||
type: 'object',
|
||
properties,
|
||
required: schema.required as string[] | undefined,
|
||
};
|
||
}
|
||
|
||
/** 单个参数字段递归转换(P2-7 schema 保真) */
|
||
private convertSchemaField(prop: Record<string, unknown>): MetonaParamField {
|
||
const field: MetonaParamField = {
|
||
type: (prop.type as MetonaParamField['type']) ?? 'string',
|
||
description: (prop.description as string) ?? '',
|
||
};
|
||
if (Array.isArray(prop.enum)) {
|
||
field.enum = prop.enum.map((v) => String(v));
|
||
}
|
||
if (prop.items && typeof prop.items === 'object') {
|
||
field.items = this.convertSchemaField(prop.items as Record<string, unknown>);
|
||
}
|
||
if (prop.properties && typeof prop.properties === 'object') {
|
||
const nested: Record<string, MetonaParamField> = {};
|
||
for (const [k, v] of Object.entries(
|
||
prop.properties as Record<string, Record<string, unknown>>,
|
||
)) {
|
||
nested[k] = this.convertSchemaField(v);
|
||
}
|
||
field.properties = nested;
|
||
}
|
||
if (Array.isArray(prop.required)) {
|
||
field.required = prop.required.map((v) => String(v));
|
||
}
|
||
for (const combinator of ['anyOf', 'oneOf'] as const) {
|
||
if (Array.isArray(prop[combinator])) {
|
||
field[combinator] = (prop[combinator] as Array<Record<string, unknown>>).map((v) =>
|
||
this.convertSchemaField(v),
|
||
);
|
||
}
|
||
}
|
||
if (prop.default !== undefined) {
|
||
field.default = prop.default;
|
||
}
|
||
return field;
|
||
}
|
||
}
|
||
|
||
// ===== MCP Manager =====
|
||
|
||
export class MCPManager {
|
||
private servers = new Map<string, MCPServerState>();
|
||
|
||
// ===== v0.7.3 P4-2: 自动重连状态 =====
|
||
/** 总开关(mcp.autoReconnect,默认 true;main.ts 启动时注入,配置变更联动) */
|
||
private autoReconnect = true;
|
||
/** 各 server 的重连定时器(disconnect/shutdown 时必须清理) */
|
||
private reconnectTimers = new Map<string, NodeJS.Timeout>();
|
||
/** 各 server 已尝试的自动重连次数(成功连接后清零) */
|
||
private reconnectAttempts = new Map<string, number>();
|
||
/** 待重连的配置快照(重连时从原配置重建连接,避免读 DB 中间态) */
|
||
private reconnectConfigs = new Map<string, MCPServerConfig>();
|
||
|
||
/**
|
||
* 工具集合变更回调(v0.5.3)
|
||
*
|
||
* connectServer / disconnectServer 完成后触发,供调用方同步引擎工具列表。
|
||
* 背景:setToolsAll 只作用于已存在的引擎;MCP 工具在运行中增删时,
|
||
* 已打开会话的引擎不会自动感知 — 不回调同步则已有引擎持有失效工具定义
|
||
* (调用报 Unknown tool)或缺失新工具(与 README"无需重启"的宣称不符)。
|
||
* 懒创建引擎由 createEngine 从 registry 实时拉取(v0.5.2),无需此回调。
|
||
*/
|
||
private toolsChangedCallback: (() => void) | null = null;
|
||
|
||
constructor(
|
||
private getDB: () => Database.Database,
|
||
private toolRegistry: ToolRegistry,
|
||
) {}
|
||
|
||
// ===== v0.7.3 P4-2: 自动重连 =====
|
||
|
||
/**
|
||
* 设置自动重连开关(main.ts 启动时按 mcp.autoReconnect 注入;
|
||
* 配置变更经 shared.ts applyConfigSideEffects 联动)。
|
||
* 关闭时立即取消所有已排程的重连并清零计数(用户显式意图优先)。
|
||
*/
|
||
setAutoReconnect(enabled: boolean): void {
|
||
this.autoReconnect = enabled;
|
||
if (!enabled) {
|
||
this.cancelAllReconnects();
|
||
}
|
||
log.debug(`[MCPManager] autoReconnect = ${enabled}`);
|
||
}
|
||
|
||
/** 查询某 server 的重连状态(测试与诊断用) */
|
||
getReconnectInfo(name: string): { attempts: number; scheduled: boolean } | null {
|
||
const attempts = this.reconnectAttempts.get(name);
|
||
const scheduled = this.reconnectTimers.has(name);
|
||
if (attempts === undefined && !scheduled) return null;
|
||
return { attempts: attempts ?? 0, scheduled };
|
||
}
|
||
|
||
/** 取消某 server 的重连排程(用户显式断开/移除时调用) */
|
||
private cancelReconnect(name: string): void {
|
||
const timer = this.reconnectTimers.get(name);
|
||
if (timer) {
|
||
clearTimeout(timer);
|
||
this.reconnectTimers.delete(name);
|
||
}
|
||
this.reconnectAttempts.delete(name);
|
||
this.reconnectConfigs.delete(name);
|
||
}
|
||
|
||
/** 取消全部重连排程(shutdown / 开关关闭时调用) */
|
||
private cancelAllReconnects(): void {
|
||
for (const timer of this.reconnectTimers.values()) {
|
||
clearTimeout(timer);
|
||
}
|
||
this.reconnectTimers.clear();
|
||
this.reconnectAttempts.clear();
|
||
this.reconnectConfigs.clear();
|
||
}
|
||
|
||
/**
|
||
* 连接失败后排程指数退避重连(5s/15s/60s,最多 3 次)。
|
||
* 状态机进入 'reconnecting'(设置页可见);重试耗尽停留 'error'。
|
||
* 仅记住传入配置快照 —— 重连时按原配置重建,不读 DB 中间态。
|
||
*/
|
||
private scheduleReconnect(name: string, config: MCPServerConfig): void {
|
||
if (!this.autoReconnect) return;
|
||
|
||
const attempts = (this.reconnectAttempts.get(name) ?? 0) + 1;
|
||
if (attempts > MAX_RECONNECT_ATTEMPTS) {
|
||
log.warn(
|
||
`[MCPManager] "${name}" reconnect exhausted (${MAX_RECONNECT_ATTEMPTS} attempts) — staying in error state`,
|
||
);
|
||
this.reconnectAttempts.delete(name);
|
||
this.reconnectConfigs.delete(name);
|
||
return;
|
||
}
|
||
|
||
this.reconnectAttempts.set(name, attempts);
|
||
this.reconnectConfigs.set(name, config);
|
||
const state = this.servers.get(name);
|
||
if (state) state.status = 'reconnecting';
|
||
|
||
const delay = nextRetryDelayMs(attempts);
|
||
log.info(
|
||
`[MCPManager] "${name}" reconnect scheduled in ${delay / 1000}s (attempt ${attempts}/${MAX_RECONNECT_ATTEMPTS})`,
|
||
);
|
||
const timer = setTimeout(() => {
|
||
this.reconnectTimers.delete(name);
|
||
const snapshot = this.reconnectConfigs.get(name);
|
||
if (!snapshot) return;
|
||
log.info(
|
||
`[MCPManager] "${name}" reconnecting (attempt ${attempts}/${MAX_RECONNECT_ATTEMPTS})`,
|
||
);
|
||
void this.connectServer(snapshot).catch(() => {
|
||
/* connectServer 失败路径已自行 scheduleReconnect / 记录状态 */
|
||
});
|
||
}, delay);
|
||
// 定时器不阻塞应用退出
|
||
timer.unref?.();
|
||
this.reconnectTimers.set(name, timer);
|
||
}
|
||
|
||
/** 注册工具集合变更回调(main.ts 在 AgentEngineManager 创建后注入) */
|
||
setOnToolsChanged(callback: () => void): void {
|
||
this.toolsChangedCallback = callback;
|
||
}
|
||
|
||
/** 工具集合变更后的统一通知(连接注册完成 / 断开注销完成) */
|
||
private notifyToolsChanged(): void {
|
||
try {
|
||
this.toolsChangedCallback?.();
|
||
} catch (err) {
|
||
// 回调失败不影响 MCP 主流程(仅引擎工具列表滞后,下轮 setToolsAll 兜底)
|
||
log.warn('[MCPManager] toolsChanged callback failed:', err);
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 初始化:从数据库加载已启用的 MCP Server 并连接
|
||
*
|
||
* P1-11: 等待所有连接完成(或 30s 超时)后才返回,保证调用方广播 tools:ready
|
||
* 时 MCP 工具已实际注册完成。原实现 connectServer 不 await,工具就绪广播
|
||
* 早于工具注册,存在"广播后一段时间内 MCP 工具仍不可用"的窗口。
|
||
*
|
||
* 单个 server 连接超时(30s)用于兜底:stdio server 挂起时不阻塞整体就绪,
|
||
* 超时的 server 会显示 error 状态,用户可在设置面板查看原因。
|
||
*/
|
||
async initialize(): Promise<void> {
|
||
const db = this.getDB();
|
||
const rows = db
|
||
.prepare(
|
||
`
|
||
SELECT * FROM mcp_servers WHERE enabled = 1
|
||
`,
|
||
)
|
||
.all() as Array<{
|
||
id: string;
|
||
name: string;
|
||
transport: string;
|
||
command: string | null;
|
||
args: string | null;
|
||
url: string | null;
|
||
headers: string | null;
|
||
}>;
|
||
|
||
const connectWithTimeout = (config: MCPServerConfig): Promise<unknown> =>
|
||
Promise.race([
|
||
this.connectServer(config),
|
||
new Promise((resolve) => setTimeout(() => resolve('timeout'), 30_000)),
|
||
]);
|
||
|
||
const settled = await Promise.allSettled(
|
||
rows.map((row) => {
|
||
const config: MCPServerConfig = {
|
||
id: row.id,
|
||
name: row.name,
|
||
// v0.4.1: 支持三种传输方式(stdio / sse / streamable-http)
|
||
transport: row.transport as MCPServerConfig['transport'],
|
||
command: row.command ?? undefined,
|
||
args: row.args ? safeParseArgs(row.args) : undefined,
|
||
url: row.url ?? undefined,
|
||
// v0.7.2 P2-8: headers 列接线(此前建表即存在但从未读写)
|
||
headers: safeParseHeaders(row.headers),
|
||
enabled: true,
|
||
};
|
||
return connectWithTimeout(config).catch((err) => {
|
||
log.warn(`MCP server "${config.name}" auto-connect failed: ${err}`);
|
||
});
|
||
}),
|
||
);
|
||
|
||
const failed = settled.filter((s) => s.status === 'rejected').length;
|
||
log.info(`MCP Manager initialized: ${rows.length} server(s) configured (${failed} failed)`);
|
||
}
|
||
|
||
/**
|
||
* 连接 MCP Server
|
||
*/
|
||
async connectServer(config: MCPServerConfig): Promise<void> {
|
||
const { name } = config;
|
||
|
||
// 断开已有连接(内部拆除 —— 保留重连簿记,否则重试计数被清零、
|
||
// 退避序列永远停在第 1 次;用户显式断开走 disconnectServer)
|
||
await this.teardownConnection(name);
|
||
|
||
this.servers.set(name, {
|
||
config,
|
||
status: 'connecting',
|
||
client: null,
|
||
tools: [],
|
||
resources: [],
|
||
prompts: [],
|
||
});
|
||
|
||
try {
|
||
let transport;
|
||
|
||
if (config.transport === 'stdio' && config.command) {
|
||
// stdio 模式
|
||
const args = config.args ?? [];
|
||
// #6 修复: 命令白名单 + 参数元字符检测,防止命令注入
|
||
validateMcpCommand(config.command, args);
|
||
transport = new StdioClientTransport({
|
||
command: config.command,
|
||
args,
|
||
// #6 修复: 不透传完整 process.env,仅保留 MCP Server 运行所需的最小环境变量
|
||
env: buildSafeEnv(),
|
||
});
|
||
} else if (config.transport === 'sse' && config.url) {
|
||
// SSE 模式(远程 HTTP,旧式传输,保留向后兼容)
|
||
// v0.7.2 P2-8: 注入自定义请求头(非空时)—— requestInit.headers 同时
|
||
// 覆盖 SSE GET 流与 JSON-RPC POST(SDK _commonHeaders 统一合并点)
|
||
const headerInit = {
|
||
...(config.headers && Object.keys(config.headers).length > 0
|
||
? { requestInit: { headers: config.headers } }
|
||
: {}),
|
||
};
|
||
transport = new SSEClientTransport(new URL(config.url), headerInit);
|
||
} else if (config.transport === 'streamable-http' && config.url) {
|
||
// v0.4.1: streamable HTTP 模式(MCP 当前主流远程传输)
|
||
const headerInit = {
|
||
...(config.headers && Object.keys(config.headers).length > 0
|
||
? { requestInit: { headers: config.headers } }
|
||
: {}),
|
||
};
|
||
transport = new StreamableHTTPClientTransport(new URL(config.url), headerInit);
|
||
} else {
|
||
throw new Error(
|
||
`Unsupported transport "${config.transport}". ` +
|
||
`'stdio' requires 'command', 'sse'/'streamable-http' requires 'url'.`,
|
||
);
|
||
}
|
||
|
||
const client = new Client(
|
||
{ name: 'metona-ai-desktop', version: '1.0.0' },
|
||
{ capabilities: {} },
|
||
);
|
||
|
||
await client.connect(transport);
|
||
|
||
// 发现工具
|
||
const toolsResult = await client.listTools();
|
||
const tools = toolsResult.tools ?? [];
|
||
|
||
// v0.8.0 P2-5: Resources / Prompts 发现(可选能力 —— server 未声明或
|
||
// 请求失败一律置空数组,不阻断连接与工具注册)
|
||
const discoveredResources: MCPResourceInfo[] = [];
|
||
const discoveredPrompts: MCPPromptInfo[] = [];
|
||
try {
|
||
const res = await client.listResources();
|
||
for (const r of res.resources ?? []) {
|
||
discoveredResources.push({
|
||
uri: String(r.uri ?? ''),
|
||
name: String(r.name ?? ''),
|
||
description: typeof r.description === 'string' ? r.description : undefined,
|
||
mimeType: typeof r.mimeType === 'string' ? r.mimeType : undefined,
|
||
});
|
||
}
|
||
} catch {
|
||
/* server 未声明 resources 能力 */
|
||
}
|
||
try {
|
||
const res = await client.listPrompts();
|
||
for (const p of res.prompts ?? []) {
|
||
discoveredPrompts.push({
|
||
name: String(p.name ?? ''),
|
||
description: typeof p.description === 'string' ? p.description : undefined,
|
||
arguments: Array.isArray(p.arguments)
|
||
? p.arguments.map((a) => ({
|
||
name: String(a.name ?? ''),
|
||
description: typeof a.description === 'string' ? a.description : undefined,
|
||
required: a.required === true,
|
||
}))
|
||
: undefined,
|
||
});
|
||
}
|
||
} catch {
|
||
/* server 未声明 prompts 能力 */
|
||
}
|
||
|
||
// 注册到 ToolRegistry(v0.6.4: registerMCP 对重名冲突返回 false,此处聚合上报)
|
||
let registeredCount = 0;
|
||
let skippedCount = 0;
|
||
for (const tool of tools) {
|
||
const adapter = new MCPToolAdapter(tool, client, name);
|
||
if (this.toolRegistry.registerMCP(name, adapter)) {
|
||
registeredCount++;
|
||
} else {
|
||
skippedCount++;
|
||
}
|
||
}
|
||
|
||
// v0.5.3: 工具集合已变化 — 通知调用方同步引擎工具列表(已存在引擎热更新)
|
||
this.notifyToolsChanged();
|
||
|
||
// 更新状态
|
||
const state = this.servers.get(name)!;
|
||
state.status = 'connected';
|
||
state.client = client;
|
||
state.tools = tools;
|
||
state.connectedAt = Date.now();
|
||
state.error = undefined;
|
||
state.resources = discoveredResources;
|
||
state.prompts = discoveredPrompts;
|
||
|
||
// v0.7.3 P4-2: 连接成功 —— 清零重连计数并取消排程
|
||
this.reconnectAttempts.delete(name);
|
||
this.reconnectConfigs.delete(name);
|
||
const pendingTimer = this.reconnectTimers.get(name);
|
||
if (pendingTimer) {
|
||
clearTimeout(pendingTimer);
|
||
this.reconnectTimers.delete(name);
|
||
}
|
||
|
||
// 更新数据库
|
||
const db = this.getDB();
|
||
db.prepare(
|
||
`
|
||
UPDATE mcp_servers SET last_connected = ?, error_message = NULL WHERE name = ?
|
||
`,
|
||
).run(Date.now(), name);
|
||
|
||
log.info(
|
||
`MCP server "${name}" connected: ${registeredCount} tool(s) registered` +
|
||
(skippedCount > 0 ? `, ${skippedCount} skipped due to name conflicts` : ''),
|
||
);
|
||
} catch (error) {
|
||
const state = this.servers.get(name);
|
||
if (state) {
|
||
state.status = 'error';
|
||
state.error = (error as Error).message;
|
||
}
|
||
|
||
// 更新数据库
|
||
const db = this.getDB();
|
||
db.prepare(
|
||
`
|
||
UPDATE mcp_servers SET error_message = ? WHERE name = ?
|
||
`,
|
||
).run((error as Error).message, name);
|
||
|
||
// v0.7.3 P4-2: 失败后排程指数退避自动重连(开关关闭时 no-op)
|
||
this.scheduleReconnect(name, config);
|
||
|
||
log.error(`MCP server "${name}" connection failed:`, error);
|
||
throw error;
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 断开 MCP Server
|
||
*/
|
||
/**
|
||
* 内部连接拆除(保留重连簿记)—— connectServer 重连前的清理动作。
|
||
* 与 disconnectServer 的区别:不清 reconnectAttempts/Timers/Configs,
|
||
* 否则自动重连的每次重试都会把自己的计数清零(退避序列永远停在第 1 次)。
|
||
*/
|
||
private async teardownConnection(name: string): Promise<void> {
|
||
const state = this.servers.get(name);
|
||
if (!state) return;
|
||
this.toolRegistry.unregisterMCPTools(name);
|
||
if (state.client) {
|
||
try {
|
||
await state.client.close();
|
||
} catch {
|
||
// 忽略关闭错误
|
||
}
|
||
}
|
||
state.client = null;
|
||
state.tools = [];
|
||
state.resources = [];
|
||
state.prompts = [];
|
||
this.notifyToolsChanged();
|
||
}
|
||
|
||
/**
|
||
* 断开 MCP Server(用户显式语义:取消重连排程 + 拆除连接)
|
||
*/
|
||
async disconnectServer(name: string): Promise<void> {
|
||
// v0.7.3 P4-2: 用户显式断开/移除 —— 取消重连排程(用户意图优先于自动重试)。
|
||
// 无论是否存在连接态(error/reconnecting 态的 server 也可能被移除)都执行。
|
||
this.cancelReconnect(name);
|
||
await this.teardownConnection(name);
|
||
|
||
const state = this.servers.get(name);
|
||
if (state) {
|
||
state.status = 'disconnected';
|
||
}
|
||
|
||
log.info(`MCP server "${name}" disconnected`);
|
||
}
|
||
|
||
/**
|
||
* 切换 Server 启用/禁用
|
||
*/
|
||
async toggleServer(name: string, enabled: boolean): Promise<void> {
|
||
const db = this.getDB();
|
||
db.prepare(
|
||
`
|
||
UPDATE mcp_servers SET enabled = ?, updated_at = ? WHERE name = ?
|
||
`,
|
||
).run(enabled ? 1 : 0, Date.now(), name);
|
||
|
||
if (enabled) {
|
||
const row = db.prepare('SELECT * FROM mcp_servers WHERE name = ?').get(name) as
|
||
| {
|
||
id: string;
|
||
name: string;
|
||
transport: string;
|
||
command: string | null;
|
||
args: string | null;
|
||
url: string | null;
|
||
headers: string | null;
|
||
}
|
||
| undefined;
|
||
if (row) {
|
||
await this.connectServer({
|
||
id: row.id,
|
||
name: row.name,
|
||
transport: row.transport as MCPServerConfig['transport'],
|
||
command: row.command ?? undefined,
|
||
args: row.args ? safeParseArgs(row.args) : undefined,
|
||
url: row.url ?? undefined,
|
||
headers: safeParseHeaders(row.headers),
|
||
enabled: true,
|
||
});
|
||
}
|
||
} else {
|
||
await this.disconnectServer(name);
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 添加新的 MCP Server
|
||
*
|
||
* v0.4.1: 校验 transport 与对应字段匹配(stdio→command,sse/streamable-http→url)
|
||
* v0.7.2 P2-8: headers 持久化(JSON 列;仅远程传输消费,stdio 忽略)
|
||
*/
|
||
async addServer(config: Omit<MCPServerConfig, 'id'>): Promise<void> {
|
||
const db = this.getDB();
|
||
const id = `mcp_${nanoid(8)}`;
|
||
|
||
// 校验传输方式与必填字段
|
||
if (config.transport === 'stdio' && !config.command) {
|
||
throw new Error('stdio transport requires "command"');
|
||
}
|
||
if ((config.transport === 'sse' || config.transport === 'streamable-http') && !config.url) {
|
||
throw new Error(`${config.transport} transport requires "url"`);
|
||
}
|
||
|
||
db.prepare(
|
||
`
|
||
INSERT INTO mcp_servers (id, name, transport, command, args, url, headers, enabled)
|
||
VALUES (?, ?, ?, ?, ?, ?, ?, 1)
|
||
`,
|
||
).run(
|
||
id,
|
||
config.name,
|
||
config.transport,
|
||
config.command ?? null,
|
||
config.args ? JSON.stringify(config.args) : null,
|
||
config.url ?? null,
|
||
// v0.7.4 P2-5: 敏感键值加密后落库(Authorization/x-api-key 等不再明文存 DB)
|
||
config.headers && Object.keys(config.headers).length > 0
|
||
? serializeHeadersForStorage(config.headers)
|
||
: null,
|
||
);
|
||
|
||
if (config.enabled !== false) {
|
||
await this.connectServer({ ...config, id, enabled: true });
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 移除 MCP Server
|
||
*/
|
||
async removeServer(name: string): Promise<void> {
|
||
await this.disconnectServer(name);
|
||
const db = this.getDB();
|
||
db.prepare('DELETE FROM mcp_servers WHERE name = ?').run(name);
|
||
this.servers.delete(name);
|
||
log.info(`MCP server "${name}" removed`);
|
||
}
|
||
|
||
/**
|
||
* 获取所有 Server 状态
|
||
*
|
||
* v0.8.2 P2-7 根治: 已配置但**禁用**的 server 此前不出现在列表中 ——
|
||
* initialize() 只连 enabled=1,getServerStates 只映射 servers Map,设置面板
|
||
* 无法展示"已配置但禁用"的完整清单。现从 DB 补齐缺失条目(status=disconnected、
|
||
* enabled=false),并为每个条目标注 enabled。
|
||
*/
|
||
getServerStates(): Array<{
|
||
name: string;
|
||
status: MCPServerStatus;
|
||
toolCount: number;
|
||
error?: string;
|
||
/** v0.7.3 P4-2: reconnecting 状态下的已尝试次数(第 N/3 次排程) */
|
||
reconnectAttempt?: number;
|
||
/** v0.8.2 P2-7: 是否为启用状态(DB enabled=1)—— 禁用 server 以 disconnected 呈现 */
|
||
enabled: boolean;
|
||
}> {
|
||
// DB 全量配置(enabled 标注 + 补齐禁用条目);DB 不可用时退回仅已连接集合
|
||
let enabledMap = new Map<string, boolean>();
|
||
try {
|
||
const db = this.getDB();
|
||
const rows = db.prepare('SELECT name, enabled FROM mcp_servers').all() as Array<{
|
||
name: string;
|
||
enabled: number;
|
||
}>;
|
||
enabledMap = new Map(rows.map((r) => [r.name, r.enabled === 1]));
|
||
} catch (err) {
|
||
log.warn(`[MCPManager] getServerStates DB lookup failed: ${(err as Error).message}`);
|
||
}
|
||
|
||
const states = Array.from(this.servers.values()).map((s) => ({
|
||
name: s.config.name,
|
||
status: s.status,
|
||
toolCount: s.tools.length,
|
||
error: s.error,
|
||
reconnectAttempt:
|
||
s.status === 'reconnecting' ? this.reconnectAttempts.get(s.config.name) : undefined,
|
||
enabled: enabledMap.get(s.config.name) ?? true,
|
||
}));
|
||
|
||
for (const [name, enabled] of enabledMap) {
|
||
if (!enabled && !states.some((s) => s.name === name)) {
|
||
states.push({
|
||
name,
|
||
status: 'disconnected' as MCPServerStatus,
|
||
toolCount: 0,
|
||
error: undefined,
|
||
reconnectAttempt: undefined,
|
||
enabled: false,
|
||
});
|
||
}
|
||
}
|
||
return states;
|
||
}
|
||
|
||
/**
|
||
* 获取单个 Server 状态
|
||
*/
|
||
getServerState(name: string): {
|
||
name: string;
|
||
status: MCPServerStatus;
|
||
toolCount: number;
|
||
error?: string;
|
||
reconnectAttempt?: number;
|
||
} | null {
|
||
const state = this.servers.get(name);
|
||
if (!state) return null;
|
||
return {
|
||
name: state.config.name,
|
||
status: state.status,
|
||
toolCount: state.tools.length,
|
||
error: state.error,
|
||
reconnectAttempt:
|
||
state.status === 'reconnecting' ? this.reconnectAttempts.get(state.config.name) : undefined,
|
||
};
|
||
}
|
||
|
||
/**
|
||
* v0.8.0 P2-5: 获取 server 发现的 Resources / Prompts(未连接返回空集合)。
|
||
*/
|
||
getServerContents(name: string): { resources: MCPResourceInfo[]; prompts: MCPPromptInfo[] } {
|
||
const state = this.servers.get(name);
|
||
if (!state || state.status !== 'connected') {
|
||
return { resources: [], prompts: [] };
|
||
}
|
||
return { resources: state.resources ?? [], prompts: state.prompts ?? [] };
|
||
}
|
||
|
||
/**
|
||
* v0.8.1 P1-4: 获取 MCP Prompt 渲染结果(prompts/get)—— 供 ChatInput 斜杠
|
||
* 菜单填充输入框。messages 展平为文本(text 块拼接;role 前缀保留多消息语义)。
|
||
*/
|
||
async getPrompt(
|
||
serverName: string,
|
||
promptName: string,
|
||
args?: Record<string, string>,
|
||
): Promise<{ text: string; description?: string } | null> {
|
||
const state = this.servers.get(serverName);
|
||
if (!state || state.status !== 'connected' || !state.client) {
|
||
throw new Error(`MCP server '${serverName}' is not connected`);
|
||
}
|
||
const res = await state.client.getPrompt({
|
||
name: promptName,
|
||
arguments: args,
|
||
});
|
||
const parts: string[] = [];
|
||
for (const m of res.messages ?? []) {
|
||
// MCP PromptMessage.content 是单个 ContentBlock(text/image/audio/resource_link
|
||
// 联合)—— 仅提取文本块;经 unknown 中转以匹配 SDK 联合类型
|
||
const content = m.content as unknown as { type?: string; text?: string } | undefined;
|
||
const text = content?.type === 'text' ? (content.text ?? '') : '';
|
||
parts.push(text);
|
||
}
|
||
return { text: parts.filter((t) => t.length > 0).join('\n\n'), description: res.description };
|
||
}
|
||
|
||
/**
|
||
* v0.8.1 P1-4: 读取 MCP Resource 内容(resources/read)—— 供 @mcp 提及注入。
|
||
* 文本类内容展平返回;二进制(blob)内容拒绝(与附件管线"二进制拒绝"口径一致)。
|
||
*/
|
||
async readResource(
|
||
serverName: string,
|
||
uri: string,
|
||
): Promise<{ text: string; mimeType?: string } | null> {
|
||
const state = this.servers.get(serverName);
|
||
if (!state || state.status !== 'connected' || !state.client) {
|
||
throw new Error(`MCP server '${serverName}' is not connected`);
|
||
}
|
||
const res = await state.client.readResource({ uri });
|
||
const contents = res.contents ?? [];
|
||
const first = contents[0];
|
||
if (!first) return null;
|
||
if ('blob' in first && typeof first.blob === 'string') {
|
||
throw new Error(`Resource '${uri}' is binary content — inline injection is not supported`);
|
||
}
|
||
return {
|
||
text: ('text' in first ? first.text : '') ?? '',
|
||
mimeType: first.mimeType,
|
||
};
|
||
}
|
||
|
||
/**
|
||
* 关闭所有连接
|
||
*/
|
||
async shutdown(): Promise<void> {
|
||
// v0.7.3 P4-2: 退出前取消全部重连排程(timer 已 unref,此处幂等清理)
|
||
this.cancelAllReconnects();
|
||
const names = Array.from(this.servers.keys());
|
||
await Promise.allSettled(names.map((n) => this.disconnectServer(n)));
|
||
log.info('MCP Manager shut down');
|
||
}
|
||
}
|