Main Chain B · Step 10

llm/stream:请求 → 适配器 → 块流

09 页造出的冻结请求在这里离开 harness,进入 LLM 适配器。这一页看 packages/llm/llm/src/index.ts(1153 行)里的 LlmRuntime——seam 的核心:llm/stream waterfall 如何找到 adapter、adapterStream 如何做代际绑定与失败归一化。0.1.3-alpha.2 起它还负责「文件引用 → 文本」的输入投影(fileRequestText,PR #3109),此前版本已有「图片 → 纯文本模型」投影与 TypertRemoteService @Remote 远程方法。

09 buildRequest:570 prepareCall · 11 step:418 stream(request) → llm/src/index.ts:1116 streamWithRegistration → :1017 adapterStream → (分叉) adapter.prepareCall / dispatch → llm-deepseek/src/adapter.ts → (下一步) 11 step 消费 chunk 流

示例本次示例:step 1 的流实况

示例轨迹 10-1 · 一个 chunk 从 adapter 到日志
# 11 页 346 行:stream = preparedCall.stream(request)
#   → 本页 streamWithRegistration(989 行)→ llm/stream waterfall
#   → 无中间件 → next() → adapterStream(898 行)
#   → adapter.prepareCall(910 行,代际绑定)→ dispatch(937 行)→ DeepSeekAdapter
# DeepSeek 返回 SSE 流 → translate 成 harness chunk:
#   { "type":"text-delta", "index":0, "text":"我来" }
#   { "type":"text-delta", "index":0, "text":"帮你" } ...
# 11 页逐 chunk 落盘(350 行):
#   { "type":"assistant/attempt", "seq":4, "data":{ "turn":1,"step":1,
#       "chunk":{ "type":"text-delta","index":0,"text":"我来" } } }
#   { "type":"assistant/attempt", "seq":5, ... } ...
# 每个 chunk 同时喂给 BlockAssembler(351 行)
来源:10/11 页行级解读;attempt 行格式与真实 fixture 的 assistant/attempt 行一致(06 页示例轨迹 06-1)
示例轨迹 10-2 · 如果 DeepSeek 在流中途断连
# adapter 的 iterator.next() reject(断连)
# → adapterStream 949-955 行:completed=true
#   yield adapterFailureChunk(error):
#     { "type":"finish", "reason":{ "kind":"error",
#         "failure":{ "message":"fetch failed","code":"NETWORK" } } }
# → 生成器正常结束(不 throw)
# → 11 页 assembler.finish.kind === 'error' → agent/request-error waterfall 仲裁
来源:10 页 946-956 行 + 11 页 373-389 行

§1LlmRuntime:服务与注册表

packages/llm/llm/src/index.ts服务骨架与 adapter 解析337-343, 969-973
337export class LlmRuntime extends TypertRemoteService {
337  private adapters = new Map<string, AdapterRegistration>()
338  private directory = new Map<string, LlmConfigurableProvider>()
339  private discoveries = new Map<
340    string,
341    (request: LlmModelDiscoveryRequest, signal?: AbortSignal) => Promise<readonly LlmDiscoveredModel[]>
342  >()
969  private registration(provider: string): AdapterRegistration {provider 解析的唯一路径
970    const registration = this.adapters.get(provider)
971    if (!registration) throw new LlmError(`no adapter registered for provider "${provider}"`, 'NO_ADAPTER')
973    return registration
974  }
336-342

LlmRuntime extends TypertRemoteService(0.1.3-alpha.2 起;仍是 Service 族)且 super(ctx, 'llm')——ctx.llm 由此而来(不是 ctx.provide)。TypertRemoteService 支持 @Remote 标注的方法(如 remoteDiscoverModels)跨进程调用——LLM 服务可被远程载波。两张表:adapters(provider route → 注册的适配器)与 directory(可配置 provider 目录,供模型发现)。

336

super(ctx, 'llm') 在构造器里(本页 §2 的 1106 行附近是公开 API)。全库的 seam 都长这样:class X extends Service { super(ctx, 'name') }。

968-972

provider 解析的唯一路径:找不到就是 NO_ADAPTER(LlmError)——09 页特别宽容的就是这个错误码。

§2streamWithRegistration:waterfall 的终局

packages/llm/llm/src/index.tsstream 的公共入口与 waterfall 包装1112-1126
1112  stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
1113    return this.streamWithRegistration(options)
1114  }
1116  private streamWithRegistration(
1117    options: GenerateOptions,
1118    prepared?: PreparedDispatch,
1119  ): AsyncIterable<StreamChunk> {
1120    return this.ctx.waterfall(
1121      this,
1122      'llm/stream',
1123      options,
1124      () => this.adapterStream(options, prepared),
1125    )
1126  }
1110-1112

公共 stream 不带 prepared——这就是 11 页 agent.ts:390 的 fallback 路径(preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request))。prepared 路径(prepareCall 返回的闭包,本文件 954 行)直接调 streamWithRegistration 并带上 PreparedDispatch。

1118-1123

llm/stream 是 waterfall——这是全 harness 拦截模型请求的扩展点。返回值是权威流:监听器可以短路(不调 next 直接返回自己的 AsyncIterable,如 llm-replay 的回放流)或包裹(调 next 后改写 chunk 流,如 llm-retry、invariant 校验)。fallback 是 adapterStream——next() 到达这里就是「真的调用模型」。

§3adapterStream:代际绑定 + 失败归一化边界

packages/llm/llm/src/index.ts终局适配器边界1017-1099
1017  private async * adapterStream(
1018    options: GenerateOptions,
1019    prepared?: PreparedDispatch,
1020  ): AsyncGenerator<StreamChunk> {
1021    let iterator: AsyncIterator<StreamChunk>
1022    try {
1023      const registration = prepared?.registration ?? this.registration(options.provider)
1024      const adapter = registration.adapter
1025      let modelInfo: LlmResolvedModelInfo
1026      let resolvedConfig: LlmCallConfig
1027      let dispatch: (options: GenerateOptions) => AsyncIterable<StreamChunk>
1028      if (prepared === undefined) {   // 普通路径也走代际绑定代际绑定
1029        const adapterCall = await adapter.prepareCall(options.provider, options.model, options.signal)
1030        modelInfo = this.normalizeModelInfo(registration, options.model, adapterCall.model)
1031        resolvedConfig = this.resolveCallWithInfo(options, modelInfo).config
1032        dispatch = options => adapterCall.stream(options)
1033      } else {
1034        modelInfo = prepared.modelInfo
1035        resolvedConfig = prepared.config
1036        dispatch = prepared.dispatch
1037      }
1038      if (prepared !== undefined && !callConfigEquals(options, resolvedConfig)) {
1039        throw new LlmError(
1040          'prepared LLM call config changed before adapter dispatch',
1041          'INVALID_PREPARED_CALL',
1042        )
1043      }
1044      const resolvedOptions = callConfigEquals(options, resolvedConfig)
1045        ? options
1046        : Object.isFrozen(options)
1047          ? deepFreeze({ ...options, ...resolvedConfig })
1048          : { ...options, ...resolvedConfig }
1049      // Files are never dispatched natively: every route receives handle text.先投影,再派发
1050      let projectedMessages: readonly Message[] = resolvedOptions.messages
1051      if (projectedMessages.some(message => contentHasFile(message.content))) {
1052        projectedMessages = projectFilesToText(projectedMessages, ref => this.fileReadPath(ref))
1053      }
1054      if (modelInfo.inputModalities !== undefined
1055        && !modelInfo.inputModalities.includes('image')
1056        && projectedMessages.some(message => contentHasImage(message.content))) {
1057        projectedMessages = projectImagesForTextModel(projectedMessages)
1058      }
1059      const projectedOptions = projectedMessages === resolvedOptions.messages
1060        ? resolvedOptions
1061        : Object.isFrozen(resolvedOptions)
1062          ? deepFreeze({ ...resolvedOptions, messages: projectedMessages as Message[] })
1063          : { ...resolvedOptions, messages: projectedMessages as Message[] }
1064      const stream = dispatch(this.forAdapter(projectedOptions, adapter))
1065      iterator = stream[Symbol.asyncIterator]()
1066    } catch (error: unknown) {
1067      yield adapterFailureChunk(error, options.signal)
1068      return
1069    }
1071    let completed = false
1072    try {
1073      while (true) {
1076          const next = await iterator.next()
1080        } catch (error: unknown) {
1081          completed = true
1082          yield adapterFailureChunk(error, options.signal)
1083          return
1084        }
1085        if (item.done) {
1086          completed = true
1087          return
1088        }
1091        yield item.value
1093    } finally {
1094      if (!completed) {
1095        const close = iterator.return?.bind(iterator)
1096        if (close) await close()
1097      }
1098    }
1099  }
1021-1035

代际绑定结构:两条路径都收敛到同一组三元(modelInfo / resolvedConfig / dispatch)。普通路径也先 adapter.prepareCall(...)(LlmAdapter.prepareCall,247 行)——模型解析与最终 stream 调用绑定到同一个 adapter 代际,设置变更发生在两者之间时不可能拼错。prepared 路径直接复用 09 页冻结的三元。

1036-1041

prepared 路径的第二个守卫(第一个在 09 页的 stream 闭包里):到 adapter 边界前再验一次 config 没漂移——中间件可能已经改写过 options。

1042-1046

如果 adapter 物化的 config 与 options 不同(普通路径),合并出一个 resolvedOptions;如果 options 已冻结(loop 请求都是 frozen)则 deepFreeze 新对象——冻结性传播。

1047-1056

输入投影(0.1.5 起是两段):先 文件投影(1049-1051)——「files are never dispatched natively,每条路由都收到 handle text」,用 projectFilesToText(llm/src/content.ts:196)把文件引用换成可读路径文本;再 图片投影(1052-1055)——模型模态声明不含 image 而消息里有图片时,用 projectImagesForTextModel(llm/src/content.ts:302)把图片块投影成文本模型可接受的形态。两者都靠代际捕获的 modelInfo.inputModalities 判定——不是靠猜。模态信息来自代际捕获的 modelInfo.inputModalities——不是靠猜。

1062-1063

真的调用模型:dispatch(...)——prepared 路径是 prepareCall 时捕获的同一代际 stream 入口;普通路径是 adapterCall.stream。传入前先过 forAdapter(§4)。

1064-1067

选择/派发阶段的失败归一化:任何错误(NO_ADAPTER、resolve 失败、dispatch 同步抛)都变成一个终态 finish chunk 后正常结束——生成器不 throw,消费者看到的是流协议内的错误。

1069-1096

迭代循环:每次 iterator.next() 的 reject 同样归一化为终态 chunk。finally(1091-1096)里 !completed 时才 iterator.return()——早退(消费者 break)要显式关闭上游;正常结束或已归一的失败路径不再 close。这是「失败边界的另一半」:中间件与消费者错误保持 throw,不被这里的 try 吃掉。这是失败边界的另一半。

1128-1135

adapterFailureChunk(index.ts:1127):normalizeLlmFailure 把任意抛出值(adapter 是第三方边界,什么都可能抛)归一化为可序列化 LlmFailure;signal 已 abort 或 code 是 ABORTED 时 reason 为 aborted,否则 error。

失败归一化边界是 LlmRuntime 的设计要点(类 JSDoc):adapter 选择、派发、迭代失败 → 终态 finish {kind:'error'|'aborted', failure}(模型请求结果);而 llm/stream 中间件、嵌套调用、cleanup、消费者错误 → 原样 throw(插件/消费者失败,不属于模型请求结果)。分清这两类失败,后面 11 页的 request-error 逻辑才读得懂。

§4forAdapter:replay 状态的归属检查

packages/llm/llm/src/index.tsforAdapter976-990
975  /** Remove replay state whose historical route is owned by another adapter. */
977  private forAdapter(options: GenerateOptions, adapter: LlmAdapter): GenerateOptions {
978    const messages: Message[] = options.messages.map((message) => {
979      const source = message.source
980      if (source?.kind !== 'model' || source.replayState === undefined) return message
981      if (this.adapters.get(source.provider)?.adapter === adapter) return message
982      return freezeMessage({
983        ...message,
984        source: { kind: 'model', provider: source.provider, model: source.model },
985      })
986    })
987    if (messages.every((message, index) => message === options.messages[index])) return options
988    const filtered = { ...options, messages }
989    return Object.isFrozen(options) ? deepFreeze(filtered) : filtered
990  }
970-980

replayState 是某 adapter 专属的回放状态(如缓存的 provider 内部状态)。切换 adapter 后,历史消息里属于其他 adapter的 replayState 必须剥掉——否则新 adapter 会拿着别人的内部状态。属于当前 adapter 的保留。

981-983

优化与一致性:没剥任何消息就原样返回 options(不分配新对象);剥了则新建,且若原 options 已冻结则新对象也 deepFreeze——冻结性贯穿请求链。