s12: Streaming — 把 push 事件变成 async iterator
一次
for await拿到全部 token,一次await result()拿到聚合结果。半个 JSON 也能解析。
问题
s08 的 provider.complete() 是”一次性返回”——fetch 发出请求,等 LLM 算完,整个 MessagesResponse 一次性回来。但真实 LLM API 是 SSE(Server-Sent Events)流式推送:模型一边生成 token,一边把增量 chunk 推给你,用户看到的是打字机效果而不是干等。
把 s08 改成流式,会冒出两个难题:
-
怎么把 push 事件变成 async iterator? SSE 是 push 模型(服务端推),而 agent 主循环想用
for await拉取。中间需要一个桥梁:生产者随时push(event),消费者随时for await取。两者速度不一致——消费者快了要在那等,生产者快了要先把事件存起来。而且消费者除了逐事件迭代,往往还想要一个”最终聚合结果”(比如完整的 assistant 消息),这不该让消费者自己累加。 -
LLM 输出的 JSON 工具参数可能不完整或非法。 工具调用的
arguments是 JSON 字符串,但它是一块一块流过来的:先是{"comm,再是and":"ls"}。中间任何一刻拿到的都是残缺 JSON,JSON.parse直接抛。更糟的是模型偶尔会吐出控制字符(U+0000~U+001F)、非法转义(\z这种)、甚至行尾裸反斜杠——这些都不是合法 JSON,但 LLM 就是会生成。
s08 的非流式调用完全不用面对这些:等完整响应回来再 JSON.parse 一次就行。流式必须边收边解析,容错是刚需。
解决方案
一个 EventStream 基类同时解决两件事:把 push 变成 async iterator,并额外提供 result() 拿最终聚合值。再加一个 repairJson 修复器 + parseStreamingJson 四层降级,处理残缺和非法 JSON。
{"a":"x<CTRL>"}{{{"a":"x\u0001"}核心是 EventStream 的双输出模型:
const stream = new EventStream<AssistantMessageEvent, AssistantMessage>(isComplete, extractResult);
// 生产者侧
stream.push(event);
stream.end();
// 消费者侧——路径 1:逐事件迭代
for await (const event of stream) { handle(event); }
// 消费者侧——路径 2:拿最终聚合结果
const final = await stream.result();
两条消费路径互不干扰:for await 适合”逐 token 更新打字机”,result() 适合”等完整消息再执行工具”。生产者只管 push,不知道有谁在消费、怎么消费。这就是 无背压生产者-消费者:生产者快则事件入 queue,消费者快则消费者入 waiting 队列、事件到达直接 deliver。没有”生产者等消费者”的阻塞——事件先存着,迟早被取走。
JSON 容错则是另一个独立问题。repairJson 用 单布尔 inString 状态机 逐字符扫描:字符串外的字符原样放行,字符串内的控制字符转义成 \uXXXX、非法转义把反斜杠加倍、行尾裸反斜杠也加倍。parseStreamingJson 在此基础上叠四层 try/catch 瀑布——JSON.parse 原文失败就修复再试,修复失败就补全括号再试,全失败就返回 {}。永远返回一个对象,不抛异常。
工作原理
打开 code.ts,重点看两块。
第 1 块:MiniEventStream 的双输出。 字段只有六个:queue(生产者快时存事件)、waiting(消费者快时存 resolve 函数)、done、finalResultPromise + resolveFinalResult、以及构造时传入的两个回调 isComplete / extractResult。
push(event) 的优先级是关键——先看 waiting,再看 queue:
push(event) {
if (this.done) return;
if (this.isComplete(event)) { this.done = true; this.resolveFinalResult(this.extractResult(event)); }
const waiter = this.waiting.shift();
if (waiter) waiter({ value: event, done: false });
else this.queue.push(event);
}
有消费者在 waiting 里等着,就直接 deliver(绕过 queue,零拷贝、零延迟);没有等待者才入队。这保证了”消费者快时事件不堆积”。isComplete 命中时(如 SSE 的 [DONE])标记 done 并 resolveFinalResult——注意此时事件仍然会投递给消费者,消费者能在迭代里看到这个完成事件。
[Symbol.asyncIterator]() 是三分支循环:
if (this.queue.length > 0) yield this.queue.shift()!; // 分支 1:queue 有货,直接取
else if (this.done) return; // 分支 2:已结束,退出
else { // 分支 3:都没,挂起等待
const result = await new Promise(resolve => this.waiting.push(resolve));
if (result.done) return;
yield result.value;
}
分支 1 处理”生产者已经 push 过了”的积压;分支 2 处理”流已关闭”的收尾;分支 3 处理”消费者先到”的挂起——把 resolve 推进 waiting,等 push 来唤醒。三分支覆盖了所有时序组合,既不丢事件也不死锁。
result() 只是返回 finalResultPromise——这个 Promise 在构造时就建好,只有 isComplete 命中或 end(result) 显式调用时才 resolve。所以 await stream.result() 可以在 for await 结束前就 resolve(完成事件一到就 resolve),也可以在结束后再 await(Promise 已 settled,立即返回)。双输出的含义就在这:迭代和聚合两条路,共享同一个事件流,互不阻塞。
第 2 块:repairJson 单布尔状态机 + parseStreamingJson 四层降级。
repairJson 只用一个布尔 inString 就跟踪了”当前在不在字符串里”。主循环对每个字符分四种情况:
- 字符串外:原样输出,遇
"切进字符串。 - 字符串内遇
":输出并切出字符串。 - 字符串内遇
\:看后继字符nextChar——undefined(行尾裸反斜杠)则加倍成\\;u且后跟 4 位 hex 则保留合法的\uXXXX;属于VALID_ESCAPES集合则保留合法转义;其余一律把反斜杠加倍成\\(让非法转义变成”转义的反斜杠 + 原字符”)。 - 字符串内其余字符:是控制字符(U+0000~U+001F)就转义成
\uXXXX,否则原样输出。
为什么”反斜杠加倍”是安全的修复?JSON 里 \\ 表示一个字面反斜杠。把非法的 \z 改写成 \\z,解析出来就是字符串 \z(反斜杠+z)——语义没丢,只是从”非法转义”变成了”字面反斜杠后跟 z”。比直接删字符或替换更保守。
parseStreamingJson 把 repairJson 嵌进四层 try/catch 瀑布:
try { return JSON.parse(partial); } // 层 1:原文恰好合法
catch {
try { return JSON.parse(repairJson(partial)); } // 层 2:修复后合法
catch {
try { return JSON.parse(closeAllOpen(repairJson(partial))); } // 层 3:补全括号后合法
catch { return {} as T; } // 层 4:兜底空对象
}
}
层 3 的 closeAllOpen 是 partial-json 库的简化桩——扫描括号栈和字符串开闭,结尾补上缺的 " } ]。pi 用 partial-json 库做更聪明的残缺补全(处理 trailing comma、未闭合 key 等),本章只做最小桩。层 4 的 {} 兜底是”绝不抛异常”的承诺——上层 agent 循环拿到 {} 就当”工具参数还没攒够”,等下一个 chunk 再试。
demo 里三个 token 演示了这条瀑布:{"a":1(残缺→层 3 补全)、,"b":"x\x01"(控制字符→层 2 修复)、,"c":"a\z"}(非法转义→层 2 修复),每步都能解析出对象,不会因为半个 chunk 崩掉。
运行
node code.ts
观察模拟 SSE chunk 逐块到达时的流式 JSON 解析:每步打印累积文本和解析结果。控制字符 \x01 被修复成 \u0001,非法转义 \z 被修复成 \\z,残缺括号被 closeAllOpen 补全。最后 await result() 拿到 [DONE] 标记——双输出两条路径都走通了。
前置概念清单
本章只引入五个新概念:
- 异步事件流:把 push 模型(SSE 服务端推)包装成 async iterator(
for await拉取),用Symbol.asyncIterator让对象可迭代 - 生产者-消费者:生产者
push、消费者for await,两者解耦、速度可不一致;本章是无背压版本(事件先存队列,不阻塞生产者) - 双输出模型:同一个流既能
for await逐事件迭代(打字机),又能await result()拿最终聚合值(完整消息);两条路径共享一个finalResultPromise - JSON 修复状态机:单布尔
inString逐字符扫描,处理控制字符转义、非法转义加倍、行尾裸反斜杠——只加字符不删字符,语义保守 - 四层降级:
JSON.parse→repairJson+parse→closeAllOpen+repair+parse→{},层层兜底,永不抛异常
源码锚点
mini-pi 的 EventStream 是教学级简化。pi 的流式基础设施分布在两处,先读再回答问题:
第一处:packages/ai/src/utils/event-stream.ts(88 行)
| 关键行 | 内容 | 本章对应 |
|---|---|---|
| 第 4-19 行 | EventStream 基类字段与构造,finalResultPromise 在构造时建好 | MiniEventStream 构造 |
| 第 21-36 行 | push:isComplete 命中后仍投递事件;waiting.shift() 优先 deliver | MiniEventStream.push |
| 第 38-48 行 | end:唤醒所有等待者返回 {done:true} | MiniEventStream.end |
| 第 50-62 行 | [Symbol.asyncIterator] 三分支 | MiniEventStream 迭代器 |
| 第 64-66 行 | result() 返回 finalResultPromise | MiniEventStream.result |
| 第 69-83 行 | AssistantMessageEventStream 特化类 | 本章省略(见妥协清单) |
读后回答:第 30-35 行 push 在 isComplete 命中后没有 return,仍然会走 deliver/queue——为什么?(提示:消费者需要看到 done 事件来收尾)
第二处:packages/ai/src/utils/json-parse.ts(124 行)
| 关键行 | 内容 | 本章对应 |
|---|---|---|
| 第 3 行 | VALID_JSON_ESCAPES 集合 | VALID_ESCAPES |
| 第 5-25 行 | isControlCharacter + escapeControlCharacter | 同名函数 |
| 第 32-83 行 | repairJson 单布尔 inString 状态机 | repairJson |
| 第 85-95 行 | parseJsonWithRepair:parse 失败则修复再试 | 层 1+2 |
| 第 104-124 行 | parseStreamingJson:parseJsonWithRepair → partialParse → partialParse(repairJson) → {} | 四层降级(层 3 用 partial-json 库) |
读后回答:pi 的 parseStreamingJson 层 3 用 partial-json 库的 partialParse,本章用 closeAllOpen 简化桩替代——partial-json 能处理哪些 closeAllOpen 处理不了的残缺?(提示:trailing comma、未闭合的 key、数组缺元素)
延伸阅读:packages/ai/src/utils/overflow.ts(158 行)
isContextOverflow 用 22 个正则模式(第 33-56 行的 OVERFLOW_PATTERNS)识别各家供应商的上下文溢出错误。本章不实现它,但流式基础设施要能正确传递这类错误事件——EventStream 的 isComplete 会把 error 事件也当完成事件,让 result() resolve 出错误对象,供上层判断是否需要压缩历史(接 s06 的 compaction)。这就是为什么 AssistantMessageEventStream 的 isComplete 同时认 done 和 error。
妥协清单
mini-pi 的流式基础设施比pi 简单很多,以及为什么省略是安全的:
| 省略项 | pi 的做法 | 为什么本章可以省 |
|---|---|---|
partial-json 库 | 层 3 用 partialParse 智能补全残缺 JSON(处理 trailing comma、未闭合 key 等) | closeAllOpen 简化桩够演示降级思路,残缺时返回 {} 语义安全 |
overflow.ts | 22 个正则识别上下文溢出,触发 s06 的压缩 | 本章只做流式管道,错误处理留给上层;error 事件走 isComplete 路径已在源码锚点说明 |
AssistantMessageEventStream 特化类 | 固定 isComplete=(e)=>e.type==="done"||"error",extractResult 按 type 取 message/error | 用泛型 MiniEventStream<T,R> + 两个回调等效表达,少一层继承 |
| 背压 | pi 的 EventStream 同样是无界数组队列(private queue: T[] = []),生产者过快不会阻塞,依赖上游节流 | 教学场景事件量小,无背压版本更易理解 push/await 握手;与 pi 行为一致 |
| 多消费者 | pi 允许多个 for await 并发消费同一流 | 单消费者已能讲清双输出,多消费者只是 waiting 数组多几项 |
| SSE 解析 | pi 按行解析 data: ...\n\n 帧、处理 event: 类型 | demo 直接 push 字符串 token,跳过 SSE 帧解析这层无关细节 |
end(result) 带参 | pi 的 end 可显式传入最终结果(非 complete 事件触发的结束) | demo 用 complete 事件触发 resolve,end() 只做唤醒,简化签名 |
| parseStreamingJson 层级 | pi 层 1 是 parseJsonWithRepair(内部 try JSON.parse 失败 → repairJson),层 2 是 partialParse,层 3 是 partialParse(repairJson) | 教学版层 1 是 JSON.parse,层 2 是 repairJson + JSON.parse,把pi 的层 1 拆成了两层;读pi 源码时会困惑层 1 怎么是 parseJsonWithRepair |
| NON_OVERFLOW_PATTERNS | pi 有 3 个排除正则(overflow.ts 第 67-71 行),排除 rate limit / throttling 等非 overflow 错误 | 教学版无排除逻辑,可能误判限流为溢出 |
这张表是 mini-pi 与pi 的差距地图。读完 pi 的 event-stream.ts 和 json-parse.ts 再回头看,你会看到每一项省略背后都有一道pi 在防的故障——尤其是 partial-json 那层,pi 能解析出半个对象的已有字段,本章只能返回 {}。
默写验收
合上 code.ts 和本 README,打开 practice.ts,凭记忆补全:
MiniEventStream.push()——等待者优先 deliver,否则入队(4 行核心逻辑)MiniEventStream [Symbol.asyncIterator]()——三分支:queue 优先 / done 退出 / new Promise 等待repairJson()——单布尔inString状态机主循环(四类字符分支)parseStreamingJson()——四层降级 try/catch 瀑布
通过标准:node practice.ts 能打印出与 code.ts 一致的流式累积解析结果(三个 token 逐步解析 + result() 返回 [DONE])。
写不出 push 说明”等待者优先 deliver”的握手没进脑子,回到「工作原理」第 1 块重读 waiting.shift() 那行。写不出 repairJson 说明”反斜杠加倍”的修复策略没记住——这是保守修复的核心:不删字符,只把非法转义变成字面反斜杠。写不出四层降级说明”层层兜底、永不抛异常”的承诺没烙下。
十二章走完,mini-pi 的流式管道接通了。 从 s08 的一次性返回,到 s12 的逐 token 流式,agent harness 离pi 又近了一步。EventStream 的双输出让”打字机渲染”和”工具执行”两条路径各取所需,repairJson 让 LLM 的脏输出不再炸掉解析器。