Repository navigation
Expand file tree
/
Copy pathnode.go
More file actions
331 lines (304 loc) · 13.9 KB
/
Copy pathnode.go
File metadata and controls
331 lines (304 loc) · 13.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
package ruleflow
import (
"encoding/json"
"fmt"
"sync"
"time"
)
// 节点类型
const (
TypeTrigger = "trigger" // 触发器
TypeCondition = "condition" // 条件
TypeAction = "action" // 动作
TypeBlock = "block" // 组合节点(内含子流)
TypeFunction = "function" // 函数(自定义脚本)
TypeSwitch = "switch" // 多路分支(按输入值匹配多个 case,各 case 有独立下游)
)
// 多输入模式(节点顶层 inputMode):控制多个输入参数何时凑齐、何时执行。
//
// 三者是一条递进链,只有两个独立开关:**是否设初始就绪门**、**凑齐后是否清空缓存**。
// - any:无门、不清空 —— 任一到达即执行
// - pair(缺省):有门、清空 —— 凑齐执行一次,然后重新等齐
// - continuous:有门、不清空 —— 凑齐后任一到达都用最新值重算
//
// 所以 continuous = any + 初始就绪门:所有参数都被喂过一次之后,两者行为完全相同,
// 差异只在"首次"。选哪个的判据就一句话:**首次判断能不能用缺省值凑**——
// 能就 any,不能就 continuous,每轮都要成套新值就 pair(缺省)。
//
// 缺省是 pair 而非 any:多输入节点(AND/OR、比较、算术)语义上要求所有参数都有值,
// 若缺失参数静默取 nil 缺省值,AND 恒假、`!=` 恒真(两边转不成数字就退化成字符串比较),
// 且不报错、不断流,排查成本极高。宁可"等不齐不执行"也不要"拿 nil 算出个假结果"。
const (
InputModeAny = "any" // 任一上游到达即执行;未流入的参数取缓存最新值,无缓存则取 value 缺省值
InputModePair = "pair" // 所有参数就绪才执行,执行前清空缓存、等下一次重新凑齐(缺省)
InputModeContinuous = "continuous" // 所有参数就绪才首次执行,此后任一上游再到达都用最新值重算
)
// inputModes 合法取值集合。空串等同 pair(缺省),不在此表内。
var inputModes = map[string]bool{
InputModeAny: true,
InputModePair: true,
InputModeContinuous: true,
}
// mode 本节点生效的多输入模式:空串归一为缺省的 pair。
// 取值合法性由加载时的 validateInputMode 保证,此处不再兜底未知值。
func (n *Node) mode() string {
if n.InputMode == "" {
return InputModePair
}
return n.InputMode
}
// Input 节点的一个输入参数(形参)声明。与上游是谁无关——同一个节点可被复制、
// 组合节点可被多处复用,各实例参数签名相同而上游各异,所以路由信息放在上游 outputs 里
// (见 Exit)。
// - Key:参数标识(节点内唯一)。这是参数唯一的引用方式,四处都用它:上游连线的
// Exit.Key 匹配它、脚本里的变量名、宿主动作 inputs map 的键、子流按 key 匹配宿主参数槽。
// 硬约束只有两条:**节点内唯一**、**不含点号**——含点号的 key 在 payload.set 的
// isVar 表达式里永远取不到值(varName/resolveVar 在第一个点号处切分表达式),
// 而这是静默取 nil,不报错。
// 另有一条软约束:key 会作 goja 全局变量名注入(见 evalScript),非合法 JS 标识符
// (如纯数字 "0")注入不报错但脚本里写不了裸名字,只能用 inputs[下标] 或 this["0"]。
// 用 script 节点时值得挑个标识符,纯 map 取值的场合(动作 inputs / bind / 连线匹配 /
// payload.set)任何合法 JSON 键都可以。
// - Label:显示文案,仅供编辑器展示,引擎不解释。可随意改名而不影响任何引用。
// - Value:缺省值。上游未流入或流入 nil 时取它;对 pair/continuous 而言,
// 声明了 value 的参数视为已就绪。
//
// 为何标识与显示分家:Key 被连线、脚本、payload.set 引用,改它要连带改引用处
// (payload.set 的变量名对不上会导致整个规则链加载失败),所以它是稳定标识而非文案。
// 想改给用户看的名字时改 Label——它不参与任何寻址。
type Input struct {
Key string `json:"key"`
Label string `json:"label"`
Value any `json:"value"`
}
// Exit 一条出口连线:触发下游 ID,数据落到下游 Key 指向的输入参数。
// JSON 支持两种写法:{"id":"a1","key":"temp"} 与简写 "a1"(key 省略=落第一个参数)。
//
// 例外:condition / switch 的出口不落槽,Key 被忽略(见 Node.routesOnly)。
//
// 字段名是 key 而非 input:它的值必须与下游某个 inputs[].key 一字不差
// (slotIndex 只按 key 匹配,从来没有按下标的分支)。这个字段曾叫 input,
// 编辑器与人手都会误读成"落第几个输入"而写下标,且下标一旦与某个参数 key 同形
// (参数 key 本该是 JS 标识符,但引擎不强制,现场确有叫 "0" 的)就会意外命中,
// 把上游的 value 写进一个不该写的槽——不报错、不断流,只是结果不对。
// 所以 key 只收字符串:写成数字在加载时就报错,而不是留到运行期变成一次错命中。
type Exit struct {
ID string `json:"id"`
Key string `json:"key"`
}
// UnmarshalJSON 支持出口项的联合类型:字符串简写 或 {id,key} 对象。
func (e *Exit) UnmarshalJSON(b []byte) error {
var s string
if err := json.Unmarshal(b, &s); err == nil {
e.ID, e.Key = s, ""
return nil
}
var obj struct {
ID string `json:"id"`
Key string `json:"key"` // 只收字符串,见 Exit 的说明
}
if err := json.Unmarshal(b, &obj); err != nil {
return fmt.Errorf("出口项须为节点 id 字符串或 {id,key} 对象(key 是下游参数名,须为字符串): %w", err)
}
e.ID, e.Key = obj.ID, obj.Key
return nil
}
// Output 节点输出。三个出口各表示一种含义,互不复用:
// - IDs:正常出口。每一项是一条连线,数据落到下游哪个输入参数由该项的 Key 决定。
// - ElseIDs:condition 节点专用的"假分支"——条件为真触发 IDs、为假触发 ElseIDs
// (无 ElseIDs 时为假则断流)。其它节点类型忽略 ElseIDs。
// - ErrorIDs:异常出口。本节点执行失败(动作回调报错、脚本异常、子流失败…)时触发,
// 下游收到的 msg 保留失败前的 payload/value,并带 error/errorId。
// 未接 ErrorIDs 时失败只断本分支(不连坐兄弟分支),并上报全局 ErrorHandler。
//
// 还有一个可选字段:
// - Value:仅含 subflow 的节点(block / 自定义 condition / function / action)——
// 取哪个子流节点的 msg 作为本节点的输出。
type Output struct {
IDs []Exit `json:"ids"`
ElseIDs []Exit `json:"elseIds,omitempty"`
ErrorIDs []Exit `json:"errorIds,omitempty"`
Value string
}
// outputAlias 承接 outputs 的原始 JSON:Value 允许写成非字符串(旧配置的固定值),
// 统一取其字符串形式。
type outputAlias struct {
IDs []Exit `json:"ids"`
ElseIDs []Exit `json:"elseIds"`
ErrorIDs []Exit `json:"errorIds"`
Value any `json:"value"`
}
func (o *Output) UnmarshalJSON(b []byte) error {
var a outputAlias
if err := json.Unmarshal(b, &a); err != nil {
return err
}
o.IDs, o.ElseIDs, o.ErrorIDs = a.IDs, a.ElseIDs, a.ErrorIDs
o.Value = anyToString(a.Value)
return nil
}
// SetItem payload.set 的一行写入。KeyName 是写入流作用域的点路径(空串=不写,
// 见 Msg.setPath);
// Value 是要写的内容——IsVar 为假时按字面量原样写入(可为任意 JSON),为真时
// 当作变量表达式求值:首段为变量名,其余为点路径(见 resolveVarPath)。
type SetItem struct {
KeyName string `json:"keyName"`
Value any `json:"value"`
IsVar bool `json:"isVar"`
}
// PayloadWrite 节点对流作用域的写入:Set 按声明顺序逐行写入(后行可覆盖前行,且能读到
// 前行的结果)、Remove 在全部 Set 之后按点路径逐项删除。都在节点执行之后生效。
type PayloadWrite struct {
Set []SetItem `json:"set"`
Remove []string `json:"remove"`
}
// Node 规则节点。Config 为节点类型相关的配置(反序列化时直接解析为 map,按类型解释)。
type Node struct {
ID string `json:"id"`
Name string `json:"name"`
Type string `json:"type"`
Config map[string]any `json:"config"`
Inputs []Input `json:"inputs"`
InputMode string `json:"inputMode"`
Payload *PayloadWrite `json:"payload"`
Outputs Output `json:"outputs"`
// 含 subflow 的节点(block / 自定义 condition/function/action)的子流编译缓存:
// 首次执行时编译一次,之后复用(子流结构只读,执行状态在 execContext 与场景 store 里)。
subOnce sync.Once
subFlow *Flow
subErr error
// debounce/throttle 阀门节点(function 节点 fn=debounce/throttle)的状态,跨触发存活。
gateMu sync.Mutex
debTimer *time.Timer // debounce 当前待触发定时器
debFirst time.Time // debounce 本轮首次输入时刻(maxWait 起点)
debMsg *Msg // debounce 最后一次输入的整条 msg
debIn *inbound // debounce 最后一次输入的入参(放行时供 payload.set 读参数变量)
lastFire time.Time // throttle 上次放行时刻
}
// isGate 该 function 节点是否为 debounce/throttle 阀门。
func (n *Node) isGate() bool {
if n.Type != TypeFunction {
return false
}
fn := n.ConfigString("fn")
return fn == "debounce" || fn == "throttle"
}
// waits 该节点是否设初始就绪门(pair/continuous):所有参数就绪前不执行。
// inputs 为空时没有参数可等,恒为 false。
func (n *Node) waits() bool {
if len(n.Inputs) == 0 {
return false
}
m := n.mode()
return m == InputModePair || m == InputModeContinuous
}
// caches 该节点是否缓存各参数的最新值。三种模式都缓存——any 靠它让"未流入的参数"
// 取到上次的值而不是退回缺省值,pair/continuous 靠它跨触发凑齐。
// inputs 为空时没有参数可缓存,恒为 false。
//
// 三种模式都缓存,所以 store 是通用设施而非 pair/continuous 专用。
func (n *Node) caches() bool {
return len(n.Inputs) > 0
}
// consumes 该节点是否在执行前清空缓存(仅 pair):凑齐即消费,下一轮重新等齐。
// any/continuous 保留最新值,所以恒为 false。
func (n *Node) consumes() bool {
if len(n.Inputs) == 0 {
return false
}
return n.mode() == InputModePair
}
// routesOnly 本节点只选路线、不产出自己的值:condition 与 switch。
// 它们的出口只起触发作用,msg 照常向下游流转(value/payload 都在),但**不把 value
// 落进下游的参数槽**,也不写入场景缓存(见 Flow.deliver)。
//
// 为何必须这样:condition 不修改 msg,流下去的 value 是**上游某个节点的值**,
// 而多输入条件在 pair 模式下 msg 取"最后到达的那条"(见 execNode)——谁最后到达取决于
// 上游 outputs.ids 的数组次序,也就是编辑器的画线顺序。若让它落槽,下游一个声明了
// 缺省值的参数会被这个"别人的、次序决定的"值静默顶掉:现场那条
// `AND → randInt(min:15,max:20)` 的线就是把 min 顶成了上游 randInt(0,1) 的 1,
// 节点名写着 15-20、实际跑 1-20,不报错、不断流。
//
// 这不是新规则,是把既有原则补齐:trigger 同样是不产值节点,定时/cron 触发器的
// msg.value 恒为 nil,store.feed / mergeSlots / resolveInputs 三处都专门写了
// "nil 视为上游没有值可给,不顶掉缺省值、不标就绪"。condition 因为转发的是别人的
// 非 nil 值,从这条原则的后门溜了进去。
//
// 要把判据传给下游就用 payload.set(isVar 能读 value/inputs/参数 key)——那是显式的。
func (n *Node) routesOnly() bool {
return n.Type == TypeCondition || n.Type == TypeSwitch
}
// slotIndex 把连线的 key 解析为本节点的参数下标:按 Input.Key 匹配。
// key 为空 → 第一个参数。匹配不到返回 -1(忽略数据、仍触发)。
func (n *Node) slotIndex(key string) int {
if len(n.Inputs) == 0 {
return -1
}
if key == "" {
return 0
}
for i, ip := range n.Inputs {
if ip.Key == key {
return i
}
}
return -1
}
// allExits 本节点的全部出口连线(正常/假分支/异常 + switch 的各 case 与 default)。
// 供子流入口推导与静态校验用。
func (n *Node) allExits() []Exit {
exits := make([]Exit, 0, len(n.Outputs.IDs)+len(n.Outputs.ElseIDs)+len(n.Outputs.ErrorIDs))
exits = append(exits, n.Outputs.IDs...)
exits = append(exits, n.Outputs.ElseIDs...)
exits = append(exits, n.Outputs.ErrorIDs...)
if n.Type == TypeSwitch {
if cases, err := parseSwitchCases(n.Config["cases"]); err == nil {
for _, c := range cases {
exits = append(exits, c.IDs...)
}
}
exits = append(exits, parseExits(n.Config["defaultIds"])...)
}
return exits
}
// ConfigString 取 config 里的字符串字段。
func (n *Node) ConfigString(key string) string {
if v, ok := n.Config[key]; ok {
if s, ok := v.(string); ok {
return s
}
}
return ""
}
// hasConfig 判断 config 是否含某键(用于 subflow 之类的存在性判断)。
func (n *Node) hasConfig(key string) bool {
_, ok := n.Config[key]
return ok
}
// anyToString 把 JSON 里可能写成数字/布尔的标量统一取字符串形式;非标量与 nil 返回空串。
func anyToString(v any) string {
switch x := v.(type) {
case nil:
return ""
case string:
return x
case float64:
// 整数不带小数点,便于 "value":1 与 "id":"1" 互认
if x == float64(int64(x)) {
return fmt.Sprintf("%d", int64(x))
}
return fmt.Sprintf("%v", x)
case bool:
return fmt.Sprintf("%v", x)
default:
return ""
}
}
// ParseNodes 解析节点数组 JSON。config 随反序列化直接解析为 map,无需二次解析。
func ParseNodes(data []byte) ([]*Node, error) {
var nodes []*Node
if err := json.Unmarshal(data, &nodes); err != nil {
return nil, err
}
return nodes, nil
}