my-pi-agent--extension机制与mcp
pi 的 extension 机制
extension 就是一个 TypeScript 模块,默认导出一个工厂函数,拿到一个 ExtensionAPI 对象,往里面注册东西。
定义在 core/extensions/types.ts:1193 的
ExtensionAPI,分六大类:
| 类别 | API | 作用 |
|---|---|---|
| 事件订阅 | pi.on(event, handler) |
30+
种事件:会话(session_start/session_compact…)、Agent(turn_start/message_end/agent_end…)、工具(tool_call/tool_result,可就地改参数或
block)、模型、输入、context(改发给 LLM 的消息) |
| 工具注册 | pi.registerTool(...) |
注册 LLM 能调用的工具(subagent 用的就是它) |
| 命令/快捷键/flag | registerCommand / registerShortcut / registerFlag |
加 /mycmd、按键、CLI flag |
| 渲染 | registerMessageRenderer / registerMarkdownTransformer / registerEntryRenderer |
自定义消息、Markdown、条目的 TUI 渲染 |
| 动作 | sendMessage / sendUserMessage / appendEntry / setSessionName / exec / setModel / getActiveTools… |
主动驱动 agent、持久化状态、切模型 |
| Provider | registerProvider |
注册/覆盖模型供应商 |
机制原理
pi 的扩展系统采用了高度解耦的 “静态注册面(ExtensionAPI) + 动态运行上下文(ExtensionContext)” 设计:
1 | ┌────────────────────────────────────────────────────────┐ |
- ExtensionAPI(pi 实体):
- 传给扩展入口函数的参数;
- 负责声明“这个扩展有什么能力”(注册了哪些工具、订阅了哪些事件、扩展了哪些斜杠命令)。
- ExtensionContext(ctx 实体):
- 当事件触发、命令执行或工具被调用时,作为运行时参数动态传入 Handler;
- 负责提供“当前环境的即时上下文与交互能力”(当前会话树、TUI 交互接口、Abort 取消信号等)。
生命周期拦截
1 | 用户输入 (User Prompt) |
input:截获用户原始输入。返回{ action: "handled" }可直接由扩展自行处理而不触发 LLM;返回{ action: "transform", text: "..." }可对 Prompt 进行预处理改写。before_agent_start:在 Agent 循环启动前触发,支持动态追加上下文消息,或修改本轮调用的systemPrompt。context:在每次调用 LLM 前触发,提供对发往模型的messages进行过滤与裁剪(非破坏性视图)。tool_call(核心拦截点):在工具实际执行前触发。支持原地修改event.input(参数修补),或返回{ block: true, reason: "..." }阻断高危命令执行。tool_result:在工具执行后触发,支持链式修改返回给模型的content或供前端消费的details。
自定义工具体系:pi.registerTool
扩展是如何把一个普通函数变成大模型能调用的工具的?
- 自动契约合成: @api.tool 内部复用框架底层的 Pydantic 动态建模,自动从 Python 函数签名、类型标注和 docstring 提取生成标准的 OpenAI / Anthropic Function Calling Schema。
- 后加载覆盖机制(Overriding Power): 在 Agent.__init__
的装配时序中: 注册内置工具与用户工具 → 加载 Extension 扩展工具
因为扩展是在最后阶段被调用的,所以如果扩展中注册了一个同名工具(如
read、bash),它会静默覆盖掉默认的内置实现。
- 架构价值:开发者无需修改核心代码,就能通过扩展将原生的本地文件读取工具替换为“带权限审计的只读沙箱工具”。
本地命令调度与反射路由(Bypass-LLM Dispatching)
为了向用户提供 0 Token 消耗的本地交互通道,ExtensionAPI 提供了命令路由系统:
1
2
3
4
5
6
7
8
9
10
11用户输入: "/echo hello"
│
▼ CLI 前置拦截 (不调大模型)
ExtensionManager.handle_command("echo", "hello")
│
├─ 1. 查表获取 handler
├─ 2. 反射探测形参: inspect.signature(handler).parameters
│ ├─ 若 len == 0: 执行 handler()
│ └─ 若 len > 0: 执行 handler("hello")
▼
直接向控制台返回结果 (0 Token 消耗,不污染会话历史)
通过 Python 的 inspect 反射能力,实现了对无参命令(如 /stats、/compact)和带参命令(如 /echo foo)的参数自动适配。
案例
在实际开发中,一个扩展往往会同时组合使用这三项能力,形成功能闭环:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30# .agents/extensions/todo_extension.py
"""一个完整的待办事项(Todo)扩展:
- 注册 todo_add 工具供大模型调用
- 注册 /todos 命令供用户在终端随时查看
- 监听 AgentEnd 事件在任务结束时统计待办
"""
from my_agent_core.events import AgentEnd
# 模块级共享状态
TODOS: list[str] = []
def extension(api):
# 1. 【工具注册】:让大模型可以在思考中记录待办
def todo_add(item: str) -> str:
TODOS.append(item)
return f"已成功添加待办: {item}"
# 2. 【命令调度】:用户在终端输入 /todos 时直接打印,不花大模型 Token
def cmd_todos():
if not TODOS:
return "当前没有待办事项。"
return "\n".join(f"{i+1}. {t}" for i, t in enumerate(TODOS))
# 3. 【事件订阅】:每次对话轮次彻底结束时,若有待办则在终端输出提示
def notify_on_end(event: AgentEnd, api):
if TODOS:
print(f"\n[Todo 扩展提示] 当前还有 {len(TODOS)} 项待办未完成,输入 /todos 查看。")
异步消息导流与插队机制(Message Steering)
当 Agent 正在忙碌地运行一个多轮 ReAct 循环时,扩展如何向模型传达新指令?
ExtensionAPI 定义了两种核心的异步分发模式:
1
2
3
4
51. 实时插队纠偏 (deliverAs: "steer"):
LLM 生成 -> 执行工具 1 -> [扩展插队注入 user 消息] -> 下一轮 LLM 生成 (立即生效纠偏!)
2. 静默排队跟进 (deliverAs: "followUp"):
LLM 生成 -> 执行工具 1 -> 执行工具 2 -> 任务彻底完成 -> [触发下一轮新任务]
- steer(方向盘):当前轮次的工具一跑完,下一次调模型前强行把消息插入队列,实现对模型的“实时急刹车与路线矫正”;
- followUp(待办队列):等整个 Agent 任务全部跑完空闲后,才开启下一轮新对话。
如何利用 extension 实现subagent
subagent 实现 = 注册一个名叫 subagent 的”执行器/转接器”工具,具体 agent 是它的运行时配置数据,不是各自的工具。
工具表里永远只有一个 subagent,它的参数是
{ agent: "某个名字", task: "..." }(index.ts:448 的
schema)。模型填的是字符串,然后 execute 去 agents 目录里查对应的 md
文件。这也解释了为什么要 --append-system-prompt
写临时文件——agent 的 system prompt
是运行时读到的字符串,不是编译进工具的代码。
pi 的 subagent 不是框架里的一等公民,而是用一个通用扩展点
pi.registerTool 实现出来的。 全部代码就一个文件
index.ts,主 agent 的模型把”委派任务”当成一次普通工具调用来完成。
入口:只注册了一个工具
1 | // index.ts:460 |
如何利用 extension 实现mcp
MCP(Model Context Protocol)连接器同样不需要作为框架核心的硬编码逻辑,而是通过 Extension 机制以“外部工具动态桥接”的方式实现:
1 | MCP 扩展启动 (extension(api)) |
- 协议转换:Extension 在
initialize握手后,通过tools/list拿到 MCP Server 声明的所有工具与 Schema; - 工具注册:遍历 MCP
工具,将每个外部工具动态包装为本地的
Tool,调用api.register_tool注入 Agent; - 调用转发:当模型调用该工具时,包装函数的
execute拦截调用,将参数打包为tools/call请求通过 Stdio/SSE 转发给 MCP Server 进程,并把执行结果原样喂回给大模型。
extensions.py
在我们的 Python 框架中,提炼了 pi 的核心设计,通过
extensions.py 实现了极简的框架级扩展三件套:
1 | ┌──────────────────────────────────────────────────────────────┐ |
ExtensionAPI
在 my-pi-agent 框架中,ExtensionAPI 是整个扩展机制的 “开发者契约面(Developer Surface)”。
当外部写一个扩展(如 mcp.py 或 .agents/extensions/my_tool.py)时,入口函数接收到的唯一参数就是 api: ExtensionAPI:
1
2
3def extension(api: ExtensionAPI):
# 扩展开发者只跟 api 对象打交道
...
ExtensionAPI 的定位是 “将 Agent 内部复杂的事件系统(HookRegistry)、工具注册表(ToolRegistry)和命令系统,收敛为最简单 直观的三套对外接口”。
核心能力 1:事件订阅系统 —— api.on
1 | class ExtensionAPI: |
统一干预数据结构:HookResult 与五大决策点
所有生命周期拦截点共享统一强类型的 HookResult
数据结构:
1 |
|
框架在 Agent.run() 中完整对齐了 Pi 的
五大生命周期决策拦截点: 1.
UserInput(input):在用户输入进入
Session 之前触发,支持 block=True 阻断输入或
updated_input 改写输入; 2.
AgentStart(before_agent_start):在准备好系统消息后触发,支持
updated_system_prompt 动态更新首条 system 消息; 3.
BeforeModelCall(context):在
_ctx.prepare() 产出 view
后、调用大模型前触发,支持 updated_messages
临时改写送给大模型的视图; -
关键架构不变式:updated_messages
仅临时修改当前这次发给 LLM 的 view
变量,真实 self.messages 与 Session 磁盘 JSONL
保持绝对纯净,实现 零历史污染; 4.
ToolExecutionStart(tool_call):工具执行前触发,支持
block=True 拦截危险调用或 updated_args
改写入参; 5.
ToolExecutionEnd(tool_result):工具执行后触发,支持
updated_result 改写返回给模型的出参。
流式 Token 级实时熔断(MessageUpdate):除上述 5 大决策点外,
MessageUpdate同样继承了Interceptable。在大模型逐 Token 流式生成的过程中,扩展如果检测到危险输出片段,返回HookResult(block=True)即可瞬间掐断生成,并且框架会直接丢弃未完成的半截文本(不写入 Session 树),彻底避免模型在下一轮对话中产生“续写断句”的严重幻觉。
核心能力 2:工具注册系统 —— api.tool 与 api.register_tool
1
2
3
4
5
6
7
8
9
10
11
12
13
14def register_tool(self, tool: Tool) -> None:
"""注册工具 → agent.registry(撞名静默覆盖,registry 语义)。"""
self.agent.registry.register(tool)
def tool(self, **kwargs: Any):
"""@api.tool(description=...) 装饰器:@tool 包装 + register_tool。"""
def decorator(func) -> Tool:
# 复用底层 pydantic 动态建模装饰器
t = _tool(**kwargs)(func)
self.register_tool(t)
return t
return decorator
- @api.tool(…):面向普通的 Python 业务函数。扩展开发者写一个原生 Python 函数,加上 @api.tool,内部自动通过 Pydantic 提取函数签名、类型标注和 docstring,生成 OpenAI/Anthropic 兼容的 JSON Schema,并注册进 Agent。
- api.register_tool(tool: Tool):面向复杂的动态工具(如 MCP 远程工具、Subagent 委派工具)。直接传入已经构造好的 Tool 实体对象(例如带有 raw_schema 的 MCP 工具)。
核心能力 3:命令注册系统 —— api.command 与 api.register_command
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18def register_command(
self, name: str, handler: CommandHandler, description: str = ""
) -> None:
"""注册命令(name 不含 /)。"""
self._commands[name] = handler
def command(self, name: str, description: str = ""):
"""@api.command("now") 装饰器。"""
def decorator(func: CommandHandler) -> CommandHandler:
self.register_command(name, func, description)
return func
return decorator
def get_commands(self) -> dict[str, CommandHandler]:
"""已注册命令的拷贝(name → handler)。"""
return self._commands.copy()
- 将用户输入的命令名(如 “mcp” 或 “now”)与对应的 Python 处理函数映射在 self._commands 字典中;
- get_commands() 返回字典的浅拷贝(.copy()),保护内部字典不被外部意外修改。
ExtensionManager
ExtensionManager 就是面向 Agent 框架内部的“扩展总管与调度仓储(Repository)”。
它对标了框架内的 SkillManager 与 SubagentManager,专门负责: 1.
解析扫描目录(三态语义:默认探测 / 显式禁用 / 外部指定); 2.
递归发现磁盘上的 .py 扩展文件(自动跳过 _
开头的私有辅助文件); 3. 通过 importlib 动态加载模块并执行
extension(api) 握手(原生支持 async def
协程与同步 def 入口,支持异步长连接初始化); 4.
单点故障隔离(坏插件隔离保护:try...except
捕获异常,单个扩展崩溃仅告警、不影响主 Agent 启动); 5.
用户斜杠命令的反射分发与调度(自动适配 0 参 / 1 参 Handler)。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
211. 启动装配期 (Agent.__init__)
┌──────────────────────────────────────────────────────────┐
│ Agent.__init__() │
│ ├── self._register_tools(tools) │
│ ├── self.extension_manager = ExtensionManager(self, ..)│
│ └── self.extension_manager.load() │
│ │ │
│ ▼ (遍历发现 .py 文件) │
│ ExtensionManager.load_extension("mcp.py") │
│ │ │
│ ▼ (执行插件入口) │
│ mcp.extension(self.api) │
│ ├── api.register_tool(t) ──► 注入 Agent.registry │
│ ├── api.on(Event) ──► 注入 Agent.hooks │
│ └── api.command("mcp") ──► 注入 api._commands │
└──────────────────────────────────────────────────────────┘
2. 运行交互期
├── [模型调用工具] ──► 查表分发至 Extension 注入的工具 (ReAct 循环)
├── [生命周期事件] ──► 触发 Extension 注册的 Hook (洋葱拦截)
└── [用户输入 /cmd] ──► extension_manager.handle_command("cmd", args)
生命周期与方法全景图
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40 【启动装配阶段】
1. __init__(agent, extension_dirs)
│ ├─ 绑定 agent
│ ├─ 创建 ExtensionAPI(agent)
│ └─ 解析目录路径 (三态: 默认探测 / 显式禁用 / 指定目录)
▼
2. load() ─── 批量加载与单点故障隔离保护
│
├─► 3. discover(directory) ─── 递归扫描文件
│ │ 扫描目录及子目录下的 **/*.py
│ └─ 自动跳过 _ 开头的私有文件 (_helper.py 等)
│
│ ┌─ 针对发现的每个 .py 文件逐个调用 ──────────┐
│ │ (用 try...except 隔离: 坏扩展只告警不崩溃) │
▼ ▼ │
4. load_extension(path) ─── 动态加载单个扩展 │
│ ├─ importlib 动态载入模块 │
│ ├─ 智能识别入口函数: extension(api) │
│ └─ 执行 extension_func(self.api) │
│ ├── 注册工具 ──► 注入 Agent.registry │
│ ├── 订阅事件 ──► 注入 Agent.hooks │
│ └── 注册命令 ──► 存入 self.api._commands │
└────────────────────────────────────────────────────┘
────────────────────────────────────────────────────────────────────────────
【运行交互阶段 (0 Token 消耗)】
用户在终端输入: "/mcp status"
│
▼
5. handle_command("mcp", "status") ─── 斜杠命令反射调度
│ ├─ 查表: 从 self.api._commands 找到对应 handler
│ ├─ 参数自适应:
│ │ ├─ 若为无参函数 cmd() ──► 执行 handler()
│ │ └─ 若为带参函数 cmd(args) ──► 执行 handler("status")
│ └─ 未知命令抛出清晰错误提示
▼
向终端控制台直接输出结果
mcp
mcp.py 是框架的 MCP (Model Context Protocol) 客户端内置扩展。
分层说明:在最新的架构重构中,MCP 客户端以 Extension 插件的形式置于产品层(
packages/my-coding-agent/src/my_coding_agent/mcp.py),使框架层my-agent-core保持极简纯净,同时产品层可随时按需插拔加载。
它不仅是一个完整的 MCP 客户端实现,同时也是一个遵循框架标准规范的 Extension 插件(包含工具注册与命令注册)。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20 .mcp.json (工作区配置)
│
▼
MCPClientManager (多 Server 管理)
│
┌───────────┴───────────┐
▼ ▼
MCPConnection (Server A) MCPConnection (Server B)
[AsyncExitStack 管理] [AsyncExitStack 管理]
├── stdio_client 子进程 ├── stdio_client 子进程
└── ClientSession (MCP) └── ClientSession (MCP)
│ │
└───────────┬───────────┘
▼
包装为内部 Tool 对象
(raw_schema=远程Schema, is_parallel_safe=True)
│
▼
ExtensionAPI.register_tool() ──► 注入 Agent.registry
ExtensionAPI.command("mcp") ──► 注册 /mcp 状态命令
配置数据类:MCPServerConfig
1
2
3
4
5
6
class MCPServerConfig:
name: str
command: str
args: list[str]
env: dict[str, str] | None = None
对应 .mcp.json 中配置的单个 Server 条目,如:
1
2
3
4
5
6
7
8{
"mcpServers": {
"filesystem": {
"command": "npx",
"args": ["-y", "@modelcontextprotocol/server-filesystem", "D:/data"]
}
}
}
单连接管理器:MCPConnection
这是与单个 MCP 子进程打交道的连接实体。
1
2
3
4
5class MCPConnection:
def __init__(self, config: MCPServerConfig):
self.config = config
self._session: ClientSession | None = None
self._exit_stack = contextlib.AsyncExitStack()
核心亮点:AsyncExitStack 上下文管理 MCP 官方 SDK 的 stdio_client 和 ClientSession 都是异步上下文管理器(async with)。在不退出当前方法的情况下长期保持连接,使用 AsyncExitStack 可以将多个异步上下文压栈保存,在需要关闭时 调用 aclose() 一键安全释放。
建立连接与握手 (start):
1
2
3
4
5
6
7
8
9
10
11
12
13
14async def start(self) -> None:
params = StdioServerParameters(
command=self.config.command,
args=self.config.args,
env=server_env,
)
read_stream, write_stream = await self._exit_stack.enter_async_context(
stdio_client(params)
)
session = await self._exit_stack.enter_async_context(
ClientSession(read_stream, write_stream)
)
self._session = session
await session.initialize() # 完成 MCP 初始化握手协议调用工具 (call_tool):
1
async def call_tool(self, name: str, arguments: dict[str, Any]) -> ToolResult:
- Never-Throw 保障:捕获所有异常,绝不向上抛出导致 Agent 崩溃,而是包装为 ToolResult(ok=False, error=…) 让大模型自我修正。
- 内容格式化:将 MCP 返回的多段 content 抽取为纯文本,并识别 MCP 协议中的 isError 标记。
1 | ┌──────────────────────────────────────────────────────────┐ |
先有第一扇门,才有第二扇门: 如果没有 stdio_client 启动进程并提供 read/write_stream,ClientSession 就根本不知道该向哪里读写协议数据。
关闭时必须倒序退出:
- 先关 ClientSession(告诉对方我们要结束会话了,把未完成的 RPC 请求取消);
- 再关 stdio_client(彻底杀掉子进程,释放操作系统进程句柄与内存)。
1
2
3
4
5
6
7
8
9
10
11
12
13
14┌───────────────────────────────────────────────────────────────┐
│ 第一步:stdio_client(params) 【物理与传输层:打通管道】 │
│ 操作系统的进程与字节流 (Process & Pipes) │
│ 输入: command="npx", args=[...] │
│ 产出: 裸管道 (read_stream, write_stream) │
└───────────────────────────────┬───────────────────────────────┘
│ 传递裸管道
▼
┌───────────────────────────────────────────────────────────────┐
│ 第二步:ClientSession(read, write) 【协议与业务层:听懂语言】 │
│ JSON-RPC 2.0 与 MCP 业务对象 │
│ 输入: 裸管道 (read_stream, write_stream) │
│ 产出: 高级会话对象 session (能调 list_tools/call_tool) │
└───────────────────────────────────────────────────────────────┘
mcp的底层通信协议:JSON-RPC 2.0
MCP 的所有消息交互全部建立在 JSON-RPC 2.0 规范之上。客户端和服务端通过互相发送结构化的 JSON 文本进行交流。
常见的几种通信帧:
#### 1. 工具列表发现(tools/list)
- Client 发出请求:
1
2
3
4
5
6{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/list",
"params": {}
} - Server 返回结果(带 JSON Schema 参数契约):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19{
"jsonrpc": "2.0",
"id": 1,
"result": {
"tools": [
{
"name": "calculate_tax",
"description": "计算个人所得税",
"inputSchema": {
"type": "object",
"properties": {
"income": { "type": "number", "description": "税前收入" }
},
"required": ["income"]
}
}
]
}
}
#### 2. 工具调用执行(tools/call)
- Client 发出执行指令:
1
2
3
4
5
6
7
8
9{
"jsonrpc": "2.0",
"id": 2,
"method": "tools/call",
"params": {
"name": "calculate_tax",
"arguments": { "income": 20000 }
}
} - Server 返回执行结果:
1
2
3
4
5
6
7
8
9
10{
"jsonrpc": "2.0",
"id": 2,
"result": {
"content": [
{ "type": "text", "text": "应纳个税: 1590 元" }
],
"isError": false
}
}
传输层模式:Stdio vs SSE/HTTP
MCP 规范定义了两种主要的物理传输通道:
1
2
3
4
5
6
7
8
9
101. Stdio(标准输入输出,本地子进程模式 —— 我们框架采用的模式):
Client (Python Agent)
│ stdin (向子进程写 JSON)
├───► [MCP Server 子进程 (如 Node.js / Python)]
◄───┘ stdout (从子进程读 JSON)
优点:零网络端口暴露,启动即用,极高安全性,适合本地工具。
2. SSE / HTTP(Server-Sent Events 远程流式模式):
Client ─── HTTP POST / SSE ───► 远程云端 MCP Server (企业内网服务)
优点:适合连接跨机器的大型企业数据库或云端 API。
多服务协调器:MCPClientManager
如果把单个 MCPConnection 比作 “专线接线员”,那么 MCPClientManager 就是 “接线总调度中心”。
在实际项目中,用户通常会在 .mcp.json 中配置多个 MCP 服务(比如一个 GitHub 查 PR、一个 Postgres 查数据库、一个 Filesystem 读文件)。 MCPClientManager 负责管理所有这些服务的生命周期编排与工具适配。
一、连接池状态:init
1
2
3class MCPClientManager:
def __init__(self):
self.connections: dict[str, MCPConnection] = {}
- 维护一个全局字典 self.connections,以服务名称为 key(如 “github”、“filesystem”),存储对应的 MCPConnection 实例;
- 方便后续按服务名查找连接、状态检查(/mcp 命令)和统一关闭。
### 二、配置解析:load_config(path)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21def load_config(self, path: Path | str) -> list[MCPServerConfig]:
p = Path(path)
if not p.exists():
return []
try:
data = json.loads(p.read_text(encoding="utf-8"))
except Exception as exc:
raise ValueError(f"Invalid JSON in {p}: {exc}") from exc
servers = data.get("mcpServers", {})
configs = []
for name, srv in servers.items():
configs.append(
MCPServerConfig(
name=name,
command=srv.get("command", ""),
args=srv.get("args", []),
env=srv.get("env"),
)
)
return configs
- 契约对齐:完全对齐 Claude Desktop / Cursor / Pi 的标准 .mcp.json
格式:
1
2
3
4
5
6
7
8
9{
"mcpServers": {
"sqlite": {
"command": "uvx",
"args": ["mcp-server-sqlite", "--db-path", "test.db"],
"env": {"DEBUG": "1"}
}
}
} - 将 JSON 转换为强类型的 MCPServerConfig 数据类列表。
### 三、协议适配与工具转换:connect_server(config) —— 最核心方法
这是整个类最精妙的部分,它完成了从 “MCP 远程服务” 到 “Agent 内部标准 Tool 对象” 的桥接转换:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31async def connect_server(self, config: MCPServerConfig) -> list[Tool]:
# 1. 建立长连接并完成握手
conn = MCPConnection(config)
await conn.start()
self.connections[config.name] = conn
# 2. 动态获取远程支持的所有工具
mcp_tools = await conn.list_tools()
wrapped_tools = []
# 3. 逐个包装成 Agent 框架的标准 Tool 对象
for t in mcp_tools:
tool_name = t.name
schema = getattr(t, "input_schema", getattr(t, "inputSchema", {}))
# 闭包生成执行函数
def _make_handler(target_conn: MCPConnection, target_name: str):
async def _handler(args: dict[str, Any]) -> ToolResult:
return await target_conn.call_tool(target_name, args)
return _handler
wrapped = Tool(
func=_make_handler(conn, tool_name),
name=tool_name,
description=t.description or "",
raw_schema=schema, # 关键点 1:直接透传远程 Schema
timeout=120.0,
is_parallel_safe=True, # 关键点 2:标记为只读并发安全
)
wrapped_tools.append(wrapped)
return wrapped_tools
这里解决了三个核心工程难点:
Python 循环中的闭包陷阱(Closure Late-Binding Trap): 如果在 for 循环里直接写 async def _handler(args): return await conn.call_tool(t.name, args),由于 Python 延迟绑定的特性,所有工具最终都 会调成最后一个工具! 通过 _make_handler(conn, tool_name) 独立工厂函数,确保每个工具在内存中精确绑定属于自己的 conn 和 tool_name。
raw_schema 机制解耦 Python 函数签名: 普通工具需要 Python 函数类型注解(如 def add(a: int) -> int)来推导 JSON Schema;但 MCP 工具的代码在远程,框架通过 raw_schema=schema 直 接透传远程给的 JSON Schema 字典,大模型能立刻看懂入参结构。
is_parallel_safe=True 赋予并发加速能力: 包装出的工具自动具备前文讲过的只读并发能力(asyncio.gather 批执行),当模型同时调用多个 MCP 工具时可以并行执行。
### 四、安全清理:close_all()
1
2
3
4
5
6async def close_all(self) -> None:
"""异步关闭所有连接。"""
for conn in self.connections.values():
with contextlib.suppress(Exception):
await conn.close()
self.connections.clear()
- 资源防泄漏:遍历所有连接逐个调用 close()(进而触发 AsyncExitStack.aclose() 回收 Stdio 子进程);
- 容错关闭:使用 contextlib.suppress(Exception),即使某一个子进程已经意外退出了,也不会中断其他正常子进程的回收;
- 重置连接池:清空字典。
### 总结:MCPClientManager 的定位
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15 .mcp.json
│
▼ (load_config)
MCPServerConfig 列表
│
▼ (connect_server)
┌────────────────────────────────────────┐
│ MCPClientManager │
│ │
│ • conn1 (github) ──► 产出 Tool 列表│──► 统一注入 Agent 注册表
│ • conn2 (filesystem) ──► 产出 Tool 列表│
└──────────────────┬─────────────────────┘
│
▼ (close_all)
一键安全回收全部子进程
它是一个非常干净的装配器(Assembler)+ 门面(Facade),屏蔽了多进程管理的复杂性,对外只提供配置加载、批量连接和统一关闭能力。