Pi Agent · Book
冬瓜 热衷于拆解 AI 工程的博主
M07

第7章:事件驱动 —— Agent 的神经系统

5829字 · 含 306 行代码 · 约 30 分钟
Python 转写 · 原作 TypeScript

前六章里,有一个东西反复出现但我们始终没深入——事件

第 3 章说”Agent Loop 每做一步都发事件让 UI 实时更新”。第 5 章说”工具执行时发出 tool_execution_starttool_execution_updatetool_execution_end 事件”。第 6 章里事件到处携带 AgentMessage

但我们始终没回答:事件到底是怎么从 Agent 内部传到外部的?谁在监听?为什么 Agent 发完事件后对一类监听器”等”,对另一类”不等”?

这一章就打开 Agent 的”神经系统”。

本章最重要的一句话:Pi 有两套并行的监听机制——session.subscribe(只读观察,Agent 不等你)和扩展系统的 pi.on(能拦截、能改写,Agent 会等你)。 它们共享同一批事件源,但”Agent 等不等你的 listener”是两者最根本的分水岭。如果你只学一套,一定会踩”代码写了却静默不生效”的坑。

本章起为进阶章节。前六章建立了对 Pi-Agent 运行机制的整体理解,从这里开始深入工程化议题。

关于本章的代码:Pi 的参考实现是 TypeScript(repo/ 仓库)。本章所有源码剖析针对 TS 源码;Python 版块的代码是概念改写——把 TS 的写法翻译成 Python 等价表达,帮助 Python 开发者理解。Python SDK 的真实 API 命名请以官方文档为准。


一、为什么需要事件系统?

一个直觉:从外卖追踪说起

你在美团上点了一份外卖。下单后,App 会给你推送一连串状态更新:“商家已接单” → “骑手已取餐” → “骑手距你 500 米” → “已送达”。每一个状态更新就是一个事件——它告诉你”发生了一件事”。你不需要一直盯着骑手的位置看,只需要在收到事件时看一眼。

Pi-Agent 的事件就是这个意思:Agent 运行过程中不断产生”发生了某事”的快照——消息开始了、消息更新了、工具开始执行了——然后把这些快照推给所有关心它的人。

不用事件会怎样?

假设你要给 Agent 加一个”工具调用日志”功能:每次调工具时打印一行 [LOG] 调用了 read,参数:main.ts

不用事件系统:你得改 Agent 源码,在 tool.execute() 前后各加一行 console.log。然后 Pi 更新了,你 merge 上游代码时发现冲突——你加的日志和上游新增的逻辑撞在一起了。手动解决冲突,下周又更新,又冲突……

用事件系统

# ============================================================
# 【Python 改写】订阅事件做工具调用日志
# 原文 TS:
#   session.subscribe((event) => {
#       if (event.type === "tool_execution_end") {
#           console.log(`[LOG] 调用了 ${event.toolName},结果:${event.isError ? "失败" : "成功"}`);
#       }
#   });
# ============================================================

# 概念对照:TS 的箭头函数 → Python 的 lambda 或 def;
# TS 的三元运算符 → Python 的条件表达式

def listener(event):
    if event.type == "tool_execution_end":
        result_str = "失败" if event.is_error else "成功"
        print(f"[LOG] 调用了 {event.tool_name},结果:{result_str}")

session.subscribe(listener)

六行代码。不碰 Agent 一行源码。Agent 更新你只需要更新包,日志逻辑不受影响。

这就是事件驱动最核心的价值:把”发生了什么”和”谁关心什么”彻底分离。 Agent 只管发事件,它不知道也不关心谁在听。

发布-订阅 vs 直接调用

用编程术语说,事件驱动实现的是发布-订阅模式。和直接函数调用做个对比:

直接调用(打电话):
  Agent ──调用──→ 终端渲染
       ──调用──→ 文件存储
       ──调用──→ 日志记录
  Agent 需要知道所有消费者的存在,每加一个新功能就要改 Agent

发布-订阅(广播):
  Agent ──emit事件──→ 📡 事件总线
                          ├──→ 终端渲染(订阅了)
                          ├──→ 文件存储(订阅了)
                          ├──→ 日志记录(订阅了)
                          └──→ (新功能只需订阅,Agent 不需要知道)

一句话:直接调用是”我亲自找你”;发布-订阅是”我对着空气喊了一声,谁听到算谁的”。 Pi 的事件系统就是发布-订阅——Agent 发出事件,关心它的人各自订阅、各自处理。但 Pi 有一个关键特点:它有两条订阅管道,而且两条的能力很不一样。这正是本章要讲的核心,下一节就铺开。


二、两条管道的全貌:从事件源到两类监听器

这一节把 Pi 事件体系的全景铺开——2.1 看事件源,2.2 认识两条管道并讲清它们的核心差别。后面第三、四节会分别深入两条管道的用法和源码。

2.1 事件源:10 种 AgentEvent

Agent 内核层定义了 10 种 AgentEvent,它们构成了 Agent 运行的完整”脉搏”:

10 种事件 4 层嵌套
10 种事件 4 层嵌套

配图说明:从外到内 4 层嵌套——Agent(Trace)→ Turn → Message → Tool Execution。每层都是”开始 → 更新(×N)→ 结束”配对。注意 Turn 2 没有 ToolCall 所以没有 Layer 4 嵌套。底部图例标注每层的事件数(2+2+3+3=10 种)。

# ============================================================
# 【Python 改写】AgentEvent 联合类型(10 种)
# 原文 TS:
#   export type AgentEvent =
#     | { type: "agent_start" }
#     | { type: "agent_end"; messages: AgentMessage[] }
#     | { type: "turn_start" }
#     | { type: "turn_end"; message: AgentMessage; toolResults: ToolResultMessage[] }
#     | { type: "message_start"; message: AgentMessage }
#     | { type: "message_update"; message: AgentMessage; assistantMessageEvent: AssistantMessageEvent }
#     | { type: "message_end"; message: AgentMessage }
#     | { type: "tool_execution_start"; toolCallId: string; toolName: string; args: any }
#     | { type: "tool_execution_update"; toolCallId: string; toolName: string; args: any; partialResult: any }
#     | { type: "tool_execution_end"; toolCallId: string; toolName: string; result: any; isError: boolean };
# ============================================================

# 概念对照:TS 的 tagged union `| { type: "X", ... }` → Python 用 dataclass + 共同 type 字段,
# 或 pydantic 的 discriminated union。这里用 dict + Literal["type"] 简化表达。

# 4 层嵌套 10 种事件:
AgentEvent = Union[
    # 第 1 层:Agent 生命周期(整个运行)
    {"type": Literal["agent_start"]},
    {"type": Literal["agent_end"],       "messages": list[AgentMessage]},

    # 第 2 层:Turn 生命周期(一轮模型调用 + 工具执行)
    {"type": Literal["turn_start"]},
    {"type": Literal["turn_end"],        "message": AgentMessage, "tool_results": list[ToolResultMessage]},

    # 第 3 层:Message 生命周期(一条消息)
    {"type": Literal["message_start"],   "message": AgentMessage},
    {"type": Literal["message_update"],  "message": AgentMessage, "assistant_message_event": AssistantMessageEvent},
    {"type": Literal["message_end"],     "message": AgentMessage},

    # 第 4 层:Tool Execution 生命周期(一次工具执行)
    {"type": Literal["tool_execution_start"],  "tool_call_id": str, "tool_name": str, "args": Any},
    {"type": Literal["tool_execution_update"], "tool_call_id": str, "tool_name": str, "args": Any, "partial_result": Any},
    {"type": Literal["tool_execution_end"],    "tool_call_id": str, "tool_name": str, "result": Any, "is_error": bool},
]

10 种看着不少,但规律很清楚——它们是 4 层嵌套的生命周期,每层都有”开始→更新→结束”的配对:

Agent 运行
├── agent_start ───────────────────── Agent 开始

├── Turn 1(第3章讲过:一次模型调用 + 它触发的工具执行)
│   ├── turn_start ────────────────── Turn 开始
│   │
│   ├── Message(LLM 的响应)
│   │   ├── message_start
│   │   ├── message_update ×N ────── 流式增量(逐 token 更新)
│   │   └── message_end
│   │
│   ├── Tool Execution(工具执行)
│   │   ├── tool_execution_start
│   │   ├── tool_execution_update ×N  工具进度(如 Bash 的输出)
│   │   └── tool_execution_end
│   │
│   └── turn_end ──────────────────── Turn 结束

├── Turn 2 ...

└── agent_end ──────────────────────── Agent 结束

回忆第 3 章的概念:一个 Turn = 一次模型调用 + 这次调用触发的所有工具执行。 turn_start 到 turn_end 之间,模型被调用了恰好一次。

为什么要 4 层? 因为不同消费者关心不同粒度。TUI(终端界面)需要逐 token 渲染文字,所以它订阅 message_update;而 Session 管理器只关心一轮对话结束了没有,所以它只看 turn_end。4 层嵌套让每种消费者都能在刚刚好的粒度上响应。

这 10 种是事件源,是两条管道的共同水源。 两类监听器都消费这同一批事件——管道 A 直接收,管道 B 收翻译版。所以这 10 种不属于任何一条管道,是它们共享的源头。两条管道具体是什么、差别在哪,下一节展开。

2.2 两条管道:subscribe 与扩展 pi.on

Pi 的事件有两条订阅管道,它们喂的是同一个事件源,但能力差很多。这一节先把两条管道介绍清楚,再讲它们的核心差别。后面第三、四节会分别深入。

管道 A:session.subscribe

你在外部脚本里注册监听器(Web 服务器、CLI 工具)。拿到 session 对象后调 session.subscribe(listener),事件就会流进你的 listener。

它的特点是只能看,不能改——listener 没有返回值(或者说返回了也被丢弃),Agent 不会因为你的监听器改变任何行为。典型用途:流式渲染、打日志、把事件转发给浏览器。

管道 B:扩展系统的 pi.on

你把代码写进一个扩展(一段被框架加载的插件),在里面调 pi.on("事件名", handler)。handler 带返回值,Agent 会读。

它的特点是能改 Agent 的行为——handler 可以返回 {"block": True} 拦掉一次工具调用,可以返回新的消息列表改写给 LLM 的上下文。典型用途:安全策略、审计、注入动态信息。

核心差别:Agent 对 B 等,对 A 不等

两条管道最关键的差别不在返回值,而在 Agent 等不等你:

  • 管道 B 的 handler,Agent 会等它返回。因为 Agent 要读返回值才能决定下一步——你说拦,它才拦。源码里这行带 await
await self._emit_extension_event(event)   # 扩展:等。Agent 要读 handler 的返回值
  • 管道 A 的 listener,Agent 不等。通知完就继续,listener 返回什么 Agent 都不读。源码里这行不带 await,而且 _emit 本身就是个同步函数:
self._emit(event)                          # subscribe:不等。同步调用,返回值丢弃

因果关系很直接:扩展要读返回值,所以必须等;subscribe 不读返回值,等了也没用。“能改 Agent 行为”是”等 + 读返回值”的结果,不是单独赋予的能力。

两条管道——subscribe 与扩展 pi.on
两条管道——subscribe 与扩展 pi.on

配图说明:同一个事件源分流到两条管道——左管道 A(subscribe,只读广播,Agent 不等),右管道 B(扩展 pi.on,能拦截改写,Agent 等你回话)。两条管道喂的是同一个事件,但 Agent 对一个等、对另一个不等。

容易搞混的两个事件名:tool_calltool_execution_start

这两个名字像,但走不同管道、发生在不同时刻:

LLM 决定调一个工具

   ▼  管道 B:tool_call(执行前)
   │   扩展可以 return {"block": True} 拦掉它;一旦 block,下面都不发生
   │   注意:tool_call 不是 2.1 那 10 种之一,是 SDK 在执行前主动触发的扩展独占事件

   ▼  (没被拦)tool.execute() 开跑

   ▼  管道 A + B 都收到:tool_execution_start(已开跑,拦不住了)
   │   这是 2.1 那 10 种之一,由 Agent 内核发出,两条管道都收

   ▼  tool_execution_end

tool_call 是执行前的安检门(能拦,管道 B 独占),tool_execution_start 是开跑后的广播(拦不了,两条管道都收)。管道 A 根本收不到 tool_call——你在 subscribe 里写 if event["type"] == "tool_call" 不会报错,但这个分支永远命中不了。

除了 tool_call,扩展还独占另外 4 个决策点(inputbefore_agent_startcontexttool_result),它们都是”Agent 要停下来读返回值”的位置,管道 A 一律收不到。这 5 个事件加起来,构成了管道 B 能干预 Agent 的全部入口。

一个细节:工具进度更新可以不等。 生命周期事件(start/end 这类低频、不能错的)Agent 会逐个等扩展处理完。但工具执行时会刷出大量进度(Bash 每一行输出都是一个 tool_execution_update),逐个等会卡住 Agent。Pi 对这类高频事件开了口子——先攒着,最后一次性等完(用 asyncio.gather 批量收口)。原则是:越重要的事件等得越严格。


三、管道 A:session.subscribe

2.2 介绍了管道 A 的特点:Agent 不等它,所以它只能看、不能改。这一节落到源码,看清”不等”和”只能看”是怎么实现的。

3.1 怎么用:注册、签名、注销

# ============================================================
# 【Python 改写】订阅事件:注册、签名、注销
# 原文 TS:
#   const unsubscribe = session.subscribe((event, signal) => { ... });
#   unsubscribe();
# ============================================================

# 注册:传入一个 listener,返回一个注销函数
def listener(event, signal):
    if event["type"] == "message_update":
        # 流式打字(assistant_message_event 里带 delta)
        sys.stdout.write(event["assistant_message_event"]["delta"])

unsubscribe = session.subscribe(listener)
# 不用了就调注销函数
unsubscribe()

listener 的签名是关键:

# 概念对照:TS 的 (event, signal) => void → Python 的 (event, signal) -> None
AgentSessionEventListener = Callable[[AgentSessionEvent, Optional[AbortSignal]], None]
#                                                                  ↑↑↑↑
#                                                            注意返回值:None

返回值是 None——就算你在 listener 里 return {"block": True},Agent 也不读、不用。这是管道 A”只能看不能改”在类型层面的体现,也是和管道 B 最根本的差别。

3.2 “不等”对管道 A 意味着什么

Agent 对 subscribe 监听器”通知一声就走”,不等你。落到实战,这个”不等”带来三个直接后果:

  • 你可以在监听器里干异步重活(比如 await 一个慢请求),不会拖慢 Agent——Agent 早就走下一步了,你的请求在后台慢慢跑。
  • 但也正因为不等,你的异步结果传不回去——Agent 已经发下一个事件了,不在乎你算出了什么。
  • 所以管道 A 适合”我慢慢干我的,不打扰 Agent”的场景:写日志、推 SSE、更新外部状态。不适合需要”先等我处理完再继续”的场景——那必须走管道 B。

一个隐藏的坑:因为不等,async 监听器里的错误不会冒泡到 Agent。如果你在 async 监听器里 await 一个会失败的操作,失败会被悄悄吞掉,你连错在哪都不知道。务必在 async 监听器里自己 try-except——Agent 不会替你兜底。

3.3 管道 A 能收到哪些事件

管道 A 收到的是 AgentSessionEvent,它比内核 10 种事件多出几个产品级事件——这些事件反映的不是”Agent 内核跑了一步”,而是”产品层做了某件事”(压缩、重试、排队等)。内核根本不知道这些概念,所以它们只出现在产品层。

大致有个印象就行,遇到具体场景再查:

事件什么时候触发典型用途
agent_settled一次 prompt() 彻底跑完(含重试/压缩/队列全部处理完)可靠的结束信号:写库收尾、推 SSE done
compaction_start / compaction_end上下文窗口快满时,自动压缩历史(第9章详讲)UI 显示”正在压缩…”提示
auto_retry_start / auto_retry_endLLM 调用失败,自动重试UI 显示重试次数、告警
queue_updatesteering / followUp 消息队列变化UI 更新”排队中”状态
session_info_changed会话名称等元信息变化UI 刷新标题
thinking_level_changed切换思考深度UI 联动显示

加上内核的 10 种生命周期事件,就是管道 A 能收到的全部。

但管道 A 收不到管道 B 独占的 5 个决策点(input/before_agent_start/context/tool_call/tool_result)——这 5 个 Agent 内核根本没有,是 SDK 在决策点主动调用扩展系统时产生的,只能走管道 B。

3.4 什么时候用管道 A

一句话:纯观察、不改 Agent 行为、不需要 Agent 等你的场景。 典型用途——流式渲染(TUI 逐 token 打字)、日志记录、SSE 转发给浏览器、统计 token 用量。这些场景的共同点是”Agent 干它的,我在旁边看一眼、或慢慢做我自己的事”,不需要 Agent 配合。

💡 落库 / 审计 / 日志也是「纯观察」,首选管道 A。因为管道 A 不 await 你的监听器——你在里面 await db.insert(),派发方 _emit 调一下就走(2.2 讲的”不读返回值”),I/O 在后台跑,Agent 不被拖慢。

只有当你要存的数据来自 tool_call / tool_result / context管道 A 收不到的 5 个决策点时,才被迫走管道 B——但管道 B 的 handler 被 await(2.2 讲的”Agent 等你”),落库必须 fire-and-forget:handler 里把数据推进队列后立刻返回,真正的写库交给独立 worker 异步处理,别让 await 链绑住 Agent 主循环。

一个易踩的坑:message_end 看着像”一条消息存一笔”的好时机,但它一轮 prompt() 会触发多次(每条 assistant 消息结束都发,含中间要调工具的那些)。拿它当整轮收尾会重复落库 / 重复推 done。整轮的可靠收尾用 agent_settled——它每 prompt 只发一次。


四、管道 B:扩展系统 pi.on

这一节是本章的重头戏。管道 B 的能力远超管道 A——它能拦截工具、改写上下文、替换系统提示词。这些能力的源头,是 SDK 在关键决策点停下来等扩展的返回值(2.2 讲的”Agent 等”)。

Python 版提醒:下面的 @pi.on(...) 装饰器写法是概念改写,对应 TS 的 pi.on("event", handler)。Python SDK 的真实扩展注册 API 请以官方文档为准——可能是装饰器,也可能是 pi.on("event", handler) 的函数调用形式。

4.1 怎么用:写扩展、挂载、注册

管道 B 的代码不在”外部脚本”里,而是写在扩展里——一个接收 pi 对象的工厂函数:

# ============================================================
# 【Python 改写】一个扩展:工厂函数 + pi.on 注册
# 原文 TS:
#   function myGuardExtension(pi) {
#       pi.on("tool_call", async (event, ctx) => {
#           if (event.toolName === "delete_table") {
#               return { block: true, reason: "生产环境禁止删除操作" };
#           }
#           return undefined;
#       });
#   }
# ============================================================

# 概念对照:TS 的扩展工厂函数 → Python 同样是接收 pi 对象的可调用对象;
# 返回 dict 对应 TS 的 return 对象字面量
# 装饰器形式注册 handler(Python SDK 风格,概念改写)

def my_guard_extension(pi):
    @pi.on("tool_call")
    async def handler(event, ctx):
        if event["tool_name"] == "delete_table":
            return {"block": True, "reason": "生产环境禁止删除操作"}   # ← 有返回值,能拦
        return None   # 放行

挂载通过 loader 的 extension_factories(TS 是 DefaultResourceLoaderextensionFactories):

# 概念改写:对应 TS 的 new DefaultResourceLoader({ extensionFactories: [...] })
loader = create_resource_loader(
    cwd=os.getcwd(),
    agent_dir=get_agent_dir(),
    extension_factories=[my_guard_extension],   # ← 你的扩展塞这儿
)
await loader.reload()

session = await create_agent_session(model=model, resource_loader=loader)

注意两个关键点:

  • handler 有返回值return {"block": True})——这是管道 B 能干预 Agent 的根本,也是”Agent 必须等你”的原因(要读返回值)。
  • handler 多一个 ctx 参数——扩展上下文,能力比管道 A 的 event + signal 强(见 4.4)。

4.2 源码:pi.on 只是往 Map(dict)里 push

管道 B 的注册实现极简。pi.on 做的事就是往一个 dict 里塞 handler(TS 源码是 Map<event, handler[]>):

# ============================================================
# 【Python 改写】pi.on 的实现(往 dict 里 append)
# 原文 TS (loader.ts createExtensionAPI):
#   on(event: string, handler: HandlerFn): void {
#       runtime.assertActive();
#       const list = extension.handlers.get(event) ?? [];
#       list.push(handler);
#       extension.handlers.set(event, list);
#   }
# ============================================================

# 概念对照:TS 的 Map<string, HandlerFn[]> → Python 的 dict[str, list[Handler]]

def on(self, event: str, handler: HandlerFn) -> None:
    self._runtime.assert_active()
    lst = self._extension.handlers.get(event, [])   # 取出该事件的 handler 列表
    lst.append(handler)                              # 塞进去
    self._extension.handlers[event] = lst            # 放回

每个扩展对象内部有一个 handlers: dict[事件名, list[handler]]pi.on("tool_call", h) 就是在 "tool_call" 这个 key 下追加 h。派发时遍历这个 dict。

4.3 源码:extension_runner 的两条派发路径

真正的差别在派发。extension_runner 有两类派发方法,对应”通知型”和”决策型”事件——但两类都 await handler(这是管道 B”Agent 等你”的实现):

路径 1:通知型 emit()(runner.ts:796)——处理 message_updateturn_start 等只读事件:

# ============================================================
# 【Python 改写】通知型 emit(串行 await、try-except 隔离、忽略返回值)
# 原文 TS (runner.ts:796-828):
#   for (const ext of this.extensions) {
#       const handlers = ext.handlers.get(event.type);
#       for (const handler of handlers) {
#           try {
#               await handler(event, ctx);   // 通知型:不读返回值
#           } catch (err) { this.emitError({...}); }   # ★ try-catch 隔离
#       }
#   }
# ============================================================

async def emit(self, event) -> None:
    ctx = self.create_context()
    for ext in self.extensions:                       # 遍历扩展
        handlers = ext.handlers.get(event["type"], [])
        for handler in handlers:
            try:
                await handler(event, ctx)             # ★ await 每个 handler
                # 通知型事件:不读返回值(session_before_* 例外,读 cancel)
            except Exception as err:                  # ★ try-except 隔离
                self.emit_error({"extension_path": ext.path, "event": event["type"], "error": str(err)})

特征:串行 await、try-except 隔离、忽略返回值。注意即便忽略返回值,仍然 await——这是为了”同步屏障”(等扩展处理完才发下一个事件,保证状态一致)。这条路径处理的是 2.1 那 10 种内核事件翻译过来的。

路径 2:决策型 emit_tool_call()(runner.ts:927)——处理 tool_call 拦截:

# ============================================================
# 【Python 改写】决策型 emit_tool_call(读返回值、block 短路、★无 try-except)
# 原文 TS (runner.ts:927-948):
#   for (const ext of this.extensions) {
#       const handlers = ext.handlers.get("tool_call");
#       for (const handler of handlers) {
#           const handlerResult = await handler(event, ctx);   // ★ 读返回值
#           if (handlerResult) {
#               result = handlerResult;
#               if (result.block) return result;               // ★ block 短路
#           }
#           // 注意:这里没有 try-catch!
#       }
#   }
# ============================================================

async def emit_tool_call(self, event) -> Optional[ToolCallEventResult]:
    ctx = self.create_context()
    result = None
    for ext in self.extensions:
        handlers = ext.handlers.get("tool_call", [])
        for handler in handlers:
            handler_result = await handler(event, ctx)        # ★ await + 读返回值
            if handler_result is not None:
                result = handler_result
                if result.get("block"):
                    return result                             # ★ block 立即短路
            # 注意:这里没有 try-except!这是刻意的 fail-closed 设计——
            # 扩展在 tool_call 里崩了,宁可拦掉工具也不放行(安全优先)
    return result

特征:await + 读返回值 + block 短路 + 无 try-except。tool_call 是所有派发方法里唯一不包 try-except 的——扩展抛错会冒泡,导致这次工具调用被 block(fail-closed:宁可错杀,不放行可能危险的操作)。

路径 3:链式 transform 型——还有一批决策事件走”链式 transform”,每个 handler 接收上一个的输出继续改。比如:

  • emit_context()(runner.ts:979):第一个 handler 拿到原始 messages,改完传给第二个,最后一个的输出就是真正发给 LLM 的消息列表。返回类型 list[AgentMessage]
  • emit_tool_result()(runner.ts:872):链式修改工具结果,返回 ToolResultEventResult
  • emit_before_agent_start()(runner.ts:1076):链式覆盖系统提示词 + 收集注入消息。
  • emit_input()(runner.ts:1191):链式改写用户输入,action: "handled" 短路。

这批方法和 emit_tool_call 一样 await + 读返回值,但都包 try-except(错误转发 emit_error)。只有 emit_tool_call 是裸的。

4.4 ctx:扩展上下文

handler 签名 (event, ctx) 里的 ctxExtensionContext,能力远强于管道 A 的 event + signal

# ============================================================
# 【Python 改写】ExtensionContext 能力清单(概念改写)
# 原文 TS (extensions/types.ts:307-347):
#   interface ExtensionContext {
#       ui: ExtensionUIContext;  mode;  cwd: string;
#       sessionManager: ReadonlySessionManager;   // ★ 只读!
#       modelRegistry;  model;  thinkingLevel;
#       signal: AbortSignal | undefined;  abort(): void;
#       isIdle(): boolean;  compact(options?): void;  getSystemPrompt(): string;
#   }
# ============================================================

@dataclass
class ExtensionContext:
    ui: ExtensionUIContext                  # select/confirm/input/notify 等 UI 能力
    mode: str                               # "tui" | "rpc" | "json" | "print"
    cwd: str                                # 工作目录
    session_manager: ReadonlySessionManager # ★ 只读!能读历史但不能直接写
    model_registry: ModelRegistry           # 模型注册表
    model: Optional[Model]                  # 当前模型
    thinking_level: Optional[ThinkingLevel]
    signal: Optional[AbortSignal]           # 中断信号
    # 方法:abort() / is_idle() / compact() / get_system_prompt() / ...

一个容易踩的坑session_managerReadonlySessionManager——能读会话历史(get_entries()),但不能直接写。扩展要写 session 得走 pi.send_message() / pi.append_entry() 等 action 方法,或在命令上下文(ExtensionCommandContext,能力更全)里操作。

4.5 扩展独占事件的触发位置

管道 B 独占的事件(包括 5 个决策点,以及 provider 类事件),都不在 _emit_extension_event 的翻译列表里(2.2 提到,那个方法只翻译 2.1 的流式生命周期事件),而是 SDK 在各自的决策点主动调用的。它们的触发位置分布在两个文件:

扩展独占事件触发位置挂在哪个内核 hook 上
tool_call(执行前)agent-session.ts:468agent.beforeToolCall
tool_result(执行后)agent-session.ts:490agent.afterToolCall
input(用户输入后)agent-session.ts:1131sendUserMessage 流程,skill/template 展开前
before_agent_start(开跑前)agent-session.ts:1224sendUserMessage 流程,agent.run
context(发 LLM 前)sdk.ts:350agent.transformContext
before_provider_request(HTTP 发出前)sdk.ts:331onPayload 回调
after_provider_response(收到响应)sdk.ts:338onResponse 回调

注:上表是 TS 源码的文件名和行号。Python 版块的读者:这些是 SDK 内部实现位置,理解触发时机即可,不需要去翻 TS 文件。

这张表回答了 2.2 那个”管道 A 为什么收不到这些”的问题:它们不是 Agent 内核 emit 出来的(内核根本没这些 type),而是 SDK 在上述位置主动调用 extension_runner.emit_xxx() 的产物。 管道 A 监听的是 _emit(被 _handle_agent_event 调用),自然听不到这些。

4.6 什么时候用管道 B

一句话:需要改变 Agent 行为的场景。 拦截危险工具调用、改写发给 LLM 的上下文、替换系统提示词、修改工具返回值、过滤用户输入——这些”干预”类需求,管道 A 做不到(不等你就意味着返回值丢弃),只能写扩展走管道 B。


五、实战:四个场景,各走哪条管道

理解了两条管道的机制,就能判断每个需求该用哪条。下面四个场景覆盖最常见的集成需求,按管道分组——这是本章最重要的实战判断。

5.1 管道 A 实战组(session.subscribe,只读观察、Agent 不等)

场景1:实时观测 Agent 在干什么

# ============================================================
# 【Python 改写】订阅事件做实时观测(管道 A:广播)
# 原文 TS:
#   session.subscribe((event) => {
#       if (event.type === "tool_execution_start") {
#           console.log(`🔧 ${event.toolName}(${JSON.stringify(event.args).slice(0, 50)})`);
#       }
#       if (event.type === "tool_execution_end") {
#           console.log(`   └─ ${event.isError ? "❌ 失败" : "✅ 成功"}`);
#       }
#   });
# ============================================================

import json

def listener(event):
    if event["type"] == "tool_execution_start":
        args_str = json.dumps(event["args"])[:50]
        print(f"🔧 {event['tool_name']}({args_str})")
    if event["type"] == "tool_execution_end":
        result_str = "❌ 失败" if event["is_error"] else "✅ 成功"
        print(f"   └─ {result_str}")

session.subscribe(listener)

为什么走管道 A:纯观察,不改变 Agent 行为。Agent 不等你也无所谓——你要做的只是”看到事件、打印一行”,同步即可完成。Pi 的 TUI 本身就是通过订阅事件实现的观测面板——你看到的所有终端输出都来自管道 A 的消费。

场景2:流式转发到 Web 前端(SSE)

# ============================================================
# 【Python 改写】订阅事件流转到 Web 前端 SSE(管道 A:广播)
# 原文 TS:
#   session.subscribe((event) => {
#       if (event.type === "message_update") {
#           res.write(`data: ${JSON.stringify({ type: "delta", text: extractText(event.message) })}\n\n`);
#       }
#       if (event.type === "agent_end") { res.end(); }
#   });
# ============================================================

import json

def listener(event):
    if event["type"] == "message_update":
        delta = {"type": "delta", "text": extract_text(event["message"])}
        res.write(f"data: {json.dumps(delta)}\n\n")
    if event["type"] == "agent_end":
        res.end()

session.subscribe(listener)

为什么走管道 A:纯转发,不改 Agent。Agent 运行在服务器上,用户通过浏览器访问,订阅事件流通过 SSE 推给浏览器——这就是 Web 集成的核心。

5.2 管道 B 实战组(扩展 pi.on,能干预、Agent 等你)

场景3:工具调用拦截

# ============================================================
# 【Python 改写】工具调用拦截(管道 B:扩展 hook)
# 原文 TS:
#   function guardExtension(pi) {
#       pi.on("tool_call", async (event) => {
#           if (event.toolName === "delete_table") {
#               return { block: true, reason: "生产环境禁止删除操作" };
#           }
#           return undefined;
#       });
#   }
# ============================================================

def guard_extension(pi):
    @pi.on("tool_call")
    async def handler(event):
        if event["tool_name"] == "delete_table":
            return {"block": True, "reason": "生产环境禁止删除操作"}
        return None   # 放行

# 挂载:extension_factories=[guard_extension]

为什么必须走管道 B:要拦截、要返回 {"block": True}——管道 A 的 listener 不被等、返回值被丢弃,根本拦不住。而且 tool_call 这个事件只在管道 B 派发,管道 A 收不到。第 5 章讲的五步管道中第 3 步 before_tool_call,底层就是这条管道实现的。

场景4:上下文预处理

# ============================================================
# 【Python 改写】上下文预处理(管道 B:扩展 hook)
# 原文 TS:
#   function contextExtension(pi) {
#       pi.on("context", async (event) => {
#           return { messages: [{ role: "user", content: `当前时间:${new Date()}` }, ...event.messages] };
#       });
#   }
# ============================================================

from datetime import datetime

def context_extension(pi):
    @pi.on("context")
    async def handler(event):
        # 在 LLM 调用前,往消息列表里注入当前时间
        now_msg = {"role": "user", "content": f"当前时间:{datetime.now()}"}
        return {"messages": [now_msg, *event["messages"]]}

为什么必须走管道 B:要改写发给 LLM 的消息列表——这是”改变 Agent 行为”,必须返回新列表让 Agent 采用(Agent 会等你返回)。第 6 章讲的 transform_context 钩子是同一条管道的实现(内核层的 transform_context 配置项和扩展层的 context hook 是同一件事的两面,一个走配置、一个走扩展)。

5.3 怎么选管道?一句话判断

你的代码要不要改变 Agent 的行为?

  • 要(拦截、改参数、改消息、存状态)→ 写扩展,走管道 B(pi.on。Agent 会等你、读你的返回值。
  • 不要(打日志、推前端、记统计)→ 用 管道 A(subscribe 更轻。Agent 不等你、不读返回值。

5.4 小结

这四个场景的共同点是:都不需要修改 Agent 内核。 无论走哪条管道,新增功能都只是”挂一段自己的代码”。

事件驱动架构的真正威力不是”通知机制”(那只是管道 A),而是开放扩展机制——尤其是管道 B,它让第三方在不碰内核源码的前提下,能拦截、能改写、能干预。两条管道合起来,才构成 Pi 完整的”神经系统”。


六、总结:两条管道

Pi 的事件系统有一条核心分叉:Agent 对扩展等、对 subscribe 不等。这个差别决定了两条管道的全部行为。

管道 A(session.subscribe —— Agent 不等你的监听器,返回值丢弃。所以你只能观察(渲染、日志、转发),改不了 Agent 的行为。它轻量、注册简单,适合在外部脚本里用。async 监听器的错误不会被 Agent 捕获,要自己 try-except。

管道 B(扩展 pi.on —— Agent 等你的 handler 返回,读返回值。所以你能干预 Agent 的下一步:拦工具、改上下文、换提示词。它还独占 5 个决策点事件(input/before_agent_start/context/tool_call/tool_result),管道 A 收不到。写法上是把代码塞进扩展、挂到 loader。扩展 handler 的异常多数被框架隔离(单个扩展崩了不连累别人),但 tool_call 是例外——它不隔离,扩展出错就拦掉工具(fail-closed:宁可错杀,不放行危险操作)。

判断用哪条管道,只问一句:你的代码要不要改变 Agent 的行为——要,写扩展走管道 B(Agent 会等你);不要,subscribe 走管道 A(Agent 不等你)。


七、下一站

本章我们看到,事件系统让 Agent 和外部世界彻底解耦——UI、日志、持久化、扩展,全部通过订阅事件工作。两条管道(subscribe + pi.on)合起来,覆盖了从”纯观察(不等)“到”深度干预(等+返回值)“的全部需求。

但有一个和事件密切相关的机制我们只提了一句:transform_context(管道 B 的 context hook 在内核层的对应物)。第 6 章讲消息系统时说它在 convert_to_llm 之前执行,负责裁剪旧消息、注入外部上下文。当对话越来越长,消息越来越多,最终会超出模型的上下文窗口。这时候 transform_context 需要做一件更激进的事——压缩对话历史

接下来两章我们就打开 Pi 的上下文工程全貌。第 8 章先讲全景——从输入侧的工具输出截断、系统提示词组装,到历史侧的 Compaction 与分支摘要,让你看清 Pi 在多个环节布置的防线;第 9 章再深入其中最核心的压缩算法(Compaction),看 Pi 怎么在上下文窗口快满时把 50 轮对话压缩成一段结构化摘要,让 Agent 继续”记住”之前发生了什么。


本章关键源码索引

  • packages/agent/src/types.ts:422-437 — 10 种 AgentEvent 定义(事件源)
  • packages/agent/src/agent-loop.ts:25AgentEventSink 类型(emit 签名)
  • packages/agent/src/agent.ts:529-576processEvents(内核同步屏障,await 汇入口)
  • packages/agent/src/agent.ts:173,243listeners Set 和 subscribe(内核层)
  • packages/agent/src/agent-loop.ts:666-707executePreparedToolCall(update 特殊处理)
  • packages/coding-agent/src/core/agent-session.ts:393agent.subscribe(this._handleAgentEvent)(汇入口注册)
  • packages/coding-agent/src/core/agent-session.ts:548-552_emit(管道 A 实体,同步不等
  • packages/coding-agent/src/core/agent-session.ts:595-666_handleAgentEvent(两条管道的分叉点:619 行 await 管道 B,622 行不等管道 A)
  • packages/coding-agent/src/core/agent-session.ts:800-807AgentSession.subscribe(注册到 _eventListeners
  • packages/coding-agent/src/core/agent-session.ts:139-181AgentSessionEvent(Session 层事件)
  • packages/coding-agent/src/core/agent-session.ts:712-793_emitExtensionEvent(翻译给管道 B)
  • packages/coding-agent/src/core/extensions/types.ts:1190-1231pi.on 的 30 个重载(全部扩展事件名)
  • packages/coding-agent/src/core/extensions/types.ts:1180ExtensionHandler 签名 (event, ctx) => Result
  • packages/coding-agent/src/core/extensions/loader.ts:238-243pi.on 实现(往 Map 里 push)
  • packages/coding-agent/src/core/extensions/runner.ts:796-828emit()(通知型派发,try-catch 隔离)
  • packages/coding-agent/src/core/extensions/runner.ts:927-948emitToolCall()(决策型,读返回值,★无 try-catch)
  • packages/coding-agent/src/core/extensions/runner.ts:979-1010emitContext()(链式 transform)
  • packages/coding-agent/src/core/extensions/runner.ts:668-746createContext()(ctx 惰性 getter)
  • packages/coding-agent/src/core/extensions/types.ts:307-347ExtensionContext(ctx 能力,sessionManager 只读)
  • packages/coding-agent/src/core/agent-session.ts:468-517tool_call/tool_result 触发(agent hooks)
  • packages/coding-agent/src/core/agent-session.ts:1131,1224input/before_agent_start 触发
  • packages/coding-agent/src/core/sdk.ts:331,338,350context/before_provider_*/after_provider_response 触发