# Step 32: 流式事件系统
一句话导读:step30 用一组回调(
onText/onToolResult/ …)把引擎内部事件抛给外部,能用但控制被「反转」、难记录难重放。step32 把这堆回调统一成一条 async-generator 事件流——run()从普通异步函数变成async *生成器,yield出一连串EngineEvent,消费方在一个for await+switch里夺回控制权。这正是真实 Claude Code 的做法。
# 一、这一步做了什么(What)
不加新功能,把「回调」换成「事件流」,改三个地方:
- 新增
events.ts:定义EngineEvent判别联合(事件协议)。 QueryEngine.run()改成async *run(): AsyncGenerator<EngineEvent>:删掉EngineHooks,所有hooks?.onX()改成yield { type: ... }。api.ts的callWithTools(onChunk)→async *streamTurn():流式文本作为text事件yield,整轮结束return出CallResult。
消费方从「传一堆 onXxx 回调」变成:
for await (const ev of engine.run(input)) {
switch (ev.type) {
case "text": process.stdout.write(ev.text); break;
case "tool_calls": showCalling(ev.names); break;
case "tool_result": showResult(ev); break;
case "round": showRoundStats(ev); break;
case "debug": appendDebug(ev.debug); break;
// ...
}
}
2
3
4
5
6
7
8
9
10
# 二、面试官视角:为什么要做?(Why)
面试题:step30 的回调机制已经能把事件抛出去了,为什么还要大动干戈换成 async generator?两者到底差在哪?
差在控制权归谁,以及事件能不能被当成「数据」二次加工。
回调是「控制反转」(Inversion of Control):你把十几个 onText / onToolResult / onRound 交给引擎,引擎在它想调的时候回调你。于是你的渲染逻辑被拆散在十几个孤立的函数里,读代码得在回调之间跳来跳去,没有一处能看到「一轮对话的完整时序」。而且回调是「转瞬即逝的通知」——没有一个统一的『事件』对象,你想记录、重放、转发(转网页端 / 写日志 / 测试断言)都别扭:回调发生时不处理就没了。
async generator 把这些「通知」变成一条可迭代的事件序列。控制权回到调用方手里:for await 由你驱动,所有事件在一个 switch 里集中处理,一眼看到完整时序。更关键的是,一条序列天然是数据——你可以 map / filter / push 进数组做断言 / 转发到别处。这三件事(记录、重放、测试)回调几乎做不到,事件流几乎免费。
| 维度 | 回调 | 事件流(generator) |
|---|---|---|
| 控制权 | 被反转,逻辑拆散在十几个回调里 | 调用方拿回,集中在一个循环 |
| 可组合 | 难 | 是一条序列,可 map/filter/收集/转发 |
| 记录/重放/测试 | 难 | push 进数组就能断言/回放 |
| 暂停/取消 | 别扭 | 停止迭代即可 |
# 三、原理:它是怎么工作的(How)
# 机制一:EngineEvent —— 一条判别联合的事件协议
events.ts 把所有事件定义成一个判别联合,靠 type 区分:
export type EngineEvent =
| { type: "stream_start" }
| { type: "text"; text: string }
| { type: "stream_end" }
| { type: "tool_calls"; names: string[] }
| { type: "tool_result"; name: string; ok: boolean; elapsed: number; error?: string }
| { type: "denied"; name: string }
| { type: "round"; model: string; inT: number; outT: number; cost: number }
| { type: "budget_warn"; percent: number }
| { type: "budget_exceeded" }
| { type: "error"; msg: string; recoverable: boolean }
| { type: "debug"; debug: { type: string; data: any } };
2
3
4
5
6
7
8
9
10
11
12
消费方一个 switch (ev.type) 全覆盖,且 TS 会在每个 case 里把 ev 收窄——case "text" 里能访问 ev.text,别的字段访问会报错。事件协议本身就是 step31 判别联合思想的又一次应用。
# 机制二:run() 变成 async generator,用 yield 代替回调
async *run(userInput: string): AsyncGenerator<EngineEvent, void> {
this.messages.push(userText(userInput));
for (let round = 0; round < maxRounds; round++) {
if (budget.isExceeded()) { yield { type: "budget_exceeded" }; return; }
yield { type: "stream_start" };
const result = yield* streamTurn(this.messages, registry.toAPI(), sys, model); // ★ yield*
yield { type: "stream_end" };
if (result.toolCalls.length === 0) {
yield { type: "round", model, inT: result.inT, outT: result.outT, cost: cost.stats().total };
return;
}
yield { type: "tool_calls", names: result.toolCalls.map(t => t.name) };
// ...权限、并行执行、yield tool_result...
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
# 机制三:yield* 委托 —— 子生成器的事件自动汇入父流
这是全篇最关键的语法细节。streamTurn 本身也是个 generator:它一路 yield 文本事件,return 出 CallResult。
// streamTurn 的签名:yield 文本事件,return 结果
export async function* streamTurn(...): AsyncGenerator<EngineEvent, CallResult> {
for await (const ev of stream) {
if (ev.type === "content_block_delta" && d?.type === "text_delta") {
fullText += d.text;
yield { type: "text", text: d.text }; // 每个流式片段 yield 出去
}
}
return { text: fullText, contentBlocks: final.content, toolCalls, inT, outT };
}
2
3
4
5
6
7
8
9
10
run() 里一行 const result = yield* streamTurn(...) 同时做了两件事:
yield*把streamTurn内部yield的每个text事件原样转发到最外层消费者——子过程的事件流自动汇入主流,不用手写「把 onChunk 再往上抛一层」;- 把
streamTurn的return值(CallResult)接到result变量。
一行代码,既转发了子流、又取到了返回值。这就是 yield* 委托相比回调的优雅:多层生成器嵌套时,事件会自动逐层冒泡,无需任何手动转发。
# 数据流
消费方 for await ──驱动──> engine.run()
│ yield stream_start
│ yield* streamTurn() ──┐
│ │ yield text (逐片段) ← 经 yield* 冒泡
│ │ yield text ... ← 到最外层
│ └ return CallResult ──> result
│ yield stream_end
│ yield tool_calls / tool_result / round ...
▼
消费方 switch(ev.type) 集中处理(渲染 / 落盘 / 断言)
2
3
4
5
6
7
8
9
10
# 四、深入追问(面试常见 follow-up)
Q:yield 和 yield* 到底有什么区别?为什么转发子生成器要用 yield*?
A:yield x 抛出一个值。yield* gen 是「委托」:它遍历 gen,把 gen 里每次 yield 的值逐个替你抛出去,并在 gen 结束时把它的 return 值作为整个 yield* 表达式的值。如果只用 yield,你得手写 for await (const ev of streamTurn(...)) yield ev; 再单独想办法取返回值——yield* 一行搞定转发 + 取返回值。
Q:权限判断(canUseTool)为什么没有做成事件,还是留着回调?
A:因为权限是**「请求-响应」语义——引擎问一句「能调 Bash 吗」,然后必须停下来等一个 y/n 答复才能继续。而事件流是单向广播**:yield 出去就往下走了,表达「等一个回答再继续」很别扭(你得反过来让消费方往 generator 里 next(value) 塞值,代码立刻绕)。所以 canUseTool 保持注入回调,引擎在循环里直接 await 它。真实 Claude Code 同样是「事件流广播 + canUseTool 回调」的混合设计——单向通知用事件,请求-响应用回调。
Q:async generator 相比回调,为什么就「天然可测试」了?
A:因为一次 run() 的完整行为可以被收集成一个数组:const evts = []; for await (const ev of engine.run(input)) evts.push(ev);,然后直接 expect(evts).toContainEqual({ type: "tool_calls", names: ["Bash"] })。事件序列是纯数据,断言、快照、回放都成立。回调版要测同样的东西,你得 mock 十几个 onXxx、记录调用顺序,繁琐且脆。
Q:这次改动是纯重构,怎么保证和 step30 行为一致?
A:跑同样输入,终端输出逐字对齐 step30。回调 onText(chunk) 和事件 yield {type:"text", text:chunk} 是一一对应的机械替换,语义等价。判据仍是「外部可观测行为不变」。
# 五、踩坑 / 设计权衡
- 不是所有交互都适合事件流:这一步的分寸感在于——识别出「权限」是请求-响应、不该硬塞进单向事件流。盲目「全部事件化」会把简单的
await callback逼成难懂的双向 generator 协议。广播用事件,问答用回调。 - 生成器的返回值容易被忽略:
AsyncGenerator<EngineEvent, CallResult>的第二个类型参数是return类型。很多人只记得 generator 能yield,忘了它还能return一个值供yield*接收——streamTurn正是靠这个把「事件」和「最终结果」分离在两条通道上。 - 消费方必须驱动到底:generator 是惰性的,
for await不跑它就不动。如果消费方提前break,引擎循环也随之停在半途——这既是「可取消」的便利,也意味着必须留意别过早退出导致状态不一致。
# 六、与真实源码的对照
| 我们的实现 | Claude Code 源码 |
|---|---|
run() 是 async * 生成器 | 真实 QueryEngine.submitMessage() 同样是 generator,yield 一连串事件 |
EngineEvent 判别联合 | 真实版丰富的 stream event 类型(message_start / delta / tool_use…) |
yield* streamTurn() 委托 | 真实版对底层 API 流的逐层 yield 委托 |
for await + switch 消费 | 真实版的事件流消费 + 中间件 |
canUseTool 仍是回调 | 真实版 canUseTool 也是回调(请求-响应语义) |
# 七、一句话总结
把一堆回调统一成一条 async-generator 事件流:run() 变 async *、yield 出 EngineEvent,消费方在一个 for await + switch 里夺回控制权;yield* streamTurn() 一行既转发子生成器的文本事件、又接住它的返回结果。事件流相比回调的本质优势是「事件成了可 map/filter/记录/重放的数据」——而权限这种请求-响应语义仍保留回调,广播用流、问答用回调。
# 下一节预告
事件流就位后,下一步可以让 thinking 块顺着这条流出来——step33 解析并展示 Extended Thinking 思考模式。