-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathflow.go
More file actions
986 lines (925 loc) · 36.1 KB
/
Copy pathflow.go
File metadata and controls
986 lines (925 loc) · 36.1 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
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
package ruleflow
import (
"context"
"errors"
"fmt"
"maps"
"strconv"
"sync"
"time"
)
// Flow 规则链:由节点数组加载而成,负责编排与执行流转。库内部使用;
// 对外入口是包级 Start/Restart/Stop/StopAll/Remove。
type Flow struct {
mu sync.RWMutex
nodes map[string]*Node // id -> 节点
order []string // 加载顺序(用于遍历触发器)
mqttPool *mqttPool // 第三方 MQTT 触发器共享连接池(懒初始化)
lifeCtx context.Context // 场景生命周期 ctx(StartTriggers 存入),供 debounce 延迟触发取消用
// store 场景作用域:inputMode 输入缓存等引擎内部状态,跨触发存活、场景停止即清空。
// 子流共享宿主的 store,靠 scope 前缀区分实例。
store *store
scope string // 本流在场景内的作用域前缀:主流为空,子流为 "宿主节点id/"
// sceneID 所属场景 id:TraceGate/TraceHandler 需要,子流继承宿主的。
// 为空是合法状态(测试里直连 newFlow() 的流),此时由宿主的 TraceGate 决定(约定返回 false)。
sceneID string
}
// newFlow 创建空规则链(内部)。
func newFlow() *Flow {
return &Flow{nodes: map[string]*Node{}, store: newStore()}
}
// load 加载节点数组 JSON 并做静态校验。
// 校验失败时不改动 f 的既有节点——半加载的规则链比加载失败更难排查
// (调用方若忽略 error,节点已在表里、会带着非法配置照跑)。
func (f *Flow) load(nodesJSON []byte) error {
nodes, err := ParseNodes(nodesJSON)
if err != nil {
return fmt.Errorf("解析规则链失败: %w", err)
}
f.mu.Lock()
defer f.mu.Unlock()
staged := make(map[string]*Node, len(nodes))
order := make([]string, 0, len(nodes))
for _, n := range nodes {
if n.ID == "" {
return fmt.Errorf("节点缺少 id")
}
staged[n.ID] = n
order = append(order, n.ID)
}
prevNodes, prevOrder := f.nodes, f.order
f.nodes, f.order = staged, order
if err := f.validate(); err != nil {
f.nodes, f.order = prevNodes, prevOrder // 回滚,不留半加载状态
return err
}
return nil
}
// validate 加载时静态校验(需持有 f.mu):
// - 出口指向不存在的节点 id;
// - 出口指向 trigger 节点(触发器只能由消息或自身定时激活,不能被连线触发);
// - outputs.value 指向 subflow 中不存在的节点 id;
// - payload.set 的 isVar 行,变量名不是本节点参数 key 也不是保留变量名;
// - inputMode 不是合法取值(空串 / any / pair / continuous)。
//
// 为何未知 inputMode 要报错而非忽略:缺省是 pair,未知值若静默落到缺省,
// 旧配置里的 "waitAll" 会从"任一更新即重算"退化成"凑齐即消费"——行为倒退且完全无痕。
//
// 不校验但会经 ErrorHandler 报告的情形(画图中间态,见 warnDeadInputs):
// 有就绪门的节点上,某参数既无入边又无缺省值 → 永不就绪、整条下游静默死掉。
//
// 不校验的情形(画图中间态或有意留白):连线的 key 匹配不到下游参数、
// isVar 变量的后续路径段取不到值(取 null)。
func (f *Flow) validate() error {
for _, id := range f.order {
n := f.nodes[id]
if err := validateInputMode(n); err != nil {
return err
}
exits := map[string][]Exit{
"ids": n.Outputs.IDs,
"elseIds": n.Outputs.ElseIDs,
"errorIds": n.Outputs.ErrorIDs,
}
if n.Type == TypeSwitch {
cases, err := parseSwitchCases(n.Config["cases"])
if err != nil {
return fmt.Errorf("节点 %s 的 cases 非法: %w", n.ID, err)
}
for i, c := range cases {
exits["cases["+strconv.Itoa(i)+"].ids"] = c.IDs
}
exits["defaultIds"] = parseExits(n.Config["defaultIds"])
}
for field, list := range exits {
for _, e := range list {
down, ok := f.nodes[e.ID]
if !ok {
return fmt.Errorf("节点 %s 的 %s 指向不存在的节点: %s", n.ID, field, e.ID)
}
if down.Type == TypeTrigger {
return fmt.Errorf("节点 %s 的 %s 不能指向触发器节点: %s", n.ID, field, e.ID)
}
}
}
// outputs.value 是子流节点 id:只对含 subflow 的节点有意义,须指向子流内存在的节点。
if n.Outputs.Value != "" && n.hasConfig("subflow") {
sub, err := n.subNodes()
if err != nil {
return fmt.Errorf("节点 %s 的 subflow 非法: %w", n.ID, err)
}
if _, ok := sub[n.Outputs.Value]; !ok {
return fmt.Errorf("节点 %s 的 outputs.value 指向 subflow 中不存在的节点: %s", n.ID, n.Outputs.Value)
}
}
if err := validatePayloadVars(n); err != nil {
return err
}
}
f.warnDeadInputs()
return nil
}
// validateInputMode 校验节点的 inputMode 取值。空串合法(归一为缺省的 pair)。
func validateInputMode(n *Node) error {
if n.InputMode == "" || inputModes[n.InputMode] {
return nil
}
return fmt.Errorf("节点 %s 的 inputMode 非法: %q(可用 any / pair / continuous,留空即 pair)", n.ID, n.InputMode)
}
// warnDeadInputs 报告"永不就绪"的参数:有就绪门的节点(pair/continuous)上,
// 某参数既没有任何上游连线指向它、又没有 value 缺省值 —— store.take 的门控恒不通过,
// 该节点及其整条下游静默死掉,运行期毫无痕迹。
//
// 为何只告警不报错:这也是合法的画图中间态(连了一半的多输入节点),
// 报错会让编辑器存不下草稿。但它同时是最典型的配置事故——连线的 key 写错、
// 或加了槽忘了连线,都会落到这里,所以必须留个痕。
//
// 只看本流:子流内的节点可按参数 key 直接引用宿主的槽(见 execContext.hostSlots),
// 无入边不代表拿不到值,所以子流节点不在此检查范围内(它们不在 f.nodes 里)。
//
// condition/switch 的出口不算"喂了值"——它们只触发不落槽(见 Flow.deliver)。
// 判据必须与 deliver 逐条一致:若这里把它们算进来,"某参数只由一条条件出口喂"的节点
// 就会漏报,而运行期门控同样不放行,等于把一个静默死亡换成另一个。
func (f *Flow) warnDeadInputs() {
if ErrorHandler == nil {
return
}
// 汇总每个节点被连线指向的参数槽
fed := map[string]map[int]bool{}
for _, id := range f.order {
if f.nodes[id].routesOnly() {
continue // 只选路线的节点不给下游落值
}
for _, e := range f.nodes[id].allExits() {
down, ok := f.nodes[e.ID]
if !ok {
continue
}
if slot := down.slotIndex(e.Key); slot >= 0 {
if fed[e.ID] == nil {
fed[e.ID] = map[int]bool{}
}
fed[e.ID][slot] = true
}
}
}
for _, id := range f.order {
n := f.nodes[id]
if !n.waits() {
continue
}
for i, ip := range n.Inputs {
if fed[id][i] || ip.Value != nil {
continue
}
ErrorHandler(id, fmt.Errorf("参数[%d] %q 既无上游连线又无 value 缺省值,inputMode=%s 下永不就绪,本节点及其下游不会执行", i, ip.Key, n.mode()))
}
}
}
// validatePayloadVars 校验 payload.set 各 isVar 行的变量名:只认本节点声明了 key 的参数
// 与五个保留变量。写错变量名会静默取 null(路径读不到即 null 是刻意的宽松),
// 所以在加载时把"变量名"这一层拦住——它是可静态判定的,不该留到运行期变成空值。
func validatePayloadVars(n *Node) error {
if n.Payload == nil {
return nil
}
for i, item := range n.Payload.Set {
if !item.IsVar {
continue
}
expr, ok := item.Value.(string)
if !ok || expr == "" {
return fmt.Errorf("节点 %s 的 payload.set[%d] 标了 isVar,value 须为变量表达式字符串", n.ID, i)
}
name := varName(expr)
switch name {
case scriptInputsVar, scriptValueVar, scriptPayloadVar, scriptErrorVar, scriptErrorIDVar:
continue
}
found := false
for _, ip := range n.Inputs {
if ip.Key != "" && ip.Key == name {
found = true
break
}
}
if !found {
return fmt.Errorf("节点 %s 的 payload.set[%d] 变量 %q 既非本节点参数 key,也非保留变量(inputs/value/payload/error/errorId)", n.ID, i, name)
}
}
return nil
}
// subNodes 解析本节点 config.subflow 的节点 id 集合(仅供加载校验)。
func (n *Node) subNodes() (map[string]struct{}, error) {
bs, err := marshal(n.Config["subflow"])
if err != nil {
return nil, err
}
nodes, err := ParseNodes(bs)
if err != nil {
return nil, err
}
ids := make(map[string]struct{}, len(nodes))
for _, s := range nodes {
ids[s.ID] = struct{}{}
}
return ids, nil
}
// node 取节点。
func (f *Flow) node(id string) *Node {
f.mu.RLock()
defer f.mu.RUnlock()
return f.nodes[id]
}
// TriggerInput 一次入口触发喂给起始节点的数据。三个字段各对应引擎里的一条通道,
// 与真实运行时下游从上游收到的东西一一对应:
//
// Value → msg.value 上游节点算出来的输出值
// Inputs → 参数槽 上游连线按 key 落进来的参数
// Payload → msg.payload 流作用域
//
// 全部可为零值,等同"只触发、不喂数据"。nil 的 *TriggerInput 与 &TriggerInput{} 等价。
//
// 存在的理由是**从链路中间起跑**(调试):那里本该有个上游,没有上游就得把它造出来。
// 只造 Value 是不够的——中间节点大量依赖上游写进 payload 的字段、依赖多个参数槽,
// 缺了它们节点行为与真实运行不一致,调试就失去意义。
type TriggerInput struct {
// Value 落 msg.value;节点声明了输入参数时同时落**第一个**参数槽。
// "顺带落第一个槽"是入口层的通用语义,emit 与定时触发器都依赖它——
// 它们不知道目标节点的参数 key,只能按位置喂(见 entryInbound)。
Value any
// Inputs 按 Input.Key 精确落槽,覆盖 Value 落进第一个槽的那个值。
// 只有手动触发/调试会用:多参数节点(pair/continuous)从中间起跑时,
// 光靠 Value 只能喂一个参数,其余会取到**上次运行残留的缓存值**而不可重现。
// key 匹配不到参数则忽略该项,与 deliver 的落槽规则一致。
Inputs map[string]any
// Payload 作 msg.payload 的初值(内部会拷一份,调用方的 map 不会进流)。
// 不给则为空对象,与其它触发来源一致。
Payload map[string]any
// TraceID 指定本次触发的追踪 id(见 TraceEvent.TraceID),留空则引擎自行生成。
//
// 用途是让调用方**在触发之前**就知道这次的 id,从而按 id 过滤掉同场景其它触发来源
// (设备上报/定时器/cron/emit/第三方订阅)的事件。Trigger 是同步的,而追踪事件在执行
// 过程中就已发出,等它返回再拿 id 已经晚了。
//
// 只决定追踪的**命名**,不决定是否追踪——后者仍由 TraceGate 说了算,
// 否则任何调用方都能绕过宿主的开关强行打开追踪。
//
// 引擎不校验内容(库不管宿主的输入):它会原样进 TraceEvent 并被宿主序列化发出去,
// 长度与字符集由调用方负责。
TraceID string
}
// slots 按 n 的参数声明把 Value/Inputs 解成槽位表。nil 值一律不落槽——
// 与引擎各处对 nil 的一贯处理相同(nil = 上游没有值可给,见 resolveInputs)。
func (t *TriggerInput) slots(n *Node) map[int]any {
slots := map[int]any{}
if t == nil {
return slots
}
if len(n.Inputs) > 0 && t.Value != nil {
slots[0] = t.Value
}
for k, v := range t.Inputs {
if v == nil {
continue
}
if i := n.slotIndex(k); i >= 0 {
slots[i] = v // 显式点名优先于 Value 落的第一个槽
}
}
return slots
}
// Trigger 从指定节点开始执行(手动/调试)。in 为 nil 表示只触发、不喂数据。
func (f *Flow) Trigger(nodeID string, in *TriggerInput) error {
n := f.node(nodeID)
if n == nil {
return fmt.Errorf("节点不存在: %s", nodeID)
}
return f.execNode(newExecContext(), n, f.entryInbound(n, in))
}
// entryInbound 构造入口节点(触发器/手动触发)的入参:新建一条 msg,
// value 与各参数槽按 in 解析(见 TriggerInput)、payload 默认为空对象。
// in 为 nil 或各字段为零值时即"无数据流入",参数照常取自己的 value 缺省值。
//
// 注意这里**不**把值写进 store 缓存(deliver 会写):手动触发是调试注入,
// 不该污染 pair/continuous 跨触发的缓存——否则调一次的残留会在下一次真实
// 触发时才显现,症状与原因隔着十万八千里。就绪门看的是本次流入(take 的 live 参数),
// 不写缓存不影响本次执行。
func (f *Flow) entryInbound(n *Node, in *TriggerInput) *inbound {
slots := in.slots(n)
var value any
if in != nil {
value = in.Value
}
msg := newMsg(value)
if in != nil && len(in.Payload) > 0 {
msg.Payload = maps.Clone(in.Payload)
}
var traceID string
if in != nil {
traceID = in.TraceID // 空串即"未指定",traceMeta 自行生成
}
// 追踪与否在这里一次性决定:这是全部触发来源的唯一漏斗
// (emitWithin / 定时器 fire / Trigger / 第三方订阅都走它)。
// Debug 模式下 newMsg 已盖过章,不重复盖——但调用方**指定**了 id 时要改写,
// 否则 Debug 一开,调用方按自己那个 id 过滤就一条都匹配不上(且无从察觉)。
if f.traced() && (msg.Meta == nil || traceID != "") {
msg.traceMeta(traceID)
}
return newInbound(msg, slots)
}
// execNode 执行单个节点,并按出口触发下游。单节点失败返回 error 但不 panic。
// in 携带触发本次执行的 msg 与各参数流入值。
func (f *Flow) execNode(ctx *execContext, n *Node, in *inbound) error {
if n == nil {
return nil
}
// 环检测按传播路径:本节点已在当前路径上则拦住(菱形扇入不受影响)。
if ctx.onPath(n.ID) {
return nil
}
if in == nil {
in = newInbound(newMsg(nil), nil)
}
// 子流内:本节点声明的参数若与宿主槽同名,直接引用宿主的值(上游连线流入优先)。
// 必须在就绪门之前——子流节点按参数 key 直接取宿主的槽、没有连线,
// 所以 store 里不会有它们的值,门控若只看 store 会把整个子流饿死。
in.slots = ctx.hostSlots(n, in.slots)
// 多输入模式:按 waits/consumes 判就绪与清空,按 caches 用缓存的最新值补齐未流入的槽。
// 缓存在场景作用域里、跨触发存活(见 store)。
// - waits()(pair/continuous):参数未全就绪则本次不执行,等下一个参数到来;
// - caches()(三种模式):未流入的槽取缓存最新值,本次流入的槽优先(见 resolveInputs 三级优先级);
// - consumes()(pair):取到快照后即清空,执行完就开始重新等齐。
if n.caches() {
ready, snap := f.store.take(f.key(n), n, in.slots, n.waits(), n.consumes())
if !ready {
// 就绪门未放行是静默断流,调试时必须可见——否则"为什么这个节点不执行"无从回答。
// missing 的计算只在追踪中发生(带 Meta 才进去),正常路径零成本。
if in.msg != nil && in.msg.Meta != nil {
f.trace(n, in.msg, in, TracePhaseBlocked, nil,
f.store.missing(f.key(n), n, in.slots))
}
return nil // 未凑齐:本次不执行,等下一个参数到来
}
in = newInbound(in.msg, mergeSlots(in.slots, snap)) // msg 取最后到达的那条
}
ctx.push(n.ID)
defer ctx.pop(n.ID)
msg := in.msg
msg.stamp(n)
switch n.Type {
case TypeCondition:
ok, err := f.execCondition(n, in)
if err != nil {
// 求值异常(脚本报错等):走异常出口,不等同于"条件为假"
return f.nodeError(ctx, n, msg, err)
}
// condition 不修改 msg:判断结果已由"走哪个出口"表达,payload/value 原样流下。
// 但流下的 value 不落下游参数槽——它是上游某个节点的值,不是本节点算出来的
// (见 Node.routesOnly)。要给下游传值用 payload.set。
if !ok {
return f.finish(ctx, n, msg, in, n.Outputs.ElseIDs) // 假分支(无则断流)
}
return f.finish(ctx, n, msg, in, n.Outputs.IDs)
case TypeSwitch:
// switch 同样不修改 msg、同样不给下游落槽,只决定走哪些 case 出口
exits, err := f.switchExits(n, in)
if err != nil {
return f.nodeError(ctx, n, msg, err)
}
return f.finish(ctx, n, msg, in, exits)
case TypeAction:
if actionKind(n) == "emit" {
// emit:引擎内部向本场景投递一条消息,激活监听它的触发器(场景内联动)。
// 被激活的链路是一次新的传播(payload 从 {} 开始),不回调宿主。
// 投递后本节点继续向下游流转——emit 是"额外投递一条消息",不是终止节点。
f.emitMessage(ctx, n, in)
return f.nodeDone(ctx, n, msg, in, nil, nil)
}
// 自定义动作:config 含 subflow 时跑子流(不回调宿主),子流输出作为本节点输出。
if n.hasConfig("subflow") {
out, err := f.runSubflow(ctx, n, in)
return f.nodeDone(ctx, n, msg, in, out, err)
}
out, err := f.execAction(n, in)
return f.nodeDone(ctx, n, msg, in, out, err)
case TypeFunction:
// debounce/throttle 阀门:不产出值,缓存整条 msg,由阀门逻辑控制下游触发时序
if n.isGate() {
return f.execGate(n, in)
}
out, err := f.execFunction(ctx, n, in)
return f.nodeProduced(ctx, n, msg, in, out, err)
case TypeBlock:
out, err := f.runSubflow(ctx, n, in)
return f.nodeProduced(ctx, n, msg, in, out, err)
case TypeTrigger:
// 触发器被激活:不产值,msg 原样向下游流转
return f.nodeDone(ctx, n, msg, in, nil, nil)
default:
// 自定义节点类型(经 Register 注册)
if fn, ok := customNodes[n.Type]; ok {
byKey, _ := resolveInputs(n, in)
out, err := fn(n, byKey)
return f.nodeDone(ctx, n, msg, in, out, err)
}
return f.nodeDone(ctx, n, msg, in, nil, nil)
}
}
// key 节点在场景作用域里的缓存键:子流实例带宿主前缀,避免同一子流被两处引用时串扰。
func (f *Flow) key(n *Node) string {
return f.scope + n.ID
}
// nodeDone 节点执行完毕的统一收尾,是 action/function/block/自定义类型的共同出口:
// - 失败 → 走异常出口 errorIds(msg 保留失败前的 payload/value,带上 error/errorId);
// - 成功 → 有返回值则写入 msg.value(无返回值则 msg 原样流下,副作用节点不掐断数据流),
// 再按 payload.set/remove 写流作用域、清空 error/errorId,最后走正常出口。
//
// 失败不向上冒泡——否则一个节点失败会连带掐掉兄弟分支与后续所有节点。
func (f *Flow) nodeDone(ctx *execContext, n *Node, msg *Msg, in *inbound, out any, err error) error {
if err != nil {
return f.nodeError(ctx, n, msg, err)
}
if out != nil {
msg.Value = out
}
return f.finish(ctx, n, msg, in, n.Outputs.IDs)
}
// nodeProduced 与 nodeDone 同,但 out 为 nil 也照写 msg.value:
// function/block 是算值节点,null 是有效计算结果,不该退回成"透传上游值"。
func (f *Flow) nodeProduced(ctx *execContext, n *Node, msg *Msg, in *inbound, out any, err error) error {
if err != nil {
return f.nodeError(ctx, n, msg, err)
}
msg.Value = out
return f.finish(ctx, n, msg, in, n.Outputs.IDs)
}
// finish 节点收尾:先写流作用域(payload.set → payload.remove),再清空 error/errorId,
// 回报给子流观察者,最后沿给定出口向下游流转。
//
// 为何先写后清:异常分支上的节点要能用 isVar 把 error/errorId 落进 payload
// (如 {keyName:"fault.msg", value:"error", isVar:true}),清早了就只剩空串。
// 清空放在写入之后,所以下游收到的 msg 仍是"成功即无错"。
func (f *Flow) finish(ctx *execContext, n *Node, msg *Msg, in *inbound, exits []Exit) error {
f.applyPayload(n, msg, in)
msg.clearError()
if ctx.watch != nil {
ctx.watch(n.ID, msg, nil)
}
// 放在 applyPayload 之后:上报的 payload 是本节点写完的最终值,与下游收到的一致。
// 放在 fanoutExits 之前:事件顺序即执行顺序,编辑器不必自己排。
f.trace(n, msg, in, TracePhaseDone, nil, nil)
return f.fanoutExits(ctx, n, exits, msg)
}
// applyPayload 执行节点的流作用域写入:payload.set 按声明顺序逐行写、
// payload.remove 在全部 set 之后逐项删。没有 payload 字段 = 不写。
//
// 每行的 keyName 是写入的点路径(空串=不写,见 Msg.setPath);IsVar 为真时 Value 作变量表达式求值,
// 否则原样写入。变量表 vars 逐行重建——payload 是实时的,后行能读到前行的写入结果。
func (f *Flow) applyPayload(n *Node, msg *Msg, in *inbound) {
if n.Payload == nil {
return
}
var byKey map[string]any
var ordered []any
if hasVar(n.Payload.Set) {
byKey, ordered = resolveInputs(n, in)
}
for _, item := range n.Payload.Set {
v := item.Value
if item.IsVar {
expr, _ := item.Value.(string)
v = resolveVar(scriptVars(msg, byKey, ordered), expr)
}
msg.setPath(item.KeyName, v)
}
for _, path := range n.Payload.Remove {
msg.removePath(path)
}
}
// hasVar 是否有需要求值的变量行(没有就不必构造变量表)。
func hasVar(set []SetItem) bool {
for _, item := range set {
if item.IsVar {
return true
}
}
return false
}
// nodeError 节点执行失败的统一处理:宿主接线错误直接冒泡;
// 否则走异常出口 errorIds,并上报 ErrorHandler(无论是否接了 errorIds 都上报,便于排障)。
// 下游收到的 msg 保留失败前的 payload 与 value,附带 error(错误信息)与 errorId(失败节点 id)。
func (f *Flow) nodeError(ctx *execContext, n *Node, msg *Msg, err error) error {
if errors.Is(err, ErrNoActionHandler) {
return err // 宿主未接线:直接冒泡,不当作节点执行失败
}
if ErrorHandler != nil {
ErrorHandler(n.ID, err)
}
// in 传 nil:失败时 resolveInputs 可能本身就是失败原因(参数非法),不重复求值。
f.trace(n, msg, nil, TracePhaseError, err, nil)
msg.fail(n.ID, err)
if ctx.watch != nil {
ctx.watch(n.ID, msg, err)
}
return f.fanoutExits(ctx, n, n.Outputs.ErrorIDs, msg)
}
// fanoutExits 沿给定出口把 msg 派发给各下游并触发执行。
// 出口为空即空转(断本分支)。多个下游时先标记 payload 共享——各分支的写入互不可见(COW)。
func (f *Flow) fanoutExits(ctx *execContext, n *Node, exits []Exit, msg *Msg) error {
if len(exits) == 0 {
return nil
}
if len(exits) > 1 {
msg.markShared()
}
for _, e := range exits {
down := f.node(e.ID)
if down == nil {
continue // 下游不存在则跳过,不中断(加载已校验,此处防御)
}
if err := f.execNode(ctx, down, f.deliver(n, down, e, msg)); err != nil {
return err
}
}
return nil
}
// deliver 构造下游一次执行的入参:msg 派生一份(payload 按 COW 共享),
// value 落到连线 key 指向的参数槽。
// - key 省略 → 落第一个参数;
// - key 匹配不到参数、或下游无参数 → 忽略数据,连线仍起触发作用。
// 但"触发"不等于"执行":下游有就绪门(pair/continuous)且某参数因此拿不到值、
// 又没有 value 缺省值时,门控不放行 → 不执行(加载时已由 warnDeadInputs 报告)。
// - 上游是 condition/switch(up.routesOnly) → 一律不落槽,连线的 key 被忽略。
// msg 仍原样流下(value/payload 都在),只是不写下游参数槽、不入缓存。
// 理由见 Node.routesOnly;warnDeadInputs 的"被连线指向"统计与此处逐条一致,
// 改一处必须改两处。
//
// 下游有输入参数时,同时把该参数值写入场景作用域缓存(跨触发存活)——三种 inputMode
// 都用缓存,所以这里不再按模式区分。nil 不入缓存,由 store.feed 拦掉。
func (f *Flow) deliver(up, down *Node, e Exit, msg *Msg) *inbound {
next := msg.fork()
if len(down.Inputs) > 0 {
msg.markShared() // 派生出去的那份与本分支共享 payload
next.markShared()
}
slot := -1
if !up.routesOnly() {
slot = down.slotIndex(e.Key)
}
slots := map[int]any{}
if slot >= 0 {
slots[slot] = msg.Value
if down.caches() {
f.store.feed(f.key(down), slot, msg.Value)
}
}
return newInbound(next, slots)
}
// execCondition 求值条件节点,返回布尔。
func (f *Flow) execCondition(n *Node, in *inbound) (bool, error) {
relation := n.ConfigString("relation")
byKey, ordered := resolveInputs(n, in)
switch relation {
case OpAnd:
for _, v := range ordered {
if !truthy(v) {
return false, nil
}
}
return len(ordered) > 0, nil
case OpOr:
for _, v := range ordered {
if truthy(v) {
return true, nil
}
}
return false, nil
case OpInRange, OpNotInRange:
in, err := evalRange(f.subject(n, in, ordered), n.ConfigString("range"))
if err != nil {
return false, err
}
if relation == OpNotInRange {
return !in, nil
}
return in, nil
case OpInSet, OpNotInSet:
ok, err := evalSet(f.subject(n, in, ordered), n.ConfigString("set"))
if err != nil {
return false, err
}
if relation == OpNotInSet {
return !ok, nil
}
return ok, nil
case OpCustom:
// 自定义条件:优先 subflow(子流编排),否则 script(goja 脚本)。取输出的布尔值。
var res any
var err error
if n.hasConfig("subflow") {
res, err = f.runSubflow(newExecContext(), n, in)
} else {
res, err = evalScript(n.ConfigString("script"), scriptVars(in.msg, byKey, ordered))
}
if err != nil {
return false, err
}
return truthy(res), nil
default:
// 二元比较:取前两个输入参数
var a, b any
if len(ordered) > 0 {
a = ordered[0]
} else {
a = f.subject(n, in, ordered)
}
if len(ordered) > 1 {
b = ordered[1]
}
return evalRelation(relation, a, b)
}
}
// subject 取"待判值":声明了输入参数时用第一个参数,未声明参数时用 msg.value。
// 供 condition 的区间/集合/二元比较与 switch 复用。
func (f *Flow) subject(n *Node, in *inbound, ordered []any) any {
if len(ordered) > 0 {
return ordered[0]
}
if in != nil && in.msg != nil {
return in.msg.Value
}
return nil
}
// ErrNoActionHandler 未注入 ruleflow.ActionHandler 就执行动作节点。
// 属于宿主接线错误(而非动作执行失败),不走异常出口、直接冒泡给调用方。
var ErrNoActionHandler = errors.New("动作需设置 ruleflow.ActionHandler")
// ErrorHandler 是全局节点异常回调:节点执行失败时上报(nodeID, err),供宿主记日志/告警。
// 节点失败只断本分支不冒泡,若不接此回调,脚本异常等问题会静默无痕。
// 调用方在启动场景前赋值:ruleflow.ErrorHandler = func(id string, err error) {...}。
var ErrorHandler func(nodeID string, err error)
// execAction 执行动作节点(调全局 ActionHandler)。
// 透传 msg 的各项给宿主:config / inputs / value / payload,外加 error / errorId
// (接在异常出口上的告警动作靠它们报告"哪一步失败了",正常分支恒为空串),
// 以及 sceneId / nodeId(动作来源标识)。
// 回调返回值非 nil 则写入 msg.value;回调里对 payload 的修改不回流——
// 要写流作用域用节点的 payload.set。
func (f *Flow) execAction(n *Node, in *inbound) (any, error) {
if ActionHandler == nil {
return nil, ErrNoActionHandler
}
byKey, _ := resolveInputs(n, in)
params := map[string]any{
"config": n.Config,
"inputs": byKey,
"value": in.msg.Value,
"payload": in.msg.Payload,
"error": in.msg.Error,
"errorId": in.msg.ErrorID,
// 动作来源:宿主记审计日志时要答"哪个场景的哪个节点改了设备"。
// 这两项不随 Debug 开关变化(msg.meta 只在 Debug 时才有 nodeId)——审计是常态需求。
// sceneId 子流继承宿主的;直连 newFlow() 的测试流为空串,属合法状态。
"sceneId": f.sceneID,
"nodeId": n.ID,
}
return ActionHandler(actionKind(n), params)
}
// execFunction 执行函数节点:config.fn 命中内置函数则调内置;
// 否则 subflow(子流) 优先于 script(goja 脚本)。
func (f *Flow) execFunction(ctx *execContext, n *Node, in *inbound) (any, error) {
byKey, ordered := resolveInputs(n, in)
if fn := n.ConfigString("fn"); fn != "" {
args := ordered
if len(args) == 0 && in.msg != nil && in.msg.Value != nil {
args = []any{in.msg.Value} // 无声明参数:用流入的 value 作单参数
}
return CallFunction(fn, args...)
}
if n.hasConfig("subflow") {
return f.runSubflow(ctx, n, in)
}
return evalScript(n.ConfigString("script"), scriptVars(in.msg, byKey, ordered))
}
// runSubflow 把节点 config.subflow 作为子流运行(block 与 subflow 自定义条件/函数/动作共用)。
// 子流结构只编译一次并缓存在节点上;子流共享宿主场景的 store(靠 scope 前缀隔离实例)。
//
// 数据边界:
// - 输入:子流内节点按 key 直接引用宿主的输入参数(无需连线);
// 子流节点自己声明了同名参数时以子流自己的为准。
// - payload:子流继承宿主的 payload(同样 COW)。
// - 入口:子流中没有入边的节点为入口,有多个则全部执行。
// - 输出:由 outputs.value(子流节点 id)指定,取该节点的 msg 作为本节点输出。
// 指定节点未执行 → value 为 nil、payload 取宿主流入那份;执行多次 → 取最后一次;
// 省略 outputs.value → value 为 nil 且子流对 payload 的修改不外泄;
// 指定节点执行失败 → 本节点走 errorIds,保留子流内失败节点的 error/errorId。
func (f *Flow) runSubflow(_ *execContext, n *Node, in *inbound) (any, error) {
if !n.hasConfig("subflow") {
return nil, nil
}
sub, err := f.subflowOf(n)
if err != nil {
return nil, err
}
byKey, _ := resolveInputs(n, in)
want := n.Outputs.Value
var picked *Msg // outputs.value 指定的子流节点最后一次执行后的 msg
var subErr error
// 子流用独立的传播上下文:子流内的出口只在子流内寻址,环检测也自成一路。
// host 带上宿主的输入参数,子流内任意节点都能按 key 引用。
subCtx := newExecContext()
subCtx.host = byKey
if want != "" {
subCtx.watch = func(id string, m *Msg, e error) {
if id == want {
picked, subErr = m, e // 执行多次则取最后一次
}
}
}
// 子流内每个入口都以宿主的 msg(派生一份)启动。先标记共享:子流内的写入按 COW
// 落到自己的副本上——多入口之间互不可见,且省略 outputs.value 时对宿主也不外泄。
// 需要外泄时由下面把 picked 的 payload 接回宿主。
entries := sub.entryNodes()
hostMsg := in.msg
hostMsg.markShared()
for _, e := range entries {
entryMsg := hostMsg.fork()
entryMsg.markShared()
if err := sub.execNode(subCtx, e, newInbound(entryMsg, nil)); err != nil {
return nil, err
}
}
if want == "" {
return nil, nil // 无输出:子流对 payload 的修改不外泄
}
if subErr != nil {
return nil, subErr // 指定节点执行失败:本节点走异常出口
}
if picked == nil {
return nil, nil // 指定节点本次没执行:value 为 nil,payload 保持宿主流入那份
}
// 取该节点的 msg 作为本节点输出:payload 外泄回宿主,value 作为节点输出值
in.msg.Payload = picked.Payload
in.msg.shared = picked.shared
return picked.Value, nil
}
// subflowOf 取(或首次编译)节点的子流。子流共享宿主的 store,scope 加宿主节点 id 前缀。
func (f *Flow) subflowOf(n *Node) (*Flow, error) {
n.subOnce.Do(func() {
bs, err := marshal(n.Config["subflow"])
if err != nil {
n.subErr = err
return
}
sub := newFlow()
sub.store = f.store
sub.scope = f.scope + n.ID + "/"
sub.sceneID = f.sceneID // 追踪事件要报到宿主场景名下
if err := sub.load(bs); err != nil {
n.subErr = err
return
}
sub.mu.Lock()
sub.lifeCtx = f.lifeCtx
sub.mu.Unlock()
n.subFlow = sub
})
if n.subErr != nil {
return nil, n.subErr
}
return n.subFlow, nil
}
// entryNodes 子流入口:没有任何入边的节点(按加载顺序)。有多个则全部执行。
func (f *Flow) entryNodes() []*Node {
f.mu.RLock()
defer f.mu.RUnlock()
hasIn := map[string]bool{}
for _, id := range f.order {
for _, e := range f.nodes[id].allExits() {
hasIn[e.ID] = true
}
}
var entries []*Node
for _, id := range f.order {
if !hasIn[id] {
entries = append(entries, f.nodes[id])
}
}
return entries
}
// execGate 处理 debounce/throttle 阀门节点:控制下游触发时序,本身不产出值。
// 缓存并放行**整条 msg**——被丢弃的那几次连 payload 一起丢。
// - throttle:距上次放行 ≥ interval(毫秒) 才放行本次,否则丢弃。
// - debounce:每次输入重置 delay(毫秒)定时器,到点用最后一次输入触发下游;
// 配 maxWait 时从本轮首次输入起超过 maxWait 强制触发一次。
//
// 放行是**一次新的传播**(脱离原触发的调用栈),环检测路径从阀门重新计起。
func (f *Flow) execGate(n *Node, in *inbound) error {
msg := in.msg
if n.ConfigString("fn") == "throttle" {
interval := millis(n.Config["interval"])
n.gateMu.Lock()
now := time.Now()
if n.lastFire.IsZero() || now.Sub(n.lastFire) >= interval {
n.lastFire = now
n.gateMu.Unlock()
return f.release(n, msg, in) // 放行
}
n.gateMu.Unlock()
// 限流窗口内丢弃:用 blocked 表示"本次被阀门挡住、不往下走",
// 与多输入未凑齐(execNode 里的 blocked)同属"本次不执行"。
f.trace(n, msg, in, TracePhaseBlocked, nil, nil)
return nil // 限流窗口内:丢弃
}
// debounce
delay := millis(n.Config["delay"])
if delay <= 0 {
delay = 300 * time.Millisecond
}
maxWait := millis(n.Config["maxWait"]) // 0 = 不限
f.mu.RLock()
lifeCtx := f.lifeCtx
f.mu.RUnlock()
if lifeCtx == nil {
lifeCtx = context.Background()
}
n.gateMu.Lock()
// 本次覆盖前,若本轮已有一条等待中的输入(debFirst 已置位而未被消费),
// 它会被下面的赋值取代——不论接下来走 maxWait 强制触发还是重置定时器,
// 这条旧输入都不会再单独触发下游,属于被阀门丢弃。
discardedMsg, discardedIn := n.debMsg, n.debIn
discarded := !n.debFirst.IsZero()
n.debMsg, n.debIn = msg, in
if n.debFirst.IsZero() {
n.debFirst = time.Now()
}
// maxWait 到点:强制立即触发一次
if maxWait > 0 && time.Since(n.debFirst) >= maxWait {
if n.debTimer != nil {
n.debTimer.Stop()
n.debTimer = nil
}
fire, fireIn := n.debMsg, n.debIn
n.debFirst = time.Time{}
n.gateMu.Unlock()
if discarded && discardedMsg != nil {
// 被新消息取代而丢弃:与上面 throttle 分支同理,用 blocked 表示本次不执行。
f.trace(n, discardedMsg, discardedIn, TracePhaseBlocked, nil, nil)
}
return f.release(n, fire, fireIn)
}
// 重置 delay 定时器
if n.debTimer != nil {
n.debTimer.Stop()
}
n.debTimer = time.AfterFunc(delay, func() {
select {
case <-lifeCtx.Done():
return // 场景已停止,不触发
default:
}
n.gateMu.Lock()
fire, fireIn := n.debMsg, n.debIn
n.debTimer = nil
n.debFirst = time.Time{}
n.gateMu.Unlock()
_ = f.release(n, fire, fireIn)
})
n.gateMu.Unlock()
if discarded && discardedMsg != nil {
// 被新消息取代而丢弃:与上面 throttle 分支同理,用 blocked 表示本次不执行。
f.trace(n, discardedMsg, discardedIn, TracePhaseBlocked, nil, nil)
}
return nil
}
// release 阀门放行:以新的传播上下文把缓存的整条 msg 原样发给下游。
// in 为放行的那次输入(阀门不产值,其 payload.set 里的 value 仍是流入值)。
func (f *Flow) release(n *Node, msg *Msg, in *inbound) error {
if msg == nil {
return nil
}
ctx := newExecContext()
ctx.push(n.ID)
defer ctx.pop(n.ID)
f.applyPayload(n, msg, in)
// 补一条 done 事件:与 finish 里那条同构,理由见那里(上报的 payload 是本节点写完的
// 最终值;放在 fanoutExits 之前使事件顺序即执行顺序)。release 绕过 finish,
// 若不补此处,阀门节点会在调试窗口里永远空白,而其下游正常出事件。
f.trace(n, msg, in, TracePhaseDone, nil, nil)
return f.fanoutExits(ctx, n, n.Outputs.IDs, msg)
}
// millis 把 config 值(数字/数字字符串,单位毫秒)转为 time.Duration。
func millis(v any) time.Duration {
switch x := v.(type) {
case float64:
return time.Duration(x) * time.Millisecond
case int:
return time.Duration(x) * time.Millisecond
case int64:
return time.Duration(x) * time.Millisecond
case string:
if n, err := strconv.Atoi(x); err == nil {
return time.Duration(n) * time.Millisecond
}
}
return 0
}