Skip to content

Latest commit

 

History

1 Commit

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

简体中文 | English

le5le-com/rule-flow

Go Reference

Golang实现的规则引擎。
Golang implementation of rule engine.

1. 基础概念

场景联动

通过创建联动规则将智能设备和场景应用联接在一起,从而实现自动协同工作。联动规则主要由触发器、条件、动作等组成,其形式为一种特殊语言,是能够在规则引擎执行一组指令。

规则引擎

通过规则节点和规则链表达联动规则,从而完成预定义的业务需求。其中,规则节点和规则链由专门的语言来表达和定义,以后统一将这个语言成为规则引擎语法。

规则节点

规则链中的最小执行单元,后面也会简称节点。每个节点负责完成一项独立任务,例如“定时触发”、“设备属性判断”或“下发控制指令”。根据其用途,节点可以划分为触发节点、判断节点、执行节点等类型。

每个节点可选定义输入/输出。通过输入/输出关联一组上下游的节点,每个节点自身业务执行完成后,自动触发输出指向的下游节点。触发过程中,如果输出包含数据,则触发的同时把数据传给下游节点。节点的执行通过触发器条件或上游节点触发执行。

规则链

由规则节点 + 输入/输出 + 条件等关联组成的可编排规则。 规则链支持串行(顺序执行)、并行分支(拓扑上并行、执行上按分支顺序展开) 和条件分支(根据结果选择路径) 等编排方式。

2. 基本组成

2.1. 触发器

定时器触发器、倒计时触发器、设备属性上报、设备事件、设备状态(连接/断开/停止运行)、消息、手动触发(临时调试)、自定义(通过定时器或者长链接请求第三方系统API封装的满足条件发送消息的一种触发器,例如天气)

2.2. 条件

关系条件、数学集合条件、自定义条件等 自定义条件指:是一种返回true/false输出的一种特殊条件。用户通过脚本或多个节点编排组合编排,将多个基础条件或复杂逻辑组合而成的逻辑判断。它允许用户实现平台内置条件无法满足的特殊业务规则,最终输出一个布尔值(true/false)。

2.3. 动作

当被触发时,系统执行的具体业务操作或脚本或函数。其中,如果执行具体业务,结果通常是改变物理设备的状态(如开灯)、向用户发送通知(如告警短信)等;如果执行脚本,结果为脚本输出;如果执行函数,则为函数输出。触发通常由自身(触发类型节点)或上游节点触发,例如上游条件节点。动作执行完成后,检测是否有下游节点,有则触发下游节点的执行。

3. 规则引擎语法

3.1. 节点

{
  "id": "001",
  "name": "节点",
  "type": "trigger/condition/switch/action/function/block",
  "config": {},                        // 节点类型相关配置,见 3.1「各节点类型 config 字段速查」
  "inputs": [                          // 输入参数(槽)声明。可为空数组=本节点无输入参数
    // key=参数标识(寻址与脚本变量名都用它),label=显示文案,value=缺省值
    { "key": "temp", "label": "温度", "value": 26 }
  ],
  "inputMode": null,                   // null(默认,等价 "pair") / "any" / "continuous",见 3.6
  "payload": {                         // 可选:本节点对流作用域数据的写入,见 3.6
    "set": [                           // 有序数组,逐行写入,后行可覆盖前行
      { "keyName": "a.b", "value": 1 },                      // 字面量
      { "keyName": "a.c", "value": "value", "isVar": true }  // 变量(含点路径)
    ],
    "remove": ["tmp"]                  // 在全部 set 之后,逐项按点路径删除
  },
  "outputs": {
    // 出口的每一项是一条连线:id=下游节点id,key=落到下游的哪个输入参数(可选)
    "ids": [{ "id": "下游节点id", "key": "参数key" }],  // 正常出口
    "elseIds": [],                     // 条件为假(仅 condition)
    "errorIds": [],                    // 本节点执行失败(异常出口),见 3.7
    "value": "子流节点id"               // 仅含 subflow 的节点:取哪个子流节点的 msg 作输出,见 3.6
  }
}

连线与数据路由

节点用 outputs.ids 声明触发关系:上游执行完,触发 ids 里的每个下游。每一项就是一条连线:

"ids": [
  { "id": "a1", "key": "temp" }, // 触发 a1,数据落到 a1 的 temp 参数
  { "id": "a2" }                  // 触发 a2,数据落到 a2 的第一个输入参数
]
  • inputs 是本节点的输入参数(形参)声明,与上游是谁无关。
    • key 为参数标识,也是参数唯一的引用方式:连线的 key、脚本里的变量名、payload.set 的变量名、动作回调 inputs map 的键、子流引用宿主槽,全都按它。硬约束只有两条:节点内唯一、不含点号(含点号的 key 在 payload.set 的 isVar 表达式里永远取不到值,因为表达式在第一个点号处切分成「变量名 + 点路径」,而且是静默取 null)。另有一条软约束:key 会作 goja 全局变量名注入,非合法 JS 标识符(如纯数字 "0")注入不报错,但脚本里写不了裸名字,只能用 inputs[下标] 或 this["0"] ——用 script 节点时值得挑个标识符,纯按名取值的场合(动作 inputs / bind / 连线匹配 / payload.set)任何合法 JSON 键都可以。
    • label 为显示文案,只给编辑器看、引擎不解释——可随意改名而不影响任何引用。
    • value 为缺省值(无上游数据流入、或流入 null 时取它)。
  • 落哪个槽由连线的 key 决定:匹配下游 inputs[].key,只按名字匹配,不是下标。写成数字("key": 0)在加载时报错——它是最常见的误用,且下标一旦与某个参数 key 同形就会意外命中,把上游 value 写进不该写的槽。
  • key 省略 → 落 inputs[0](第一个参数)。
  • key 指向下游不存在的参数、或下游 inputs 为空 → 忽略数据、仍触发下游执行。连线的首要含义是触发,数据传递是次要的。这也是"只触发、不传值"的写法:给一个下游没有的参数名即可。
  • 上游是 condition / switch → 一律不落槽,key 被忽略(见 3.3)。这两类节点只选路线、不产自己的值。
  • 多个上游连到同一下游的不同参数(扇入),互不撞车;连到同一参数则后到的覆盖。
  • elseIds、errorIds、switch 的 cases[].ids 与 defaultIds 均按同一格式。
  • outputs.ids 等出口不能指向 trigger 节点(触发器只能由消息或自身定时激活),加载时报错。

为什么槽是形参而不是"上游 id":节点在编辑器里可任意复制、组合节点可被多处复用,此时各实例的输入参数签名相同、上游各不相同。把连线信息集中在上游 outputs 里,画线只改一处,节点定义可原样复制。

  • 触发器节点
{
  "id": "001",
  "name": "每分钟触发一次",
  "type": "trigger",
  "config": { "cron": "* * * * *" },
  "outputs": { "ids": [{ "id": "c1" }] }
}
  • 条件节点
{
  "id": "c1",
  "name": "温度超限",
  "type": "condition",
  "config": { "relation": ">" },
  "inputs": [
    { "key": "temp" },          // 无 value:等上游流入
    { "key": "limit", "value": 30 } // 有 value:缺省阈值
  ],
  "inputMode": null,
  "outputs": {
    "ids": [{ "id": "a1" }],                // 条件为真
    "elseIds": [{ "id": "a2" }]             // 条件为假(可选,省略则为假断流)
  }
}

条件为真触发 outputs.ids、为假触发 outputs.elseIds。condition 只决定走哪个出口,不修改 msg(见 3.6);它不产值,所以出口不给下游落参数槽,连线的 key 会被忽略。

  • 多路分支节点
{
  "id": "sw", "name": "按档位分流", "type": "switch",
  "config": {
    "mode": "first",                        // first(默认,命中即停) | all(触发所有命中的 case)
    "cases": [
      { "relation": "=",  "value": "low",   "ids": [{ "id": "a1" }] },
      { "relation": ">",  "value": 10,      "ids": [{ "id": "a2" }] },
      { "relation": "[)", "range": "[0,5)", "ids": [{ "id": "a3" }] },
      { "relation": "[]", "set": "[1,2,3]", "ids": [{ "id": "a4" }] }
    ],
    "defaultIds": [{ "id": "a5" }]          // 无命中走默认(可选)
  },
  "inputs": [{ "key": "level" }]
}
  • 动作节点
{
  "id": "a1",
  "name": "通知",
  "type": "action",
  "config": { "kind": "http", "http": "url地址", "method": "POST", "body": {} },
  "inputs": [{ "key": "text", "label": "告警内容" }],
  "outputs": { "ids": [], "errorIds": [{ "id": "a9" }] }
}
  • 组合节点
{
  "id": "blk",
  "name": "组合节点完成条件判断",
  "type": "block",
  "config": {
    "flowId": "子流定义在宿主库中的 id",    // 引擎不使用;宿主库中只存它,加载时展开成 subflow
    "subflow": [{}, {}, {}]                 // 引擎只认 subflow(宿主加载期填入)
  },
  "inputs": [
    { "key": "temp", "label": "温度", "value": 26 },
    { "key": "humi", "label": "湿度", "value": 60 },
    { "key": "load", "label": "负载", "value": 0.3 }
  ],
  "outputs": {
    "ids": [{ "id": "a1" }],
    "errorIds": [],
    "value": "s3"                           // 取 subflow 中 s3 节点的 msg 作为本节点输出
  }
}

组合节点被多处引用时,各处 inputs 参数签名(各参数的 key)相同、id/outputs 各不相同,内部 subflow 可原样复制。 block 的输出由 outputs.value 指定的子流节点决定,详见 3.6「组合节点的数据边界」。

各节点类型 config 字段速查

节点行为完全由 type + config 决定。下表列出每种类型识别的 config 字段(未列字段引擎忽略)。

trigger(触发器) —— 按 config 里出现的字段决定触发方式,优先级 cron > interval > countdown > source > message:

字段 说明
cron cron 表达式定时触发(5 字段标准 / 6 字段带秒)。自驱动,仅执行点/leader 跑。
interval 周期定时(秒)。自驱动。
countdown 倒计时(秒),到点触发一次。自驱动。
source + url/broker/topic/username/password + script + interval 第三方订阅:source 取 http/mqtt/sse/websocket;http/sse/websocket 用 url,mqtt 用 broker+topic(+username/password);script 返回 true 才触发下游;http 轮询用 interval(秒)。自驱动。
message 监听的消息名(device-report/device-connect/click/自定义…)。靠 Emit 驱动,非自驱动。
deviceId 可选。与 message 配合:只接该设备的消息(读消息体的 deviceId 过滤;多设备场景区分来源,见 3.6 inputMode)。

condition(条件) —— relation 选运算符:

字段 说明
relation = != > >= < <=(取前两个输入比较);[)/![) 配 range;[]/![] 配 set;AND/OR(各输入真值);自定义(配 script 或 subflow,返回 bool)。
range 区间表达式,如 [0,100)。relation 为 [)/![) 时用。
set 集合表达式,如 [1,20,30..50,65]。relation 为 []/![] 时用。
script goja 脚本,relation=自定义 时用;return 布尔。脚本变量见 3.5「脚本如何读取数据」。
subflow 子流(节点数组),relation=自定义 时用;取 outputs.value 指定子流节点的 msg.value 作布尔。

真值判断:AND/OR 及自定义脚本的返回值按 JS 真值语义判断。 假值只有 false、0、""、null、字段缺失;空数组 [] 与空对象 {} 为真。

假分支:condition 节点的 outputs.elseIds(与 ids 平级)为"条件为假"时触发的下游; 真走 ids、假走 elseIds。省略 elseIds 时为假即断流(旧行为不变)。 脚本报错等求值异常走 errorIds(异常出口,见 3.7),不会被当成"条件为假"走假分支。

switch(多路分支) —— 按输入值匹配多个 case,各 case 有独立下游(编辑器里每个 case = 一个输出口):

字段 说明
cases 分支数组,每项 {relation,value/range/set,ids}:relation 复用 condition 运算符(=默认/!=/>/>=/</<= 配 value;[)/![) 配 range;[]/![] 配 set),命中则触发该 case 的 ids。
mode first(默认,命中第一个即停) / all(触发所有命中的 case)。
defaultIds 全不命中时触发的下游(可选)。

action(动作) —— kind 选动作类型(缺省按字段推断):

字段 说明
kind 动作类型:见 3.4 全清单(device-*/http/mqtt/db/email/sms/voice/alert/message/emit…)。
emit 专用 message + data kind=emit 时向本场景投递该消息(内联动,不回调宿主)。message 为消息名、data 为静态字段,再并入本节点输入参数,见 3.4「emit 的消息结构」。
subflow 自定义动作:跑子流(不回调宿主)。
各动作业务字段 由宿主 ActionHandler 解释(如 http 的 http/method/body、mqtt 的 topic、device-set-property 的 deviceId/properties/property 等),引擎原样透传 config / inputs / value / payload 给回调(见 3.6)。

动作出口:成功走 outputs.ids、失败走 outputs.errorIds(见 3.7)。 handler 有返回值则写入 msg.value;无返回值则 msg 原样流向下游(副作用节点不掐断数据流)。详见 3.4。

function(函数) —— fn/script/subflow 三选一,优先级 fn > subflow > script:

字段 说明
fn 内置函数名(见 3.5:uuid7/avg/add/get/merge… 及阀门 debounce/throttle)。
fn=debounce + delay/maxWait 防抖阀门(毫秒):静默 delay 后用最后一次输入触发下游;maxWait 强制上限。
fn=throttle + interval 节流阀门(毫秒):每 interval 最多放行一次。
subflow 自定义函数:跑子流算输出(输出由 outputs.value 指定的子流节点决定)。
script goja 脚本自定义函数;return 值写入 msg.value。

block(组合节点) —— subflow(节点数组)作为子流运行。子流内节点按参数 key 直接引用宿主的输入参数; 输出由 outputs.value(子流节点 id)指定,详见 3.6「组合节点的数据边界」。 另有 flowId:子流定义在宿主库中的 id(引擎不使用,宿主用它反查"哪些流引用了这个子流"); flowId 与 subflow 并存时引擎只认 subflow。

宿主侧建议区分两种形态:库中只存 flowId(单一数据源,改子流定义即改所有引用它的实例, 无需级联重写各流的 nodes),加载时按 flowId 递归展开成 subflow 再交给引擎—— 引擎只见展开后那份。同一子流在一个流中出现多次时各实例各展开一份, subflowOf 按宿主节点 id 加 scope 前缀隔离各实例的 store。 le5le 场景联动即如此:组合节点必须带 flowId(保存时校验),采集端加载时展开。

inputMode(多输入节点) —— 节点顶层字段(与 inputs 平级),控制多输入何时执行: 缺省(=pair)=所有参数到齐才执行一次、执行后清空重等;any=任一上游到达即执行、 缺失的参数取缓存最新值或缺省值;continuous=到齐后首次执行,此后任一更新即用最新值重算。详见 3.6。

各节点类型 JSON 配置示例

每种类型一个可直接使用的示例(与库测试/示例场景一致)。

  • 触发器 · cron 定时器
{
  "id": "t1",
  "type": "trigger",
  "config": { "cron": "*/1 * * * * *" },
  "outputs": { "ids": [{ "id": "a1" }] }
}
  • 触发器 · 倒计时
{
  "id": "t1",
  "type": "trigger",
  "config": { "countdown": 60 },
  "outputs": { "ids": [{ "id": "a1" }] }
}
  • 触发器 · 设备属性上报(按设备过滤,并把整包上报数据写入流作用域)
{
  "id": "t1",
  "type": "trigger",
  "config": { "message": "device-report", "deviceId": "dev-1" },
  "payload": { "set": [{ "keyName": "report", "value": "value", "isVar": true }] },
  "outputs": { "ids": [{ "id": "c1", "key": "temp" }] }
}
  • 触发器 · 第三方 HTTP 轮询
{
  "id": "t1",
  "type": "trigger",
  "config": {
    "source": "http",
    "url": "http://x/api",
    "interval": 5,
    "script": "return body.temp > 30"
  },
  "outputs": { "ids": [{ "id": "a1" }], "errorIds": [{ "id": "a9" }] }
}
  • 条件 · 关系比较(温度 > 阈值;真走 a1、假走 a2)
{
  "id": "c1",
  "type": "condition",
  "config": { "relation": ">" },
  "inputs": [
    { "key": "temp" },
    { "key": "limit", "value": 30 }
  ],
  "outputs": { "ids": [{ "id": "a1" }], "elseIds": [{ "id": "a2" }] }
}
  • 条件 · 区间
{
  "id": "c1",
  "type": "condition",
  "config": { "relation": "[)", "range": "[0,100)" },
  "inputs": [{ "key": "v" }],
  "outputs": { "ids": [{ "id": "a1" }], "elseIds": [{ "id": "a2" }] }
}
  • 条件 · 集合
{
  "id": "c1",
  "type": "condition",
  "config": { "relation": "[]", "set": "[10,20,30..50]" },
  "inputs": [{ "key": "v" }],
  "outputs": { "ids": [{ "id": "a1" }], "elseIds": [{ "id": "a2" }] }
}
  • 条件 · 自定义脚本(按参数 key 读值,也可用 inputs[0] 或 payload)
{
  "id": "c1",
  "type": "condition",
  "config": { "relation": "自定义", "script": "return temp > 32" },
  "inputs": [{ "key": "temp" }],
  "outputs": { "ids": [{ "id": "a1" }], "elseIds": [{ "id": "a2" }] }
}
  • switch · 多路分支(按档位分流,命中即停,无命中走默认)
{
  "id": "sw",
  "type": "switch",
  "config": {
    "mode": "first",
    "cases": [
      { "relation": "=", "value": "low", "ids": [{ "id": "a1" }] },
      { "relation": ">", "value": 10, "ids": [{ "id": "a2" }] },
      { "relation": "[)", "range": "[0,5)", "ids": [{ "id": "a3" }] }
    ],
    "defaultIds": [{ "id": "a4" }]
  },
  "inputs": [{ "key": "level" }]
}
  • 条件 · 自定义 subflow(取 s1 的 msg.value 作布尔)
{
  "id": "c1",
  "type": "condition",
  "config": {
    "relation": "自定义",
    "subflow": [
      {
        "id": "s1",
        "type": "function",
        "config": { "script": "return v > 32" }
      }
    ]
  },
  "inputs": [{ "key": "v" }],
  "outputs": {
    "ids": [{ "id": "a1" }],
    "elseIds": [{ "id": "a2" }],
    "value": "s1"
  }
}

子流内的 s1 未声明 inputs,脚本里的 v 直接引用宿主 c1 中 key 为 v 的参数(见 3.6「组合节点的数据边界」)。

  • 动作 · 设置设备属性(固定值)
{
  "id": "a1",
  "type": "action",
  "config": {
    "kind": "device-set-property",
    "deviceId": "dev-2",
    "properties": { "t": 30 }
  }
}
  • 动作 · 设置设备属性(上游动态值写入单个属性)
{
  "id": "a1",
  "type": "action",
  "config": {
    "kind": "device-update-property",
    "deviceId": "dev-1",
    "property": "fahrenheit"
  },
  "inputs": [{ "key": "v", "label": "值" }]
}
  • 动作 · HTTP 通知
{
  "id": "a1",
  "type": "action",
  "config": {
    "kind": "http",
    "http": "http://x/notify",
    "method": "POST",
    "body": { "k": 1 }
  }
}
  • 动作 · 告警
{
  "id": "a1",
  "type": "action",
  "config": {
    "kind": "alert",
    "level": "warning",
    "desc": "温度超限",
    "deviceId": "dev-1"
  }
}
  • 动作 · emit(场景内投递消息)
{
  "id": "a1",
  "type": "action",
  "config": {
    "kind": "emit",
    "message": "relay-on",
    "data": { "deviceId": "dev-1" }
  },
  "outputs": { "ids": [{ "id": "a2" }] }
}

投递出去的消息为 {"name":"relay-on","deviceId":"dev-1"},见 3.4「emit 的消息结构」。 a1 自身继续向 a2 流转(emit 不是终止节点)。

  • 动作 · 自定义 subflow
{
  "id": "a1",
  "type": "action",
  "config": {
    "subflow": [
      {
        "id": "s1",
        "type": "action",
        "config": { "kind": "mqtt", "topic": "cmd" },
        "outputs": { "ids": [{ "id": "s2" }] }
      },
      {
        "id": "s2",
        "type": "action",
        "config": { "kind": "message", "text": "完成" }
      }
    ]
  }
}
  • 函数 · 内置函数(两路温度均值)
{
  "id": "fn",
  "type": "function",
  "config": { "fn": "avg" },
  "inputs": [
    { "key": "a" },
    { "key": "b" }
  ],
  "inputMode": "continuous",
  "outputs": { "ids": [{ "id": "a1" }] }
}
  • 函数 · 自定义脚本(摄氏转华氏)
{
  "id": "fn",
  "type": "function",
  "config": { "script": "return celsius * 9 / 5 + 32" },
  "inputs": [{ "key": "celsius" }],
  "outputs": { "ids": [{ "id": "a1" }] }
}
  • 函数 · 无输入参数(生成 uuid)
{
  "id": "fn",
  "type": "function",
  "config": { "fn": "uuid7" },
  "inputs": [],
  "outputs": { "ids": [{ "id": "a1" }] }
}
  • 函数 · 自定义 subflow
{
  "id": "fn",
  "type": "function",
  "config": {
    "subflow": [
      {
        "id": "s1",
        "type": "function",
        "config": { "fn": "add" },
        "inputs": [{ "key": "a" }, { "key": "b", "value": 100 }]
      }
    ]
  },
  "inputs": [{ "key": "a" }],
  "outputs": { "ids": [{ "id": "a1" }], "value": "s1" }
}
  • 函数 · 防抖阀门
{
  "id": "fn",
  "type": "function",
  "config": { "fn": "debounce", "delay": 500, "maxWait": 2000 },
  "inputs": [{ "key": "v" }],
  "outputs": { "ids": [{ "id": "a1" }] }
}
  • 函数 · 节流阀门
{
  "id": "fn",
  "type": "function",
  "config": { "fn": "throttle", "interval": 1000 },
  "inputs": [{ "key": "v" }],
  "outputs": { "ids": [{ "id": "a1" }] }
}
  • 组合节点 block
{
  "id": "blk",
  "type": "block",
  "config": {
    "flowId": "子流定义在宿主库中的 id",
    "subflow": [
      {
        "id": "b1",
        "type": "function",
        "config": { "script": "return temp + 100" },
        "outputs": { "ids": [{ "id": "b2" }] }
      },
      { "id": "b2", "type": "action", "config": { "kind": "message" } }
    ]
  },
  "inputs": [{ "key": "temp", "label": "温度" }],
  "outputs": { "ids": [{ "id": "a1" }], "value": "b1" }
}
  • inputMode · continuous 多输入(两设备温度均值 > 50 才触发;两设备各自上报都重算一次)
{
  "id": "cond",
  "type": "condition",
  "config": { "relation": "自定义", "script": "return (t1.t + t2.t)/2 > 50" },
  "inputs": [
    { "key": "t1" },
    { "key": "t2" }
  ],
  "inputMode": "continuous",
  "outputs": { "ids": [{ "id": "alert" }], "elseIds": [{ "id": "ok" }] }
}

两个上游分别连到 t1、t2 两个参数:

{ "id": "trig1", "type": "trigger", "config": { "message": "device-report", "deviceId": "dev-1" },
  "outputs": { "ids": [{ "id": "cond", "key": "t1" }] } }
{ "id": "trig2", "type": "trigger", "config": { "message": "device-report", "deviceId": "dev-2" },
  "outputs": { "ids": [{ "id": "cond", "key": "t2" }] } }
  • payload · 写入与读取(上游写入流作用域,下游跨节点读取)
{
  "id": "fn",
  "type": "function",
  "config": { "script": "return celsius * 9 / 5 + 32" },
  "inputs": [{ "key": "celsius" }],
  "payload": { "set": [{ "keyName": "temp.fahrenheit", "value": "value", "isVar": true }] },
  "outputs": { "ids": [{ "id": "a1" }] }
}
{
  "id": "a1",
  "type": "action",
  "config": { "kind": "alert", "level": "warning" },
  "inputs": [],
  "payload": { "set": [{ "keyName": "alerted", "value": true }], "remove": ["tmp"] }
}

fn 把计算结果(变量 value)写入 payload.temp.fahrenheit;a1 无输入参数、不接收 value,但仍能读到 payload。

3.2. 触发器

定时器:cron 表达式(config.cron)。支持 5 字段标准 cron(分 时 日 月 周,最小粒度 1 分钟)与 6 字段带秒(秒 分 时 日 月 周)。例:{"cron":"0 8 * * *"} 每天 8:00;{"cron":"*/1 * * * * *"} 每秒。
倒计时:config.countdown(秒),到点触发一次。
设备属性上报:监听消息“device-report”
设备事件上报:监听消息“device-event”
设备连接:监听消息“device-connect”
设备断开:监听消息“device-disconnect”
设备停用:监听消息“device-disable”
消息:自定义消息名(config.message)
单击:监听消息“click”。节点单击被单击时,发送此消息
第三方应用:订阅第三方系统(http、mqtt、sse、websocket + JS脚本),脚本return true/false , 触发下游条件或动作节点即可

触发器产出的 msg:每次触发新建一条 msg,payload 为 {}、value 为触发数据(见 3.6):

触发方式 msg.value
cron / interval null
countdown null
消息类 整个消息对象(含 name、业务字段)
第三方订阅 拉到/收到的原始数据(http 为响应体等)

需要把触发数据带进流作用域给全链路读,在触发器上配 payload.set(见 3.1 示例)。

3.3. 条件

= :等于
!= :不等于
> :大于
> = :大于等于
< :小于
<= :小于等于
[) :包含,与数学上的前闭后开相同,例如:[0, 100) 指0100,包含0、不含100; [0,100]指 0100,包含0和100
![) : 不包含 ,上面包含的反集 []: 属于。例如:[1,20,30..50,65],其中30..50指30~50之间的值(包含30和50)、 [1,20,'aaa']
![]:不属于
AND :逻辑与。各输入按 JS 真值语义判断
OR :逻辑或
switch :switch case 多路分支。它是独立节点类型(type:"switch")而非 relation,见 3.1 示例
自定义 :执行一段 js 脚本(config.script)或子流(config.subflow),返回 true/false

真值语义:假值只有 false、0、""、null、字段缺失;空数组 [] 与空对象 {} 为真(同 JS)。

分支出口:条件为真触发 outputs.ids;为假触发 outputs.elseIds(假分支,可选;省略则为假断流)。 需要"按取值走多条不同路径"时用 switch 节点。

condition 与 switch 不修改 msg:判断结果已由"走哪个出口"表达,payload 与 value 原样流向下游。 所以「温度超限 → 告警」里,告警节点仍能拿到被判断的那个温度值(动作回调的 value、脚本里的 value)。详见 3.6。

但它们的出口不给下游落参数槽:连线上的 key 被忽略,value 不写入下游的输入参数、也不入输入缓存。 因为 condition/switch 只选路线、不产自己的值——流下去的 value 是上游某个节点的值, 多输入条件下还取决于"哪条支线最后到达"(即上游 outputs.ids 的数组次序)。 若让它落槽,下游一个声明了缺省值的参数会被这个值静默顶掉: AND → randInt(min:15, max:20) 的连线若指向 min,min 会被条件流下的值顶掉, 节点名写着 15-20、实际跑的是别的区间,不报错也不断流。

要把值传给下游有两条正路:payload.set(isVar 可读 value/inputs/参数 key,见 3.5), 或从产值的那个节点直接再拉一条连线到下游参数,条件只作闸门。

副作用:某个参数只由条件出口喂、且没有 value 缺省值时,它与"无入边"等价—— 缺省的 pair 就绪门永不放行。这种情况在加载时经 ErrorHandler 报"永不就绪"告警(见 3.7), 不会静默死掉。

求值异常走 errorIds,不会被当成"条件为假"走 elseIds(见 3.7)。

3.4. 动作

动作通过全局回调 ActionHandler 监听动作类型与数据参数(config/inputs/payload/value)执行副作用;emit 由引擎内部完成、不回调宿主。

• device-connect :连接/启动设备
• device-disconnect :设备断开连接
• device-disable :设备停用
• device-update-property :更新平台内部设备属性值(仅写平台,不下发远程设备)
• device-set-property :更新平台内部设备属性值,并发送 set-property 指令更新远程设备属性
• device-exec-function :调用设备服务
• http:调用一个HTTP(S)请求
• mqtt:发送mqtt消息
• db:发送到数据库
• email:发送邮件,邮件内容支持字符串模板参数替换
• sms:发送短信
• voice:发送电话语音
• alert:发送告警
• emit:向本场景投递一条消息,触发监听该消息的触发器(场景内联动,引擎内部完成、不回调宿主)
• message:发送平台消息(按 config.topic 路由到宿主的具体业务,如生成工单、摄像头录像/报警;topic 清单由宿主平台维护)
• 自定义:config 含 subflow 时执行一个子流(引擎内部完成、不回调宿主);否则作为其它自定义动作类型回调宿主 ActionHandler

动作出口:成功走 outputs.ids、失败走 outputs.errorIds(异常出口,见 3.7)。

  • handler 有返回值 → 写入 msg.value 流向下游;无返回值 → msg 原样流下。 动作常是链路中间一环(收到属性 → 下发设备 → 记日志),副作用节点不掐断数据流。
  • emit 投递消息后同样继续向下游流转(它是"额外投递一条消息",不是终止节点)。
{
  "id": "a1",
  "type": "action",
  "config": { "kind": "http", "http": "http://x/notify" },
  "outputs": {
    "ids": [{ "id": "成功后续id" }],
    "errorIds": [{ "id": "失败告警id" }]
  }
}

宿主 handler 建议返回业务本体(不要再包一层 {ok,data,error},成功/失败已由出口表达): http → {status, body};db → 影响行数;device-exec-function → 设备函数返回值; device-set-property/update-property → 实际下发的属性 map;通知类(sms/voice/email)→ nil(msg 原样流下)。

emit 的消息结构

emit 投递的消息是一个扁平对象,name 为消息名,其余为业务字段:

{ "name": "relay-on", "deviceId": "dev-1", "temp": 36.5 }

字段按顺序合并,后者覆盖前者:

  1. name = config.message
  2. config.data 的各字段(静态内容)
  3. 本节点各输入参数(按参数 key 索引,key 为空的参数跳过)——动态内容从这里带出

监听该消息名的 trigger 被激活,整个对象作为它的 msg.value;trigger 的 config.deviceId 读 value.deviceId 做过滤。 需要把内容带进流作用域时,在 trigger 上配 payload.set。

被 emit 激活的链路是一次新的传播:payload 从 {} 开始,不继承发送方的流作用域数据。 这也是跨分支传数据的推荐方式——把要传的内容放进 emit 消息体。

3.5. 函数

内置函数:

  • 随机数
    uuid7():生成一个v7版本的UUID
    rand(len):随机生成指定长度的字符串
    randHex(len):随机生成指定长度的16进制数字格式的字符串
    randInt(min,max):随机生成指定范围的整数
    randFloat(min,max):随机生成指定范围的浮点数

  • 聚合统计 avg(num, ...) :平均值
    sum(num, ...) :求和
    max(num, ...) :最大值
    min(num, ...) :最小值
    stddev(num, ...) :标准差

  • 数学运算
    add/sub/mul/div(a,b) :加减乘除
    abs(x) :绝对值
    ceil/floor/round(x) :取整
    sqrt/pow(x,y) :根/幂运算
    mod(x,y) :取模

  • 字符串 len(s)
    contains(s, sub)
    startsWith(s, prefix)
    endsWith(s, suffix)
    matches(s, regex)
    replace(s, old, new)
    split(s, sep)
    trim(s)
    substring(s, start, end) concat(a,b,c...):字符串拼接
    upper() / lower():大小写转换

  • 时间函数 timestamp() 获取当前 Unix 时间戳
    timestamp_ms() 毫秒级时间戳
    formatTime(t, fmt)
    parseTime(s, fmt)
    before(t)
    after(t)

  • 数据操作 merge(obj1, obj2, ...):浅合并多个对象为新对象,后者覆盖前者
    get(obj, path, default):按 "a.b.c" 路径取嵌套值,缺失返回 default;数字段可作数组下标(见 3.6「点路径的规则」)
    pick(obj, key1, key2, ...):从对象选取指定键构成新对象
    debounce(delay, maxWait):防抖阀门——连续触发时只在停止 delay 毫秒后用最后一次输入触发下游;配 maxWait 则最长等待毫秒数强制触发一次
    throttle(interval):节流阀门——interval 毫秒内最多放行一次到下游

    debounce/throttle 是控制下游触发节奏的阀门,不产出值;作为 function 节点的 fn 使用,如 {"type":"function","config":{"fn":"debounce","delay":500}}。

  • 自定义函数
    接收一组输入,通过 js 脚本(config.script) 或 子流(config.subflow) 计算得出一个输出结果,用于作为下游节点的输入

脚本如何读取数据

自定义条件 / 自定义函数的 js 脚本里可用以下变量:

变量 含义
参数 key 输入参数按其 key 注入为同名变量,如 temp、limit
inputs 保留数组变量,inputs[0]、inputs[1]… 按声明顺序对应各参数
value 本次触发本节点的那条 msg 的 value(多输入时为最后到达的那条)
payload 流作用域数据,可读可写(payload.x = 1 生效,见 3.6)
error 上游失败信息,非异常分支为 ""(见 3.7)
errorId 失败节点 id,非异常分支为 ""

参数 key 是首选读法。key 直接作 JS 变量名,所以声明 key 时就要保证它是合法标识符; key 为空的参数不注入变量,只能用 inputs[下标] 读。 inputs、value、payload、error、errorId 为保留变量名,参数 key 与之冲突时以保留变量为准。

脚本 return 的值写入 msg.value 传给下游。

3.6. 数据流转

节点之间流转的数据统一封装为 msg:

{
  "payload": {},      // 流作用域:一次触发内贯穿全链路,可写
  "value": null,      // 传递作用域:落到下游输入参数的临时值
  "error": "",        // 仅异常出口非空:错误信息
  "errorId": "",      // 仅异常出口非空:失败节点 id
  "meta": {}          // 仅 debug 模式存在:上游身份/来源等调试信息
}

作用域

流作用域(payload) —— 一次触发内贯穿全链路的业务数据。初始为 {},需要时由节点写入 (见下方「payload 的写入」)。任何节点都能读,与输入参数无关。

分支之间写隔离(Copy-On-Write):A→(B,C) 时 B 对 payload 的修改不会被 C 看到。 同一分支内连续传递不产生拷贝。需要跨分支传数据时,用 emit 投递场景消息(见 3.4)。

实现上:一条 msg 被发给多个下游时标记为共享;共享状态下首次写入时深拷贝一次 payload(并清除共享标记),之后本分支的连续写入不再拷贝。所以"分支多"不等于"拷贝多", 只写才拷贝,且每分支至多拷一次。

传递作用域(value) —— 上游给下游的临时值,即上游函数/动作/组合等产值节点的计算结果。 落到下游哪个输入参数,由连线的 key 决定(见 3.1「连线与数据路由」)。 condition / switch 不产值,value 照常随 msg 流下(动作回调的 value、脚本里的 value 都能读到), 但不落下游的参数槽(见 3.3)。

场景作用域:引擎内部用于 inputMode 的输入缓存等状态,按场景隔离、场景停止即清空, 不在规则语法中暴露。

执行模型

  • 触发层并发:定时器、倒计时、Emit、第三方订阅、阀门延迟触发各自独立, 同一场景可同时有多次触发在跑。
  • 传播层顺序:单次触发内按深度优先展开——A→(B,C) 时 B 整条支线执行完才开始 C。 "并行分支"指拓扑上的并行,不是并发执行。
  • 环检测按传播路径:同一节点不能在同一条传播路径上出现两次。 菱形结构(A→B、A→C、B→D、C→D)中 D 会被触发两次——想合并成一次执行就让 D 用缺省的 pair(两条支线各喂一个参数,凑齐后只执行一次,然后清空重等)。

节点对 msg 的影响

节点类型 对 msg 的影响
trigger 新建 msg:payload 为 {}、value 为触发数据(消息类为整个消息对象)
condition / switch 不修改,只决定走哪个出口
function 返回值写入 value
action handler 有返回值则写入 value,无返回值则 msg 原样流下
block 取 outputs.value 指定子流节点的 msg(见下方「组合节点的数据边界」)
阀门(debounce/throttle) 不产值:缓存整条 msg,放行时原样发给下游(被丢弃的那几次连 payload 一起丢)
全部 成功执行后清空 error / errorId;再按 payload.set / payload.remove 写入

阀门放行时是一次新的传播(脱离原触发的调用栈,在自己的定时器 goroutine 里), 环检测路径从阀门重新计起。

payload 的写入

节点上的 payload 字段(与 inputs 平级),在节点执行之后生效:

{
  "payload": {
    "set": [                                                      // 有序,逐行执行
      { "keyName": "alarm", "value": true },                       // ① 字面量
      { "keyName": "cfg", "value": { "unit": "℃" } },              // ② 字面量可为任意 JSON
      { "keyName": "temp.f", "value": "value", "isVar": true },    // ③ 变量:本节点结果
      { "keyName": "dev.id", "value": "value.deviceId", "isVar": true }, // ④ 变量+点路径
      { "keyName": "peak", "value": "value.records.1.t", "isVar": true } // ⑤ 穿过数组
    ],
    "remove": ["tmp", "cfg.src"]                                   // ⑥ 全部 set 之后
  }
}

set 的每一行是一次写入:keyName 是写入 payload 的点路径,value 是要写的内容, isVar 表示 value 是变量表达式而不是字面量(省略/false = 字面量)。 同一次 set 内后行覆盖前行,且后行能读到前行的写入结果(payload 变量是实时的)。 remove 在所有 set 之后执行,路径不存在时静默忽略。脚本内 payload.x = 1 直接赋值同样生效。

变量可以是什么

isVar: true 时,value 的第一段是变量名,其余为点路径。可用的变量名与脚本一致(见 3.5):

变量名 含义
参数 key 本节点声明了 key 的输入参数,如 "limit"、"report.site"
inputs 参数数组,按声明顺序,如 "inputs.0.records.0.t"
value 求值那一刻的 msg.value——见下方时机说明
payload 流作用域数据本身,可自引用,如 "payload.report.site"
error 本节点的失败信息,非异常分支为 ""(见 3.7)
errorId 失败节点 id,非异常分支为 ""

value 的时机:统一规则是"value 恒指求值那一刻的 msg.value"。脚本在节点执行前求值, 读到的是流入值;payload.set 在节点执行后求值,读到的是本节点的结果。 condition / switch / 阀门不产值,其 payload.set 里的 value 仍是流入值。

第一段必须是上表之一或本节点的参数 key,否则加载时报错(见 3.7);后续路径段不校验,取不到即为 null。

点路径的规则

  • 数字段的含义由当前值决定:当前值是数组则为下标(records.1.t),是对象则为键名 "1"。
  • 读不到就是 null:路径中间断了、下标越界、类型不匹配,一律静默取 null,不报错。
  • 写入不扩容数组:keyName 的路径遇到数组只改已存在的元素,越界静默忽略; 中间层级缺失时一律新建对象(即使该段是数字)——a.0.b 写进空 payload 得到 {"a":{"0":{"b":1}}}。
  • keyName: "" 不写:空路径的含义是"没指定写哪里",整行忽略(与 remove 的空路径一致)。 编辑器里未填写的空 set 行因此无害。要整体替换 payload 就写到一个具名字段下 (如 {"keyName":"report","value":"value","isVar":true},下游读 payload.report)。

没有 payload 字段 = 不写 payload。

inputMode 多输入

节点顶层字段,控制多个输入参数何时凑齐、何时执行:

取值 行为
"pair"(缺省) 所有参数到齐才执行一次,执行前清空缓存,然后重新等齐
"any" 任一上游到达即执行。未流入的参数取缓存的最新值,无缓存则取 value 缺省值
"continuous" 所有参数到齐才首次执行;此后任一上游再到达都用最新值重算(last-value)

留空 / 不写 = pair。未知取值加载时报错(不会静默落到缺省)。

三者是一条递进链,只有两个独立开关:是否设初始就绪门、凑齐后是否清空缓存。 所以 continuous = any + 初始就绪门——所有参数都被喂过一次之后,两者行为完全相同, 差异只在「首次」。选哪个就一句话:首次判断能不能用缺省值凑?能就 any, 不能就 continuous,每轮都要成套新值就用缺省的 pair。

取值优先级(三级):本次流入的非 nil 值 > 缓存的最新值 > value 缺省值。

  • 就绪 = 该参数本次流入、或此前被上游喂过、或声明了 value 缺省值。 所以「温度 > 阈值」里只有 value 的阈值参数不会导致死等;要让某参数必须等上游,就省略它的 value。
  • 流入 nil 视为「上游没有值可给」:不覆盖缺省值、不入缓存、不算就绪 (定时/cron 触发器的 msg.value 恒为 nil,否则定时链路会整条静默失效)。
  • pair 清空时只清有上游流入的参数,声明了 value 的参数填回缺省值。
  • 上游积压(同一参数连续到达多次)时 last-value 覆盖,中间值会丢弃。
  • inputs 为空的节点没有参数可等,inputMode 对其无意义,每次触发即执行。

为何缺省是 pair:多输入节点(AND/OR、比较、算术)语义上要求所有参数都有值。 若缺失参数静默取 nil,AND 恒假、!= 恒真(两边转不成数字就退化成字符串比较), 且不报错、不断流,排查成本极高。宁可「等不齐不执行」也不要「拿 nil 算出个假结果」。

⚠️ 参数永不就绪 → 整条下游静默死掉:有就绪门的节点(pair / continuous)上, 某参数既没有任何上游连线指向它、又没有 value 缺省值时,门控恒不通过。 最常见的原因是连线的 key 写错、或加了参数槽忘了连线。 condition / switch 的出口不算"有连线指向"——它们只触发不落槽(见 3.3)。 加载时会经 ErrorHandler 报一条告警(但不阻断加载——这也是合法的画图中间态)。

缓存跨触发:三种模式都用场景作用域的输入缓存,不随单次触发结束而清空—— 两个独立触发器分别喂 1、2 两个参数正是靠这一点凑齐的。场景 stop 时清空。 每个参数只留一个最新值,内存上界是「节点数 × 参数数」,不随触发次数增长。

执行时用哪条 msg:凑齐后取最后到达的那条 msg(它的 payload、value、error 生效), 其余上游只贡献自己那个参数的值。所以多输入节点要用 payload 时,注意它来自最后到达的分支。

已知限制:pair 靠「参数都被喂过」判凑齐,不保证这些值来自同一次触发。 触发器各自起 goroutine 并发点火,所以两条上游支线快慢不均时, 可能把「本轮的 a」和「上一轮的 b」配成一对。要严格同轮对齐需要给 msg 加触发批次标记, 目前未实现。

空输入参数的节点

inputs 为空数组表示本节点无输入参数,如 {"fn":"uuid7"}、 {"kind":"device-set-property","properties":{"t":30}}。

这类节点照样能作下游:它接收完整的 msg(payload、error 都在),只是没有参数可落, 所以 value 不落槽、脚本里 inputs 为空数组。需要上游数据时从 payload 读。

组合节点的数据边界

block(及 config 含 subflow 的 condition / function / action):

  • 输入:子流内的节点按参数 key 直接引用宿主的输入参数(无需连线)。 子流节点自己声明了同 key 的参数时,以子流自己的为准。
  • payload:子流继承宿主的 payload(同样 COW)。
  • 入口:子流中没有入边的节点为入口;有多个则全部执行。
  • 输出:由 outputs.value(子流节点 id)指定,取该节点的 msg 作为 block 的输出。
    • 指定节点本次没执行(在未走到的分支上)→ value 为 null,payload 取宿主流入那份
    • 指定节点执行多次 → 取最后一次
    • 省略 outputs.value → value 为 null,且子流对 payload 的修改不外泄
    • 指定节点执行失败 → block 走 errorIds,保留 error / errorId(子流内失败节点的 id)
    • 指向子流中不存在的 id → 加载时报错
  • 取到 msg 后,再执行 block 自身的 payload.set / payload.remove。
  • 出口只认宿主的:子流内节点的 outputs.ids 只在子流内寻址,指不到宿主流的节点; 子流跑完后由宿主 block 的 outputs 继续向下游流转。
  • 子流内的 inputMode 缓存按宿主实例隔离:同一子流被两个 block 引用时,各自独立计数。

宿主如何读取

ActionHandler 的 params 包含:

键 内容
config 节点 config 原样
inputs 各输入参数的值(按参数 key 索引,key 为空的参数不出现)
value 本次流入的临时值
payload 流作用域数据
error 上游失败信息,非异常分支为 ""(见 3.7)
errorId 失败节点 id,非异常分支为 ""
sceneId 动作所属场景 id;子流内的动作报宿主场景 id
nodeId 动作节点 id

sceneId / nodeId 是动作来源标识,供宿主记审计日志时回答"哪个场景的哪个节点做了这次副作用"。 与调试模式的 msg.meta.nodeId 不同,这两项不随 Debug 开关变化——审计是常态需求。

回调返回 (any, error):返回值非 nil 则写入 msg.value,error 非 nil 则本节点走异常出口(见 3.7)。 回调里对 payload 的修改不回流——要写流作用域用节点的 payload.set。

调试模式 meta

ruleflow.Debug = true 时,msg 额外带 meta,含 nodeId、nodeType、timestamp、traceId (一次触发的唯一 id,用于串联日志)。非调试模式下 msg 里没有 meta 字段—— 脚本不要依赖它,否则调试期能跑、生产期为 undefined。

需要定位失败节点时用 errorId(异常出口恒有,见 3.7)。

3.7. 出口与异常

outputs 的三个出口各表示一种含义,互不复用:

出口 含义 适用节点
ids 正常出口 全部
elseIds 条件为假(值分支) 仅 condition(其它类型忽略)
errorIds 本节点执行失败(异常出口) action / function / block / switch / condition / 自定义类型

需要"按取值走多条不同路径"用 switch 节点,不要拿 errorIds 当值分支用。

异常出口 errorIds

节点执行失败(动作回调报错、js 脚本异常、子流加载失败、区间表达式写错…)时触发 errorIds:

{
  "id": "f1",
  "type": "function",
  "config": { "script": "return payload.t * 2" },
  "outputs": {
    "ids": [{ "id": "正常后续id" }],
    "errorIds": [{ "id": "异常告警id" }]
  }
}

下游收到的 msg 保留失败前的 payload 与 value,并带上:

  • error:错误信息
  • errorId:失败节点 id(errorIds 的下游常被多个节点共用,靠它区分是哪一步失败)

所以异常分支能拿到"哪一步失败、当时数据是什么",可写告警、降级或转其它处理路径。 error / errorId 在下游节点成功执行后被清空——异常出口是它们唯一的来源。

四条关键语义:

  1. 失败只断本分支,不冒泡。errorIds 为空即空转,兄弟分支与后续节点照常执行—— 发短信失败不该连带掐掉同级的写库。
  2. 求值异常 ≠ 条件为假。condition 的脚本报错走 errorIds,不会走 elseIds; 否则"脚本挂了"会被误当成"条件不成立"而静默走假分支。
  3. 不做自动重试。引擎不重投失败节点,异常出口的下游自行处理(告警/降级/结束本分支)。
  4. 唯一冒泡的例外是未设置 ruleflow.ActionHandler——属宿主接线错误而非节点失败, 直接返回 ErrNoActionHandler 给调用方,可用 errors.Is 判断。

只有 source 类 trigger 有异常出口:第三方订阅(http/mqtt/sse/websocket)的连接失败与 脚本报错走 errorIds(msg 的 payload 为 {})。定时/倒计时/消息类 trigger 只做数据流转、 没有失败可能,不设异常出口。

全局异常回调 ErrorHandler

因为失败不冒泡,不接 errorIds 时问题会静默无痕(脚本 typo 就是这样丢失的)。 注入全局回调兜底,无论是否接了 errorIds 都会上报:

ruleflow.ErrorHandler = func(nodeID string, err error) {
    log.Warn().Str("node", nodeID).Err(err).Msg("场景节点执行异常")
}

errorIds 是每节点的异常出口(拿得到"哪一步失败、当时数据是什么",能在图里编排处理逻辑), ErrorHandler 是全流的兜底上报(保证可观测)。两者互补,建议都用。

加载时校验

以下情形在加载规则链时直接报错(属静态错误,早发现早修):

  • 节点缺少 id
  • 出口(ids/elseIds/errorIds/cases[].ids/defaultIds)指向不存在的节点 id
  • 出口指向 trigger 节点——触发器只能由消息或自身定时激活,不能被连线触发
  • outputs.value 指向 subflow 中不存在的节点 id
  • payload.set 某行 isVar: true,但 value 的第一段既不是本节点的参数 key、 也不是 inputs/value/payload/error/errorId 之一(写错变量名会静默取 null,故前置拦住)

以下情形不报错(画图中间态或有意留白):连线的 key 匹配不到下游参数(忽略数据、仍触发)、 参数既无上游连线也无 value 缺省值(取 null)、isVar 变量的后续路径段取不到值(取 null)。

3.8. 场景联动

由规则链组成的一组或多组业务流,即节点组成规则链,一个或多个规则链组成场景联动。数据格式为节点数组 JSON [...]。

安装

go get github.com/le5le-com/rule-flow

使用

package main

import (
	"time"

	ruleflow "github.com/le5le-com/rule-flow"
)

func main() {
	// 1. 设置全局动作处理器:所有场景的动作副作用都回调它
	ruleflow.ActionHandler = func(msg string, params map[string]any) (any, error) {
		// msg 为动作类型;params 含 "config"(节点配置)、"inputs"(输入参数)、
		// "value"(msg.value)、"payload"(流作用域数据)、"error"/"errorId"(异常出口信息)、
		// "sceneId"/"nodeId"(动作来源,供记审计日志),详见 3.6
		return "ok", nil
	}

	// 2. 注册并启动场景(nodesJSON 为节点数组 JSON);含触发器则注册后自动后台运行
	nodesJSON := []byte(`[
		{"id":"trig","type":"trigger","config":{"cron":"*/1 * * * * *"},"outputs":{"ids":[{"id":"act"}]}},
		{"id":"act","type":"action","config":{"kind":"message","text":"每秒执行一次"}}
	]`)
	if err := ruleflow.Start("scene-1", nodesJSON); err != nil {
		panic(err)
	}
	time.Sleep(3 * time.Second)

	// 3. 生命周期管理:
	//    ruleflow.Start("scene-1")    // 不传 nodes:启动已注册但已停止的场景
	//    ruleflow.Restart("scene-1")  // 先停再用注册 nodes 重启
	//    ruleflow.StopAll()           // 停止所有运行中的场景
	//    ruleflow.Remove("scene-1")   // 停止并从注册表删除
	//    ruleflow.Emit("scene-1", "device-report", data) // 向场景投递消息,触发监听该消息的触发器
	//    ruleflow.TriggerNode("scene-1", "n-manual", nil) // 从指定节点手动跑一次(调试,可从链路中间起跑)
	//    // 从中间起跑时那里本该有个上游,用 TriggerInput 把它造出来(三个字段各对应一条通道):
	//    ruleflow.TriggerNode("scene-1", "n-calc", &ruleflow.TriggerInput{
	//    	Value:   42,                                // → msg.value(节点有参数时顺带落第一个槽)
	//    	Inputs:  map[string]any{"b": 8},            // → 按参数 key 精确落槽,覆盖 Value 落的那个
	//    	Payload: map[string]any{"deviceId": "d-1"}, // → 流作用域初值,还原上游写进 payload 的现场
	//    	TraceID: "my-trace-1",                      // → 指定追踪 id,留空则引擎自行生成
	//    })
	ruleflow.Stop("scene-1")
}

License

MIT

About

Go实现的一套JSON语法的规则引擎。A rule engine implemented in Go that follows a set of JSON syntax rules.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages