一条 Matrix 消息从「路由入队」到「LLM 流式返回」,中间到底发生了什么?
本文以 AgentTeams 场景下的一段真实运行日志为线索,逐行定位到 OpenClaw 源码,还原 Embedded Agent Runner 的完整执行链路。
一、背景
某天晚上,我给 ClawForge 的 Code Analyst Worker 发了一条消息。Worker 很快回复了。
见《openclaw源码解读(15)》
log2.log日志内容(《源码解读15》中的日志log.log的续篇)
日志来自 AgentTeams 多 Agent 框架ClawForge里的一个 Worker(code-analyst,运行时是 OpenClaw),日志分两段:
log.log:消息从 Matrix 通道进来 → 路由 → 入队 → 触发 hooks(《源码解读(15)》、《源码解读(16)》的日志内容)log2.log:Embedded Agent Runner 真正启动,组装 prompt、调 LLM、流式返回(《源码解读(17)》的日志内容)
log2.log 是 log.log 的续篇——log.log 停在「路由 + 消息入队 + reply_dispatch 钩子」,log2.log 从这里接着往下走,进入Embedded Agent Runner 真正调 LLM 的执行链路(对应你四课体系里的第三、四课)。时间上 log.log 最后是04:39:16.5,log2.log 从04:39:18.7开始,中间那 ~2 秒就是 prompt 组装阶段。
二、全景:一条消息的生命周期
Matrix 消息到达 │ ▼ resolve-route.ts 路由解析:把消息路由到具体 agent(buildAgentSessionKey 等) │ ▼ diagnostic.ts 诊断事件:message received → queued → session state=processing │ ▼ hooks.ts reply_dispatch 钩子(面向切面,runClaimingHook) │ ▼ runs.ts 【Admission 并发控制】登记 active run,reason=run_started │ ▼ attempt-prompt-assembly.ts 组装 prompt + 修复孤儿消息 │ ▼ preemptive-compaction.ts 上下文溢出预检(fits / compact / truncate) │ ▼ openai-completions-transport.ts 钳制 max_tokens │ ▼ provider-transport-fetch.ts POST 请求 → SSE 流式返回 │ ▼ (AgentTeams)mc mirror 会话状态落盘 → mirror 回 MinIO三、逐行拆解:从日志到源码
1. 日志行1. 注入 extraParams:把配置塞进请求
[agent/embedded] applying extraParams to agent streamFn for agentteams-gateway/qwen3.7-plus源码:src/agents/embedded-agent-runner/extra-params.ts:819
function applyPrePluginStreamWrappers(ctx: ApplyExtraParamsContext): void { //... const wrappedStreamFn = createStreamFnWithExtraParams( //... ctx.agent.streamFn, streamParams, ctx.provider, ctx.model, ); if (wrappedStreamFn) { log.debug(`applying extraParams to agent streamFn for ${ctx.provider}/${ctx.modelId}`); ctx.agent.streamFn = wrappedStreamFn; } }逻辑:把配置里给 provider/model/agent 三层设的采样参数(temperature、topP、maxTokens、thinking 等)通过createStreamFnWithExtraParams()包一层,在真正调模型时作为请求options的默认值展开注入({...streamParams, ...options}),请求级 options 可再覆盖。优先级:请求级 > agent 级 > model 级 > provider 默认。
注意区分:这里是options 展开(不改 payload 对象);真正改 payload 的是另一类
streamWithPayloadPatchwrapper(extra_body注入、store删除等),在applyPostPluginStreamWrappers里。
亮点:日志里的provider=agentteams-gateway/qwen3.7-plus说明——OpenClaw 不是直连大模型厂商,而是把请求发给了AgentTeams 的 AI 网关,由网关再代理到 qwen。这是多 Agent 框架常见的「网关统一鉴权/限流/路由」设计。
2. 日志行2-3:登记 active run:Admission 并发控制的落地
2026-08-11T04:39:18.712+00:00 [diagnostic] session state: sessionId=... sessionKey=agent:main:main prev=processing new=processing reason="run_started" queueDepth=1 2026-08-11T04:39:18.712+00:00 [diagnostic] run registered: sessionId=... totalActive=1源码:src/agents/embedded-agent-runner/runs.ts:832-862
export function setActiveEmbeddedRun(sessionId, handle, sessionKey?, sessionFile?) { const previousHandle = ACTIVE_EMBEDDED_RUNS.get(sessionId); const wasActive = previousHandle !== undefined; // ... ACTIVE_EMBEDDED_RUNS.set(sessionId, handle); // ... logSessionStateChange({ sessionId, sessionKey, sessionFile, state: "processing", reason: wasActive ? "run_replaced" : "run_started", }); // ... diag.debug(`run registered: sessionId=${sessionId} totalActive=${ACTIVE_EMBEDDED_RUNS.size}`); }逻辑:ACTIVE_EMBEDDED_RUNS是一个按sessionId索引的 Map,每个 session 同一时刻最多一个 active run。首次登记reason=run_started,若已有 run 在跑则是run_replaced。totalActive=并发 embedded run 数 = Admission 并发控制的落地:一个 session 同一时刻只一个 run。
亮点:这就是 OpenClaw 的Admission 并发控制——防止同一个 session 的多个消息并发触发多个 run 导致上下文串扰。totalActive=1是全局并发 run 数。
3. 日志行4:组装 prompt:embedded run prompt start
[agent/embedded] embedded run prompt start: runId=... sessionId=... provider=agentteams-gateway api=openai-completions endpoint=custom route=proxy-like policy=none源码:src/agents/embedded-agent-runner/run/attempt-prompt-assembly.ts:218
const routingSummary = describeProviderRequestRoutingSummary({ provider: attempt.provider, api: attempt.model.api, baseUrl: attempt.model.baseUrl, capability: "llm", transport: "stream", }); log.debug(`embedded run prompt start: runId=${attempt.runId} sessionId=${attempt.sessionId} ${routingSummary}`);逻辑:进入 prompt 组装阶段。做的事情包括拼 system prompt(含 model identity 行、缓存边界)、启动 prompt cache 观测,最后打出路由摘要。(在代码L76-216行,详情见《源码解读(18)》)
亮点:route=proxy-like是关键——它告诉下游传输层「这是个自定义代理端点」,会触发后面第 8 节 max_tokens 的特殊钳制逻辑。
4. 日志行5:修复孤儿消息:避免连续 user turn
[agent/embedded] Removed already-queued orphaned user message to prevent consecutive user turns. ... trigger=user源码:src/agents/embedded-agent-runner/run/attempt-prompt-assembly.ts:232
const leafEntry = input.orphanRepair?.messageEntry; if (leafEntry && input.orphanRepair) { const messageMergeStrategy = input.orphanRepair.strategy; const orphanPromptMerge = messageMergeStrategy.mergeOrphanedTrailingUserPrompt({ prompt: effectivePrompt, trigger: attempt.trigger, leafMessage: leafEntry.message, }); // ... const action = input.orphanRepair.removeLeaf ? orphanPromptMerge.merged ? "Merged and removed" : "Removed already-queued" : "Preserved"; }逻辑:当一条新 user 消息进来,但会话叶子节点已经有一条「trailing 的 user 消息」时,直接把它合并/移除。
背景:很多 LLM 的 chat 协议要求user/assistant严格交替,出现两条连续的 user turn 会被拒绝。所以 OpenClaw 在组装阶段就「修复」这种结构。
5. 日志行6:上下文诊断:[context-diag] pre-prompt
[agent/embedded] [context-diag] pre-prompt: ... messages=6 roleCounts=assistant:3,toolResult:1,user:2 historyTextChars=1539 maxMessageTextChars=644 systemPromptChars=26492 promptChars=1711 ...源码:src/agents/embedded-agent-runner/run/attempt-prompt-observability.ts:169
const sessionSummary = summarizeSessionContext(input.sessionMessages); // ... log.debug( `[context-diag] pre-prompt: sessionKey=... messages=${input.sessionMessages.length} ` + `roleCounts=${sessionSummary.roleCounts} historyTextChars=${sessionSummary.totalTextChars} ...`, );逻辑:summarizeSessionContext()汇总会话上下文(消息数、角色分布、文本字符数、图片块数),发出context.assembled诊断事件 + debug 日志。
亮点:几个数字很有意思——
| 字段 | 值 | 解读 |
|---|---|---|
systemPromptChars=26492 | ~26KB | 系统 prompt 巨长(agent 的 system 配置 + skills) |
promptChars=1711 | — | 真正要发的用户 prompt 很短 |
sessionFile=sqlite:.../code-analyst/.openclaw/.../sessions.json | — | 会话状态持久化到 worker 的本地文件系统 |
6. 日志行7:上下文溢出预检:编译期预算检查
[agent/embedded] [context-overflow-precheck] ... route=fits estimatedPromptTokens=10209 promptBudgetBeforeReserve=130000 overflowTokens=0 toolResultReducibleChars=0 reserveTokens=20000 effectiveReserveTokens=20000 contextTokenBudget=150000 messages=6 ...源码:src/agents/embedded-agent-runner/run/preemptive-compaction.ts:409
let route: PreemptiveCompactionRoute = "fits"; if (overflowTokens > 0) { if (toolResultReducibleChars <= 0) { route = "compact_only"; } else if (toolResultReducibleChars >= truncateOnlyThresholdChars) { route = "truncate_tool_results_only"; } else { route = "compact_then_truncate"; } }逻辑:这是「编译期」的预算检查,在真正发请求前估算输入 token,决定要不要压缩上下文。决策公式:
contextTokenBudget = 150000 (模型上下文上限) reserveTokens = 20000 (预留给输出) promptBudgetBeforeReserve = 150000 - 20000 = 130000 estimatedPromptTokens = 10209 10209 < 130000 → overflowTokens = 0 → route = fits(不压缩)亮点:这是 OpenClaw 防上下文溢出(OOM)的第一道防线。如果超了,有三条降级路线:
compact_only:直接压缩历史truncate_tool_results_only:只截断 tool 结果(因为 tool result 通常最占空间、最可压缩)compact_then_truncate:先压历史再截 tool result
7. 日志行8:生命周期事件:agent start
[matrix] embedded run agent start: runId=...源码:src/agents/embedded-agent-subscribe.handlers.lifecycle.ts:40
export function handleAgentStart(ctx: EmbeddedAgentSubscribeContext) { ctx.log.debug(`embedded run agent start: runId=${ctx.params.runId}`); emitAgentEvent({ ... stream: "lifecycle", data: { phase: "start", startedAt: Date.now() } }); }逻辑:发出一个lifecycle流的phase=start事件,标记 agent 生命周期开始。日志前缀[matrix]说明这条是从 Matrix 通道订阅链路进来的。
8. 日志行9:钳制 max_tokens:clamp_max_tokens
[openai-transport] [completions] clamp_max_tokens provider=agentteams-gateway api=openai-completions model=qwen3.7-plus requested=128000 output=126533 effectiveContext=150000 estimatedInput=23466源码:src/agents/openai-completions-transport.ts:1804
if (compatDetection.capabilities.usesExplicitProxyLikeEndpoint && clampedMaxTokens !== undefined && effectiveContextTokens !== undefined) { const estimatedInputTokens = estimateOpenAICompletionsInputTokens(params); const remainingBudget = Math.max(1, effectiveContextTokens - estimatedInputTokens - 1); if (clampedMaxTokens > remainingBudget) { clampedMaxTokens = remainingBudget; // 打出 clamp_max_tokens 日志 } }逻辑:这是第二个钳制分支,专门针对proxy-like端点(第 3 节里的route=proxy-like):
remainingBudget = 150000 - 23466 - 1 = 126533 requested(128000) > remainingBudget(126533) → 钳到 126533亮点:注意estimatedInput=23466比第 6 节的预检值10209大了不少——因为这里是发送前的最终估算,已经包含了 tools 定义、system prompt 等所有实际会塞进请求的内容。
9. 日志行10-11:真正发请求:[model-fetch]
[provider-transport-fetch] [model-fetch] start provider=agentteams-gateway api=openai-completions model=qwen3.7-plus method=POST url=http://agentteams-controller:8080/v1/chat/completions ... [provider-transport-fetch] [model-fetch] response provider=agentteams-gateway api=openai-completions model=qwen3.7-plus status=200 elapsedMs=1247 contentType=text/event-stream; charset=utf-8源码:src/agents/provider-transport-fetch.ts:866 / 894
emitModelTransportDebug(log, `[model-fetch] start provider=${model.provider} api=${model.api} model=${model.id} ` + `method=${...} url=${formatModelTransportDebugUrl(rawUrl)} ...`); // ... fetch ... emitModelTransportDebug(log, `[model-fetch] response ... status=${response.status} elapsedMs=... ` + `contentType=${response.headers.get("content-type") ?? ""}`);逻辑:真正发 HTTP 请求的地方。1.2 秒后收到status=200,contentType=text/event-stream(SSE 流式返回)。
四、核心知识点:上下文预算的三次估算
整条链路里最值得记住的,是上下文预算的「三次估算」层层递进:
| 阶段 | 位置 | 估算值 | 作用 |
|---|---|---|---|
| 预检 | preemptive-compaction.ts | estimatedPromptTokens=10209 | 决定要不要压缩(fits/compact/truncate) |
| 发送前 | openai-completions-transport.ts | estimatedInput=23466 | 含 tools/system,用于钳 max_tokens |
| 钳制 | openai-completions-transport.ts | remainingBudget=126533 | 保证「输入 + 输出」不超上下文 |
这三个数字关系:预检最粗(只算消息),发送前最细(算全量),钳制是兜底(强制不超过预算)。这种「粗筛 → 精算 → 兜底」的防御式设计,是工程上非常成熟的防 OOM 思路。
五、日志行12-25:番外:AgentTeams 的 MinIO 数据流
log2.log的尾巴有一段不是 OpenClaw 的日志,而是MinIO 客户端mc的mirror输出:
/root/agentteams-fs/agents/code-analyst/HEARTBEAT.md -> agentteams/agentteams- storage/agents/code-analyst/HEARTBEAT.md /root/agentteams-fs/agents/code-analyst/.openclaw/agents/main/agent/openclaw-agent.sqlite- wal -> agentteams/agentteams-storage/...源码:agentteams-controller/internal/oss/minio.go:137
func (c *MinIOClient) Mirror(ctx context.Context, src, dst string, opts MirrorOptions) error { // ... args := []string{"mirror", src, dst} if opts.Overwrite { args = append(args, "--overwrite") } _, err := c.runMC(ctx, args...) return err }逻辑:AgentTeams 的 Controller 在 reconcile 时,把 worker 的本地文件系统/root/agentteams-fs/agents/<worker>/镜像回 MinIO 的agentteams/agentteams-storage/agents/<worker>/。
设计意图:OpenClaw 的会话/记忆状态(SQLite + WAL)落在 worker 本地文件系统,AgentTeams 再把它 mirror 到 MinIO——这样 worker 是「无状态」的,挂了之后可以从 MinIO 恢复重建。这是多 Agent 框架里经典的「本地状态 + 对象存储持久化」架构。
六、总结
本文从一段真实日志出发,还原了 OpenClaw Embedded Agent Runner 的执行链路。核心结论:
入站到 run 启动是严格串行的:一条
dispatchInboundMessage流水线贯穿始终,日志里的函数大多是「诊断探针」(被主流程直接调用),模块间用 hooks(观察者模式)解耦——函数之间没有调用关系 ≠ 多线程并行,lane队列和 SQLite 锁正是串行化的证据。Admission 并发控制:
ACTIVE_EMBEDDED_RUNS按 sessionId 去重,一个 session 同一时刻只一个 run。防御式上下文管理:三次预算估算(预检 → 精算 → 兜底钳制)+ 孤儿消息修复,保证请求结构合法、不超上下文。
网关透明代理:
route=proxy-like表明 OpenClaw 可把请求发往自定义 AI 网关(AgentTeams Controller),而非直连厂商。状态持久化分离:OpenClaw 管执行,AgentTeams 管状态(MinIO mirror)。
对做多 Agent 框架的同学,这条链路里最有参考价值的是串行化设计(lane 队列 + Admission 控制)和上下文预算的三次估算——两者都是「单 Agent 执行引擎」的精华,可以直接平移复用到多 Agent 编排层。
本文是《OpenClaw 源码解读》系列的一篇。
ClawForge是由 AgentTeams + OpenClaw 搭建的,日志前缀([matrix]、[oss]、[watcher])是 AgentTeams 的那套 Go 代码
ClawForge 里 Code Analyst 那个"定位根因"的 Agent,本质干的就是这件事:从日志/报错反查源码找根因。
正在规划《OpenClaw源码解读》书籍,欢迎出版社编辑交流