Skip to content

1.2 事件流机制

LLM 回复是流式的:模型一个字一个字地吐出来。pi 用 EventStream 这个通用类型承载异步迭代,它是整个系统"数据流动"的血管。这一节我们不只讲"它是什么",而是逐行拆解它的实现、每个方法、以及数据如何在生产者和消费者之间流转

先建立心智模型

把它想成一条有缓冲的传送带

  • 生产者(Producer):用 push() 把事件放到传送带上。
  • 消费者(Consumer):用 for await...of 逐个取走事件。
  • 最终结果(Result):用 result() 等待"最后一件包裹"。
  • 结束(End):生产者用 end() 通知传送带停运。

一、它要解决的核心问题

一个流式回复,其实有两个彼此冲突的诉求

诉求想要的 API例子
① 逐一看到增量异步迭代打字机效果:每个 text_delta 渲染一行
② 等最终完整消息一个 Promiseresult() 拿到整条 AssistantMessage

EventStream 同时满足两者:它既是 AsyncIterable(给①),又暴露 result() Promise(给②)。这是它区别于普通队列的关键。

二、内部数据结构(两条"暂存区")

packages/ai/src/utils/event-stream.ts:4EventStream<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,就先记在这里。

关键:queuewaiting 永远不会同时非空。 因为:

  • push 时:若有人在等(waiting 非空)→ 直接交给等待者;否则入 queue
  • 消费时:若 queue 有货 → 直接取;否则加入 waiting 等待。

生产者产得快、消费者消费慢 → 货积累在 queue;消费者等得急、生产者产得慢 → 人积累在 waiting。两者互补,构成一个"零拷贝交接"的缓冲。

三、构造器:两个"规则函数"

ts
constructor(
	isComplete: (event: T) => boolean,   // 判断:哪个事件代表"流结束"
	extractResult: (event: T) => R,      // 提取:从结束事件里取出"最终结果"
) { ... }
  • isComplete流什么时候算结束。比如"事件类型是 doneerror 就算结束"。
  • 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 它

  1. push() 收到一个 isComplete 事件时(见 4.1)。
  2. 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

时间消费者生产者queuewaiting
t1for await:queue 空、done=否[][]
t2→ 加入 waiting[][C]
t3⏳ 挂起push(E1) → 有等待者→直给[][C]→[]
t4yield E1[][]
t5再取:queue 空 → 再等push(E2) → 直给[][C]
t6yield E2,再等push(E3, isComplete)[][C]→[]
t7yield E3(此时 result() 已提前结算)[][]
t8再取:queue 空、done=是 → return

场景 B:生产者先 push,消费者后消费

时间生产者消费者queuewaiting
t1push(E1) → 无等待者→入队[E1][]
t2push(E2) → 入队[E1,E2][]
t3push(E3)=结束 → 标记 done[E1,E2,E3][]
t4result() 已结算)for await:queue 非空→yield[E2,E3][]
t5yield E2[E3][]
t6yield 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.tsstreamAssistantResponse(见 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
真实源码位置
  • EventStreampackages/ai/src/utils/event-stream.ts:4-67
  • AssistantMessageEventStream:同文件 :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:为什么 queuewaiting 会"永不同时非空"?这个设计解决了什么问题? 因为它俩是互补的:push 时有人等就直给、否则入队;消费时队里有货就取、否则登记等待。无论生产者消费者谁先谁后,事件都不会丢、顺序不乱。这是流式传输的核心正确性保证 —— 如果只保留一个数组,先生产后消费会丢事件(没 queue),或先消费后生产会死锁(没 waiting)。

Q3:为什么 push 一个"结束事件"时要提前结算 result(),而不等 end() 因为结束事件(done/error)本身就携带最终结果。提前结算,让"只 push 结束事件、不调 end()"的用法也能拿到结果(见 Demo 场景 C),也让 result() 不依赖迭代器是否被消费。

Q4:为什么事件流规定必须以 doneerror 结束,且不能 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 —— 如何把不同厂商统一起来。