数据截至 (上游 commit 7538cc96774b)
模型接入层:一个请求怎么发出去,以及订阅引擎这条岔路
30 秒导读: 02 讲的 loop 每一步都要问模型一次。这一章讲那"一次"到底怎么落地——
TurnItem数组怎么变成一个 HTTP body、SSE 字节流怎么变回结构化的 item;以及当用户选了 Claude 订阅时,kun 怎么不跑自己的 loop,改把整轮交给官方 Agent SDK,再把 SDK 的输出翻译回自己的事件契约。
1. 这一章解决的问题(零基础也能懂)
agent 的每一步都是同一个动作:把"到目前为止发生了什么"打包发给模型,听模型说下一步干什么。
听起来只是一次 fetch。真做起来,这一层要同时解决四件事:
| 问题 | 白话 | 谁负责 |
|---|---|---|
| 供应商说不同的方言 | 同一个"发消息 + 带工具"的意思,OpenAI、Anthropic、各家中转各写各的 JSON | CompatModelClient |
| 回来的是字节流不是对象 | 模型边想边吐,工具参数是一个字符一个字符拼出来的 | SSE 解析 + 增量累积 |
| 网络和供应商都不可靠 | 502、429、流卡死、参数被截断成半个 JSON | 重试 / 空闲超时 / 参数修复 |
| 有些用户根本不想按 token 付费 | 手里有 Claude Pro/Max 订阅,想用订阅额度跑 | 订阅引擎(Agent SDK) |
一句话直觉:这一层是 kun 的"出 口海关"。 里面流通的是 kun 自己的 TurnItem;出关要换成对方国家的护照(三种线上协议之一);入关要把对方的东西再换回 TurnItem。而订阅引擎那条路,相当于直接借别人的车队跑一趟,回来把行程单翻译成自己的账本。
2. 顶层全景:一个端口,两条出口
kun 的 loop 只认识一个接口 ModelClient(kun/src/ports/model-client.ts:113),它只有三个成员:provider、model、stream(request)。loop 不知道下面接的是 HTTP 还是别的东西。
真正的分叉发生在两个层次:
AgentLoop.runTurn
|
┌───────────────┴────────────────┐
│ ① 线程 provider 是 agent-sdk? │ kun/src/loop/agent-loop-turn-lifecycle.ts:88-97
└───────┬───────────────┬────────┘
是 否
│ │
┌───────────▼──────┐ ┌────▼──────────────────┐
│ 订阅引擎 │ │ MultiProviderModelClient│ 按 providerId 选客户端
│ SDK 拥有 loop │ └────┬──────────────────┘
│ kun 注入大脑 │ │
└───────────┬──────┘ ┌────▼──────────────────┐
│ │ CompatModelClient │ 自己拼 HTTP / 解 SSE
│ └────┬──────────────────┘
└───────┬───────┘
▼
同一份 kun 运行时事件
怎么读这张图:从上往下是"一个 turn 的出口选择",两条路最后汇到同一种事件——GUI 完全看不出这一轮是自己发的 HTTP 还是 SDK 跑的。
部件职责:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
ModelClient | 唯一的出口端口,loop 只依赖它 | kun/src/ports/model-client.ts:113 |
MultiProviderModelClient | 按请求上的 providerId 挑一个 HTTP 客户端,挑不到就用默认的 | kun/src/adapters/model/multi-provider-model-client.ts:16 |
CompatModelClient | 主力 HTTP 客户端:拼 body、发请求、解流、算钱 | kun/src/adapters/model/compat-model-client.ts:33 |
ModelEndpointFormat | 三种线上形态的枚举与路径推导 | kun/src/contracts/model-endpoint-format.ts:1 |
AgentSdkRuntime | 订阅引擎:把整轮委托给官方 Agent SDK | kun/src/runtime/agent-sdk/agent-sdk-runtime-core.ts:56 |
SdkEventMapper | 把 SDK 消息流反投影成 kun 运行时事件 | kun/src/runtime/agent-sdk/sdk-event-mapper.ts:113 |
装配点分散在 runtime-composition 系列:kun/src/server/runtime-factory-model.ts:127 造默认客户端;kind: 'agent-sdk' 的 provider 不造 HTTP 客户端,由 agentSdkProviderIdsForOptions(runtime-factory-model.ts:303)记进名单;各 HTTP 客户端包进 MultiProviderModelClient(kun/src/server/runtime-composition-model.ts:228);订阅引擎在 kun/src/server/runtime-composition-agent.ts:216 构造。
3. 协议层:端口长什么样
3.1 请求:ModelRequest 是一份"这一轮的全部原料"
ModelRequest(kun/src/ports/model-client.ts:22)不是一个 messages 数组,而是分好层的原料,让适配器自己决定怎么摆:
| 字段 | 是什么 | 为什么单列出来 |
|---|---|---|
systemPrompt | 字节稳定的系统提示 | 必须逐字节不变,否则缓存前缀作废(见 03) |
modeInstruction | 模式说明(如 Plan 模式) | 紧跟 systemPrompt 之后,不污染前缀本身 |
prefix / history | 不可变前缀 + 会话历史 | 前缀在前、历史在后,顺序即缓存边界 |
contextInstructions | 每轮都变的动态指令(目标预算、待办、记忆) | 必须排在历史之后,否则每次计数器一动就打穿整段缓存 |
tools | 本轮广告出去的工具目录 | 见 04 |
providerId | 可选的 provider 覆盖 | 让 workflow / 定时任务在同一个 kun 进程里换供应商 |
abortSignal | 中断信号 | 用户点停、空闲超时都走它 |
3.2 响应:ModelStreamChunk 只有七种
kun/src/ports/model-client.ts:9-16 定义了适配器唯一被允许吐出的东西:
| chunk | 含义 |
|---|---|
assistant_text_delta | 正文增量 |
assistant_reasoning_delta | 思维链增量 |
tool_call_delta | 工具调用参数的增量字符串 |
tool_call_complete | 一个工具调用拼完了,参数已是对象 |
usage | token / 缓存 / 成本快照 |
completed | 本次响应结束,带 stopReason |
error | 出错了,带可选 code |
这七种就是全部契约。 不论下面是 chat completions 的 delta.tool_calls、Responses 的 response.function_call_arguments.delta、还是 Anthropic 的 input_json_delta,到了 loop 眼里都长一个样。
3.3 三种线上形态,加一个"你自己填全路径"
kun/src/contracts/model-endpoint-format.ts:1 只有四个值:
| 形态 | 请求体长相 | 路径后缀 |
|---|---|---|
chat_completions | { model, messages, tools:[{type:'function',...}] } | chat/completions |
responses | { model, input, tools:[{type:'function', name, ...}] } | responses |
messages | { model, system, messages, tools:[{name, input_schema}] } | messages |
custom_endpoint | 不是第四种协议,是"baseUrl 就是完整 URL" | 原样使用 |
custom_endpoint 是个巧妙的偷懒:很多中转站的地址不是标准的 /v1/...(例如智谱 Coding Plan 的 https://open.bigmodel.cn/api/coding/paas/v4/chat/completions)。这时用户直接把完整 URL 填进 baseUrl,resolveModelEndpointFormat(kun/src/contracts/model-endpoint-format.ts:71)再从 URL 尾巴反推真实协议:
// 示意,非源码 —— inferModelEndpointFormatFromUrl 的核心判断
if (path.endsWith('/chat/completions') || path.endsWith('/completions')) return 'chat_completions'
if (path.endsWith('/responses')) return 'responses'
if (path.endsWith('/messages')) return 'messages'
return null // 反推不出来 → 直接报错,不猜
反推失败会在 streamInner 开头就吐 error 并附上"必须以 /chat/completions、/completions、/responses 或 /messages 结尾"的原话(compat-model-client.ts:62-69)。真实实现见 model-endpoint-format.ts:62 的 inferModelEndpointFormatFromUrl。
URL 拼装另有一段针对现实的补丁:buildModelEndpointUrl(kun/src/adapters/model/compat-model-support.ts:44)会识别 baseUrl 结尾是 /beta(DeepSeek 的 beta 域)或已经带了 /v1、/v2,避免拼出 /v1/v1/chat/completions 这种废 URL。
3.4 协议还能按"单个模型"覆盖
endpointFormatForModel(kun/src/adapters/model/compat-model-client-base.ts:93)先问模型的能力元数据要 endpointFormat,再退回 provider 级配置。注释点明了动机:同一个供应商(如 OpenCode Go)可能一部分模型走 chat completions、另一部分走 Anthropic Messages。协议选择因此是 per-model 的,不是 per-provider 的。
4. 出站:一次请求怎么被拼出来
4.1 消息顺序就是缓存策略
collectMessages(kun/src/adapters/model/compat-model-client-base.ts:281)现在是个薄壳,真正的分段在 CompatMessageProjector.project(kun/src/adapters/model/compat-message-projector.ts:39-73),顺序本身是设计:
① systemPrompt (字节稳定,缓存锚点)
② threadProfileInstruction (线程 persona,仍是 system)
③ modeInstruction (模式说明,仍在前段)
④ prefix + history → repairModelHistoryItems() → itemsToMessages()
⑤ contextInstructions (每轮都变的:目标/待办/记忆/技能)
⑥ 附件挂到最后一条 user 消息,最后 healToolMessagePairs + normalizeThinkingAssistantMessages
第 ⑤ 段的位置是本层最值得记住的一条经验。源码注释直白写着:请求级上下文刻意按时序排在历史之后,不进供应商稳定前缀(compat-message-projector.ts:57-60)——它每轮都变,一旦放在历史之前,整段对话的供应商前缀缓存每一步都作废。(顺带:目标指令里那个"已用 token"计数器如今干脆不存在了——用量被刻意移出指令以保住缓存,见 02 §7.2;"压到历史尾巴"这条纪律对剩下的易变内容依然成立。)
Anthropic 形态下同一条规则再实现一遍:messagesToAnthropic(kun/src/adapters/model/compat-request-builder.ts:109)发现一条 system 消息出现在已有对话之后,就不往顶层 system 块里塞,而是用 appendTrailingInstruction(:208)挂进最后一个 user 轮次里(:119-131)。
4.2 同一份"历史修复"服务两个目的
第 ③ 段调用的 repairModelHistoryItems(kun/src/domain/model-history-repair.ts:11)是 03 里为缓存服务的那份修复函数,这里被原封不动复用为 400 防御。
原因是同一条事实:kun 的 TurnItem 里混着 GUI 才关心的东西(审批、用户输入、思维块),而供应商 API 严格得多——每个 assistant 的 tool_call 块后面必须紧跟数量一致的 tool_result。修复函数把不合规的组合剔掉,缓存一致性和"不吃 400"于是是同一件事的两面。
出站路径上还叠了两道同构的保险:
toolCallBlockToMessages(kun/src/adapters/model/compat-message-projector.ts:135)把连续的多个 tool_call 折成一条 assistant 消息 + N 条 tool 消息;只要有一个 callId 没配到结果,整块返回null被丢弃(:193-195)。healToolMessagePairs(:438)在消息层面再扫一遍:孤儿role:'tool'消息直接丢,配不齐的 assistant 工具消息连同结果一起丢。
为什么要做两遍? 一遍在 item 层(懂 kun 的语义,知道哪些是"桥接项"),一遍在 message 层(懂线上协议的硬约束)。任一层单独都堵不死。
4.3 三种 body 的分岔
buildRequestBody(kun/src/adapters/model/compat-model-client-base.ts:249)把消息拼好后交给 CompatRequestCodecs.build(kun/src/adapters/model/compat-request-codecs.ts:101)分三路:
| 形态 | 入口 | 关键差异 |
|---|---|---|
| chat completions | chatCompletions compat-request-codecs.ts:110 | stream_options:{include_usage:true} 才拿得到 usage(:121) |
| responses | case 'responses' :103 | input 而非 messages;max_output_tokens;工具是扁平的 {type,name,...} |
| messages | case 'messages' :105,body 在 :261-290 | max_tokens 必填;system 单独成块并打 cache_control |
Anthropic 那一路藏着一个真实教训。DEFAULT_MESSAGES_MAX_TOKENS = 8192 与 DEFAULT_MESSAGES_REASONING_MAX_TOKENS = 32_768(compat-request-codecs.ts:89-90)写明:思考 token 和输出 token 共用同一份预算,旧的 4096 默认值会让模型想完之后没剩几个 token,把工具调用的参数截断成非法 JSON。所以开了 thinking 就换更大的默认上限(:266-268)。
4.4 显式缓存断点
Anthropic 系协议的缓存是显式的:只有 cache_control 断点之前的内容才进缓存,而且每请求最多 4 个。applyAnthropicCacheControl(kun/src/adapters/model/compat-request-builder.ts:233)的做法很省:
// 示意,非源码 —— 从后往前打两个断点
let breakpoints = 0
for (let i = messages.length - 1; i >= 0 && breakpoints < 2; i -= 1) {
const content = messages[i].content
if (typeof content === 'string' || content.length === 0) continue
content[content.length - 1].cache_control = { type: 'ephemeral' } // 打在最后一个块上
breakpoints += 1
}
加上 system 块自带的那一个(kun/src/adapters/model/compat-request-codecs.ts:275-277),一共三个断点:一个盖住 system + 工具定义,两个盖住最近两条消息——于是下一步的请求正好能命中上一步写下的前缀缓存。
4.5 思考模式:一个字段翻译成五种方言
reasoningEffort 从 loop 传下来时只是 'off' | 'low' | ... | 'max' 这样的抽象档位。落到线上要看模型能力元数据里的 requestProtocol:
| requestProtocol | 落成什么 | 实现 |
|---|---|---|
deepseek-chat-completions | reasoning_effort +(仅官方 host)thinking:{type} | kun/src/adapters/model/compat-request-reasoning.ts:176 |
glm-chat-completions | GLM 自己的 thinking 开关 | :193 |
mimo-chat-completions | 小米 MiMo 的写法 | :205 |
openai-responses | body.reasoning(走 responsesReasoningForEffort) | :7 |
anthropic-thinking | thinking 块 + 输出预算联动 | :220 |
none | 什么都不加 | applyReasoningEffort :40 的默认支 |
有一条防御值得单拎出来。requiresReasoningRoundTrip(:314)的注释指出:thinking 字段是 DeepSeek 私有扩展,第三方 OpenAI 兼容中转(SiliconFlow、OpenRouter、llama.cpp)看到它会 400 或返回空(issue #26)。所以自动开启只在官方 DeepSeek host 上发生(isDeepSeekHost,kun/src/adapters/model/model-error-probe.ts:7);用户显式选档位才强制走这条路。
Azure 同理:isAzureOpenAiEndpoint(compat-request-reasoning.ts:298)命中就把 thinking 整个关掉(生效处在 compat-request-codecs.ts:124 的 includeThinking 判定)。
开了思考模式还有一个副作用:DeepSeek 协议要求每条 assistant 消息都带 reasoning_content。normalizeThinkingAssistantMessages(kun/src/adapters/model/compat-message-projector.ts:418)给缺失的补一个空格 ' '——因为空字符串会被判非法(reasoningContentOrSpace,:385)。这类"补一个空格"的细节是兼容层的日常。
5. 入站:字节流怎么变回 item
5.1 SSE 循环的骨架
streamSse(kun/src/adapters/model/compat-model-client-stream.ts:335)是一个手写的 SSE 解析器,不依赖任何库:
读一块字节 ──► 拼进 buffer ──► 找 "\n\n" 帧边界
│
┌──────────────┴──────────────┐
data: [DONE] ? JSON.parse
│ │
收尾退出 consumeStreamPayload(按协议分派)
│
累积 pending 工具参数 / 产出 chunk
三个细节:
- 帧边界只认
\n\n,一帧里所有data:行拼接后再解析(:390起)——多行 data 的 SSE 也吃得下。 - JSON 解析失败直接跳过(
:398一带), 不让一个坏帧毁掉整轮。 - 每读一块都过一次看门狗(下一节)。
5.2 空闲超时:比总超时更实用的那种
readStreamChunk(kun/src/adapters/model/compat-model-support.ts:330)把三件事塞进一个 Promise.race:读到数据、被 abort、空闲计时器到点。
// 示意,非源码 —— 三选一的看门狗
const result = await Promise.race([
reader.read().then(r => ({ kind: 'chunk', ...r })),
abortPromise, // 用户点了停
new Promise(res => setTimeout(() => res({ kind: 'timeout' }), idleTimeoutMs))
])
if (result.kind === 'timeout') await reader.cancel('model stream idle timeout')
默认 450 秒(DEFAULT_STREAM_IDLE_TIMEOUT_MS,compat-model-support.ts:12),可由运行时配置覆盖(normalizeStreamIdleTimeoutMs,:187)。
妙在计的是"两块数据之间的间隔",不是整轮总时长。 一个思考 5 分钟但持续吐 token 的模型不会被误杀;一个 TCP 连着但再也不说话的僵死连接会在到点后被判死,并吐出带 code: 'stream_idle_timeout' 的 error(compat-model-client-stream.ts:372-380)。
5.3 工具调用增量:三种协议,一张同构状态表
三种协议吐工具调用的方式完全不同,但客户端用同一组状态容器接住:pendingArguments: Map<callId, {index, name, arguments}> 和 pendingByIndex: Map<index, callId>。
| 阶段 | chat completions | OpenAI Responses | Anthropic Messages |
|---|---|---|---|
| 声明调用 | delta.tool_calls[].function.name | item.type === 'function_call' | content_block_start 里 type:'tool_use' |
| 参数增量 | function.arguments 字符串片段 | response.function_call_arguments.delta | input_json_delta.partial_json |
| 宣告完成 | finish_reason === 'tool_calls'(整批一起) | response.output_item.done(逐个) | content_block_stop(逐个) |
| 实现位置 | chat-completions-stream-decoder.ts | responses-stream-decoder.ts | anthropic-messages-stream-decoder.ts |
最麻烦的是 callId 认领。有的供应商第一帧不给 id 只给 index,有的中途才补上真 id。三个解码器如今共用同一个认领函数 resolvePendingToolCall(kun/src/adapters/model/tool-call-stream-identity.ts:18):
- 认领顺序:显式 id → 按 index 查
pendingByIndex→ "只有一个待定就是它" → 都没有就合成一个__kun_stream_tool_call_index_<index>(tool-call-stream-identity.ts:25-35、:92-98)。 - 真 id 中途到货时,
migratePendingCallId(:59-76)把之前按 index 建的临时条目迁移到真 id 下,不丢已累积的参数;若迟到的真 id 撞上另一个 pending 调用,直接抛ModelStreamProtocolError而不是猜。 - 畸形 id 不再静默兜底:null/空串视为"没给"(继续走 index 或唯一待定认领),非字符串、带控制字符、超长的 id 直接抛
ModelStreamProtocolError(:78-90)。
5.4 收尾兜底:宁可给个坏参数,也不能凭空吞掉一次调用
流结束后有一段"安全网"(kun/src/adapters/model/compat-model-client-stream.ts:494-520),注释写得很清楚:chat completions 分支只在 finish_reason === 'tool_calls' 时结算工具调用;如果供应商用 stop、length 或干脆一个裸 [DONE] 收尾,而参数还挂在 pending 里,这次调用就静悄悄消失了。
于是收尾时强制 flush 所有有名字、未结算的 pending;并且只要 flush 出过东西,就把 stopReason 从 stop 纠正成 tool_calls(:535-546)——"供应商标错了 stop reason"这件事在这里被当作已知事实处理。
5.5 参数修复:半个 JSON 也要救回来
parseToolArguments(kun/src/adapters/model/compat-model-client-base.ts:315)转手给 repairToolArguments(kun/src/adapters/model/tool-argument-repair.ts:6),四级降级:
直接 JSON.parse
└─失败→ 剥掉 ```json 代码围栏
└─失败→ 括号配平提取第一个 {...}
└─失败→ 括号配平提取第一个 [...]
└─全失败→ { __raw: 原始字符串 }
括号配平那步(extractBalanced,:71)会正确跳过字符串里的括号与转义(:79-91),不是幼稚的 indexOf('}')。
最后一档 { __raw }(:29)是刻意的:它不是"修好了",而是把一个模型能看懂的错误递给工具层。工具执行失败 → 报错回到模型 → 模型重试。比抛异常炸掉整轮温和得多。
5.6 usage 与算钱:两套 token 语义
mapUsage(kun/src/adapters/model/compat-model-client-base.ts:306)转手给 normalizeCompatUsage(kun/src/adapters/model/compat-usage-normalizer.ts:6),最核心的一段判定解释了一个容易算错的差异:
- OpenAI 系:
prompt_tokens是总数,缓存命中的那部分在prompt_tokens_details.cached_tokens里另标。 - Anthropic 系:
input_tokens不含缓存读写,真实提示大小要input + cache_read + cache_creation。
代码用 anthropicUsage 这个判定(compat-usage-normalizer.ts:29-32)区分两套语义再统一成 UsageSnapshot(kun/src/contracts/usage.ts:11)。搞错了会让缓存命中率显示成假的。
成本估算是"供应商报了就用报的,没报才自己算"(同一文件内),自算走 estimateDeepseekCost(deepseek-pricing.ts:80,分 hit/miss 两档单价)或 estimateMiniMaxCost(minimax-pricing.ts:116,分 input/cacheRead/cacheWrite/output 四档,并有 512K 长上下文加价阈值)。
6. 韧性:重试、错误分类与代理
6.1 两种性质完全不同的重试
| 场景 | 触发条件 | 做法 | 实现 |
|---|---|---|---|
| 网关抖动 | 配置的状态码(默认 [429, 503],kun/src/config/kun-config-runtime.ts:61-66) | 指数退避重发同一 body,默认最多 5 次 | kun/src/adapters/model/compat-model-client.ts:161-174,退避在 compat-retry-policy.ts |
| 参数不被支持 | 400/422 且报文里提到 stream_options/include_usage | 去掉 stream_options 重发一次 | shouldRetryWithoutStreamUsage compat-model-support.ts:172,分支 compat-model-client.ts:262-272 |
第一种能安全重发的理由写在注释里:此时还没有任何响应体被流出去,重发是幂等的。退避期间被 abort 会立刻中止(sleepWithAbort,compat-retry-policy.ts:56)。
第二种是典型的"探测式兼容":很多兼容中转不认 stream_options。与其在配置里让用户勾一个"你的供应商支持 include_usage 吗",不如先试一次,被拒了就降级重发——代价只有一次失败请求,收益是零配置。
6.2 错误分类:给用户可执行的下一步
classifyHttpError(薄壳在 kun/src/adapters/model/compat-model-client-base.ts:218,实现 classifyCompatHttpError 在 kun/src/adapters/model/compat-http-diagnostics.ts:28)不是简单地把状态码丢出来:
| 状态 | 附加信息 | code |
|---|---|---|
| 404 | 加一句"检查 Base URL 与 Endpoint format" | http_404 |
| 429 | 标为限流 | rate_limited |
| 5xx + DeepSeek host | 另发一个探测 请求看端点是否还活着 | deepseek_http_5xx / deepseek_unreachable |
探测逻辑在 probeDeepSeekReachable(model-error-probe.ts:16):把 baseUrl 尾部的 beta/vN 段剥掉再拼 /v1/models 去 GET(:41-56)。区分"DeepSeek 挂了"和"你的网络到不了 DeepSeek"——这两种情况用户要做的事完全不同。
日志侧有对应的卫生要求:logHttpFailure(compat-model-client-base.ts:226)打日志前把 URL 过一遍 redactUrlForLog(compat-http-diagnostics.ts:132),query 里带 key/token/secret/signature/auth/password 的参数一律替换成 [redacted];响应体过 summarizeForLog(:152)截到 1000 字。
6.3 代理只包模型请求
createProxyFetch(kun/src/adapters/model/proxy-fetch.ts:6)在配了 modelProxyUrl 时返回一个替代 fetch,内部用 node:http/node:https + ProxyAgent 手工发请求,再把 Node 的响应流 Readable.toWeb 成 Web ReadableStream 包进 Response(:46-51)——这样上面的 SSE 解析代码一行都不用改。
装配顺序是 config.fetchImpl ?? createProxyFetch(...) ?? fetch(kun/src/adapters/model/compat-model-client-base.ts:77):测试注入优先,其次代理,最后全局。
一个体贴的细节:发请求捕获异常时会判断这是不是 AbortError,只有真的传输失败才追加"