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

LLM Service、流式 chunk 与 BlockAssembler

从请求就绪到模型响应完整组装为事件的全流程追踪。LlmRuntime 是适配器注册表而非 API 客户端;prepareCall 锁定 adapter registration 防止 HMR 竞态;markAgentLoopRequest 深冻结请求保证可重建;llm/stream waterfall 提供 pre/around/post 拦截点;BlockAssembler 增量组装 chunk 为 ContentBlock;每个 chunk 成为 session log 中的 assistant/chunk 事件。

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

触发时刻:请求就绪,下一步是什么

你在上一章看到 agent loop 的 buildRequest 把 system prompt、messages、tools、config 全部组装完成。现在请求对象就绪了——一个 frozen 的 GenerateOptions,带着 markAgentLoopRequest 的身份标记。接下来发生什么?

答案不是”直接 fetch 模型 API”。请求会穿过三个关键层:先是 prepareCall 锁定 adapter registration 并解析模型能力,然后进入 llm/stream waterfall 接受插件拦截,最终到达 adapter 的 stream() 方法产出原始 chunk。这些 chunk 逐个喂给 BlockAssembler,同时每个 chunk 都作为 assistant/chunk 事件写入 session log。当流结束,assembler 组装出完整的 ContentBlock 数组和 Message,以 assistant/message 事件提交。

我们一跳一跳追踪这条路径。

第一跳:prepareCall 与 adapter 锁定

LlmRuntime 是注册表,不是客户端

先建立一个关键认知。看 LlmRuntime 的构造函数和私有字段:

export class LlmRuntime extends Service {
  private adapters = new Map<string, AdapterRegistration>()
  private directory = new Map<string, LlmConfigurableProvider>()
  private discoveries = new Map<string, ...>()

  constructor(ctx: Context) {
    super(ctx, 'llm')
  }
}

三个 Map,零个 HTTP 客户端。LlmRuntime 的职责是:注册 adapter、路由请求、提供 waterfall 拦截点。所有真正的 HTTP 调用都发生在具体的 LlmAdapter 实现内部(比如 DeepSeek adapter 或 pi-ai adapter)。

Agent loop 调用 prepareCall

agent.tsbuildRequest 方法中,loop 拿到 waterfall 解析后的 proposedConfig,然后:

preparedCall = await this.loopCtx.llm.prepareCall(proposedConfig, signal)
config = preparedCall.config

prepareCall 是这条链路真正收口的地方。它做五件事:

第一步:查找当前 adapter registration。 通过 this.registration(config.provider)adapters Map 中取出对应的 AdapterRegistration 对象(包含 adapter 实例、provider 元数据、retryPolicy)。如果没注册,立即抛 NO_ADAPTER

第二步:解析模型能力。 调用 resolveCallFor(registration, config, signal),内部通过 adapter.resolveModel() 获取 contextWindow、reasoning efforts 列表、defaultMaxTokens。验证 requested reasoningEffort 是否在 adapter 声明的 efforts 列表中——不支持就抛 UNSUPPORTED_REASONING_EFFORT,不做静默降级。

第三步:Materialize adapter defaults。 如果用户没传 maxTokens 但模型有 defaultMaxTokens,填上它。如果用户没传 reasoningEffort 但模型有 defaultEffort,填上它。同时记录 adapterDefaults 标记——哪些字段是 adapter 填的而不是用户传的。

第四步:深冻结。 对 resolved config 和 context 做 deepFreeze(structuredClone(...)):先 structuredClone 脱离原始对象的引用,再递归冻结使其不可修改。deepFreeze 的实现是迭代式的(非递归调用),用 WeakSet 防环,并特意跳过 AbortSignal(因为 signal 是活的取消通道,冻结它会破坏 abort 功能)。

第五步:构造 one-shot PreparedLlmCall。 返回一个 Object.freeze() 的对象,其中 stream(options) 方法通过闭包捕获了 registration

为什么需要锁定 adapter

考虑 HMR(热模块替换)场景:

  1. Loop 调用 prepareCall({provider:'deepseek', model:'chat'}),解析出 defaultMaxTokens: 8192。
  2. 此时 adapter 插件热更新,新版本的 resolveModel 返回 defaultMaxTokens: 4096。
  3. Loop 调用 preparedCall.stream(request) 发请求——request 里带着 maxTokens: 8192。

如果 stream 重新从 Map 查找 adapter,拿到的新 adapter 可能拒绝 8192 的 maxTokens。但 PreparedLlmCall 不重新查 Map——它用闭包捕获的 registration 对象直接调 streamWithRegistration。旧 adapter 仍然在内存中(因为有闭包引用),它知道自己返回过 8192,所以请求一致性得到保证。

同时 dispatched 标志位保证了一次性语义:每次 step 都要重新 prepareCall,因为 step 之间模型配置可能经过 agent/request waterfall 变化了。

markAgentLoopRequest:请求的可重建性保证

buildRequest 的最后一步:

const request = markAgentLoopRequest(deepFreeze({
  ...header.config,
  messages: boundaryMessages,
  ...header.system !== undefined ? { system: header.system } : {},
  ...header.tools !== undefined ? { tools: header.tools } : {},
  sessionId: this.session.id,
  signal,
}))

markAgentLoopRequest 把请求对象加入一个 process-local 的 WeakSetAGENT_LOOP_REQUESTS)。isAgentLoopRequest(request) 可以在 waterfall listener 中检测到这个标记。

这个标记传达的信息是:这个请求的内容完全是 session log 的纯函数——给定相同的 session events,可以精确重建出相同的请求。因此 waterfall listener 应该只读它,不尝试修改它。deepFreeze 把这个”应该”变成”不能”——strict mode 下修改 frozen 对象会 throw TypeError。

注意:deepFreeze 的实现在 call-config.ts 中是一个迭代式遍历,用栈模拟 DFS,不依赖 JavaScript 调用栈深度。它显式跳过 AbortSignal 实例,因为 signal 是请求的活通道。

第二跳:stream waterfall 与 chunk 处理

waterfall 的入口

请求准备好了,现在要实际调模型。在 step() 方法中:

const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)

如果有 preparedCall(正常情况),用它的 stream——这会走到 streamWithRegistration,传入捕获的 registration。如果 prepareCall 因为 NO_ADAPTER 失败了(middleware 可能服务一个未注册的路由),退化到 this.loopCtx.llm.stream(request)——这时 waterfall listener 必须提供 chunks,否则 adapterStream 里会再次抛 NO_ADAPTER

streamWithRegistration 的实现只有三行:

private streamWithRegistration(options, prepared?) {
  return this.ctx.waterfall(this, 'llm/stream', options, () => this.adapterStream(options, prepared))
}

ctx.waterfall 是 Cordis 框架提供的 waterfall 事件分发。它的语义是:所有注册了 llm/stream 事件的 listener 形成一个链,每个 listener 接收 (options, next) 参数。Listener 可以:

  • 调用 next() 获取下游的 AsyncIterable<StreamChunk>,包装后返回
  • 不调用 next(),直接返回自己的 AsyncIterable<StreamChunk> 短路整个链
  • 调用 next() 但在产出 chunk 之前/之后插入自己的逻辑

最终的”终端” next 是 () => this.adapterStream(options, prepared)——到达真正的 adapter。

llm/stream 事件声明

waterfall 事件的声明写得很讲究:

interface Events {
  'llm/stream'(
    this: LlmRuntime,
    options: GenerateOptions,
    next: () => AsyncIterable<StreamChunk>
  ): AsyncIterable<StreamChunk>
}

注释中有一段关键说明:“A LOOP-built request carries the process-local markAgentLoopRequest identity and arrives deep-frozen (mutation throws): its content is a pure function of the session log (the reconstructability Agent Note), so listeners read it, never rewrite it. Hand-built calls do not carry that marker; their messages already obey the immutable creation contract.”

这意味着 waterfall listener 面对两种请求:

  • Loop 请求:frozen,有标记,不能修改
  • Hand-built 请求:不 frozen,无标记,但遵循 immutable creation contract(创建后不改)

两种情况下 listener 都不应该修改 options 对象本身。如果你需要”改写请求”,正确做法是创建一个新对象传给下游,而不是 mutate 传入的 options。

adapterStream:adapter 边界的错误隔离

adapterStream 是 waterfall 的终端,也是 adapter 和框架之间的边界。它的设计核心是错误隔离

private async * adapterStream(options, prepared?) {
  let iterator
  try {
    const registration = prepared?.registration ?? this.registration(options.provider)
    // ... resolve config, get adapter, call adapter.stream()
    iterator = stream[Symbol.asyncIterator]()
  } catch (error) {
    yield adapterFailureChunk(error, options.signal)
    return
  }

  try {
    while (true) {
      let item
      try {
        const next = await iterator.next()
        item = next.done ? { done: true } : { done: false, value: next.value }
      } catch (error) {
        yield adapterFailureChunk(error, options.signal)
        return
      }
      if (item.done) return
      yield item.value  // Consumer/middleware failures remain thrown
    }
  } finally {
    if (!completed) await iterator.return?.()
  }
}

三层错误处理:

  1. Adapter 选择/创建阶段(try 块 1):如果 registration 查找失败或 adapter.stream() 构造失败,转为 terminal failure chunk。
  2. Iterator 迭代阶段(内部 try):如果 iterator.next() throw(adapter 内部网络错误等),转为 terminal failure chunk。
  3. Yield 之后:consumer 消费 chunk 时如果 throw,或者 middleware 在包装时 throw,这些错误不被捕获——它们作为正常异常向上传播。

adapterFailureChunk 把任何错误规范化为:

{ type: 'finish', reason: signal?.aborted ? { kind: 'aborted', failure } : { kind: 'error', failure } }

StreamChunk 协议

adapter 产出的 chunk 遵循 StreamChunk 联合类型:

  • block-start:声明一个新块(文本/推理/工具调用),携带 indexblockType
  • text-delta / reasoning-delta:文本增量,关联到 index
  • tool-call-delta:工具调用参数增量,携带 id、可选 nameargumentsDelta
  • block-end:块完成,携带组装好的完整 ContentBlock
  • usage:token 用量报告
  • finish:终端 chunk,携带 FinishReason 和可选的 replayState

协议规则:“Adapters emit usage before the terminal finish and nothing afterward.” 即 usage 在 finish 之前,finish 之后不应再有 chunk。

Agent loop 中的 chunk 处理循环

回到 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)
}

每个 chunk 做两件事:

  1. 写入 session log 作为 assistant/chunk 事件(保存原始 chunk 以支持 replay fidelity),记录返回的序列号
  2. 喂给 BlockAssembler 做增量组装

注意 signal.throwIfAborted() 出现在循环头部和每次迭代之后——这保证了 cancel 信号能及时中断流式处理,不会等到下一个 chunk 到达才退出。

第三跳:BlockAssembler 与事件提交

BlockAssembler 的增量组装算法

BlockAssembler 是”the single canonical assembly algorithm”——唯一的规范组装实现。它内部维护:

  • partials: Map<number, PartialBlock>——按 index 跟踪每个正在构建的块
  • order: number[]——块的出现顺序
  • _usage / _finish / _replayState——终端状态

push(chunk) 方法是一个大的 switch:

block-start:如果 index 未见过,创建一个新的 PartialBlock 并记录顺序。如果已见过(重复 block-start),忽略。

text-delta / reasoning-delta:调用 ensure(index, blockType) 获取或创建 partial。如果 partial.block 已存在(block-end 已到达),忽略这个 straggler delta——“malformed stream cannot grow memory or corrupt a completed block”。否则追加 partial.text += chunk.text

tool-call-delta:类似,追加 partial.toolCallArguments += chunk.argumentsDelta,更新 id 和 name。

block-end:设置 partial.block = chunk.block。“First close wins”——后续重复的 block-end 被忽略。

usage:记录 this._usage = chunk.usage

finish:记录 this._finish = chunk.reasonthis._replayState = chunk.replayState

容忍性设计

注意 ensure 方法:如果收到一个 delta 但之前没有 block-start,它会自动创建 partial。这意味着 adapter 可以省略 block-start 直接发 delta——“Tolerant of delta-only protocols (no block-start/end)“。这是对接不同提供商协议差异的容忍性设计。

同样,如果 block-end 之后又来了同 index 的 delta,不会追加到已完成的块上。这防止了一个行为异常的 adapter 无限增长内存或破坏已完成的数据。

blocks() 方法和 max-tokens 截断

当流结束后,调用 assembler.blocks() 获取组装结果:

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
}

如果 finish reason 是 max-tokens(模型因 token 上限截断),tool-call 类型的块会被过滤掉。原因:截断的 tool call 参数是不完整的 JSON,执行它们不安全。文本块可以截断显示,但 tool call 不行。

事件提交:从 chunks 到 message

流结束后,agent loop 检查 finish reason:

const finish = assembler.finish
if (finish.kind === 'error' || finish.kind === 'aborted') {
  // 走错误处理/重试逻辑
}

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 },
)

关键点:

  1. message source 携带 replayState——adapter 私有的不透明数据,用于在下一次请求中实现高保真重放(不需要重新解析完整响应)。
  2. surfaceOp: ‘append’——这个消息是面向用户的追加操作,UI 层会渲染它。
  3. sourceEventSeqs: chunkSeqs——assistant/message 事件通过序列号数组引用了所有产出它的 assistant/chunk 事件。这建立了一个溯源链:任何时候你都能从最终消息回溯到它的原始 chunk 流。

错误重试路径

如果 assembler.finisherroraborted,流程进入 agent/request-error waterfall:

const action = await this.dispatch.waterfall(
  'agent/request-error', { turn, step, provider, failure, retryPolicy, signal },
  () => Promise.resolve(undefined),
)
if (action?.kind !== 'retry') {
  throw new LlmError(finish.failure.message, finish.failure.code, finish.failure)
}
continue  // 回到 while(true) 重新 buildRequest + stream

注意这里传入了 preparedCall?.retryPolicy——在 prepareCall 时捕获的策略。retry 插件(如 dsh-llm-retry)读这个 policy 决定是否重试、延迟多久。如果 waterfall 没有返回 {kind: 'retry'},错误被抛出终止 step。

时间线图

整个流程的时间线:

buildRequest
  |
  |-- [1] waterfall agent/request -> proposedConfig
  |
  |-- [2] llm.prepareCall(proposedConfig, signal)
  |         |-- registration = adapters.get(provider)
  |         |-- resolveModel -> contextWindow, reasoning, defaults
  |         |-- deepFreeze(structuredClone(config))
  |         +-- return frozen PreparedLlmCall { stream: closure(registration) }
  |
  |-- [3] markAgentLoopRequest(deepFreeze({ config, messages, system, tools, signal }))
  |
  +-- return { request, preparedCall }

step (消费流)
  |
  |-- [4] preparedCall.stream(request)
  |         |-- dispatched = true (one-shot)
  |         |-- callConfigEquals 验证
  |         +-- streamWithRegistration(request, { registration, config })
  |
  |-- [5] ctx.waterfall('llm/stream', request, () => adapterStream(...))
  |         |-- listener A: 缓存检查 -> miss -> next()
  |         |-- listener B: metrics wrap -> next()
  |         +-- terminal: adapterStream
  |               |-- adapter.stream(options) -> AsyncIterator
  |               +-- yield chunks (errors -> failure chunk)
  |
  |-- [6] for await (chunk of stream)
  |         |-- session.append('assistant/chunk', { turn, step, chunk }) -> seq
  |         +-- assembler.push(chunk)
  |
  |-- [7] assembler.finish -> check error/aborted
  |
  +-- [8] session.append('assistant/message', { message, usage }, { surfaceOp, sourceEventSeqs })

时间线上有两个关键的冻结时刻:

  • 时刻 2:prepareCall 冻结了 config 和 context,锁定了 adapter registration。从这里开始,任何外部变化(HMR、配置更新)都不影响这次调用。
  • 时刻 3:markAgentLoopRequest 冻结了完整请求。从这里开始,waterfall listener 无法修改请求内容,保证了 session log 的可重建性。

两次冻结之间(时刻 2 到 3),loop 还做了 header logging 和 context logging——这些 session events 记录了 prepareCall 解析出的能力信息,为后续 replay 和审计提供数据。

读者容易走错的路

错误路线一:“LlmRuntime 就是个 HTTP 客户端封装”

不是。LlmRuntime 完全不知道 HTTP、SSE、WebSocket 或任何网络协议。它是一个纯粹的注册表 + 分发器。registerAdapter 注册了 provider 路由到 adapter 的映射;stream 通过 waterfall 把请求分发到匹配的 adapter。adapter 内部可以用 fetch、可以用 SDK、可以读本地文件(mock),LlmRuntime 不关心。

如果你想理解真正的 HTTP 调用是怎么发生的,要去看具体的 adapter 包(如 packages/llm/llm-deepseekpackages/llm/llm-pi-ai),不在本章范围内。

错误路线二:“adapter 替换有时间间隙,请求可能看到 NO_ADAPTER”

不会。AdapterRegistrationHandle.replace() 的实现是:先 prepareRoutes 全量验证新路由集(不修改任何状态),通过后调用 commitRoutes 在一个同步代码段内完成 delete + set。JavaScript 单线程保证了这个同步段不会被打断——没有 await,没有 yield,中间不可能有其他代码执行。

而且即使 adapter 被 replace 了,已经通过 prepareCall 锁定了 registration 的调用完全不受影响——它用的是闭包引用,不再查 Map。

错误路线三:“在 llm/stream listener 里修改 options 就能改写请求”

对于 loop 请求:不能。它是 deep-frozen 的,修改会 throw(strict mode)或静默失败。

正确做法:如果你的 waterfall listener 需要”改写”请求(比如注入额外的 system message),你应该创建一个新对象并在调用 next() 时传入——但实际上 waterfall 的 next() 签名不接受新 options(它是 () => AsyncIterable<StreamChunk>),所以 listener 只能包装/短路流本身,不能改写传给下游 adapter 的 options。这是有意的限制。

错误路线四:“BlockAssembler 按 block-start 计数块数量”

不完全对。BlockAssemblerblock-start容忍缺失的。如果 adapter 直接发 text-delta 而没有先发 block-startensure() 方法会自动创建 partial。块的真正数量由 order 数组决定,而不是 block-start 的数量。这意味着你不能假设每个块都有 block-start/block-end 对。

错误路线五:“adapter 抛错会直接中断 for await 循环”

不会。adapterStream 内部有 try/catch 包裹 iterator.next(),adapter 的任何 throw 都被转为 {type: 'finish', reason: {kind: 'error', failure}} chunk yield 出去。for await 循环正常收到这个 chunk,push 给 assembler,assembler 记录 finish reason。循环正常结束后,loop 检查 assembler.finish.kind === 'error' 才进入错误处理。

middleware 抛的错误就不同了。如果一个 llm/stream listener 在 yield 之后 throw,这个错误会作为 generator 异常传播到 for await 循环,直接中断并 throw。这是 adapterStream 注释里说的 “consumer/middleware failures remain thrown” 的含义。

错误路线六:“preparedCall 可以重试——多次调 stream()”

不行。dispatched 标志位在第一次调用时设为 true,第二次调用会抛 INVALID_PREPARED_CALL。这是 one-shot 语义。如果需要重试,loop 回到 while(true) 的顶部,重新执行 buildRequest——意味着重新 prepareCall、重新解析能力、重新冻结请求。每次重试都是一次完整的新调用准备。

错误路线七:“replayState 是给用户看的调试数据”

不是。replayState 是 adapter 私有的不透明数据——由 adapter 在 finish chunk 中产出,存储在 message source 中,下次请求时通过 message.source.replayState 传回同一个 adapter(forAdapter 方法会在 adapter 不匹配时剥离它)。它的用途是高效重放:adapter 可以存储自己需要的状态,避免下次请求时重新解析完整的历史消息。用户和外部插件不应该读或依赖它的内容。


现在你知道了:一个请求从 buildRequest 出发,经过 prepareCall 锁定 adapter、markAgentLoopRequest 深冻结、llm/stream waterfall 分发、adapter 产出 chunks、BlockAssembler 逐 chunk 组装,最终以 assistant/message 事件提交到 session log,同时通过 sourceEventSeqs 引用所有原始 chunk 事件。这条路径的每个设计决策都指向同一个目标:请求的可重建性和调用的一致性,即使在 HMR、插件热加载、adapter 替换的动态环境中也不例外。

tool-call 块真正执行之后,工具结果还要回到 session log,并决定下一轮 step 怎么继续。