ZeroClaw SOP Fan-In:用 AMQP 投递触发标准作业流程
【免费下载链接】zeroclawFast, small, and fully autonomous AI personal assistant infrastructure, any OS, any platform — deploy anywhere, swap anything 🦀项目地址: https://gitcode.com/gh_mirrors/ze/zeroclaw
导读
本文讲解 ZeroClaw 中将 AMQP 0-9-1(RabbitMQ 及其兼容 broker)消息投递转换为 SOP(Standard Operating Procedure,标准作业流程)运行触发器的完整链路:从[channels.amqp.<alias>]的 broker 连接配置、dispatch路由模式、routing key 通配匹配与condition表达式,到投递被"提升"为 SOP 事件后的引擎分发、背压重投递与审批观察。读完本文,你将能够把任意 AMQP 主题上的业务事件(发布监控、CI 通知、告警流等)接入 ZeroClaw SOP 引擎,让事件自动驱动需要人工审批的作业流程。
一、AMQP 在 SOP Fan-In 中的角色
ZeroClaw 的fan-in是指"外部事件源启动 SOP 运行"的机制:每个事件源通过dispatch_sop_event将事件交给 SOP 引擎,引擎把事件与所有已加载 SOP 的触发器逐一匹配,命中即启动运行。一个实例可以同时绑定多个 fan-in——MQTT 主题、文件系统路径、AMQP routing key 可以同时喂给同一个引擎,无需独立进程(见 fan-in 总览)。
AMQP 是其中一类事件源:当某个 alias 以 SOP 派发模式运行时,AMQP 消费者会把每条投递"提升"为一个 SOP 事件——routing key 成为事件 topic,消息体成为事件 payload——然后交给引擎匹配。也就是说:
- 传输侧(broker 连接、队列、交换器、TLS)由 AMQP channel 负责;
- 触发侧(routing key 匹配、condition 求值)由本文对应的 SOP 触发器负责;
- 决定一条投递驱动 agent 循环、SOP 引擎还是两者兼顾的,是 channel 的
dispatch字段。
从源码看,这一职责划分非常清晰:crates/zeroclaw-channels/src/amqp.rs中AmqpChannel持有dispatch: SopDispatch、engine: Option<Arc<Mutex<SopEngine>>>与audit: Option<Arc<SopAuditLogger>>三组状态,route_delivery(amqp.rs 第 142 行)根据dispatch的值分别把投递送入 agent 循环(ChannelMessage)或 SOP 引擎(dispatch_untrusted_fan_in)。
安全前置条件:
dispatch为sop或sop_and_agent_loop时,若构造 channel 时没有提供 SOP engine/audit 句柄,AmqpChannel::new会直接bail拒绝启动(amqp.rs 第 86-95 行),避免出现"确认了投递却没人派发"的静默丢消息。测试new_rejects_sop_dispatch_without_handles专门守护了这一 fail-closed 行为。
二、传输侧配置:[channels.amqp.<alias>]
要接入 SOP 触发,先要有一个可用的 AMQP channel。完整的字段清单(来自运行时 schema,schema.rs 第 17734 行起):
| 字段 | 默认值 | 作用 |
|---|---|---|
enabled | false | 是否启用该 channel。运行时只加载enabled = true的 channel;默认关闭,避免粘贴半截配置就意外上线 |
amqp_url | 无 | broker 地址。明文用amqp://,TLS 用amqps://(如amqps://fedora:@rabbitmq.fedoraproject.org/%2Fpublic_pubsub),属于 secret 字段 |
exchange | 无 | 消费者队列要绑定的交换器(如amq.topic),必填 |
routing_keys | [] | 要绑定的 routing key 列表。建议收窄到感兴趣的主题;绑定#会消费整个交换器,几乎永远不是你想要的,必填至少一个 |
queue | 无 | 队列名。留空则使用服务端生成的临时独占队列;仅在需要跨重连持久投递时才设置稳定名称(推荐 UUID) |
ca_cert | 无 | amqps://连接用的 CA 证书 bundle 路径 |
client_cert/client_key | 无 | 双向 TLS(mTLS)的客户端证书与私钥,必须成对出现(Fedora Messaging 就要求客户端证书) |
sender_label | "amqp" | 写入每条投递ChannelMessage.sender的标识(如anitya),供编排器自环防护与按 channel 路由识别来源 |
content_template | 空 | 入站消息内容模板,{field}占位符从 JSON 投递体顶层键插值;为空时原样使用投递体 |
thread_id_field | 空 | 指向 JSON 投递体的点分路径,其值作为消息thread_ts用于关联回复(如message.project.name);空则关闭线程关联 |
durable_ack | true | 确认模式:true时投递只在消息被可靠移交给 agent 循环后才 ack(至少一次语义,崩溃会重投);false时 broker 派发即确认(至多一次),仅适合无副作用、可丢弃的消费者 |
dispatch | agent_loop | 投递路由去向:驱动 agent 回合(默认)、派发到 SOP 引擎(sop)、或两者同时(sop_and_agent_loop) |
excluded_tools | [] | 不向该 channel 的工具规格暴露的工具列表 |
配置校验(AmqpConfig::validate,schema.rs 第 17826 行)会强制以下规则:
amqp_url必须以amqp://或amqps://开头;- 使用
amqps://时ca_cert必须提供(否则报错 "amqps:// requires ca_cert to verify the broker"); client_cert与client_key必须同时设置或同时缺省(mTLS 成对校验);exchange不能为空;至少配置一个 routing key。
一个最小化的 SOP 触发用配置示例:
[channels.amqp.release] enabled = true amqp_url = "amqps://fedora:@rabbitmq.example.org/%2Fpublic_pubsub" exchange = "amq.topic" routing_keys = ["org.release-monitoring.prod.anitya.project.version.update"] ca_cert = "/etc/zeroclaw/certs/ca.pem" # 将投递交给 SOP 引擎(而不是 agent 循环) dispatch = "sop"Dispatch 三种模式
dispatch字段决定一条投递做什么(channels/amqp.md):
agent_loop(默认):投递作为一条消息交给 agent 循环,保持原有行为,存量消费者不受影响;sop:投递被提升为 SOP 事件(routing key → 事件 topic,消息体 → payload)并派发到 SOP 引擎;sop_and_agent_loop:每条投递同时执行上述两者。
TLS 与 mTLS
- TLS:
amqp_url指向amqps://端点并配置ca_cert; - 双向 TLS:额外设置
client_cert与client_key。
不配置这些就是明文连接,不要把明文消费者暴露在不可信网络上。底层实现上,AmqpChannel::connect(amqp.rs 第 255 行)会把 PEM 格式的客户端证书与私钥转换为内存中的 PKCS#12 bundle(pem_to_pkcs12_der,amqp.rs 第 373 行)交给 rustls 客户端认证路径;client_cert/client_key只设置其一会在构造期直接报错。
三、触发侧:AMQP 触发器与 routing key 匹配
SOP 的 AMQP 触发器在SOP.toml的[[triggers]]中声明。从测试代码可见其结构(amqp.rs 第 912 行):
SopTrigger::Amqp { routing_key: "anitya.update".into(), condition: None, }对应到SOP.toml:
[[triggers]] type = "amqp" routing_key = "anitya.update" # condition = "$.value > 85" # 可选routing_key采用AMQP topic-exchange 语义:
- key 以
.分隔为多个词(word); *精确匹配一个词;#匹配零个或多个词。
例如org.release-monitoring.prod.anitya.project.version.update这类多点 key 可直接作为触发器模式;发布方的 routing key 与触发器模式匹配即视为命中候选。
匹配时,消息体被原样转发进 SOP 事件 payload,供可选的触发器condition求值;而进入 step 上下文的则是经过截断、清洗、框架化的受限形式(详见下文"安全默认")。一个 JSON-path 的condition(如$.value > 85)要求发布方发送 JSON 消息体。
四、Condition 条件表达式
触发器的condition与 step 的when:守卫共用同一套表达式语法(sop/syntax.md),触发条件针对事件 payload 求值。求值失败即关闭:非法条件、缺失 payload、无法解析的 JSON path、以及两侧不是数字的直接数值比较,一律判为不匹配;空条件无条件匹配。
JSON Path 形式
以$开头,比较 JSON payload 内的某个值:$.path.to.field <op> <value>。常用示例(来自 syntax 参考的官方匹配表):
| 表达式 | Payload | 是否匹配 |
|---|---|---|
$.value > 85 | {"value":90} | 是 |
$.status == "critical" | {"status":"critical"} | 是 |
$.data.sensor.value > 85 | {"data":{"sensor":{"value":87.3}}} | 是 |
$.readings.1 == 20 | {"readings":[10,20,30]} | 是 |
$.nonexistent > 0 | {"value":90} | 否 |
路径规则:
- 使用点分隔的段;数组元素用数字段,如
$.readings.1,不支持方括号语法; - 缺失键、越界索引、非法 JSON、空 payload 一律失败关闭;
- 没有通配符、过滤器、递归下降或内置变量。
直接数值形式
无前导$的条件把整个 payload 当作数值比较,适合标量事件 payload:如> 0、>= 5、== 42、> 3.14。若任何一侧解析不出数字则无匹配。
运算符与注意事项
- 支持的运算符:
==、!=、>、>=、<、<=; - 解析器按最长优先匹配运算符 token;JSON path 比较先尝试数值比较——两侧都能解析为数字则按数值比较,否则按字符串比较;
- 比较值两侧的英文双引号会被剥离,所以字符串字面量要加引号:
$.status == "critical"; - JSON 布尔值会被转换为字符串
true/false,因此用引号字符串比较:$.active == "true"; - 一个条件只允许单个比较,不支持
AND/OR/NOT逻辑组合。
五、开火:把事件送入 SOP 引擎
端到端触发一条 AMQP 驱动的 SOP 运行,步骤为:
- 确认
sops_dir已配置。SOP 定义从sops_dir下的子目录加载,该字段默认未设置(运行时 SOP 关闭),需要显式开启:相对路径相对于安装根目录(config.toml所在目录)解析,文档化值为shared/sops(即<install>/shared/sops),也可用绝对路径或~前缀路径。 - 设置 channel 的
dispatch为 SOP 模式(sop或sop_and_agent_loop)。 - 加载 SOP:在
<sops_dir>/<name>/下放置SOP.toml(含[[triggers]] type = "amqp")与可选的SOP.md,用zeroclaw sop validate <name>验证。 - 发布消息:向配置的 exchange 发布一条 routing key 与触发器匹配的消息。消费者把投递提升为事件(routing key → topic,body → payload)并派发。每个已加载且 routing key 模式匹配、
condition(若有)对 body 成立的 SOP 都会启动一个运行。
如果什么都不启动,按顺序排查(也见 fan-in 总览的故障排查表):
dispatch是否是 SOP 模式(sop或sop_and_agent_loop)而非默认的agent_loop;- 队列是否正确绑定,使 routing key 真正到达消费者(检查
exchange、routing_keys与发布方实际发射值是否一致); condition是否对 payload 成立(可用$.value > 85这类带 JSON body 的发布做快速验证)。
底层派发链路
消费者循环(AmqpChannel::listen,amqp.rs 第 444 行)逐条取出投递后调用route_delivery;当routes_sop为真时,调用dispatch_untrusted_fan_in(dispatch.rs 第 1209 行)——这是一个兼容包装,内部走SopIngress统一入口,执行:
- 用
cap_untrusted对 topic 与 payload 做长度截断(上限由sop.untrusted_payload_max_bytes控制,默认8192字节,按 UTF-8 字符边界截断,0表示不设上限); - 规范化、prompt-guard 筛查、加"不可信内容"框(framing),再进入匹配与上下文;
- 每个事件对所有已加载 SOP 的触发器求值,命中即启动运行,并经由
SopAuditLogger持久化启动审计。
六、背压与重投递:Deferred 的两种命运
SOP 引擎有并发准入控制(SOP.toml中的max_concurrent、admission_policy、max_pending_approvals)。当触发器到达但执行槽位或待审批池已满时,投递被标记为Deferred(背压),其恢复方式依赖传输(本版本引擎内没有持久化待触发队列,见 sop/syntax.md 的准入小节):
- AMQP SOP-only 派发(
dispatch = "sop"且durable_ack = true):投递被nack(requeue = true),broker 会在有空位后重投——不会把触发器 ack 掉丢失。这正是DeliveryOutcome::Deferred分支(amqp.rs 第 477-494 行)的行为,由results_need_redelivery(dispatch.rs 第 1196 行)判定:结果集中存在 Deferred 且没有任何Started 时才要求重投。若想延迟重试,可在 broker 侧配置带 TTL 的死信交换器。 - AMQP 组合模式(
sop_and_agent_loop):agent 侧已经消费了这条投递,若再 requeue 会向 agent 循环重复投递同一条消息、双重执行其副作用。因此 SOP 侧溢出时大声记录日志后直接 ACK(不重投),避免 agent 侧被跑两遍。源码中该分支会发出WARN级日志(amqp.rs 第 219-242 行),测试combined_mode_acks_sop_overflow_and_agent_gets_exactly_one_message专门守护"agent 恰好收到一次"这一回归点。
准入策略(admission_policy,snake_case)包括:parallel(默认,无法准入则延迟,永不静默丢弃)、hold(串行化,仅当该 SOP 无运行进行中或停靠时准入)、coalesce(把并发触发器折叠到已在飞的运行上)、drop(显式选择的历史 fire-and-forget)。
恰好一次的语义:message_id 去重
AMQP 投递的message_id被用作每消息幂等键(按 channel alias 加命名空间,amqp:<alias>:<id>),且只有 broker 标记为redelivered的确认重投才会被折叠;全新投递即使复用 message_id 也总是会派发(amqp.rs 第 206-216 行)。因此:
- 发布方应为每条逻辑消息设置唯一且稳定的
message_id以获得恰好一次语义; - 没有或空白 message_id 的投递完全不参与去重(每条都启动自己的运行,绝不猜测性 ACK);
- 复用 message_id(违反约定)时,重投仍可能把不同触发器折叠——这是文档化的至多一次边界。对应测试覆盖了
route_delivery_coalesces_only_a_redelivery_of_the_same_message_id与route_delivery_fresh_deliveries_reusing_a_message_id_both_start两个方向。
七、安全默认:不可信输入处理
AMQP 事件源属于"活的"外部输入,默认安全机制(fan-in 总览安全表):
| 关注点 | 机制 |
|---|---|
| 不可信触发输入 | topic 与 payload 文本在进入模型上下文前截断、规范化、prompt-guard 筛查并加框(framing 始终开启;可隐藏提示文字,但绝不把外部原始文本直接插值进模型上下文) |
| 不安全触发块 | sop.untrusted_input_guard = "block"直接拒绝不安全的不可信事件(BlockedUnsafe);默认warn是审计并放行 |
| 头部受限上下文 | 无 agent 循环时,process_headless_results将ExecuteStep动作记为 pending 而不是静默执行 |
[sop]下与不可信输入相关的字段还包括:untrusted_guard_sensitivity(默认0.7,prompt-guard 筛查灵敏度)、untrusted_frame_warning(默认true,不可信内容框中的警告文字)、untrusted_outbound_redact(默认true,SOP 内容安全消费者的出站脱敏)。
八、审批与观察:checkpoint 停靠后的处理
AMQP 触发的运行一旦到达 checkpoint(kind: checkpoint/requires_confirmation: true)即暂停为WaitingApproval(待审批)。审批与观察有两种途径:
CLI 命令:
zeroclaw sop list # 列出运行(含停靠在审批处的运行) zeroclaw sop approve # 批准停靠的运行Gateway API(out-of-band,通过 gateway API):
GET /admin/sop/pending— 查看待审批运行;POST /admin/sop/approve— 批准;POST /admin/sop/deny— 拒绝。
审批门可以绑定策略(SOP.mdstep 上的- policy: prod),策略定义在[sop.approval.policies.*],包含required_group、quorum(法定人数)、escalation_route等;需要跨渠道通知时还可以配置request_route(如discord.ops:123456789012345678),把审批请求路由到 Discord 等渠道(仅 daemon 路径生效)。
运行持久化:
sop.persist_runs默认true,停靠在 HITL 审批或确定性 checkpoint 的运行会跨 daemon 重启存活(后端sqlite,写入runs.db);如希望引擎为纯内存、非持久化,可显式设false。
九、故障排查速查表
| 症状 | 可能原因 | 修复 |
|---|---|---|
| 没有消息被消费 | exchange 或 routing keys 与发布方不匹配 | 核对exchange与routing_keys和发布方实际发射值一致 |
| TLS 握手失败 | amqps://未配ca_cert,或证书与密钥不匹配 | 提供ca_cert;mTLS 校验client_cert/client_key配对 |
| 投递到达但无 SOP 启动 | dispatch还是agent_loop,或 SOP 触发器不匹配 | 将dispatch设为sop或sop_and_agent_loop;检查触发器的 routing key 模式与condition |
| SOP 启动了但某步未执行 | 无活跃 agent 循环的 headless 触发 | 为ExecuteStep运行 agent 循环,或把运行设计为停靠在审批处 |
进一步资料:AMQP channel 传输侧细节见 channels/amqp.md;fan-in 各事件源统一机制见 fan-in 总览;SOP.toml/SOP.md完整格式、准入控制与条件语法见 SOP 语法参考。
【免费下载链接】zeroclawFast, small, and fully autonomous AI personal assistant infrastructure, any OS, any platform — deploy anywhere, swap anything 🦀项目地址: https://gitcode.com/gh_mirrors/ze/zeroclaw
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考