事件流是 Pi 的神经系统——Agent 内部发生的一切,
都以事件的形式传递到 UI、存储、扩展等外部消费者。
Node.js 有 EventEmitter,浏览器有 CustomEvent,RxJS 有 Observable。Pi 为什么要自己实现?
因为 Pi 的需求很特殊:
| 需求 | 现有方案的问题 |
|---|---|
| 流式消费 + 最终结果 | EventEmitter 没有"结果"概念 |
| 背压处理 | EventEmitter 是同步的,不处理背压 |
| async iteration | 需要同时支持 for await 和回调 |
| 部分 JSON 解析 | 工具参数在流式传输中需要实时解析 |
Pi 的 EventStream 是一个轻量级异步可迭代流,同时支持事件消费和最终结果提取。
// packages/ai/src/utils/event-stream.ts
class EventStream<TEvent, TResult> implements AsyncIterable<TEvent> {
// 推入事件
push(event: TEvent): void;
// 结束流并设置最终结果
end(result: TResult): void;
// 获取最终结果(流结束后可用)
result(): Promise<TResult>;
// 支持 for-await-of
[Symbol.asyncIterator](): AsyncIterator<TEvent>;
}构造时指定两个判断函数:
const stream = new EventStream<AgentEvent, AgentMessage[]>(
// 1. 什么事件算"结束"
(event) => event.type === "agent_end",
// 2. 如何从结束事件提取最终结果
(event) => event.type === "agent_end" ? event.messages : [],
);const stream = agentLoop(prompts, context, config, signal, streamFn);
for await (const event of stream) {
switch (event.type) {
case "message_update":
// 实时渲染 LLM 输出
renderDelta(event.assistantMessageEvent);
break;
case "tool_execution_start":
// 显示工具执行状态
showSpinner(event.toolName);
break;
case "tool_execution_end":
// 显示工具结果
showResult(event.result);
break;
}
}
// 流结束后,获取所有新消息
const newMessages = await stream.result();Agent 类封装了 subscribe 模式:
agent.subscribe((event, signal) => {
// 与 for-await-of 相同的事件类型
if (event.type === "message_update") {
renderDelta(event);
}
});区别:subscribe 的监听器在事件发射时同步等待(awaited),确保监听器有序执行。
Agent Loop 发射的事件分为三组:
| 事件 | 时机 | 携带数据 |
|---|---|---|
agent_start |
循环启动 | 无 |
agent_end |
循环结束 | messages: AgentMessage[] |
turn_start |
每轮开始 | 无 |
turn_end |
每轮结束 | message, toolResults |
| 事件 | 时机 | 携带数据 |
|---|---|---|
message_start |
消息开始 | message: AgentMessage |
message_update |
消息流式更新 | assistantMessageEvent, message |
message_end |
消息完成 | message: AgentMessage |
| 事件 | 时机 | 携带数据 |
|---|---|---|
tool_execution_start |
工具开始执行 | toolCallId, toolName, args |
tool_execution_update |
工具进度更新 | partialResult |
tool_execution_end |
工具执行完毕 | result, isError |
@startuml
skinparam backgroundColor transparent
concise "Agent Loop" as Loop
concise "Events" as E
@0
Loop is "启动"
E is "agent_start"
@1
Loop is "第1轮"
E is "turn_start"
@2
Loop is "用户消息"
E is "message_start → message_end"
@3
Loop is "LLM 流式输出"
E is "message_start → message_update... → message_end"
@4
Loop is "工具执行"
E is "tool_execution_start → tool_execution_end"
@5
Loop is "工具结果"
E is "message_start → message_end"
@6
Loop is "第1轮结束"
E is "turn_end"
@7
Loop is "第2轮"
E is "turn_start → ... → turn_end"
@8
Loop is "结束"
E is "agent_end"
@enduml除了 Agent Loop 的基础事件,AgentHarness(产品层)还会发射更多事件:
type AgentHarnessOwnEvent =
| QueueUpdateEvent // Steering/Follow-up 队列变化
| SavePointEvent // 会话保存点
| AbortEvent // 中止事件
| SettledEvent // 完全停止
| BeforeAgentStartEvent // Agent 启动前(可修改 system prompt)
| ContextEvent // 上下文快照
| BeforeProviderRequestEvent // Provider 请求前
| ToolCallEvent // 工具调用(Harness 级别)
| ToolResultEvent // 工具结果(Harness 级别)
| SessionBeforeCompactEvent // 压缩前
| SessionCompactEvent // 压缩后
| SessionTreeEvent // 会话树变化
| ModelUpdateEvent // 模型切换
| ThinkingLevelUpdateEvent // 思维级别变化
| ToolsUpdateEvent // 工具列表变化
| ResourcesUpdateEvent; // 资源(Skills 等)变化这些事件使得扩展可以在几乎任何时机介入 Agent 的行为。
当 LLM 流式返回工具调用参数时,参数 JSON 是逐 token 到达的:
toolcall_start: { name: "read" }
toolcall_delta: '{"pa'
toolcall_delta: 'th":'
toolcall_delta: ' "fo'
toolcall_delta: 'o.ts'
toolcall_delta: '"}'
toolcall_end
Pi 在 toolcall_end 时才做 JSON 解析和校验,但通过 message_update 事件让 UI 可以实时显示部分参数——这就是"实时不完全 JSON 解析"。
如果事件消费者处理速度慢于生产速度(比如 UI 渲染跟不上 LLM 输出),EventStream 的 push() 会自动排队。消费者通过 for await 按自己的速度消费:
Producer: push(e1) push(e2) push(e3) push(e4)
↓ ↓ ↓ ↓
Buffer: [e1, e2, e3, e4]
↓
Consumer: await e1 ... await e2 ... await e3 ... await e4
- 理解 Pi 为什么自己实现 EventStream
- 知道两种事件消费方式(for-await / subscribe)
- 能列出三组事件类型的代表
- 理解流式 JSON 解析的意义
- 知道背压是什么以及 EventStream 如何处理
上一章:第五章 · 消息系统
下一章:第七章 · 上下文工程 — Agent 的"记忆管理"