Skip to content

Latest commit

 

History

History
239 lines (179 loc) · 6.45 KB

File metadata and controls

239 lines (179 loc) · 6.45 KB

第六章 · 事件流

事件流是 Pi 的神经系统——Agent 内部发生的一切,
都以事件的形式传递到 UI、存储、扩展等外部消费者。


6.1 为什么自己造 EventStream

Node.js 有 EventEmitter,浏览器有 CustomEvent,RxJS 有 Observable。Pi 为什么要自己实现?

因为 Pi 的需求很特殊:

需求 现有方案的问题
流式消费 + 最终结果 EventEmitter 没有"结果"概念
背压处理 EventEmitter 是同步的,不处理背压
async iteration 需要同时支持 for await 和回调
部分 JSON 解析 工具参数在流式传输中需要实时解析

Pi 的 EventStream 是一个轻量级异步可迭代流,同时支持事件消费和最终结果提取。

6.2 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 : [],
);

6.3 两种消费方式

方式 A:for-await-of(流式消费)

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();

方式 B:subscribe(回调消费)

Agent 类封装了 subscribe 模式:

agent.subscribe((event, signal) => {
  // 与 for-await-of 相同的事件类型
  if (event.type === "message_update") {
    renderDelta(event);
  }
});

区别:subscribe 的监听器在事件发射时同步等待(awaited),确保监听器有序执行。

6.4 Agent 事件类型全表

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

6.5 Harness 扩展事件

除了 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 的行为。

6.6 流式 JSON 解析

当 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 解析"。

6.7 背压处理

如果事件消费者处理速度慢于生产速度(比如 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

6.8 本章检查清单

  • 理解 Pi 为什么自己实现 EventStream
  • 知道两种事件消费方式(for-await / subscribe)
  • 能列出三组事件类型的代表
  • 理解流式 JSON 解析的意义
  • 知道背压是什么以及 EventStream 如何处理

上一章第五章 · 消息系统
下一章第七章 · 上下文工程 — Agent 的"记忆管理"