从 0 开始实现一个 Coding Agent · 第 3 章:流式输出
/ 7 min read
第 3 章 · 流式输出:SSE、统一事件词汇与组装器
对应代码:code/03-streaming.ts(约 230 行,零依赖)
学习目标
- 手写 SSE 解析,理解
data:帧与[DONE]终止符。 - 把 wire 协议的增量碎片翻译成一组统一的事件词汇(text-delta / tool-call-delta / usage / finish)。
- 写一个组装器,把事件流拼回一条可入历史的完整消息——尤其是工具调用参数的分段拼接。
概念:三层结构
流式不是“把打印语句加进去”那么简单。工程上它分三层,每层职责单一,这也是真实 harness 的分层:
HTTP/SSE 字节流 ──解析──▶ 增量事件(StreamChunk)──组装──▶ 完整消息- 解析层:读响应体,按行剥出
data:载荷,处理分包(一个 SSE 帧可能被切成两次网络读)。 - 词汇层:把协议特有的字段(
choices[0].delta.content等)翻译成中立事件。上层代码从此不认识任何具体协议。 - 组装层:把事件流折叠回
{ role, content, tool_calls }的完整消息。只有组装结果才进入历史。
关键代码
SSE 解析器。要点:网络分包随时可能把一行切成两半,所以维护一个 buffer,indexOf('\n') 逐行消费,剩余的留给下一轮:
async function* sseData(res: Response): AsyncIterable<string> { const reader = res.body!.getReader() const decoder = new TextDecoder() let buffer = '' while (true) { const { done, value } = await reader.read() if (done) break buffer += decoder.decode(value, { stream: true }) let newline = buffer.indexOf('\n') while (newline >= 0) { const line = buffer.slice(0, newline).replace(/\r$/, '') buffer = buffer.slice(newline + 1) if (line.startsWith('data:')) yield line.slice(5).trim() newline = buffer.indexOf('\n') } }}事件词汇。这是本章真正的主角——一组中立的、协议无关的流事件:
type StreamChunk = | { type: 'text-delta'; text: string } | { type: 'tool-call-delta'; index: number; id?: string; name?: string; argumentsDelta: string } | { type: 'usage'; usage: Usage } | { type: 'finish'; reason: 'stop' | 'tool-calls' | 'length' | 'error' }工具调用参数的拼接。流式下,一次工具调用被拆成很多帧:第一帧带 id 和 name,后续帧只带 index 和一小段 arguments 增量。按 index 聚拢、逐段累加,finish 到来时得到完整 JSON 字符串:
case 'tool-call-delta': { const call = calls.get(chunk.index) ?? { id: '', name: '', arguments: '' } if (chunk.id !== undefined) call.id = chunk.id if (chunk.name !== undefined) call.name = chunk.name call.arguments += chunk.argumentsDelta calls.set(chunk.index, call) break}组装器遍历全部事件,产出完整消息。循环(第 2 章)从此改为:先收集事件 → 组装 → 推历史,流式只是“边收集边给用户看”的副作用。这个顺序很重要:历史里只进完整消息,不进碎片。
运行:
DEEPSEEK_API_KEY=sk-... node --experimental-strip-types 03-streaming.ts设计要点与坑
- 不要把协议类型泄漏到上层。如果循环代码里出现
choices[0].delta,你就被这个协议绑架了。词汇层是留给“换模型厂商”的保险。 finish_reason: 'length'要当回事。它意味着输出被 max-tokens 截断——如果截断发生在一个工具调用的参数中间,那条调用就是废的,组装器应丢弃不完整调用,循环要把本轮标记为“因长度收束”,否则下一轮 API 会因为“调用没有配对结果”拒绝请求。本章代码已处理:assemble过滤参数没收完的调用,循环在length时收束本轮并告警。- usage 也要流式拿(OpenAI 兼容协议用
stream_options: { include_usage: true },在最后一个 chunk 上送达)。没有它,你无法做第 7 章的上下文预算。
对照 DeepSeek Harness
- 事件词汇几乎同名。真实实现的
StreamChunk是七种事件的闭式联合:block-start / text-delta / reasoning-delta / tool-call-delta / block-end / usage / finish(packages/llm/llm/src/types.ts:452-464)。教程版少了block-start/block-end——那是因为真实消息由 ContentBlock 组成,流按“块”开合;理解了本章的index聚拢,就理解了块机制的一半。 - 组装器是同一个算法。
BlockAssembler按 index 聚拢增量、block-end权威收口、max-tokens 时丢弃未完成的 tool call(packages/llm/llm/src/assembler.ts:69-76, 135-141)。你刚写的就是它的简化版。 - SSE 解析不手写。仓库用
eventsource-parser/stream的EventSourceParserStreampipe 解帧(packages/llm/llm-deepseek/src/sse.ts:13-28),并校验 SSE event 名与帧内event.type一致;error 帧统一转成结构化的 provider 错误(AUTH/RATE_LIMIT/CONTEXT_WINDOW_EXCEEDED 等,transport.ts:22-44)。教程手写一遍是为了看清它只是“按行切 data:”,生产实现换成解析器库即可。 - 失败与重试是一等公民。真实实现把一次流封装成 attempt:失败时先把已收到的碎片作为
assistant/attempt落账(仅日志、不进历史),再触发agent/request-error瀑布;重试插件llm-retry按策略指数退避(packages/llm/llm-retry/src/index.ts:59),且重试复用同一次已渲染的请求,不重复提交用户消息。教程的llm.ts现在有一个最小版:只重试流开始前的 429/5xx 与网络错误(指数退避 + 抖动),流开始后的失败不重试——已展示给用户的增量无法撤销。 - 取消时保留用户已看到的内容。用户按 Ctrl+C 打断流,已到达的前缀会以
interrupted: true的 assistant 消息持久化(packages/core/agent-loop/src/agent.ts:448-465)——“用户看到过的,必须留下”,这是模型可见性与 UI 一致性的要求。
练习
- 在词汇层加一种新事件
reasoning-delta(deepseek-reasoner 模型的delta.reasoning_content字段),组装器忽略它、UI 把它打成灰色——体验“词汇层隔离协议”的好处。 - 故意在组装后丢弃
finish_reason,然后让模型生成一个超长工具调用,观察第 2 章的循环如何死掉。 - (思考)如果两个工具调用的帧交错到达(
index: 0与index: 1交替),你的组装器还能工作吗?
本章留下的问题
现在 agent 能跑、能看、能干活了,但它手里只有一把“跑命令”的锤子。下一章给它一套真正的工具:带行号的读、防误伤的写、唯一性校验的编辑,以及一个像样的 bash 执行器(超时、截断、退出码)。