Block 开源 buzz 多 Agent 通信平台:中继内部三层转发
本文基于 buzz 仓库 commit 8342dfc(2026-08-04)梳理,该项目仍在高频迭代,具体行为以仓库 https://github.com/block/buzz 最新代码与文档为准。
读 buzz 中继源码最容易走错的第一步,是把 crates/buzz-relay/src/router.rs 当成消息路由表——它其实是 axum 的 HTTP 路由与 WebSocket 升级入口,真正决定”这条消息发给谁”的代码在 subscription.rs 的几张索引表里。 名字对不上位置,后面读什么都是错位的。
先交代两个协议名词,后面全篇都要用。buzz 建在 Nostr 之上:Nostr 里的”事件”(event)是一条带作者公钥、时间戳、kind 编号、标签和内容的 JSON 对象,由作者私钥签名,任何人拿到都能独立验签;“中继”(relay)是接收、存储、按订阅条件转发这些事件的服务端进程,它不持有用户私钥,也无法伪造别人的事件。这个仓库里的 crates/buzz-relay 就是这样一个中继实现,Apache-2.0 许可(Copyright 2026 Block, Inc.)。
站内相邻的几篇分工不同:Agent 并发编排讲的是多个 Agent 的任务怎么拆怎么派,opencode 的事件总线讲的是单进程内部事件如何串起来,hermes 的监控与可观测讲的是跑起来之后看什么指标;本篇只钉在一层——一条消息在 buzz 中继进程内部,从 socket 连上来到帧写回去的那条路。
一、三层各在哪个文件:先把名字对上位置
buzz 中继把”收一条消息、转一条消息”这件事切成了职责相当清楚的几块。下面这张表里的路径都是我在仓库里实际打开读过的文件。
| 组成部分 | 它负责什么 | 仓库位置 | 你什么时候会碰到它 |
|---|---|---|---|
| HTTP 路由与升级 | 声明所有 HTTP 路由、按 Accept 头决定返 NIP-11 文档还是升级成 WebSocket、挂 CORS 与请求体上限 | crates/buzz-relay/src/router.rs | 加一个新端点、排查 404 或跨域 |
| 连接生命周期 | 名额许可、发认证挑战、收/发/心跳三个循环、背压与断连清理 | crates/buzz-relay/src/connection.rs | 客户端莫名被踢、发送卡顿 |
| 协议解析 | 把一帧 JSON 数组解析成 ClientMessage,把出站消息格式化成 NIP-01 数组 | crates/buzz-relay/src/protocol.rs | 客户端收到 ["NOTICE","invalid message: …"] |
| 订阅注册表 | 维护 conn→sub→过滤器的映射与多张扇出索引,回答”这条事件命中哪些订阅” | crates/buzz-relay/src/subscription.rs | 订阅收不到消息、想知道扇出复杂度 |
| REQ 处理 | 鉴权与作用域校验、注册订阅、按过滤器回放历史、发 EOSE | crates/buzz-relay/src/handlers/req.rs | 历史消息拉不全、订阅被 CLOSED |
| EVENT 处理与扇出 | 验签、入库、发布到跨节点通道、本地扇出并复核访问权 | crates/buzz-relay/src/handlers/event.rs | 消息发出去别人没收到 |
| 签名校验 | 校验事件 id 哈希与 Schnorr 签名 | crates/buzz-core/src/verification.rs | 客户端拿到 invalid: 开头的 OK |
| 过滤器匹配 | NIP-01 过滤器的逐字段匹配、按 kind 的读取授权 | crates/buzz-core/src/filter.rs | 过滤器语义存疑时 |
给个规模感:整个 crates/ 目录下有 28 个 crate(共 413 个文件,ls crates 加 find crates -type f | wc -l 自己就能数出来),中继只是其中一个,仓库另有 desktop/ 2359 个文件、mobile/ 467 个、web/ 65 个、admin-web/ 14 个;docs/nips/ 放着 15 份 NIP 规范 md 和 2 份 fixtures json,另有 crates/buzz-core/src/pairing/NIP-AB.md 是第 16 份规范。NIP 是 Nostr 的规范提案编号,作用类似 RFC——一个编号对应一组事件 kind 与字段约定。
二、协议解析层:能早拒绝的一律早拒绝
protocol.rs 只做一件事:把一段字符串变成五种枚举之一,或者失败。五种是 Event、Req、Close、Count、Auth,分别对应 NIP-01 的 EVENT/REQ/CLOSE、NIP-45 的 COUNT、NIP-42 的 AUTH。解析逻辑是纯函数式的:先 serde_json::from_str 成 Value,要求是非空数组,第一个元素必须是字符串,然后按类型分支。
值得单独记住的是它把两个上限硬编码在了这一层:
/// NIP-11 advertised limit: subscription IDs longer than this are rejected.
const MAX_SUB_ID_LENGTH: usize = 256;
/// NIP-11 advertised limit: REQ messages with more filters than this are rejected.
const MAX_FILTERS_PER_REQ: usize = 10;
这两个数字不是随手写的,crates/buzz-relay/src/nip11.rs 里对外宣告的 max_subid_length 是 256、max_filters 是 10、max_subscriptions 是 1024,宣告值和执行值是同一套。NIP-11 是中继自描述文档的规范:客户端用 Accept: application/nostr+json 请求根路径,就能拿到这份能力与限制声明。而 max_message_length 直接取配置里的帧上限,默认值在 crates/buzz-relay/src/config.rs 里是 DEFAULT_MAX_FRAME_BYTES = 512 * 1024,可用环境变量 BUZZ_MAX_FRAME_BYTES 覆盖。
出站方向对称地放在 RelayMessage 里,一组返回 String 的静态方法:auth_challenge、event、notice、eose、ok、closed、count。这个设计让”回什么”这件事只有一处能改,也让被拒绝的请求有精确回执——connection.rs 里 request_rejection_message 会看有没有 sub_id:有就发 CLOSED(客户端知道是哪个订阅挂了),没有就退化成 NOTICE。做过 WebSocket 网关的人应该能体会这个区别的价值:一条笼统的 NOTICE 会让客户端不知道该重试哪个订阅。
三、订阅注册表:几张索引表决定这条事件扇给谁
这是三层里最值得细看的一层。SubscriptionRegistry 的字段结构基本把设计意图写在脸上(下面是裁剪后的定义):
pub struct SubscriptionRegistry {
/// Maps conn_id → sub_id → (filters, community_id, channel_id).
subs: DashMap<ConnId, HashMap<SubId, SubEntry>>,
channel_kind_index: DashMap<(CommunityId, IndexKey), Vec<(ConnId, SubId)>>,
/// Subscriptions with a channel_id but no kind filter — need to receive ALL kinds.
channel_wildcard_index: DashMap<(CommunityId, Uuid), Vec<(ConnId, SubId)>>,
/// Global subscriptions indexed by kind — avoids O(all_subs) scan for global events.
global_kind_index: DashMap<(CommunityId, Kind), Vec<(ConnId, SubId)>>,
/// Global subscriptions indexed by both kind and `#p` recipient.
global_p_kind_index: DashMap<GlobalPKindIndexKey, Vec<(ConnId, SubId)>>,
/// Global subscriptions with no kind filter — wildcard, receives all global events.
global_wildcard_index: DashMap<CommunityId, Vec<(ConnId, SubId)>>,
}
subs 是权威数据,五张索引都是为了避免全表扫描而存在的加速结构。注册时(register_scoped)走一次判断:订阅有没有绑定到某个频道(channel_id),过滤器里有没有 kinds 约束,全局订阅是否被 #p 收件人完全约束住。三个维度组合出来就是上面那五张表。
有两个语义细节容易写错,这份代码把它们分开处理了。第一,某个过滤器完全没有 kinds 字段,按 NIP-01 的 OR 语义整条订阅就成了通配(extract_kinds_from_filters 返回 None),进 wildcard 索引。第二,kinds: [] 是”不匹配任何 kind”,这种订阅哪张索引都不进,注释写得很直白——它永远收不到事件,扇出时 filters_match 也会兜住。把这两种情况混同,就是”我明明写了空数组却收到全部消息”这类线上事故的来源。
扇出入口是 fan_out_scoped(community_id, event)。它先看事件有没有频道作用域:有就查 channel_kind_index 加 channel_wildcard_index;没有(全局事件)就先按事件上的每个 #p 标签查 global_p_kind_index,再查 global_kind_index 与 global_wildcard_index。命中的候选不会直接投递,而是都过一遍 push_match:回 subs 拿权威记录,重新比对 community 与 channel 作用域,再跑 filters_match,最后用一个 seen 集合去重。为什么要复核?源码注释说得很清楚——同 ID 订阅被替换时,索引里的候选快照可能已经过期,不复核就可能把 A 频道的事件投给已经换到 B 频道的同名订阅。仓库里为此专门留了 test_stale_candidate_snapshot_does_not_cross_subscription_scope 这样的回归测试。
还有一条对称的隔离不变量,注释里明确标注:全局订阅收不到频道作用域的事件,频道订阅也收不到全局事件。前者防频道内容外泄,后者防成员变更这类全局基础设施事件漏进频道流。
四、一条消息从连上来到扇出去
把三层串起来,路径大致是这样,每一步都能在文件里找到对应位置。
连接建立。 请求打到 /,nip11_or_ws_handler 先看是不是 admin 主机(是则短路,绝不让它落到公共前端或 WebSocket 入口),再看 Accept 是不是 application/nostr+json(是则返 NIP-11 文档)。接着是注释里叫”row zero”的一步:用请求 Host 调 crate::tenant::bind_community 把连接绑定到一个社区,在 WebSocket 升级之前完成,失败一律返回一句不区分原因的 relay: no community is configured for this host——不回显 host、不区分”未映射”和”查询失败”,免得未认证的调用方拿它探测部署里有哪些社区。升级时 limit_relay_websocket 把 max_message_size 与 max_frame_size 都设成配置的帧上限,让 tungstenite 在组装完整消息之前就拒掉超大帧;进程处于 shutting_down 时直接回 503 relay restarting。
连接就位。 handle_active_connection 先抢 conn_semaphore 名额(BUZZ_MAX_CONNECTIONS 默认 10000),生成 NIP-42 挑战串发给客户端,然后建两个 channel:数据通道容量取 BUZZ_SEND_BUFFER(默认 1000),控制通道固定容量 8。随后起三个任务——发送循环、心跳循环、认证超时任务,AUTH_TIMEOUT 是 5 秒,超时未完成认证就取消连接。
收帧。 recv_loop 拿到文本帧后,先在应用层再查一遍长度(注释说这是 defense in depth,解析器那层已经查过),超限回一条 NOTICE 并断开;二进制帧尝试按 UTF-8 解码后当文本处理,注释说明这是对某些 Nostr 客户端的兼容而非 NIP-01 要求。Ping 走控制通道回 Pong,控制通道满意味着写端彻底卡死,按终止处理。
分派。 handle_text_message 先解析,再过 enforce_ws_admission 做按主体的配额检查(EVENT 还要额外查一次消息配额),然后按类型分派:AUTH 与 CLOSE 在循环里同步处理,EVENT/REQ/COUNT 各抢一个 handler_semaphore 许可后 spawn 出去,并且是先建好 tracing span 再 spawn——裸 tokio::spawn 会丢掉链路上下文,这行注释值得抄走。抢不到许可就回”rate-limited: too many concurrent requests”。
验签与入库。 事件路径上会把 verify_event 丢进 spawn_blocking。crates/buzz-core/src/verification.rs 的模块注释直接写了原因:Schnorr 验签是 CPU 密集的,在异步上下文里绝不能直接调。Schnorr 是 Nostr 事件所用的签名算法,验一次就是一次完整的密码学计算,纯吃 CPU、没有任何 IO 等待——异步运行时的工作线程被这种计算占住,同一线程上其它连接的读写就一起停摆,所以要挪到阻塞线程池里去。验签先比对事件 id 的哈希,再验签名,失败回 ["OK", id, false, "invalid: …"]。
扇出。 落库之后的投递集中在 dispatch_persistent_event:先把审计入队(这一步是 await 的,保留背压),再 spawn 出去做剩下的事——按事件的 channel_id 选 EventTopic::Channel 或 EventTopic::Global 发到跨节点通道,标记本地回声,然后 fan_out_scoped 求收件人集合。关键是紧跟着的 filter_fanout_by_access:私有频道的每个收件人都要重新查一次成员资格,查不到就丢弃,查询出错也丢弃(注释写明 fail closed)。它守的不变量是”注册过订阅不等于有权收到”——即使某个节点上残留了一条过期订阅,投递这一刻还会再校验一次。之后按 sub_id 建帧缓存(同一 sub_id 的帧只序列化一次),交给连接管理器写入各自的发送通道。
写回。 send_loop 每轮先把控制帧排空,再用 biased select 保证取消优先于控制帧、控制帧优先于数据帧;数据帧最多攒 MAX_WS_SEND_BATCH = 64 个 feed 再一次 flush。发送通道满时不会阻塞,ConnectionState::send 用 try_send,满了就累加背压计数,达到 slow_client_grace_limit(默认 15)才取消连接。跨节点方向由 fan_out_pubsub_event 承接,用 (community_id, event_id) 做本地回声去重,避免自己发的事件绕一圈再投一遍。
收尾。 连接结束后逐条摘掉订阅并释放对应的跨节点 topic,注销连接,如果这个公钥在该社区已无其它连接就清理在线状态。
五、边界与代价
这套设计的取舍相当明确,用之前得认。
索引是按 (community, channel, kind) 这些维度建的,所以受益最大的是”订阅明确指定了频道和 kind”的形态。反过来,大量通配订阅会退化:wildcard 索引里的每个条目都要过一遍 push_match,事件多、通配订阅多的时候这是实打实的成本。
作用域隔离是双向硬隔离,不提供”一条订阅同时吃频道和全局”的能力。handlers/req.rs 里的 extract_channel_id_from_filters 有一条容易踩的规则:只要任一过滤器没有可解析成 UUID 的 #h 标签,或者多个过滤器给出了不同的频道 UUID,整条订阅就退化成全局订阅——它不会报错,只是走了另一条索引路径。
它明确不管的事:不替你保管私钥。Nostr 的身份就是密钥对,中继只验签不持钥,私钥丢了等于身份丢了,别人也无法替你补发——这是密钥自持的代价,不是能靠运维补救的缺口。它也不承诺”OK 返回即已投递”,源码注释写得很明白:NIP-01 的 OK 表示事件已被持久接受,不代表跨节点发布、本地扇出、工作流触发已经完成。
自建一个中继,你要清楚暴露面。router.rs 里除了 WebSocket 入口,还挂着媒体上传(PUT /upload、PUT /media/upload)与读取、git smart HTTP 路由、NIP-05 的 /.well-known/nostr.json、健康探针、以及可选的 admin 子路由(/api/admin/v1,只在 admin 主机上响应,其它情况一律 404)。事件本体落在 Postgres(migrations/ 下 27 份 SQL 就是这套表结构),跨节点扇出与在线状态依赖 Redis——readiness 探针会同时 ping 这两者,2 秒超时,任一不通就报 not_ready。CORS 默认行为要注意:BUZZ_CORS_ORIGINS 未设置时是 permissive,设置了但一个都解析不出来时,代码宁可返回一个空 CORS 层也不回落到 permissive,日志里会有明确的 error。
还有一件与 Agent 有关的事必须说清:在这套模型里,人和 Agent 都是持有密钥的主体,Agent 拿到密钥就能以自己的身份发事件、被别人订阅到。仓库里对 Agent 主体有独立的配额档位(agent_standard_messages_per_min 与 human_messages_per_min 是两个不同的限额键),说明设计上就假定 Agent 的发送行为需要单独约束。给 Agent 发凭据之前,先想清楚它能往哪些频道写、这些事件谁会收到。
六、上手与避坑清单
别拿 router.rs 找消息路由。 会踩是因为名字和直觉冲突:多数网关项目里 router 就是消息分发。避法是记住入口顺序——router.rs 只到”升级成 WebSocket”,之后的路在 connection.rs,收件人在 subscription.rs。
订阅收不到消息,先确认它进了哪张索引,而不是先怀疑过滤器。 会踩是因为 #h 标签不合法、或者多个过滤器写了不同频道时,订阅会静默退化成全局订阅,而全局订阅按不变量收不到频道事件。避法是照着 extract_channel_id_from_filters 的规则逐条核对过滤器,确保每个过滤器都带同一个可解析的频道 UUID。
kinds: [] 不是通配。 会踩是因为很多客户端库把空数组当”不限制”。避法是不需要 kind 约束时直接不写 kinds 字段,别写空数组——写了就等于把这条订阅静默作废。
REQ 之前必须先完成认证。 会踩是因为 NIP-42 的挑战是连接建立后由中继主动发的,客户端如果不处理这条 AUTH 帧就直接发 REQ,只会拿到 auth-required 的 NOTICE 加 CLOSED;再拖过 5 秒还会被认证超时踢掉。避法是把”收挑战、签名、回 AUTH”做成连接建立流程的一部分,不要当可选项。
帧太大会在两处被拒。 会踩是因为客户端往往只根据一处报错调参:解析器层(升级时设的 max_message_size)直接不交给应用,应用层再查一遍并回 NOTICE。避法是按 NIP-11 文档里的 max_message_length 分片,别拿”我这条 JSON 不算大”估。
别把慢客户端当无害。 会踩是因为发送通道满时中继不阻塞、只累加计数,看起来一切正常,直到累计到宽限上限连接被取消。避法是客户端及时读、订阅收窄,服务端侧观察背压相关的断连指标。
读源码时注意主机名的作用。 会踩是因为本地起服务习惯用 localhost 随便访问,而这套实现里 Host 是社区绑定的权威选择器,未映射的 Host 在升级前就被拒。避法是先确认 Host 到社区的映射配好了,再排查协议层。
别把 VISION 文档当已实现清单。 会踩是因为根目录 8 份 VISION 开头的 md 读起来很像功能说明,但那是项目自己的愿景定位;docs/ 下另有 53 个文件。避法是任何行为判断都回代码核,文档只用来理解意图。想补协议侧的背景,可以看协议是什么这类基础篇和Agent 协议生态的取向差异,本篇不做横向排序。
接下来该读哪个文件,取决于你卡在哪一步:客户端连不上或拿不到预期响应,从 router.rs 的 nip11_or_ws_handler 往下读;发出去别人没收到,先读 handlers/event.rs 的 filter_fanout_by_access,再回 subscription.rs 看索引选择;历史消息拉不全,读 handlers/req.rs 里按过滤器逐个建查询、去重、发 EOSE 的那段。三个问题、三条路径,别混着查。
本篇属于一个把开源多 Agent 通信平台 buzz逐层拆开讲的系列,整体地图见 buzz 是什么:Block 开源的多 Agent 通信平台全景图;沿着这条线往下,还可以看 Block 开源多 Agent 通信平台 buzz:engram 把记忆做成事件之后,检索边界在哪 和 谁能连上你的 buzz 中继:Block 开源多 Agent 通信平台的三层门禁。