-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmsg.go
More file actions
196 lines (179 loc) · 6.63 KB
/
Copy pathmsg.go
File metadata and controls
196 lines (179 loc) · 6.63 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
package ruleflow
import (
"strconv"
"time"
)
// Debug 调试模式。为 true 时每条 msg 额外带 Meta(nodeId/nodeType/timestamp/traceId),
// 供宿主串联日志;为 false 时 msg 里没有 meta,脚本不应依赖它。
// 在 Start 之前设置。
var Debug = false
// Msg 是节点之间流转的数据载体。一次触发新建一条,沿规则链传递、按分支派生。
//
// - Payload:流作用域。一次触发内贯穿全链路的业务数据,任何节点都能读,
// 与输入参数无关;分支之间写隔离(Copy-On-Write,见 fork/writable)。
// - Value:传递作用域。上游给下游的临时值(上游节点的计算结果),
// 落到下游哪个输入参数由连线的 input 决定(见 flow.go deliver)。
// - Error/ErrorID:仅异常出口非空——错误信息与失败节点 id。
// 节点成功执行后清空,异常出口是它们唯一的来源。
// - Meta:仅 Debug 模式存在。
type Msg struct {
Payload map[string]any `json:"payload"`
Value any `json:"value"`
Error string `json:"error"`
ErrorID string `json:"errorId"`
Meta map[string]any `json:"meta,omitempty"`
// shared 标记 Payload 这棵树可能被多个分支共享:写入前须先拷贝一份(COW)。
// 扇出到多个下游时置位;本分支首次写入时整棵拷贝并清零,之后连续写入不再拷贝。
//
// 为何首次写入拷整棵而不只拷写入路径:只拷路径的话,同一分支第二次写另一条
// 嵌套路径时会直接改到仍与兄弟分支共享的那层 map,隔离就漏了。
// 拷贝只发生在"本分支第一次写"这一刻,纯传递与后续写入都不拷贝。
shared bool
}
// newMsg 新建一条 msg(一次触发的起点):payload 恒为空对象,value 为触发数据。
func newMsg(value any) *Msg {
m := &Msg{Payload: map[string]any{}, Value: value}
if Debug {
m.Meta = map[string]any{
"traceId": newTraceID(),
"timestamp": time.Now().UnixMilli(),
}
}
return m
}
// fork 派生一条交给下游的 msg:结构浅拷贝,Payload 仍指向同一张 map。
// 真正的隔离由 shared 标记 + writable 的按需拷贝完成,所以"分支多"不等于"拷贝多"。
func (m *Msg) fork() *Msg {
if m == nil {
return newMsg(nil)
}
c := *m
if m.Meta != nil {
meta := make(map[string]any, len(m.Meta))
for k, v := range m.Meta {
meta[k] = v
}
c.Meta = meta
}
return &c
}
// markShared 标记 payload 进入共享状态:扇出到多个下游前调用。
func (m *Msg) markShared() {
if m != nil {
m.shared = true
}
}
// writable 取可写的 payload:共享状态下先整体拷贝一份并清除共享标记,
// 之后本分支的连续写入不再拷贝。只写才拷贝,只传不拷贝。
func (m *Msg) writable() map[string]any {
if m.Payload == nil {
m.Payload = map[string]any{}
}
if m.shared {
m.Payload = copyMap(m.Payload)
m.shared = false
}
return m.Payload
}
// traceMeta 给本次触发的 msg 盖上追踪标记。是否追踪由触发入口一次性决定,
// 之后整条链路(含子流——Meta 随 fork 深拷)只看 Meta 是否存在,不再查全局状态。
//
// id 非空时用调用方指定的值,否则自行生成。指定的意义在于调用方能**在触发之前**
// 就知道这次的 traceId:手动触发是同步的(Trigger 跑完整条链才返回),而追踪事件是
// 执行过程中发出去的,等返回值才拿到 id 就只能先缓冲事件再回溯过滤。
// 指定 id 只决定"这次追踪叫什么名字",不决定"要不要追"——后者仍归 TraceGate。
func (m *Msg) traceMeta(id string) {
if m == nil {
return
}
if id == "" {
id = newTraceID()
}
m.Meta = map[string]any{
"traceId": id,
"timestamp": time.Now().UnixMilli(),
}
}
// stamp 记录本次流经的节点。只对带 Meta 的 msg 生效——
// 即 Debug 模式,或本次触发被追踪(见 Flow.traced)。是否追踪由触发入口一次性决定,
// 节点层不再查全局状态,所以这里也不再看 Debug。
func (m *Msg) stamp(n *Node) {
if m == nil || m.Meta == nil || n == nil {
return
}
m.Meta["nodeId"] = n.ID
m.Meta["nodeType"] = n.Type
m.Meta["timestamp"] = time.Now().UnixMilli()
}
// fail 置异常信息:保留失败前的 payload 与 value,供异常出口的下游读取。
func (m *Msg) fail(nodeID string, err error) {
if m == nil || err == nil {
return
}
m.Error = err.Error()
m.ErrorID = nodeID
}
// clearError 清空异常信息(节点成功执行后调用)。
func (m *Msg) clearError() {
if m != nil {
m.Error = ""
m.ErrorID = ""
}
}
// setPath 按点路径写入 payload("a.b.c" 表示嵌套),规则见 path.go。
// 路径为空时不写(与 removePath 一致)——空路径的含义是"没指定写哪里",
// 不是"目标是整个 payload"。
//
// 曾经空路径是整包替换,为此埋掉过一个场景:可视化编辑器在 payload.set 里
// 留下未填写的空行({keyName:"",value:"",isVar:false}),加载校验放行(它是合法的
// 画图中间态,见 flow.go validate),运行期把整个 payload 换成 {"value":""},
// 上游各节点写入的业务数据全丢且无痕。破坏性最强的操作不该绑在最"空"的那一档输入上——
// 与 path.go 的一贯态度(读不到即 null、写入越界静默忽略)也不一致。
//
// 整包替换没有留替代拼写:{"keyName":"report","value":"value","isVar":true}
// 加下游读 payload.report 完全等价,只多一层嵌套且意图明确。
func (m *Msg) setPath(path string, value any) {
if path == "" {
return
}
setInPath(m.writable(), path, value) // writable 已保证本分支独占,逐层直接写即可
}
// removePath 按点路径删除 payload 里的字段;路径不存在时静默忽略。
func (m *Msg) removePath(path string) {
if path == "" || m.Payload == nil {
return
}
removeInPath(m.writable(), path)
}
// copyMap 深拷贝 map(嵌套 map/slice 逐层拷贝,标量原样)。
// 仅在 COW 首次写入时调用一次。
func copyMap(src map[string]any) map[string]any {
dst := make(map[string]any, len(src))
for k, v := range src {
dst[k] = copyValue(v)
}
return dst
}
func copyValue(v any) any {
switch x := v.(type) {
case map[string]any:
return copyMap(x)
case []any:
arr := make([]any, len(x))
for i, e := range x {
arr[i] = copyValue(e)
}
return arr
default:
return v
}
}
// newTraceID 生成一次触发的唯一 id(串联调试日志)。uuid7 不可用时退化为时间戳。
func newTraceID() string {
if v, err := CallFunction("uuid7"); err == nil {
if s, ok := v.(string); ok {
return s
}
}
return strconv.FormatInt(time.Now().UnixNano(), 36)
}