12 · 工作流引擎
buzz-workflow 把频道事件、reaction、diff、schedule 和 webhook 变成 YAML 自动化。它拥有 schema、表达式/模板、并发控制、run trace 与部分 action;但当前 approval、send DM 和设置频道 topic 明确未完成。理解这个边界比只看 schema 支持列表更重要。
1. 模型
name: Incident alert
trigger:
on: message_posted
filter: 'str_contains(trigger_text, "P1")'
steps:
- id: notify
action: send_message
text: 'P1: {{trigger.text}}'Workflow definition 是频道绑定的 parameterized replaceable event;数据库索引可执行定义,run/trace 则进入关系表,不要求把每一步都重放成 Nostr event。
2. Trigger
Schema 支持:
message_postedreaction_addeddiff_postedschedule(cron 或 interval,interval 至少 60 秒)webhook
on_event 会排除 workflow 自己的 execution kinds,避免 workflow 产出的消息再次形成无界自触发环。具体 filter 还可匹配文本、emoji 等 trigger context。
3. 缓存与 authority recheck
Engine 使用约 10 秒 workflow cache 降低每事件 DB 查询,但运行前重新检查 owner authority:
- owner 必须仍是频道成员。
- 有外泄能力的
call_webhook要求 owner/admin 级别。 - community、workflow、channel ID 查询全都带租户坐标。
缓存只决定候选定义,不决定最终授权。成员在缓存有效期内被撤权,也不能继续启动高权限动作。
4. 并发与状态
WorkflowConfig 默认全局最大并发 100、step timeout 300 秒。Engine 用 semaphore 限制 run,并把 (community, workflow) 作为 last-fired/cache key,防不同租户同 UUID 互相影响。
运行状态大致:
Pending → Running → Completed
├→ Failed
└→ (设计上 WaitingApproval,当前主动 Failed)每个 step 生成 trace,finalize_run 是 executor result 到 DB status 的单一映射点,减少 manual/event/approval resume 三条入口状态漂移。
5. Action 完成度
| Action | 当前状态 | 说明 |
|---|---|---|
send_message | 已实现 | 经 Relay ActionSink 构造、签名、持久化、审计与正常 fan-out |
call_webhook | 已实现 | owner/admin gate、SSRF 防护、10s timeout、1 MiB response cap |
add_reaction | 已实现但依赖部署 | reqwest feature + Relay HTTP/API token 环境;否则显式 skipped |
delay | 已实现 | 最长 270 秒,长延迟留给未来 scheduled resume |
send_dm | 未实现 | executor 返回 NotImplemented("SendDm") |
set_channel_topic | 未实现 | executor 返回 NotImplemented |
request_approval | 基础结构存在,闭环未实现 | 生成 token/suspended,finalizer 主动把 run 标 Failed |
6. Relay-signed send_message
工作流不能持有用户私钥。RelayActionSink 使用 Relay/服务端签名发布消息,同时保留:
- workflow 来源 tag。
- owner/actor attribution。
hchannel。- 解析
@Name后的pmentions。
它验证目标频道存在、未归档、owner 当前有权,并走 DB insert + 正常 post-commit dispatch。这样 workflow 消息不会绕过频道权限或实时分发。
为避免循环依赖 AppState → WorkflowEngine → ActionSink → AppState 泄漏,sink 持有 Weak<AppState>;Relay shutdown 时 upgrade 失败并返回明确错误。
7. call_webhook 的 SSRF 防护
实现不是只检查字符串 localhost:
- 解析 URL/host/port。
- DNS resolve。
- 拒绝 private/reserved IP。
- 将已检查的 IP pin 到 reqwest client,避免 DNS rebinding TOCTOU。
no_proxy,避免系统代理重新解析原 host。- 禁止 redirect,避免跳转内网。
- 10 秒超时,response 分块读取并限制 1 MiB。
这是很成熟的外呼边界,但仍需 egress firewall 作为第二层,并限制用户可设置的 header/secret 来源。
8. Approval 为什么不能算完成
Executor 的 RequestApproval 会生成 OS CSPRNG-backed UUID token 并返回 Suspended;数据库也有 approval table、run status,CLI 也有 approve 命令。但 finalize_run 明确记录 WF-08 未完成,把遇到 gate 的 run 标为 Failed,避免创建永远无法恢复的 WaitingApproval。
这是合理的 fail-explicit:比 UI 看似“等待审批”但永不继续更诚实。
完整闭环还需:
- 事务内保存 token、run、step index/trace。
- 发送 approval request event。
- token 单次消费、租户/审批人/过期验证。
- approved/denied 分支恢复。
- 崩溃后 scheduler 扫描和幂等继续。
9. 调度与交付语义
Schedule 需要防多 pod 重复触发。Engine/DB 使用 last-fired/claim 语义把执行记录租户化;即便触发重复,run/actions 也应有幂等边界。delay 只允许短睡眠,是因为把数小时暂停放在 Tokio task 中无法跨进程重启恢复。
10. 源码入口
crates/buzz-workflow/src/schema.rs:YAML schema 与验证。crates/buzz-workflow/src/lib.rs:Engine、cache、semaphore、authority 与 finalizer。crates/buzz-workflow/src/executor.rs:step executor 与 action 完成度。crates/buzz-workflow/src/action_sink.rs:副作用接口。crates/buzz-relay/src/workflow_sink.rs:Relaysend_message实现。crates/buzz-db/src/workflow.rs:definition/run/approval 状态存储。