一条 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),日志分两段:

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 三层设的采样参数(temperaturetopPmaxTokens、thinking 等)通过 createStreamFnWithExtraParams() 包一层,在真正调模型时作为请求 options 的默认值展开注入({...streamParams, ...options}),请求级 options 可再覆盖。优先级:请求级 > agent 级 > model 级 > provider 默认。

注意区分:这里是 options 展开(不改 payload 对象);真正改 payload 的是另一类 streamWithPayloadPatch wrapper(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=200contentType=text/event-stream(SSE 流式返回)。


四、核心知识点:上下文预算的三次估算

整条链路里最值得记住的,是上下文预算的「三次估算」层层递进

阶段位置估算值作用
预检preemptive-compaction.tsestimatedPromptTokens=10209决定要不要压缩(fits/compact/truncate)
发送前openai-completions-transport.tsestimatedInput=23466含 tools/system,用于钳 max_tokens
钳制openai-completions-transport.tsremainingBudget=126533保证「输入 + 输出」不超上下文

这三个数字关系:预检最粗(只算消息),发送前最细(算全量),钳制是兜底(强制不超过预算)。这种「粗筛 → 精算 → 兜底」的防御式设计,是工程上非常成熟的防 OOM 思路。


五、日志行12-25:番外:AgentTeams 的 MinIO 数据流

log2.log 的尾巴有一段不是 OpenClaw 的日志,而是 MinIO 客户端 mcmirror 输出

/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 的执行链路。核心结论:

  1. 入站到 run 启动是严格串行的:一条 dispatchInboundMessage 流水线贯穿始终,日志里的函数大多是「诊断探针」(被主流程直接调用),模块间用 hooks(观察者模式)解耦——函数之间没有调用关系 ≠ 多线程并行,lane 队列和 SQLite 锁正是串行化的证据。

  2. Admission 并发控制ACTIVE_EMBEDDED_RUNS 按 sessionId 去重,一个 session 同一时刻只一个 run。

  3. 防御式上下文管理:三次预算估算(预检 → 精算 → 兜底钳制)+ 孤儿消息修复,保证请求结构合法、不超上下文。

  4. 网关透明代理route=proxy-like 表明 OpenClaw 可把请求发往自定义 AI 网关(AgentTeams Controller),而非直连厂商。

  5. 状态持久化分离:OpenClaw 管执行,AgentTeams 管状态(MinIO mirror)。

对做多 Agent 框架的同学,这条链路里最有参考价值的是串行化设计(lane 队列 + Admission 控制)上下文预算的三次估算——两者都是「单 Agent 执行引擎」的精华,可以直接平移复用到多 Agent 编排层。


本文是《OpenClaw 源码解读》系列的一篇。

ClawForge是由 AgentTeams + OpenClaw 搭建的,日志前缀([matrix][oss][watcher])是 AgentTeams 的那套 Go 代码

ClawForge 里 Code Analyst 那个"定位根因"的 Agent,本质干的就是这件事:从日志/报错反查源码找根因

正在规划《OpenClaw源码解读》书籍,欢迎出版社编辑交流

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐