跳到主要内容

数据截至 (上游 commit 8d6cbee1b527)

回复流水线:指令、运行队列与分块投递

30 秒导读: 上一章把入站消息归一化并算出了会话键(03)。本章接手之后的全部编排:这条消息要不要进模型(指令快路径)、现在能不能跑(每会话运行队列)、模型吐出来的东西怎么变成一条条聊天消息(分块流式投递)、投失败了怎么办(pending final 兜底)。这是 src/auto-reply/ 这个子系统的本体。


1. 先把坐标定下来

1.1 本章负责的那一段

一条消息从进门到出门要过三个大关,本章是中间那一关:

[03 入站与会话] >>> 本章 04 <<< [05 智能体运行时]
通道归一化 dispatch 编排 模型循环
安全闸 指令层 失败转移
会话路由 运行队列 上下文压缩
──────────► SessionKey ──────────► 提示词 ──────────►
◄────────── 文本流 ◄──────────
分块 / 去重 / 投递

1.2 一句话定义

回复流水线 = 一个「一进多出」的编排器。

进来的是一条入站消息;出去的可能是 0 条、1 条,也可能是几十条——工具进度、流式文本块、最终答案、TTS 音频,全都是独立的出站消息。这个「一变多」以及「多条之间的顺序、去重、失败重试」,就是本章要讲的全部难点。

1.3 三个必须先记住的名词

这三个词在源码里一词一义,后文不换说法:

名词是什么定义处
dispatcher出站投递器。持有一条串行发送链,暴露 sendToolResult / sendBlockReply / sendFinalReply 三个入口src/auto-reply/reply/reply-dispatcher.types.ts:53 ReplyDispatcher
ReplyPayload通道无关的一条出站消息(文本 + 媒体 + 呈现 + 投递偏好)src/auto-reply/reply-payload.ts:33
run(运行)一次模型调用的完整生命周期,按 sessionId 注册在运行注册表里src/auto-reply/reply/reply-run-registry.ts

2. 顶层全景

2.1 主链路

从左到右是调用顺序;虚线箭头是回流(模型产出经 dispatcher 反向出站)。

通道插件 / 网关 RPC / 心跳


① dispatch 入口层 src/auto-reply/dispatch.ts
裸 / 缓冲 / 指定 三选一


② dispatchReplyFromConfig reply/dispatch-from-config.ts(装配)+ dispatch-from-config.*.ts(分段实现)
钩子、抑制策略、回调装配


③ getReplyFromConfig reply/get-reply.ts
指令快路径 → 指令解析 → 会话状态


④ runPreparedReply reply/get-reply-run.ts
队列裁决(跑 / 插队 / 排队 / 丢)


⑤ runReplyAgent ─────────► [05 模型循环] (reply/agent-runner-run.ts)

╎ 流式回调

⑥ 分块流水线 reply/block-reply-pipeline.ts
切分 → 合并 → 去重


⑦ dispatcher.sendXxx ──► deliver() ──► 通道

2.2 部件一句话职责

部件干什么在哪个文件
dispatch 入口建/收 dispatcher,装前置钩子,排前台回复租约src/auto-reply/dispatch.ts
dispatch-from-config抑制策略、进度回调装配、最终投递与计数(已拆成 dispatch-from-config.*.ts 一族)src/auto-reply/reply/dispatch-from-config.ts(门面)
get-reply指令快路径、指令解析、会话状态初始化src/auto-reply/reply/get-reply.ts
get-reply-run提示词组装 + 队列裁决 + 交给 agent runnersrc/auto-reply/reply/get-reply-run.ts
命令注册表斜杠命令的唯一真源(原生 + 文本别名)src/auto-reply/commands-registry*.ts
跟进队列会话忙时把消息排队、去重、限额、排水src/auto-reply/reply/queue/
reply-dispatcher串行发送链、人类延迟、计数、幂等src/auto-reply/reply/reply-dispatcher.ts
分块流水线块切分、合并、去重键、超时熔断src/auto-reply/reply/block-reply-*.ts

2.3 走一遍主线(不进代码)

  1. 通道收到「帮我看下这个 PR」,调 dispatchInboundMessageWithBufferedDispatcher
  2. 入口层现建一个带 typing 的 dispatcher,并给它挂上 message_sending 等插件钩子。
  3. getReplyFromConfig 先问:这是不是纯斜杠命令?是就直接返回文本,不碰模型
  4. 不是。于是初始化会话、解析行内指令(/think high 之类),拼提示词。
  5. runPreparedReply 查运行注册表:这个会话正在跑吗?跑着就 steer 或排队;空闲就开跑。
  6. 模型开始吐字。每个文本块经切分器 → 合并器 → 去重 → dispatcher.sendBlockReply()
  7. 模型结束。最终答案经 sendFinalReply() 出站;出站前写一份 pendingFinalDelivery 存根,成功后清掉。
  8. withReplyDispatcherfinally 排空发送链,finalizeDispatchAndAudit 校正计数,返回。

3. dispatch 三层入口

3.1 三个入口的分工

三个导出函数是同一条路的三种起点,区别只在「dispatcher 谁来造、造成什么样」:

入口谁传 dispatcher附加能力典型调用方
dispatchInboundMessage调用方自己传现成的无(最底层)已有 dispatcher 的复用场景
dispatchInboundMessageWithDispatcher入口按 options 现造一个裸的全局发送钩子简单通道
dispatchInboundMessageWithBufferedDispatcher入口现造一个带 typingtyping 控制 + 前台回复租约 + onSettled 兜底绝大多数聊天通道

依据:src/auto-reply/dispatch.ts:199 dispatchInboundMessage:452 dispatchInboundMessageWithDispatcher:371 dispatchInboundMessageWithBufferedDispatcher

3.2 共同的骨架:不管怎么进,都要排空

三个入口最后都汇进 dispatchInboundMessage,而它把真正的工作包在 withReplyDispatcher 里(src/auto-reply/dispatch-dispatcher.ts:46):

try {
return await params.run();
} finally {
await settleReplyDispatcher(params); // markComplete() + waitForIdle()
}

这三行是整个投递层的可靠性地基:无论 run() 正常返回还是抛异常,dispatcher 一定被标记完成并等到发送链排空。没有这个 finally,一次模型异常就会把已入队但未发出的块永久卡住。

3.3 前台回复租约:同一场对话的可见投递按 FIFO 排队

缓冲入口独有的一个机制。场景:用户连发两条消息,两次 dispatch 并发跑,旧的那次慢一拍才把答案发出来——用户会看到过期答案盖在新答案后面

解法是一把按会话复合键的 FIFO 租约(不是锁,不串行化计算,只串行化「可见投递」这一小段):

同一个 [通道, 账号, 会话键, 聊天类型, 目标] 组成租约 key
┌────────────────────────────────────────────────┐
│ dispatch A → 租约排到队尾,拿到自己的完成位 │
│ dispatch B → 排在 A 后面 │
│ │
│ B 的 beforeDeliver: await A 的租约先释放 │
│ A 投递完 → release() → 才轮到 B 的可见内容出场 │
└────────────────────────────────────────────────┘

关键实现三处:

  • foregroundReplyLeasessrc/auto-reply/dispatch.ts:54)是全局单例租约注册表;复合键由 resolveForegroundReplyOrderKey 用 JSON 序列化拼出(:68),避免各段里含分隔符时切错;
  • reserveForegroundReplyLease:95)在派发前领号;
  • 真正等待发生在包了一层的 beforeDeliver 里:await foregroundReplyLease?.wait():325),收尾 release():364)。

租约原语本身极简:createKeyedFifoLeaseRegistrysrc/shared/keyed-fifo-lease.ts:19)只是把每个 key 的「尾 promise」串成链(:34-46)。算的年代替了旧实现里的「代数比较 + 取消过期投递」:不再取消任何回复,而是让顺序本身不可能乱

3.4 计数:queuedFinal 只是「进了队」

enqueue 返回 true 只代表「进了队」,不代表「发出去了」。被 beforeDeliver 取消或投递抛错的载荷,仍然计在 queued 里。所以收尾时要按「实际投递/取消/失败」校正——这部分逻辑集中在 finalizeDispatchAndAuditsrc/auto-reply/reply/dispatch-from-config.finalize.ts:37),它逐路合并各段的 queuedFinal 标志(如 :149:168),并在 final 全部被观察到后才认账。

queuedFinal 是上游判断「要不要走无回复兜底」的信号;如果不做这层校正,一次被吞掉的投递会被误报成「已回复」,用户就彻底收不到东西了。


4. 指令层:能不进模型就不进模型

4.1 为什么要有指令层

/status/new/compact 这类命令的答案完全由本地状态决定,调模型既慢又贵还可能被模型改写。所以流水线在模型前面放了三道快路径,命中即返回。

4.2 命令注册表:一个定义,两种形态

一条命令同时要能被「Telegram 原生斜杠菜单」调用,和被「用户手打 /status」调用。注册表用元组式的 defineBuiltinCommand 一次描述两者(src/auto-reply/commands-registry.shared.ts:175defineChatCommand:60):

// src/auto-reply/commands-registry.shared.ts:259 /status 的定义
defineBuiltinCommand("status", "Show current status.", "status", "essential", {})
// key ─┘ 描述 ─┘ 归类 ─┘ 层级 ─┘

注册表本身按「通道插件注册表版本」缓存(getChatCommands,src/auto-reply/commands-registry.data.ts:45),因为 dock 类命令(/dock_telegram 之类)是按已加载的通道插件动态生成的(defineDockCommand,:19)。

本章点名的几条命令落点:

命令定义行归类谁执行
/statuscommands-registry.shared.ts:259statusbuildStatusReply(经 commands-status.ts)
/new:557sessionmaybeHandleResetCommand(reply/commands-reset.ts:35)
/reset:554session同上
/compact:567sessionreply/commands-compact.ts(extractCompactInstructions :36)
/think:574options行内指令层(见 §4.5)
/verbose:583options行内指令层

/think 的候选值不是写死的字符串数组,而是一个依赖当前 provider/model/catalog 的函数:577 choices: ({ provider, model, catalog, agentRuntime }) => …),这样原生菜单弹出来的档位永远和当前模型能力一致。

4.3 命令鉴权:三个独立的门

resolveCommandAuthorizationsrc/auto-reply/command-auth.ts:458)返回的不是一个布尔值,而是一组事实:

字段含义影响什么
senderIsOwner发信人是 owner(身份匹配、operator.admin 作用域、或 allowAll)高危命令、trace 授权
isAuthorizedSender能执行命令所有斜杠命令
ownerList解析出的 owner 清单传给 agent 运行参数

三个门是分开的,因为它们的来源不同:owner 来自通道插件配置 + cfg.commands.ownerAllowFrom;命令授权可以由 cfg.commands.allowFrom 单独放开;而原生命令回合有一条捷径——通道已经在传输层验过身,所以在没有强制 owner 要求时直接放行(command-auth.ts:535-536):

const nativeCommandAuthorized =
commandAuthorized && isNativeCommandTurn(resolveCommandTurnContext(ctx)) && !requireOwner;

未授权时的行为不是报错,是静默丢弃。 整条消息就是一个控制命令、但发信人没权限时,流水线清 typing、return undefined——这避免了把「本机器人支持哪些命令」泄漏给陌生人。

4.4 快路径一:原生斜杠命令

最靠前的一道,在 workspace 都还没准备好之前就跑(调用点 src/auto-reply/reply/get-reply.ts:482):

// src/auto-reply/reply/get-reply-native-slash-fast-path.ts:78
function shouldRunNativeSlashCommandFastPath(ctx: MsgContext): boolean {
const commandName = resolveNativeSlashCommandName(ctx);
return Boolean(
commandName &&
commandName !== "new" && // /new /reset 需要完整会话生命周期
commandName !== "reset" &&
(isNativeCommandTurn(commandTurn) ||),
);
}

命中后走一个轻量会话状态initFastReplySessionStateget-reply-fast-path.ts:178,只读会话存储、不做 workspace 引导),/status 甚至直接短路到 buildStatusReply,其他命令交给 handleCommands(get-reply-native-slash-fast-path.ts:352)。整个过程不加载模型目录、不建沙箱。

/new/reset 被明确排除——它们要真正轮换 sessionId、触发重置钩子、可能还要重发系统提示,属于会话生命周期操作,必须走完整路径。

4.5 快路径二:只含指令的消息

/think high 单独一条消息发过来时,它只是改状态,没有要问模型的内容。判据很直接——有指令,且剥掉指令和 @提及后正文为空isDirectiveOnly,src/auto-reply/reply/directive-handling.directive-only.ts:8):

if (!directives.hasThinkDirective && !directives.hasVerboseDirective && /* … */) {
return false; // 压根没指令
}
const stripped = stripStructuralPrefixes(cleanedBody ?? "");
const noMentions = isGroup ? stripMentions(stripped, ctx, cfg, agentId) : stripped;
return noMentions.length === 0; // 剥干净后没剩下东西 = 纯指令

纯指令消息由专门的 handleDirectiveOnlysrc/auto-reply/reply/directive-handling.impl.ts:60)生成确认回执(directiveAck);而「指令 + 正文」的混合消息走另一条路——把指令应用到本轮但仍要跑模型(applyInlineDirectiveOverrides,src/auto-reply/reply/get-reply-directives-apply.ts:118)。

4.6 快路径三:命令处理器链

handleCommandssrc/auto-reply/reply/commands-core.ts:34)是一条惰性加载的处理器链

const resetResult = await maybeHandleResetCommand(params); // 重置优先(:51)
if (resetResult) return normalizeCommandHandlerResult(resetResult);

for (const handler of HANDLERS) {
const result = await handler(params, allowTextCommands);
if (result) return normalizeCommandHandlerResult(result); // 命中即停
}
return { shouldContinue: true }; // 没人认领 → 继续走模型

normalizeCommandHandlerResult:20)做了一件容易忽略的事:清掉命令回复的 replyToId / replyToCurrent。命令回执是对系统说话,不该挂在用户那条消息下面变成引用回复。

4.7 四条路径的汇合图

入站消息

├─► 原生斜杠快路径? ──是──► buildStatusReply / handleCommands ──► 返回
│ (非 /new /reset)
├─► 指令解析
│ │
│ ├─► 纯指令消息? ──是──► handleDirectiveOnly ──► directiveAck ──► 返回
│ │
│ └─► 混合消息 ──► 应用指令到本轮,继续

├─► handleInlineActions (含 handleCommands) ──命中──► 返回

└─► runPreparedReply ──► 进模型

5. 每会话运行队列

5.1 问题:模型跑着的时候又来消息了

聊天场景下这几乎是常态。可选的处理有四种,OpenClaw 把它做成了队列模式(queue mode)

模式行为适合
steer把新消息注入正在跑的那次运行(转向)默认;对话式交互
followup排队,当前运行结束后逐条跑顺序敏感的任务
collect排队,结束后合并成一条再跑用户连发碎片消息
interrupt中断当前运行,立刻跑新的强制打断

默认值是 steer(解析在 resolveQueueSettings,src/auto-reply/reply/queue/settings.ts:25 起,默认落在 :35)。解析优先级从高到低:行内 /queue 指令 → 会话持久化设置 → 按通道配置 → 全局配置 → 默认。

其他默认量:debounce 500ms、cap 20 条、溢出策略 summarizesrc/auto-reply/reply/queue/state.ts:51-53)。

5.2 裁决函数:短短几行定了所有事

// src/auto-reply/reply/queue-policy.ts:8 resolveActiveRunQueueAction
if (!params.isActive) return "run-now"; // 没在跑 → 直接跑
if (params.isHeartbeat) return "drop"; // 心跳不排队,忙就丢
if (params.resetTriggered) return "run-now"; // /new /reset 必须立刻生效
if (params.shouldFollowup) return "enqueue-followup";
return "run-now";

四条规则的顺序就是优先级。两个细节值得记:

  • 心跳被丢弃而不是排队。 心跳是「定时看看有没有事」,排队等半小时后再跑毫无意义,还会堆积。
  • 重置命令强制 run-now,且上游会把队列模式改写成 interrupt——用户说「重开」就得立刻重开。

5.3 前置布尔量与线程守卫

shouldSteer / shouldFollowup 的裁决在 src/auto-reply/reply/get-reply-run-admission.ts(裁决入口 resolveActiveRunQueueAction 的调用点 :567 区域)。其中有一条通道特定守卫:Slack 的私聊也能开线程,两个线程其实是两个对话,不能互相转向——判据是 isSlackDirectRoutedThreadTurnget-reply-run-admission.ts:457)+ routeThreadIdsMatch:474,实现 get-reply-run-helpers.ts:117)。

shouldFollowupsteer 模式下也为真——因为转向可能失败(模型已经不在流式阶段了),那时要退回排队。

5.4 run-now 但确实在跑:先等它死透

activeRunQueueAction === "run-now"isActive(典型是 interrupt 模式)时,不能直接开新的,得先收尸(resolvePreparedReplyQueueState,src/auto-reply/reply/get-reply-run-queue.ts:16):

interrupt? → abortActiveRun() 先请它中止
→ waitForActiveRunEnd() 再等它结束
→ refreshPreparedState() 状态可能已变,重算
→ 还在跑? → 回复 REPLY_RUN_STILL_SHUTTING_DOWN_TEXT

等不到就明确告诉用户「上一轮还在关闭中,稍后再试」(常量 :13,产生点 :49),而不是静默失败或强杀。

5.5 转向通道

转向只在模型正在流式输出时可行。注入走 queueEmbeddedAgentMessageWithOutcomeAsyncsrc/agents/embedded-agent-runner/runs.ts:475),裁决与分支在 runReplyAgentsrc/auto-reply/reply/agent-runner-run.ts:64effectiveShouldSteer 的合成在 :138、注入分支在 :265 起)。

转向成功时本次 dispatch 不产生任何出站消息,答案会从被转向的那轮里出来;注入被拒则记日志、落到排队路径。

5.6 排队与去重

enqueueFollowupRunsrc/auto-reply/reply/queue/enqueue.ts:124)叠了两层去重

机制防什么
近期消息 id 缓存全局 TTL 缓存,5 分钟 / 1 万条(recent-message-ids.ts:5-6通道重投同一条消息
队内比对isRunAlreadyQueuedenqueue.ts:77)按「路由身份 + messageId」比对同一条消息重复入队

去重键不是简单字符串拼接,而是 JSON 元组序列化enqueue.ts:67-74,注释直说是为了避开 to/accountId 里含 | 时的分隔符冲突)。

队列满了按 dropPolicy 处理,默认 summarize:被挤掉的消息不是丢弃,而是压成一行摘要留在队里,还带一个「省略计数」结构(summaryElisionsenqueue.ts:242:254),保证用户至少知道「刚才还说过 N 句」。

5.7 排水

scheduleFollowupDrainsrc/auto-reply/reply/queue/drain.ts,导出经 queue.ts:5)的要点:

const existingQueue = FOLLOWUP_QUEUES.get(key);
if (existingQueue?.draining) {
rememberFollowupDrainCallback(key, runFollowup); // 已在排水:只更新回调(drain.ts:114)
return;
}

为什么要更新回调:正在排水的循环沿用自己的回调,但延迟重试要用刚结束那轮提供的最新会话/运行时上下文

排水循环里 collect 模式最复杂:先检查队里的消息是否跨通道hasCrossChannelItems,调用点 drain.ts:1334),跨了就退化成逐条处理以保住各自的回复路由;不跨才按投递上下文分组合并成一条提示词。

enqueueFollowupRun 最后一段处理一个竞态:排水刚结束并删了队列对象,新消息又创建了一个 draining: false 的新队列,但没人给它排水——所以要 kickFollowupDrainIfIdleenqueue.ts:120:304:362)。

5.8 重启后的排水

进程重启会丢掉内存队列,但被打断的那一轮必须恢复。这由 src/agents/main-session-recovery/ 一族负责:扫会话存储,找 pendingFinalDelivery.kind === "replayable" 的条目(main-session-recovery-state.ts:524),用 sanitizePendingFinalDeliveryText 清洗后重建恢复消息(buildResumeMessage,main-session-restart-dispatch.ts:71),并用 MAX_RECOVERY_RETRIES 做有限次重试(main-session-recovery-state.ts:271:550)。


6. 出站表达:ReplyPayload

6.1 字段族

ReplyPayloadsrc/auto-reply/reply-payload.ts:33)是通道无关的,字段按用途分五族:

字段说明
内容text / mediaUrl / mediaUrls最基本的文本与媒体
呈现presentation / channelData富呈现;通道不支持就降级
投递delivery / replyToId / replyToTag / replyToCurrent置顶、引用回复目标
语音audioAsVoice / spokenText / ttsSupplement语音条 vs 音频文件;TTS 补充
分类标记isError / isReasoning / isReasoningSnapshot / isCommentary / isCompactionNotice / isFallbackNotice / isStatusNotice决定通道要不要显示、要不要计入 TTS(reply-payload.ts:84-95 的注释)

6.2 TTS 补充:一段被延迟到最后的语音

ttsSupplement 标记的载荷含义特殊:文本已经通过流式发出去了,这条只是配套的音频。所以投递前要把可见字段摘掉(buildTtsSupplementMediaPayload,src/auto-reply/reply-payload.ts:223)。

getReplyPayloadTtsSupplement:175)还加了一道守卫:没有媒体就不算 TTS 补充——没音频的「补充音频」是矛盾状态。

对应的生成时机在 dispatch 尾部:如果流式发过块、但最终没有 final 载荷,就用累积的可见块文本合成一条纯音频(needsTtsFallback,src/auto-reply/reply/dispatch-from-config.finalize.ts:34;ACP 路径的累积缓冲见 dispatch-acp-delivery.ts:173-176)。

6.3 元数据:不污染线上形状的旁路

有些信息只给内部用,不能出现在发给通道的对象上。解法是 WeakMap 旁路(src/auto-reply/reply-payload.ts:312):

const replyPayloadMetadata = new WeakMap<object, ReplyPayloadMetadata>();

ReplyPayloadMetadata:240)里几个承重字段:

字段作用
assistantMessageIndex标识这个块属于第几条助手消息;合并器据此不跨消息合并
pendingFinalDeliveryCompletion这条载荷对应的 pending-final 存根坐标(见 §8)
replyDelivery / replyDeliverySource回复策略 + 产生它的路由身份,用于拒绝跨路由的过期策略

代价是:任何克隆载荷的地方都必须手动搬运元数据。所以源码里到处是 copyReplyPayloadMetadata(source, next)

6.4 通道能力降级

降级发生在三个层次:

① 分类抑制 isReasoning === true → 没有专用 reasoning 通道的通路直接跳过
reply-payload.ts:86 的注释 + dispatch-from-config.execute.ts:442

② 长度裁剪 resolveTextChunkLimit(cfg, provider, accountId)
chunk.ts:60 —— 通道配置 > 账号配置 > 插件 textChunkLimit > 默认

③ 切分模式 resolveChunkMode chunk.ts:103
"length"(默认) 只在超限时切;"newline" 优先在段落边界切

7. 流式体验

这一节是本章最密的部分。先看它要同时满足的四个约束:

┌── 要快 ── 用户想尽早看到字

├── 要少 ── 不能一个 token 一条消息刷屏

├── 要有序 ── 工具进度 → 文本块 → 最终答案,顺序不能乱

└── 不能重 ── 流式发过的内容,最终答案里不能再发一遍

7.1 typing 指示:策略与模式两级

第一级(策略) 决定这一轮能不能出 typing(resolveRunTypingPolicy,src/auto-reply/reply/typing-policy.ts:21):

const typingPolicy = params.isHeartbeat ? "heartbeat"
: params.originatingChannel === INTERNAL_MESSAGE_CHANNEL ? "internal_webchat"
: params.systemEvent ? "system_event"
: (params.requestedPolicy ?? "auto");

三类非用户可见的回合一律压掉——心跳时让用户看到「对方正在输入」是纯粹的噪音。

第二级(模式) 决定什么时候点亮(resolveTypingMode,src/auto-reply/reply/typing-mode.ts:24):

模式触发时机什么时候默认用它
never不点心跳 / 系统事件 / 被压制
instant立刻点私聊;群里被 @;message-tool-only
message等到真有消息内容才点群聊未被 @

关灯逻辑有兜底。 正常路径是 dispatcher 排空后触发 onIdle。但如果 dispatcher 那边出了岔子,markRunComplete 里的定时器会强制清理(src/auto-reply/reply/typing.ts:263):

log?.("typing: dispatch idle not received after run complete; forcing cleanup");
cleanup();

7.2 进度草稿:把工具过程折成一条会更新的消息

支持编辑消息的通道(Slack、Telegram)可以不发一串进度消息,而是维持一条不断被改写的草稿。这由 createChannelProgressDraftCompositorsrc/channels/progress-draft-compositor.ts:62)负责:内部维护 lines 数组,渲染成文本,和 lastRenderedText 比对,变了才调 update()

难点是什么时候开始画这条草稿。太早,一次两秒的调用会闪一下就消失;太晚,用户干等。createChannelProgressDraftGatesrc/channels/streaming.ts:635)的规则:

第 1 个工作事件 ──► 起 1.5 秒定时器(DEFAULT_PROGRESS_DRAFT_INITIAL_DELAY_MS,streaming.ts:81)
└─ 1.5 秒内没有第 2 个事件 → 定时器到点才开画
第 2 个工作事件 ──► 立即开画(多步任务,注定要跑一会儿)

闸门还用一个 startPromise 单例挡住并发:定时器、显式启动、第二个工作事件三条路径可能同时到,if (startPromise) { … }streaming.ts:660-685)保证只创建一条草稿。

另有 createStatusReactionControllersrc/channels/status-reactions.ts:202)走的是表情反应这条更轻的通道:给用户那条消息贴 emoji 表示状态,还带软/硬停滞定时器,卡住太久会自动换成「停滞」表情(:285-298)。

7.3 块级切分与合并:两层,别搞混

这是最容易看错的地方。切分(chunking)和合并(coalescing)是两个方向相反的操作,串在一起

模型 token 流


① 切分 chunking 按语义边界把流切成"逻辑块"
minChars / maxChars resolveBlockStreamingChunking
breakPreference block-streaming.ts:161

▼ (逻辑块,可能很碎)
② 合并 coalescing 把碎块攒起来再发,减少消息条数
minChars / maxChars createBlockReplyCoalescer
idleMs block-reply-coalescer.ts:16
joiner: "\n\n" | "\n" | " "

▼ (出站块,条数少)
③ 去重 + 串行发送 createBlockReplyPipeline

两组参数由 resolveEffectiveBlockStreamingConfigblock-streaming.ts:103)统一算出,并强制合并上限不超过切分上限:136)、切分上限被钳到通道文本上限内(:119-126)。joiner 也跟着 breakPreference 走:段落用 \n\n、换行用 \n、句子用空格(:235)。

合并器只在「能合」的时候才把碎块并进缓冲:canMergeBufferedTextWithMediablock-reply-coalescer.ts:109-124)要求两边都不是语音条、都不是 reasoning/commentary/状态通知、且 replyToId 一致——不满足就先冲刷再开新缓冲。还有一个 flushOnEnqueue 逃生口(:26),给「确实需要每条单发」的传输用。

7.4 两个去重键:为什么要两个

流水线维护两套 key(src/auto-reply/reply/block-reply-pipeline.ts:57:72):

key含谁用途
createBlockReplyPayloadKey文本 + 媒体 + 呈现 + channelData + replyToId + 是否状态通知防止同一条被发两次
createBlockReplyContentKey同上但不含 replyToId流式发过之后,压掉内容相同的 final

第二个 key 的动机:流式那条挂了引用回复、最终那条没挂,两者 replyToId 不同但内容一样——不忽略 replyToId 就压不住,用户会看到同一段话两遍。

hasSentPayload:325)还有第三层兜底:如果精确 key 没命中,就把所有已发文本片段拼起来、去掉全部空白再比对(:337-340normalize)。这防的是「流式分了 5 块、最终是 1 整段,只有换行位置不同」的情况。

7.5 有序与熔断

发送顺序靠一条 promise 链保序(block-reply-pipeline.ts:113:147 sendChain = sendChain.then(...))。一旦某块投递超时,整条链熔断(:200):

block reply delivery timed out after {timeoutMs}ms;
skipping remaining block replies to preserve ordering

选择很明确:宁可少发,不能乱序。超时之后剩下的块全部跳过,让最终答案兜底。

7.6 工具事件回显

工具进度走 onToolResult 类回调(装配在 src/auto-reply/reply/dispatch-from-config.execute.ts)。它做的顺序处理值得看(:203 起):

await waitForPendingDirectBlockReplyDelivery(...); // 等直发块落地
// …(已产生副作用后,入站去重不可重放)
await flushPendingCommentaryProgress(); // 先把缓冲的旁白发掉
// …然后才是工具摘要本身

flushPendingCommentaryProgress 出现在多个地方:162:211:433):旁白(模型在工具之间说的话)被缓冲起来,时机到了才冲刷,保证它永远排在它所解释的那件事前面。

块投递路径上还累积两份文本:accumulatedBlockText(判断有没有实质输出)和 accumulatedBlockTtsText(喂给 §6.2 的尾部 TTS),定义在 dispatch-acp-delivery.ts:173-176

7.7 dispatcher 的串行发送链

createReplyDispatchersrc/auto-reply/reply/reply-dispatcher.ts:254)是三个 send 方法的共同底座(sendToolResult / sendBlockReply / sendFinalReply 都汇到同一个 enqueue:601-602)。enqueue 的流程:

归一化(normalizeReplyPayloadInternal,:215)──空──► 返回 false,触发 onSkip

├─► queuedCounts[kind]++ ; pending++

└─► 接到 sendChain 尾部:
人类延迟(仅 block,且非第一块)
beforeDeliver 钩子 ──返回 null──► cancelledCounts[kind]++
deliver()
.catch → failedCounts[kind]++
.finally → onDeliverySettled ; pending--

pending 计数用了一个「预约位」技巧:259-261 注释):初始值是 1 而不是 0。这样在还没有任何回复入队时,dispatcher 也不会被判定为空闲——防止网关在模型刚开始跑时就重启。markComplete():580)通过一次微任务延迟来清预约位(:587 的注释),给正在进行中的 enqueue() 调用留出机会先把 pending 加上去。

人类延迟只对 block 生效且跳过第一块(getHumanDelay 调用点 :473)——第一块要尽快出,后续块之间才需要节奏感。


8. 交付兜底

8.1 pendingFinalDelivery:一份可重放的存根

模型算完了,答案发不出去(网络断、通道限流、进程被杀),这份答案就永久丢了。解法是在投递前先把它写进会话存储。写法在 completeReplyAgentRun 里(src/auto-reply/reply/agent-runner-result-complete.ts:365-405):

await updateSessionEntry({ storePath, sessionKey }, (entry) =>
entry.sessionId === expectedSessionId
? {
pendingFinalDelivery: {
...(resolvedPendingText
? { kind: "replayable", text: resolvedPendingText } // 可重放:带文本
: { kind: "transport-only" }), // 只补投递,不重放文本
intentId, // 这次意图的唯一 id
deliveries: [{ id: deliveryId, state: "prepared" }], // 每条载荷一条投递存根
context: pendingFinalDeliveryContext, // 投递路由快照
createdAt: Date.now(),
},
updatedAt: Date.now(),
}
: null, // 会话中途被 /new 换代 → 不写,旧运行的存根绝不让新会话继承
{ skipMaintenance: true, takeCacheOwnership: true });

三个设计点:

  • 只存策略上可见的文本。 被抑制(sourceReplyPolicy.suppressDelivery)时根本不建存根(:360-363sendableFinalPayloads 过滤);心跳回合还要先剥掉 HEARTBEAT_OK 这类 ack token——纯 ack 没有重放价值。
  • 每条载荷带 intentId / deliveryId:365-382),写进载荷的 WeakMap 元数据(pendingFinalDeliveryCompletion),于是投递成功时可以按 id 精确结算存根,而不是靠文本比对。
  • 写完还要回读校验:发现 sessionId 或 intentId 对不上,直接抛错(:412-417 附近)——宁可报错也不留一份指向错误会话的存根。

8.2 成功后立刻清

投递全部结算后清存根(clearPendingFinalDeliveryAfterSuccess,src/auto-reply/reply/dispatch-from-config.pending-final.ts:25;编排调用点 dispatch-from-config.finalize.ts:173:184)。

顺序是修 bug 修出来的(注释里带 issue #89115,dispatch-from-config.finalize.ts:187):就算外层操作已被中止,清理也必须照跑——否则 pendingFinalDelivery 会永久留下,而重投短路会静默挡掉这个会话之后所有的入站消息

8.3 心跳与二次投递

心跳既是重放的执行者,也要避免和重放打架。

避让: 有活跃的 pending final 且会话 30 秒内动过 → 心跳直接跳过(src/infra/heartbeat-runner-execution.ts:321-325):

const HEARTBEAT_DEFER_WINDOW_MS = 30_000;
// recentSessionEntry.pendingFinalDelivery?.kind === "replayable" 且刚动过 → skipped
return { kind: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT };

判断只看 kind === "replayable"——transport-only 的存根没有用户可见文本,不值得为它推迟心跳。

归属: 心跳跑完要清 pending,但只能清自己产生的那份。判据是时间戳而非文本比对(heartbeatRunOwnsPendingFinalDelivery,src/infra/heartbeat-runner-delivery.ts:49):

const createdAt = entry?.pendingFinalDeliveryCreatedAt;
return typeof createdAt === "number" && createdAt >= runStartedAt;

投递: 心跳里对 pending 文本的清洗和恢复路径共用 sanitizePendingFinalDeliveryTextsrc/auto-reply/reply/pending-final-delivery.ts:168)——它剥内部元数据、剥静默 token,全静默就返回空串。

8.4 完整的兜底时序

模型产出 final


写 pendingFinalDelivery(kind/intentId/deliveries/context)


sendFinalPayload 进程被杀 / 长时间失败
├─ 路由到源通道? 路由投递 重启恢复扫描
└─ 否 → dispatcher.sendFinalReply main-session-recovery-state.ts:491
│ 找 kind === "replayable" 的条目
├─ 成功 → 按 deliveryId 结算存根 → 重投 / 计次(MAX_RECOVERY_RETRIES) / 记错
│ → clearPendingFinalDeliveryAfterSuccess
└─ 失败 → 存根留着


下一次心跳 / 恢复扫描接手

8.5 无可见回复的兜底信号

如果整轮下来什么都没发出去,dispatch 会给上游一个标记(src/auto-reply/reply/dispatch-from-config.finalize.ts:402):

{ noVisibleReplyFallbackEligible: true }

类型在 dispatch-from-config.types.ts:24。条件缺一不可:没排上 final、没观察到任何投递、且这个会话不允许空回复当静默。通道拿到这个标记可以决定要不要发一句「我没什么要说的」。


9. 巧妙之处(可以搬走的做法)

  1. 用「租约排队」而不是「取消过期」解决乱序投递。 不发号比代、不取消旧回复,而是让同一会话的可见投递 FIFO 串行——顺序问题在结构上被消除(dispatch.ts:51:325;keyed-fifo-lease.ts:34-46)。

  2. 预约位计数。 pending = 1 起步(reply-dispatcher.ts:259-261),把「还没开始产出」和「已经产出完」这两个都是 0 的状态区分开。这是所有「等待空闲」型 API 的通病,这个技巧只花一行。

  3. 两个去重键,一个含路由一个不含。 一个防重发、一个防「流式 + 最终」重复(block-reply-pipeline.ts:57 / :72)。合并成一个键必然二选一失效。

  4. 超时熔断整条链而不是跳过单条。 保序优先于完整性(block-reply-pipeline.ts:200)。聊天场景里乱序比缺失更让人困惑。

  5. 进度草稿的「二次事件立即启动」。 单个快工具不闪草稿,多步任务立刻可见(streaming.ts:660-685)。用一个计数器就区分了「短活」和「长活」,不需要预测执行时间。

  6. 元数据走 WeakMap 而非载荷字段。 内部标记不进线上形状,也不需要在每个通道适配器里剥字段(reply-payload.ts:307)。代价是克隆必须显式搬运,源码用一个 copyReplyPayloadMetadata 统一收口。

  7. 清存根不被中止打断。 顺序颠倒会造成永久性会话阻塞,注释里挂了 issue 号(dispatch-from-config.finalize.ts:187)。这类「顺序即正确性」的地方值得学它把理由写在代码旁边。

  8. 未授权命令静默丢弃而非报错(§4.3)。不向陌生人泄漏命令表面。


10. 边界与本章不覆盖

本章不讲:

这套设计的固有代价:

代价表现
状态分散在内存和存储两处跟进队列在内存(重启即失),pending final 在会话存储(能恢复)。两者不一致时以存储为准,但内存队列的丢失是真丢失
元数据依赖显式搬运任何 {...payload} 展开都会丢元数据,除非配 copyReplyPayloadMetadata
超时熔断会丢块一次慢投递会让后续所有块被跳过,只剩最终答案
转向依赖「正在流式」模型在工具执行阶段时转向注入会被拒,退回排队
心跳在忙时被丢弃长时间运行的会话可能连续错过多个心跳周期

代码里看不出来的: dispatch-from-config.*.ts 一族合计数千行,本章只走通了主干(入口 → 回调装配 → 最终投递 → 计数);其中 ACP 会话绑定、插件自有绑定绕过、卡死会话恢复等分支未展开。


11. 代码地图

主题文件路径符号名
裸入口src/auto-reply/dispatch.tsdispatchInboundMessage
缓冲入口src/auto-reply/dispatch.tsdispatchInboundMessageWithBufferedDispatcher
指定入口src/auto-reply/dispatch.tsdispatchInboundMessageWithDispatcher
前台回复租约src/auto-reply/dispatch.tssrc/shared/keyed-fifo-lease.tsforegroundReplyLeasesresolveForegroundReplyOrderKeycreateKeyedFifoLeaseRegistry
排空保证src/auto-reply/dispatch-dispatcher.tswithReplyDispatcher / settleReplyDispatcher
编排装配src/auto-reply/reply/dispatch-from-config.ts(门面)dispatchReplyFromConfig
收尾与计数src/auto-reply/reply/dispatch-from-config.finalize.tsfinalizeDispatchAndAuditneedsTtsFallback
存根清除src/auto-reply/reply/dispatch-from-config.pending-final.tsclearPendingFinalDeliveryAfterSuccess
回复解析src/auto-reply/reply/get-reply.tsgetReplyFromConfig
原生命令快路径src/auto-reply/reply/get-reply-native-slash-fast-path.tsmaybeResolveNativeSlashCommandFastReply
纯指令判定src/auto-reply/reply/directive-handling.directive-only.tsisDirectiveOnly
纯指令回执src/auto-reply/reply/directive-handling.impl.tshandleDirectiveOnly
指令应用(混合消息)src/auto-reply/reply/get-reply-directives-apply.tsapplyInlineDirectiveOverrides
命令注册表src/auto-reply/commands-registry.shared.tsbuildBuiltinChatCommands / defineChatCommand
注册表缓存src/auto-reply/commands-registry.data.tsgetChatCommandsdefineDockCommand
原生规格导出src/auto-reply/commands-registry.tslistNativeCommandSpecs
命令鉴权src/auto-reply/command-auth.tsresolveCommandAuthorization
命令分派src/auto-reply/reply/commands-core.tshandleCommands
/statussrc/auto-reply/reply/commands-status.tsbuildStatusReply
/compactsrc/auto-reply/reply/commands-compact.tsextractCompactInstructions
队列裁决src/auto-reply/reply/queue-policy.tsresolveActiveRunQueueAction
忙时等待src/auto-reply/reply/get-reply-run-queue.tsresolvePreparedReplyQueueStateREPLY_RUN_STILL_SHUTTING_DOWN_TEXT
轮次准备src/auto-reply/reply/get-reply-run.tsrunPreparedReply
转向与排队src/auto-reply/reply/agent-runner-run.tsrunReplyAgent
队列门面src/auto-reply/reply/queue.tsenqueueFollowupRun / scheduleFollowupDrain / resolveQueueSettings
入队去重src/auto-reply/reply/queue/enqueue.tsrecent-message-ids.tsenqueueFollowupRun / isRunAlreadyQueued
队列排水src/auto-reply/reply/queue/drain.tsscheduleFollowupDrainrememberFollowupDrainCallback
队列设置src/auto-reply/reply/queue/settings.tsstate.tsresolveQueueSettingsDEFAULT_QUEUE_*
/queue 指令src/auto-reply/reply/queue/directive.tsextractQueueDirective
跟进执行src/auto-reply/reply/followup-runner.tscreateFollowupRunner
出站契约src/auto-reply/reply-payload.tsReplyPayload / ReplyPayloadMetadata
TTS 补充src/auto-reply/reply-payload.tsbuildTtsSupplementMediaPayload / getReplyPayloadTtsSupplement
投递器src/auto-reply/reply/reply-dispatcher.tscreateReplyDispatcher / createReplyDispatcherWithTyping
投递器契约src/auto-reply/reply/reply-dispatcher.types.tsReplyDispatcher
流式参数src/auto-reply/reply/block-streaming.tsresolveEffectiveBlockStreamingConfig / resolveBlockStreamingChunking
块合并src/auto-reply/reply/block-reply-coalescer.tscreateBlockReplyCoalescer
块流水线src/auto-reply/reply/block-reply-pipeline.tscreateBlockReplyPipeline / createBlockReplyContentKey
块投递处理器src/auto-reply/reply/reply-delivery.tscreateBlockReplyDeliveryHandler / normalizeReplyPayloadDirectives
typing 策略src/auto-reply/reply/typing-policy.tsresolveRunTypingPolicy
typing 模式src/auto-reply/reply/typing-mode.tsresolveTypingMode / createTypingSignaler
typing 控制器src/auto-reply/reply/typing.tscreateTypingController
文本切分src/auto-reply/chunk.tsresolveTextChunkLimit / resolveChunkMode
进度草稿src/channels/progress-draft-compositor.tscreateChannelProgressDraftCompositor
草稿启动闸src/channels/streaming.tscreateChannelProgressDraftGate
状态表情src/channels/status-reactions.tscreateStatusReactionController
存根清洗src/auto-reply/reply/pending-final-delivery.tssanitizePendingFinalDeliveryText
完成投递策略src/auto-reply/reply/completion-delivery-policy.tscompletionRequiresMessageToolDelivery
心跳 tokensrc/auto-reply/heartbeat.tsstripHeartbeatToken
心跳避让与归属src/infra/heartbeat-runner-execution.tsheartbeat-runner-delivery.tsheartbeatRunOwnsPendingFinalDelivery
重启恢复src/agents/main-session-recovery/resumeMainSessionbuildResumeMessage