Skip to content

Latest commit

 

History

History
1953 lines (1554 loc) · 79.8 KB

File metadata and controls

1953 lines (1554 loc) · 79.8 KB

MCPilot 通用AI工具接入平台 系统设计文档

文档信息

项 内容
版本 v2.1.0
日期 2026年8月18日
技术栈 Python 3.11+ / FastAPI / Vue 3 / PostgreSQL
文档状态 已定稿

1. 项目背景与目标

1.1 项目背景

当前AI应用生态中,MCP(Model Context Protocol)服务逐渐成为工具调用的标准协议,但多个MCP服务分散管理,缺乏统一入口。服务更新通常需要重启主服务,影响在线用户的正常使用。同时,运维人员缺乏可视化监控和调试工具,无法统一掌握服务状态与调用统计。

MCPilot 通用AI工具接入平台应运而生,旨在提供一个统一的MCP服务管理与接入网关。平台采用适配器架构,支持多种AI系统的协议接入,通过WebSocket网关实现实时通信,并具备热插拔能力,使配置变更秒级生效而无需重启主服务。

1.2 核心目标

编号 目标 说明
1 热插拔能力 MCP服务动态注册与注销,无需重启主服务
2 统一管理 一站式平台管理所有MCP服务的生命周期
3 可视化操作 Web界面管理,降低使用门槛
4 协议兼容 兼容MCP标准协议,支持多种AI系统接入
5 易扩展 模块化适配器设计,便于接入新的AI类型

1.3 设计决策记录

决策项 决策内容 理由
项目名称 MCPilot(智控领航) 体现控制与领航含义,呼应飞行员角色
开发语言 Python 生态丰富、开发速度快、团队熟悉度高
Web框架 FastAPI 高性能异步、自动文档、类型安全、WebSocket原生支持
架构模式 适配器模式 支持多AI系统统一接入,易于扩展新AI类型
进程管理 asyncio子进程 Python原生异步支持、跨平台
配置方式 YAML文件 + 环境变量 易于版本管理、支持容器化部署
热插拔实现 文件监听 + 动态注册 无需重启主服务,配置变更秒级生效

2. 系统架构

2.1 整体架构图

┌─────────────────────────────────────────────────────────┐
│                    AI 客户端                             │
│              (WebSocket JSON-RPC 接入)                   │
└────────────────────────────┬────────────────────────────┘
                             │
┌────────────────────────────▼────────────────────────────┐
│                 AI 适配层 (Adapters)                     │
│    协议转换 · 认证解析 · 请求/响应编解码                 │
└────────────────────────────┬────────────────────────────┘
                             │
┌────────────────────────────▼────────────────────────────┐
│              WebSocket 网关层 (Gateway)                  │
│    连接管理 · 消息分发 · 心跳保活                        │
└────────────────────────────┬────────────────────────────┘
                             │
┌────────────────────────────▼────────────────────────────┐
│              路由与分发层 (Router)                       │
│    工具名匹配 · 服务实例路由                             │
└────────────────────────────┬────────────────────────────┘
                             │
┌────────────────────────────▼────────────────────────────┐
│            MCP 进程管理层 (ProcessManager)               │
│    子进程生命周期 · stdio管道通信 · 健康检查             │
└────────────────────────────┬────────────────────────────┘
                             │
        ┌────────────────────┴────────────────────┐
        │                                         │
┌───────▼──────────┐                     ┌────────▼──────────┐
│  MCP 服务 A      │                     │  MCP 服务 B       │
│  (音乐播放)      │                     │  (新闻资讯)       │
└──────────────────┘                     └───────────────────┘
        │                                         │
┌───────▼──────────┐                     ┌────────▼──────────┐
│  MCP 服务 C      │                     │  MCP 服务 D       │
│  (智能家居)      │                     │  (知识库问答)     │
└──────────────────┘                     └───────────────────┘

2.2 分层设计

系统采用四层架构,自上而下依次为AI适配层、网关层、核心业务层和API层。各层职责清晰分离,通过定义良好的接口交互,支持独立演进与测试。

AI 适配层 (Adapters Layer)

处理不同AI系统的协议差异,将各类协议统一转换为平台内部的标准格式。该层包含适配器基类 BaseAIAdapter、适配器工厂 AdapterFactory 以及各AI系统的具体适配器实现。适配器通过 @register_adapter 装饰器自动注册到工厂,新AI类型的接入只需实现基类接口并添加装饰器即可。

网关层 (Gateway Layer)

负责WebSocket连接管理、消息分发和心跳保活。连接管理器维护 client_id 到 (websocket, adapter) 的映射,消息分发器将请求路由到对应适配器处理,并支持心跳保活机制。消息循环处理包含五个步骤:接收消息、适配器解析请求、调用工具、适配器编码响应、发送回客户端。

核心业务层 (Core Layer)

承载进程管理、服务路由和热插拔三大核心能力。ProcessManager 管理MCP子进程的完整生命周期,包括启动、停止、重启以及健康检测和自动重启。Router 维护工具名称到服务ID的映射表,支持动态注册和注销。Hotplug 通过文件系统监听实现配置变更检测,驱动服务的自动发现和注册。

API 层 (API Layer)

提供RESTful接口管理和WebSocket端点。REST API覆盖服务管理(列表、启停、重启)和工具调用(列表、调用)两大功能域。WebSocket端点路径为 /mcp/ws,通过查询参数传递Token和AI类型信息。

2.3 核心组件职责表

层级 组件 职责 关键接口
适配层 BaseAIAdapter 定义统一适配器接口规范 parse_connection / parse_request / encode_response
适配层 AdapterFactory 创建和管理适配器实例 create_adapter / register_adapter
网关层 ConnectionManager WebSocket连接生命周期管理 connect / disconnect / send_response
网关层 MessageDispatcher 消息分发与请求处理 handle_request
核心层 ProcessManager MCP子进程生命周期管理 start_service / stop_service / get_service
核心层 MCPProcess 单个MCP进程封装与IO通信 start / stop / send_request
核心层 ToolRouter 工具名称到服务的路由映射 register_service / route / get_all_tools
核心层 HotplugManager 配置文件监听与热更新 start / stop
API层 ServicesAPI 服务管理REST接口 GET/POST /api/v1/services
API层 ToolsAPI 工具调用REST接口 POST /api/v1/services/tools/call

3. 技术选型

3.1 后端技术栈

组件 技术 版本 选型理由
语言 Python 3.11+ 生态丰富,团队熟悉,3.11性能大幅提升
Web框架 FastAPI 0.110+ 高性能、类型安全、自动生成API文档
ASGI服务器 Uvicorn 0.30+ WebSocket原生支持,性能优秀
WebSocket websockets 12+ FastAPI原生支持,标准库级实现
ORM SQLAlchemy 2.0+ 异步支持完善,Python生态最成熟
数据库 PostgreSQL 17+ 生产级关系型数据库,高并发支持完善
配置管理 Pydantic Settings 2.5+ 类型安全配置,支持环境变量覆盖
文件监听 Watchdog 4.0+ Python最成熟的文件系统监听库
日志 structlog 24.4+ 结构化日志,JSON输出,与logging生态兼容
HTTP客户端 httpx 0.27+ 异步HTTP客户端,支持HTTP/2
任务队列 asyncio Queue - Python标准库,满足初期需求

可选组件根据规模扩展按需引入。当单机MCP服务超过50个时引入Redis缓存,生产环境部署时接入Supervisor进程监控,需要可视化监控时集成Prometheus指标采集,异步任务队列需求出现时引入Celery + Redis。

3.2 前端技术栈

组件 技术 版本 选型理由
框架 Vue 3.4+ 渐进式框架、学习曲线平缓、生态成熟
构建工具 Vite 5+ 超快冷启动、热更新体验优秀
语言 TypeScript 5+ 类型安全、IDE体验好
UI组件库 Element Plus 2.6+ 组件丰富、文档完善、国内生态好
HTTP客户端 Axios 1.6+ 拦截器、取消请求等功能完整
状态管理 Pinia 2.1+ Vue官方推荐,比Vuex更简洁

3.3 开发工具链

# 代码质量
pip install ruff       # 超快速Lint + Format
pip install mypy       # 类型检查

# 测试框架
pip install pytest     # 测试框架
pip install pytest-asyncio  # 异步测试支持
pip install httpx      # HTTP客户端测试

# 包管理
pip install poetry     # 依赖管理与打包

3.4 项目目录结构

mcpilot/
├── backend/                          # FastAPI后端
│   ├── app/
│   │   ├── main.py                   # 应用入口
│   │   ├── config.py                 # 配置管理
│   │   ├── core/                     # 核心业务逻辑
│   │   │   ├── process_manager.py    # 进程管理器
│   │   │   ├── router.py             # 工具路由引擎
│   │   │   └── hotplug.py            # 热插拔管理器
│   │   ├── adapters/                 # AI适配器
│   │   │   ├── base.py               # 适配器基类
│   │   │   └── factory.py            # 适配器工厂
│   │   ├── gateway/                  # WebSocket网关
│   │   │   ├── server.py             # 连接管理
│   │   │   └── protocol.py           # JSON-RPC协议实现
│   │   ├── api/                      # REST API
│   │   │   ├── services.py           # 服务管理API
│   │   │   └── tools.py              # 工具调用API
│   │   ├── models/                   # 数据模型
│   │   │   ├── schemas.py            # Pydantic请求/响应模型
│   │   │   └── database.py           # SQLAlchemy ORM模型
│   │   └── db/                       # 数据库
│   │       ├── base.py               # 数据库连接
│   │       └── migrations.py         # 数据迁移
│   ├── configs/                      # 配置文件目录
│   │   ├── config.yaml               # 主配置
│   │   └── mcp-services/             # MCP服务配置目录
│   ├── tests/                        # 测试目录
│   │   ├── unit/                     # 单元测试
│   │   └── integration/              # 集成测试
│   ├── pyproject.toml                # Python依赖配置
│   └── Dockerfile                    # 后端Dockerfile
│
├── frontend/                         # Vue3前端
│   ├── src/
│   │   ├── views/                    # 页面组件
│   │   ├── components/               # 公共组件
│   │   ├── api/                      # API客户端
│   │   └── types/                    # TypeScript类型定义
│   ├── package.json
│   └── Dockerfile                    # 前端Dockerfile
│
├── docker-compose.yml                # 一键启动
├── samples/                          # 核心示例代码
│   ├── adapters/                     # AI适配器示例
│   │   ├── base.py                   # 适配器基类
│   │   ├── factory.py                # 适配器工厂
│   │   ├── xiaozhi.py                # 小智AI适配器
│   │   └── openai.py                 # OpenAI适配器
│   └── gateway/
│       └── multi_ai_gateway.py       # 多AI统一网关
├── docs/                             # 文档目录
│   ├── 系统设计.md              # 系统设计文档
│   ├── 开发计划.md           # 开发计划
│   ├── API参考.md              # API参考文档
│   ├── 开发者指南.md          # 开发者指南
│   └── 部署指南.md           # 部署指南
└── README.md                         # 项目说明

4. 核心模块设计

4.1 配置管理模块

配置管理模块基于 Pydantic Settings 实现,支持YAML文件加载与环境变量覆盖。配置按功能域划分为服务器、数据库、MCP服务和日志四组,每组独立定义默认值与类型约束。环境变量通过 __ 嵌套分隔符覆盖YAML配置,例如 SERVER__TOKEN 覆盖 server.token 字段,满足容器化部署的配置注入需求。

from pydantic_settings import BaseSettings, SettingsConfigDict
from pydantic import Field
from typing import Optional
import yaml


class ServerSettings(BaseSettings):
    http_host: str = "0.0.0.0"
    http_port: int = 8000
    token: str = "your-secret-token-here"
    debug: bool = False


class DatabaseSettings(BaseSettings):
    type: str = "postgresql"
    url: str = "postgresql+asyncpg://postgres:abc123@192.168.31.210:5432/mcpilot"
    echo: bool = False


class MCPSettings(BaseSettings):
    services_dir: str = "./configs/mcp-services"
    health_check_interval: int = 30
    health_check_timeout: int = 10


class LoggingSettings(BaseSettings):
    level: str = "INFO"
    format: str = "json"


class Settings(BaseSettings):
    model_config = SettingsConfigDict(
        env_file=".env",
        env_nested_delimiter="__",
        case_sensitive=False,
        extra="ignore"
    )

    server: ServerSettings = Field(default_factory=ServerSettings)
    database: DatabaseSettings = Field(default_factory=DatabaseSettings)
    mcp: MCPSettings = Field(default_factory=MCPSettings)
    logging: LoggingSettings = Field(default_factory=LoggingSettings)

    @classmethod
    def from_yaml(cls, config_path: str = "configs/config.yaml") -> "Settings":
        try:
            with open(config_path, 'r', encoding='utf-8') as f:
                config_data = yaml.safe_load(f)
            return cls(**config_data)
        except FileNotFoundError:
            return cls()


# 全局配置实例
settings = Settings.from_yaml()

主配置文件示例如下:

server:
  http_host: "0.0.0.0"
  http_port: 8000
  token: "your-secret-token-here"
  debug: false

database:
  type: postgresql
  url: "postgresql+asyncpg://postgres:abc123@192.168.31.210:5432/mcpilot"
  echo: false

mcp:
  services_dir: "./configs/mcp-services"
  health_check_interval: 30
  health_check_timeout: 10

logging:
  level: INFO
  format: json

4.2 MCP协议模块

MCP协议模块封装JSON-RPC 2.0协议的编解码逻辑,定义MCP服务配置、工具定义和传输配置的数据模型。所有模型基于Pydantic BaseModel,在构造时自动校验类型与格式。编解码函数负责JSON序列化与换行符拼接,确保与MCP子进程的stdio通信格式一致。

from pydantic import BaseModel, Field, field_validator
from typing import Any, Optional, Dict, List
import json


class JSONRPCError(BaseModel):
    code: int
    message: str
    data: Optional[Any] = None


class JSONRPCRequest(BaseModel):
    jsonrpc: str = "2.0"
    id: Optional[Any] = None
    method: str
    params: Optional[Dict[str, Any]] = None

    @field_validator('jsonrpc')
    @classmethod
    def check_jsonrpc_version(cls, v):
        if v != "2.0":
            raise ValueError("Only JSON-RPC 2.0 is supported")
        return v


class JSONRPCResponse(BaseModel):
    jsonrpc: str = "2.0"
    id: Optional[Any] = None
    result: Optional[Any] = None
    error: Optional[JSONRPCError] = None


class MCPToolInputSchema(BaseModel):
    type: str = "object"
    properties: Dict[str, Any] = Field(default_factory=dict)
    required: List[str] = Field(default_factory=list)


class MCPTool(BaseModel):
    name: str
    description: str
    inputSchema: MCPToolInputSchema


class MCPTransportConfig(BaseModel):
    type: str = "stdio"  # stdio / http / websocket
    command: Optional[str] = None
    args: List[str] = Field(default_factory=list)
    env: Dict[str, str] = Field(default_factory=dict)
    url: Optional[str] = None
    timeout: int = 30


class MCPServiceConfig(BaseModel):
    id: str
    name: str
    version: str
    transport: MCPTransportConfig
    tools: List[MCPTool] = Field(default_factory=list)
    status: str = "stopped"  # registered / starting / running / error / stopped
    pid: Optional[int] = None
    started_at: Optional[float] = None


# 编解码函数
def encode_request(request: JSONRPCRequest) -> bytes:
    return (json.dumps(request.model_dump(), ensure_ascii=False) + "\n").encode('utf-8')


def decode_response(data: bytes) -> JSONRPCResponse:
    return JSONRPCResponse.model_validate_json(data.strip())


def decode_request(data: bytes) -> JSONRPCRequest:
    return JSONRPCRequest.model_validate_json(data.strip())


def encode_response(response: JSONRPCResponse) -> bytes:
    return (json.dumps(response.model_dump(), ensure_ascii=False) + "\n").encode('utf-8')

4.3 进程管理器

进程管理器是平台的核心组件,负责MCP子进程的完整生命周期管理。MCPProcess 封装单个子进程的启动、停止、请求发送和IO监控,通过asyncio子进程实现非阻塞的stdin/stdout/stderr管道通信。ProcessManager 作为全局单例管理所有进程实例,提供并发安全的启动、停止和查询接口。

MCPProcess 内部维护请求ID与Future的映射表,实现请求-响应的异步匹配。stdout读取协程持续监听子进程输出,将JSON-RPC响应通过ID匹配到对应的Future并设置结果。所有请求均带有超时保护,超时后自动清理Pending状态。

class MCPProcess:
    def __init__(self, config: MCPServiceConfig):
        self.config = config
        self.process: Optional[asyncio.subprocess.Process] = None
        self._lock = asyncio.Lock()
        self._request_id = 0
        self._pending_requests: Dict[Any, asyncio.Future] = {}
        self._stdout_task: Optional[asyncio.Task] = None
        self._stderr_task: Optional[asyncio.Task] = None

    @property
    def is_running(self) -> bool:
        return self.process is not None and self.process.returncode is None

    @property
    def pid(self) -> Optional[int]:
        return self.process.pid if self.process else None

    async def start(self):
        """启动MCP子进程"""
        if self.is_running:
            logger.warning("Process already running", service_id=self.config.id)
            return

        cmd = self.config.transport.command
        args = self.config.transport.args
        full_cmd = [cmd] + args

        env = os.environ.copy()
        env.update(self.config.transport.env)

        try:
            self.process = await asyncio.create_subprocess_exec(
                *full_cmd,
                stdin=asyncio.subprocess.PIPE,
                stdout=asyncio.subprocess.PIPE,
                stderr=asyncio.subprocess.PIPE,
                env=env,
                limit=1024 * 1024,  # 1MB行缓冲区
            )

            self._stdout_task = asyncio.create_task(self._read_stdout())
            self._stderr_task = asyncio.create_task(self._read_stderr())

            self.config.status = "running"
            self.config.pid = self.process.pid
            self.config.started_at = datetime.now().timestamp()

        except Exception as e:
            self.config.status = "error"
            logger.error("Failed to start MCP service",
                        service_id=self.config.id, error=str(e))
            raise

    async def stop(self, graceful: bool = True, timeout: int = 5):
        """停止MCP子进程"""
        if not self.is_running:
            return

        if graceful:
            try:
                self.process.terminate()
                await asyncio.wait_for(self.process.wait(), timeout=timeout)
            except asyncio.TimeoutError:
                self.process.kill()
                await self.process.wait()
        else:
            self.process.kill()
            await self.process.wait()

        # 清理IO任务
        if self._stdout_task and not self._stdout_task.done():
            self._stdout_task.cancel()
        if self._stderr_task and not self._stderr_task.done():
            self._stderr_task.cancel()

        # 清理未完成的请求
        for future in self._pending_requests.values():
            if not future.done():
                future.set_exception(RuntimeError("Process terminated"))
        self._pending_requests.clear()

        self.config.status = "stopped"
        self.config.pid = None

    async def send_request(self, method: str, params: dict = None) -> JSONRPCResponse:
        """发送JSON-RPC请求并等待响应"""
        if not self.is_running:
            raise RuntimeError("Process not running")

        async with self._lock:
            self._request_id += 1
            req_id = self._request_id

            request = JSONRPCRequest(id=req_id, method=method, params=params or {})
            future = asyncio.Future()
            self._pending_requests[req_id] = future

            try:
                data = encode_request(request)
                self.process.stdin.write(data)
                await self.process.stdin.drain()

                response = await asyncio.wait_for(
                    future,
                    timeout=self.config.transport.timeout
                )
                return response
            finally:
                self._pending_requests.pop(req_id, None)

ProcessManager 提供全局进程管理接口,通过asyncio.Lock保证并发安全:

class ProcessManager:
    def __init__(self):
        self._processes: Dict[str, MCPProcess] = {}
        self._lock = asyncio.Lock()

    async def start_service(self, config: MCPServiceConfig) -> MCPProcess:
        """启动一个MCP服务"""
        async with self._lock:
            if config.id in self._processes:
                process = self._processes[config.id]
                if process.is_running:
                    return process
                del self._processes[config.id]

            process = MCPProcess(config)
            await process.start()
            self._processes[config.id] = process
            return process

    async def stop_service(self, service_id: str, graceful: bool = True):
        """停止一个MCP服务"""
        async with self._lock:
            process = self._processes.get(service_id)
            if process:
                await process.stop(graceful)
                del self._processes[service_id]

    def get_service(self, service_id: str) -> Optional[MCPProcess]:
        return self._processes.get(service_id)

    def list_services(self) -> Dict[str, MCPServiceConfig]:
        return {
            service_id: process.config
            for service_id, process in self._processes.items()
        }

    async def stop_all(self):
        """停止所有服务"""
        async with self._lock:
            tasks = [process.stop() for process in self._processes.values()]
            if tasks:
                await asyncio.gather(*tasks, return_exceptions=True)
            self._processes.clear()


# 全局进程管理器实例
process_manager = ProcessManager()

4.4 路由引擎

路由引擎维护工具名称到服务ID的映射关系,是请求分发的核心决策组件。ToolRouter 内部维护三张映射表:工具名到服务ID的正向映射、服务ID到工具列表的映射、服务ID到配置的映射。注册与注销操作通过asyncio.Lock保证线程安全,确保高并发下的路由一致性。

class ToolRouter:
    def __init__(self):
        self._tool_to_service: Dict[str, str] = {}  # tool_name -> service_id
        self._service_tools: Dict[str, List[MCPTool]] = {}  # service_id -> [tools]
        self._services: Dict[str, MCPServiceConfig] = {}  # service_id -> config
        self._lock = asyncio.Lock()

    async def register_service(self, config: MCPServiceConfig):
        """注册一个MCP服务"""
        async with self._lock:
            self._services[config.id] = config
            self._service_tools[config.id] = config.tools

            for tool in config.tools:
                self._tool_to_service[tool.name] = config.id

            logger.info("Service registered",
                       service_id=config.id,
                       tools_count=len(config.tools))

    async def unregister_service(self, service_id: str):
        """注销一个MCP服务"""
        async with self._lock:
            if service_id in self._services:
                del self._services[service_id]

            if service_id in self._service_tools:
                for tool in self._service_tools[service_id]:
                    self._tool_to_service.pop(tool.name, None)
                del self._service_tools[service_id]

    def route(self, tool_name: str) -> Optional[str]:
        """根据工具名路由到MCP服务"""
        return self._tool_to_service.get(tool_name)

    def get_all_tools(self) -> List[MCPTool]:
        """获取所有已注册的工具"""
        tools = []
        for service_tools in self._service_tools.values():
            tools.extend(service_tools)
        return tools

    def get_service_config(self, service_id: str) -> Optional[MCPServiceConfig]:
        return self._services.get(service_id)

    def list_services(self) -> List[MCPServiceConfig]:
        return list(self._services.values())


# 全局路由实例
router = ToolRouter()


def load_service_config(file_path: Path) -> MCPServiceConfig:
    """从YAML文件加载MCP服务配置"""
    with open(file_path, 'r', encoding='utf-8') as f:
        data = yaml.safe_load(f)
    return MCPServiceConfig.model_validate(data)

4.5 热插拔管理器

热插拔管理器基于Watchdog文件系统监听实现MCP服务的动态注册与注销。HotplugManager 在启动时先加载目录中所有现有配置文件,随后注册文件事件监听器,持续监控配置目录的变化。MCPConfigHandler 作为事件处理器,将文件创建、修改、删除事件转换为异步任务,通过 asyncio.run_coroutine_threadsafe 桥接到事件循环。

防抖机制通过0.1秒延迟等待和 _processing 集合去重,避免编辑器保存时产生的重复事件触发多次加载。更新操作采用先停旧服务、再更新路由、最后启新服务的顺序,确保切换过程的原子性。

class MCPConfigHandler(FileSystemEventHandler):
    def __init__(self, loop: asyncio.AbstractEventLoop):
        self.loop = loop
        self._processing: Set[str] = set()

    def on_created(self, event: FileSystemEvent):
        if not event.is_directory and event.src_path.endswith(('.yaml', '.yml')):
            asyncio.run_coroutine_threadsafe(
                self._handle_create(event.src_path), self.loop
            )

    def on_modified(self, event: FileSystemEvent):
        if not event.is_directory and event.src_path.endswith(('.yaml', '.yml')):
            asyncio.run_coroutine_threadsafe(
                self._handle_modify(event.src_path), self.loop
            )

    def on_deleted(self, event: FileSystemEvent):
        if not event.is_directory and event.src_path.endswith(('.yaml', '.yml')):
            asyncio.run_coroutine_threadsafe(
                self._handle_delete(event.src_path), self.loop
            )

    async def _handle_create(self, file_path: str):
        if file_path in self._processing:
            return
        self._processing.add(file_path)

        try:
            await asyncio.sleep(0.1)  # 防抖等待

            config = load_service_config(Path(file_path))
            await router.register_service(config)
            await process_manager.start_service(config)

        except Exception as e:
            logger.error("Failed to load new service",
                         file=file_path, error=str(e))
        finally:
            self._processing.discard(file_path)

    async def _handle_modify(self, file_path: str):
        if file_path in self._processing:
            return
        self._processing.add(file_path)

        try:
            await asyncio.sleep(0.1)

            config = load_service_config(Path(file_path))

            # 停止旧服务
            await process_manager.stop_service(config.id, graceful=True)
            # 更新路由
            await router.unregister_service(config.id)
            await router.register_service(config)
            # 启动新服务
            await process_manager.start_service(config)

        except Exception as e:
            logger.error("Failed to update service",
                         file=file_path, error=str(e))
        finally:
            self._processing.discard(file_path)

    async def _handle_delete(self, file_path: str):
        try:
            service_id = Path(file_path).stem
            await process_manager.stop_service(service_id)
            await router.unregister_service(service_id)
        except Exception as e:
            logger.error("Failed to unload service",
                         file=file_path, error=str(e))


class HotplugManager:
    def __init__(self, services_dir: str):
        self.services_dir = services_dir
        self.observer = Observer()
        self._started = False

    async def start(self):
        """启动热插拔管理器"""
        if self._started:
            return

        await self._load_existing_services()

        loop = asyncio.get_running_loop()
        event_handler = MCPConfigHandler(loop)
        self.observer.schedule(event_handler, self.services_dir, recursive=False)
        self.observer.start()

        self._started = True

    async def stop(self):
        """停止热插拔管理器"""
        if self._started:
            self.observer.stop()
            self.observer.join()
            self._started = False

    async def _load_existing_services(self):
        """加载目录中所有现有的服务配置"""
        dir_path = Path(self.services_dir)
        dir_path.mkdir(parents=True, exist_ok=True)

        config_files = list(dir_path.glob("*.yaml")) + list(dir_path.glob("*.yml"))

        for config_file in config_files:
            try:
                config = load_service_config(config_file)
                await router.register_service(config)
                await process_manager.start_service(config)
            except Exception as e:
                logger.error("Failed to load service config",
                             file=str(config_file), error=str(e))

4.6 WebSocket网关

WebSocket网关是AI客户端与平台交互的入口,负责连接管理、Token认证和消息分发。ConnectionManager 维护活跃连接集合,提供连接接受、断开清理和响应发送能力。消息处理流程通过 REQUEST_HANDLERS 映射表实现方法分发,支持initialize、tools/list、tools/call三个核心方法。

Token认证在连接建立前完成,校验失败时以 WS_1008_POLICY_VIOLATION 关闭连接。所有异常均被捕获并编码为JSON-RPC错误响应,确保连接不会因单次请求异常而中断。

class ConnectionManager:
    def __init__(self):
        self.active_connections: Set[WebSocket] = set()
        self._lock = asyncio.Lock()

    async def connect(self, websocket: WebSocket):
        await websocket.accept()
        async with self._lock:
            self.active_connections.add(websocket)

    async def disconnect(self, websocket: WebSocket):
        async with self._lock:
            self.active_connections.discard(websocket)

    async def send_response(self, websocket: WebSocket, response: JSONRPCResponse):
        try:
            data = encode_response(response)
            await websocket.send_text(data.decode('utf-8'))
        except Exception as e:
            logger.error("Failed to send response", error=str(e))


manager = ConnectionManager()


async def websocket_endpoint(websocket: WebSocket, token: str):
    """WebSocket接入端点"""
    from app.config import settings

    if token != settings.server.token:
        await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
        return

    await manager.connect(websocket)

    try:
        while True:
            data = await websocket.receive_text()

            try:
                request = decode_request(data.encode('utf-8'))
                await handle_request(websocket, request)
            except json.JSONDecodeError:
                error_resp = JSONRPCResponse(
                    error=JSONRPCError(code=-32700, message="Parse error")
                )
                await manager.send_response(websocket, error_resp)
            except Exception as e:
                error_resp = JSONRPCResponse(
                    error=JSONRPCError(code=-32603, message=str(e))
                )
                await manager.send_response(websocket, error_resp)

    except WebSocketDisconnect:
        await manager.disconnect(websocket)


async def handle_request(websocket: WebSocket, request: JSONRPCRequest):
    """处理WebSocket请求"""
    handler = REQUEST_HANDLERS.get(request.method)

    if handler:
        try:
            result = await handler(request.params or {})
            response = JSONRPCResponse(id=request.id, result=result)
            await manager.send_response(websocket, response)
        except Exception as e:
            error_resp = JSONRPCResponse(
                id=request.id,
                error=JSONRPCError(code=-32603, message=str(e))
            )
            await manager.send_response(websocket, error_resp)
    else:
        error_resp = JSONRPCResponse(
            id=request.id,
            error=JSONRPCError(code=-32601, message="Method not found")
        )
        await manager.send_response(websocket, error_resp)


async def handle_initialize(params: dict) -> dict:
    return {
        "protocolVersion": "2024-11-05",
        "capabilities": {"tools": {}},
        "serverInfo": {"name": "MCPilot", "version": "1.0.0"}
    }


async def handle_list_tools(params: dict) -> dict:
    tools = router.get_all_tools()
    return {
        "tools": [tool.model_dump() for tool in tools],
        "nextCursor": None
    }


async def handle_call_tool(params: dict) -> dict:
    tool_name = params.get("name")
    arguments = params.get("arguments", {})

    if not tool_name:
        raise ValueError("Tool name is required")

    service_id = router.route(tool_name)
    if not service_id:
        raise ValueError(f"Tool not found: {tool_name}")

    process = process_manager.get_service(service_id)
    if not process or not process.is_running:
        raise RuntimeError(f"Service not available: {service_id}")

    response = await process.send_request("tools/call", {
        "name": tool_name,
        "arguments": arguments
    })

    if response.error:
        raise RuntimeError(f"MCP error: {response.error.message}")

    return response.result


REQUEST_HANDLERS = {
    "initialize": handle_initialize,
    "tools/list": handle_list_tools,
    "tools/call": handle_call_tool,
    "notifications/initialized": lambda params: None
}

4.7 REST API

REST API提供服务管理和工具调用的HTTP接口,前端控制台通过这些接口完成可视化操作。服务管理API覆盖列表查询、详情查询、启动、停止、重启和工具列表六个端点。工具调用API通过POST请求接收工具名和参数,内部复用路由引擎和进程管理器完成调用。

from fastapi import APIRouter, HTTPException
from typing import List

api_router = APIRouter(prefix="/api/v1/services", tags=["services"])


@api_router.get("", response_model=List[ServiceInfo])
async def list_services():
    """列出所有MCP服务"""
    services = []
    for config in process_manager.list_services().values():
        services.append(ServiceInfo(
            id=config.id,
            name=config.name,
            version=config.version,
            status=config.status,
            pid=config.pid,
            tools_count=len(config.tools)
        ))

    # 补充已注册但未启动的服务
    for config in router.list_services():
        if not any(s.id == config.id for s in services):
            services.append(ServiceInfo(
                id=config.id,
                name=config.name,
                version=config.version,
                status="stopped",
                tools_count=len(config.tools)
            ))

    return services


@api_router.get("/{service_id}", response_model=ServiceInfo)
async def get_service(service_id: str):
    """获取单个服务详情"""
    config = router.get_service_config(service_id)
    if not config:
        raise HTTPException(status_code=404, detail="Service not found")

    process = process_manager.get_service(service_id)
    return ServiceInfo(
        id=config.id,
        name=config.name,
        version=config.version,
        status=process.config.status if process else "stopped",
        pid=process.pid if process else None,
        tools_count=len(config.tools),
        tools=[ToolInfo(name=t.name, description=t.description) for t in config.tools]
    )


@api_router.post("/{service_id}/start")
async def start_service(service_id: str):
    """启动服务"""
    config = router.get_service_config(service_id)
    if not config:
        raise HTTPException(status_code=404, detail="Service not found")
    try:
        process = await process_manager.start_service(config)
        return {"status": "started", "service_id": service_id, "pid": process.pid}
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))


@api_router.post("/{service_id}/stop")
async def stop_service(service_id: str):
    """停止服务"""
    try:
        await process_manager.stop_service(service_id)
        return {"status": "stopped", "service_id": service_id}
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))


@api_router.post("/{service_id}/restart")
async def restart_service(service_id: str):
    """重启服务"""
    config = router.get_service_config(service_id)
    if not config:
        raise HTTPException(status_code=404, detail="Service not found")
    try:
        await process_manager.stop_service(service_id)
        process = await process_manager.start_service(config)
        return {"status": "restarted", "service_id": service_id, "pid": process.pid}
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))


@api_router.post("/tools/call", response_model=CallToolResponse)
async def call_tool(request: CallToolRequest):
    """调用MCP工具"""
    service_id = router.route(request.tool_name)
    if not service_id:
        raise HTTPException(status_code=404, detail=f"Tool not found: {request.tool_name}")

    process = process_manager.get_service(service_id)
    if not process or not process.is_running:
        raise HTTPException(status_code=503, detail=f"Service not available: {service_id}")

    try:
        response = await process.send_request("tools/call", {
            "name": request.tool_name,
            "arguments": request.arguments or {}
        })

        if response.error:
            raise HTTPException(status_code=500, detail=response.error.message)

        return CallToolResponse(
            tool_name=request.tool_name,
            result=response.result,
            success=True
        )
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))

FastAPI主应用通过lifespan上下文管理器统一管理启动与关闭流程。启动时创建数据库表并启动热插拔管理器,关闭时停止文件监听并清理所有子进程:

@asynccontextmanager
async def lifespan(app: FastAPI):
    """应用生命周期管理"""
    await create_tables()
    await hotplug_manager.start()
    yield
    await hotplug_manager.stop()
    await process_manager.stop_all()


app = FastAPI(title="MCPilot API", version="1.0.0", lifespan=lifespan)

app.add_middleware(CORSMiddleware, allow_origins=["*"],
                   allow_credentials=True, allow_methods=["*"], allow_headers=["*"])

app.include_router(api_router)


@app.websocket("/mcp/ws")
async def websocket_route(websocket: WebSocket, token: str):
    await websocket_endpoint(websocket, token)


@app.get("/health")
async def health_check():
    return {
        "status": "ok",
        "services_count": len(process_manager.list_services()),
        "version": "1.0.0"
    }

4.8 数据模型

数据模型分为Pydantic API模型和SQLAlchemy ORM模型两组。API模型用于请求校验和响应序列化,ORM模型用于数据库持久化。两者职责分离,API模型关注接口契约,ORM模型关注存储结构。

Pydantic API模型定义如下:

from pydantic import BaseModel, Field
from typing import Optional, List, Any, Dict


class ToolInfo(BaseModel):
    name: str
    description: str


class ServiceInfo(BaseModel):
    id: str
    name: str
    version: str
    status: str  # registered / starting / running / error / stopped
    pid: Optional[int] = None
    tools_count: int = 0
    tools: Optional[List[ToolInfo]] = None


class ServiceCreate(BaseModel):
    id: str
    name: str
    version: str
    transport_config: Dict[str, Any]


class ServiceUpdate(BaseModel):
    name: Optional[str] = None
    version: Optional[str] = None


class CallToolRequest(BaseModel):
    tool_name: str
    arguments: Optional[Dict[str, Any]] = None


class CallToolResponse(BaseModel):
    tool_name: str
    result: Any
    success: bool
    error: Optional[str] = None

此外,适配器层定义了统一的工具调用请求与响应模型,用于适配器与网关之间的数据传递:

@dataclass
class ToolCallRequest:
    request_id: str           # 请求ID(用于追踪)
    tool_name: str            # 工具名称
    arguments: Dict[str, Any] # 调用参数
    context: Optional[Dict]   # 上下文信息(方法名等)
    stream: bool = False      # 是否启用流式输出


@dataclass
class ToolCallResponse:
    request_id: str           # 对应请求ID
    success: bool             # 调用是否成功
    result: Any = None        # 返回结果
    error: Optional[str] = None  # 错误信息
    metadata: Dict = None     # 元数据(耗时等)

SQLAlchemy ORM模型负责服务状态的持久化存储:

from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
from sqlalchemy.orm import declarative_base, sessionmaker
from sqlalchemy import Column, Integer, String, DateTime, func

DATABASE_URL = settings.database.url
engine = create_async_engine(DATABASE_URL, echo=settings.database.echo)

AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
Base = declarative_base()


class MCPServiceDB(Base):
    __tablename__ = "mcp_services"

    id = Column(Integer, primary_key=True, index=True)
    service_id = Column(String, unique=True, index=True)
    name = Column(String)
    version = Column(String)
    config_path = Column(String)
    status = Column(String, default="stopped")
    last_started_at = Column(DateTime)
    last_stopped_at = Column(DateTime)
    created_at = Column(DateTime, server_default=func.now())
    updated_at = Column(DateTime, server_default=func.now(), onupdate=func.now())


async def create_tables():
    async with engine.begin() as conn:
        await conn.run_sync(Base.metadata.create_all)


async def get_db() -> AsyncSession:
    async with AsyncSessionLocal() as session:
        yield session

5. 数据流与调用链

5.1 完整调用流程

以下为AI客户端通过WebSocket调用MCP工具的完整数据流,涵盖从连接建立到响应返回的九个步骤。每个步骤对应架构中的一个组件,数据在各组件间以JSON-RPC格式传递。

  AI 客户端
      │
      ▼ 1. WebSocket 连接
  ws://host/mcp/ws?token=xxx
      │
      ▼ 2. 适配器认证
  ┌─────────────────────┐
  │ Adapter             │
  │ parse_connection()  │
  └──────────┬──────────┘
             │
      ┌──────▼──────┐
      │   连接建立    │
      └──────┬──────┘
             │
             ▼ 3. 接收工具调用请求
  {
    jsonrpc: "2.0",
    id: "123",
    method: "tools/call",
    params: {
      name: "music.search",
      arguments: { keyword: "周杰伦" }
    }
  }
             │
             ▼ 4. 适配器解析请求
  ┌────────────────────────────┐
  │ Adapter.parse_request()    │
  │ -> ToolCallRequest 对象     │
  └──────────────┬─────────────┘
                 │
                 ▼ 5. 服务路由
  ┌────────────────────────────┐
  │ Router.route(tool_name)    │
  │ -> service_id = "music"    │
  └──────────────┬─────────────┘
                 │
                 ▼ 6. 获取进程
  ┌────────────────────────────┐
  │ ProcessManager.get_service()│
  │ -> MCPProcess 实例          │
  └──────────────┬─────────────┘
                 │
                 ▼ 7. 调用 MCP 服务
  ┌────────────────────────────┐
  │ process.send_request()      │
  │ 通过 stdin 发送到子进程      │
  │ 等待 stdout 返回结果         │
  └──────────────┬─────────────┘
                 │
                 ▼ 8. 适配器编码响应
  ┌────────────────────────────┐
  │ adapter.encode_response()   │
  │ 转换回 AI 期望格式           │
  └──────────────┬─────────────┘
                 │
                 ▼ 9. 返回给客户端
  {
    jsonrpc: "2.0",
    id: "123",
    result: { songs: [...] }
  }

5.2 消息循环处理

网关层的消息循环处理包含五个阶段,形成完整的请求-响应闭环。接收阶段从WebSocket读取原始数据,解析阶段通过适配器将原始格式转换为统一的ToolCallRequest,调用阶段通过路由引擎定位服务并转发请求到子进程,编码阶段将ToolCallResponse转换回AI期望的格式,发送阶段将响应写回WebSocket连接。整个流程在单个asyncio事件循环中运行,所有IO操作均为非阻塞。


6. 热插拔机制设计

6.1 配置变更检测流程

热插拔机制通过文件系统监听驱动配置变更的检测与生效。Watchdog Observer监控配置目录的文件创建、修改和删除事件,事件处理器通过防抖延迟后加载并校验新配置。配置校验通过后,根据操作类型执行相应的服务生命周期管理操作。

  配置文件目录 (configs/mcp-services/)
      │
      │ 文件创建/修改/删除事件
      ▼
  ┌─────────────────────────────┐
  │   Watchdog Observer         │
  │   - FileSystemEventHandler  │
  └──────────────┬──────────────┘
                 │ 防抖延迟处理 (0.1s)
                 ▼
  ┌─────────────────────────────┐
  │  1. 加载新配置文件           │
  │  2. 解析 YAML               │
  │  3. 校验配置格式            │
  └──────────────┬──────────────┘
                 │
         ┌───────┴───────┐
         │ 配置合法?     │
         └───────┬───────┘
           NO    │    YES
                 │ NO
  ┌──────────▼──────────┐
  │ 记录警告日志        │
  │ 保持原有服务运行    │
  └─────────────────────┘
                 │ YES
         ┌───────▼────────┐
         │ 判断操作类型    │
         └───────┬────────┘
                 │
      ┌──────────┼──────────┐
      ▼          ▼          ▼
  ┌────────┐ ┌─────────┐ ┌─────────┐
  │ 新增服务│ │ 更新服务│ │ 删除服务│
  └───┬────┘ └───┬─────┘ └───┬─────┘
      │          │            │
      └──────────┼────────────┘
                 ▼
  ┌─────────────────────────────────┐
  │ 1. 停止旧进程                    │
  │ 2. 从路由器注销工具              │
  │ 3. 启动新进程                    │
  │ 4. 注册新工具到路由器             │
  │ 5. 发送广播通知(服务变更事件)   │
  └─────────────────────────────────┘

6.2 关键特性

特性 说明
原子性操作 配置更新过程中,要么全部成功,要么保持原状
优雅降级 更新失败时,自动回滚到上一个可用版本
流量无损 更新过程中,正在处理的请求不会中断
快速生效 配置变更在1秒内生效
版本追踪 记录每次配置变更历史
防抖机制 0.1秒延迟等待,避免编辑器重复保存事件

6.3 服务状态机

MCP服务在生命周期中经历五个状态,状态转换由进程管理器和热插拔管理器驱动。各状态之间的转换遵循严格规则,非法转换将被拒绝并记录告警日志。

                 register
    ┌──────────────────────────┐
    │                          ▼
┌─────────┐  start()  ┌──────────┐  stop()  ┌─────────┐
│registered│─────────▶│ starting │────────▶│ stopped │
└─────────┘           └────┬─────┘         └─────────┘
                           │                            │
                    ready  │                            │ restart()
                           ▼                            │
                      ┌─────────┐                       │
                      │ running │◀──────────────────────┘
                      └────┬────┘
                           │
                    crash  │
                           ▼
                      ┌─────────┐
                      │  error  │──auto_restart──▶ running
                      └─────────┘

7. 适配器架构设计

7.1 基类定义

适配器架构的核心是 BaseAIAdapter 抽象基类,定义了所有AI适配器必须实现的统一接口规范。基类声明了五个核心方法:provider和protocol_type属性返回AI提供商与协议类型枚举,parse_connection解析并验证连接认证信息,parse_request将AI系统的原始请求转换为统一的ToolCallRequest,encode_response将ToolCallResponse编码为AI系统期望的格式,list_tools编码工具列表格式。

from abc import ABC, abstractmethod
from enum import Enum
from typing import Any, List, Optional
from dataclasses import dataclass


class AIProvider(Enum):
    """AI提供商枚举"""
    XIAOZHI = "xiaozhi"
    OPENAI = "openai"
    DOUBAO = "doubao"
    CUSTOM = "custom"


class ProtocolType(Enum):
    """协议类型枚举"""
    JSON_RPC = "json_rpc"
    FUNCTION_CALLING = "function_calling"
    CUSTOM = "custom"


class BaseAIAdapter(ABC):
    """AI适配器抽象基类"""

    @property
    @abstractmethod
    def provider(self) -> AIProvider:
        """返回AI提供商枚举"""
        pass

    @property
    @abstractmethod
    def protocol_type(self) -> ProtocolType:
        """返回协议类型枚举"""
        pass

    @abstractmethod
    async def parse_connection(self, websocket) -> bool:
        """
        解析并验证连接
        - 从 headers/query params 获取认证信息
        - 验证 token/api_key 是否合法
        - 返回 True 表示连接验证通过
        """
        pass

    @abstractmethod
    async def parse_request(self, raw_data: Any) -> ToolCallRequest:
        """
        解析AI系统发来的原始请求
        - 提取 request_id/tool_name/arguments
        - 转换为统一的 ToolCallRequest 格式
        - 处理特定AI的消息格式差异
        """
        pass

    @abstractmethod
    async def encode_response(self, response: ToolCallResponse) -> Any:
        """
        编码响应返回给AI系统
        - 将 ToolCallResponse 转换为AI期望的格式
        - 处理错误信息编码
        - 支持流式响应编码
        """
        pass

    @abstractmethod
    async def list_tools(self, tools: List[Any]) -> Any:
        """
        编码工具列表格式
        - 将统一的 ToolDefinition 列表转换为
          该AI系统期望的工具格式
        """
        pass

7.2 工厂模式

AdapterFactory 采用工厂模式管理适配器实例的创建与注册。适配器通过 @register_adapter 装饰器在模块加载时自动注册到工厂,网关层根据AI类型参数从工厂获取对应适配器实例。这种设计使新AI类型的接入对现有代码零侵入,只需编写适配器类并添加装饰器即可。

from typing import Dict, Type, Optional


class AdapterFactory:
    """适配器工厂"""

    def __init__(self):
        self._adapters: Dict[AIProvider, Type[BaseAIAdapter]] = {}
        self._instances: Dict[AIProvider, BaseAIAdapter] = {}

    def register(self, provider: AIProvider, adapter_class: Type[BaseAIAdapter]):
        """注册适配器类"""
        self._adapters[provider] = adapter_class

    def create_adapter(self, provider: AIProvider) -> Optional[BaseAIAdapter]:
        """创建或获取适配器实例(单例)"""
        if provider not in self._instances:
            adapter_class = self._adapters.get(provider)
            if adapter_class is None:
                return None
            self._instances[provider] = adapter_class()
        return self._instances[provider]

    def list_supported(self) -> list:
        """列出所有已注册的AI提供商"""
        return list(self._adapters.keys())


# 全局工厂实例
adapter_factory = AdapterFactory()


def register_adapter(provider: AIProvider):
    """适配器注册装饰器"""
    def decorator(cls):
        adapter_factory.register(provider, cls)
        return cls
    return decorator

适配器开发模板示例,新AI类型的适配器只需继承基类并实现抽象方法:

@register_adapter(AIProvider.CUSTOM)
class CustomAIAdapter(BaseAIAdapter):
    """自定义AI适配器开发模板"""

    @property
    def provider(self) -> AIProvider:
        return AIProvider.CUSTOM

    @property
    def protocol_type(self) -> ProtocolType:
        return ProtocolType.CUSTOM

    async def parse_connection(self, websocket) -> bool:
        # 从 headers/query params 获取认证信息并验证
        pass

    async def parse_request(self, raw_data: Any) -> ToolCallRequest:
        # 提取 request_id/tool_name/arguments 并转换格式
        pass

    async def encode_response(self, response: ToolCallResponse) -> Any:
        # 将统一响应转换为该AI系统期望的格式
        pass

    async def list_tools(self, tools: List[Any]) -> Any:
        # 将工具列表转换为该AI系统期望的格式
        pass

7.3 支持AI列表

AI提供商 协议类型 适配器类 状态 说明
小智AI JSON-RPC 2.0 XiaoZhiAdapter 已支持 WebSocket JSON-RPC协议接入
OpenAI Function Calling OpenAIAdapter 已支持 OpenAI Function Calling协议接入
豆包 Function Calling DoubaoAdapter 规划中 字节跳动豆包AI接入
自定义 自定义协议 CustomAIAdapter 可扩展 通过继承BaseAIAdapter实现

7.4 适配器开发最佳实践

适配器开发需遵循五项实践准则。错误处理方面,适配器内部应捕获所有异常,确保不影响网关主流程。日志记录方面,详细记录解析与编码过程,便于调试。版本兼容方面,需考虑AI系统协议版本差异,保持向后兼容。性能优化方面,避免解析与编码过程中的冗余计算。测试覆盖方面,为每个适配器编写完整的单元测试和集成测试。


8. 安全设计

8.1 多层安全防护架构

平台采用纵深防御策略,从传输层、认证层、权限层到服务层逐级加固安全防线。传输层强制TLS 1.2+加密,认证层验证Token与API Key,权限层校验JWT角色权限,服务层通过沙箱隔离限制子进程行为。

┌─────────────────────────────────────────────────────┐
│                    传输层安全                         │
│             TLS 1.2+ (HTTPS/WSS)                    │
└──────────────────────┬──────────────────────────────┘
                       │
        ┌──────────────┴───────────────┐
        │                              │
┌───────▼──────────┐        ┌────────▼──────────┐
│  连接认证层       │        │   权限控制层       │
│                   │        │                   │
│ - Token 验证      │        │ - JWT Token       │
│ - API Key 验证    │        │ - 角色权限校验     │
│ - 请求签名校验    │        │ - 操作审计日志     │
└──────────┬────────┘        └────────┬──────────┘
           │                           │
┌──────────▼───────────────────────────▼──────────┐
│             服务层安全                              │
│                                                    │
│ - MCP 进程沙箱隔离                                │
│ - 文件系统访问白名单                                │
│ - 子进程权限最小化                                  │
│ - 网络访问控制(仅允许特定 outbound)              │
└────────────────────────────────────────────────────┘

8.2 认证授权

WebSocket连接在握手阶段验证Token,Token从查询参数中获取并与配置中的预设值比对,校验失败时以策略违规码关闭连接。REST API通过API Key验证访问权限,后续阶段将引入JWT Token实现细粒度的角色权限控制。

# WebSocket认证流程
async def websocket_endpoint(websocket: WebSocket, token: str):
    from app.config import settings
    if token != settings.server.token:
        await websocket.close(code=status.WS_1008_POLICY_VIOLATION)
        return
    await manager.connect(websocket)
    # ...

8.3 数据安全

MCP子进程以低权限用户运行,通过操作系统层面的用户隔离降低安全风险。进程资源限制包括CPU和内存上限,防止单个服务异常消耗系统资源。文件系统访问通过白名单机制控制,子进程仅能访问指定工作目录。网络访问控制限制子进程的出站连接,仅允许特定目标地址。

配置安全方面,敏感配置如API密钥等加密存储,配置文件权限限制为600(仅所有者可读写)。环境变量注入敏感信息时,容器编排层面通过Secrets管理,避免明文暴露。

8.4 审计日志

平台记录完整的操作审计日志,覆盖服务启停、配置变更、工具调用等关键操作。审计日志采用结构化JSON格式,包含时间戳、操作者、操作类型、目标对象、操作结果和耗时等字段。日志通过structlog输出,支持JSON和文本两种格式,便于接入ELK等日志分析平台。

# 结构化日志格式
log_entry = {
    'timestamp': '2026-08-06T10:30:00Z',
    'level': 'INFO',
    'component': 'process_manager',
    'service_id': 'music-service',
    'tool_name': 'music.search',
    'request_id': 'req_12345',
    'message': 'Tool call completed',
    'duration_ms': 45,
    'error': None
}

9. 性能设计

9.1 优化策略

优化维度 策略 说明
并发模型 asyncio事件循环 所有IO操作异步化,单线程支撑高并发连接
进程通信 1MB行缓冲区 stdin/stdout管道设置大缓冲区,减少IO次数
请求匹配 Future映射表 请求ID与Future的O(1)匹配,避免轮询
路由查询 哈希表查找 工具名到服务ID的Dict查找,O(1)复杂度
连接管理 Set去重 活跃连接使用Set存储,O(1)增删
内存管理 生成器流处理 大数据流处理使用生成器,避免全量加载
超时保护 全链路超时 所有await操作设置超时,防止协程泄漏
防抖优化 延迟去重 文件事件防抖0.1秒,避免重复加载
子进程管理 进程组管理 子进程加入进程组,主进程退出时统一清理

9.2 基准指标

指标 目标值 说明
单节点支持MCP服务数 100+ 单进程管理100个以上MCP子进程
并发WebSocket连接数 1000+ 同时维持1000+活跃WebSocket连接
工具调用平均延迟 <100ms 从收到请求到返回响应的端到端延迟
配置更新生效时间 <1s 从配置文件变更到服务状态切换完成
服务启动时间 <3s MCP子进程从启动到就绪的时间
内存占用(核心服务) <500MB 平台主进程的常驻内存占用

9.3 监控指标

平台通过Prometheus暴露核心监控指标,覆盖连接、工具调用、服务状态和进程资源四个维度。连接指标包括按AI类型统计的活跃连接数和总连接次数。工具调用指标包括按服务与工具名统计的调用总数、成功数、错误数和延迟分布。服务状态指标包括运行中服务数和异常服务数。进程资源指标包括按服务ID统计的CPU使用率和内存占用。

# Prometheus 指标定义
metrics = {
    # 连接指标
    'active_connections': Gauge('活跃连接数', ['ai_type']),
    'connection_total': Counter('总连接次数', ['ai_type']),

    # 工具调用指标
    'tool_calls_total': Counter('工具调用总数', ['service_id', 'tool_name']),
    'tool_calls_success': Counter('调用成功数', ['service_id', 'tool_name']),
    'tool_calls_error': Counter('调用错误数', ['service_id', 'tool_name', 'error_type']),
    'tool_call_duration': Histogram('调用延迟分布', ['service_id', 'tool_name']),

    # 服务状态指标
    'services_running': Gauge('运行中服务数'),
    'services_error': Gauge('异常服务数'),

    # 进程资源指标
    'process_cpu_usage': Gauge('进程CPU使用率', ['service_id']),
    'process_memory_usage': Gauge('进程内存占用', ['service_id']),
}

9.4 Python特有风险与应对

风险 概率 影响 应对措施
GIL导致的并发瓶颈 中 中 使用asyncio充分利用IO并发;CPU密集型任务移到子进程;单worker部署
子进程孤儿泄漏 高 高 所有子进程加入进程组;主进程退出时SIGTERM整个组;定期巡检孤儿进程
asyncio死锁 中 高 严格遵循async/await规范;使用asyncio.Lock而非threading.Lock;添加调试超时
内存泄漏 中 中 tracemalloc定期快照对比;长连接定期清理;对象生命周期严格管理
Windows子进程问题 高 中 注意Windows平台asyncio子进程限制;提供WSL2部署方案;充分测试
Python版本兼容 低 低 锁定3.11+版本;使用类型注解和mypy检查

10. 测试策略

10.1 单元测试

单元测试覆盖协议编解码、进程管理、路由逻辑等核心模块,目标覆盖率分别为协议层90%以上、进程管理85%以上、路由引擎80%以上。测试框架使用pytest搭配pytest-asyncio支持异步测试,通过httpx提供HTTP客户端测试能力。

# 运行所有单元测试
pytest tests/unit/ -v

# 生成覆盖率报告
pytest tests/unit/ --cov=app --cov-report=html

重点测试模块与覆盖率要求:

模块 测试重点 覆盖率目标
app/gateway/protocol.py JSON-RPC编解码、配置校验、边界条件 >90%
app/core/process_manager.py 进程启停、请求匹配、超时处理、异常清理 >85%
app/core/router.py 注册注销、路由匹配、并发安全 >80%
app/core/hotplug.py 文件事件处理、防抖去重、状态切换 >75%
app/gateway/server.py 连接管理、消息分发、错误处理 >75%

10.2 集成测试

集成测试验证各模块协作的完整流程,重点覆盖WebSocket端到端调用链路。测试通过FastAPI TestClient建立WebSocket连接,依次发送initialize、tools/list、tools/call请求,验证响应格式与返回结果的正确性。

import pytest
from fastapi.testclient import TestClient
from app.main import app


@pytest.mark.asyncio
async def test_websocket_tool_call():
    """测试完整的WebSocket工具调用流程"""
    client = TestClient(app)

    with client.websocket_connect("/mcp/ws?token=your-secret-token-here") as ws:
        # 1. Initialize
        ws.send_json({
            "jsonrpc": "2.0",
            "id": 1,
            "method": "initialize",
            "params": {}
        })
        response = ws.receive_json()
        assert response["jsonrpc"] == "2.0"
        assert response["id"] == 1

        # 2. List tools
        ws.send_json({
            "jsonrpc": "2.0",
            "id": 2,
            "method": "tools/list"
        })
        response = ws.receive_json()
        assert "tools" in response["result"]

        # 3. Call tool
        ws.send_json({
            "jsonrpc": "2.0",
            "id": 3,
            "method": "tools/call",
            "params": {
                "name": "echo",
                "arguments": {"message": "Hello World"}
            }
        })
        response = ws.receive_json()
        assert response["result"]["echo"] == "Hello World"

10.3 压力测试

压力测试使用vegeta工具对REST API施加持续负载,评估系统在高并发场景下的吞吐量、延迟和错误率。测试场景包括服务列表查询的高频访问和工具调用的并发请求,通过调整请求速率和持续时间观察系统的性能拐点。

# 安装vegeta
go install github.com/tsenart/vegeta@latest

# 压力测试API
echo "GET http://localhost:8000/api/v1/services" | vegeta attack -rate=100/s -duration=30s | vegeta report

# 工具调用压力测试
echo '{"method":"POST","url":"http://localhost:8000/api/v1/services/tools/call","body":"{\"tool_name\":\"echo\",\"arguments\":{\"message\":\"test\"}}"}' | vegeta attack -rate=50/s -duration=60s | vegeta report

10.4 混沌测试

混沌测试通过主动注入故障验证系统的容错能力与恢复机制。测试场景覆盖进程崩溃、网络断连、配置异常和并发冲突四种典型故障模式,每种场景均定义明确的验证目标和通过标准。

测试场景 故障注入方式 验证目标
随机杀死MCP子进程 kill -9 子进程PID 自动检测异常并重启服务
断开WebSocket连接 强制关闭客户端连接 连接清理无泄漏,支持重连
格式错误的配置文件 写入非法YAML到配置目录 优雅降级,保持原有服务运行
高并发下更新路由 压测同时触发配置变更 无请求丢失,路由切换无损
子进程stdout无响应 注入不返回响应的MCP服务 超时保护生效,请求不阻塞
磁盘空间不足 填满日志目录 日志写入失败告警,服务不崩溃

附录

Echo示例服务配置

以下为Echo示例MCP服务的完整配置,该服务通过stdio方式与平台通信,提供echo工具将输入消息原样返回。此配置可用于开发环境的快速验证。

id: echo-service
name: Echo Service
version: 1.1.0
transport:
  type: stdio
  command: python
  args:
    - "-c"
    - |
      import sys
      import json

      def main():
          for line in sys.stdin:
              try:
                  req = json.loads(line)
                  if req["method"] == "initialize":
                      resp = {
                          "jsonrpc": "2.0",
                          "id": req["id"],
                          "result": {
                              "protocolVersion": "2024-11-05",
                              "capabilities": {"tools": {}},
                              "serverInfo": {"name": "echo", "version": "1.0.0"}
                          }
                      }
                  elif req["method"] == "tools/list":
                      resp = {
                          "jsonrpc": "2.0",
                          "id": req["id"],
                          "result": {
                              "tools": [{
                                  "name": "echo",
                                  "description": "Echo back the input message",
                                  "inputSchema": {
                                      "type": "object",
                                      "properties": {"message": {"type": "string"}},
                                      "required": ["message"]
                                  }
                              }]
                          }
                      }
                  elif req["method"] == "tools/call":
                      args = req["params"]["arguments"]
                      resp = {
                          "jsonrpc": "2.0",
                          "id": req["id"],
                          "result": {"echo": args.get("message", "")}
                      }
                  else:
                      resp = {
                          "jsonrpc": "2.0",
                          "id": req["id"],
                          "error": {"code": -32601, "message": "Method not found"}
                      }
                  print(json.dumps(resp), flush=True)
              except Exception as e:
                  error_resp = {
                      "jsonrpc": "2.0",
                      "id": req.get("id"),
                      "error": {"code": -32603, "message": str(e)}
                  }
                  print(json.dumps(error_resp), flush=True)

      if __name__ == "__main__":
          main()
  timeout: 30

架构设计原则

原则 说明
KISS原则 保持简单,避免过度设计
面向接口 依赖抽象而非具体实现
开闭原则 对扩展开放,对修改关闭
容错设计 任何组件失败都不会导致整个系统崩溃
可观测性 完善的监控、日志、追踪体系
默认安全 最小权限原则,安全配置为默认选项

扩展点

平台预留五个扩展点,支持后续功能演进。AI适配器扩展通过装饰器注册新AI类型,协议扩展支持HTTP、SSE、gRPC等更多传输协议,存储扩展支持从PostgreSQL切换到MySQL等其他关系型数据库,消息队列扩展支持Kafka或RabbitMQ实现异步调用链,认证扩展支持OAuth2、LDAP、SSO等更多认证方式。

部署架构

单机部署适用于开发和测试环境,FastAPI应用与所有MCP子进程运行在同一Docker容器内,PostgreSQL作为数据存储。集群部署适用于生产环境,Nginx负载均衡将流量分发到多个MCPilot节点,Redis集群提供分布式状态共享和锁服务,PostgreSQL主备架构保障数据高可用。

集群部署架构:

              ┌──────────────┐
              │   负载均衡    │
              │   (Nginx)    │
              └──────┬───────┘
                     │
        ┌────────────┼────────────┐
        │            │            │
   ┌────▼────┐  ┌────▼────┐  ┌────▼────┐
   │ MCPilot │  │ MCPilot │  │ MCPilot │
   │ Node 1  │  │ Node 2  │  │ Node 3  │
   └────┬────┘  └────┬────┘  └────┬────┘
        └────────────┼────────────┘
                     │
              ┌──────▼──────┐
              │ Redis 集群  │
              │ (状态/锁)   │
              └──────┬──────┘
                     │
              ┌──────▼──────┐
              │ PostgreSQL  │
              │ (主备+只读) │
              └─────────────┘

附录 B:实现状态(阶段1-4)

已实现模块

模块 文件 实现状态 说明
配置管理 backend/app/config.py 完成 Pydantic Settings + YAML + 环境变量嵌套覆盖
协议模型 backend/app/gateway/protocol.py 完成 JSON-RPC 2.0 模型 + 编解码函数
进程管理器 backend/app/core/process_manager.py 完成 asyncio subprocess + Future 匹配 + 崩溃检测 + 优雅停止
路由引擎 backend/app/core/router.py 完成 工具名路由 + 动态工具发现 + 并发安全
热插拔管理器 backend/app/core/hotplug.py 完成 Watchdog 文件监听 + 防抖 + 自动注册/注销
适配器层 backend/app/adapters/ 完成 小智 AI(JSON-RPC)+ OpenAI(Function Calling)
WebSocket 网关 backend/app/gateway/server.py 完成 多 AI 统一网关 + 消息循环 + 心跳保活 + 错误容错
REST API backend/app/api/ 完成 服务 CRUD + 工具调用 + 适配器查询 + 健康检查
数据模型 backend/app/models/ 完成 Pydantic API 模型 + SQLAlchemy ORM
数据库连接 backend/app/db/base.py 完成 PostgreSQL 异步引擎 + 建表 + 会话注入
数据迁移 backend/migrations/versions/0000_v2_1_0_baseline.py 完成 单一基线迁移(17 张表),启动时 create_all + 自动 stamp;backend/scripts/init_database.sql 为等价全量脚本
前端服务列表 frontend/src/views/Services.vue 完成 表格 + 统计卡片 + 自适应刷新 + 错误重试
前端服务详情 frontend/src/views/ServiceDetail.vue 完成 信息展示 + 工具列表 + 操作按钮 + 日志占位
前端工具测试 frontend/src/views/ToolTester.vue 完成 服务/工具选择 + JSON 输入 + 结果展示
Docker 部署 backend/Dockerfile + frontend/Dockerfile 完成 后端 Python slim + 前端 Nginx 多阶段构建
Docker Compose docker-compose.yml 完成 前后端一键启动 + 卷挂载 + 健康检查

阶段1代码审查修复

编号 问题 修复方式
C1 热插拔删除事件用文件名 stem 推断 service_id 遍历已注册服务匹配 id
C2 OpenAI 适配器响应 name 硬编码为 "tool" ToolCallResponse 增加 tool_name 字段
C3 Echo 服务异常处理二次抛错 try 前初始化 req = None
C4 子进程崩溃后 status 永久停留 running 新增 _on_process_exit 更新状态
C5 前端健康检查用错 axios 实例 改用带 baseURL 的 http 实例
C6/C7 ServiceDetail 轮询闪烁 + 路由切换不刷新 isPoll 标志位 + watch props.serviceId
W1-W14 锁粒度、异步阻塞、JSON 容错、认证、id 类型等 详见 AGENT.md

与设计文档的偏差

设计项 实际实现 原因
app/db/migrations.py 已实现(create_all) 设计文档列出但原实现遗漏,已补充
backend/tests/ 已实现(109+ 测试:单元+集成+分布式) 阶段2-3已完成全面测试覆盖
健康检查自动重启 已实现(后台探活30s间隔 + 自动重启 + 最大5次重试) 阶段2验收项已完成
Prometheus 指标 已实现(metrics.py + /api/v1/metrics 端点) 阶段3已完成:service_count/process_pid/tool_calls/uptime 等
实时日志推送 已实现(Redis pub/sub + /ws/logs/{id} + LogPanel.vue) 阶段3已完成:Redis 广播 + WebSocket 端点 + 前端日志面板
服务配置 CRUD API 已实现(POST/PUT/DELETE + YAML同步) 阶段2验收项已完成
Redis 分布式支持 已实现(分布式锁 + 状态共享 + pub/sub 日志) 阶段3引入 Redis 解决分布式与性能问题
Docker 容器化 已实现(Dockerfile + docker-compose + Redis) 阶段3已完成:资源限制 + 日志轮转 + .dockerignore
部署脚本 已实现(start.sh + stop.sh + .env.example) 阶段3已完成
用户认证与RBAC 已实现(API Key + JWT + 三角色权限矩阵) 阶段4已完成
操作审计 已实现(审计表 + best-effort 写入 + 分页查询 API) 阶段4已完成
服务配置管理页面 已实现(YAML 编辑器 + CRUD + 热加载) 阶段4已完成
MCP 服务市场 已实现(内置目录 + 一键安装) 阶段4已完成

文档版本: v1.5 最后更新: 2026年8月7日 维护者: MCPilot Team