my-pi-agent--异步支持
为什么需要异步,哪里需要异步
1 | 1. 模型层 (packages/my-agent-llm) 【已就绪!】 |
## 价值 1:工具并发执行(Parallel Tool Execution)
- 场景:当 Agent 在做代码重构或审查时,大模型经常会一次性输出 5 个文件读取请求(read(“a.py”), read(“b.py”), read(“c.py”)…)。
- 价值:同步只能逐个读;异步可以直接用 asyncio.gather 同时并发读取 5 个文件或调用 5 个外部 API,耗时直接从 O(N) 降到 O(1)。
## 价值 2:打字机流式输出与实时事件(Streaming UX)
- 场景:大模型生成一段长篇代码(可能需要 15 秒)。
- 价值:
- 同步只能等 15 秒后一次性拿到完整文本;
- 异步通过我们在 events.py 中预留的 MessageUpdate 与 ToolExecutionUpdate,可以做到 Token 级别的毫秒级流式回显,在 终端或 Web 前端呈现丝滑的打字机动画。
## 价值 3:即时可取消性与信号响应(Graceful Cancellation)
- 场景:大模型写出了一个死循环 bash 命令,或者正在进行一个错误的超长推理。
- 价值:
- 在异步事件循环中,用户按一次 Esc 或 Ctrl+C,框架可以在事件循环中触发 asyncio.CancelledError,毫秒级掐断网络连接 、杀掉运行中的子进程,并把当前已产生的上下文安全存盘,绝不破坏 Session 树。
## 价值 4:多用户并发服务能力(Web / Server 部署)
- 场景:未来如果你把这个 Agent 包装成一个 FastAPI 接口,或者部署到 Web 平台供多人同时使用。
- 价值:
- 如果是同步代码:一个用户提问(耗时 10 秒),这个 Python 进程线程就被占死,其他所有用户全都在排队等待;
- 如果是异步代码:单个 Python 进程可以轻松并发服务成百上千个用户同时对话。
工具并发
如何防止工具并发冲突与因果时序倒置
第一道机制:is_parallel_safe 声明与“一票否决”因果保护
#### 1. 要解决的核心问题:防止“并发写冲突”与“因果时序倒置(Causal Inversion)”
假设大模型在一轮推理中,同时发起了 2 个有先后因果关系的工具调用: - tool_0: write(“config.json”, “port=8080”)(写操作,有副作用) - tool_1: read(“config.json”)(只读,目的是确认刚才写入的新配置)
如果粗暴地把只读工具(tool_1)提前拿去并发执行,就会导致
read 先于 write
发生,读出未修改前的旧数据,发生严重的因果时序倒置
Bug!
#### 2. 机制实现原理:Pi 风格的“一票否决制”
在工具定义时,每个工具声明自己是否是“并发安全的(只读无副作用)”: -
@tool(is_parallel_safe=True):如
read、get_weather、db_query(只读); - 默认
is_parallel_safe=False:如
write、edit、bash(有写操作/副作用)。
在 ToolRegistry.execute_batch
执行前,采用一票否决判定: -
全员只读放行并发:当且仅当批次中的每一个工具均为
is_parallel_safe=True 时,才使用
asyncio.gather 全并发执行; -
含写全批保序串行:只要批次中包含任何一个有写副作用的工具(或未知工具),整批工具立即放弃并发,严格按照大模型输出的原始先后顺序依次串行执行!
1
2
3
4
5
6
7
8
9场景 A (全只读): [ 0: read(A), 1: read(B), 2: grep(C) ]
│
▼ (全员 is_parallel_safe=True)
【asyncio.gather 全员并发加速 ⚡】
场景 B (含写入): [ 0: write(A), 1: read(A) ]
│
▼ (检测到 write: 一票否决降级)
【严格按 0 ➔ 1 原始因果顺序串行执行 🛡️】
────────────────────────────────────────────────────────────────────────────────
### 第二道机制:asyncio.gather 并发(极速执行)
#### 1. 要解决的核心问题:消灭累加的“网络与 I/O 等待耗时”
传统的同步模式下,读 3 个远程 API 耗时 2s + 2s + 2s = 6s。
#### 2. 机制实现原理:
当批次判定为全员只读安全时,使用 asyncio.gather
同时把所有协程扔进事件循环,让操作系统在后台同时发起网络/文件
I/O。耗时取决于最慢的那个,直接从 O(N) 降到 O(1)。
1
全只读批次 ──► asyncio.gather(task_0, task_1, task_2) ──► 多个只读工具在后台同时跑 ⚡
────────────────────────────────────────────────────────────────────────────────
### 第三道机制:严格保序回填(协议对齐)
#### 1. 要解决的核心问题:大模型对 tool_call_id 的强顺序依赖
大模型在发请求时,它心里的顺序是:[0号工具, 1号工具, 2号工具]。 - 在并发执行时,2 号工具可能 0.1 秒就跑完了,而 0 号工具花了 3 秒; - 如果我们按“谁先跑完谁先写回”,消息列表会变成 [2号结果, 1号结果, 0号结果] → 大模型 API 校验直接报错,或者把 2 号的结果误当成 0 号的结果!
#### 2. 机制实现原理:“结果索引保序对齐”
无论工具是并发执行还是串行回退执行,返回的 ToolResult
列表严格与入参 tool_calls 的索引位置完全对齐:
1
2
3
4大模型 tool_calls 列表: [ 0: call_0, 1: call_1, 2: call_2 ]
│
▼
生成的 results 结果列表: [ res_0, res_1, res_2 ] (索引 100% 完美对应!)
然后,框架按照这个保序的数组逐条生成 role: "tool"
消息追加到 messages 与 session
中,安全喂给下一轮大模型。
────────────────────────────────────────────────────────────────────────────────
### 总结工具批调度的决策全景图
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19大模型返回工具批次: [ tool_0, tool_1, ... ]
│
▼
【检查是否包含任何写操作/非安全工具?】
/ \
[ 是 ] [ 否 ]
/ \
▼ ▼
【一票否决: 串行降级】 【全员并发: asyncio.gather】
严格按 0 ➔ 1 顺序执行 全员只读工具同时发起 I/O ⚡
│ │
└──────────┬──────────┘
│
▼
【严格保序回填入库】
与原始 tool_call_id 严格对齐
│
▼
安全、合规、零时序倒置地喂回大模型!
withFileMutationQueue 细粒度文件锁
1. 思考一个问题:写文件工具(edit / write)到底能不能并发?
粗暴的思路:“写文件有副作用,所以写工具必须串行!”
但是请看这个真实场景: 大模型决定重构项目,在同一轮里发起了两个操作:
- ToolCall 1: edit(path=“src/login.py”)(修改登录逻辑)
- ToolCall 2: edit(path=“src/user.py”)(修改用户逻辑)
这两个工具修改的是两个完全不同的文件!它们有任何冲突吗?没有! 如果强行串行,修改 10 个文件就要等 10 倍时间;如果并发修改,耗时直接除以 10!
────────────────────────────────────────────────────────────────────────────────
### 2. 那什么情况下会冲突?
只有当大模型在同一轮里,同时发起两个针对“同一个文件”的编辑时才会冲突: - ToolCall 1: edit(path=“src/app.py”, old=“v1”, new=“v2”) - ToolCall 2: edit(path=“src/app.py”, old=“v3”, new=“v4”) 如果这两个并发跑,就会互相踩踏、文件内容被覆盖损毁。
────────────────────────────────────────────────────────────────────────────────
### 3. Pi 的神级解法:工具内置“文件级互斥锁”
Pi 并没有在外层粗暴地把 edit 标记为全局串行,而是做了两层防线:
- 第一层(外层):7 个内置工具默认全都可以参与并发(executionMode: “parallel”);
- 第二层(工具内部):在 edit.ts 内部,维护一个以文件绝对路径为 Key 的锁字典(withFileMutationQueue)。
1
2
3
4
5
6
7
8
9场景 A:同时修改不同文件 (login.py 与 user.py)
login.py ──► 获取 login.py 的文件锁 ──► 执行编辑 ──► 释放锁 ┐
├─► 真正完全并发!⚡ (极速)
user.py ──► 获取 user.py 的文件锁 ──► 执行编辑 ──► 释放锁 ┘
场景 B:同时修改同一个文件 (app.py 和 app.py)
操作 1 (app.py) ──► 抢到 app.py 锁 ──► 正在编辑 app.py...
操作 2 (app.py) ──► 发现 app.py 锁被占用 ──► 在队列中排队等待
操作 1 完成 ──────► 释放锁 ──► 唤醒操作 2 ──► 执行编辑 ──► 安全无损!🛡️
取消控制
为什么异步能进行取消控制
一、为什么“同步代码”根本无法被优雅叫停?
要理解异步为什么能取消,先看看同步代码在操作系统底层发生了什么。
#### 1. 同步的底层机制:操作系统级阻塞(Blocked Thread)
当你执行同步代码 resp = requests.post(…) 或 client.chat(…) 时:
1
[你的 Python 代码] ──► 调用操作系统网络底层 recv() ──► 线程进入 SLEEP 状态
- CPU 视角的同步: 操作系统直接把你的 Python 线程“冻结”了,告诉 CPU:“这个线程在等网卡数据,先不要给它分配任何算力”。
- 为什么叫不醒? 因为当前线程已经被冻结了,它根本无法执行下一行 Python 代码! 你哪怕在旁边定义了一个 stop() 函数,只要程序卡在同步网络 I/O 里,你的 stop() 函数就连被执行的机会都没有。
- 唯一能叫醒它的方式: 只能靠极其暴力的手段——用户狂按 Ctrl+C 发送操作系统的 SIGINT 中断信号,或者直接在任务管理器里 kill 掉整个进程。这会导致所有内存状态瞬间撕裂,文件损坏。
────────────────────────────────────────────────────────────────────────────────
### 二、为什么“异步代码”随时可以被叫停?
异步(asyncio)彻底改变了程序等待网络的方式:它使用的是 非阻塞 I/O(Non-blocking I/O) + 事件循环(Event Loop)。
#### 1. 核心密码:await 的“主动交权机制”
在异步代码中,每当你写下一个 await(例如 await stream 或 async for chunk in stream:):
1
2async for chunk in self.llm.achat_stream(...): # <-- 这里的 async 背后就是一个 await!
...
它的底层真实动作是: 1. “交出控制权”:Python 告诉操作系统:“我发起了一个非阻塞网络请求,但我不冻结线程,我把 CPU 控制权交还给事件循环(Event Loop)”; 2. “事件循环插队”:在网卡收到下一个 Token 的这 20~50 毫秒空档期里,事件循环可以从容地执行其他事情(例如:检查用户有没有按 Esc、执行用户的 /stop 指令、检查 self._aborted 状态); 3. “随时变卦”:如果在等待下一个 Token 的过程中,事件循环检测到中止信号,事件循环可以在下一次恢复该协程时,直接给它注入一个取消指令,不再继续读网卡!
────────────────────────────────────────────────────────────────────────────────
### 三、图解:同步 vs 异步的取消对比
#### 同步场景(高速公路上刹不住):
1
2
3
4[发起网络请求] ──────────────────────────────────────────► [等了 10 秒拿结果]
▲
│ 线程被操作系统完全冻结 (无法执行任何代码)
用户按 Esc: "我想停!" ──► 操作系统: "线程冻结中,听不见!"
#### 异步场景(每个 Token 都有服务区可以下高速):
1
2
3
4
5
6
7
8
9[发请求] ──► await ──► [第 1 个字] ──► await ──► [第 2 个字] ──► await ──► ...
│ │ │
▼ ▼ ▼
【交出 CPU】 【交出 CPU】 【交出 CPU】
事件循环检查状态: 事件循环检查状态: 事件循环检查状态:
一切正常,继续跑~ 一切正常,继续跑~ ⚡ 发现用户调了 abort()!
│
▼ 立即 break!
【毫秒级下高速,掐断连接】
取消控制的机制
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【触发源 A:用户按 Esc / 点击停止】 【触发源 B:安全扩展实时监控】
│ │
▼ ▼
调用 agent.abort() ─────────────────┐ @api.on(MessageUpdate)
│ │ 发现模型正在输出高危指令
▼ │ │
内部标记: self._aborted = True │ ▼
│ │ 返回 HookResult(block=True)
│ │ │
└─────────────────────────┼────────────────────┘
│
▼
┌───────────────────────────────────┐
│ 流式循环在下一个毫秒瞬间捕获中断: │
│ async for chunk in achat_stream: │
│ if self._aborted: break │
└─────────────────┬─────────────────┘
│
▼
┌───────────────────────────────────────────────────────────┐
│ 核心四连动作与自愈保护 (The 4 Actions): │
│ 1. 掐断网络连接,停止烧 Token │
│ 2. 丢弃未完成的文本半截内容 (防止模型断句幻觉) │
│ 3. 悬空断头调用自愈:若已产生 tool_calls,立即自动合成 │
│ "Tool call interrupted by user" 补齐,杜绝 API 400 死锁│
│ 4. 广播 AgentEnd(cancelled) 并原子落盘 │
└─────────────────────────────┬─────────────────────────────┘
│
▼
优雅返回 "(cancelled)",会话历史 100% 具备自愈完备性
机制 1:双轨触发通道与断头调用自愈(谁可以叫停大模型?)
你的设计提供了两条平行的触发通道,兼顾了 “人类用户交互” 与 “程序自动化安全拦截”,并在底层织入了严密的转录本自愈网:
通道 ①:面向用户交互 —— agent.abort() 与断头自愈 (tool_history.py)
- 场景:用户在终端按下了 Esc 键、Web 前端点击了“停止生成”按钮、或者输入了 /stop 命令;
- 做法:调用
agent.abort(),内部触发取消令牌CancellationToken.cancel()并清理后台进程树;
1. CancellationToken(协作式取消令牌)解决的核心痛点
为什么不能直接使用系统底层的 asyncio.Task.cancel()
粗暴杀死协程? -
粗暴强杀的致命陷阱:Task.cancel()
会在不可预测的任意一行 Python 代码上直接抛出
asyncio.CancelledError。如果恰好在大模型刚刚流式输出了
tool_calls: [{"id": "call_123"}]、但本地尚未开始执行工具并落盘结果的“微秒级时间窗口”内强杀了协程:
-
本地会话历史中就会永久留下一条“声称发起了工具调用、却永远没有工具返回结果”的断头
Assistant 消息; -
下一轮用户再发起任何提问,这段畸形历史被直接发给 OpenAI / Anthropic
时,云端 API 会当场抛出致命的 HTTP 400 校验错误: >
400 Bad Request: An assistant message with 'tool_calls' must be followed by tool messages responding to each 'tool_call_id'
- 整个 Session
历史彻底报废并陷入永久死锁,用户除了清空会话别无他法!
2. CancellationToken 的协作式拉手刹机制
CancellationToken 坚守协作式取消原则(Cooperative
Cancellation),不搞暴力撕裂: 1. 信号打标:宿主调用
token.cancel(),仅将内部布尔状态原子化标记为
_cancelled = True; 2. 微内核安全检查点(Safety
Checkpoints)轮询:调度循环在关键时序节点主动检测
if signal.is_cancelled():: -
流式吐字阶段:每收到一个 Token
增量检测一次,发现取消立刻 break
掐断流式读取,丢弃未完半截文本; -
工具执行前夕:在批量执行工具前再次检测,发现取消坚决不再调用真实工具(防止误跑高危或不可逆命令);
- 断头自愈与优雅终结:由微内核集中调用
_synthesize_interrupted_tool_calls,统一将大模型声称要调用的工具自动合成一条带有
[INTERRUPTED] Tool call interrupted by user(is_error=True)的虚拟工具消息追加落盘,发射配对闭合的
TurnEnd,最后以 stop_reason="cancelled"
优雅收尾。 -
收益:既实现了毫秒级响应用户的取消操作,又绝对不污染会话树,保证转录本拓扑时刻合法。
3. _provider_context(大模型上游上下文防卫清洗器)的两重净化
本地磁盘的
Session(JSONL)为了审计与追溯,必须如实记录一切发生的事实(包括网络抖动断流、报错轮次与用户中途取消的轮次)。但
OpenAI / Claude 等上游 API 拥有极其严苛、容错率为零的格式校验规则(拒绝
role="assistant", content=""
空消息、拒绝断头调用与孤儿结果)。
loop.py 中的 _provider_context
在每次调用大模型前 1ms 筑牢了两道净化门: 1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16def _provider_context(messages: Sequence[Message]) -> list[Message]:
# 第一重净化:把因取消、abort 或报错残留的、正文为空的空白 assistant 消息全部剥离!
replayable = tuple(
m
for m in messages
if not (
m.role == "assistant"
and bool(
m.metadata
and m.metadata.get("stop_reason") in {"error", "aborted", "cancelled"}
)
and not m.content
)
)
# 第二重净化:串联 repair_tool_history 拓扑自愈引擎,自动修剪孤儿 ToolResult、补齐断头调用
return list(repair_tool_history(replayable).messages)
#### 通道 ②:面向扩展安全监控 —— MessageUpdate(Interceptable)
- 场景:扩展或安全看门狗在模型吐字的过程中,实时检测内容(如发现模型正在输出 rm -rf / 或包含敏感词);
- 做法:
1
2
3
4
5
def safety_guard(event: MessageUpdate, api):
if "rm -rf" in event.message.content:
# 发现危险,直接返回 HookResult 紧急熔断!
return HookResult(block=True, reason="触发安全规则紧急熔断") - 优点:复用了框架底层的 HookRegistry 拦截机制,完全零新增概念。
────────────────────────────────────────────────────────────────────────────────
### 机制 2:毫秒级响应与即时掐断(怎么停下来的?)
在大模型流式生成(achat_stream)过程中,每当大模型吐出一个 Token(约 20~50 毫秒):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16async for chunk in self.llm.achat_stream(...):
# 1. 检测通道 ① (是否被 agent.abort 叫停)
if self._aborted:
break
if chunk.content:
content_acc += chunk.content
# 2. 检测通道 ② (发射 MessageUpdate 并检查是否有 Hook 拦截)
hook = await self._emit(MessageUpdate(
message=Message(role="assistant", content=content_acc),
chunk=chunk
))
if isinstance(hook, HookResult) and hook.block:
self._aborted = True
break
- 毫秒级刹车:无论哪个通道触发,在下一个 Token 出来的瞬间,循环 break,异步流式连接被底层的 Python 异步生成器自动关闭(关闭 Socket 避免继续消耗 Token)。
中途中断后,云端 API 还会继续跑吗?会浪费 Token 扣费吗?
结论:在“异步流式(Streaming)”下只要正确关闭连接,云端 API 会立刻停止生成,绝不会继续浪费后续的 Token!
### 底层通信与计费机制深度揭秘
1
2
3
4
5
6
7
8
9
10
11
12
13[用户按 Esc 中断] ──► 你的客户端代码 break 退出 async for 循环
│
▼ (关键动作: Python 自动触发 aclose())
【客户端主动切断底层 TCP 连接 (发送 TCP FIN 包)】
│
▼
【OpenAI / Anthropic 云端网关检测到连接断开 (Broken Pipe)】
│
▼
【云端 GPU 推理引擎立即杀掉该 Generation Task】
│
▼
【计费停止:只结算断开前已生成的少量 Token】
- 流式(Streaming)计费规则:
大模型服务商(OpenAI、Anthropic、DeepSeek)的计费引擎是按实际生成的
Token 计费的。当你主动断开流式连接时,服务商网关捕捉到 Socket
关闭,会立即终止后续几千个 Token 的推理。
- 比如原本要生成 4000 个 Token(价值 0.1 美元),在生成到第 100 个 Token 时你按了取消;
- 云端立刻刹车,你只需要付这 100 个 Token 的费用,后 3900 个 Token 彻底停止,零扣费!
- 什么时候会白白浪费 Token?(反模式警告): 如果你只在前端 UI
层面把文字隐藏了,但底层的 Python/Node
进程没有真正退出流式循环(没有调用 stream.close()),那么底层 HTTP
连接依然连着,云端就会傻傻地把 4000 个 Token 全跑完并计入 账单。
- 在我们的异步方案中:async for chunk in achat_stream: 一旦被 break 或取消,Python 异步生成器会自动触发 aexit 关闭底层 HTTP 连接,从源头确保云端停止计费。
多agent的并发支持
## 同步 vs 异步:执行模式对比
### 1. 在同步模式下(现在的串行阻塞)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15[主 Agent] ──► 派发 task_0 (code-reviewer)
│
▼ 【等待 10 秒...】 (主线程完全卡死)
子代理 0 跑完 3 轮 ReAct 循环,返回审查报告
│
▼
接着派发 task_1 (researcher)
│
▼ 【再等待 10 秒...】 (主线程继续卡死)
子代理 1 跑完 3 轮 ReAct 循环,返回调研报告
│
▼
[主 Agent 汇总回答]
⏳ 总耗时: 10s + 10s = 20 秒!
────────────────────────────────────────────────────────────────────────────────
### 2. 在异步并发模式下(改造后的多代理并行协同)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18[主 Agent] ──► 一次性发起 2 个委派任务!
│
▼ 【asyncio.gather 并发启动】
┌────────────┴────────────┐
│ │
▼ 【后台并发执行】 ▼ 【后台并发执行】
[子代理 0: code-reviewer] [子代理 1: researcher]
- 独立 Session 1 - 独立 Session 2
- 独立 ReAct 循环 - 独立 ReAct 循环
- 独立调 LLM & 工具 - 独立调 LLM & 工具
(耗时 10 秒) (耗时 10 秒)
│ │
└────────────┬────────────┘
│ ⚡ 两者在后台同时跑完 (总耗时仅 10 秒!)
▼
[主 Agent 收到两个结果,保序汇总输出]
⚡ 总耗时: max(10s, 10s) = 10 秒 (时间直接减半!)!
────────────────────────────────────────────────────────────────────────────────
## 底层为什么能安全并发?(3 大隔离基石)
你可能会问:两个子 Agent 同时在后台跑,会不会把会话搞乱?会不会互相冲突?
答案是:绝对不会!因为我们在之前的阶段已经打下了极其完美的 3 大隔离地基:
### 1. 会话文件物理级隔离(Session Isolation)
在 TaskManager._run 中,每个子代理拥有独立的 Session 文件:
1
2
3
4
5.my_agent_core/sessions/
├── parent_session.jsonl <-- 父会话文件(完全不被子代理的消息污染)
└── subagents/
├── agent-task_00000001.jsonl <-- 子代理 0 的独立 ReAct 完整记录
└── agent-task_00000002.jsonl <-- 子代理 1 的独立 ReAct 完整记录
- 两个子代理各自向不同的磁盘文件追加 JSONL,完全零文件写锁冲突!
### 2. 内存上下文隔离(Fresh Context)
- 每个子 Agent 都是一个全新的 Agent 类实例;
- 它们拥有各自独立的 self.messages = [],各自独立的 system_prompt;
- 它们在内存中没有共享的可变状态,天然满足并发安全性(Thread/Coroutine-Safe)。
### 3. 工具标记天然安全(is_parallel_safe=True)
在 tools/builtin/task.py 中:
1
2
3def make_task_tool(manager: SubagentManager, parent: "Agent") -> Tool:
...
return Tool(func=task, name="task", is_parallel_safe=True)
- 因为子代理拥有上述独立的上下文与会话,所以 task 工具被天然标记为 is_parallel_safe=True(并发安全);
- 主 Agent 的 ToolRegistry.execute_batch 看到 task 是只读/并发安全的,就会自动把多个 task 放入 asyncio.gather 同时并发执行!
────────────────────────────────────────────────────────────────────────────────
## 代码执行链路一览
在 TaskManager 与 Agent 异步化后,整个调用链变得极度丝滑:
1
2
3
4
5
6
7
8
9
10
11
12
13# 1. 主 Agent 收到模型发来的 2 个 task 调用
# 2. ToolRegistry.execute_batch 检测到它们都是 is_parallel_safe=True
# 3. 自动并发调度:
results = await asyncio.gather(
task_manager.start_task(prompt="审查 api.py", subagent_type="code-reviewer"),
task_manager.start_task(prompt="调研迁移方案", subagent_type="researcher"),
)
# 4. 在后台:
# child_0.run() 和 child_1.run() 在同一个 asyncio 事件循环中并发运转
# 5. 返回结果:
# results[0] 对应审查报告
# results[1] 对应调研报告