第 10 章
LLM 适配层:twin adapters 与流式词汇
场景还原
给一个 agent harness 接新的模型供应商,是那种看起来简单、做起来翻车的活。
你第一版可能这么写:在 agent-loop 里直接调 fetch,把供应商的流式响应 JSON.parse 出来,拼进会话。第一个供应商很顺利,因为它返回的 content 是字符串。第二个供应商来了,它的增量里多了 reasoning_content,tool call 的 arguments 是对象而第一个是字符串,错误是流中间的一个事件而第一个是 HTTP 状态码。于是循环里开始长 if (provider === ...)。等到第三个供应商,你已经说不清这段代码在干什么了。
本章回答两个问题。适配一个模型供应商,最少要写什么?供应商的流式响应,怎么变成 session 里的事件?dsh 的答案是:一个抽象类上一个必需方法,加一套七个词的 chunk 协议,加两条错误路径规则。
逐行精读
责任声明:llm 包是什么
先看根 AGENTS.md 第 17 行怎么定义 llm 包。
17 llm/ LLM capability: Service Definition/Consumer + DeepSeek providers这一行说 llm 包同时承担两个角色:对 agent-loop 它是 Service Definition 与 Consumer,对 DeepSeek 它是 provider 实现。包目录 packages/llm/ 下有五个子包:llm(词汇表与运行时)、llm-deepseek(直接 fetch 实现)、llm-pi-ai(库实现)、llm-retry(重试策略)、token-meter。本章只看前三个。
第 4 章讲 turn flow 时,docs/architecture.md 第 76 行画过一句话的请求链路。看 73 至 81 行这一段:
73 step/start74 append entered messages as user/message75 derive model history from the log76 agent/request -> llm/stream -> assistant/chunk* -> assistant/message77 tool/call* -> tools/pre-execute -> tools/execute -> tools/post-execute -> tool/result*78 step/end79 tools owe another request, or next-step input arrived -> claim -> next step80 -> agent/turn-stopping81turn/end注意第 76 行的形状:agent/request -> llm/stream -> assistant/chunk* -> assistant/message。一次模型请求从 agent/request 出发,经过 llm/stream,一路以 assistant/chunk 事件流回,最后落成一条 assistant/message。本章顺着这条链走:契约、注册、分派、翻译、事件化、装配。
契约:唯一必需方法是 stream
packages/llm/llm/src/index.ts 里 LlmAdapter 是适配器的基类。看 238 至 260 行:
238 /**239 * Bind exact model metadata and the eventual request dispatch to one adapter generation.240 * Dynamic adapters override this so settings changes between preparation and241 * dispatch cannot combine one generation's capabilities with another's endpoint.242 * @param provider - registered provider route.243 * @param model - exact model id.244 * @param signal - cancellation for model resolution.245 * @returns model metadata and a one-generation stream entry point.246 */247 async prepareCall(provider: string, model: string, signal?: AbortSignal): Promise<PreparedAdapterCall> {248 return {249 model: await this.resolveModel(provider, model, signal),250 stream: options => this.stream(options),251 }252 }253254 /**255 * Stream one model call as raw chunks. The only required method.256 * @param options - the fully-assembled request; implementations must honor `options.signal`.257 * @returns the chunk stream, obeying the adapter contract documented on `StreamChunk`.258 */259 abstract stream(options: GenerateOptions): AsyncIterable<StreamChunk>第 255 行注释写得很直白:The only required method。适配器实现者被要求写的只有一个方法:输入组装好的请求,输出一叠 StreamChunk。其余方法(providerInfo、listModels、resolveModel、prepareCall)都有默认实现,是给注册表做元数据用的。
prepareCall 值得单独看。它把「这次调用的模型元数据」和「这次调用的流入口」绑在同一代解析上:第 249 行先解析模型信息,第 250 行让 stream 闭包直接落回这个适配器。配置在请求中途变了,正在飞的请求不会换模型。
注册:一次注册绑定路由与适配器
适配器怎么进注册表?LlmRuntime.registerAdapter,365 至 394 行:
365 registerAdapter(providers: string[], adapter: LlmAdapter): AdapterRegistrationHandle {366 // The routes this registration currently holds; `replace` rewrites it, and367 // the disposer releases whatever it holds at disposal time.368 const owned = new Set<string>()369 // The disposer has run: `owned` being empty cannot say so on its own,370 // because `replace([])` legally leaves a live registration holding none.371 let released = false372 const dispose = this.ctx.effect(function* (this: LlmRuntime) {373 if (providers.length === 0) throw new LlmError('an adapter must register at least one provider', 'INVALID_ADAPTER')374 this.commitRoutes(owned, this.prepareRoutes(providers, adapter, owned))375 yield () => {376 released = true377 for (const provider of owned) this.adapters.delete(provider)378 owned.clear()379 this.emitAdaptersUpdated()380 }381 }.bind(this), 'llm.registerAdapter()')⋯ // ... 中间两行是效果注册的收尾,见正文 ...384 const handle = (() => void dispose()) as AdapterRegistrationHandle385 handle.replace = (next: string[]): void => {386 // Registering here would leak: the effect's disposer already ran, so387 // nothing remains to release whatever this call would put in the map.388 if (released) {389 throw new LlmError('a disposed adapter registration cannot replace its routes', 'REGISTRATION_DISPOSED')390 }391 this.commitRoutes(owned, this.prepareRoutes(next, adapter, owned))392 }393 return handle394 }三件事。第一,注册走 ctx.effect,随生命周期自动回收,disposer 把路由从表里删掉并通知观察者(379 行 emitAdaptersUpdated)。第二,路由校验全部或全无:prepareRoutes(401 行起)逐条检查名字非空、路由未被他方占用、元数据合法,任何一条不过就整体拒绝,注册表保持原样。第三,返回的句柄带 replace,可以原子换路由集。注意 373 行:空路由集直接拒绝,换到零路由要走 replace([]),不能靠注册一个空数组。
然后,分派入口 stream,985 至 999 行:
985 stream(options: GenerateOptions): AsyncIterable<StreamChunk> {986 return this.streamWithRegistration(options)987 }988989 private streamWithRegistration(990 options: GenerateOptions,991 prepared?: PreparedDispatch,992 ): AsyncIterable<StreamChunk> {993 return this.ctx.waterfall(994 this,995 'llm/stream',996 options,997 () => this.adapterStream(options, prepared),998 )999 }ctx.waterfall 是 Cordis 的瀑布机制:先过所有 llm/stream 监听者,每个监听者可以调 next() 继续往下,也可以自己 yield chunk 短路。这行代码把「每个模型调用必经的路口」变成扩展点。后面会看到 replay、session 标题、检查点策略都挂在这里。
归一化:adapter 抛的错变成 finish chunk
瀑布的末端是 adapterStream(898 行)。它做两件关键的事。第一,把适配器选择、prepare、迭代的所有抛错统一转成流里最后一个 finish chunk,错误不会从 for await 外面冒出来。看 1002 至 1011 行的收尾函数:
1002/** Convert one adapter throw into the stream protocol's terminal outcome. */1003function adapterFailureChunk(error: unknown, signal?: AbortSignal): StreamChunk {1004 const failure = normalizeLlmFailure(error)1005 return {1006 type: 'finish',1007 reason: signal?.aborted || failure.code === 'ABORTED'1008 ? { kind: 'aborted', failure }1009 : { kind: 'error', failure },1010 }1011}信号已中止且错误代码是 ABORTED 时,产生 aborted finish,其余一律 error finish。消费者只需要处理两种终态。这正是契约里两个错误路径的运行时实现:adapter 可以抛(在这里被接住),也可以主动结束于 finish {kind:'error'|'aborted'}。
词汇表:七个词闭起来
adapter 返回的 StreamChunk 在 packages/llm/llm/src/types.ts 304 至 324 行:
304/**305 * Raw streaming protocol emitted by adapters.306 * Block indexes correlate interleaved deltas, and `block-end` carries the307 * assembled block. Adapters emit usage before the terminal finish and nothing308 * afterward; tool arguments remain raw JSON strings. An adapter implementation309 * may throw, but `LlmRuntime.stream()` normalizes that failure to a terminal310 * `error` or `aborted` finish before exposing it to consumers.311 */312export type StreamChunk =313 | { type: 'block-start'; index: number; blockType: ContentBlockType }314 | { type: 'text-delta'; index: number; text: string }315 | { type: 'reasoning-delta'; index: number; text: string }316 | { type: 'tool-call-delta'; index: number; id: CallId; name?: string; argumentsDelta: string }317 | { type: 'block-end'; index: number; block: ContentBlock }318 | { type: 'usage'; usage: TokenUsage }319 | {320 type: 'finish'321 reason: FinishReason322 /** Replay metadata for a successful response; see {@link ReplayEnvelope}. */323 replayState?: ReplayEnvelope324 }七个变体:block-start、三种 delta、block-end、usage、finish。这是闭的 discriminated union,docs/subsystems/llm-streaming.md 第 158 行写明:a switch over type ends with assertNever,新增一种变体,没处理它的消费方编译期就报错。注意第 306 至 308 行的三条约定已经写进注释:block-end 携带完整块;usage 在 finish 前、之后什么都没有;tool arguments 全程是原始 JSON 字符串。
块的形状由 ContentBlockMap 定义,95 至 105 行:
95/**96 * Merge-extensible content blocks keyed by `type`. New core blocks must land97 * with adapter, UI, and compaction support.98 */99export interface ContentBlockMap {100 'text': TextBlock101 'reasoning': ReasoningBlock102 'image': ImageBlock103 'tool-call': ToolCallBlock104 'tool-result': ToolResultBlock105}和 chunk 相反,块词汇表是 merge-extensible 的:插件可以往 ContentBlockMap 里加新 type,switch 时对未知 type 落默认分支。第 96 至 97 行的注释定了新核心块的成本:adapter、UI、compaction 三处都要跟上。一个是协议(谁都不能漏处理),一个是内容(可扩展),这是两个层级的界限。
翻译:DeepSeek wire 怎么变成 chunk
第一个实现是 dsh-llm-deepseek:直接 fetch,SSE 帧解析交给 eventsource-parser 库(packages/llm/llm-deepseek/src/sse.ts 第 28 至 40 行的 parseSse),格式翻译在仓库内完成。
wire 上的 finish_reason 是供应商词,要先映射成 harness 词。packages/llm/llm-deepseek/src/translate.ts 31 至 43 行:
31export function mapFinishReason(reason: string): FinishReason {32 switch (reason) {33 case 'stop': return { kind: 'stop' }34 case 'tool_calls': return { kind: 'tool-calls' }35 case 'length': return { kind: 'max-tokens' }36 default:37 // content_filter, insufficient_system_resource, future additions.38 return {39 kind: 'error',40 failure: { message: `model stopped: ${reason}`, code: reason.toUpperCase() },41 }42 }43}默认分支是显式写出来的:供应商新增的终止原因(content_filter、insufficient_system_resource 都在第 37 行注释里点名了)不转成某种合理猜测,而是变成 error finish,code 用原文大写。
然后是主体翻译函数 translate,86 至 118 行:
86export async function* translate(payloads: AsyncIterable<string>): AsyncGenerator<StreamChunk> {87 let nextIndex = 088 let textBlock: OpenBlock | undefined89 let reasoningBlock: OpenBlock | undefined90 const toolBlocks = new Map<number, OpenBlock>()91 const order: OpenBlock[] = []92 let pendingFinish: FinishReason | undefined93 let pendingUsage: TokenUsage | undefined9495 function open(kind: OpenBlock['kind']): OpenBlock {96 const block: OpenBlock = { index: nextIndex++, kind, text: '' }97 order.push(block)98 return block99 }100101 for await (const payload of payloads) {102 if (payload === DONE) {103 for (const block of order) {104 yield { type: 'block-end', index: block.index, block: closeBlock(block) }105 }106 if (pendingUsage) yield { type: 'usage', usage: pendingUsage }107 const reason = pendingFinish ?? { kind: 'stop' as const }108 yield {109 type: 'finish',110 reason: reason.kind === 'stop' && order.length === 0111 ? {112 kind: 'error',113 failure: { message: 'model returned a completed response with no content', code: EMPTY_RESPONSE_CODE },114 }115 : reason,116 }117 return118 }翻译器维护三类开着的块:text、reasoning、tool-call,以及一个统一出场顺序表 order。[DONE] 哨兵到达时统一收尾:按顺序发 block-end、发累积的 usage、发 finish,然后 return。注意 110 至 115 行:finish reason 是 stop 但一个块都没开出来,变成 EMPTY_RESPONSE 错误,当不了空成功。
再看 delta 到 chunk 的逐段映射,130 至 150 行:
130 // Reasoning first: thinking mode interleaves it before text. The131 // empty-string first chunk must not open a block.132 const reasoning = delta?.reasoning_content133 if (typeof reasoning === 'string' && reasoning.length > 0) {134 if (!reasoningBlock) {135 reasoningBlock = open('reasoning')136 yield { type: 'block-start', index: reasoningBlock.index, blockType: 'reasoning' }137 }138 reasoningBlock.text += reasoning139 yield { type: 'reasoning-delta', index: reasoningBlock.index, text: reasoning }140 }141142 const content = delta?.content143 if (typeof content === 'string' && content.length > 0) {144 if (!textBlock) {145 textBlock = open('text')146 yield { type: 'block-start', index: textBlock.index, blockType: 'text' }147 }148 textBlock.text += content149 yield { type: 'text-delta', index: textBlock.index, text: content }150 }先说明:130 至 131 行是注释。thinking 模式会把 reasoning 交织在 text 前面,所以先处理 reasoning;但首块可能是个空字符串,空字符串不能开块。所以第 133 行的判断同时查了类型和长度,第 143 行对 text 用同一模式。block 只开一次,之后每个 delta 追加到同一块的累积文本,同时发一条带 index 的增量。
翻译:pi-ai 事件怎么变成 chunk
第二个实现是 dsh-llm-pi-ai:同一端点走 @earendil-works/pi-ai 库。库自带事件词汇表(text_start、text_delta、thinking_delta、toolcall_start 等),adapter 不需要手写 SSE 解析,只做事件到 chunk 的映射。packages/llm/llm-pi-ai/src/stream.ts 127 至 156 行:
127export async function* toStreamChunks(128 events: AsyncIterable<AssistantMessageEvent>,129 contextWindow?: number,130): AsyncGenerator<StreamChunk> {131 // pi-ai contentIndex ↔ our block index map 1:1 (both count blocks from 0132 // in stream order), but we track ids per index for tool calls.133 const toolIds = new Map<number, { id: string; name: string }>()134135 for await (const event of events) {136 switch (event.type) {137 case 'start':138 break139 case 'text_start':140 yield { type: 'block-start', index: event.contentIndex, blockType: 'text' }141 break142 case 'text_delta':143 yield { type: 'text-delta', index: event.contentIndex, text: event.delta }144 break145 case 'text_end':146 yield { type: 'block-end', index: event.contentIndex, block: { type: 'text', text: event.content } }147 break148 case 'thinking_start':149 yield { type: 'block-start', index: event.contentIndex, blockType: 'reasoning' }150 break151 case 'thinking_delta':152 yield { type: 'reasoning-delta', index: event.contentIndex, text: event.delta }153 break154 case 'thinking_end':155 yield { type: 'block-end', index: event.contentIndex, block: { type: 'reasoning', text: event.content } }156 break事件的名字直接对映 chunk 的名字:text_start 对 block-start,text_delta 对 text-delta,thinking_* 对 reasoning-*。pi-ai 的 contentIndex 和 harness 的块 index 一一对应(131 至 132 行注释)。映射接近直译,因为两边的模型一致。差别在细节:pi-ai 的事件流是「开始、增量、结束」三段式,block-end 可以直接携带完整块,而 deepseek 的 wire 只有增量,完整块要等 [DONE] 才组装。这是同一词汇表下的两种翻译风格。
pi-ai 的另一个特点在 191 至 198 行:终态事件处理。
191 case 'done':192 yield { type: 'usage', usage: mapUsage(event.message.usage) }193 yield {194 type: 'finish',195 reason: mapStopReason(event.message, contextWindow),196 replayState: toPiReplayState(event.message),197 }198 returnpi-ai 不在流中间抛错,失败以 error 事件到达,所以 done 和 error 两个事件都映射成 usage 加 finish 的结尾(202 至 204 行)。两个 adapter 殊途同归:deepseek 的收尾由 [DONE] 触发,pi-ai 的收尾由终态事件触发,最后产出的都是同一串 chunk,usage 在前 finish 在后。
消费:chunk 变成 session 事件
现在看消费端。packages/core/agent-loop/src/agent.ts 的 step 循环,339 至 354 行:
339 while (true) {340 const { request, preparedCall } = await this.buildRequest(341 turn, step, assembly.tools, system, this.session.deriveMessages(), signal,342 )343 const assembler = new BlockAssembler()344 const chunkSeqs: number[] = []345 try {346 const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)347 signal.throwIfAborted()348 for await (const chunk of stream) {349 signal.throwIfAborted()350 chunkSeqs.push(this.session.append('assistant/chunk', { turn, step, chunk }).seq)351 assembler.push(chunk)352 }353 signal.throwIfAborted()这里回答「流如何映射成 session 事件」。第 346 行选流:prepare 过的调用走 preparedCall.stream(request),否则直接 ctx.llm.stream(request)。第 348 至 352 行的 for-await 里,每个 chunk 干两件事:第 350 行把原始 chunk 写进 session 事件日志(assistant/chunk),拿到事件序号;第 351 行把同一个 chunk 喂给装配机。
流结束后,第 392 至 409 行落成消息事件:
392 const message = createAssistantMessage({393 content: assembler.blocks(),394 source: {395 provider: request.provider,396 model: request.model,⋯ ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {},· },· })· this.session.append(· 'assistant/message',· {· turn,· step,· message,⋯ ...assembler.usage === undefined ? {} : { usage: assembler.usage },407 },408 { surfaceOp: 'append', sourceEventSeqs: chunkSeqs },409 )消息来自装配机的 blocks(),source 带上 provider、model 和(如果 finish 带了)replayState。第 408 行的 sourceEventSeqs: chunkSeqs 把这条消息和它由哪些 assistant/chunk 事件拼成这件事记录在案,UI 和回放可以靠它把增量重新串起来。
两个事件类型的定义在 packages/core/session/src/types.ts 265 至 277 行:
265 /** Raw stream chunk — token-level replay fidelity. */266 'assistant/chunk': { turn: number; step: number; chunk: StreamChunk }267 /**268 * Assembled assistant message for one step (derived history uses this).269 * Carries the step's `usage` when the adapter reported token accounting, so270 * the model output and its accounting travel together (there is no separate271 * usage record). `usage` is absent when the adapter reported none. A turn272 * cancelled mid-stream finalizes its delivered text/reasoning prefix as this273 * event with `interrupted: true`; undispatched tool calls are absent. The274 * marker distinguishes that prefix without re-deriving interruption from turn275 * boundaries. An aborted turn with no such event streamed no visible content.276 */277 'assistant/message': { turn: number; step: number; message: AssistantMessage; usage?: TokenUsage; interrupted?: true }第 265 行的注释只有一句话:Raw stream chunk — token-level replay fidelity。chunk 事件存在的理由是回放保真,token 粒度。第 268 至 276 行注释解释消息事件的几条约束:usage 与消息同行;中断的 turn 以 interrupted: true 消息收尾,未派发的 tool call 不出现在里面;中断但没有任何可见内容的 turn 不产生这条事件。
装配:一个共享的折叠算法
BlockAssembler 是唯一的装配实现,packages/llm/llm/src/assembler.ts。docs/subsystems/llm-streaming.md 289 至 293 行说明它的地位:
289## `BlockAssembler`290291`BlockAssembler` ([`packages/llm/llm/src/assembler.ts`](../../packages/llm/llm/src/assembler.ts)) is the single shared implementation that folds a `StreamChunk` stream back into `ContentBlock`s, usage, finish reason, and replay state. The loop logs the raw chunks while feeding the same chunks through an assembler, then stores the assembled assistant content with the provider and model that produced it. A consumer that needs the assembled result without re-implementing the fold uses this.292293One keep/drop decision covers content and metadata together: a `max-tokens` finish drops every tool call because a truncated call is unsafe to execute, and the same decision prunes the replay envelope's per-block entry at each dropped position. `blocks()` and `replayState` therefore cannot disagree, whatever assembly removes.single shared implementation:所有消费方共用同一个折叠算法。第 293 行的关键决策:keep/drop 同时作用于内容与元数据。max-tokens 收尾时每个 tool call 都被丢弃,因为截断的调用不能安全执行;同一决策同步修剪 replay envelope 里对应位置的条目,所以 blocks() 和 replayState 永远一致。
看 push 方法,48 至 95 行,整个闭 union 的 switch:
48 push(chunk: StreamChunk): void {49 switch (chunk.type) {50 case 'block-start': {51 if (!this.partials.has(chunk.index)) {52 this.order.push(chunk.index)53 this.partials.set(chunk.index, {54 blockType: chunk.blockType,55 text: '',56 toolCallArguments: '',57 })58 }59 return60 }61 case 'text-delta':62 case 'reasoning-delta': {63 const partial = this.ensure(chunk.index, chunk.type === 'text-delta' ? 'text' : 'reasoning')64 if (partial.block) return // closed by block-end; ignore stragglers65 partial.text += chunk.text66 return67 }68 case 'tool-call-delta': {69 const partial = this.ensure(chunk.index, 'tool-call')70 if (partial.block) return // closed by block-end; ignore stragglers71 partial.toolCallId = chunk.id72 if (chunk.name) partial.toolCallName = chunk.name73 partial.toolCallArguments += chunk.argumentsDelta74 return75 }76 case 'block-end': {77 const partial = this.ensure(chunk.index, chunk.block.type)78 // First close wins; ignoring re-close stragglers keeps streamed output79 // and the final assembled block in agreement.80 if (partial.block) return81 partial.block = chunk.block82 return83 }84 case 'usage': {85 this._usage = chunk.usage86 return87 }88 case 'finish': {89 this._finish = chunk.reason90 this._replayState = chunk.replayState91 return92 }93 default: return assertNever(chunk, 'BlockAssembler.push')94 }95 }装配机对每种 chunk 的响应:block-start 建 partial;delta 追加到对应 partial,但块已被 block-end 关闭后的迟到 delta 直接忽略(第 64 和 70 行注释:closed by block-end; ignore stragglers),防止异常 adapter 撑大内存或污染已完成的块;block-end 以第一次关闭为准;usage 和 finish 各存一份;default 走 assertNever,闭 union 的编译期保障落地到运行时。
第 134 至 149 行是 keep/drop 的核心:
134 private assembled(): { blocks: ContentBlock[]; replay: ReplayEnvelope | undefined } {135 const all = this.order.map(index => this.assemble(this.mustGet(index), index))136 const kept = this.finish.kind === 'max-tokens'137 ? all.map(block => block.type !== 'tool-call')138 : undefined139 const blocks = kept === undefined ? all : all.filter((_, position) => kept[position])140 const envelope = this._replayState141 if (envelope?.blocks === undefined) return { blocks, replay: envelope }142 if (envelope.blocks.length !== all.length) return { blocks, replay: undefined }143 return {144 blocks,145 replay: kept === undefined || blocks.length === all.length146 ? envelope147 : { response: envelope.response, blocks: envelope.blocks.filter((_, position) => kept[position]) },148 }149 }第 136 至 137 行:只有 finish 是 max-tokens 时才做过滤,all.map(block => block.type !== 'tool-call') 把每个 tool-call 位置标成 false。第 142 行有个防御:envelope 的块数和实际块数对不上,整个 envelope 丢弃,replay 元数据宁可不要也不能错位。
中断时的投影是另一条路径,168 至 178 行:
168 interruptedBlocks(): ContentBlock[] {169 return this.order170 .map((index) => {171 const partial = this.mustGet(index)172 const type = partial.block?.type ?? partial.blockType173 if (type !== 'text' && type !== 'reasoning') return undefined174 return this.assemble(partial, index)175 })176 .filter((block): block is ContentBlock =>177 (block?.type === 'text' || block?.type === 'reasoning') && block.text.trim() !== '')178 }只保留文字与推理块,且内容非空白。tool call 被排除的原因写在第 164 至 165 行的 JSDoc 里:中断发生在派发之前,保留一个 tool call 需要伪造结果。这个前缀由 agent.ts 356 至 368 行的 catch 分支消费。
消息与来源:不可变的价值类型
最后看消息本身。packages/llm/llm/src/message.ts 128 至 138 行:
128/** One immutable message representation shared by delivery, durable history, and model requests. */129export interface Message {130 /** Stable identity preserved across every representation boundary. */131 readonly id: MessageId132 /** Provider-neutral conversation role. */133 readonly role: 'system' | 'user' | 'assistant'134 /** Exact model-facing blocks. */135 readonly content: ContentBlock[]136 /** Required source fields supplied by the producer. */137 readonly source: MessageSource138}第 128 行注释:同一份表示,贯穿投递、持久历史与模型请求三处。不可变靠构造函数保证,169 至 185 行:
169export function freezeMessage<T extends Message>(message: T): T {170 return deepFreeze(structuredClone(message))171}172173/**174 * Create one identified message and freeze it before publication.175 * @param input - complete role, content, and source for a new message.176 * @returns an immutable message with a fresh stable identity.177 */178export function createMessage<T extends NewMessage>(179 input: T & { readonly id?: never },180): T & Pick<Message, 'id'> {181 return freezeMessage({⋯ ...input,183 id: MessageId(crypto.randomUUID()),184 })185}deepFreeze 加 structuredClone,发布前冻结与脱离。任何拿到 Message 的人改不动它。
设计决策分析
twin adapters:一个契约配两个实现
为什么有两个 DeepSeek adapter?.agents/notes/implemented/architecture/2026-06-13-twin-llm-adapters.md 13 至 18 行:
13Ship **two** adapters against the one contract from the start, deliberately built on different internals:1415- `dsh-llm-deepseek` — direct `fetch` + in-repo translation against the DeepSeek API; SSE framing is delegated to `eventsource-parser` ([the archived SSE-parser swap](../../archived/simplification/2026-07-26-eventsource-parser-for-deepseek-sse.md)). The twin identity is owning the fetch/translate internals rather than delegating to a full provider SDK, not hand-rolling transport plumbing.16- `dsh-llm-pi-ai` — the same endpoint through the `@earendil-works/pi-ai` library (its own event vocabulary).1718The rule they enforce: **anything the StreamChunk vocabulary cannot express for BOTH implementations is a core-vocabulary bug**, caught immediately rather than at the next provider. The pair pinned down conventions now documented on `StreamChunk` in `dsh-llm/src/types.ts`: usage emitted before finish, nothing after finish, tool-call `arguments` as raw JSON strings end-to-end, and the two sanctioned error paths (throw from `stream()` *or* end with `finish {kind:'error'|'aborted'}`) that a consumer must handle on both sides — a divergence the library-backed adapter surfaced that a single direct-fetch adapter would have hidden.第 18 行的规则是这条 note 的题眼:StreamChunk 词汇表如果两个实现都表达不了,那是核心词汇的 bug,当下就暴露,等不到第三个供应商来撞。注意它验证出来的东西:usage 在 finish 前、finish 后无内容、arguments 保持原始 JSON、两条错误路径,这四条约定都是从双实现对峙中长出来的,写进了 StreamChunk 的注释。第 13 至 16 行还明确了两者的差异点:一个自己持 fetch 与翻译(SSE 帧解析交给库),一个把整个传输交给 provider SDK。
note 第 22 至 23 行记录了两个备选方案为何被否。单 adapter 会让「DeepSeek-via-fetch 的假设」悄悄编进词汇表,抽象未经检验;mock 第二 adapter 便宜但复现不了真实供应商的线上怪癖。双实现的代价在 note 第 27 行:adapter 与 key-gated e2e 维护翻倍,换来持续的中立性验证。
契约条款:把失败变成词汇的一部分
docs/subsystems/llm-streaming.md 227 至 239 行是适配器必须遵守的契约:
231- **`usage` before `finish`, nothing after `finish`.** Defer both to the provider's end-of-stream marker so a trailing usage-only chunk can't violate the ordering.232- **Tool-call `arguments` stay raw JSON strings end-to-end.** Partial fragments stream via `argumentsDelta`; a provider that hands back parsed objects re-stringifies at `block-end`.233- **Two sanctioned error paths, one `LlmFailure` type.** A failure may either THROW from `stream()` (transport/protocol errors) **or** end the stream with `finish {kind:'error'|'aborted', failure}` (provider in-band errors, for adapters that can't throw mid-stream). `LlmError.failure` carries the same `LlmFailure`. After the call selects its adapter, the stream preserves the exact thrown `Error` object and associates immutable facts plus the serving registration's immutable retry policy with that call; the agent loop closes the failed step and offers the error, facts, immutable prior-retried facts, serving policy, and turn signal to `agent/request-error`. A handling listener returns `{ kind: 'retry' }` after its awaited repair; absent recovery the structured failure becomes the turn error, and no normal assistant message or tool side effect is committed for that attempt.234- **One adapter call is one provider attempt.** Adapters disable library retries. Agent-level recovery opens another durable numbered turn; direct `ctx.llm.stream()` callers remain single-attempt.235- **Provider stalls are bounded at the transport.** Both shipping remote adapters expose positive finite `streamIdleTimeoutMs` with a five-minute default. The watchdog arms only while iterator `next()` is outstanding, uses one stable signal for the whole request, maps its own expiry to `TIMEOUT`, and keeps an earlier caller abort as `ABORTED`.236- **Context overflow has one canonical code.** Both DeepSeek adapters classify explicit provider detail through `isContextWindowExceededError()` and surface `CONTEXT_WINDOW_EXCEEDED`, whether the failure arrives as a thrown HTTP `LlmError` or an in-band finish error. Consumers route on the code, never provider text.237- **An empty completion is a retryable error, not a silent success.** Both adapters map a terminal `stop` finish that carried no content blocks to `finish {kind:'error'}` with the canonical `EMPTY_RESPONSE` code, and `dsh-llm-retry` retries it by default; see [empty model responses are retryable](../../.agents/notes/implemented/bug-fix/2026-07-24-empty-model-response-is-retryable.md).238- **Every provider HTTP request carries the app-attribution header.** Adapters send `attributionHeaders()` (below) - the `User-Agent` baseline - and prove it with a wire-level test.239- **Replay state is adapter-owned; its split is shared.** A successful `finish` may carry a `ReplayEnvelope`: opaque response-level metadata plus optional per-block entries aligned with the emitted block sequence. The alignment is the harness's vocabulary — when assembly drops a block it drops the entry at the same position, so stored metadata always describes stored content. The loop stores the pruned envelope with the assembled assistant message. On a later request, `LlmRuntime` passes the state only when the historical provider and target provider are currently registered to the exact same adapter instance. That adapter validates the state and owns any cross-model or cross-provider conversion; other adapters receive the provider-neutral content plus provider/model fields without the private state. Durable content stays authoritative: a stored state the reading adapter cannot use degrades that one message to provider-neutral conversion with a diagnostic instead of failing the request.这七条是消费者可以依赖的硬约束。第 233 条值得拆开:两个错误路径都合法,但失败类型只有一个 LlmFailure(packages/llm/llm/src/types.ts 39 至 51 行:message、code、status?、providerRetryAfterMs?、requestId?)。agent-loop 只认 finish 里的 failure 字段,agent/request-error 瀑布(agent.ts 373 至 390 行)拿到它可以决定 retry 或放弃。第 234 条把重试边界画清楚:adapter 层一次调用等于一次 provider attempt,库自带的重试要关掉(pi-ai adapter.ts 第 126 至 127 行:注释「one adapter call is one SDK attempt」后紧跟 maxRetries: 0),重试是 agent 层的职责,开新的一轮 durable turn。
不这样做会出什么事?把供应商细节放进 agent-loop 和工具层,事故长这样:循环里出现 provider 名判断,工具层为兼容某家供应商的 arguments 对象形态做二次解析,重试逻辑被每家供应商的错误形状牵着走。这些在 dsh 里都被契约挡住。
词汇表为什么闭一个开一个
StreamChunk 闭、ContentBlockMap 开,随手定的?chunk 是协议,所有消费方(agent-loop、session、UI、回放、装配机)都要 switch 它,新增变体时编译期让所有消费方都亮红灯,防止有人静默漏处理。块是内容,插件按领域扩展。第 96 至 97 行注释给新核心块定了三条线的成本,adapter、UI、compaction,等于把「加内容类型」的价格标签公开标好。
装配机为什么只有一个
chunk 从 adapter 出来经过装配变成 message,这个折叠算法只有一份实现,并且住在 llm 包里。模块头注释(assembler.ts 2 至 4 行)写明:This is the single canonical assembly algorithm used by the agent loop to build an assistant message from a chunk stream while logging the raw chunks for replay fidelity。agent-loop 和任何想复用的消费方调同一个类,日志里的 message 和 UI 折叠出来的 message 不会分裂。keep/drop 决策同时作用内容与元数据,把「截断了就该丢 tool call」和「replay 元数据也要跟着丢」绑成同一个决策。
瀑布作为扩展点
llm/stream 瀑布把「每个模型调用必经的路口」做成可插入点。真实监听者:packages/session/session-title/src/index.ts 第 332 行挂 ctx.on('llm/stream', ...) 观察主请求;packages/test-support/llm-replay/src/index.ts 第 784 行用 llm/stream 直接短路整个调用链,用录制好的回放代替真实请求;packages/session/session-checkpoint-policy/src/index.ts 第 64 行挂检查点策略。retry 包 dsh-llm-retry 则挂在 agent/request-error(packages/llm/llm-retry/src/index.ts 第 9 行引入 RequestErrorAction),因为重试发生在失败之后,属于 agent 层恢复。适配器本身保持薄:只管把一次调用翻译成 chunk,重试、回放、标题、检查点这些横切关注点全在瀑布上叠加。
边界条件剖析
边界一:SSE 流在 [DONE] 之前断掉。deepseek 侧,parseSse(sse.ts 28 至 40 行)在流结束还没看到 [DONE] 时抛 LlmError('SSE stream ended without [DONE]', 'STREAM_CLOSED')(第 39 行);translate 还有兜底:payload 源违反契约时第 184 行抛同样的 STREAM_CLOSED。pi-ai 侧对称:toStreamChunks 事件流没等到 done/error 就结束,第 210 行抛同样的 STREAM_CLOSED。两个翻译器各自防守,因为输入协议不同,但错误词汇相同。
边界二:模型返回 stop 却一个块都没有。translate.ts 110 至 115 行把这种情况转成 finish {kind:'error', code: EMPTY_RESPONSE};pi-ai 侧在 stream.ts 92 至 104 行(mapStopReason 里 stop 且 content 为空同样映射 EMPTY_RESPONSE)。契约第 237 行说明这是可重试的错误:dsh-llm-retry 默认重试它。
边界三:max-tokens 截断。finish reason 是 max-tokens 时,assembler.ts 第 136 至 138 行把每个 tool-call 块标成丢弃,第 139 行过滤;replay envelope 在 145 至 147 行按同一掩码修剪。为什么丢 tool call:截断的调用参数不全,执行它不安全,llm-streaming.md 第 293 行有说明。这同时解释了 agent.ts 第 410 行:if (finish.kind === 'max-tokens') return { kind: 'max-tokens' },本轮结束、不派发。
边界四:请求中途被取消。agent.ts 354 至 371 行的 catch 分支:信号已中止时取 assembler.interruptedBlocks(),只保留非空白的 text/reasoning(assembler.ts 168 至 178 行),以 interrupted: true 的 assistant/message 收尾(第 365 行);一个可见内容都没流出的取消连这条事件都没有(session/types.ts 第 275 行注释:An aborted turn with no such event streamed no visible content)。tool call 不进中断前缀,因为中断先于派发,留一个就得伪造结果。
横向对比
对比点:供应商 API 的抽象厚度。同一个问题,claude-code 与 dsh 各有一个答案。
dsh 侧:抽象是一个明确的契约层。两个 adapter 对同一端点、同一流式响应的两种映射就是活的对照实验。deepseek adapter(packages/llm/llm-deepseek/src/adapter.ts)自己持 fetch:429 至 431 行的 stream() 直接调 streamWithConnection,SSE 帧解析外包(sse.ts 28 至 40 行),wire 词到 chunk 词的翻译在仓库内(translate.ts)。pi-ai adapter(packages/llm/llm-pi-ai/src/adapter.ts)把传输交给库:318 至 320 行的 stream() 调 streamWithSnapshot;310 至 316 行的 prepareCall 先捕获当前 snapshot,stream 闭包绑定这一代;367 行 streamSimple 返回库自己的事件流,映射在 stream.ts 的 toStreamChunks。两个实现要表达同一个响应,都落在同一套 chunk 上;两边的差异(段落式事件对增量式 wire、库内错误事件对 HTTP 错误)都被翻译层吸收。
claude-code 侧:上游公开仓库没有主产品源码,检索不到它的流式消费实现。检索关键词:anthropic SDK、messages.stream、@anthropic-ai/sdk、content_block_delta、client.messages。仓库里只有 plugins/、examples/ 和 scripts/ 下的仓库维护脚本,没有 packages/ 目录;plugins/agent-sdk-dev 只是 SDK 开发文档。能拿到的是 CHANGELOG 里的行为记录。claude-code/CHANGELOG.md 第 19 行有一条:
19- Fixed Bedrock streaming behind proxies that strip the response Content-Type header, which silently doubled billed API calls by re-running every turn non-streaming同文件里还有:282 行修流空闲超时在 Bedrock、Vertex、网关部署上失败;315 行给网关流式响应加 SSE keepalive 心跳;216 行把「API 返回空或畸形响应」的错误改成说清内容类型、body 种类、大小、request ID;530 行修 claude -p 在 turn 中途死在流错误时丢掉已经产出的答案。
这是同一个问题的两种答案。上游的答案是把流处理握在产品自己手里,直接对 anthropic SDK 与各家网关的怪癖逐个打补丁(推断:产品源码不在本仓库,内部结构无法确认)。dsh 的答案是把这些问题钉进一个可见的契约:idle timeout 是契约条款(第 235 行),空响应是契约条款(第 237 行),错误归一化是契约条款(第 233 行),每个新 adapter 都要过一遍。代价也对称:dsh 每接一个供应商都要写一段翻译(通常几十行);上游不必维护词汇表,但每个部署形态的怪癖都要产品本体处理,282 与 315 两条同一问题的反复修补就是证据。
本章涉及网络与事件流,但 SSE 帧解析(eventsource-parser)、idle watchdog(dsh-timeout)、fetch 都是跨平台实现,三个平台行为一致,与平台无关。
互动演示设计
形态:格式实验台。一句话结论:供应商的流式字节,经过翻译、事件化、装配三步,变成 UI 能显示、日志能回放的东西。
舞台比喻是海关报关线。轨道一是供应商原文,各国报关单格式各异;翻译机把任何格式译成同一套货物清单,也就是 chunk;轨道二是事件流,assistant/chunk 逐条落档案;装配机把清单折叠成一张完整的入库单,也就是 assistant/message;轨道三是投影面板,UI 与回放都从事件流取数。
逻辑轨迹面板的伪代码行,右侧标真实行号,随动画步进高亮:
wire payload 到达翻译机 translate.ts:101
按 delta 开块并发出增量 translate.ts:133-149
[DONE] 收尾:block-end、usage、finish translate.ts:102-117
循环收到 chunk,追加 assistant/chunk agent.ts:348-352
同一个 chunk 喂给装配机 agent.ts:351
正常结束:blocks 折叠,追加消息事件 agent.ts:392-409
max-tokens:装配机丢弃 tool-call assembler.ts:136-148
中断:只留文字前缀,标 interrupted agent.ts:356-368分六步,每步一句字幕:
第一步:一截供应商字节滑上轨道一,高亮显示这是 DeepSeek 的 SSE payload。字幕:「供应商的字节到了,先不认格式」。
第二步:翻译机逐个词过:reasoning_content 变 reasoning-delta,content 变 text-delta。轨道二同步出现带 index 的 chunk 卡片。字幕:「增量被切成七种 chunk 之一」。
第三步:chunk 卡片在分叉口一分为二,一份进事件档案,一份进装配机。字幕:「同一份 chunk,一边落档案一边进装配机」。
第四步:[DONE] 到达,装配机吐出折叠好的 message,事件档案追加 assistant/message,投影面板出现完整回答。字幕:「块闭合,增量折叠成完整内容」。
第五步:切换供应商到 pi-ai,轨道一换成库事件流,轨道二、三的形状不变。字幕:「换供应商只换轨道一」。
第六步:拨动异常开关。max-tokens 截断时装配机把 tool-call 卡牌翻成灰色丢弃;中断时消息卡牌只保留文字前缀并盖 interrupted 章;把 [DONE] 从序列里抽走时轨道二末端标红 STREAM_CLOSED。字幕:「截断、中断、断流,各有各的投影」。
读者能操作:切换供应商观察轨道一的差异;拖动截断点看 tool-call 何时消失;按中断键看前缀消息;抽走 [DONE] 看错误。
下面是三张图。第一张回答词汇分层:供应商格式经过翻译变成 chunk,chunk 落成 session 事件,事件喂给投影。
第二张是时序图,回答一个 step 里 chunk 从 adapter 到 session 事件的完整旅程。
第三张是流程图,回答 finish reason 五个分支各自的处置。
可迁移结论
一,适配层的最小形态是「一个接口加一个协议枚举」,这个形态不依赖 TypeScript,任何语言都能照做:接口只有一个方法,输入完整请求,输出增量事件的异步序列;协议是七到十个枚举变体,加一条「错误只有两种表达方式」的规则。最小成本形态是约三十行的接口声明加一页协议文档。
二,值得抄的是「翻译与消费分离」:供应商格式只活在 adapter 里,agent-loop、日志、UI 永远消费同一词汇。这对应本章开头的翻车场景:第一个供应商正常、第二个乱套,根因就是翻译和消费写在了一起。
三,超量设计要认得出来。twin adapters 双实现验证、session 级 chunk 保真、ReplayEnvelope 的 per-block 对齐,都是这个体量才合算的。双实现是刻意用双倍维护换词汇表中立性,个人项目或单供应商产品直接省掉;chunk 级事件日志是回放与 UI 共用才需要的粒度,只有一个 UI 的话直接消费装配结果就够了;per-block replay 元数据只在「一个 adapter 要跨模型回放自己产出的响应」时才派得上用场。
思考题
一,为什么 usage 必须在 finish 前发?如果某 adapter 在 finish 之后补发一条 usage,消费方会怎样?提示:装配机的 push 对 usage 不校验顺序(assembler.ts 84 至 87 行直接存),agent-loop 在 finish 后就 append 了 assistant/message 并带走 assembler.usage(agent.ts 400 至 409 行),迟到的 usage 就丢了;派生历史拿不到本次调用的记账。契约原文在 llm-streaming.md 第 231 行。
二,StreamChunk 是闭 union(switch 以 assertNever 收尾),ContentBlockMap 是 merge-extensible,为什么一个闭一个开?提示:chunk 是协议,每个消费方都要 switch 全量处理,新变体必须让所有消费方编译期知道;块是内容词汇,插件按领域扩展,switch 时对未知 type 走默认分支。如果反过来,chunk 可扩展而块封闭,会有什么后果?
三(动手题):不依赖 harness,写一个三十行的翻译器骨架验证错误路径。用你熟悉的语言,把 translate.ts 的骨架复制过来:输入是字符串数组(模拟 SSE payloads,最后一项是 [DONE]),输出逐个打印 chunk 的 type 与关键字段。跑三个用例:正常序列(空 reasoning 开头、text 增量、DONE),看开块与收尾顺序;把 [DONE] 从数组里去掉,看你的实现抛什么,对照 sse.ts 第 39 行的 STREAM_CLOSED;payload 全是空内容后接 DONE,看 finish 的 reason 是什么,对照 translate.ts 110 至 115 行的 EMPTY_RESPONSE。运行命令 node script.mjs 或 python3 script.py。