Skip to main content

EventStream:同一个对象怎样交付下一项和最终结果

你将得到什么

上一章中的事件已经放在数组里,读取它们时不用等待。真实模型会分段返回内容:第一段 文字到达时,第二段可能还在网络上;界面想立刻显示第一段,Agent 又要等整条回复结束 后取得 stop reason、usage 和错误信息。 这一章始终使用同一个 EventStream 实例 events。生产者调用 push() 交出事件, 消费者调用异步迭代器的 next() 取下一项,也可以调用 result() 等最终结果。我们会 观察两种时间顺序:事件先到,以及消费者先等待。
这里不连接网络,也不定义 Agent 消息。取消需要生产者观察 AbortSignal,不在 checkpoint 02 的实现范围内。我们只建立一条规则:
同一个结束动作要让迭代器停止,也要让 result() 得到最终值。

先手动读取同一个 events

先看完成后的调用方式。流里的事件只有两种:delta 携带过程片段,done 携带最终 字符串。
构造器收到两个函数。第一个函数识别哪种事件会结束这条流;这种事件叫终态事件。 第二个函数从终态事件中取出 result() 要返回的值。EventStream 因此不需要认识 deltadone 的业务含义。 先不用 for await...of。手动取得迭代器,能直接看到每一次 next() 的结果:
稳定输出是:
输出里出现了两种容易混淆的“完成”。事件 { type: "done" } 仍是流里的一项,所以 读取它的那次 next() 返回 done: false。下一次 next() 才用 IteratorResult.done === true 表示迭代器已经结束。与此同时,终态事件中的 "AB" 完成了 events.result() for await...of 只是替我们反复调用 next(),直到它返回 done: true
接下来把刚才的四次调用放慢,观察 events 内部的值怎样变化。
Checkpoint 02 · 让 push、next 与 result 汇合模式: 重建起终点: parent 是第 01 章完成后的起点;target 是 2 项聚焦测试通过的终点。教学文件: packages/pi-course/src/event-stream.ts动手前只需知道: 事件先到时,events 把它放进一个数组;next() 先到时, events 保存这次等待的 resolve;终态事件既是一项可读取的事件,也会给出最终值。第一次红灯: parent 中还没有 event-stream.ts。首次 build 会报告 Cannot find module '../src/event-stream.js'。测试中的 .js 是编译后的导入路径; 你要创建对应的 .ts 源文件。第一步: 先不看 target diff。运行 build,确认第一条错误;随后读“事件先到”这条 时间线,声明完整公共接口,只实现 queue 路径。聚焦测试: packages/pi-course/test/02-event-stream.test.ts定位命令: npm run checkpoint -w @pi/course -- 02练习目录: npm run practice -w @pi/course -- 02聚焦运行: npm run build -w @pi/course,然后运行 node --test packages/pi-course/dist/test/02-*.test.js通过证据: 2 项测试分别固定“事件先到”和“等待先到”的输出;终态事件仍能被 next() 读取,result() 也从同一次结束中取得值。

事件先到:它留在 queue 里

新建 events 时,还没有事件,也没有消费者在等待:
生产者先交出第一项:
push() 检查出这不是终态。此时没有一项 next() 正在等待,events 只能暂存这条 事件。保存“已经到达、还没有被读取的事件”的数组叫 queue,也就是队列:
消费者随后调用 next()。迭代器先看 queue,从头取出 delta "A",立刻返回:
再让终态事件先于下一次 next() 到达:
这一次 push() 做两件事。它用 "AB" 完成最终 Promise,并把 done 设为 true,表示以后不再接受新事件。终态事件本身仍然进入 queue:
现在内部已经是 done = true,queue 却还有一项。因此迭代器每次循环要按照固定顺序 检查:
如果先检查 done,消费者会漏掉已经排队的终态事件。正确顺序下,接下来的两次 next() 分别得到:
现在 queue = [done "AB"],同时 done = true。若迭代器先检查 done,再检查 queue,消费者会观察到什么?events.result() 又会得到什么?答案迭代器会直接结束,漏掉终态事件;result() 仍会得到 "AB"。同一个流由此出现两 份不一致的观察结果。迭代器要先清空 queue,再根据 done 结束。

events 保存两条交付队列和一条完成状态

刚才的时间线落在五个字段上:两个数组、一个布尔值,以及构造时配成一对的 Promise 与 resolver:
queue 保存已经到达、尚未读取的事件;waiting 保存已经调用 next()、尚未得到事件 的 resolver。done 是流的结束标记。构造器创建 finalResult 时,把它的 resolve 函数 保存为 resolveFinalResultresult() 始终返回这个 finalResult,所以多个调用者等待 的是同一个最终值。 push(event) 按函数入口时的旧状态决定是否接收本次事件:
第 2 步把 done 改成 true,只影响下一次 push() 和 queue 清空后的 next()。当前这 条终态事件已经通过入口检查,因此仍执行第 3 步:它会交给 waiter 或进入 queue,不会被 刚写入的新 done 状态吞掉。
实践 2.1 · 让先到的事件排队目标: 复现 push(delta) → push(done) → next() → next(),同时取得最终值。文件: packages/pi-course/src/event-stream.ts动作:
  1. 声明构造器、push()end(result)result() 和异步迭代器的完整公共接口。
  2. 创建 queue、waiting、done 与 finalResult;暂未进入的分支可以抛出 new Error("not implemented in lab 2.1")
  3. 实现 push()。终态写入 done、完成 finalResult,随后继续走本次事件的交付路径。
  4. 实现迭代器的 queue 与 done 分支,只运行第一项聚焦测试。
运行:
领域输出: 手动调用 next() 时,先读到 delta "A",再读到 done "AB", 第三次才得到迭代结束;result() 返回 "AB"测试证据: 局部测试应为 1/1,并比较事件顺序和最终字符串。

next 先到:它把 resolve 留在 waiting 里

重新创建一个空的 events,这次消费者先行动:
queue 为空,流也没有结束,所以 next() 还不能返回。迭代器创建一个 Promise,把它的 resolve 放进 waiting
waiting 中的这一个 resolve 就是一名 waiter。它代表一项具体请求:“请把下一条 事件交给我。” 生产者随后推入 delta "A"
push() 取出最早的 waiter,直接用这条事件完成它。事件已经交到消费者手里,无需再 进入 queue:
所以每次 push() 只走一条交付路径:有 waiter 就直接交付,没有 waiter 才排队。 同一条事件不能既唤醒 waiter,又进入 queue,否则消费者会读到两次。 异步迭代器等待下一项的核心动作是:
这段 await 暂停当前迭代器。它没有阻塞 JavaScript 继续执行;生产者仍能调用 events.push(...),再由 push() 唤醒这一次等待。

end() 在没有终态事件时关闭流

有时生产者已经在别处算出最终结果,没有一条终态事件需要交付。end(result) 提供这 条显式关闭路径:
它完成 finalResult,把 done 设为 true,并把仍在 waiting 中的每一项请求唤醒为:
end("A") 不会制造 { type: "done" } 事件。下一次 next() 只观察到迭代结束, result() 则返回 "A"。这与 push({ type: "done", value: "AB" }) 的区别是: 两条路径共享同一个 finalResult 和同一个 doneend() 还会唤醒已经登记的 waiter, 防止它们在流关闭后一直等待。
实践 2.2 · 让先等待的 next 被唤醒目标: 复现 next() → push(delta) → end(result),不借助定时器。文件: packages/pi-course/src/event-stream.ts动作:
  1. queue 为空且流未结束时,创建 Promise,把它的 resolve 放进 waiting。
  2. 下一次 push() 取出最早的 waiter,直接交付事件。
  3. 实现 end(result):完成最终值、标记结束,并把剩余 waiter 唤醒为 done: true
  4. 删除 Lab 2.1 的临时异常,运行两项聚焦测试。
运行:
领域输出: 第一项 next()push(delta "A") 后返回该事件;end("A") 以后, 下一项 next() 返回 done: trueresult() 返回 "A"测试证据: 完整聚焦测试应为 2/2。第二项测试直接比较 pending next 被唤醒后的 IteratorResult

四个公开调用各改哪一份状态

前面的两条时间线可以收成一张状态查找表: 当前 EventStream 的直接范围是一名事件消费者:queue 与 waiting 由这个消费者使用;两个 异步迭代器会争抢同一批事件,而不是各得一份副本。第 03 章把 T/R 换成 Agent 消息并 加入错误终态,第 04、05 章接入模型行为与传输,第 09 章再由运行控制器管理取消。
与当前上游 Pi 对照固定提交 8479bd8packages/ai/src/utils/event-stream.ts 也让 EventStream<T, R> 实现 AsyncIterable<T>,并提供 result(): Promise<R>。 上游的 AssistantMessageEventStream 在这个通用容器上把 doneerror 识别为 终态,并从两者提取最终 AssistantMessage。通用类不认识消息业务类型;第 03 章会沿用 同一分工完成课程里的消息特化。

完成正常路径后再做一次诊断

诊断 · result 完成了,迭代器漏掉终态临时让终态完成 result() 后立刻返回:
运行聚焦测试后,events.result() 仍得到 "AB",事件序列却只剩 delta "A"。这项 差异直接定位到多出的 return:终态没有继续进入 waiter 或 queue。删除它以后,事件 序列恢复为 delta → done,两项测试回到 2/2

本章验收

Checkpoint 02 · 两种到达顺序都能结束运行:
结果应为 2/2:第一项同时观察 queue 中的 delta → done 和最终值 "AB";第二项观察 waiting 中的 next()push(delta) 唤醒,并由 end("A") 结束同一个流。确认只修改 packages/pi-course/src/event-stream.ts。重新定位起终点可运行 npm run checkpoint -w @pi/course -- 02;创建新的隔离练习目录可运行 npm run practice -w @pi/course -- 02 <新目录>。第 03 章会把这里的通用 TR 换成 Agent 的消息与最终回复。

可选迁移练习

迁移 · 用同一容器传送字节块定义 ByteEvent = chunk | completechunk 携带 Uint8Arraycomplete 携带总字节 数。复用现有 EventStream<ByteEvent, number>,分别写出“两个 chunk 先到”和 “next 先等待”两条测试。两条路径都要观察事件顺序、迭代结束和 result();不加入 取消、网络或新的流实现。

小结

EventStream 在一个生命周期上提供两个观察接口:异步迭代器读取每条已接收事件, result() 等待唯一的最终值。queue 与 waiting 解决事件和消费者谁先到,done 与 finalResult 共同收束结束。下一章会保持这套时间协议,只把 T/R 换成 Agent 的消息与 最终回复。