| 项 | 内容 |
|---|---|
| 版本 | v2.1.0 |
| 日期 | 2026年8月18日 |
| 技术栈 | Python 3.11+ / FastAPI / Vue 3 / PostgreSQL |
| 文档状态 | 已定稿 |
当前AI应用生态中,MCP(Model Context Protocol)服务逐渐成为工具调用的标准协议,但多个MCP服务分散管理,缺乏统一入口。服务更新通常需要重启主服务,影响在线用户的正常使用。同时,运维人员缺乏可视化监控和调试工具,无法统一掌握服务状态与调用统计。
MCPilot 通用AI工具接入平台应运而生,旨在提供一个统一的MCP服务管理与接入网关。平台采用适配器架构,支持多种AI系统的协议接入,通过WebSocket网关实现实时通信,并具备热插拔能力,使配置变更秒级生效而无需重启主服务。
| 编号 | 目标 | 说明 |
|---|---|---|
| 1 | 热插拔能力 | MCP服务动态注册与注销,无需重启主服务 |
| 2 | 统一管理 | 一站式平台管理所有MCP服务的生命周期 |
| 3 | 可视化操作 | Web界面管理,降低使用门槛 |
| 4 | 协议兼容 | 兼容MCP标准协议,支持多种AI系统接入 |
| 5 | 易扩展 | 模块化适配器设计,便于接入新的AI类型 |
| 决策项 | 决策内容 | 理由 |
|---|---|---|
| 项目名称 | MCPilot(智控领航) | 体现控制与领航含义,呼应飞行员角色 |
| 开发语言 | Python | 生态丰富、开发速度快、团队熟悉度高 |
| Web框架 | FastAPI | 高性能异步、自动文档、类型安全、WebSocket原生支持 |
| 架构模式 | 适配器模式 | 支持多AI系统统一接入,易于扩展新AI类型 |
| 进程管理 | asyncio子进程 | Python原生异步支持、跨平台 |
| 配置方式 | YAML文件 + 环境变量 | 易于版本管理、支持容器化部署 |
| 热插拔实现 | 文件监听 + 动态注册 | 无需重启主服务,配置变更秒级生效 |
┌─────────────────────────────────────────────────────────┐
│ AI 客户端 │
│ (WebSocket JSON-RPC 接入) │
└────────────────────────────┬────────────────────────────┘
│
┌────────────────────────────▼────────────────────────────┐
│ AI 适配层 (Adapters) │
│ 协议转换 · 认证解析 · 请求/响应编解码 │
└────────────────────────────┬────────────────────────────┘
│
┌────────────────────────────▼────────────────────────────┐
│ WebSocket 网关层 (Gateway) │
│ 连接管理 · 消息分发 · 心跳保活 │
└────────────────────────────┬────────────────────────────┘
│
┌────────────────────────────▼────────────────────────────┐
│ 路由与分发层 (Router) │
│ 工具名匹配 · 服务实例路由 │
└────────────────────────────┬────────────────────────────┘
│
┌────────────────────────────▼────────────────────────────┐
│ MCP 进程管理层 (ProcessManager) │
│ 子进程生命周期 · stdio管道通信 · 健康检查 │
└────────────────────────────┬────────────────────────────┘
│
┌────────────────────┴────────────────────┐
│ │
┌───────▼──────────┐ ┌────────▼──────────┐
│ MCP 服务 A │ │ MCP 服务 B │
│ (音乐播放) │ │ (新闻资讯) │
└──────────────────┘ └───────────────────┘
│ │
┌───────▼──────────┐ ┌────────▼──────────┐
│ MCP 服务 C │ │ MCP 服务 D │
│ (智能家居) │ │ (知识库问答) │
└──────────────────┘ └───────────────────┘
系统采用四层架构,自上而下依次为AI适配层、网关层、核心业务层和API层。各层职责清晰分离,通过定义良好的接口交互,支持独立演进与测试。
处理不同AI系统的协议差异,将各类协议统一转换为平台内部的标准格式。该层包含适配器基类 BaseAIAdapter、适配器工厂 AdapterFactory 以及各AI系统的具体适配器实现。适配器通过 @register_adapter 装饰器自动注册到工厂,新AI类型的接入只需实现基类接口并添加装饰器即可。
负责WebSocket连接管理、消息分发和心跳保活。连接管理器维护 client_id 到 (websocket, adapter) 的映射,消息分发器将请求路由到对应适配器处理,并支持心跳保活机制。消息循环处理包含五个步骤:接收消息、适配器解析请求、调用工具、适配器编码响应、发送回客户端。
承载进程管理、服务路由和热插拔三大核心能力。ProcessManager 管理MCP子进程的完整生命周期,包括启动、停止、重启以及健康检测和自动重启。Router 维护工具名称到服务ID的映射表,支持动态注册和注销。Hotplug 通过文件系统监听实现配置变更检测,驱动服务的自动发现和注册。
提供RESTful接口管理和WebSocket端点。REST API覆盖服务管理(列表、启停、重启)和工具调用(列表、调用)两大功能域。WebSocket端点路径为 /mcp/ws,通过查询参数传递Token和AI类型信息。
| 层级 | 组件 | 职责 | 关键接口 |
|---|---|---|---|
| 适配层 | 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 |
| 组件 | 技术 | 版本 | 选型理由 |
|---|---|---|---|
| 语言 | 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。
| 组件 | 技术 | 版本 | 选型理由 |
|---|---|---|---|
| 框架 | Vue | 3.4+ | 渐进式框架、学习曲线平缓、生态成熟 |
| 构建工具 | Vite | 5+ | 超快冷启动、热更新体验优秀 |
| 语言 | TypeScript | 5+ | 类型安全、IDE体验好 |
| UI组件库 | Element Plus | 2.6+ | 组件丰富、文档完善、国内生态好 |
| HTTP客户端 | Axios | 1.6+ | 拦截器、取消请求等功能完整 |
| 状态管理 | Pinia | 2.1+ | Vue官方推荐,比Vuex更简洁 |
# 代码质量
pip install ruff # 超快速Lint + Format
pip install mypy # 类型检查
# 测试框架
pip install pytest # 测试框架
pip install pytest-asyncio # 异步测试支持
pip install httpx # HTTP客户端测试
# 包管理
pip install poetry # 依赖管理与打包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 # 项目说明
配置管理模块基于 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: jsonMCP协议模块封装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')进程管理器是平台的核心组件,负责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()路由引擎维护工具名称到服务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)热插拔管理器基于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))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
}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"
}数据模型分为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以下为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: [...] }
}
网关层的消息循环处理包含五个阶段,形成完整的请求-响应闭环。接收阶段从WebSocket读取原始数据,解析阶段通过适配器将原始格式转换为统一的ToolCallRequest,调用阶段通过路由引擎定位服务并转发请求到子进程,编码阶段将ToolCallResponse转换回AI期望的格式,发送阶段将响应写回WebSocket连接。整个流程在单个asyncio事件循环中运行,所有IO操作均为非阻塞。
热插拔机制通过文件系统监听驱动配置变更的检测与生效。Watchdog Observer监控配置目录的文件创建、修改和删除事件,事件处理器通过防抖延迟后加载并校验新配置。配置校验通过后,根据操作类型执行相应的服务生命周期管理操作。
配置文件目录 (configs/mcp-services/)
│
│ 文件创建/修改/删除事件
▼
┌─────────────────────────────┐
│ Watchdog Observer │
│ - FileSystemEventHandler │
└──────────────┬──────────────┘
│ 防抖延迟处理 (0.1s)
▼
┌─────────────────────────────┐
│ 1. 加载新配置文件 │
│ 2. 解析 YAML │
│ 3. 校验配置格式 │
└──────────────┬──────────────┘
│
┌───────┴───────┐
│ 配置合法? │
└───────┬───────┘
NO │ YES
│ NO
┌──────────▼──────────┐
│ 记录警告日志 │
│ 保持原有服务运行 │
└─────────────────────┘
│ YES
┌───────▼────────┐
│ 判断操作类型 │
└───────┬────────┘
│
┌──────────┼──────────┐
▼ ▼ ▼
┌────────┐ ┌─────────┐ ┌─────────┐
│ 新增服务│ │ 更新服务│ │ 删除服务│
└───┬────┘ └───┬─────┘ └───┬─────┘
│ │ │
└──────────┼────────────┘
▼
┌─────────────────────────────────┐
│ 1. 停止旧进程 │
│ 2. 从路由器注销工具 │
│ 3. 启动新进程 │
│ 4. 注册新工具到路由器 │
│ 5. 发送广播通知(服务变更事件) │
└─────────────────────────────────┘
| 特性 | 说明 |
|---|---|
| 原子性操作 | 配置更新过程中,要么全部成功,要么保持原状 |
| 优雅降级 | 更新失败时,自动回滚到上一个可用版本 |
| 流量无损 | 更新过程中,正在处理的请求不会中断 |
| 快速生效 | 配置变更在1秒内生效 |
| 版本追踪 | 记录每次配置变更历史 |
| 防抖机制 | 0.1秒延迟等待,避免编辑器重复保存事件 |
MCP服务在生命周期中经历五个状态,状态转换由进程管理器和热插拔管理器驱动。各状态之间的转换遵循严格规则,非法转换将被拒绝并记录告警日志。
register
┌──────────────────────────┐
│ ▼
┌─────────┐ start() ┌──────────┐ stop() ┌─────────┐
│registered│─────────▶│ starting │────────▶│ stopped │
└─────────┘ └────┬─────┘ └─────────┘
│ │
ready │ │ restart()
▼ │
┌─────────┐ │
│ running │◀──────────────────────┘
└────┬────┘
│
crash │
▼
┌─────────┐
│ error │──auto_restart──▶ running
└─────────┘
适配器架构的核心是 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系统期望的工具格式
"""
passAdapterFactory 采用工厂模式管理适配器实例的创建与注册。适配器通过 @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| AI提供商 | 协议类型 | 适配器类 | 状态 | 说明 |
|---|---|---|---|---|
| 小智AI | JSON-RPC 2.0 | XiaoZhiAdapter | 已支持 | WebSocket JSON-RPC协议接入 |
| OpenAI | Function Calling | OpenAIAdapter | 已支持 | OpenAI Function Calling协议接入 |
| 豆包 | Function Calling | DoubaoAdapter | 规划中 | 字节跳动豆包AI接入 |
| 自定义 | 自定义协议 | CustomAIAdapter | 可扩展 | 通过继承BaseAIAdapter实现 |
适配器开发需遵循五项实践准则。错误处理方面,适配器内部应捕获所有异常,确保不影响网关主流程。日志记录方面,详细记录解析与编码过程,便于调试。版本兼容方面,需考虑AI系统协议版本差异,保持向后兼容。性能优化方面,避免解析与编码过程中的冗余计算。测试覆盖方面,为每个适配器编写完整的单元测试和集成测试。
平台采用纵深防御策略,从传输层、认证层、权限层到服务层逐级加固安全防线。传输层强制TLS 1.2+加密,认证层验证Token与API Key,权限层校验JWT角色权限,服务层通过沙箱隔离限制子进程行为。
┌─────────────────────────────────────────────────────┐
│ 传输层安全 │
│ TLS 1.2+ (HTTPS/WSS) │
└──────────────────────┬──────────────────────────────┘
│
┌──────────────┴───────────────┐
│ │
┌───────▼──────────┐ ┌────────▼──────────┐
│ 连接认证层 │ │ 权限控制层 │
│ │ │ │
│ - Token 验证 │ │ - JWT Token │
│ - API Key 验证 │ │ - 角色权限校验 │
│ - 请求签名校验 │ │ - 操作审计日志 │
└──────────┬────────┘ └────────┬──────────┘
│ │
┌──────────▼───────────────────────────▼──────────┐
│ 服务层安全 │
│ │
│ - MCP 进程沙箱隔离 │
│ - 文件系统访问白名单 │
│ - 子进程权限最小化 │
│ - 网络访问控制(仅允许特定 outbound) │
└────────────────────────────────────────────────────┘
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)
# ...MCP子进程以低权限用户运行,通过操作系统层面的用户隔离降低安全风险。进程资源限制包括CPU和内存上限,防止单个服务异常消耗系统资源。文件系统访问通过白名单机制控制,子进程仅能访问指定工作目录。网络访问控制限制子进程的出站连接,仅允许特定目标地址。
配置安全方面,敏感配置如API密钥等加密存储,配置文件权限限制为600(仅所有者可读写)。环境变量注入敏感信息时,容器编排层面通过Secrets管理,避免明文暴露。
平台记录完整的操作审计日志,覆盖服务启停、配置变更、工具调用等关键操作。审计日志采用结构化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
}| 优化维度 | 策略 | 说明 |
|---|---|---|
| 并发模型 | asyncio事件循环 | 所有IO操作异步化,单线程支撑高并发连接 |
| 进程通信 | 1MB行缓冲区 | stdin/stdout管道设置大缓冲区,减少IO次数 |
| 请求匹配 | Future映射表 | 请求ID与Future的O(1)匹配,避免轮询 |
| 路由查询 | 哈希表查找 | 工具名到服务ID的Dict查找,O(1)复杂度 |
| 连接管理 | Set去重 | 活跃连接使用Set存储,O(1)增删 |
| 内存管理 | 生成器流处理 | 大数据流处理使用生成器,避免全量加载 |
| 超时保护 | 全链路超时 | 所有await操作设置超时,防止协程泄漏 |
| 防抖优化 | 延迟去重 | 文件事件防抖0.1秒,避免重复加载 |
| 子进程管理 | 进程组管理 | 子进程加入进程组,主进程退出时统一清理 |
| 指标 | 目标值 | 说明 |
|---|---|---|
| 单节点支持MCP服务数 | 100+ | 单进程管理100个以上MCP子进程 |
| 并发WebSocket连接数 | 1000+ | 同时维持1000+活跃WebSocket连接 |
| 工具调用平均延迟 | <100ms | 从收到请求到返回响应的端到端延迟 |
| 配置更新生效时间 | <1s | 从配置文件变更到服务状态切换完成 |
| 服务启动时间 | <3s | MCP子进程从启动到就绪的时间 |
| 内存占用(核心服务) | <500MB | 平台主进程的常驻内存占用 |
平台通过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']),
}| 风险 | 概率 | 影响 | 应对措施 |
|---|---|---|---|
| GIL导致的并发瓶颈 | 中 | 中 | 使用asyncio充分利用IO并发;CPU密集型任务移到子进程;单worker部署 |
| 子进程孤儿泄漏 | 高 | 高 | 所有子进程加入进程组;主进程退出时SIGTERM整个组;定期巡检孤儿进程 |
| asyncio死锁 | 中 | 高 | 严格遵循async/await规范;使用asyncio.Lock而非threading.Lock;添加调试超时 |
| 内存泄漏 | 中 | 中 | tracemalloc定期快照对比;长连接定期清理;对象生命周期严格管理 |
| Windows子进程问题 | 高 | 中 | 注意Windows平台asyncio子进程限制;提供WSL2部署方案;充分测试 |
| Python版本兼容 | 低 | 低 | 锁定3.11+版本;使用类型注解和mypy检查 |
单元测试覆盖协议编解码、进程管理、路由逻辑等核心模块,目标覆盖率分别为协议层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% |
集成测试验证各模块协作的完整流程,重点覆盖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"压力测试使用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混沌测试通过主动注入故障验证系统的容错能力与恢复机制。测试场景覆盖进程崩溃、网络断连、配置异常和并发冲突四种典型故障模式,每种场景均定义明确的验证目标和通过标准。
| 测试场景 | 故障注入方式 | 验证目标 |
|---|---|---|
| 随机杀死MCP子进程 | kill -9 子进程PID | 自动检测异常并重启服务 |
| 断开WebSocket连接 | 强制关闭客户端连接 | 连接清理无泄漏,支持重连 |
| 格式错误的配置文件 | 写入非法YAML到配置目录 | 优雅降级,保持原有服务运行 |
| 高并发下更新路由 | 压测同时触发配置变更 | 无请求丢失,路由切换无损 |
| 子进程stdout无响应 | 注入不返回响应的MCP服务 | 超时保护生效,请求不阻塞 |
| 磁盘空间不足 | 填满日志目录 | 日志写入失败告警,服务不崩溃 |
以下为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 │
│ (主备+只读) │
└─────────────┘
| 模块 | 文件 | 实现状态 | 说明 |
|---|---|---|---|
| 配置管理 | 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 |
完成 | 前后端一键启动 + 卷挂载 + 健康检查 |
| 编号 | 问题 | 修复方式 |
|---|---|---|
| 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