青雲的博客
深入浅出 DeepSeek Harness 第二部:一句话的旅程——输入怎样变成模型请求 第 11 章

流式响应:chunk 怎样变成消息

追踪一个 LLM 流式 chunk 从 adapter 发出到最终成为 assistant/message 事件的完整路径。涵盖 BlockAssembler 增量组装算法、七种 chunk 协议、delta-only 兼容、straggler 容错、max-tokens 截断过滤、finish 默认值,以及 sourceEventSeqs provenance 链的建立过程。

源码版本
47f943859bef60e4160492346772ded9b24f765a
验证日期
Commit
47f943859bef60e4160492346772ded9b24f765a

你可能以为模型把所有内容都吐完了,系统才统一处理、一次性拼成一条消息存入日志。

不是这样。流式响应中,每个 chunk 一到就立刻写入事件日志——一个 token 对应一条 assistant/chunk 事件。同时,BlockAssembler 在内存中增量维护组装状态。你屏幕上看到文字逐字蹦出来,不是等完后才渲染,而是每个 chunk 都触发一次状态更新。如果流在中间断了,已到达的 chunk 都在日志里,不会丢失——下次 replay 可以完整重建。

这一章就顺着一个 StreamChunk 走到底:它从 adapter 产出,先被写成事件,再被 BlockAssembler 拼成 assistant/message。顺着这条路看下去,你会明白为什么 Harness 把“先持久化、再组装”当成硬顺序,也会看到 BlockAssembler 怎么处理 delta-only 协议、straggler delta、max-tokens 截断这些边角情况。


全景:从 stream 到 message 的数据流

在深入细节之前,先看整条路径的全景:

adapter.stream(request)
  → AsyncIterable<StreamChunk>
    → for-await 循环逐个消费
      → session.append('assistant/chunk', {turn, step, chunk}) [durable-first]
      → assembler.push(chunk) [in-memory state update]
    → stream 结束
  → assembler.finish 检查结束原因
    → error/aborted → agent/request-error waterfall(可能重试)
    → 正常 → assembler.blocks() 组装 ContentBlock[]
      → createAssistantMessage({content, source})
      → session.append('assistant/message', {...}, {surfaceOp:'append', sourceEventSeqs: chunkSeqs})
        → 有 tool-call blocks → executeToolCalls()
        → 无 tool-call → step completed

关键点:chunk 持久化和 assembler 状态更新是流消费循环的两步原子操作。chunk 丢失等于消息不完整,所以持久化在前。


StreamChunk 协议:七种 chunk 类型

LLM adapter 向 agent loop 产出的每个 chunk 都是一个 StreamChunk 联合类型。协议设计支持 interleaved blocks——一个流中可以同时推进多个块(比如 text 和 tool-call 交织),通过 index 字段关联同一个块的多个 delta。

七种类型:

类型用途关键字段
block-start声明一个新块开始index, blockType
text-delta文本块增量内容index, text
reasoning-delta推理块增量内容index, text
tool-call-delta工具调用块增量index, id, name?, argumentsDelta
block-end闭合一个块index, block: ContentBlock
usagetoken 用量统计usage: TokenUsage
finish流结束信号reason: FinishReason, replayState?

注意 finishreason 字段是 FinishReason 类型——一个 merge-extensible 的联合:

  • {kind: 'stop'} — 正常完成
  • {kind: 'tool-calls'} — 模型要调用工具
  • {kind: 'max-tokens'} — 输出被 token 上限截断
  • {kind: 'error', failure} — provider 返回错误
  • {kind: 'aborted', failure} — 被取消

adapter 保证在正常情况下 usagefinish 之前到达。但 assembler 不假设这个顺序——它允许 usage 在任何时候到达。


agent loop 的流消费循环

BlockAssembler 的使用点在 ReactLoopAgent.step() 方法内。看核心循环:

const assembler = new BlockAssembler()
const chunkSeqs: number[] = []
const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)
signal.throwIfAborted()
for await (const chunk of stream) {
  signal.throwIfAborted()
  chunkSeqs.push(this.session.append('assistant/chunk', { turn, step, chunk }).seq)
  assembler.push(chunk)
}

三个设计决策在这几行里

  1. durable-first:先 session.append 拿到 seq,再 push 进 assembler。如果进程在两步之间崩溃,日志里有这个 chunk 但 assembler 没更新——replay 时从 chunk 事件重建即可。绝不会出现 assembler 有状态但日志没记录的情况。

  2. abort 检查在 for-await 顶部:每处理一个 chunk 前先检查 signal,保证取消信号能及时中断流消费。

  3. chunkSeqs 收集所有 seq:包括内容 delta、usage、finish——所有 chunk 不分类型一律记录。最终 assistant/message 事件的 sourceEventSeqs 就是这个数组,形成完整的 provenance 链。


BlockAssembler:增量组装状态机

BlockAssembler 是一个 确定性有限状态机,它唯一的输入是按流顺序到达的 StreamChunk,唯一的输出是通过 blocks() / finish / usage / replayState 等 getter 暴露的最终组装结果。

内部数据结构

class BlockAssembler {
  private partials = new Map<number, PartialBlock>()  // index → 块状态
  private order: number[] = []                         // 记录块出现顺序
  private _usage: TokenUsage | undefined
  private _finish: FinishReason | undefined
  private _replayState: unknown = undefined
}

PartialBlock 是每个块的增量状态:

interface PartialBlock {
  blockType: string          // 'text' | 'reasoning' | 'tool-call'
  text: string               // text-delta / reasoning-delta 累加
  toolCallId?: CallId        // tool-call-delta 的 id
  toolCallName?: string      // tool-call-delta 的 name
  toolCallArguments: string  // tool-call-delta 的 argumentsDelta 累加
  block?: ContentBlock       // block-end 设置,一旦设置该 partial 被冻结
}

关键不变量:order 数组中的每个 index 在 partials Map 中一定有对应条目mustGet() 方法强制这个不变量——违反则立刻抛错。

push() 方法:七路分发

push(chunk) 是 assembler 唯一的写入接口。一个 switch 根据 chunk.type 分发到七条处理路径:

block-start:如果该 index 还没有 partial,在 order 中记录顺序,在 partials 中创建新的 PartialBlock(blockType 来自 chunk,text 和 toolCallArguments 初始化为空字符串)。如果已有 partial(重复的 block-start),静默忽略——不重复创建,不抛错。

text-delta / reasoning-delta:调用 ensure(index, blockType) 找到或创建 partial。如果 partial.block 已设置(被 block-end 冻结),直接 return——straggler 忽略。否则 partial.text += chunk.text

tool-call-delta:同样 ensure + 检查冻结。设置 toolCallId、更新 toolCallName(有就覆盖)、追加 toolCallArguments += chunk.argumentsDelta

block-end:ensure 确保 partial 存在,然后 partial.block = chunk.block。如果已经被关闭过,忽略重复 close(first close wins)。

usage:直接 this._usage = chunk.usage。后到达的 usage 覆盖先前的。

finishthis._finish = chunk.reasonthis._replayState = chunk.replayState

defaultassertNever(chunk) 编译期穷尽检查,运行时抛错。


delta-only 协议兼容

不是所有 LLM provider 都发送 block-startblock-end。很多 OpenAI 兼容 API 只发 content delta 和 tool_calls delta,没有显式块边界。BlockAssembler 通过 ensure() 方法兼容这种协议:

private ensure(index: number, blockType: string): PartialBlock {
  let partial = this.partials.get(index)
  if (!partial) {
    partial = { blockType, text: '', toolCallArguments: '' }
    this.partials.set(index, partial)
    this.order.push(index)
  }
  return partial
}

text-delta 到达但该 index 还没有 partial 时,ensure 自动创建一个——blockType 由 delta 类型推断(text-delta 推断为 'text'reasoning-delta 推断为 'reasoning'tool-call-delta 推断为 'tool-call')。order 数组也记录该 index。

这意味着:delta-only 协议不需要 block-start 也能工作。assembler 看到第一个 delta 就知道开了个新块。同样,流可以不发 block-end 就结束——assemble() 方法会从累计的 delta 数据组装块(而不是要求 partial.block 必须存在)。


straggler delta 容错

这是一个专门为 misbehaving adapter 设计的防御机制。如果 adapter 在发送 block-end 之后又发了同一个 index 的 delta(malformed stream),这些 delta 会被 静默忽略

检查点在 text-deltatool-call-delta 的处理中:

case 'text-delta':
case 'reasoning-delta': {
  const partial = this.ensure(chunk.index, ...)
  if (partial.block) return  // ← straggler 忽略
  partial.text += chunk.text
  return
}

partial.blockblock-end 设置。一旦设置,该块就是 “frozen” 状态。后续 delta 不追加、不抛错——静默 return。

为什么不抛错?源码注释说得清楚:

“deltas arriving for an index already closed by block-end are ignored (malformed stream) so a misbehaving adapter cannot grow memory or corrupt a completed block.”

两个安全保证:

  1. 内存安全:不能通过无限发 straggler delta 让 partial.text 无限增长。
  2. 数据一致性:streamed output 到 block-end 为止就是最终内容。block-end 的 chunk.block 是权威的 ContentBlock,不会被后续 delta 修改。

注意 block-end 自身也有 first-close-wins 语义:如果同一个 index 收到多个 block-end,只有第一个生效。


blocks():最终组装

流消费结束后,调用 assembler.blocks() 获取 ContentBlock[]

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
}

两步:

  1. 按 order 顺序组装:对每个 index,如果 partial.block 存在(被 block-end 冻结),直接用它;否则从累计的 delta 数据构造 ContentBlock。

  2. max-tokens 过滤:如果 finish reason 是 max-tokens(输出被 token 上限截断),过滤掉所有 tool-call 类型的块

assemble() 的降级组装

对于没有 block-end 的块(delta-only 协议或流提前结束),assemble() 从 partial 的累计数据构造:

private assemble(partial: PartialBlock, index: number): ContentBlock {
  if (partial.block) return partial.block  // 已关闭:直接用权威块
  switch (partial.blockType) {
    case 'text': return { type: 'text', text: partial.text }
    case 'reasoning': return { type: 'reasoning', text: partial.text }
    case 'tool-call': return {
      type: 'tool-call',
      id: partial.toolCallId ?? CallId(`call-${index}`),
      name: partial.toolCallName ?? '',
      arguments: partial.toolCallArguments,
    }
    default: throw new Error(...)
  }
}

注意 tool-call 的 id 有兜底:如果 delta 从未设置过 toolCallId,用 call-${index} 生成一个。name 缺失则为空字符串。但未知的 blockType 会抛错——这是不应该发生的 invariant 违反。

为什么 max-tokens 要丢弃 tool-call

当输出被 token 上限截断时,正在输出的 tool-call 几乎一定是不完整的——它的 arguments 字段是一段截断的 JSON。把这样的 tool-call 交给工具执行,JSON.parse(arguments) 会失败(或者更糟:parse 成功但语义不完整)。

所以 blocks() 在 max-tokens 时直接过滤掉 所有 tool-call 块——不只是最后一个可能被截断的,而是全部。为什么?因为一旦输出被截断,你无法确定模型的意图是否完整。可能模型打算输出三个 tool-call,但第三个截断了;如果只执行前两个,语义可能不正确。

text 和 reasoning 块保留——即使被截断,文本内容可能仍有参考价值。


finish getter:没有 finish chunk 的默认行为

正常的流以 finish chunk 结尾,告诉 assembler 结束原因。但如果流结束了却没收到 finish chunk(adapter bug、网络异常导致流提前关闭但没报错)呢?

get finish(): FinishReason {
  return this._finish ?? { kind: 'stop' }
}

默认是 {kind: 'stop'}——当作正常完成。这是容错设计:比起因为缺少 finish chunk 就抛错丢掉整轮响应,把它当作正常 stop 更安全。

但代价是:如果你依赖 finish 里的 replayState(adapter 私有的续传状态),没有 finish chunk 就没有 replayState。adapter 在下次续传时可能需要重新获取。


流结束后:从 assembler 到 assistant/message 事件

流消费循环结束后,agent.ts 处理 finish reason 决定后续路径:

const finish = assembler.finish
if (finish.kind === 'error' || finish.kind === 'aborted') {
  // → agent/request-error waterfall:让插件决定是否重试
  const action = await this.dispatch.waterfall('agent/request-error', {...})
  if (action?.kind !== 'retry') {
    throw new LlmError(finish.failure.message, finish.failure.code, finish.failure)
  }
  continue  // retry: 回到 while(true) 重新 buildRequest + stream
}

重试路径:如果 finish 是 error 或 aborted,不立刻组装消息,而是通过 agent/request-error waterfall 让插件(如重试插件)决定下一步。如果插件返回 {kind: 'retry'},整个 step 从 buildRequest 开始重来(新的 assembler、新的 chunkSeqs)。

正常路径(stop / tool-calls / max-tokens):

const message = createAssistantMessage({
  content: assembler.blocks(),
  source: {
    provider: request.provider,
    model: request.model,
    ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {},
  },
})
this.session.append(
  'assistant/message',
  {
    turn, step, message,
    ...assembler.usage === undefined ? {} : { usage: assembler.usage },
  },
  { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
)

sourceEventSeqs:provenance 链

sourceEventSeqs: chunkSeqs 让你从一条 assistant/message 追溯到每一个原始 chunk(包括 delta、usage、finish)。surface layer 通过 surfaceOp: 'append' 知道展示消息,通过 sourceEventSeqs 关联之前的 chunk 事件(支持 replay、审计、增量 UI 合并)。

createAssistantMessagesource 字段记录 provider、model 和可选的 replayState(adapter 私有,如 DeepSeek adapter 的对话 ID,用于续传)。


后续分支:finish reason 决定 step 走向

组装并 append assistant/message 之后,根据 finish reason 和消息内容决定 step 是否结束:

if (finish.kind === 'max-tokens') return { kind: 'max-tokens' }

const toolCalls = message.content.filter(block => block.type === 'tool-call')
if (toolCalls.length === 0) return { kind: 'completed' }
const { concluded } = await executeToolCalls(...)
return concluded ? { kind: 'completed' } : null

三条路径:

  1. max-tokens:直接返回 {kind: 'max-tokens'}。上层 turn 循环会标记这个 turn 的 TurnEndReason 为 max-tokens(sticky——一旦某个 step 触发 max-tokens,后续 step 正常完成也不会降级这个 turn outcome)。

  2. 无 tool-call:step completed,控制权回到 turn 循环。

  3. 有 tool-call:调用 executeToolCalls() 执行所有工具。工具执行完后,如果任何工具结果标记 concludesTurn,返回 completed;否则返回 null,turn 循环会继续进入下一个 step(等 inbox 里有工具结果作为 context)。


完整时序:一个典型的多块流

把所有部分串起来,看一个模型返回 text + tool-call 的典型场景:

sequenceDiagram
    participant A as Adapter
    participant L as Agent Loop
    participant S as Session Log
    participant B as BlockAssembler

    Note over L: new BlockAssembler(), chunkSeqs=[]
    A->>L: chunk{block-start, index:0, blockType:text}
    L->>S: append('assistant/chunk') → seq 101
    L->>B: push → create partial[0]
    A->>L: chunk{text-delta, index:0, text:"让我帮你查一下"}
    L->>S: append('assistant/chunk') → seq 102
    L->>B: push → partial[0].text += ...
    A->>L: chunk{block-end, index:0, block:{type:text,...}}
    L->>S: append('assistant/chunk') → seq 103
    L->>B: push → partial[0].block = frozen
    A->>L: chunk{block-start, index:1, blockType:tool-call}
    L->>S: append('assistant/chunk') → seq 104
    L->>B: push → create partial[1]
    A->>L: chunk{tool-call-delta, index:1, id:"call_abc", name:"search", args:'{"q":"天气"}'}
    L->>S: append('assistant/chunk') → seq 105
    L->>B: push → partial[1] updated
    A->>L: chunk{block-end, index:1, block:{type:tool-call,...}}
    L->>S: append('assistant/chunk') → seq 106
    L->>B: push → partial[1].block = frozen
    A->>L: chunk{usage, usage:{inputTokens:50, outputTokens:20}}
    L->>S: append('assistant/chunk') → seq 107
    L->>B: push → _usage set
    A->>L: chunk{finish, reason:{kind:tool-calls}}
    L->>S: append('assistant/chunk') → seq 108
    L->>B: push → _finish set
    Note over L: stream ends
    L->>L: blocks() → [text, tool-call]
    L->>S: append('assistant/message', {sourceEventSeqs:[101..108]})
    L->>L: executeToolCalls([tool-call-block])

错误和重试路径

如果 adapter 在流中间遇到错误(比如 HTTP 500),它会产出一个 {type: 'finish', reason: {kind: 'error', failure: {...}}} chunk(或者直接 throw,由 LlmRuntime.stream() 包装为 error/aborted finish)。

agent loop 不会为错误的流组装消息。它走一条不同的路:

  1. chunk 仍然被持久化为 assistant/chunk 事件(包括 error finish chunk)
  2. assembler.finish 返回 error/aborted
  3. 进入 agent/request-error waterfall——插件链决定:
    • 返回 {kind: 'retry'} → 循环 continue,重新 buildRequest + 新 stream
    • 返回 undefined → throw LlmError,step 失败

重试时,新的 BlockAssembler 和 chunkSeqs 在 while(true) 循环的下一轮迭代中重新创建。失败的那批 chunk 仍然留在日志中(provenance 完整),但不会影响新一次 stream 的组装。


容易踩的坑

坑一:以为 assembler 等流结束才处理

BlockAssembler 是增量的——每个 push() 立刻更新状态。你在流进行中的任意时刻调用 blocks() 都能拿到当前已组装的块。这不是”收集完所有 chunk 再批量处理”的模式。

坑二:以为 block-start 必须存在

delta-only 协议下,ensure() 在第一个 delta 到达时自动创建 partial。如果你写 adapter 时忘了发 block-start,assembler 仍然能工作——只是 blockType 由 delta 类型推断而非显式声明。

坑三:max-tokens 后以为 tool-call 会被执行

不会。blocks()finish.kind === 'max-tokens' 时过滤掉所有 tool-call。后续 message.content.filter(block => block.type === 'tool-call') 拿到空数组,不会走 executeToolCalls 路径。

坑四:以为 straggler delta 会追加到已关闭块

partial.block 存在后,后续 delta 全部忽略。不报错、不追加、不修改。block-end 的 chunk.block 是该块的权威最终值。

坑五:以为 chunk 不写入日志只存内存

每个 chunk 都作为 assistant/chunk 事件持久化。这是 replay 保真的基础——你可以从 chunk 事件流完全重建流式输出过程,包括 timing(事件有 seq 序号)。

坑六:以为 sourceEventSeqs 只含内容 delta

chunkSeqs 收集了 for-await 循环中 所有 chunk 的 seq——包括 block-start、block-end、usage、finish。assistant/message 的 provenance 链指向完整的原始流。


收口:几条不能破的顺序

原则体现
Durable-firstchunk 先 append 事件,再 push assembler
协议兼容ensure() 让 delta-only 和 full-protocol 都能工作
防御式容错straggler 忽略、first-close-wins、no-finish 默认 stop
安全截断max-tokens 过滤不完整 tool-call
完整溯源sourceEventSeqs 建立 message → chunks 引用链
重试隔离error/aborted 不组装消息,retry 创建全新 assembler

消息组装好了。如果有 tool-call,下一步是工具执行;如果没有或 finish 是 max-tokens,step 结束,控制权回到 turn 循环。