# Step 32: 流式事件系统

一句话导读:step30 用一组回调(onText / onToolResult / …)把引擎内部事件抛给外部,能用但控制被「反转」、难记录难重放。step32 把这堆回调统一成一条 async-generator 事件流——run() 从普通异步函数变成 async * 生成器,yield 出一连串 EngineEvent,消费方在一个 for await + switch 里夺回控制权。这正是真实 Claude Code 的做法。


# 一、这一步做了什么(What)

不加新功能,把「回调」换成「事件流」,改三个地方:

  1. 新增 events.ts:定义 EngineEvent 判别联合(事件协议)。
  2. QueryEngine.run() 改成 async *run(): AsyncGenerator<EngineEvent>:删掉 EngineHooks,所有 hooks?.onX() 改成 yield { type: ... }。
  3. 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;
    // ...
  }
}
1
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 } };
1
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...
  }
}
1
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 };
}
1
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) 集中处理(渲染 / 落盘 / 断言)
1
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 思考模式。