跳转到内容

规则引擎

Aether 的规则引擎位于 aether-rules 库中并在自动化内部运行。它由三个部分组成:一个解析器,将可视化编辑器的流文档转换为紧凑的执行拓扑;一个调度器,决定何时运行每个规则;以及一个执行器,用于遍历拓扑、评估条件并写入操作点。本页涵盖了发动机的机械结构;有关如何将控制策略表示为规则流(带有有效的充电状态示例),请参阅控制策略

每条规则都以两个平行列的形式保留在 SQLite rules 表中:

  • flow_json — 完整的 Vue Flow 编辑器文档:带有位置和标签、边的节点、视口、元数据。这是可视化编辑器加载以重新填充其画布的内容。
  • nodes_json — 引擎实际运行的紧凑执行拓扑:带有 start_node ID 的 RuleFlow 以及节点 ID 到执行节点的映射。解析器 (extract_rule_flow) 仅保留执行所需的内容,并丢弃位置、标签和边缘样式。

不变:两列始终由一个函数 aether_rules::flow_column_values() 一起生成,该函数解析编辑器文档并返回一个包含两个序列化字符串的 FlowColumns 结构体。该结构使用命名字段而不是元组,因此调用站点在绑定 SQL 参数时无法默默地交换两个值,并且无效流作为一个单元失败 - 没有部分输出。没有代码路径独立序列化任一列。

存在三个生产调用站点,并且 rules 表的任何新写入路径都必须经过相同的函数:

  • libs/aether-rules/src/repository.rs 中的repository::upsert_rule(由规则导入使用)
  • services/automation/src/rule_routes.rs
  • 配置中的 tools/aether/src/core/syncer.rs 处理程序tools/aether/src/core/syncer.rs (aether sync) 中的同步器

一个细微差别:POST /api/rules 创建一个仅元数据存根 — 一个空的 {} 拓扑、一个 NULL 编辑器文档和 enabled = false。流内容总是稍后通过 PUT 到达,PUT 将两列一起派生。以纯压缩形式导入的旧规则保留 NULL flow_json;它们的 nodes_json 仍然来自相同的函数。

为什么这很重要:如果列出现分歧,编辑器将显示一个逻辑,而引擎执行另一个逻辑 - 审核策略的操作员将读到谎言。通过一个生产者汇集每次写入,在结构上不可能出现这种分歧,而不是代码审查问题。

单个调度程序循环 (RuleScheduler::start) 与 tokio::select! 复用两个输入:一个周期性滴答(默认为 100 毫秒,DEFAULT_TICK_MS),以及当连接 PointWatch 事件平面时,一个点更改事件的有界通道。规则在其 trigger_config 列中声明两种触发器类型之一:

  • 间隔 ({"type": "interval", "interval_ms": 1000}) — 规则在自上次执行以来的时间已达到 interval_ms 的任何时间点上到期。没有 trigger_config 的规则默认为 1000 毫秒间隔(或者以 cooldown_ms 作为周期,如果已设置)。
  • OnChange — 规则通过 point_refs 订阅特定测量点 (M) 或操作点 (A),并在订阅值更改超出其死区时触发。

OnChange 规则由两个并行运行的路径提供服务。快速路径是事件驱动的:io 在写入订阅点时发布 PointWatch 事件,调度程序将 (channel, point) 对映射到规则 ID,调度程序立即评估这些规则。事件携带新值,因此触发决策不需要从实时数据库或共享内存中回读。滴答路径是后备:每次滴答,调度程序都会对一批中的所有订阅点进行采样(如果可用,则直接从共享内存中采样)并重新评估每个 OnChange 规则 - 这涵盖了多点规则,并在事件套接字关闭时保持规则触发。

两个死区过滤噪声,与 AND 语义相结合:

  • time_deadband_ms — 规则级频率限制;自上次触发后经过该时长,规则才会再次触发。
  • value_deadband — 绝对值(|new - last| > threshold)或百分比阈值。 如果没有,则有限值之间的任何更改都会计入。

NaN 值(Aether 的“暂时不可用”哨兵)永远不会算作更改;间隙触发一次后的第一个有限值。每次触发后,每点的“最后一个值”都会前进到执行器在执行期间实际读取的值,因此将来的比较将锚定到规则逻辑实际看到的值。

由于规则以有限并行性同时执行(默认情况下一次四个)。独立于触发器,规则可以声明 cooldown_ms;仅在成功执行并执行至少一个操作后才开始冷却,并抑制重新执行,直至其结束。

执行器从起始节点开始沿着每个节点的连线遍历紧凑拓扑:切换节点评估条件分支并选择输出连线,更改节点将值写入一个点,计算和周期增量节点计算派生值,最终节点终止

输入变量通过 SHM 支持的 RuleLiveState 读取,并且读取是严格的:如果变量的数据在本周期不可用,则评估会短路而不是替换默认值 - 丢失的读数绝不能满足像 current < threshold 这样的条件,就好像它为零一样。同样的规则也适用于写入:NaN、无限或超出范围的计算值将被拒绝,操作将记录为失败,永远不会强制为数字。

对操作点的写入采用命令路径:执行器解析实例的模型到通道路由,然后通过自动化的 HTTP 控制端点使用的相同 ActionDispatch 进行调度 — 写入共享内存命令槽(C/A 槽),然后向 io 发出 Unix 域套接字通知(请参阅共享内存)。调度程序在写入之前和之后检查共享内存写入器生成,因此会检测到并丢弃通过 io 重新启动触发的规则,而不是落在过时的槽中。

当目标不可用时(io 重新启动后没有共享内存写入器、缺少 C/A 槽或降级的通知套接字),该操作将被记录为失败并有原因,并且规则的成功标志会反映该情况。命令永远不会排队等待稍后传送:下一个周期会根据当前值重新评估,而不是针对已返回不同状态的设备重放过时的设置点。

每次执行都会生成结果记录:成功标志、执行的操作列表(每个操作都有目标、点、值及其自己的成功标志)、作为访问节点 ID 的执行路径、匹配的条件、变量值的快照,以及每个节点的执行详细信息。调度程序将此记录写入本地持久且可观察的表面:

  • 写入每个规则的日志文件,因此每个规则都有独立的、可 grep 的历史记录;
  • 写入 SQLite rule_history,它为 API/WebSocket 使用方保留结构化执行结果和错误。

规则执行和结果观察都不需要外部

POST /api/scheduler/reload 调用 RuleScheduler::reload_rules,它会从 SQLite 重新读取所有已启用的规则,并在写入锁定下批量替换调度程序的内存中规则集 — 运行执行将根据其开始时的规则完成。当连接 PointWatch 句柄(生产模式)时,重新加载会根据新的规则集重建事件平面作为一个单元:告诉 io 哪些点要发布事件的订阅位图,以及将 (channel, point) 对映射到规则 ID 的调度程序索引。新添加或重新定位的 OnChange 规则立即开始接收事件,无需重新启动服务。

实际上很少需要端点:规则 CRUD 端点(创建、更新、删除、启用、禁用)每个在数据库写入后都会触发相同的重新加载。显式重新加载对于带外写入很重要 - 从配置中批量导入或推送规则文件 - 调度程序将不知道表已更改。