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 远程方法。
示例本次示例:step 1 的流实况
# 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 行)
# 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 仲裁
§1LlmRuntime:服务与注册表
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 }
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 目录,供模型发现)。
super(ctx, 'llm') 在构造器里(本页 §2 的 1106 行附近是公开 API)。全库的 seam 都长这样:class X extends Service { super(ctx, 'name') }。
provider 解析的唯一路径:找不到就是 NO_ADAPTER(LlmError)——09 页特别宽容的就是这个错误码。
§2streamWithRegistration:waterfall 的终局
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 }
公共 stream 不带 prepared——这就是 11 页 agent.ts:390 的 fallback 路径(preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request))。prepared 路径(prepareCall 返回的闭包,本文件 954 行)直接调 streamWithRegistration 并带上 PreparedDispatch。
llm/stream 是 waterfall——这是全 harness 拦截模型请求的扩展点。返回值是权威流:监听器可以短路(不调 next 直接返回自己的 AsyncIterable,如 llm-replay 的回放流)或包裹(调 next 后改写 chunk 流,如 llm-retry、invariant 校验)。fallback 是 adapterStream——next() 到达这里就是「真的调用模型」。
§3adapterStream:代际绑定 + 失败归一化边界
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 }
代际绑定结构:两条路径都收敛到同一组三元(modelInfo / resolvedConfig / dispatch)。普通路径也先 adapter.prepareCall(...)(LlmAdapter.prepareCall,247 行)——模型解析与最终 stream 调用绑定到同一个 adapter 代际,设置变更发生在两者之间时不可能拼错。prepared 路径直接复用 09 页冻结的三元。
prepared 路径的第二个守卫(第一个在 09 页的 stream 闭包里):到 adapter 边界前再验一次 config 没漂移——中间件可能已经改写过 options。
如果 adapter 物化的 config 与 options 不同(普通路径),合并出一个 resolvedOptions;如果 options 已冻结(loop 请求都是 frozen)则 deepFreeze 新对象——冻结性传播。
输入投影(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——不是靠猜。
真的调用模型:dispatch(...)——prepared 路径是 prepareCall 时捕获的同一代际 stream 入口;普通路径是 adapterCall.stream。传入前先过 forAdapter(§4)。
选择/派发阶段的失败归一化:任何错误(NO_ADAPTER、resolve 失败、dispatch 同步抛)都变成一个终态 finish chunk 后正常结束——生成器不 throw,消费者看到的是流协议内的错误。
迭代循环:每次 iterator.next() 的 reject 同样归一化为终态 chunk。finally(1091-1096)里 !completed 时才 iterator.return()——早退(消费者 break)要显式关闭上游;正常结束或已归一的失败路径不再 close。这是「失败边界的另一半」:中间件与消费者错误保持 throw,不被这里的 try 吃掉。这是失败边界的另一半。
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 状态的归属检查
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 }
replayState 是某 adapter 专属的回放状态(如缓存的 provider 内部状态)。切换 adapter 后,历史消息里属于其他 adapter的 replayState 必须剥掉——否则新 adapter 会拿着别人的内部状态。属于当前 adapter 的保留。
优化与一致性:没剥任何消息就原样返回 options(不分配新对象);剥了则新建,且若原 options 已冻结则新对象也 deepFreeze——冻结性贯穿请求链。