my-pi-agent--extension机制与mcp

pi 的 extension 机制

extension 就是一个 TypeScript 模块,默认导出一个工厂函数,拿到一个 ExtensionAPI 对象,往里面注册东西。

定义在 core/extensions/types.ts:1193ExtensionAPI,分六大类:

类别 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
               ┌────────────────────────────────────────────────────────┐
│ Extension Factory Entry │
│ export default function (pi: ExtensionAPI) │
└───────────────────────────┬────────────────────────────┘

┌─────────────────────────────┴─────────────────────────────┐
▼ (启动装配阶段) ▼ (运行时触发阶段)
┌────────────────────────────────────────┐ ┌────────────────────────────────────────┐
│ ExtensionAPI (pi 对象) │ │ ExtensionContext (ctx 对象) │
│ 【静态注册面】 │ │ 【动态运行上下文】 │
│ - pi.registerTool(...) │ │ - ctx.ui (弹窗/选择/通知/底部状态栏) │
│ - pi.on(event, handler) │ ──触发事件注入──► │ - ctx.sessionManager (会话树只读访问) │
│ - pi.registerCommand(...) │ │ - ctx.modelRegistry (模型与认证) │
│ - pi.registerProvider(...) │ │ - ctx.signal (中止信号 AbortSignal) │
│ - pi.sendMessage() / sendUserMessage()│ │ - ctx.cwd / ctx.mode / ctx.hasUI │
└────────────────────────────────────────┘ └────────────────────────────────────────┘
  1. ExtensionAPI(pi 实体):
    • 传给扩展入口函数的参数;
    • 负责声明“这个扩展有什么能力”(注册了哪些工具、订阅了哪些事件、扩展了哪些斜杠命令)。
  2. ExtensionContext(ctx 实体):
    • 当事件触发、命令执行或工具被调用时,作为运行时参数动态传入 Handler;
    • 负责提供“当前环境的即时上下文与交互能力”(当前会话树、TUI 交互接口、Abort 取消信号等)。

生命周期拦截

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
用户输入 (User Prompt)

├─► [1. input 事件] ──────────► 扩展可直接 handled (不调LLM) 或 transform (改写Prompt)
├─► [2. before_agent_start] ──► 扩展可动态注入上下文消息、改写系统提示词 SystemPrompt


进入 Agent ReAct 循环

├──► [3. context 事件] ───────► 发送给模型前,扩展可非破坏性过滤/压缩历史 messages
├──► [4. before_provider_request] ──► 拦截并修改发往 OpenAI/Anthropic 的原始 HTTP 请求 Payload

│ LLM 返回 tool_call:
│ ├──► [5. tool_call 事件] ──► 扩展可原地修改 args 参数,或返回 { block: true } 拦截!
│ ├──► (执行工具真实逻辑)
│ └──► [6. tool_result 事件] ─► 扩展可后置篡改返回给模型的 content / details

└──► [7. turn_end / agent_end]
  • 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

扩展是如何把一个普通函数变成大模型能调用的工具的?

  1. 自动契约合成: @api.tool 内部复用框架底层的 Pydantic 动态建模,自动从 Python 函数签名、类型标注和 docstring 提取生成标准的 OpenAI / Anthropic Function Calling Schema。
  2. 后加载覆盖机制(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. 【工具注册】:让大模型可以在思考中记录待办
@api.tool(description="添加一条待办事项")
def todo_add(item: str) -> str:
TODOS.append(item)
return f"已成功添加待办: {item}"

# 2. 【命令调度】:用户在终端输入 /todos 时直接打印,不花大模型 Token
@api.command("todos", description="查看所有待办列表")
def cmd_todos():
if not TODOS:
return "当前没有待办事项。"
return "\n".join(f"{i+1}. {t}" for i, t in enumerate(TODOS))

# 3. 【事件订阅】:每次对话轮次彻底结束时,若有待办则在终端输出提示
@api.on(AgentEnd)
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
5
1. 实时插队纠偏 (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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// index.ts:460
export default function (pi: ExtensionAPI) {
pi.registerTool({
name: "subagent",
label: "Subagent",
description: "Delegate tasks to specialized subagents with isolated context...",
parameters: SubagentParams, // TypeBox schema
async execute(_toolCallId, params, signal, onUpdate, ctx) {
// 1. 查找指定子代理定义 (agents/<name>.md)
// 2. 创建独立的子 Session 实例 (Fresh Context)
// 3. 运行子 Agent 独立 ReAct 循环
// 4. 返回子 Agent 的最终总结文本
},
renderCall(...) { /* 怎么画调用 */ },
renderResult(...) { /* 怎么画结果 */ },
});
}

如何利用 extension 实现mcp

MCP(Model Context Protocol)连接器同样不需要作为框架核心的硬编码逻辑,而是通过 Extension 机制以“外部工具动态桥接”的方式实现:

1
2
3
4
5
6
MCP 扩展启动 (extension(api))
├── 1. 读取配置文件 (.mcp.json),启动 MCP Server 子进程 (Stdio / SSE 连接)
├── 2. 与 MCP Server 进行握手与协议初始化 (initialize)
├── 3. 调用 tools/list 获取 MCP Server 暴露的所有外部工具定义
├── 4. 将 MCP 工具批量翻译为本地 Tool 对象
└── 5. 循环调用 api.register_tool(translated_tool) 注册进 Agent!
  1. 协议转换:Extension 在 initialize 握手后,通过 tools/list 拿到 MCP Server 声明的所有工具与 Schema;
  2. 工具注册:遍历 MCP 工具,将每个外部工具动态包装为本地的 Tool,调用 api.register_tool 注入 Agent;
  3. 调用转发:当模型调用该工具时,包装函数的 execute 拦截调用,将参数打包为 tools/call 请求通过 Stdio/SSE 转发给 MCP Server 进程,并把执行结果原样喂回给大模型。

extensions.py

在我们的 Python 框架中,提炼了 pi 的核心设计,通过 extensions.py 实现了极简的框架级扩展三件套:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
┌──────────────────────────────────────────────────────────────┐
│ ExtensionManager │
│ - 目录发现 (discover: 递归扫描 **/*.py,跳过 _ 开头文件) │
│ - 动态加载 (load_extension: importlib 加载模块) │
│ - 坏扩展隔离 (load: try...except 捕获异常,仅告警不崩进程) │
│ - 命令调度 (handle_command: 参数个数自适应) │
└──────────────┬───────────────────────────────────────────────┘
│ 构造并持有

┌──────────────────────────────────────────────────────────────┐
│ ExtensionAPI │
│ 1. 事件订阅: api.on(EventCls, handler) -> (event, api) 双参 │
│ 2. 工具注册: @api.tool(...) / api.register_tool │
│ 3. 命令注册: @api.command("name") / api.register_command │
└──────────────────────────────────────────────────────────────┘

ExtensionAPI

在 my-pi-agent 框架中,ExtensionAPI 是整个扩展机制的 “开发者契约面(Developer Surface)”。

当外部写一个扩展(如 mcp.py 或 .agents/extensions/my_tool.py)时,入口函数接收到的唯一参数就是 api: ExtensionAPI:

1
2
3
def extension(api: ExtensionAPI):
# 扩展开发者只跟 api 对象打交道
...

ExtensionAPI 的定位是 “将 Agent 内部复杂的事件系统(HookRegistry)、工具注册表(ToolRegistry)和命令系统,收敛为最简单 直观的三套对外接口”。

核心能力 1:事件订阅系统 —— api.on

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
class ExtensionAPI:
@overload
def on(
self, event_cls: type[Event], handler: None = None
) -> Callable[[Callable[..., Any]], Callable[..., Any]]: ...

@overload
def on(self, event_cls: type[Event], handler: Callable[..., Any]) -> None: ...

def on(
self, event_cls: type[Event], handler: Callable[..., Any] | None = None
) -> Any:
"""注册事件 handler(可作装饰器)。handler 签名 (event, api):
返回 None=观察,返回 HookResult=干预。"""

def _register(h: Callable[..., Any]) -> Callable[..., Any]:
def wrapped(event: Event):
# 关键点:将底层的单参 (event) 包装并注入 (event, self) 双参
return h(event, self)

self.agent.hooks.register(event_cls, wrapped)
return h

if handler is not None:
_register(handler)
return None
return _register

统一干预数据结构:HookResult 与五大决策点

所有生命周期拦截点共享统一强类型的 HookResult 数据结构:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
@dataclass(frozen=True)
class HookResult:
"""Hook 回调的干预结果。返回 None = 纯观察,返回 HookResult = 干预。"""

block: bool = False
reason: str | None = None

# 决策点 1 (input): 改写用户输入文本
updated_input: str | None = None

# 决策点 2 (before_agent_start): 动态改写 System Prompt
updated_system_prompt: str | None = None

# 决策点 3 (context): 临时改写发给大模型的 messages 视图(不污染 Session 树)
updated_messages: list[Message] | None = None

# 决策点 4 (tool_call): 改写工具入参
updated_args: dict | None = None

# 决策点 5 (tool_result): 改写工具出参
updated_result: str | None = None

框架在 Agent.run() 中完整对齐了 Pi 的 五大生命周期决策拦截点: 1. UserInputinput):在用户输入进入 Session 之前触发,支持 block=True 阻断输入或 updated_input 改写输入; 2. AgentStartbefore_agent_start):在准备好系统消息后触发,支持 updated_system_prompt 动态更新首条 system 消息; 3. BeforeModelCallcontext):在 _ctx.prepare() 产出 view 后、调用大模型前触发,支持 updated_messages 临时改写送给大模型的视图; - 关键架构不变式updated_messages 仅临时修改当前这次发给 LLM 的 view 变量,真实 self.messages 与 Session 磁盘 JSONL 保持绝对纯净,实现 零历史污染; 4. ToolExecutionStarttool_call):工具执行前触发,支持 block=True 拦截危险调用或 updated_args 改写入参; 5. ToolExecutionEndtool_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
14
def 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
18
def 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
21
1. 启动装配期 (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
@dataclass
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
5
class 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
    14
    async 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
2
3
4
5
6
7
8
9
10
11
12
13
┌──────────────────────────────────────────────────────────┐
│ MCP 架构分层 │
├──────────────────────────────────────────────────────────┤
│ │
│ 【第二扇门】ClientSession (协议层) │
│ • list_tools() • call_tool() • JSON-RPC 消息解析 │
│ ─────────────────────────┬──────────────────────────── │
│ │ 依赖底层流传输 │
│ ▼ │
│ 【第一扇门】stdio_client (传输层) │
│ • 启动 Node.js/Python 进程 • 管道管理 (stdin/stdout) │
│ │
└──────────────────────────────────────────────────────────┘
  1. 先有第一扇门,才有第二扇门: 如果没有 stdio_client 启动进程并提供 read/write_stream,ClientSession 就根本不知道该向哪里读写协议数据。

  2. 关闭时必须倒序退出:

    • 先关 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
10
1. 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
3
class 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
21
def 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
31
async 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

这里解决了三个核心工程难点:

  1. 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。

  2. raw_schema 机制解耦 Python 函数签名: 普通工具需要 Python 函数类型注解(如 def add(a: int) -> int)来推导 JSON Schema;但 MCP 工具的代码在远程,框架通过 raw_schema=schema 直 接透传远程给的 JSON Schema 字典,大模型能立刻看懂入参结构。

  3. is_parallel_safe=True 赋予并发加速能力: 包装出的工具自动具备前文讲过的只读并发能力(asyncio.gather 批执行),当模型同时调用多个 MCP 工具时可以并行执行。

### 四、安全清理:close_all()

1
2
3
4
5
6
async 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),屏蔽了多进程管理的复杂性,对外只提供配置加载、批量连接和统一关闭能力。