my-pi-agent--agent类与hook系统

生命周期事件

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// Agent 生命周期
agent_start
agent_end(messages)

// Turn 生命周期(一轮 = 一次助手响应 + 工具调用/结果)
turn_start
turn_end(message, toolResults)

// 消息生命周期(user / assistant / toolResult 消息都会发)
message_start(message)
message_update(message) // 仅异步流式
message_end(message)

// 工具执行生命周期(ToolExecutionStart/End 可被 hook 干预)
tool_execution_start(toolCallId, toolName, args)
tool_execution_update(toolCallId, toolName, args, partialResult) // 仅异步流式
tool_execution_end(toolCallId, toolName, result, isError)

事件集对齐 pi 的生命周期模型(Agent/Turn/Message/Tool 四组,正常执行 start/end 成对;被拦截/畸形参数的调用不发射 End)。 每个事件自动带 timestamp(Unix 秒,实例化时刻)。

事件的作用

一、事件的作用:把「可观察」变成「可扩展」

一句话:事件机制让 Agent.run() 成为一个「可被插桩的黑盒」——循环本身写死了,但循环的每一步都广播出去,让任何新功能都能「监听而不改框架」。

先看现在 run() 里实际发生了什么:

1
2
3
4
5
6
7
8
9
10
11
12
13
Agent.run("东京天气")

├─ AgentStart
├─ MessageStart(user) → MessageEnd(user) # 你的问题进 transcript
├─ TurnStart(1)
│ ├─ MessageStart(assistant) → MessageEnd(assistant) # 模型想调工具
│ ├─ ToolExecutionStart → ToolExecutionEnd # 执行 get_weather
│ ├─ MessageStart(tool) → MessageEnd(tool) # 结果写回
├─ TurnEnd(message=assistant, tool_results=[...])
├─ TurnStart(2)
│ ├─ MessageStart(assistant) → MessageEnd(assistant) # 模型直接回答
│ └─ TurnEnd(message=assistant, tool_results=[]) # 本轮工作闭环结算
├─ AgentEnd(messages=[...], ...)

这 8 个事件(同步阶段)就是「循环的完整心电图」——每一段代码执行都对应一个事件。而事件的作用 = 把「这段代码发生了」告诉任何想知道的人。

二、核心:为什么事件 = 扩展点

关键洞察在于——你的所有后续功能,几乎都是「想观察循环在干什么」,而不是「想改循环本身」。事件机制让这些功能都变成「加一个 hook」,而不是「改 run()」。

看三个后续阶段怎么用:

扩展 1:session 落盘(阶段 3)—— Agent 显式双写(不是 hook 消费事件流)

1
2
3
4
# 实际实现:Agent.run() 里每加一条消息,同时写内存 + 落盘
# (不走 hook——session 是核心持久化,Agent 显式调用)
self.messages.append(user_msg) # 内存 transcript
self.session.add_message("user", user_input) # 树 + 原子落盘(session 必填,无条件)

(注:早期设计设想过”用 hook 消费事件流写 session”,但实际实现是 Agent 双写——session 是核心功能,显式调用比 hook 更直接、不依赖事件顺序。)

扩展 2:context 压缩(阶段 4)—— prepare 每轮调用(不是监听事件)

1
2
3
4
5
6
# 实际实现:prepare() 在每次 llm.chat 前调用,内部判断是否超阈
# (不走"监听 MessageEnd 计数"——压缩是发送前视图变换,时机在发消息前)
while ...:
view = self._ctx.prepare(self.messages) # 四层管线(超阈 → 摘要)
resp = self.llm.chat(messages=view, tools=...)
...

ContextCompacted 事件早在设计时就预留了(阶段 4 用)——这就是「事件先定义、后发射」的意义:接口提前定好,实现到时接上。

扩展 3:coding agent UI(未来)—— 渲染过程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
# 新功能:TUI 渲染 —— 监听事件更新界面
def ui_render(event):
if isinstance(event, TurnStart):
ui.show_round(event.iteration)
elif isinstance(event, ToolExecutionStart):
ui.show_tool_call(event.tool_name, event.args)
elif isinstance(event, MessageEnd):
ui.show_message(event.message.content)
return None

```python
# 观察者只读订阅(UI 渲染、日志打点)
unsubscribe = agent.subscribe(ui_render)

# 决策拦截门禁(权限检查、参数改写)
agent = Agent(llm=llm, tools=tools, hooks=[
(UserInputHook, input_sanitizer),
(ToolCallHook, permission_guard),
(ToolResultHook, output_redactor),
])
```

hook系统

架构解耦重塑(对标 Pi & Tau):在最新架构中,只读事件(Event决策门禁(*Hook实现彻底的物理与概念解耦: - 只读事件(events.pyAgentStart, TurnEnd, MessageEnd, ToolExecutionStart 等事实通知,统一通过 agent.subscribe(listener) 进行只读旁路监听; - 决策门禁(hooks.pyUserInputHook, AgentStartHook, BeforeModelCallHook, ToolCallHook, ToolResultHook 5 大专职拦截载荷,在 Agent(hooks=[(HookClass, callback), ...]) 中批量注册,负责阻断、参数改写与结果篡改。

多 hook 注册表

核心两个函数:注册(register)和触发(emit)。

注册:把 hook 回调挂到 Hook 目标名下

1
2
def register(self, target: type, callback: Callable):
self._hooks.setdefault(target, []).append(callback)
  • target = Hook 类型(如 ToolCallHook
  • callback = 门禁函数((*Hook) -> HookResult | None
  • 同一个 Hook 类型可挂载多个回调,按注册顺序保序存储

触发:到点时异步逐个调用(支持 async/sync 混合回调)

1
2
3
4
5
6
7
8
9
10
11
12
async def emit(self, event: Any) -> HookResult | None:
"""异步触发 Hook,支持协程与普通函数,返回第一个非 None 结果(短路)。"""
for cb in self._hooks.get(type(event), []):
try:
res = cb(event)
if inspect.isawaitable(res):
res = await res
if res is not None:
return res # 短路:第一个干预结果即刻生效返回
except Exception as exc:
logger.error(f"Hook callback {cb} error: {exc}", exc_info=True)
return None

HookResult 干预模型与五大决策拦截点

Hook 回调返回 None = 放行;返回 HookResult = 强力干预。

1
2
3
4
5
6
7
8
9
10
@dataclass(frozen=True)
class HookResult:
block: bool = False # 是否阻断当前执行
reason: str | None = None # block 阻断时的原因说明
terminate: bool | None = None # 三态:是否请求微内核提前退出当前批次
updated_input: str | None = None # 改写用户原始输入 (UserInputHook 用)
updated_system_prompt: str | None = None # 改写系统提示词 (AgentStartHook 用)
updated_messages: list | None = None # 临时改写送入模型的上下文 (BeforeModelCallHook 用)
updated_args: dict | None = None # 改写工具入参 (ToolCallHook 用)
updated_result: str | None = None # 改写工具执行结果 (ToolResultHook 用)

五大独立生命周期决策拦截点(hooks.py

  1. UserInputHook:截获用户原始输入,支持 block 阻断或 updated_input 前置清洗与改写;
  2. AgentStartHook:启动前拦截,支持 updated_system_prompt 动态更新首条 system 消息;
  3. BeforeModelCallHook:调 LLM 前拦截,支持 updated_messages 临时改写视图(临时 View 改写 vs 真实 Session 零污染);
  4. ToolCallHook:工具执行前拦截,支持 block 拦截高危命令或 updated_args 动态修补参数;
  5. ToolResultHook:工具执行后拦截,支持 updated_result 篡改出参,以及 terminate 提前退出整批 ReAct 循环。

用起来的样子

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
# 注册阶段(应用层,构造 Agent 后)
def guard_hook(event) -> HookResult | None:
"""权限门控(挂到 ToolExecutionStart,可拦截/改参数)"""
if isinstance(event, ToolExecutionStart):
if event.tool_name == "bash" and "rm -rf" in event.args.get("command", ""):
return HookResult(block=True, reason="危险命令被拦截")
if event.tool_name == "write_file":
return HookResult(updated_args={**event.args, "path": abs(event.args["path"])})
return None

def log_hook(event) -> HookResult | None:
"""日志(挂到 ToolExecutionStart,纯观察)"""
if isinstance(event, ToolExecutionStart):
print(f"[hook] {event.tool_name}({event.args})")
return None

agent = Agent(llm=llm, tools=tools, hooks=[
(ToolExecutionStart, guard_hook), # 同一事件挂多个
(ToolExecutionStart, log_hook),
])

机制一句话:注册表 = 事件类 → 回调列表;触发 = 到点时按序调用,短路返回。 你说「在某个事件时调用注册的 hook」——完全正确,_emit 就是那个「到点」的调用点。

Agent类(轻量 Harness 架构演进)

在经历 Tau 对齐重塑(阶段 18)后,Agent 类已经从早期“状态+循环+调度揉在一个类”的上帝类,全面演进为对标 Tau / Pi 的轻量有状态外壳(AgentHarness

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
┌─────────────────────────────────────────────────────────────────┐
│ class Agent (AgentHarness) │
│ "轻量宿主外壳:持有状态 + 暴露原生事件流" │
│ │
│ ┌───────────── __init__ 装配 ─────────────┐ │
│ │ llm ← 注入(模型边界) │ │
│ │ session ← 会话(只追加存储驱动) │ │
│ │ registry ← 工具注册表 │ │
│ │ messages ← 状态/记忆(自愈清洗) │ │
│ │ message_queue ← 干预队列 (Steer/Followup) │
│ │ hooks ← HookRegistry(拦截管线)│ │
│ └─────────────────────────────────────────┘ │
│ │
│ ┌────────── 一等公民核心驱动:prompt_stream() ──────────┐ │
│ │ 异步生成器:run_agent_loop(llm, tools, messages...) │ │
│ │ │ │ │
│ │ ▼ 逐一 yield 生命周期事件流 │ │
│ │ AgentStart ➔ TurnStart ➔ MessageUpdate ➔ TurnEnd... │ │
│ └───────────────────────────────────────────────────────┘ │
│ │
│ ┌────────── 经典便利门面:run() ──────────┐ │
│ │ async def run(user_input): │ │
│ │ async for event in prompt_stream: │ │
│ │ return final_text │ │
│ └────────────────────────────────────────┘ │
│ │
│ ┌────────── 观察者监听:subscribe() ──────┐ │
│ │ subscribe(listener) -> unsubscribe │ │
│ └────────────────────────────────────────┘ │
│ │
│ ┌────────── 外部纯函数微内核 (loop.py) ───┐ │
│ │ run_agent_loop: 彻底接管两阶段调度与双层循环 │
│ │ _provider_context: 前置剥离空异常 + 自愈断头 │
│ └─────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘

对应代码:

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
41
42
43
44
45
46
47
48
49
50
51
52
class Agent:
def __init__(
self,
*,
llm: LLM,
tools: list[Tool],
session: Session,
system_prompt: str | None = None,
max_iterations: int | None = None,
context_budget: int | None = None,
hooks: list[tuple[type[Event], Callable]] | None = None,
...
): ...

# 1. 现代化原生事件流生成器(面向 TUI / CLI 打字机 / 前端)
async def prompt_stream(self, user_input: str) -> AsyncIterator[Event]:
"""将用户输入送入微内核,实时产生事件流,同时驱动 session 落盘。"""
# 前置决策点 UserInput / AgentStart ...
async for event in run_agent_loop(
llm=self.llm,
model=self.model,
system=system_prompt,
messages=self.messages,
tools=self.registry,
context_manager=self._ctx,
get_steering_messages=self._get_steering_messages,
get_follow_up_messages=self._get_follow_up_messages,
before_model_call=self.hooks.emit,
before_tool_call=self.hooks.emit,
after_tool_call=self.hooks.emit,
):
# 同步落盘与外部观察者广播
yield event

# 2. 经典便利门面(100% 向后兼容老脚本与测试)
async def run(self, user_input: str) -> str | None:
"""运行并返回最终文本(内部消费 prompt_stream)。"""
final_text = None
async for event in self.prompt_stream(user_input):
if isinstance(event, AgentEnd):
final_text = event.final_text
return final_text

# 3. 观察者订阅 API
def subscribe(self, listener: Callable[[Event], None]) -> Callable[[], None]:
"""订阅事件流(返回注销句柄)。"""
self._subscribers.append(listener)
return lambda: self._subscribers.remove(listener)

def abort(self) -> None: ... # 协作式取消当前任务并自愈落盘
def steer(self, message: str) -> None: ... # 即时转向
def follow_up(self, message: str) -> None: ... # 排队追加

三、调度总装枢纽:prompt_stream 全流程剖析与 run() 职责解构

在现代 Agent 架构中,对外暴露的驱动接口形成了清晰的分工: - async def prompt_stream(...):系统最高阶的一等公民事件流生成器(Producer & Orchestrator); - async def run(...):仅 10 行代码的阻塞式便利消费门面(Pure Consumer)

1. prompt_stream 的六阶段全生命周期流水线

当外部调用 async for event in agent.prompt_stream("帮我重构代码"): 时,它严格按以下 6 大阶段推进:

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
用户输入 user_input


【阶段 1: 状态复位与扩展懒加载】
└─ 复位 _aborted = False,生成新的 CancellationToken 取消令牌
└─ 若插件/扩展尚未加载,触发 self.extension_manager.load()


【阶段 2: 用户输入安全门禁 (Hook 1: UserInputHook)】
└─ 审查输入是否违规或需脱敏:
├─ block=True ──> 发射 AgentEnd(stop_reason="blocked") 并直接 return 退出!
└─ 改写输入 ────> 更新 user_input 为脱敏后的新文本


【阶段 3: 会话树同步与转录本拓扑自愈】
└─ 读取当前 Session 路径历史:restored = system + get_full_history_messages()
└─ 运行 repair_tool_history(restored) 瞬间缝合所有断头调用,免疫 API 400


【阶段 4: 启动提示词门禁 (Hook 2: AgentStartHook)】
└─ 审查/改写 System Prompt (允许插件动态植入全局规则)


【阶段 5: 微内核点火启动 (run_agent_loop)】
└─ 将准备好的参数喂给 run_agent_loop,获得异步生成器 loop_gen


【阶段 6: 三向分发大循环 (async for event in loop_gen)】
├─ 1. 磁盘持久化分支 ──> 监听到 MessageEnd ➔ self.session.add_message 落盘
├─ 2. 旁路广播分支 ──> 遍历 self._subscribers 广播给第三方监听器
└─ 3. 主干流式分支 ──> yield event 实时推给外部调用方

逐阶段技术细节解密:

  1. 阶段 1(状态复位与扩展懒加载)
    为本次调用创建专属的 CancellationToken。采用懒加载策略:在第一次真正发起问答时才执行 await self.extension_manager.load(),保证 Agent 初始化零异步卡顿;
  2. 阶段 2(用户输入门禁 UserInputHook
    防守前置,零成本拦截。在用户输入触碰模型或写入历史前先行安检。若插件判定违规(block=True),直接就地发射 AgentEnd(stop_reason="blocked") 退出,不花一分钱模型费用;若需要脱敏(updated_input),则无缝篡改后放行;
  3. 阶段 3(会话树指针同步与转录本自愈)
    • 解决 /rewind(会话回滚)后的指针同步问题,严格以磁盘 Session 树当前指针为基准重新拉齐内存历史;
    • 运行 repair_tool_history,即使上次任务中途被强制杀死或崩溃,残留的断头工具调用也能被毫秒级修复,绝对杜绝大模型 API 报 400 Bad Request
  4. 阶段 4(系统提示词门禁 AgentStartHook
    允许策略插件根据当前用户身份或环境,动态向 system_prompt 中追加安全规约;
  5. 阶段 5(微内核点火启动)
    将组装完毕的消息历史、工具注册表、取消信号和队列回调打包,调用 run_agent_loop 启动纯函数引擎;
  6. 阶段 6(三向分发大循环 Tee Splitter)
    微内核吐出的每一个事件,在此处被“一鱼三吃”
    • 向磁盘写:一旦捕获到 MessageEnd(说明消息已完整终态定型),立刻调用 self.session.add_message 通过原子重命名安全落盘到 .jsonl
    • 向广播发:遍历 self._subscribers 列表,安全通知所有旁路监听器(日志、监控、状态栏);
    • 向调用方推:通过 yield event 原生推向终端屏幕或 WebSocket。

2. prompt_stream 与 run() 的职责解构

很多传统框架的代码之所以臃肿,是因为把“流式输出”和“最终结果返回”混在了同一个方法里。而在现代微内核架构中,两者被进行了极其清晰的“生产者 vs 消费者”职责解构:

维度 prompt_stream(...) run(...)
返回类型 AsyncIterator[Event](异步生成器事件流) str | None(最终纯文本答案)
角色定位 调度总装者与事件生产者(Orchestrator & Producer) 端到端纯消费者(Pure Consumer)
内部职责 负责安检门禁、会话同步、驱动微内核、写磁盘、旁路广播 完全不关心底层调度,只管迭代 prompt_stream 等到终态
面向场景 终端 TUI 彩色打字机、WebSocket 实时推送、长程任务进度监听 离线单测(0.001 秒断言输出)、自动化脚本批处理、极简命令行调用
异常体现 将错误包装为结构化 StreamErrorEventAgentEnd(stop_reason="error") 发射 遇到 stop_reason == "error" 时主动向上抛出 RuntimeError,遇到取消返回 "(cancelled)"

极简消费门面:run() 的优雅退化

正是因为 prompt_stream 把所有的脏活累活(安全拦截、自愈、驱动微内核、写文件)全部在底层包揽了,上层的便利方法 run() 才能做到前所未有的纯粹与轻薄——它退化成了仅仅 10 行代码的普通事件消费者

1
2
3
4
5
6
7
8
9
10
11
12
13
async def run(self, user_input: str) -> str | None:
"""追加 user 消息 → 内部消费 prompt_stream 事件流 → 返回最终文本。"""
final_text = None
async for event in self.prompt_stream(user_input):
if isinstance(event, AgentEnd):
if event.stop_reason == "cancelled":
return "(cancelled)"
if event.stop_reason == "blocked":
return event.final_text or "(blocked)"
if event.stop_reason == "error":
raise RuntimeError(event.final_text or "Error during model stream")
final_text = event.final_text
return final_text

这种职责解构的收益极为巨大: 1. 单一事实来源(Single Source of Truth):整个框架只有 prompt_stream 一条真正的执行流水线,任何安全补丁、持久化优化或时序更新只需改动一处,run() 自动享受到所有红利; 2. 兼得双重体验:既能满足现代交互界面对毫秒级流式响应的严苛要求,又能保留单测和简单脚本中 ans = await agent.run(...) 一行调用的极致清爽!

3. 只读事件订阅管道:subscribe 与 _notify 广播器

在生产级 Agent 中,UI 状态更新、日志埋点、监控探针等外部系统需要实时感知 Agent 状态,但绝不能侵入或打断核心推理流程Agent 类实现了对标 Pi(agent-session.js)与 Tau(harness.py)的只读轻量观察者管道:

1
2
3
4
5
# 订阅生命周期事件流(返回注销句柄)
unsubscribe = agent.subscribe(my_ui_listener)

# 业务注销(消除内存泄漏隐患)
unsubscribe()

_notify 广播器的六大工程保证:

1
2
3
4
5
6
7
async def _notify(self, event: Event) -> None:
"""将生命周期事件安全广播给所有旁路订阅者(对标 Tau AgentHarness._notify)。"""
for sub in list(self._subscribers):
with contextlib.suppress(Exception):
res = sub(event)
if inspect.isawaitable(res):
await res
  1. 快照遍历(Snapshot Iteration)list(self._subscribers) 保证在迭代过程中即使外部调用了 unsubscribe() 也绝不会引发并发修改异常;
  2. 异常完全隔离(Never-Throw Guarantee):外部监听器抛出的任何异常被 contextlib.suppress(Exception) 彻底吞默,第三方的 bug 绝对无法击垮 Agent 主循环;
  3. 同异步无缝自适应:自动检查 inspect.isawaitable(res),无论监听器写成普通同步函数还是 async def 异步协程,框架均能自适应等待;
  4. 闭包注销句柄(Teardown Token):解耦索引,注销时按对象引用精准剔除;
  5. 全量关键触点覆盖:在用户输入拦截阻断、系统提示词阻断、微内核主事件迭代循环、以及上下文压缩回写等所有生命周期转折点统一触发广播。

核心架构演进:为什么要把 run_agent_loop 抽离为独立纯函数?

在阶段 18 对标 Tau 架构的重塑中,我们做出了一个决定性的重构动作:将原先内联在 Agent 类中长达 260 行的双层循环彻底剥离,提炼为 packages/my-agent-core/src/my_agent_core/loop.py 中的独立纯函数生成器 run_agent_loop

这一解耦设计带来了两大根本性的架构质变:

1. 纯函数无状态微内核的威力(五大工程质变)

传统框架喜欢把循环写死在类方法内部(class Agent: def run(self): while True:),导致循环与类状态深度咬合。在长对话或复杂场景下,这种设计弊端尽显。而抽离为独立的 async def run_agent_loop(...) 后,系统获得了五大不可替代的红利:

  1. 彻底消除隐式状态突变(零 self,绝对确定性)
    • 在类方法中,循环随时可能修改 self.xxx,多轮对话后状态满天飞;
    • 纯函数微内核完全没有 self,零隐式状态!所有输入项(大模型门面 llm、消息列表 messages、工具注册表 tools、取消令牌 signal)全部通过参数显式传入。输入确定则输出事件流绝对确定,排查调试的确定性拉满。
  2. 通用计算发动机:多场景极致复用(Write Once, Run Everywhere)
    • run_agent_loop 就像一台通用的汽车发动机,不绑定任何特定底盘,能够随处无缝挂载:
      • 交互式 CLI:配上 Agent 外壳,驱动 Session 树持久化与终端彩色打字机;
      • Web / API 流式服务:直接挂载到 FastAPI 的 StreamingResponse 或 WebSocket,中间事件流无需中间人直接推送到前端;
      • 大规模无头评测(Headless Benchmarks):跑 SWE-bench 评测 1000 道题时,无需任何 Session 落盘与复杂插件装配,直接拿微内核裸跑,内存开销降低 80%,极速并发;
      • Subagent 子代理:轻量级临时执行手脚任务,无需背负庞大的 Agent 实例。
  3. 单元测试速度提升 100 倍(零装配、纯内存秒跑)
    • 以前测试循环逻辑,必须先装配一整台 Agent 大机器(Session 树、MemoryStore、Skills、Plugins),单测准备代码冗长且慢;
    • 现在在 tests/test_agent_loop_pure.py 中,直接传入 FakeLLM 和列表,0.001 秒测完并发工具与动态转向(11 项复杂微内核单测在 0.05 秒内全通)。
  4. 动静分离,职责正交(SRP 原则)
    • agent.py(类)只负责静态的结构与持久化(持有 Session 树、长期记忆、扩展装配);
    • loop.py(函数)只负责动态的过程推理与调度(怎么推理、怎么并发调工具、怎么转向、怎么中断)。
  5. 双层循环拓扑与动态干预毫秒级响应
    • 内层微观循环在工具执行间隙动态收割 get_steering_messages() 实现即时转向;
    • 外层宏观循环在轮次自然收尾处收割 get_follow_up_messages() 自动衔接追问任务。

2. 将中间“事件流”提升为系统级的一等公民(First-Class Stream)

大模型的运行本质上是一个基于时间轴的渐进流式过程:Token 在一个字一个字蹦出、思考块在逐步形成、工具在按需触发并耗时运行。

传统 Callback / Hook 的三大灾难(打洞插管):

在传统黑盒设计中,函数只能最后 return final_text。外部若想知道中间状态,必须在外面定义一堆回调函数(on_tokenon_tool_start)硬塞给 Agent。 - 控制权倒置:调用方不能主动拉取数据,只能被框架被动调用; - 异常黑洞:回调函数一旦抛出异常,会直接击垮主循环,或者被框架静默吞掉; - 无法优雅暂停与取消:外部用户在前端点了停止,回调依然在疯狂发射,根本停不下来。

生成器(AsyncIterator[Event])的一等公民之美:

在现代 run_agent_loop 架构下,事件流就是函数的正常返回值本身!没有任何侧路管子,没有控制权倒置!

  1. 调用方掌握 100% 绝对主动权(传送带模型): 外部调用者只需写出自然优雅的 async for 循环,站在传送带旁边接住事件:
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    # 外部使用代码极度自然、主动:
    async for event in agent.prompt_stream("帮我计算 37 * 19"):
    if isinstance(event, MessageUpdate):
    print(event.chunk.delta, end="", flush=True) # 原生打字机!
    elif isinstance(event, ToolExecutionStart):
    print(f"\n[正在调用工具: {event.tool_name}]")
    elif isinstance(event, ToolExecutionEnd):
    print(f"\n[工具执行完毕: {event.result}]")

    # 想停就停:直接 break 即可触发生成器安全清理与底层中断自愈!
    if user_clicked_stop:
    break
  2. 让上层门面实现前所未有的轻薄: 因为事件流成了一等公民,外层的 agent.run() 彻底被掏空,退化成了仅 10 行代码的纯事件消费者:
    1
    2
    3
    4
    5
    6
    7
    async def run(self, user_input: str) -> str | None:
    """自己只做普通消费者,消费事件流并返回最终文本。"""
    final_text = None
    async for event in self.prompt_stream(user_input):
    if isinstance(event, AgentEnd):
    final_text = event.final_text
    return final_text
    一套微内核底座,同时完美满足了“前端极致细粒度的响应式流渲染”与“后端一行代码调用的极简体验”!

3. 架构深思:Pi 的 ExtensionAPI 如何将“只读事件”与“拦截决策点”相辅相成?

在深入对比 Pi(@earendil-works/pi-coding-agent)与 Tau 的设计后,我们终于看清了困扰很多开发者的核心谜团:为什么在很多初学者的 Agent 代码里,事件和 Hook 会被写成一团乱麻?而 Pi 却能写得如此规整有序?

答案就在于:Pi 深刻认清了“只读事件(传信息)”与“决策拦截点(改控制流)”是两套完全正交的体系,并通过 ExtensionAPI 的三层架构实现了二者的完美相辅相成!

① 概念正交化:纯只读事件 vs 决策拦截点

体系类别 代表事件名称 核心职责 返回值意义 典型消费场景
纯只读事件流
(Agent Events)
turn_start, turn_end, message_update, tool_execution_end 客观传递信息(Stream of Facts),单向广播发生了什么 被动忽略 (void),没有资格改变控制流 驱动终端 TUI 彩色打字机、状态栏、Web 前端进度条、日志审计
五大决策拦截点
(Decision Points)
input, before_agent_start, context, tool_call, tool_result 主观裁决与改写(Control Gates),把守关键命脉卡口 决定生死!返回 { block: true } 阻断或篡改参数/上下文 拦截 rm -rf 高危命令、人机审批流、动态注入提示词、清洗脱敏敏感词

② 为什么初学者(包括我们重构初期)会陷入杂乱?

  • 概念混淆的陷阱:我们把所有的只读事件和拦截点硬塞进了同一个 HookRegistry,导致微内核里充斥着几十处毫无意义的 if hook_registry: await hook_registry.emit(ev)
  • 双通道视觉噪音:每个事件既要 yield 给外部,又要 await emit 给内部。这就好比广播电台播报“天亮了”(纯通知),却非要在广播室门口设个安检站一本正经地问安检员“天亮了需要拦截吗?”。

③ Pi 的三层解耦美学:ExtensionAPI 与微内核的精密咬合

Pi 极为高明地通过三层架构化解了这一矛盾:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
┌─────────────────────────────────────────────────────────────┐
│ 【外部插件层 (Extension Developer)】 │
│ 面向开发者提供极简统一入口:pi.on(event, handler) │
│ • 监听只读事件:pi.on("message_update", handler) │
│ • 拦截安全门禁:pi.on("tool_call", handler) │
└──────────────────────────────┬──────────────────────────────┘
│ 智能分流

┌─────────────────────────────────────────────────────────────┐
│ 【中间调度层 (ExtensionManager)】 │
│ • 只读事件 ➔ 接入单向广播观察者队列 (Observer Queue) │
│ • 决策拦截点 ➔ 串联为洋葱模型中间件流水线 (Middleware) │
└──────────────────────────────┬──────────────────────────────┘
│ 注入纯函数插槽

┌─────────────────────────────────────────────────────────────┐
│ 【底层微内核 (runAgentLoop)】 │
│ • 微内核对 ExtensionAPI 零感知!不持有全局 Hook 注册表 │
│ • 纯只读事件:微内核只管一行代码 yield event (纯单向推流) │
│ • 五大决策点:微内核仅通过入参暴露 5 个具体干净的回调插槽 │
│ (beforeToolCall, afterToolCall, transformContext...) │
└─────────────────────────────────────────────────────────────┘

这种相辅相成的设计妙在哪里?

  1. 对插件作者(API 面)极致的直觉与统一。开发者不需要记两套完全不同的 API,只要在同一个 pi.on 后面写函数,TypeScript 类型系统会自动约束当前事件“能不能拦截、可以返回什么”;
  2. 对中间层(调度面)职责清晰的智能分流。只读事件走广播,安全门禁走洋葱流水线(支持多插件按序改写参数与一票短路阻断);
  3. 对微内核(底层引擎)极致的轻盈与规整。微内核根本不需要感知插件的存在,也不需要充斥 30 处 Hook 检查样板代码;它只在 5 个固定的命脉节点上留出回调插槽,其余时间顺风顺水地跑自己的纯状态机!

这正是顶级工业级 Agent 框架的架构魅力所在:上层 API 极其人性统一,底层微内核极其纯粹高效,中间通过清晰的分流管道相辅相成!

4. 旁路广播总线:subscribe() 订阅机制的底层本质与实战全景

很多开发者在接触到 agent.prompt_stream(...) 后会产生一个经典困惑: > “既然我已经能通过 async for event in agent.prompt_stream(...) 直接拿到所有事件了,为什么还需要 agent.subscribe(...) 订阅机制?这不是多此一举吗?”

答案就在于:prompt_stream 是单通道的“独食吸管”,而 subscribe 是多方共享的“公共大喇叭”。

① 为什么必须要有 subscribe?(两大刚需场景)

  1. 配合最常用的极简门面 await agent.run(): 大部分日常业务或测试中,我们只想要最终答案:ans = await agent.run("计算 37*19")。 然而,run() 内部已经把 prompt_stream消费吞掉了,外部只能拿到最终字符串。如果你此时既想用清爽的 agent.run(),又想在控制台打印工具调用日志,或者在状态栏转圈,唯一的办法就是提前挂载 agent.subscribe(my_listener)
  2. 克服 Python 生成器的“单消费者限制”: Python 的异步生成器不能被多个 async for 同时读取。如果你的系统需要同时做三件事:
    • 主流程:在终端打印大模型的话;
    • 日志系统:把所有工具调用记录进 agent.log
    • 监控系统:把消耗的 Token 上报云端; 若没有 subscribe,你必须把日志和监控代码全揉进主流程的 async for 里,代码变成大泥球。有了 subscribe,日志和监控模块只需各自独立 agent.subscribe(...)旁路静默监听,与主业务彻底解耦

② agent.subscribe(…) 的入参到底是什么?

入参是一个“函数”(由你编写的回调函数 Callable)。
该函数必须接收一个参数(框架会把发生的 event 实体传给它)。支持三种形式:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 形式 1: 普通函数(最直观常用)
def my_listener(event: Event):
print("收到事件:", event)

agent.subscribe(my_listener) # 传入函数名,切勿带括号!

# 形式 2: 异步函数(适合写数据库或发网络请求)
async def my_async_listener(event: Event):
await websocket.send_json({"event": type(event).__name__})

agent.subscribe(my_async_listener)

# 形式 3: 匿名函数(适合一行调试)
agent.subscribe(lambda ev: logger.info("Event: %s", ev))

③ 订阅只能看到事件类型吗?(绝对不是!)

每个推送到订阅函数的 event,都是一个满载真实业务数据的强类型实体对象,你能从中拿到所有的具体数据:

  • 工具启动(ToolExecutionStart
    • event.tool_name:工具名(如 "bash""read_file");
    • event.args大模型传入的真实入参字典(如 {"path": "src/app.py"});
    • event.tool_call_id:工具调用唯一 ID。
  • 工具结束(ToolExecutionEnd
    • event.result工具执行的完整输出内容(如文件全文、命令执行结果);
    • event.is_error:执行是否出错(True / False);
    • event.terminate:是否触发了阶段 7 提前熔断标记。
  • 大模型流式吐字(MessageUpdate
    • event.delta刚刚吐出的那一个字/词(驱动终端流式打字机)。
  • 任务终结(AgentEnd
    • event.final_text:最终交付给用户的完整回答文本;
    • event.iterations:总共跑了多少轮推理;
    • event.stop_reason:为何结束("end_turn""max_iterations" 等)。
  • 上下文压缩(ContextCompacted
    • event.tokens_beforeevent.tokens_after:压缩前后的真实 Token 水位。

实战中,你可以轻松写出一个赏心悦目的终端监听器:

1
2
3
4
5
6
7
8
9
10
11
def pretty_printer(event: Event):
if isinstance(event, ToolExecutionStart):
print(f"🔧 [工具启动] {event.tool_name},参数: {event.args}")
elif isinstance(event, ToolExecutionEnd):
status = "❌ 失败" if event.is_error else "✅ 成功"
print(f"{status} [工具完成] 结果: {event.result}")
elif isinstance(event, MessageUpdate):
print(event.delta, end="", flush=True)

agent.subscribe(pretty_printer)
await agent.run("计算 37 * 19")

④ 底层实现原理(三部曲机制)

subscribe 在底层代码中仅有约 15 行,设计极其紧凑:

1
2
3
4
5
6
7
8
1. 登记通讯录 (subscribe)
agent.subscribe(listener) ──> self._subscribers.append(listener)

2. 逐个打电话广播 (prompt_stream 内部)
微内核产生新事件 ──> for sub in list(self._subscribers): sub(event)

3. 随时退订除名 (unsubscribe 闭包)
调用返回的 unsub() ──> self._subscribers.remove(listener)
  1. 闭包式退订令牌(Teardown Token)
    subscribe 返回一个私有闭包 unsubscribe()。调用方无需持有监听器列表或索引,直接调用 unsub() 即可瞬间安全退订,从机制上彻底杜绝内存泄漏
  2. list(self._subscribers) 浅拷贝快照
    在广播循环中使用快照遍历,防止监听器在执行过程中自行动态注销导致 RuntimeError: list changed size during iteration 崩溃;
  3. contextlib.suppress(Exception) 异常隔离
    第三方订阅函数内部若抛出任何未捕获异常,全部被静默隔离,绝对不影响 Agent 主调度循环的稳定性
  4. 同步/异步自适应(isawaitable
    底层通过 inspect.isawaitable(res) 自动检测返回值,若订阅函数是 async def,主线程自动进行 await 调度。

5. 微内核调度流水线独立专文导航

至此,Agent 类已经蜕变成一个高度轻量、专注会话状态管理与全系统依赖注入的 Harness 宿主外壳。

关于微内核内部的双层状态机循环(宏观 Follow-up 追问 + 微观 ReAct 转向)、单轮 8 步时序流水线、大模型推理车间(_assistant_turn)以及工业级工具执行车间(_execute_tools_turn)的深度解析,已独立拆分并收录于专属博客中:

👉 《my-pi-agent–loop微内核.md》