|DeepSeek-Harness(十五)流式输出管道:从 StreamChunk 到 UI)
一个模型输出的字符从 DeepSeek 服务端的 SSE 报文到浏览器里打字机式跳出来的文字中间要经过好几层完全独立的重组Provider 把裸字节解析成协议无关的StreamChunkLlmRuntime用一个 waterfall 中间件链包一层ReactLoopAgent.step()通过AssistantStreamAttempt边组装边压缩Host/Client 之间的传输层把这份数据推给浏览器Client 再用一个独立的累加器把它重新拼回可渲染的块。本篇按数据实际流动的顺序逐层拆开这条管道并回答一个核心设计取舍问题逐 token 的回放保真度和日志体量这两个互相冲突的目标是怎么被同时满足的。2026-09 更新这一篇讲的核心机制——尤其是落盘这一层——相比早期版本有实质性演进请重点关注第三层和为什么这两节Host/Client 之间具体的传输协议原来的FrameQueue/WebSocket 帧也发生了包名和实现方式的变化完整的 Host-Client RPC 架构以第 06 章为准这里只更新到数据传到 Client 之后要怎么再折叠这个层面。学习目标理解StreamChunk这个协议无关的中间表示长什么样以及 Provider 的 adapter 实现如何把 SSE 字节流转换成它。理解LlmRuntime的llm/streamwaterfall 中间件链的作用——它是插件比如下一篇的 checkpoint 策略介入模型请求即将发出这一时刻的唯一入口。通读BlockAssembler的真实实现理解它如何把碎片化的block-start/text-delta/tool-call-delta/block-end等七种 chunk 类型增量组装成完整的ContentBlock[]。理解AssistantStreamAttemptAssistantStreamAccumulator这套新机制为什么现在不再是每个原始 chunk 都落一条日志而是把连续的同类 delta 压缩成紧凑记录、只落一条assistant/message/assistant/attempt事件同时依然保证可以无损还原出原始的逐 chunk 时间序列。理解实时 UI 更新AssistantStreamFrame的start/chunk/end和落盘供回放这两件事现在是彻底分离的两条路径不再像早期版本那样共用同一条日志。背景与设计动机流式输出streaming本身不难——大多数 LLM SDK 都提供边生成边吐的能力。真正难的是把这件事嵌进一个需要完整历史可回放、可持久化、可在多个消费者之间转发的系统里会同时冒出几个互相牵制的需求协议要统一DeepSeek、pi-ai等不同 provider 的原始流格式各不相同SSE 字段名、分块粒度都不一样Agent 循环不该关心这些差异。既要流式展示、又要有一个最终定型的消息UI 需要边到边渲染的原始 delta但会话历史第二篇讲的 Surface只应该记一条组装完整的assistant/message不能把中间态污染进正式历史。可靠重放优先于存储效率调试一次异常输出比如模型在某个 token 处莫名其妙换了语言或者一次奇怪的工具调用参数是怎么被拼出来的最有效的手段是能看到逐 token 的原始流而不是只有组装完的最终结果。多端消费同一份流式输出既要喂给持久化的会话日志又要喂给可能同时打开的多个浏览器标签/多个客户端。dsh的解法是让每一层都做且只做自己该做的事Provider 只管协议转换LlmRuntime只管中间件编排ReactLoopAgent通过AssistantStreamAttempt同时负责实时转发给 UI和落盘一份可无损重放的紧凑记录这两件事但两条路径彼此独立见下Host/Client 之间的传输层负责把数据搬到浏览器Client 只管把收到的数据再折叠成可渲染的 UI 状态。下面按这个顺序展开。核心机制详解第一层 Provider → Adapter把 SSE 字节流转换成 StreamChunk2026-09 更新packages/llm/llm-deepseek的目录形态又变了一次——之前按协议形态拆出来的src/protocols/chat-completions/、src/protocols/messages/两层目录已经不存在包重新回到了扁平结构而且 DeepSeek 侧现在只保留一套Messages 协议实现Anthropic 风格的message_start/content_block_delta/message_stop事件流请求发往${baseURL}/messages。DeepSeekAdapter.stream()依然只是一层转发stream(options) { return this.generate(options, this.dependencies.options()) }真正的流式逻辑在私有的generate()里。packages/llm/llm-deepseek/src/sse.ts里的parseSse()负责把裸的 SSE 字节流解析成 DeepSeek Messages 协议的事件对象// packages/llm/llm-deepseek/src/sse.ts export async function* parseSse(body: ReadableStreamBufferSource, activity: () void): AsyncGeneratorRecordstring, unknown { const events body.pipeThrough(new TextDecoderStream()).pipeThrough(new EventSourceParserStream({ onComment: activity })) for await (const frame of events) { activity() let raw: unknown try { raw JSON.parse(frame.data) } catch (_invalidSseJson) { throw new LlmError(DeepSeek Messages SSE contains invalid JSON, MALFORMED_RESPONSE) } const event object(raw) if (typeof event.type ! string || (frame.event ! undefined frame.event ! event.type)) { throw new LlmError(DeepSeek Messages SSE event type mismatch, MALFORMED_RESPONSE) } if (event.type error) throw providerError(event, undefined) yield event } }它仍然把帧重组分块可能在任意字节边界断开甚至断在一个 UTF-8 多字节字符中间完全委托给eventsource-parser这个第三方库但职责比早期版本多了一层协议级校验每一帧解析成 JSON 之后必须携带字符串类型的type字段、且与 SSE 的event:行一致否则抛MALFORMED_RESPONSEtype error的帧直接转成结构化 provider 错误抛出。第二个参数activity是一个心跳回调每收到一帧就调一次专门用来喂下面要讲的空闲看门狗。早期版本那条必须显式收到[DONE]才算流正常结束的规则没有消失而是随协议迁移换了形态现在它体现在translate()packages/llm/llm-deepseek/src/translate.ts的出口处——for await正常跑完都没见到message_stop事件就throw new LlmError(DeepSeek Messages stream ended before message_stop, STREAM_CLOSED)。流看起来正常结束、但其实是网络层面被截断这类隐蔽故障的第一道防线还在只是判据从[DONE]哨兵文本换成了 Messages 协议的终止事件。再往上一层DeepSeekAdapter的generate()packages/llm/llm-deepseek/src/adapter.ts把parseSse()产出的事件流喂给translate()转换成协议无关的StreamChunkyield* translate(parseSse(response.body, activity), options.model)并且套了一层空闲超时看门狗// packages/llm/llm-deepseek/src/adapter.ts节选 private async * generate(options: GenerateOptions, connection: Connection): AsyncGeneratorStreamChunk { const consumer new AbortController() const signal options.signal undefined ? consumer.signal : AbortSignal.any([consumer.signal, options.signal]) using watchdog idleWatchdog(signal, connection.streamIdleTimeoutMs, MESSAGES_IDLE) const iterator this.request(options, connection, watchdog.signal, () { watchdog.pulse() }) try { while (true) { const next await watchdog.next(iterator) if (next.done) return yield next.value } } catch (error) { if (timeoutOf(watchdog.signal, MESSAGES_IDLE) ! undefined) throw new LlmError(DeepSeek Messages stream idle timeout, TIMEOUT, { cause: error }) if (options.signal?.aborted) throw new LlmError(DeepSeek Messages request aborted, ABORTED, { cause: error }) if (error instanceof LlmError) throw error throw new LlmError(DeepSeek Messages transport failed, TRANSPORT, { cause: error }) } finally { consumer.abort() try { await iterator.return(undefined) } catch { /* 请求已结束清理失败不改变结果 */ } } }一个信号同时服务两件事调用方主动取消options.signal和空闲看门狗超时idleWatchdog内部的watchdog.signal两者用AbortSignal.any融合成一个。catch块里的判断顺序也是精心设计的——先判断是不是看门狗超时映射成TIMEOUT再判断是不是调用方主动取消映射成ABORTED最后才是兜底的TRANSPORT——这保证了同一次失败无论真实原因是什么最终抛出的LlmError都带着一个稳定、可被后续重试逻辑第五篇用来做路由判断的code而不是一段自然语言消息。第二层 Adapter → LlmRuntime一个可被中间件插入的 waterfallLlmRuntimepackages/llm/llm/src/index.ts并不直接把 adapter 的流原样吐出去,而是包了一层 Cordis 的waterfall中间件链// packages/llm/llm/src/index.ts private streamWithRegistration( options: GenerateOptions, prepared?: { registration: AdapterRegistration; config: LlmCallConfig }, ): AsyncIterableStreamChunk { return this.ctx.waterfall( this, llm/stream, options, () this.adapterStream(options, prepared), ) }llm/stream这个事件名的类型签名是// packages/llm/llm/src/index.ts llm/stream(this: LlmRuntime, options: GenerateOptions, next: () AsyncIterableStreamChunk): AsyncIterableStreamChunk任何插件都可以监听这个事件,拿到next()代表更内层中间件或者最终 adapter 会返回的流,决定原样转发、包一层新的AsyncIterable再返回甚至完全替换掉。下一篇要讲的session-checkpoint-policy正是挂在这里——它在真正调用next()也就是真正向 provider 发起请求之前先做一次会话落盘flush实现发模型请求前必须先把请求本身持久化的 fail-closed 语义// packages/session/session-checkpoint-policy/src/index.ts function afterCheckpoint(ctx: Context, session: Session, next: () AsyncIterableStreamChunk): AsyncIterableStreamChunk { return (async function* (): AsyncIterableStreamChunk { await ctx.sessions.flush(session) yield* next() })() } ctx.on(llm/stream, (options, next): AsyncIterableStreamChunk { if (options.sessionId undefined) return next() const session ctx.sessions.get(options.sessionId) return session undefined ? next() : afterCheckpoint(ctx, session, next) })adapterStream()waterfall 链条最内层的默认实现本身还负责把adapter 选择失败迭代器构造失败迭代过程中途抛异常这三类完全不同来源的失败统一收敛成一种协议——一个{ type: finish, reason: { kind: error | aborted, failure } }的终止 chunk而不是让异常直接从AsyncGenerator里抛出来打断消费者的for await// packages/llm/llm/src/index.ts节选 function adapterFailureChunk(error: unknown, signal?: AbortSignal): StreamChunk { const failure normalizeLlmFailure(error) return { type: finish, reason: signal?.aborted || failure.code ABORTED ? { kind: aborted, failure } : { kind: error, failure }, } }这意味着ReactLoopAgent.step()消费流的时候永远只需要处理正常的 chunk和一条携带失败信息的finishchunk两种情况,不需要额外套try/catch来兜底 adapter 层面的各种抛异常方式——这也是为下一篇的重试机制铺路重试判断的输入永远是一条结构化的finishchunk不是裸的 JS 异常。第三层 LlmRuntime → Agent LoopAssistantStreamAttempt——实时转发与落盘分离这一节是本篇变化最大的地方。早期版本里step()自己持有一个BlockAssembler每个 chunk 到达时先原样落一条assistant/chunk日志、再喂给组装器——落盘和实时转发是同一件事、同一条数据路径。当前版本把这两件事拆成了两条独立路径都封装进了AssistantStreamAttemptpackages/core/agent-loop/src/assistant-stream.ts。// packages/core/agent-loop/src/assistant-stream.ts节选 push(chunk: StreamChunk): void { const timed this.accumulator.push({ time: Date.now(), chunk }) this.assembler.push(timed.chunk) this.emit({ type: chunk, attemptId: this.attemptId, revision: this.nextRevision(), index: this.index, time: timed.time, chunk: timed.chunk, }) }一个 chunk 到达时AssistantStreamAttempt.push()同时做三件事但它们服务三个完全不同的目的this.accumulator.push(...)——喂给AssistantStreamAccumulatorpackages/llm/llm/src/assistant-stream.ts这是唯一会被落盘的路径但它不是每个 chunk 一条日志而是把连续的同类 delta 就地压缩细节见下一节。this.assembler.push(...)——喂给BlockAssembler内部持有逻辑和早期版本完全一致见下方节选负责把碎片拼成完整的ContentBlock[]供这一次请求成功后组装出最终AssistantMessage。this.emit(...)——发出一个AssistantStreamFrametype: chunk通过dispatch.emit(agent/assistant-stream, { frame })广播出去。这一步完全不落盘是纯粹的进程内实时通知专门服务于UI 需要立刻看到这个 token这一个需求。一次请求结束时step()调用live.settle(assistant/message, () this.session.append(...))——只有这一刻才会真正往会话日志写一条事件而且这条事件携带的是live.streamAssistantStreamAccumulator.snapshot()的结果一份紧凑记录不再是一长串独立的assistant/chunk事件// packages/core/agent-loop/src/assistant-stream.ts节选 settle(eventType: assistant/message | assistant/attempt, append: () SessionSeq): void { let seq: SessionSeq try { seq append() } catch (error: unknown) { this.abandon(); throw error } this.terminal true this.emit({ type: end, attemptId: this.attemptId, revision: this.nextRevision(), index: this.index, outcome: { kind: committed, eventType, seq } }) }这就是为什么本篇开头说实时转发和落盘是两条彻底独立的路径即使 UI 一次帧都没收到比如没有任何浏览器连着落盘这条路径完全不受影响反过来即使日志写入失败abandon()分支之前已经发给 UI 的帧也不会被撤回——UI 展示的是这一刻模型说了什么的事实会话日志记的是这次请求最终、经得起回放验证的结果两者故意不耦合。BlockAssembler本身packages/llm/llm/src/assembler.ts依然是这条管道里真正的折叠算法核心实现和早期版本一致:// packages/llm/llm/src/assembler.ts push(chunk: StreamChunk): void { switch (chunk.type) { case block-start: { if (!this.partials.has(chunk.index)) { this.order.push(chunk.index) this.partials.set(chunk.index, { blockType: chunk.blockType, text: , toolCallArguments: }) } return } case text-delta: case reasoning-delta: { const partial this.ensure(chunk.index, chunk.type text-delta ? text : reasoning) if (partial.block) return // closed by block-end; ignore stragglers partial.text chunk.text return } case tool-call-delta: { const partial this.ensure(chunk.index, tool-call) if (partial.block) return partial.toolCallId chunk.id if (chunk.name) partial.toolCallName chunk.name partial.toolCallArguments chunk.argumentsDelta return } case block-end: { const partial this.ensure(chunk.index, chunk.block.type) if (partial.block) return partial.block chunk.block return } case usage: { this._usage chunk.usage; return } case finish: { this._finish chunk.reason; this._replayState chunk.replayState; return } default: return assertNever(chunk, BlockAssembler.push) } }StreamChunk是一个按index内容块在这条消息里的位置分片的增量协议,BlockAssembler用一个Mapnumber, PartialBlock按index维护每个内容块自己的组装状态——文本类块text/reasoning是简单的字符串累加,工具调用块是调用 id 名字 参数 JSON 片段三个字段各自累加,直到一条block-end到达把这个index钉死成最终的ContentBlock。第一次关闭生效,之后的迟到 delta 被忽略if (partial.block) return是一条专门针对畸形流的防御:如果某个 provider 的实现有 bug,在block-end之后又发来同一个index的 delta,这条防线保证最终组装结果和当时流式展示给用户看到的内容完全一致,不会因为迟到数据而产生UI 上看到的和存进历史的不一样的诡异错位。blocks()还处理了一个边界情况:// packages/llm/llm/src/assembler.ts blocks(): ContentBlock[] { const blocks this.order.map(index this.assemble(this.mustGet(index), index)) return this.finish.kind max-tokens ? blocks.filter(block block.type ! tool-call) : blocks }如果这一步是被输出 token 上限截断的max-tokens),组装结果里的工具调用块会被整体过滤掉——一个被截断的工具调用参数比如一个被截成一半的文件路径字符串)如果被当真执行,后果可能是危险的,所以宁可让模型这一步什么工具都没调用,也不要执行一个残缺的调用。AssistantStreamAccumulator紧凑记录如何做到无损上一节说落盘路径不再是每个 chunk 一条日志但依然要保证无损——这是AssistantStreamAccumulatorpackages/llm/llm/src/assistant-stream.ts真正解决的问题。它的核心思路是打包连续同类 delta// packages/llm/llm/src/assistant-stream.ts节选text-delta 分支 case text-delta: case reasoning-delta: { const type chunk.type text-delta ? text-chunks : reasoning-chunks const gap previous ! undefined previous.type type ? safeGap(previous.lastTime, time) : undefined if (previous ! undefined previous.type type previous.index chunk.index gap ! undefined) { previous.dt.push(gap) // 和上一条记录同类型、同 index追加时间差和文本不新开一条记录 previous.texts.push(chunk.text) previous.lastTime time } else { this.records.push({ type, time0: time, index: chunk.index, dt: [], texts: [chunk.text], lastTime: time }) } return timed }如果连续几十个text-delta都属于同一个内容块index相同它们不会变成几十条独立记录而是被压进一条text-chunks记录里time0第一条的时间戳dt后续每条相对上一条的时间差数组texts每条的文本片段数组。这份紧凑记录不是有损摘要——expandAssistantStream()能把它精确地展开回原始的、带时间戳的逐 chunk 序列time0累加每个dt就是每条原始 chunk 的真实到达时间,assembleAssistantStream()也能把它重新喂回一个BlockAssembler得到和当初完全一样的组装结果。只有block-start/block-end/usage/finish这类每种最多出现一次或语义上不该合并的 chunk 才会被原样存成一条{ type: chunk }记录。这是比早期版本更好的答案早期版本用存储换保真度每个 chunk 一条日志体量大但简单现在用一个不复杂的打包算法同时拿到了两者——日志体量大幅下降连续的同类 delta 从 N 条压成 1 条但expandAssistantStream()保证这依然是一份可以逐 token 精确重放的记录,不是有损压缩。第四层与第五层Host/Client 传输与再折叠架构已重组细节见第 06 章早期版本里Host 端用packages/host/apiproxy/src/api-proxy.ts里手写的FrameQueue/events.mux把session/event重新打包成 WebSocket 帧、packages/client/connection负责在 Node 侧转发、packages/client/runtime里的PartialAccumulator在浏览器侧再折叠一次。当前版本里packages/host/下已经不存在apiproxy这个包取而代之的是packages/api/session-controller这个新包下辖remote-events.ts/commands.ts/control.ts等是命令下发、远程事件订阅、断线重连快照的统一入口客户端侧的折叠器也从packages/client/runtime/src/client/sessions/partial.ts挪到了packages/client/ui-chat/src/client/conversation-nodes/partial.ts。这是一次真正的包重组不是简单改名——完整、准确的 Host-Client 传输架构RPC 契约怎么生成、连接怎么建立、重连快照怎么工作已经超出本篇流式管道的范围留给第 06 章《Host-Client 分离与 Typert RPC 生成》专门讲解。这里只需要记住一个不变的架构原则折叠这件事在 Client 侧依然是独立于服务端重新做一遍的浏览器不会盲目信任服务端已经算好的中间结果而是拿到原始的流式数据/事件后自己重新组装出可渲染的状态——这个每一层只信任自己收到的原始输入的设计原则本身没有变变的只是搬运这些数据的具体传输层代码。为什么现在只落盘紧凑记录而不是每个原始 chunk把整条链路串起来看,这个问题的答案是:回放保真度没有被牺牲:紧凑记录AssistantStreamRecord[]依然完整保留了模型输出到底是怎么一小块一小块吐出来的这一事实——expandAssistantStream()可以精确重建出原始的逐 chunk 序列包括每条的到达时间。这不是取舍而是同一份保真度用更省空间的编码方式存下来。崩溃恢复的精确性不受影响:即使进程在流式响应过程中崩溃AssistantStreamAccumulator内部维护的记录列表本身就是增量构建的只是最终settle()那一刻才整体落盘——如果崩溃发生在settle()之前这一次 attempt 本来就不会被认为已经完成,和早期版本部分 chunk 已落盘、但没有 assistant/message的中间状态相比语义更干净要么完整落盘要么这次 attempt 视为没发生。实时性和持久化解耦:UI 需要的立刻看到这个 token通过完全独立的AssistantStreamFrame广播满足不再依赖chunk 必须先落盘才能被转发这个顺序——这也是为什么 Host/Client 重连后想要恢复刚才漏掉的片段现在要靠专门的重连快照机制第 06 章而不是简单地重放一段assistant/chunk日志。早期版本的代价是显而易见的:一次几百 token 的回复可能对应几十上百条assistant/chunk事件,日志体量会明显膨胀。当前版本用一个不复杂的打包连续同类 delta算法解决了这个代价同时没有放弃逐 token 级别可精确重放这个目标——这是一个值得在自己的系统设计中借鉴的思路:当简单但浪费的方案跑了一段时间之后往往能找到一个不牺牲原有保证、但明显更省资源的编码方式。小结流式管道核心分三层Provider 协议转换parseSsetranslate当前 DeepSeek 侧只保留一套 Messages 协议实现包回到扁平目录→LlmRuntime的llm/streamwaterfall中间件可介入,如 checkpoint→ Agent Loop 的AssistantStreamAttempt把实时转发给 UIAssistantStreamFrame不落盘和落盘一份可无损重放的紧凑记录AssistantStreamAccumulator打包连续同类 delta彻底拆成两条独立路径。Host/Client 之间具体怎么传输、怎么重连恢复架构已经重组详见第 06 章。BlockAssembler按内容块index维护增量组装状态,block-end到达即钉死,之后的迟到 delta 被忽略;max-tokens截断时会把不完整的工具调用块整体过滤掉。落盘的紧凑记录AssistantStreamRecord[]通过打包连续同类 delta 大幅降低日志体量同时用expandAssistantStream()/assembleAssistantStream()保证依然可以无损还原出原始的逐 chunk 时间序列——不是用存储换保真度的取舍而是同一份保真度换了一种更省空间的编码方式。