Block 多 Agent 通信平台 buzz 的 ACP 会话池与队列
本文基于 buzz 仓库 commit 8342dfc(2026-08-04)梳理,该项目仍在高频迭代,具体行为以仓库 https://github.com/block/buzz 最新代码与文档为准。
外部 Agent 客户端接进消息网络这件事,真正难的不是”把消息喂给模型”,而是同一个频道不能被两个 Agent 同时处理、进程挂了要能自己爬起来、事件积压时要有地方放且知道什么时候该丢。 buzz 这个项目把这三件事拆成了三个互不越界的模块:池负责”谁在干活”,生命周期负责”什么时候有人干活”,队列负责”活按什么顺序排”。任何一块缺席,另外两块都会被拖垮。
这里说的 buzz,是 Block 开源的多 Agent 通信平台(仓库 https://github.com/block/buzz,许可证 Apache-2.0,Copyright 2026 Block, Inc.),不是”热度”也不是”蜂鸣器”。它把人和 Agent 放进同一张消息网络里协作,底座是 Nostr 协议。
站内已有三篇相邻的文章:Agent 并发编排 讲的是你自己写编排层时怎么切分并发单元,多会话并发冲突 讲的是多个会话抢同一份状态时的冲突面,Agent 缓存与幂等 讲的是重复投递怎么不出二次副作用。本篇不重复这些通用结论,只做一件事:把 buzz 这个具体项目里的三个真实模块摊开,看它在源码层面怎么把这些问题落成常量、状态机和日志。
一、先把底座和边界说清楚
Nostr 是一个去中心化的消息协议,理解本文只需要四个词:
- 事件(event):一条带作者公钥、时间戳、标签和签名的 JSON 消息,是网络上唯一的数据单位。
- 签名:作者用自己的私钥对事件内容签名,任何人拿公钥就能验签。仓库里
crates/buzz-core/src/event.rs的StoredEvent就是给 Nostr 事件包一层中继侧元数据(接收时间、频道归属、是否已验签),并提供is_verified()查询。 - 中继(relay):负责存储和转发事件的服务器。客户端连中继、订阅过滤条件、收事件。
- 密钥自持:身份就是那对密钥,没有”找回密码”这一说。私钥丢了,这个身份就没了;私钥泄露了,别人就能以这个身份发消息。这一条对 Agent 尤其要紧——Agent 的身份也是一对密钥。
本文拆的是 crates/buzz-acp 这个 crate。仓库 crates/ 下共 28 个 crate,buzz-acp 是其中把 Buzz 上的消息事件桥接给外部 ACP Agent 的那一个——它的命令行 about 文案(crates/buzz-acp/src/config.rs)写的是 “ACP harness that bridges Buzz events to AI agents”。ACP 客户端在这里指的是被 harness 拉起的 Agent 子进程,代码里叫 AcpClient。
三个主角分别住在:crates/buzz-acp/src/pool.rs(6998 行)、crates/buzz-acp/src/pool_lifecycle.rs(312 行)、crates/buzz-acp/src/queue.rs(4764 行)。粘合它们的主循环在 crates/buzz-acp/src/lib.rs。
二、池:所有权取出来,用完必须还回去
pool.rs 开头的模块注释把心智模型写得很直白:
AgentPool
├── agents: Vec<Option<OwnedAgent>> ← idle agents sit here
├── join_set: JoinSet<()> ← in-flight tasks
├── task_map: HashMap<Id, TaskMeta> ← panic recovery metadata
└── result_tx/rx: mpsc channel ← tasks return agents here
关键约束是下一行那句:AcpClient 不是 Clone。所以它不能”共享”,只能”转移”——try_claim() 把 OwnedAgent 从槽位里 take() 出来,槽位变成 None;任务跑完通过 result_tx 把 agent 送回主循环,return_agent() 按 agent.index 塞回原槽位。槽位索引是不变量:from_slots() 的注释专门说明,启动时失败的 agent 要保留成 None 占位,不能压缩数组,否则 index 和数组下标对不上。
try_claim(channel_id) 是两趟扫描:第一趟优先找”已经持有这个频道 session 的 agent”,第二趟才退化为”任意空闲 agent”。第一趟就是会话亲和——同一个频道尽量落回同一个 Agent 进程,省掉重建 session 的开销。主循环在调用 try_claim 之前会先用 has_session_for() 算出 affinity_hit,用于观测。
会话状态本身不在池上,在 SessionState 里:sessions: HashMap<Uuid, String> 是频道到 session id 的映射,另有每频道的 turn_counts(轮次计数,用于主动轮换 session)、core_sections、canvas_sections。这个结构被刻意从 OwnedAgent 里拆出来,注释给的理由是”不用真的拉起 Agent 子进程就能测这个状态机”。
两个细节值得记住:
其一,return_agent() 发现目标槽位已经有人时,会打一条 BUG: return_agent called for slot {idx} which is already occupied — overwriting 的 ERROR 日志然后覆盖。注释解释了为什么选覆盖而不是丢弃:丢弃会永久泄漏一个槽位,覆盖至少还能继续跑。这是典型的”错了要吵,但不能停”。
其二,live_count() 算的是”空闲的 + 被借出去的”,也就是 agents 里非 None 的数量加上 task_map 的条目数。它存在的目的是判断”是不是所有 Agent 都死光了”,好触发重启。
Agent 数量由 --agents / BUZZ_ACP_AGENTS 控制,默认 1,取值范围 1 到 32(config.rs 里 clap 的 value_parser 写死了这个区间)。crate 的 README 建议大多数部署从 N=2 起步,并提醒每个 Agent 会拉起自己的 MCP server 子进程,资源占用大致按 N 倍增长。README 还写明:N 个 Agent 用的是同一个 Nostr 身份,用户看到的是一个机器人;同一频道永远不会被两个 Agent 同时处理(这条由队列保证);N>1 时跨频道的消息顺序不做保证。
三、生命周期:懒启动的四态机,只唤醒一次
默认情况下 harness 启动就把 Agent 子进程全拉起来。但加上 --lazy-pool / BUZZ_ACP_LAZY_POOL=true 之后,它会先连中继、认证、订阅,把接到的活先攒着,等第一条真正要处理的事件到了再去启动子进程。仓库 docs/remote-agents.md 是一份自我标注为 draft 的远端 Agent 规范草案,里面给远端 pod 选择懒启动的理由写得很实在:集群里挂一堆空转的 LLM 进程是纯烧钱,而且没人在旁边等它预热。要留意的是,这份文档描述的是这套远端托管方案该长什么样,不等于它已经全部落进代码——本文引它,只引它对既有 harness 行为的观察和它自己写下的约束。
这套开关状态就是 pool_lifecycle.rs 里那个不到一百行的核心枚举:
pub(crate) enum PoolLifecycle<P> {
Listening,
Waking {
attempt: u32,
},
Ready(P),
Failed {
attempt: u32,
retry_at: Instant,
error: String,
},
}
这个模块的注释一开头就划清了边界:中继连接、订阅、事件缓冲都不归它管,它只管”这个延迟启动的池到底是没起、正在起、起好了、还是等着重试”。这种把状态机从 IO 里剥出来的写法,直接换来的是它自己文件里那一串不需要网络的单元测试。
三个设计点:
唤醒必须有活可干。 start_wake_if_due(has_pending_work, now) 第一件事就是 if !has_pending_work { return None; }。没有待办就不启动,这是懒启动的字面含义。
一次转换只发一张票。 这个函数返回 Option<u32>,只在状态真的从 Listening(或到期的 Failed)转进 Waking 时返回一次 attempt 编号。调用方把这个编号挂在唯一那个池初始化任务上,任务完成时原样带回来。连续调用两次,第二次拿到的是 None——所以不会出现”事件一多就并发拉起好几个池”。
旧结果不许覆盖新池。 complete_wake() 只接受与当前 Waking { attempt } 完全匹配的编号,否则返回 Err,错误串分别是 "wake result attempt did not match Waking attempt" 和 "wake completed while lifecycle was not Waking"。主循环收到这种错误时打一条 discarding stale pool wake result 的 warn 就跳过。这是并发编程里很容易漏的一格:第一次启动超时、第二次已经在跑了,第一次的结果这时候姗姗来迟,如果照单全收就会把新池顶掉。
失败退避是 retry_delay(attempt):初始 5 秒(INITIAL_RETRY_DELAY),每次翻倍,封顶 300 秒(MAX_RETRY_DELAY)。实现用的是 checked_shl 加 saturating_mul 再 min,把指数位先 min(63),所以 attempt 传 u32::MAX 也不会溢出——测试里就直接断言了 retry_delay(u32::MAX) == MAX_RETRY_DELAY。
take_ready() 用 std::mem::replace 把池取出来并把状态退回 Listening,取第二次返回 None。主循环里对应的是 .take_ready().expect("successful wake stores a ready pool")。
四、队列:每频道一条队,跨频道比”谁的队头最老”
queue.rs 里的 EventQueue 是三块里状态最多的一块。它的字段清单本身就是一份职责说明:
queues: HashMap<Uuid, VecDeque<QueuedEvent>>——每个频道一条队。in_flight_channels: HashSet<Uuid>——哪些频道正有一轮 prompt 在跑。in_flight_deadlines——每个在飞频道的兜底截止时间。retry_after/retry_counts——每频道的退避时刻与重试次数。cancelled_batches/cancel_reasons——被打断的那批事件,以及打断原因。withheld_native_steer——为了避开一个竞态,被临时从主队列挪走的事件。
入队。 push() 先看 dedup 模式。DedupMode::Drop 下,如果这个频道正在飞,新事件直接 debug 日志后丢弃并返回 false;DedupMode::Queue 下则照常入队。每条队列的深度上限 MAX_PENDING_PER_CHANNEL = 500,到顶就 pop_front() 丢最老的那条,并打 warn。
出队。 flush_next() 的顺序是固定的:先扫一遍 in_flight_deadlines 把过期的在飞频道强制释放,再从”队列非空 && 不在飞 && 没被 retry_after 卡住”的频道里,挑队头事件时间最老的那个。这是跨频道的公平性口径——不是轮询,是按最老的待办排。选中之后一次性 drain 出最多 MAX_BATCH_EVENTS = 50 条,合成一个 FlushBatch。
这里有一处很容易忽略的修正:drain 出来之后代码会做一次 events.sort_by_key(|be| be.event.created_at)。注释给的理由是中继回放历史事件时是 ORDER BY created_at DESC(新的在前),而下游消费者要求批次里最后一条是最新的——因为 prompt 的 scope 和回复锚点都取自最后一条。用的是稳定排序,同一秒的事件保留投递顺序。你要是自己拼这套东西,这类”上游顺序和下游假设相反”的坑,通常要到线上才发现。
完成与重试。 mark_complete(channel_id) 把频道从在飞集合里摘掉。它对 retry_counts 的处理有讲究:如果这个频道还挂着未到期的 retry_after,说明它是被 requeue 回来的,计数保留,退避序列继续;否则视为健康完成,计数清零。
requeue(batch) 是失败重投:事件用 push_front 塞回队头(且保留原始 received_at,不重置,这样频道在跨频道排序里的位置不变),然后设置退避。退避是 BASE_RETRY_DELAY_SECS = 5 起,按 1u64 << (attempt-1).min(6) 翻倍,封顶 MAX_RETRY_DELAY_SECS = 300,再乘一个 0.8 到 1.2 之间的抖动系数(熵源是当前时间的亚秒纳秒数)。重试次数超过 MAX_RETRIES = 10 就死信:打 ERROR 日志、清掉计数和 retry_after、把这批事件返回给调用方而不是重投,让上层去频道里发一条可见的失败通知。
在飞兜底。 这是队列里最像”防呆”的一段。in_flight_deadline 的取值是 max_turn_duration + IN_FLIGHT_DEADLINE_BUFFER_SECS,后者是 100 秒;不显式设置时用 DEFAULT_IN_FLIGHT_DEADLINE_SECS = 7300,正好是默认 max turn 7200 秒加这 100 秒缓冲。这个值必须严格大于单轮上限,注释说得明白:让正常跑满硬上限的那一轮有机会走 mark_complete 正常返回,兜底才不会误伤。一旦真的过期,日志是 BUG: in-flight channel expired without mark_complete — auto-releasing; N dispatched event(s) orphaned——注意措辞,已经派发给那个挂死 prompt 的事件是孤儿,不恢复;只有 withheld_native_steer 里那些根本没送到 Agent 手上的事件会被 recover_withheld_for_expired_channel() 捞回队头。这个”送到了就不重投、没送到才恢复”的区分,是幂等设计里最值钱的一句判断。
compact_expired_state() 是周期性维护:清掉已过期的 retry_after,以及那些既没队列、没退避、也没在飞的频道遗留的 retry_counts。注释特意点出在飞这个条件是关键的——一个刚被 flush 空、退避也过期的频道,可能还有一轮在跑,这时候清计数会把退避序列重置回去。
五、这三样合起来意味着什么
| 组成部分 | 它负责什么 | 仓库位置 | 你什么时候会碰到它 |
|---|---|---|---|
AgentPool | 持有 N 个 OwnedAgent(每个含一个 AcpClient 进程),取还所有权,会话亲和,槽位不变量 | crates/buzz-acp/src/pool.rs | 调 --agents、看到 pool_exhausted、排查”为什么加机器不提速” |
SessionState | 频道到 session id 的映射、每频道轮次计数、失效与轮换 | crates/buzz-acp/src/pool.rs | 用 !rotate、切模型、Agent 被移出频道后清理陈旧 session |
PoolLifecycle<P> | 懒启动池的四态:Listening / Waking / Ready / Failed,单次唤醒与退避重试 | crates/buzz-acp/src/pool_lifecycle.rs | 开 --lazy-pool、启动一直失败、看到 discarding stale pool wake result |
EventQueue | 每频道排队、跨频道按最老队头调度、批量 drain、退避重投、死信、在飞兜底 | crates/buzz-acp/src/queue.rs | 事件积压、看到队列深度告警、排查事件为什么被丢 |
| 主循环粘合层 | 唤醒判定、dispatch_pending、结果回收、30 秒维护 tick | crates/buzz-acp/src/lib.rs | 读整条链路、加观测点 |
把它们连起来看,一次正常投递是这样走的:事件进 push() → 主循环发现 has_flushable_work() 为真 → 懒启动模式下 start_wake_if_due() 拿到 attempt 并拉起池 → flush_next() 选出最老的频道并 drain 成一个批次 → try_claim() 借出一个 Agent(优先借”已经认识这个频道”的那个)→ 跑完通过 mpsc 把 Agent 送回 → return_agent() + mark_complete()。
借不到 Agent 时的处理很干脆:dispatch_pending 打一条 pool_exhausted 的 debug 日志,调 requeue_preserve_timestamps(batch) 把批次原样退回队头(保留时间戳、不设退避、不计重试次数),再 mark_complete 释放频道,然后 break 跳出派发循环。池满不是错误,不该罚这个频道。
三块的分工可以这么概括:池决定同时能有几个人在干活,生命周期决定什么时候真的有人,队列决定活按什么顺序排、排不下时怎么办。只有池没有队列,事件峰值会直接打到进程上;只有队列没有生命周期,冷启动失败就变成无限重试风暴;只有生命周期没有池,就没有”同一频道串行、跨频道并行”这个基本保证。
六、边界与代价
这套设计放弃的东西同样清楚。
跨频道顺序不保证。 README 明说 N>1 时不保证跨频道消息顺序。它保证的只有”单频道串行”。你要是依赖多个频道之间的先后关系,这层不给。
事件会被丢,而且是设计内的。 三处丢弃:DedupMode::Drop 下在飞频道的新事件直接丢;单频道超过 500 条时丢最老的;重试超过 10 次整批死信。前两处只有 debug/warn 日志,后一处是 ERROR 并把批次交还给上层去发失败通知。
在飞兜底救频道,不救事件。 上文那条 orphaned 日志说得很直白。兜底解决的是”频道被永久卡住”,不解决”这批事件到底处理没处理”。
懒启动会牵动别的东西。 上面那份规范草案记了一个具体后果,而这个后果是能回代码验证的:harness 的 30 秒维护 tick 是挂在 pool_ready 上的,而懒启动模式下 pool_ready 一开始是 false、只有被唤醒才翻真,唤醒又要求有待办(就是前面那个 has_pending_work)。所以一个从没被 @ 过的懒启动 pod,那个 tick 根本不跑。文档因此给它规划中的空闲回收逻辑立了一条约束:必须走自己的定时器,不能搭这趟车。这类”状态门控把周期任务一起门掉”的耦合,是懒启动方案的通用代价。
它明确不管的事:不管跨进程、跨主机的调度(池就是本进程里的一个 Vec);不管模型侧的配额与限流;不管 Agent 在自己那轮里干了什么。最后这条要展开一句——PromptContext 里有 cwd 字段,Agent 是在一个真实工作目录里跑的,并且被提示用 buzz messages send 这类命令往频道里发消息。也就是说,一个被接进来的 Agent 具备”以这个机器人身份对外发言”和”在 cwd 里动文件”两种能力。要不要给、给到哪一步,是部署方的决定,池和队列不替你把关。
密钥这件事没有退路。 Agent 的身份就是一对 Nostr 密钥,PromptContext 里的 agent_keys 拿的就是它。私钥丢了这个 Agent 的身份就找不回来;私钥被人拿到,对方就能以它的名义在网络上发签名事件。自建中继时还要多想一层:事件是明文存在你自己那台机器的中继存储里的,端口暴露面、备份落盘位置都得自己交代清楚。
七、上手与避坑清单
1. --dedup=drop 配默认的 --multiple-event-handling=steer,启动会直接报错。
为什么会踩:两个参数看起来正交,很容易分别按直觉设。但 steer / interrupt / owner-interrupt 这三种模式的实现是”取消当前轮 + 合并重投”,取消排空的那个窗口里,Drop 模式会把新事件丢掉——正好丢掉触发这次打断的那条消息。config.rs 里的 validate_multiple_event_handling 直接把这个组合拒了。
怎么避:要用打断类模式就配 --dedup=queue(它本来就是默认值);只有明确选了 --multiple-event-handling=queue 才考虑 Drop。
2. idle_timeout 一定要小于 max_turn_duration。
为什么会踩:这两个超时一个管”多久没动静”,一个管”总共跑多久”。把前者调大到超过后者,硬上限会先触发,空闲检测永远轮不到,等于白配。
怎么避:Config::from_args 里有一道硬校验,idle_timeout >= max_turn_duration 直接返回错误。默认值是 900 秒对 7200 秒,注释解释 900 秒是给”外层 ACP 通道安静、里面在跑长工具调用”留的余量。
3. 把 max_turn_duration 调很大,在飞兜底会跟着一起变长。
为什么会踩:兜底截止时间是 max_turn_duration + 100。你为了让某个长任务跑完把上限调到几天,代价是某个频道真挂死时,也要等这么久才会被自动释放。上限是 MAX_TURN_DURATION_CEILING_SECS = 604800(7 天),再高会被拒。
怎么避:先确认长耗时的真实来源。如果是单个工具调用慢,调 idle_timeout 比调 max_turn_duration 更对症。
4. 加 --agents 不会让单个频道变快。
为什么会踩:把 N 从 1 调到 8,看着像”并发 8 倍”,但单频道的串行约束是队列层的 in_flight_channels 保证的,跟 Agent 数量无关。同一个频道的事件永远排队。
怎么避:先看 pending_channels()——积压是分布在多个频道上,加 N 才有意义;如果全压在一个频道,加 N 只是多几个空转的子进程和 MCP 子进程。
5. 两条 ERROR 日志要单独设告警。
为什么会踩:这两条都是”系统还在跑,但有东西已经悄悄不对了”,混在普通日志里很容易冲掉。一条是 BUG: in-flight channel expired without mark_complete,意味着有一轮 prompt 没有正常收尾;一条是 dead-lettering batch after 开头那条(后面跟的是实际重试次数和被丢弃的事件条数),意味着一整批用户消息被丢了。
怎么避:按这两条日志里不随参数变化的固定前缀抓关键词告警,别只盯 ERROR 总量。这也是 Agent 可观察日志 里那套”给失败路径留可搜索指纹”的具体落法。
6. 别假设批次里的事件是按投递顺序排的。
为什么会踩:中继回放存量事件时是新的在前。如果你在下游按”数组第一个是最早的”写逻辑,回放场景下就会拿反。
怎么避:flush_next() 已经按 created_at 做了稳定排序,最后一条是最新的。你自己扩展这条链路时,要沿用同一个口径,别在中途插一个改变顺序的步骤。
7. 失败重投保留原始时间戳,这是有意的,别”顺手修掉”。
为什么会踩:requeue 里把 received_at 原样带回去,看起来像忘了刷新。实际上刷新它会让这个频道在跨频道排序里退到队尾,一个反复失败的频道就再也抢不到调度。退避的职责由 retry_after 单独承担,跟排序解耦。
怎么避:读一下 requeue 上面那段注释——它把”退避延迟来自指数退避,不来自重置 received_at”写成了一句话。这类”看起来像 bug 的正确代码”,改之前先找注释。这条跟 Agent 失败重试 里说的”退避策略要和排队策略分开”是同一个道理。
收个尾
如果你要把这套思路搬到自己的 Agent 接入层,可以拿这几条对照自查:
- 同一个逻辑会话(频道、工单、对话)能不能被两个执行器同时拿到?靠什么数据结构保证不能?
- 执行器满了的时候,任务是被丢、被罚(退避)、还是原样退回?三种语义你分清了吗?
- 有没有一个”执行器没有正常收尾”的兜底释放?它的时限是不是严格大于业务本身的最大耗时?
- 重投时哪些事件算”已送达不能重发”,哪些算”从没送到必须恢复”?这两类有没有分开的存储?
- 失败退避的上限、死信阈值、单队列深度上限,这三个数字是不是都写在了代码里且能被搜到?
接着往下读的话,建议按这个顺序:先 crates/buzz-acp/src/pool_lifecycle.rs(最短,312 行,状态机完整且自带测试),再 crates/buzz-acp/src/queue.rs 前 140 行(那段 ASCII 状态机注释把整个队列的转换写全了),最后才是 crates/buzz-acp/src/pool.rs 和主循环 crates/buzz-acp/src/lib.rs 里的 dispatch_pending。协议侧想补课,docs/nips/ 下有 15 份 NIP 规范 md 和 2 份 fixtures json,其中标注为 draft optional 的 NIP-OA.md(Owner Attestation)定义了一个可选的 auth 标签——owner 的密钥用它授权某个 Agent 密钥”以 Agent 自己的名义发事件”,正好接得上本文里 PromptContext 中 agent_keys 与 agent_owner_pubkey 这对身份字段。
本篇属于一个把开源多 Agent 通信平台 buzz逐层拆开讲的系列,整体地图见 buzz 是什么:Block 开源的多 Agent 通信平台全景图;沿着这条线往下,还可以看 buzz 多 Agent 通信:Block 开源平台里活递给谁看目录 和 Block 开源多 Agent 通信平台 buzz 的可观测最小实现。