Block 开源多 Agent 通信平台 buzz 的工作流执行器拆解
本文基于 buzz 仓库 commit 8342dfc(2026-08-04)梳理,该项目仍在高频迭代,具体行为以仓库 https://github.com/block/buzz 最新代码与文档为准。
buzz 的工作流执行器最值得抄走的一点,不是它支持哪几种动作,而是它把「自动化跑出来的东西」和「人打的字」压成了同一类对象:一条由中继签名、带 buzz:workflow 标记的 kind:9 消息事件。 自动化不是挂在系统旁边的一根管子,它的输出直接落进所有人和所有 Agent 都在读的那条时间线里。这个前提一旦定死,后面每一个设计——步骤 ID 只准用下划线、延时不能超过 270 秒、动作出口只留一个方法——都是被它逼出来的。
先说这篇跟站内其它几篇的分工。如果你要找的是编排范式本身,比如工作流该怎么拆、多个 Agent 怎么并发、任务粒度切到多细,站内另有工作流编排的写法、Agent 并发编排、任务分解粒度三篇。这一篇不重复它们,只盯 buzz 这一个仓库的具体实现,看「结果即事件」这个前提会把代码逼成什么形状。
一、地基:Nostr 这几个词先讲清楚
buzz 建在 Nostr 协议之上。这里只需要四个概念,一句话一个。
事件:网络里流通的最小单位,一条 JSON 结构,带内容、时间戳、若干标签和一个 kind 数字,kind 表示这条事件是什么类型。签名:每条事件由发布者的私钥签名,谁签的可验证,事后改不动。密钥自持:身份就是那对密钥,不存在”找回账号”这回事,私钥丢了身份就丢了。中继:负责接收、存储、转发事件的服务端;在 buzz 里中继同时也是工作空间本体。另外 NIP 是 Nostr 的协议提案编号,作用类似 RFC,仓库 docs/nips/ 下自己写了 15 份 NIP 规范 md 加 2 份 fixtures json,crates/buzz-core/src/pairing/NIP-AB.md 是第 16 份。
kind 数字在 crates/buzz-core/src/kind.rs 里集中定义。跟工作流相关的几个:KIND_STREAM_MESSAGE 是 9,就是频道里的普通消息;KIND_REACTION 是 7,表情回应;KIND_STREAM_MESSAGE_DIFF 是 40008,代码 diff 消息;工作流自身的执行事件占 46001 到 46012 这一段,其中 46010 是 KIND_WORKFLOW_APPROVAL_REQUESTED。整个仓库 crates/ 下有 28 个 crate,工作流这条线主要落在其中两个:buzz-workflow 和 buzz-relay。许可证是 Apache-2.0(Copyright 2026 Block, Inc.)。
四块源码的分工是这样:
| 组成部分 | 它负责什么 | 对应仓库位置 | 你什么时候会碰到它 |
|---|---|---|---|
| 定义层 schema | YAML 反序列化成 WorkflowDef,校验触发器、步骤 ID、cron 与 interval | crates/buzz-workflow/src/schema.rs | 写 YAML 被 InvalidDefinition 打回时 |
| 执行器 executor | 模板替换、if: 条件求值、串行派发、超时、执行轨迹 | crates/buzz-workflow/src/executor.rs | 步骤跳过了、超时了、变量没替换上时 |
| 动作出口 ActionSink | 一个 trait,把”产生副作用”这件事从执行器里抽出去 | crates/buzz-workflow/src/action_sink.rs | 想加一种新的对外动作时 |
| 中继侧接收端 | 实现 ActionSink:建事件、签名、落库、扇出 | crates/buzz-relay/src/workflow_sink.rs | 排查消息为什么没发出去、@ 没唤醒 Agent 时 |
| 引擎与调度 | 事件触发入口 on_event、cron 循环、并发闸门、权限复核 | crates/buzz-workflow/src/lib.rs | 定时任务没按点跑、被并发拒绝时 |
| 落库结构 | workflows、workflow_runs、workflow_approvals、scheduled_workflow_fires 四张表 | migrations/0001_initial_schema.sql | 想知道数据到底落在哪时 |
二、定义层:一份 YAML 能说什么,以及它故意不让你说什么
WorkflowDef 一共五个字段:name、description、trigger、steps、enabled。触发器 TriggerDef 是五选一:message_posted(可带 filter 表达式)、reaction_added(可带 emoji)、diff_posted(对应 kind:40008,可带 filter)、schedule(cron 与 interval 二选一)、webhook(HTTP POST 打到 /hooks/{id})。动作 ActionDef 是七选一:send_message、send_dm、set_channel_topic、add_reaction、call_webhook、request_approval、delay。
一份最小的定义长这样,这段 YAML 直接取自 lib.rs 里的往返测试:
name: "Test Workflow"
trigger:
on: message_posted
steps:
- id: s1
action: send_message
text: "Hello {{trigger.author}}"
真正有意思的是 validate() 拒绝了什么。步骤 ID 必须非空、长度不超过 64、且只允许字母数字和下划线。原因写在注释里:步骤 ID 会被拼成条件表达式的变量名 steps_{id}_output_{field},如果你写 my-step,表达式引擎会把它读成 steps_my 减去 step_output_field。仓库里专门有一条测试盯这个坑。
schedule 触发器要么给 cron 要么给 interval,给两个直接报错。cron 表达式先过一层归一化:这里用的 cron 库要 7 段(秒 分 时 日 月 周 年),标准 5 段会补上前面的 0 和后面的 *,6 段只补年。interval 有个硬下限——低于 60 秒直接拒,理由是调度循环本身就是 60 秒一跳,写 30s 永远不可能按点触发。把物理上做不到的配置在定义时就拦掉,而不是让它在运行时表现得时灵时不灵,这个取舍值得单独记一笔。
还有一个方法专门服务于权限判断:
pub fn requires_elevated_authority(&self) -> bool {
self.steps
.iter()
.any(|s| matches!(s.action, ActionDef::CallWebhook { .. }))
}
只要有一步是 call_webhook,整份定义就被标记为需要更高权限。注释解释得很直白:工作流是拿着所有者的长期权限在跑的,保存之后很久还在往外转发频道内容,所以普通成员身份不够,保存和运行都要频道 owner 或 admin。这个判断在 lib.rs 的 check_owner_authority 里被复用,事件触发、定时触发、webhook 触发三条入口在创建运行记录前都要再查一次所有者当前的角色,查不到就拒——注释里写的是 fail-closed。这块和站内讲的最小权限设计是同一个思路,只是它把复核点放在了每一次触发。
三、执行器:模板、条件、串行,外加一圈硬上限
执行器做四件事:解析模板变量、求值 if: 条件、顺序派发步骤、把轨迹写回数据库。
模板支持 {{trigger.X}} 和 {{steps.ID.output.FIELD}} 两种路径,两个过滤器:truncate(N) 按字符截断,npub 把十六进制公钥编成完整的 bech32 形式(truncate_pubkey 是它的旧别名)。这里有个细节值得看:npub 过滤器的注释说明了为什么不截断公钥前缀——短前缀是可以被暴力磨出来的,截断反而制造混淆。另一个行为要记住:未知变量不会报错,会原样输出。你把 {{trigger.tekst}} 拼错了,发出去的消息里就带着那对花括号。
条件用的是 evalexpr 表达式引擎。因为它不支持带点的标识符,变量名全部改成下划线:trigger.text 在条件里写作 trigger_text,步骤输出写作 steps_STEP_ID_output_FIELD。执行器还自己注册了四个字符串函数补上库里没有的能力:str_contains、str_starts_with、str_ends_with、str_len。
条件求值这段的防御写得很细。webhook 传进来的字段会先注册成 trigger_前缀 的变量,标准触发字段随后注册,后写覆盖先写——注释明确说这是为了让外部传入的字段无法冒充标准字段;键名本身以 trigger_ 或 steps_ 开头的还会被直接跳过。表达式长度上限 4096 字节,超了报错;求值本身丢到阻塞线程池并套一层 100 毫秒超时,注释坦白说明 evalexpr 不是为对抗性输入设计的。
串行派发这一层的上限更多。每个步骤的超时优先取自身的 timeout_secs,没写就用引擎配置的 default_timeout_secs,默认 300 秒。delay 动作被限制在 270 秒以内,理由写在常量旁边:必须小于默认步骤超时,否则会撞出不确定的超时失败,长等待要走调度恢复的路子。并发闸门是一个信号量,默认 max_concurrent 是 100,且是 try_acquire ——满了立刻返回 CapacityExceeded,不排队。
派发结果只有三种:Completed 带一个 JSON 输出、Suspended 带一个审批令牌、Skipped。每一步的结果都追加进 trace 数组,条件为假的步骤也会留下一条 status: skipped 的记录。失败时执行器返回的是错误加上一个 PartialProgress,让调用方能把失败前已经跑完的轨迹存下来——这是可观测性上很实在的一笔,比只记一个”失败”有用得多。
有几件事要说清楚它现在还没做完:send_dm 和 set_channel_topic 直接返回 NotImplemented;request_approval 会生成一个 UUID v4 令牌并返回 Suspended,但 finalize_run 里对这种情况的处理是记一条警告然后把运行标成失败,注释写着审批门还没实现。所以如果你要在 buzz 上做人类介入的审批环节,现在只能自己接。
四、出口:只留一个方法的 trait,和中继侧那个接收端
action_sink.rs 整个文件不到 70 行,核心就是一个 trait,一个方法:
pub trait ActionSink: Send + Sync {
fn send_message(
&self,
community_id: CommunityId,
channel_id: &str,
text: &str,
author_pubkey: &str,
) -> Pin<Box<dyn Future<Output = Result<String, ActionSinkError>> + Send + '_>>;
}
注释说明了它为什么存在:早先执行器是通过 HTTP 回环打自己中继的 REST 接口来发消息的,结果撞上 401 鉴权失败,于是改成由中继实现这个 trait,直接给执行器数据库和事件通道。错误类型只有六种,全都是这条路径上真会发生的:输入非法、频道不存在、频道已归档、事件构造失败、数据库错误、内容为空。
把副作用收敛到一个 trait 上,好处是执行器不需要知道消息是怎么变成事件的。 沉在下面的 crates/buzz-relay/src/workflow_sink.rs 才是”把自动化做成网络事件”这句话真正兑现的地方。它一步步做的事是:
先把弱引用升级成 AppState(用 Weak 是为了破掉 AppState → 引擎 → sink → AppState 的循环引用);用运行记录自带的 community 反查 host,注释特别强调不能从配置里的中继地址重新推导租户,否则 B 社区的工作流会把消息发进默认社区。然后校验内容非空、频道 UUID 合法、频道未归档,再查所有者是不是该频道成员——不是成员且频道不是 open,直接拒。
接着构造 kind:9 事件,三个基础标签:p 标签指向工作流所有者,h 标签按 NIP-29 的写法把消息限定到频道(NIP-29 是 Nostr 里描述”中继托管的群组”的那份提案,h 就是它约定的群组归属标签,这里放的是频道的规范化 UUID),buzz:workflow 标签标记这是工作流产物。
// 摘自 workflow_sink.rs,省略了每行后面的错误映射
let mut tags = vec![
Tag::parse(["p", &author_pubkey_hex]),
Tag::parse(["h", &channel_id_canonical]),
Tag::parse(["buzz:workflow", "true"]),
];
这三个标签,第三个是防递归的关键。中继侧的事件处理路径里有一条判断:如果事件是中继密钥签的并且带 buzz:workflow 标签,就跳过工作流触发。加上 46001–46012 这段执行事件本身也被排除,工作流才不会自己触发自己。
最后是 @Name 提及解析。因为客户端发消息时是从自动补全里选人、直接带上 p 标签的,而工作流只有一段纯文本,所以这里得反向解析。四条规则在注释里写得很硬:只在目标频道成员里匹配;只认完整显示名,不做前缀和模糊匹配;同名先匹配最长的且吃掉整段;同一个显示名对应不止一个公钥时,一个都不 tag。最后这条的理由是错唤醒比不唤醒更糟。测试里甚至覆盖了土耳其语大写 İ 小写化后变成两个码点导致索引错位的情况——这类 bug 一旦发生,表现是”后面所有的 @ 都突然失灵”,非常难查。
事件签好、落库带上线程元数据(工作流消息一律是顶层,depth 为 0),只有在确实插入成功时才走后续的扇出、搜索索引和审计。
五、边界与代价:它明确不管什么
签名主体不是你。 工作流发出的消息是中继密钥签的,p 标签只是归属标注。如果你带着 Nostr 那套”事件由本人私钥签名”的直觉去做审计,会把这类消息的责任主体判断错。反过来说,这也意味着工作流的发言能力不依赖所有者在线,代价是中继密钥变成了一个高价值目标。密钥自持在这里是双向的:私钥丢了,身份就没了,没有客服能帮你找回。
顺着这条往下想一层就是自建中继的账要怎么算。工作流的定义、每一次运行的轨迹、审批令牌、调度声明行,全都落在中继自己那套 migrations/ 建出来的库里,不在别处;工作流发出去的每条消息又是中继密钥签的。也就是说,谁运维这台中继,谁就同时握着这个工作空间的全量消息、全部工作流定义,以及那把能代表所有人发言的中继私钥。把中继端口暴露到公网,暴露面不只是”聊天记录可能被读”,还包括工作流定义里可能写死的 webhook 目标地址和请求头——call_webhook 的 headers 字段接受任意键值,凭据填在那里就是明文存在库里的一行 JSON。这台机器的备份策略、磁盘加密、密钥文件权限,得按”托管别人身份”的标准来定,而不是按一个聊天服务。
它不是分布式工作流引擎。 步骤是串行的,一份定义里的步骤按顺序跑完,没有分支跳转、没有并行分叉、没有循环。想表达 DAG,靠一份 YAML 表达不出来。
没有重试、没有补偿。 某一步失败就整个运行失败,轨迹留下,但不会回滚已经发出去的消息。事件已经进了网络,就不可能撤销。
定时触发不补跑。 内存里的 interval 时钟重启即丢,注释直说 MVP 阶段停机期间错过的触发不会重放。跨进程的去重靠 scheduled_workflow_fires 那张表的声明行,保证的是”至多一次”而不是”恰好一次”——运行记录插入失败时,声明行仍然保留,注释写明这是刻意选择。
外呼是要付代价的。 call_webhook 那条路径做了不少防护:先解析域名并拒绝私有和保留地址,把校验过的 IP 钉进 HTTP 客户端以防 DNS 重绑定,禁用系统代理、禁用重定向(重定向到内网就绕过检查了),10 秒超时,响应体分块读取、超过 1 MiB 就中断。这些防护本身说明了风险有多实在:一份带 call_webhook 的工作流,等于给频道内容开了一个长期有效的对外出口。
功能缺口要认。 私信、改频道主题两个动作没实现;审批门生成令牌但没接上,撞到就是失败。还有一个安静的坑:add_reaction 和 call_webhook 的真实实现都在 reqwest 这个 feature 后面,没启用时它们不报错,而是返回一个带 skipped: true 的成功输出。步骤看起来绿了,实际什么都没发生。
六、上手与避坑清单
步骤 ID 里出现连字符。 为什么会踩:YAML 里写 id: my-step 看着完全正常。怎么避:记住 ID 会被拼成表达式变量名,只用字母数字和下划线,长度控制在 64 以内,validate 会在保存时就打回来,别等到运行。
模板和条件混用两套变量写法。 为什么会踩:模板里是 {{trigger.text}} 带点,条件里必须写 trigger_text 带下划线,同一个东西两种拼法。怎么避:写条件时先在脑子里把点换成下划线;步骤输出更绕,steps.ask.output.replied 到条件里是 steps_ask_output_replied。
变量名拼错了却毫无提示。 为什么会踩:未知变量原样输出,没有任何错误。怎么避:新工作流第一次跑必须在测试频道跑,肉眼看实际发出去的文本里有没有残留的花括号。
delay 写成小时级。 为什么会踩:duration 字段接受 1h 这样的写法,看着像支持长等待。怎么避:上限是 270 秒。注意这道检查不在 validate() 里,而是在执行器真正派发这一步时才抛出(错误类型虽然叫 InvalidDefinition,却是运行期才暴露),保存的时候不会有任何提示。真要等几小时,拆成两个工作流,用调度触发第二个。
interval 写成秒级。 为什么会踩:以为调度精度可以自己定。怎么避:调度循环 60 秒一跳,低于 60s 的 interval 在校验期就被拒,别浪费时间试。
给工作流加了一步外呼,然后所有者被移出了频道。 为什么会踩:call_webhook 会让整份定义需要 owner 或 admin,而权限是在每次触发前重查的。怎么避:加外呼之前先确认所有者的角色够、并且不会被调整;权限一掉,工作流是静默停跑的,日志里只有一条 warn。
@ 了 Agent 却没唤醒它。 为什么会踩:唤醒是靠 p 标签的,而 p 标签靠精确匹配显示名生成;名字打错一个字、或者频道里有两个同名成员,结果都是零匹配。怎么避:给 Agent 起唯一且不带歧义的显示名,改名之后回头检查引用它的工作流文本。
收尾
把自动化做成网络里的事件,换来的是一致性:一条工作流发的消息和一条人打的字,在存储、订阅、扇出、审计上走的是同一条路径,谁也不需要为自动化单独造一套通道。付出的代价是所有事件系统都得付的那份——不可撤销、不可重放、必须显式防递归、必须显式防止自己把自己的输出当输入。
想继续往下读,建议按这个顺序:先 crates/buzz-workflow/src/schema.rs 看定义边界,再 executor.rs 看那几个硬上限分别在挡什么,然后 crates/buzz-relay/src/workflow_sink.rs 从头读到尾——这个文件是”结果即事件”最完整的一次落地。最后回到 crates/buzz-workflow/src/lib.rs,看 on_event 和 cron 循环各自在什么时刻做权限复核。
自检三问:你的工作流输出会不会触发另一个工作流,防递归靠的是什么?某一步失败时,已经发出去的消息你打算怎么处理?带外呼的定义,所有者权限掉了之后你从哪里能看出来它停了?这三个问题在任何”结果即事件”的系统里都得有答案,不只是在 buzz 里。
本篇属于一个把开源多 Agent 通信平台 buzz逐层拆开讲的系列,整体地图见 buzz 是什么:Block 开源的多 Agent 通信平台全景图;沿着这条线往下,还可以看 Block 多 Agent 平台 buzz:persona 包四道工序 和 Block 开源 buzz 多 Agent 平台的设备配对:一份 NIP 规范加一份形式化模型。