ZeroClaw SOP Fan-In:用 AMQP 投递触发标准作业流程
2026/9/20 14:59:39 网站建设 项目流程

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.rsAmqpChannel持有dispatch: SopDispatchengine: Option<Arc<Mutex<SopEngine>>>audit: Option<Arc<SopAuditLogger>>三组状态,route_delivery(amqp.rs 第 142 行)根据dispatch的值分别把投递送入 agent 循环(ChannelMessage)或 SOP 引擎(dispatch_untrusted_fan_in)。

安全前置条件:dispatchsopsop_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 行起):

字段默认值作用
enabledfalse是否启用该 channel。运行时只加载enabled = true的 channel;默认关闭,避免粘贴半截配置就意外上线
amqp_urlbroker 地址。明文用amqp://,TLS 用amqps://(如amqps://fedora:@rabbitmq.fedoraproject.org/%2Fpublic_pubsub),属于 secret 字段
exchange消费者队列要绑定的交换器(如amq.topic),必填
routing_keys[]要绑定的 routing key 列表。建议收窄到感兴趣的主题;绑定#会消费整个交换器,几乎永远不是你想要的,必填至少一个
queue队列名。留空则使用服务端生成的临时独占队列;仅在需要跨重连持久投递时才设置稳定名称(推荐 UUID)
ca_certamqps://连接用的 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_acktrue确认模式:true时投递只在消息被可靠移交给 agent 循环后才 ack(至少一次语义,崩溃会重投);false时 broker 派发即确认(至多一次),仅适合无副作用、可丢弃的消费者
dispatchagent_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_certclient_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_certclient_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 运行,步骤为:

  1. 确认sops_dir已配置。SOP 定义从sops_dir下的子目录加载,该字段默认未设置(运行时 SOP 关闭),需要显式开启:相对路径相对于安装根目录(config.toml所在目录)解析,文档化值为shared/sops(即<install>/shared/sops),也可用绝对路径或~前缀路径。
  2. 设置 channel 的dispatch为 SOP 模式sopsop_and_agent_loop)。
  3. 加载 SOP:在<sops_dir>/<name>/下放置SOP.toml(含[[triggers]] type = "amqp")与可选的SOP.md,用zeroclaw sop validate <name>验证。
  4. 发布消息:向配置的 exchange 发布一条 routing key 与触发器匹配的消息。消费者把投递提升为事件(routing key → topic,body → payload)并派发。每个已加载且 routing key 模式匹配、condition(若有)对 body 成立的 SOP 都会启动一个运行。

如果什么都不启动,按顺序排查(也见 fan-in 总览的故障排查表):

  • dispatch是否是 SOP 模式(sopsop_and_agent_loop)而非默认的agent_loop
  • 队列是否正确绑定,使 routing key 真正到达消费者(检查exchangerouting_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_concurrentadmission_policymax_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_idroute_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_resultsExecuteStep动作记为 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_groupquorum(法定人数)、escalation_route等;需要跨渠道通知时还可以配置request_route(如discord.ops:123456789012345678),把审批请求路由到 Discord 等渠道(仅 daemon 路径生效)。

运行持久化:sop.persist_runs默认true,停靠在 HITL 审批或确定性 checkpoint 的运行会跨 daemon 重启存活(后端sqlite,写入runs.db);如希望引擎为纯内存、非持久化,可显式设false

九、故障排查速查表

症状可能原因修复
没有消息被消费exchange 或 routing keys 与发布方不匹配核对exchangerouting_keys和发布方实际发射值一致
TLS 握手失败amqps://未配ca_cert,或证书与密钥不匹配提供ca_cert;mTLS 校验client_cert/client_key配对
投递到达但无 SOP 启动dispatch还是agent_loop,或 SOP 触发器不匹配dispatch设为sopsop_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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询