Appearance
1.2 事件流机制
LLM 回复是流式的:模型一个字一个字地吐出来。pi 用 EventStream 这个通用类型承载异步迭代,它是整个系统"数据流动"的血管。这一节我们不只讲"它是什么",而是逐行拆解它的实现、每个方法、以及数据如何在生产者和消费者之间流转。
先建立心智模型
把它想成一条有缓冲的传送带:
- 生产者(Producer):用
push()把事件放到传送带上。 - 消费者(Consumer):用
for await...of逐个取走事件。 - 最终结果(Result):用
result()等待"最后一件包裹"。 - 结束(End):生产者用
end()通知传送带停运。
一、它要解决的核心问题
一个流式回复,其实有两个彼此冲突的诉求:
| 诉求 | 想要的 API | 例子 |
|---|---|---|
| ① 逐一看到增量 | 异步迭代 | 打字机效果:每个 text_delta 渲染一行 |
| ② 等最终完整消息 | 一个 Promise | result() 拿到整条 AssistantMessage |
EventStream 同时满足两者:它既是 AsyncIterable(给①),又暴露 result() Promise(给②)。这是它区别于普通队列的关键。
二、内部数据结构(两条"暂存区")
packages/ai/src/utils/event-stream.ts:4 的 EventStream<T, R> 内部有两个数组,分别给"生产者"和"消费者"当暂存:
ts
export class EventStream<T, R = T> implements AsyncIterable<T> {
private queue: T[] = []; // ① 已推送、尚未被消费的事件
private waiting: ((v: IteratorResult<T>) => void)[] = []; // ② 已在等事件、但还没等到的人
// ...
}queue是"产品暂存架":生产者push了事件,但消费者还没来取,就先放这里。waiting是"等待者名单":消费者已经await了,但生产者还没push,就先记在这里。
关键:queue 和 waiting 永远不会同时非空。 因为:
push时:若有人在等(waiting非空)→ 直接交给等待者;否则入queue。- 消费时:若
queue有货 → 直接取;否则加入waiting等待。
生产者产得快、消费者消费慢 → 货积累在 queue;消费者等得急、生产者产得慢 → 人积累在 waiting。两者互补,构成一个"零拷贝交接"的缓冲。
三、构造器:两个"规则函数"
ts
constructor(
isComplete: (event: T) => boolean, // 判断:哪个事件代表"流结束"
extractResult: (event: T) => R, // 提取:从结束事件里取出"最终结果"
) { ... }isComplete:流什么时候算结束。比如"事件类型是done或error就算结束"。extractResult:结束事件里,哪个字段是最终结果。比如从done事件里取.message。
这两个函数让 EventStream 变成一个通用骨架:具体流(助手消息流、日志流…)只需告诉它"何时结束、结果在哪",就能复用全部机制。
四、核心方法逐个拆解
4.1 push(event) —— 生产者放货
ts
push(event: T): void {
if (this.done) return; // 流已结束,丢弃新事件
if (this.isComplete(event)) { // 这个事件是"结束事件"?
this.done = true; // 标记结束
this.resolveFinalResult(this.extractResult(event)); // 提前结算 result()
}
const waiter = this.waiting.shift(); // 有等待者吗?
if (waiter) {
waiter({ value: event, done: false }); // 有 → 直接塞给等待者(不经 queue)
} else {
this.queue.push(event); // 无 → 先放到暂存架
}
}数据流转(生产者视角):
push(event)
│
▼
有人在等(waiting) ?
├─ 是 → 直接交给等待者,不入队
└─ 否 → 压入 queue 暂存架注意一个细节:push 一个"结束事件"时,会提前 resolve result()(见 4.3)—— 因为结束事件自带最终结果,不必等 end()。
4.2 迭代器 [Symbol.asyncIterator]() —— 消费者取货
ts
async *[Symbol.asyncIterator](): AsyncIterator<T> {
while (true) {
if (this.queue.length > 0) { // queue 有货 → 直接取
yield this.queue.shift()!;
} else if (this.done) { // 没货且已结束 → 停止
return;
} else { // 没货也没结束 → 加入等待名单
const result = await new Promise<IteratorResult<T>>((resolve) => this.waiting.push(resolve));
if (result.done) return; // end() 会以 done=true 唤醒
yield result.value;
}
}
}数据流转(消费者视角):
for await(...)
│
▼
queue 有货?
├─ 是 → yield 队首事件
├─ done → return(结束)
└─ 否 → 加入 waiting,等 push小知识:await new Promise(resolve => this.waiting.push(resolve)) 是什么意思
这是 EventStream 里最"绕"的一行,但它是一个很常用的模式:把 Promise 的 resolve 函数存起来,等以后某个时刻再调用它。
1. 一个 Promise 的 resolve 是什么
new Promise((resolve) => {...}) 里的 resolve 是一个函数,你一旦调用它,这个 Promise 就"完成"了:
ts
const p = new Promise((resolve) => {
// 什么都不做,先不调用 resolve
});
await p; // 卡在这里,永远不会结束(因为 resolve 从没被调用)只要 resolve 不被调用,await p 就一直挂起。而 resolve 本身就是一个普通函数,可以把它存到数组里、传给别人、留到以后调用。
2. 拆开那一行
ts
const result = await new Promise<IteratorResult<T>>((resolve) => this.waiting.push(resolve));等价于三步:
ts
// 第 1 步:创建一个 promise,但把它的 resolve 函数存进 waiting 数组
const promise = new Promise<IteratorResult<T>>((resolve) => {
this.waiting.push(resolve); // 注意:这里只 push,不调用 resolve
});
// 第 2 步:await 这个 promise,此刻挂起
const result = await promise;
// 第 3 步:等将来某个时刻,有人调用存起来的 resolve(...),await 才继续它的作用一句话:"我(消费者)现在没货可拿,我先登记一下联系方式(resolve),等有货了别人叫我。"
3. push 那边怎么"叫人"
ts
push(event: T): void {
// ...
const waiter = this.waiting.shift(); // 从等待者名单里取出一个 resolve
if (waiter) {
waiter({ value: event, done: false }); // 调用它,把事件交给消费者
} else {
this.queue.push(event); // 没人等,就放到暂存架
}
}两边一对接:
| 消费者(迭代器) | 生产者(push) |
|---|---|
| 没货可拿 | |
this.waiting.push(resolve) → 登记联系方式 | |
await 挂起 ... | |
| 有货了! | |
← waiter(event) 调用 resolve,把事件送过来 | 用登记的 resolve 唤醒消费者 |
yield 事件 |
4. 为什么要存 resolve 而不是直接传值
waiting 里存的不是"值",而是**"如何唤醒一个被挂起的 await"。for await 的执行流是暂停的,只有拿到 resolve 才能"叫醒"它。类比餐馆等位:消费者把手机号(resolve)留给叫号员(waiting 数组),叫号员有空位时打电话(调用 resolve)**把消费者叫来。消费者不可能一直占着窗口(await 必须挂起),只能靠电话被叫醒。
5. 一句话记忆
ts
// waiting = [ (value) => { 让某个 await 继续 } , ... ]
// push 时:waiting.shift()(event) ← 取出一个,用事件数据唤醒它4.3 result() —— 消费者等最终结果
ts
result(): Promise<R> {
return this.finalResultPromise; // 一个在"结束"时被 resolve 的 Promise
}finalResultPromise 在构造器里创建,只有两个地方会 resolve 它:
push()收到一个isComplete事件时(见 4.1)。end(result)显式结束时(见 4.4)。
它不依赖迭代器是否被消费。也就是说:即使你从不 for await,只要生产者 push 了结束事件或调了 end(),result() 一样会给出最终结果。
4.4 end(result?) —— 生产者收工
ts
end(result?: R): void {
this.done = true;
if (result !== undefined) {
this.resolveFinalResult(result); // 若给了最终结果,结算 result()
}
// 唤醒所有仍在 waiting 的消费者,告诉它们"没了"
while (this.waiting.length > 0) {
const waiter = this.waiting.shift()!;
waiter({ value: undefined as any, done: true });
}
}end() 的作用:
- 标记
done = true,让迭代器在queue空时正常return。 - 若传入
result,结算result()。 - 唤醒所有还在
waiting里干等的消费者,让它们的for await可以正常结束(否则会永远挂起)。
五、完整生命周期:一次"生产→消费→结束"的时序
设生产者要发 [E1, E2, E3],其中 E3 是结束事件。
场景 A:消费者先开始 for await,生产者后 push
| 时间 | 消费者 | 生产者 | queue | waiting |
|---|---|---|---|---|
| t1 | for await:queue 空、done=否 | [] | [] | |
| t2 | → 加入 waiting | [] | [C] | |
| t3 | ⏳ 挂起 | push(E1) → 有等待者→直给 | [] | [C]→[] |
| t4 | yield E1 | [] | [] | |
| t5 | 再取:queue 空 → 再等 | push(E2) → 直给 | [] | [C] |
| t6 | yield E2,再等 | push(E3, isComplete) | [] | [C]→[] |
| t7 | yield E3 | (此时 result() 已提前结算) | [] | [] |
| t8 | 再取:queue 空、done=是 → return |
场景 B:生产者先 push,消费者后消费
| 时间 | 生产者 | 消费者 | queue | waiting |
|---|---|---|---|---|
| t1 | push(E1) → 无等待者→入队 | [E1] | [] | |
| t2 | push(E2) → 入队 | [E1,E2] | [] | |
| t3 | push(E3)=结束 → 标记 done | [E1,E2,E3] | [] | |
| t4 | (result() 已结算) | for await:queue 非空→yield | [E2,E3] | [] |
| t5 | yield E2 | [E3] | [] | |
| t6 | yield E3,再取:done→return | [] | [] |
两个场景的对比就是全部精髓:
- 场景 A:消费者"追着"生产者跑(
waiting有人)。 - 场景 B:生产者"堆着"等消费者来拿(
queue有货)。 - 无论谁先谁后,事件顺序不乱、一个不丢 —— 这就是
queue+waiting互补缓冲的意义。
六、AssistantMessageEventStream:具体的流
EventStream 是通用骨架。pi 为"模型流式回复"定义了具体子类:
ts
export class AssistantMessageEventStream
extends EventStream<AssistantMessageEvent, AssistantMessage> {
constructor() {
super(
(e) => e.type === "done" || e.type === "error", // 结束事件
(e) => (e.type === "done" ? e.message : e.error), // 最终结果
);
}
}它只做两件事:告诉骨架"何时结束、结果从哪取"。其余机制(缓冲、等待、结算)全部继承。
七、流式事件类型(AssistantMessageEvent)
packages/ai/src/types.ts:512 定义了事件联合类型,分两大类:
① 内容增量事件(都带 partial —— 当前累积到的完整部分消息):
ts
| { type: "start"; partial } // 流开始
| { type: "text_start" | "text_delta" | "text_end"; contentIndex; partial }
| { type: "thinking_start" | "thinking_delta" | "thinking_end"; ... }
| { type: "toolcall_start" | "toolcall_delta" | "toolcall_end"; ... }② 终止事件:
ts
| { type: "done"; reason; message } // 成功,message 是完整回复
| { type: "error"; reason; error } // 失败/取消,error 是带错误信息的消息为什么 partial 每次都带"完整消息"而不是只带增量 delta? 因为消费端(如 TUI)收到任何事件都能直接渲染出"当前完整状态",无需自己拼 delta。这是"增量传输 + 完整渲染"的折中:网络传增量(省流量),内存保持完整(好渲染)。
八、消费端:Agent 循环怎么用
agent-loop.ts 的 streamAssistantResponse(见 2.1)就是典型消费方式:
ts
for await (const event of response) {
switch (event.type) {
case "start":
partialMessage = event.partial; // 拿到首个 partial
context.messages.push(partialMessage); // 压入上下文
await emit({ type: "message_start", message: { ...partialMessage } });
break;
case "text_delta": case "thinking_delta": case "toolcall_delta":
case "text_start": case "text_end": /* ... */
partialMessage = event.partial; // 更新 partial
context.messages[context.messages.length - 1] = partialMessage;
await emit({ type: "message_update", assistantMessageEvent: event, message: { ...partialMessage } });
break;
case "done": case "error":
const finalMessage = await response.result(); // 拿完整消息
context.messages[context.messages.length - 1] = finalMessage; // 回填
await emit({ type: "message_end", message: finalMessage });
return finalMessage;
}
}可以看到:增量事件驱动 message_update 的流式更新,result() 提供最终消息用于回填。事件流的两大诉求(增量 + 结果)在这里被同时满足。
九、小结
EventStream= 可迭代(增量)+ 可 Promise(结果) 的事件管道。- 内部靠
queue(产品暂存架)+waiting(等待者名单)两个互补数组缓冲,两者永不同时非空。 push直接交给等待者或入队;end标记停止并唤醒所有人;result在结束事件到来时提前结算。- 无论生产者消费者谁先谁后,事件顺序不乱、一个不丢。
AssistantMessageEventStream只提供"何时结束、结果在哪"两个规则,机制全复用。- 事件约定:
start→ 内容增量(partial)→done/error。
真实源码位置
EventStream:packages/ai/src/utils/event-stream.ts:4-67AssistantMessageEventStream:同文件:69-83- 事件类型定义:
packages/ai/src/types.ts:512-528 - 消费端:
packages/agent/src/agent-loop.ts:317-361
面试角度:为什么这样设计事件流
Q1:为什么 EventStream 既要能"逐个迭代"又要能"拿到最终结果"? 因为流式回复的两个诉求本质不同:UI 要增量(打字机效果),循环要最终完整消息(回填上下文)。如果只给迭代器,循环得自己收尾;如果只给 Promise,UI 没法实时渲染。EventStream 同时满足两者,才让"生产者 push、消费者两用"成为可能。
Q2:为什么 queue 和 waiting 会"永不同时非空"?这个设计解决了什么问题? 因为它俩是互补的:push 时有人等就直给、否则入队;消费时队里有货就取、否则登记等待。无论生产者消费者谁先谁后,事件都不会丢、顺序不乱。这是流式传输的核心正确性保证 —— 如果只保留一个数组,先生产后消费会丢事件(没 queue),或先消费后生产会死锁(没 waiting)。
Q3:为什么 push 一个"结束事件"时要提前结算 result(),而不等 end()? 因为结束事件(done/error)本身就携带最终结果。提前结算,让"只 push 结束事件、不调 end()"的用法也能拿到结果(见 Demo 场景 C),也让 result() 不依赖迭代器是否被消费。
Q4:为什么事件流规定必须以 done 或 error 结束,且不能 throw? 这是 StreamFunction 契约 的一部分。Agent 循环要保证事件序列永远完整(start → … → done/error)。如果底层直接 throw,循环就崩了,无法发出正常的 agent_end 关闭流程。把错误编码进流,循环才能优雅处理。
Q5:为什么每个增量事件都带 partial(完整累积)而不是只带 delta? 为了"拿到任意事件就能直接渲染当前完整状态"。TUI 收到 text_delta 时,直接用 partial 渲染即可,不用自己维护拼接状态。这是"网络传增量(省流量)+ 内存保完整(好渲染)"的折中。
十、动手实验
配对的工程化 Demo 有两个:
① 数据流转演示(推荐先看这个)—— pi-principles/play/event-stream.ts 用三种场景验证"顺序不乱、一个不丢":
bash
cd pi-principles
bun play/event-stream.ts【场景 A】消费者先开始 for await,生产者后 push
消费到: E1, E2, <=END | result: <=END
【场景 B】生产者先 push,消费者后消费
消费到: W1, W2, <=END | result: <=END
【场景 C】结束事件提前结算 result(),无需 end()
result(): <=END② 应用到模型流 —— pi-principles/play/ai.ts 演示"消费增量 + result() 拿结果 + 中止编码进流":
bash
cd pi-principles
bun play/ai.ts思考题
如果去掉 queue(生产者直接 push 给消费者),场景 B(先生产后消费)会怎样?—— 事件会丢失。这就是 queue 缓冲存在的意义。
下一步:1.3 Model 与 Provider —— 如何把不同厂商统一起来。