DOCS v0.1.0 BUILD: STABLE GitHub ↗
章节目录
S12

Streaming

EventStream 双输出模型 + 流式 JSON 容错解析 + 四层降级。

前置概念 异步事件流生产者-消费者双输出模型JSON 修复状态机四层降级

s12: Streaming — 把 push 事件变成 async iterator

一次 for await 拿到全部 token,一次 await result() 拿到聚合结果。半个 JSON 也能解析。

...s08 → … → s11s12


问题

s08 的 provider.complete() 是”一次性返回”——fetch 发出请求,等 LLM 算完,整个 MessagesResponse 一次性回来。但真实 LLM API 是 SSE(Server-Sent Events)流式推送:模型一边生成 token,一边把增量 chunk 推给你,用户看到的是打字机效果而不是干等。

把 s08 改成流式,会冒出两个难题:

  1. 怎么把 push 事件变成 async iterator? SSE 是 push 模型(服务端推),而 agent 主循环想用 for await 拉取。中间需要一个桥梁:生产者随时 push(event),消费者随时 for await 取。两者速度不一致——消费者快了要在那等,生产者快了要先把事件存起来。而且消费者除了逐事件迭代,往往还想要一个”最终聚合结果”(比如完整的 assistant 消息),这不该让消费者自己累加。

  2. 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。

EVENTSTREAM 双输出 + REPAIRJSON 状态机
图 A · EventStream 生产者-消费者时序
生产者 push
queue0
waiting1
queue
q0
waiting
await
消费者 await next(),进入 waiting 队列
消费者 for await
await next()
双输出 · result() = (未解析)
图 B · repairJson 单布尔状态机
注入:
输入流{"a":"x<CTRL>"}
当前字符{
inStringfalse
输出流{
字符串外,原样输出
parseStreamingJson 四层降级
层 1 JSON.parse 原文失败
层 2 repairJson + parse命中
层 3 closeAllOpen + parse失败
层 4 返回 {}失败
修复后 = {"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 函数)、donefinalResultPromise + 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])标记 doneresolveFinalResult——注意此时事件仍然会投递给消费者,消费者能在迭代里看到这个完成事件。

[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”。比直接删字符或替换更保守。

parseStreamingJsonrepairJson 嵌进四层 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 的 closeAllOpenpartial-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] 标记——双输出两条路径都走通了。


前置概念清单

本章只引入五个新概念:

  1. 异步事件流:把 push 模型(SSE 服务端推)包装成 async iterator(for await 拉取),用 Symbol.asyncIterator 让对象可迭代
  2. 生产者-消费者:生产者 push、消费者 for await,两者解耦、速度可不一致;本章是无背压版本(事件先存队列,不阻塞生产者)
  3. 双输出模型:同一个流既能 for await 逐事件迭代(打字机),又能 await result() 拿最终聚合值(完整消息);两条路径共享一个 finalResultPromise
  4. JSON 修复状态机:单布尔 inString 逐字符扫描,处理控制字符转义、非法转义加倍、行尾裸反斜杠——只加字符不删字符,语义保守
  5. 四层降级JSON.parserepairJson+parsecloseAllOpen+repair+parse{},层层兜底,永不抛异常

源码锚点

mini-pi 的 EventStream 是教学级简化。pi 的流式基础设施分布在两处,先读再回答问题:

第一处:packages/ai/src/utils/event-stream.ts(88 行)

关键行内容本章对应
第 4-19 行EventStream 基类字段与构造,finalResultPromise 在构造时建好MiniEventStream 构造
第 21-36 行pushisComplete 命中后仍投递事件;waiting.shift() 优先 deliverMiniEventStream.push
第 38-48 行end:唤醒所有等待者返回 {done:true}MiniEventStream.end
第 50-62 行[Symbol.asyncIterator] 三分支MiniEventStream 迭代器
第 64-66 行result() 返回 finalResultPromiseMiniEventStream.result
第 69-83 行AssistantMessageEventStream 特化类本章省略(见妥协清单)

读后回答:第 30-35 行 pushisComplete 命中后没有 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 行parseStreamingJsonparseJsonWithRepairpartialParsepartialParse(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)识别各家供应商的上下文溢出错误。本章不实现它,但流式基础设施要能正确传递这类错误事件——EventStreamisComplete 会把 error 事件也当完成事件,让 result() resolve 出错误对象,供上层判断是否需要压缩历史(接 s06 的 compaction)。这就是为什么 AssistantMessageEventStreamisComplete 同时认 doneerror


妥协清单

mini-pi 的流式基础设施比pi 简单很多,以及为什么省略是安全的:

省略项pi 的做法为什么本章可以省
partial-json层 3 用 partialParse 智能补全残缺 JSON(处理 trailing comma、未闭合 key 等)closeAllOpen 简化桩够演示降级思路,残缺时返回 {} 语义安全
overflow.ts22 个正则识别上下文溢出,触发 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_PATTERNSpi 有 3 个排除正则(overflow.ts 第 67-71 行),排除 rate limit / throttling 等非 overflow 错误教学版无排除逻辑,可能误判限流为溢出

这张表是 mini-pi 与pi 的差距地图。读完 pi 的 event-stream.tsjson-parse.ts 再回头看,你会看到每一项省略背后都有一道pi 在防的故障——尤其是 partial-json 那层,pi 能解析出半个对象的已有字段,本章只能返回 {}


默写验收

合上 code.ts 和本 README,打开 practice.ts,凭记忆补全:

  1. MiniEventStream.push()——等待者优先 deliver,否则入队(4 行核心逻辑)
  2. MiniEventStream [Symbol.asyncIterator]()——三分支:queue 优先 / done 退出 / new Promise 等待
  3. repairJson()——单布尔 inString 状态机主循环(四类字符分支)
  4. 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 的脏输出不再炸掉解析器。


下一章:s13 有状态 Agent — 双队列与状态机