nx_wflow_engine_completion_t31.nx source
↩ module page · 1620 lines · 68252 B
1// nx_wflow_engine.nx -- R1 KEYSTONE of the Workflow Automation ladder (/compare/automation census row
2// "Unified multi-step durable workflow engine"). A WORKFLOW = DATA (law: data as data):
3// flow header = `flowid|event-kind|cond(k=v or -)`
4// step row = `flowid|idx|action|channel-or-dash|arg|max-attempts` (idx 1-based, author ascending)
5// An EVENT = `kind~k=v~...~id=EN` (the nx_crm_flow contract). The engine upgrades nx_crm_flow from
6// single-shot in-memory rules to DURABLE EXECUTION (the Temporal model, sovereign):
7// - every transition is ONE appended line in an append-only RUN LEDGER (flock + O_APPEND, torn-free):
8// `WFRUN rid=<evid>.<flowid> flow=<f> step=<n> status=START|ATT|OK|FAILSTEP|FAILED|DONE att=<k>`
9// - state is DERIVED BY REPLAY of the ledger (event-sourced; nothing mutates, history is sacred)
10// - IDEMPOTENT by construction: a (event id, flow) pair with a START line never starts twice --
11// the dedup survives crashes/restarts (stronger than the in-memory seen[] of nx_crm_flow)
12// - per-step RETRY up to max-attempts; exhausted -> run FAILED at that step, later steps never run
13// - wf_resume: any run with START but neither DONE nor FAILED (a crash mid-flight) is completed
14// from its first non-OK step; already-OK steps are NOT re-executed (proven in the selftest)
15// ACTIONS v1 (in-process, composing the send plane like nx_crm_flow): create-task | send (any nx_send
16// channel, sr_valid-checked at load) | notify | update-field | probe-fail (DIAGNOSTIC: fails while
17// attempt < arg -- the deterministic retry witness; keep it out of production flows).
18// wf_load validates the WHOLE step set up front (unknown action / bad channel / attempts<1 = LOUD -1,
19// nothing half-loaded). Production ledger path suggestion: knowledge/status/wflow_runs.log.
20// R3 CONNECTOR LAYER (exec-organ): a step can run ANY blessed organ as a workflow step, resolved through a
21// GREEN fail-closed catalog (rows `name<TAB>elf<TAB>GREEN`, the tool_allowlist.conf shape; cx[3]=catalog path,
22// 0=exec-organ refused). Child exit 0 = step OK, nonzero = step fail (so per-step retry applies to REAL organ
23// execution). Production catalog: knowledge/wflow/connectors.conf.
24// PURE CORE (no main -- the house pure-core+gate idiom): gates = _hdl_build/nx_wflow_engine_gate (wf_selftest)
25// + _hdl_build/nx_wflow_connect_gate (connector battery). license_tier: ORIGINAL
26import "nx_send.nx"
27const WF_MAGIC_1089: i64 = 1089
28const WF_MAGIC_65536: i64 = 65536
29const WF_MAGIC_16777216: i64 = 16777216
30const WF_MAGIC_1024: i64 = 1024
31const WF_MAGIC_60000: i64 = 60000
32
33// ---- tilde-event helpers (nx_crm_flow idiom, local wf_ copies) ----
34func wf_has(hay: *u8, needle: *u8) -> i64 {
35 let n: i64 = slen(hay)
36 let m: i64 = slen(needle)
37 if m == 0 { return 0 }
38 var i: i64 = 0
39 while i + m <= n {
40 var k: i64 = 0
41 var hit: i64 = 1
42 while k < m { if hay[i+k] != needle[k] { hit = 0; k = m } else { k = k + 1 } }
43 if hit == 1 { return 1 }
44 i = i + 1
45 }
46 return 0
47}
48func wf_evid(ev: *u8, out: *u8, cap: i64) -> i64 {
49 let n: i64 = slen(ev)
50 var i: i64 = 0
51 while i + 3 < n {
52 if ev[i] == (126 as u8) { if ev[i+1] == (105 as u8) { if ev[i+2] == (100 as u8) { if ev[i+3] == (61 as u8) {
53 var q: i64 = i + 4
54 var t: i64 = 0
55 while q < n { if ev[q] == (126 as u8) { break } if t < cap - 1 { out[t] = ev[q]; t = t + 1 } q = q + 1 }
56 out[t] = 0 as u8
57 return 1
58 } } } }
59 i = i + 1
60 }
61 out[0] = 0 as u8
62 return 0
63}
64
65// ---- small locals ----
66func wf_cat(dst: *u8, off: i64, s: *u8) -> i64 { var x: i64 = off; var i: i64 = 0; while s[i] != (0 as u8) { dst[x] = s[i]; x = x + 1; i = i + 1 } return x }
67func wf_catn(dst: *u8, off: i64, v: i64) -> i64 {
68 var x: i64 = off
69 var m: i64 = v
70 if m < 0 { dst[x] = 45 as u8; x = x + 1; m = 0 - m }
71 if m == 0 { dst[x] = 48 as u8; return x + 1 }
72 let t: *u8 = sys_mmap(24)
73 var k: i64 = 0
74 while m > 0 { t[k] = (48 + (m - (m/10)*10)) as u8; m = m / 10; k = k + 1 }
75 var j: i64 = 0
76 while j < k { dst[x] = t[k-1-j]; x = x + 1; j = j + 1 }
77 return x
78}
79func wf_atoi(s: *u8) -> i64 {
80 var v: i64 = 0
81 var i: i64 = 0
82 var any: i64 = 0
83 while s[i] != (0 as u8) {
84 let c: i64 = s[i] as i64
85 if c >= 48 { if c <= 57 { v = v*10 + (c - 48); any = 1 } }
86 if c < 48 { return 0 - 1 }
87 if c > 57 { return 0 - 1 }
88 i = i + 1
89 }
90 if any == 0 { return 0 - 1 }
91 return v
92}
93// raw __syscall(201) returns a broken constant on this backend (the getpid -25 class of quirk) --
94// use the proven clock_gettime wrapper so concurrent selftests never share a ledger path
95func wf_now() -> i64 { return sys_now_realtime_sec() }
96
97// ---- durable append-only ledger (flock + O_APPEND; one line per transition) ----
98func wf_append(path: *u8, line: *u8) -> i64 {
99 // openat(AT_FDCWD, path, O_WRONLY|O_CREAT|O_APPEND, 0644)
100 let fd: i64 = __syscall(257, 0 - 100, path as i64, WF_MAGIC_1089, 420, 0, 0)
101 if fd < 0 { p("WFLOW append ERROR cannot open ledger -- fail loud\n" as *u8); return 0 - 1 }
102 __syscall(73, fd, 2, 0, 0, 0, 0)
103 let n: i64 = slen(line)
104 let w: i64 = sys_write(fd, line, n)
105 __syscall(73, fd, 8, 0, 0, 0, 0)
106 __syscall(3, fd, 0, 0, 0, 0, 0)
107 if w != n { p("WFLOW append ERROR short write -- fail loud\n" as *u8); return 0 - 1 }
108 return 0
109}
110const WF_LEDCAP: i64 = 1048576
111func wf_readall(path: *u8, szp: *i64) -> *u8 {
112 let buf: *u8 = sys_mmap(WF_LEDCAP)
113 szp[0] = 0
114 let fd: i64 = __syscall(257, 0 - 100, path as i64, 0, 0, 0, 0)
115 if fd < 0 { return buf }
116 var off: i64 = 0
117 var go: i64 = 1
118 while go == 1 {
119 let r: i64 = __syscall(0, fd, (buf as i64) + off, WF_LEDCAP - off, 0, 0, 0)
120 if r <= 0 { go = 0 }
121 if go == 1 { off = off + r; if off >= WF_LEDCAP { __syscall(3, fd, 0, 0, 0, 0, 0); p("WFLOW ledger ERROR over cap -- fail loud\n" as *u8); return 0 as *u8 } }
122 }
123 __syscall(3, fd, 0, 0, 0, 0, 0)
124 szp[0] = off
125 return buf
126}
127func wf_writeall(path: *u8, buf: *u8, n: i64) -> i64 {
128 let fd: i64 = __syscall(257, 0 - 100, path as i64, 577, 420, 0, 0)
129 if fd < 0 { p("WFLOW write ERROR cannot open out -- fail loud\n" as *u8); return 0 - 1 }
130 let w: i64 = sys_write(fd, buf, n)
131 __syscall(3, fd, 0, 0, 0, 0, 0)
132 if w != n { p("WFLOW write ERROR short write -- fail loud\n" as *u8); return 0 - 1 }
133 return 0
134}
135
136// ---- ledger replay queries ----
137func wf_has_rng(buf: *u8, ls: i64, le: i64, needle: *u8) -> i64 {
138 let m: i64 = slen(needle)
139 if m == 0 { return 0 }
140 var i: i64 = ls
141 while i + m <= le {
142 var k: i64 = 0
143 var hit: i64 = 1
144 while k < m { if buf[i+k] != needle[k] { hit = 0; k = m } else { k = k + 1 } }
145 if hit == 1 { return 1 }
146 i = i + 1
147 }
148 return 0
149}
150// count ledger LINES containing BOTH substrings
151func wf_lines_with2(buf: *u8, n: i64, a: *u8, b: *u8) -> i64 {
152 var c: i64 = 0
153 var ls: i64 = 0
154 while ls < n {
155 var le: i64 = ls
156 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
157 if wf_has_rng(buf, ls, le, a) == 1 { if wf_has_rng(buf, ls, le, b) == 1 { c = c + 1 } }
158 ls = le + 1
159 }
160 return c
161}
162// parse the number after key within [ls,le); -1 if absent
163func wf_num_after(buf: *u8, ls: i64, le: i64, key: *u8) -> i64 {
164 let m: i64 = slen(key)
165 var i: i64 = ls
166 while i + m <= le {
167 var k: i64 = 0
168 var hit: i64 = 1
169 while k < m { if buf[i+k] != key[k] { hit = 0; k = m } else { k = k + 1 } }
170 if hit == 1 {
171 var q: i64 = i + m
172 var v: i64 = 0
173 var any: i64 = 0
174 while q < le { let c: i64 = buf[q] as i64; if c >= 48 { if c <= 57 { v = v*10 + (c - 48); any = 1; q = q + 1 } else { q = le } } else { q = le } }
175 if any == 1 { return v }
176 return 0 - 1
177 }
178 i = i + 1
179 }
180 return 0 - 1
181}
182// copy the space-terminated value after key within [ls,le) into out; 1 if found
183func wf_val_after(buf: *u8, ls: i64, le: i64, key: *u8, out: *u8, cap: i64) -> i64 {
184 let m: i64 = slen(key)
185 var i: i64 = ls
186 while i + m <= le {
187 var k: i64 = 0
188 var hit: i64 = 1
189 while k < m { if buf[i+k] != key[k] { hit = 0; k = m } else { k = k + 1 } }
190 if hit == 1 {
191 var q: i64 = i + m
192 var t: i64 = 0
193 while q < le { if buf[q] == (32 as u8) { break } if buf[q] == (10 as u8) { break } if t < cap - 1 { out[t] = buf[q]; t = t + 1 } q = q + 1 }
194 out[t] = 0 as u8
195 return 1
196 }
197 i = i + 1
198 }
199 out[0] = 0 as u8
200 return 0
201}
202// highest step idx with status=OK for this run token
203func wf_max_ok_step(buf: *u8, n: i64, ridtok: *u8) -> i64 {
204 var mx: i64 = 0
205 var ls: i64 = 0
206 while ls < n {
207 var le: i64 = ls
208 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
209 if wf_has_rng(buf, ls, le, ridtok) == 1 { if wf_has_rng(buf, ls, le, "status=OK" as *u8) == 1 {
210 let v: i64 = wf_num_after(buf, ls, le, " step=" as *u8)
211 if v > mx { mx = v }
212 } }
213 ls = le + 1
214 }
215 return mx
216}
217
218// ---- definition validation (whole-set, loud, nothing half-loaded) ----
219func wf_action_ok(a: *u8) -> i64 {
220 if seq(a, "create-task" as *u8) == 1 { return 1 }
221 if seq(a, "send" as *u8) == 1 { return 1 }
222 if seq(a, "notify" as *u8) == 1 { return 1 }
223 if seq(a, "update-field" as *u8) == 1 { return 1 }
224 if seq(a, "probe-fail" as *u8) == 1 { return 1 }
225 if seq(a, "exec-organ" as *u8) == 1 { return 1 }
226 if seq(a, "approve" as *u8) == 1 { return 1 }
227 if seq(a, "set-var" as *u8) == 1 { return 1 }
228 if seq(a, "for-each" as *u8) == 1 { return 1 }
229 return 0
230}
231func wf_load(steps: *i64, n: i64) -> i64 {
232 let a: *u8 = sys_mmap(64)
233 let ch: *u8 = sys_mmap(64)
234 let ix: *u8 = sys_mmap(32)
235 let ma: *u8 = sys_mmap(32)
236 var i: i64 = 0
237 while i < n {
238 let r: *u8 = steps[i] as *u8
239 pipe_field(r, 1, ix, 32)
240 if wf_atoi(ix) < 1 { p("WFLOW load ERROR bad step idx -- fail loud\n" as *u8); return 0 - 1 }
241 pipe_field(r, 2, a, 64)
242 if wf_action_ok(a) == 0 { p("WFLOW load ERROR unknown action=" as *u8); p(a); p(" -- fail loud\n" as *u8); return 0 - 1 }
243 if seq(a, "send" as *u8) == 1 {
244 pipe_field(r, 3, ch, 64)
245 if sr_valid(ch) == 0 { p("WFLOW load ERROR bad send channel=" as *u8); p(ch); p(" -- fail loud\n" as *u8); return 0 - 1 }
246 }
247 if seq(a, "exec-organ" as *u8) == 1 {
248 pipe_field(r, 3, ch, 64)
249 var chok: i64 = 1
250 if ch[0] == (0 as u8) { chok = 0 }
251 if ch[0] == (45 as u8) { if ch[1] == (0 as u8) { chok = 0 } }
252 if ch[0] == (62 as u8) { chok = 0 }
253 var gsep: i64 = 0
254 var gi2: i64 = 0
255 while ch[gi2] != (0 as u8) { if ch[gi2] == (62 as u8) { gsep = gi2 } gi2 = gi2 + 1 }
256 if gsep > 0 { if ch[gsep+1] == (0 as u8) { chok = 0 } }
257 if chok == 0 { p("WFLOW load ERROR exec-organ needs connector or connector>var -- fail loud\n" as *u8); return 0 - 1 }
258 }
259 if seq(a, "approve" as *u8) == 1 {
260 pipe_field(r, 3, ch, 64)
261 var apok: i64 = 1
262 if ch[0] == (0 as u8) { apok = 0 }
263 if ch[0] == (45 as u8) { if ch[1] == (0 as u8) { apok = 0 } }
264 if apok == 0 { p("WFLOW load ERROR approve needs an approver name -- fail loud\n" as *u8); return 0 - 1 }
265 }
266 if seq(a, "set-var" as *u8) == 1 {
267 let ag: *u8 = sys_mmap(512)
268 pipe_field(r, 4, ag, 512)
269 var haseq: i64 = 0
270 var gi: i64 = 0
271 while ag[gi] != (0 as u8) { if ag[gi] == (61 as u8) { haseq = 1 } gi = gi + 1 }
272 if haseq == 0 { p("WFLOW load ERROR set-var needs key=value -- fail loud\n" as *u8); return 0 - 1 }
273 }
274 if seq(a, "for-each" as *u8) == 1 {
275 pipe_field(r, 3, ch, 64)
276 var feok: i64 = 1
277 if ch[0] == (0 as u8) { feok = 0 }
278 if ch[0] == (45 as u8) { if ch[1] == (0 as u8) { feok = 0 } }
279 if ch[0] == (62 as u8) { feok = 0 }
280 let ag2: *u8 = sys_mmap(512)
281 pipe_field(r, 4, ag2, 512)
282 if ag2[0] == (0 as u8) { feok = 0 }
283 if feok == 0 { p("WFLOW load ERROR for-each needs connector and a list arg -- fail loud\n" as *u8); return 0 - 1 }
284 }
285 pipe_field(r, 5, ma, 32)
286 if wf_atoi(ma) < 1 { p("WFLOW load ERROR max-attempts under 1 -- fail loud\n" as *u8); return 0 - 1 }
287 // R9: optional 7th field = per-step condition, must be k=v or dash
288 let cnd: *u8 = sys_mmap(256)
289 pipe_field(r, 6, cnd, 256)
290 var cok: i64 = 1
291 if cnd[0] != (0 as u8) {
292 if cnd[0] == (45 as u8) { if cnd[1] == (0 as u8) { cok = 1 } else { cok = 0 } } else { cok = 0 }
293 if cok == 0 {
294 var hq: i64 = 0
295 var ci: i64 = 0
296 while cnd[ci] != (0 as u8) { if cnd[ci] == (61 as u8) { hq = 1 } ci = ci + 1 }
297 if hq == 1 { cok = 1 }
298 }
299 }
300 if cok == 0 { p("WFLOW load ERROR step condition must be k=v or dash -- fail loud\n" as *u8); return 0 - 1 }
301 // R10: optional 8th field = on-fail target step idx (forward-only, > own idx)
302 let onf: *u8 = sys_mmap(32)
303 pipe_field(r, 7, onf, 32)
304 var ofok: i64 = 1
305 if onf[0] != (0 as u8) {
306 var isdash: i64 = 0
307 if onf[0] == (45 as u8) { if onf[1] == (0 as u8) { isdash = 1 } }
308 if isdash == 0 {
309 let tj: i64 = wf_atoi(onf)
310 if tj < 1 { ofok = 0 }
311 pipe_field(r, 1, ix, 32)
312 if tj <= wf_atoi(ix) { ofok = 0 }
313 }
314 }
315 if ofok == 0 { p("WFLOW load ERROR on-fail target must be a LATER step idx or dash -- fail loud\n" as *u8); return 0 - 1 }
316 i = i + 1
317 }
318 return n
319}
320func wf_load_flows(flows: *i64, n: i64) -> i64 {
321 let k: *u8 = sys_mmap(128)
322 var i: i64 = 0
323 while i < n {
324 let r: *u8 = flows[i] as *u8
325 pipe_field(r, 1, k, 128)
326 if k[0] == (0 as u8) { p("WFLOW load ERROR flow without event-kind -- fail loud\n" as *u8); return 0 - 1 }
327 i = i + 1
328 }
329 return n
330}
331// does flow header match event? (kind equal + cond k=v present, or cond '-')
332func wf_match(fhdr: *u8, ev: *u8) -> i64 {
333 let fk: *u8 = sys_mmap(128)
334 let ek: *u8 = sys_mmap(128)
335 pipe_field(fhdr, 1, fk, 128)
336 var i: i64 = 0
337 var t: i64 = 0
338 var go: i64 = 1
339 while go == 1 {
340 let c: i64 = ev[i] as i64
341 if c == 0 { go = 0 }
342 if c == 126 { go = 0 }
343 if go == 1 { if t < 127 { ek[t] = ev[i]; t = t + 1 } i = i + 1 }
344 }
345 ek[t] = 0 as u8
346 if seq(fk, ek) == 0 { return 0 }
347 let cond: *u8 = sys_mmap(256)
348 pipe_field(fhdr, 2, cond, 256)
349 if cond[0] == (45 as u8) { return 1 }
350 if cond[0] == (0 as u8) { return 1 }
351 let nb: *u8 = sys_mmap(300)
352 nb[0] = 126 as u8
353 var q: i64 = 0
354 while cond[q] != (0 as u8) { nb[q+1] = cond[q]; q = q + 1 }
355 nb[q+1] = 0 as u8
356 return wf_has(ev, nb)
357}
358
359// ---- one durable transition line ----
360func wf_emit(led: *u8, rid: *u8, fid: *u8, step: i64, st: *u8, att: i64) -> i64 {
361 let ln: *u8 = sys_mmap(512)
362 var o: i64 = 0
363 o = wf_cat(ln, o, "WFRUN rid=" as *u8)
364 o = wf_cat(ln, o, rid)
365 o = wf_cat(ln, o, " flow=" as *u8)
366 o = wf_cat(ln, o, fid)
367 o = wf_cat(ln, o, " step=" as *u8)
368 o = wf_catn(ln, o, step)
369 o = wf_cat(ln, o, " status=" as *u8)
370 o = wf_cat(ln, o, st)
371 o = wf_cat(ln, o, " att=" as *u8)
372 o = wf_catn(ln, o, att)
373 ln[o] = 10 as u8
374 ln[o+1] = 0 as u8
375 return wf_append(led, ln)
376}
377
378// ---- R3 connector layer: GREEN fail-closed catalog, any blessed organ as a workflow step ----
379// catalog rows: name<TAB>elfpath<TAB>GREEN. No catalog / unknown name / flag not exactly GREEN = refuse.
380func wf_conn_resolve(cx: *i64, name: *u8, out: *u8) -> i64 {
381 let conf: *u8 = cx[3] as *u8
382 if (conf as i64) == 0 { p("WFLOW exec-organ ERROR no connector catalog loaded -- fail closed\n" as *u8); return 0 }
383 let szp: *i64 = sys_mmap(8) as *i64
384 let buf: *u8 = wf_readall(conf, szp)
385 if (buf as i64) == 0 { return 0 }
386 let n: i64 = szp[0]
387 if n == 0 { p("WFLOW exec-organ ERROR connector catalog missing or empty -- fail closed\n" as *u8); return 0 }
388 let nm: i64 = slen(name)
389 var ls: i64 = 0
390 while ls < n {
391 var le: i64 = ls
392 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
393 var t1: i64 = ls
394 while t1 < le { if buf[t1] == (9 as u8) { break } t1 = t1 + 1 }
395 var hit: i64 = 0
396 if t1 - ls == nm {
397 hit = 1
398 var k: i64 = 0
399 while k < nm { if buf[ls+k] != name[k] { hit = 0; k = nm } else { k = k + 1 } }
400 }
401 if hit == 1 {
402 var t2: i64 = t1 + 1
403 while t2 < le { if buf[t2] == (9 as u8) { break } t2 = t2 + 1 }
404 let fl0: i64 = t2 + 1
405 var okflag: i64 = 0
406 if le - fl0 == 5 {
407 okflag = 1
408 let g: *u8 = "GREEN" as *u8
409 var k2: i64 = 0
410 while k2 < 5 { if buf[fl0+k2] != g[k2] { okflag = 0; k2 = 5 } else { k2 = k2 + 1 } }
411 }
412 if okflag == 0 { p("WFLOW exec-organ ERROR connector=" as *u8); p(name); p(" not GREEN -- fail closed\n" as *u8); return 0 }
413 var q: i64 = t1 + 1
414 var o: i64 = 0
415 while q < t2 { if o < 500 { out[o] = buf[q]; o = o + 1 } q = q + 1 }
416 out[o] = 0 as u8
417 if o == 0 { return 0 }
418 return 1
419 }
420 ls = le + 1
421 }
422 p("WFLOW exec-organ ERROR connector unknown=" as *u8); p(name); p(" -- fail closed\n" as *u8)
423 return 0
424}
425func wf_rd32(b: *u8, off: i64) -> i64 {
426 return (b[off] as i64) + ((b[off+1] as i64) * 256) + ((b[off+2] as i64) * WF_MAGIC_65536) + ((b[off+3] as i64) * WF_MAGIC_16777216)
427}
428// run the resolved organ as the step: argv = arg split on ~ ('-' or empty = no args); exit 0 = step OK.
429// R8b: channel `connector>varname` CAPTURES child stdout (bounded 480B, newlines to spaces, loud over-cap)
430// into the run variable plane -- step OUTPUT becomes later steps' INPUT.
431func wf_exec_organ(cx: *i64, name0: *u8, arg: *u8, rid: *u8) -> i64 {
432 let name: *u8 = sys_mmap(128)
433 let capv: *u8 = sys_mmap(128)
434 var ni: i64 = 0
435 var ci: i64 = 0
436 var seen: i64 = 0
437 var i0: i64 = 0
438 while name0[i0] != (0 as u8) {
439 if name0[i0] == (62 as u8) { seen = 1 } else {
440 if seen == 0 { if ni < 126 { name[ni] = name0[i0]; ni = ni + 1 } } else { if ci < 126 { capv[ci] = name0[i0]; ci = ci + 1 } }
441 }
442 i0 = i0 + 1
443 }
444 name[ni] = 0 as u8
445 capv[ci] = 0 as u8
446 var docap: i64 = 0
447 if ci > 0 { docap = 1 }
448 let elf: *u8 = sys_mmap(512)
449 if wf_conn_resolve(cx, name, elf) == 0 { return 0 }
450 let av: *i64 = sys_mmap(8 * 16) as *i64
451 av[0] = elf as i64
452 var na: i64 = 1
453 var skip: i64 = 0
454 if arg[0] == (0 as u8) { skip = 1 }
455 if arg[0] == (45 as u8) { if arg[1] == (0 as u8) { skip = 1 } }
456 if skip == 0 {
457 let ab: *u8 = sys_mmap(WF_MAGIC_1024)
458 var i: i64 = 0
459 var over: i64 = 0
460 while arg[i] != (0 as u8) { if i < 1000 { ab[i] = arg[i] } else { over = 1 } i = i + 1 }
461 if over == 1 { p("WFLOW exec-organ ERROR argv too long -- fail closed\n" as *u8); return 0 }
462 ab[i] = 0 as u8
463 let alen: i64 = i
464 var st0: i64 = 0
465 var j: i64 = 0
466 while j <= alen {
467 var cut: i64 = 0
468 if j == alen { cut = 1 }
469 if cut == 0 { if ab[j] == (126 as u8) { cut = 1 } }
470 if cut == 1 {
471 ab[j] = 0 as u8
472 if na >= 14 { p("WFLOW exec-organ ERROR too many args -- fail closed\n" as *u8); return 0 }
473 av[na] = (ab as i64) + st0
474 na = na + 1
475 st0 = j + 1
476 }
477 j = j + 1
478 }
479 }
480 av[na] = 0
481 let fb: *u8 = sys_mmap(16)
482 if docap == 1 {
483 if __syscall(293, fb as i64, 0, 0, 0, 0, 0) != 0 { p("WFLOW capture ERROR pipe failed -- fail closed\n" as *u8); return 0 }
484 }
485 let pid: i64 = sys_fork()
486 if pid == 0 {
487 if docap == 1 {
488 let pw: i64 = wf_rd32(fb, 4)
489 sys_dup3(pw, 1, 0)
490 __syscall(3, wf_rd32(fb, 0), 0, 0, 0, 0, 0)
491 __syscall(3, pw, 0, 0, 0, 0, 0)
492 }
493 let envp: *i64 = sys_mmap(16) as *i64
494 envp[0] = "PATH=/usr/bin:/bin" as *u8 as i64
495 envp[1] = 0
496 sys_execve(elf, av, envp)
497 sys_exit(127)
498 }
499 if pid < 0 { p("WFLOW exec-organ ERROR fork failed -- fail closed\n" as *u8); return 0 }
500 let cbuf: *u8 = sys_mmap(WF_MAGIC_1024)
501 var total: i64 = 0
502 if docap == 1 {
503 let pr: i64 = wf_rd32(fb, 0)
504 __syscall(3, wf_rd32(fb, 4), 0, 0, 0, 0, 0)
505 var go: i64 = 1
506 while go == 1 {
507 let rr: i64 = __syscall(0, pr, (cbuf as i64) + total, 1023 - total, 0, 0, 0)
508 if rr <= 0 { go = 0 }
509 if go == 1 { total = total + rr; if total >= 1023 { go = 0 } }
510 }
511 __syscall(3, pr, 0, 0, 0, 0, 0)
512 }
513 let st: *i64 = sys_mmap(16) as *i64
514 var waited: i64 = sys_wait4(pid, st, 0)
515 while waited == (0 - 4) { waited = sys_wait4(pid, st, 0) }
516 if waited != pid { p("WFLOW exec-organ ERROR child wait failed -- fail closed\n" as *u8); return 0 }
517 let ec: i64 = wait_status_rc(st[0])
518 p("WFLOW-EXEC-ORGAN connector=" as *u8); p(name); p(" elf=" as *u8); p(elf); p(" exit=" as *u8); pn(ec); p("\n" as *u8)
519 if ec != 0 { return 0 }
520 if docap == 1 {
521 if total > 480 { p("WFLOW capture ERROR output over 480 bytes -- fail loud\n" as *u8); return 0 }
522 var t: i64 = 0
523 while t < total { if cbuf[t] == (10 as u8) { cbuf[t] = 32 as u8 } if cbuf[t] == (13 as u8) { cbuf[t] = 32 as u8 } t = t + 1 }
524 while total > 0 { if cbuf[total-1] == (32 as u8) { total = total - 1 } else { break } }
525 cbuf[total] = 0 as u8
526 let ledc: *u8 = cx[2] as *u8
527 wf_emit_var(ledc, rid, capv, cbuf)
528 p("WFLOW-CAPTURE run=" as *u8); p(rid); p(" var=" as *u8); p(capv); p(" v=" as *u8); p(cbuf); p("\n" as *u8)
529 }
530 return 1
531}
532
533// ---- R11 loops: apply-to-each over the connector plane ----
534func wf_emit_iter(led: *u8, rid: *u8, n: i64, val: *u8) -> i64 {
535 let ln: *u8 = sys_mmap(768)
536 var o: i64 = 0
537 o = wf_cat(ln, o, "WFITER rid=" as *u8)
538 o = wf_cat(ln, o, rid)
539 o = wf_cat(ln, o, " n=" as *u8)
540 o = wf_catn(ln, o, n)
541 o = wf_cat(ln, o, " item=" as *u8)
542 o = wf_cat(ln, o, val)
543 ln[o] = 10 as u8
544 ln[o+1] = 0 as u8
545 return wf_append(led, ln)
546}
547// for-each: arg = <listexpr>~<per-item template> (template default {item}). The list substitutes ONCE
548// (comma-split, empties skipped, 32-item loud cap); per item: bind {item} as a run var, emit WFITER,
549// substitute the template, fork the connector. Any failed item fails the STEP (retry/on-fail apply).
550// Empty list = zero iterations, step OK (apply-to-each on an empty collection is a no-op).
551func wf_for_each(cx: *i64, chspec: *u8, argraw: *u8, rid: *u8) -> i64 {
552 let led: *u8 = cx[2] as *u8
553 let lb: *u8 = sys_mmap(WF_MAGIC_1024)
554 let tpl: *u8 = sys_mmap(512)
555 var i: i64 = 0
556 var sep: i64 = 0 - 1
557 while argraw[i] != (0 as u8) { if sep < 0 { if argraw[i] == (126 as u8) { sep = i } } i = i + 1 }
558 let alen: i64 = i
559 var lend: i64 = alen
560 if sep >= 0 { lend = sep }
561 var li: i64 = 0
562 while li < lend { lb[li] = argraw[li]; li = li + 1 }
563 lb[li] = 0 as u8
564 if sep >= 0 {
565 var ti: i64 = 0
566 var q: i64 = sep + 1
567 while q < alen { if ti < 510 { tpl[ti] = argraw[q]; ti = ti + 1 } q = q + 1 }
568 tpl[ti] = 0 as u8
569 } else {
570 var to: i64 = 0
571 to = wf_cat(tpl, to, "{item}" as *u8)
572 tpl[to] = 0 as u8
573 }
574 let lbx: *u8 = sys_mmap(WF_MAGIC_1024)
575 if wf_subst(cx, rid, lb, lbx, WF_MAGIC_1024) < 0 { return 0 }
576 var ln2: i64 = 0
577 while lbx[ln2] != (0 as u8) { ln2 = ln2 + 1 }
578 let iv: *i64 = sys_mmap(8 * 33) as *i64
579 var ni: i64 = 0
580 var st0: i64 = 0
581 var j: i64 = 0
582 var over: i64 = 0
583 while j <= ln2 {
584 var cut: i64 = 0
585 if j == ln2 { cut = 1 }
586 if cut == 0 { if lbx[j] == (44 as u8) { cut = 1 } }
587 if cut == 1 {
588 lbx[j] = 0 as u8
589 if lbx[st0] != (0 as u8) {
590 if ni >= 32 { over = 1 } else { iv[ni] = (lbx as i64) + st0; ni = ni + 1 }
591 }
592 st0 = j + 1
593 }
594 j = j + 1
595 }
596 if over == 1 { p("WFLOW for-each ERROR over 32 items -- fail loud\n" as *u8); return 0 }
597 if ni == 0 { p("WFLOW-FOREACH empty list, zero iterations\n" as *u8); return 1 }
598 let ibuf: *u8 = sys_mmap(WF_MAGIC_1024)
599 var k: i64 = 0
600 while k < ni {
601 let itv: *u8 = iv[k] as *u8
602 wf_emit_var(led, rid, "item" as *u8, itv)
603 wf_emit_iter(led, rid, k + 1, itv)
604 if wf_subst(cx, rid, tpl, ibuf, WF_MAGIC_1024) < 0 { return 0 }
605 p("WFLOW-FOREACH n=" as *u8); pn(k + 1); p(" item=" as *u8); p(itv); p("\n" as *u8)
606 if wf_exec_organ(cx, chspec, ibuf, rid) == 0 { p("WFLOW for-each item FAILED -- step fails\n" as *u8); return 0 }
607 k = k + 1
608 }
609 return 1
610}
611
612// ---- action execution (in-process v1 + exec-organ; send-channel semantics from nx_send) ----
613func wf_exec(cx: *i64, action: *u8, ch: *u8, arg: *u8, att: i64, rid: *u8) -> i64 {
614 p("WFLOW-ACTION action=" as *u8); p(action)
615 if seq(action, "send" as *u8) == 1 { p(" channel=" as *u8); p(ch); p(" route=" as *u8); p(sr_route(ch)) }
616 p(" arg=" as *u8); p(arg)
617 p(" att=" as *u8); pn(att)
618 p("\n" as *u8)
619 if seq(action, "probe-fail" as *u8) == 1 {
620 let k: i64 = wf_atoi(arg)
621 if att < k { return 0 }
622 return 1
623 }
624 if seq(action, "exec-organ" as *u8) == 1 { return wf_exec_organ(cx, ch, arg, rid) }
625 if seq(action, "for-each" as *u8) == 1 { return wf_for_each(cx, ch, arg, rid) }
626 return 1
627}
628
629// ---- R4 human-in-the-loop approvals: WFDEC decision records on the same ledger ----
630// WFDEC rid=<rid> step=<n> decision=APPROVED|DENIED by=<who>. DENY WINS if both exist (fail-safe,
631// deny-by-default). No decision = the approve step PARKS the run (status=WAIT, in-flight); wf_resume
632// re-examines it and completes once a decision lands.
633func wf_decision(cx: *i64, rid: *u8, idx: i64) -> i64 {
634 let led: *u8 = cx[2] as *u8
635 let szp: *i64 = sys_mmap(8) as *i64
636 let buf: *u8 = wf_readall(led, szp)
637 if (buf as i64) == 0 { return 0 }
638 let tok: *u8 = sys_mmap(256)
639 var q: i64 = 0
640 q = wf_cat(tok, q, " rid=" as *u8)
641 q = wf_cat(tok, q, rid)
642 tok[q] = 32 as u8
643 tok[q+1] = 0 as u8
644 let nb: *u8 = sys_mmap(128)
645 var o: i64 = 0
646 o = wf_cat(nb, o, " step=" as *u8)
647 o = wf_catn(nb, o, idx)
648 o = wf_cat(nb, o, " decision=DENIED" as *u8)
649 nb[o] = 0 as u8
650 if wf_lines_with2(buf, szp[0], tok, nb) > 0 { return 0 - 1 }
651 var o2: i64 = 0
652 o2 = wf_cat(nb, o2, " step=" as *u8)
653 o2 = wf_catn(nb, o2, idx)
654 o2 = wf_cat(nb, o2, " decision=APPROVED" as *u8)
655 nb[o2] = 0 as u8
656 if wf_lines_with2(buf, szp[0], tok, nb) > 0 { return 1 }
657 return 0
658}
659// the operator surface: append a decision record (only APPROVED or DENIED accepted, loud else)
660func wf_decide(cx: *i64, rid: *u8, idx: i64, decision: *u8, who: *u8) -> i64 {
661 var okd: i64 = 0
662 if seq(decision, "APPROVED" as *u8) == 1 { okd = 1 }
663 if seq(decision, "DENIED" as *u8) == 1 { okd = 1 }
664 if okd == 0 { p("WFLOW decide ERROR decision must be APPROVED or DENIED -- fail loud\n" as *u8); return 0 - 1 }
665 let led: *u8 = cx[2] as *u8
666 let ln: *u8 = sys_mmap(512)
667 var o: i64 = 0
668 o = wf_cat(ln, o, "WFDEC rid=" as *u8)
669 o = wf_cat(ln, o, rid)
670 o = wf_cat(ln, o, " step=" as *u8)
671 o = wf_catn(ln, o, idx)
672 o = wf_cat(ln, o, " decision=" as *u8)
673 o = wf_cat(ln, o, decision)
674 o = wf_cat(ln, o, " by=" as *u8)
675 o = wf_cat(ln, o, who)
676 ln[o] = 10 as u8
677 ln[o+1] = 0 as u8
678 return wf_append(led, ln)
679}
680
681// ---- R8 data passing: run-scoped variables on the SAME event-sourced ledger ----
682// WFVAR rid=<rid> k=<key> v=<value-to-end-of-line>. Trigger-event fields auto-bind at fire; a set-var
683// step writes derived vars; {key} placeholders in any step arg substitute at execution; LAST write wins;
684// unknown placeholder = LOUD step failure (the send_merge law). Values must not contain ~ or newline.
685func wf_emit_var(led: *u8, rid: *u8, k: *u8, v: *u8) -> i64 {
686 let ln: *u8 = sys_mmap(768)
687 var o: i64 = 0
688 o = wf_cat(ln, o, "WFVAR rid=" as *u8)
689 o = wf_cat(ln, o, rid)
690 o = wf_cat(ln, o, " k=" as *u8)
691 o = wf_cat(ln, o, k)
692 o = wf_cat(ln, o, " v=" as *u8)
693 o = wf_cat(ln, o, v)
694 ln[o] = 10 as u8
695 ln[o+1] = 0 as u8
696 return wf_append(led, ln)
697}
698// copy value after key to END OF LINE (values may contain spaces)
699func wf_val_line_end(buf: *u8, ls: i64, le: i64, key: *u8, out: *u8, cap: i64) -> i64 {
700 let m: i64 = slen(key)
701 var i: i64 = ls
702 while i + m <= le {
703 var k: i64 = 0
704 var hit: i64 = 1
705 while k < m { if buf[i+k] != key[k] { hit = 0; k = m } else { k = k + 1 } }
706 if hit == 1 {
707 var q: i64 = i + m
708 var t: i64 = 0
709 while q < le { if t < cap - 1 { out[t] = buf[q]; t = t + 1 } q = q + 1 }
710 out[t] = 0 as u8
711 return 1
712 }
713 i = i + 1
714 }
715 out[0] = 0 as u8
716 return 0
717}
718func wf_obs_tok_fwd(tok: *u8, rid: *u8) -> i64 {
719 var q: i64 = 0
720 q = wf_cat(tok, q, " rid=" as *u8)
721 q = wf_cat(tok, q, rid)
722 tok[q] = 32 as u8
723 tok[q+1] = 0 as u8
724 return q
725}
726// latest WFVAR value for (rid,key); 1 found / 0 not
727func wf_var_get(cx: *i64, rid: *u8, key: *u8, out: *u8, cap: i64) -> i64 {
728 let led: *u8 = cx[2] as *u8
729 let szp: *i64 = sys_mmap(8) as *i64
730 let buf: *u8 = wf_readall(led, szp)
731 if (buf as i64) == 0 { return 0 }
732 let n: i64 = szp[0]
733 let tok: *u8 = sys_mmap(256)
734 wf_obs_tok_fwd(tok, rid)
735 let kt: *u8 = sys_mmap(160)
736 var q: i64 = 0
737 q = wf_cat(kt, q, " k=" as *u8)
738 q = wf_cat(kt, q, key)
739 kt[q] = 32 as u8
740 kt[q+1] = 0 as u8
741 var found: i64 = 0
742 var ls: i64 = 0
743 while ls < n {
744 var le: i64 = ls
745 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
746 if wf_has_rng(buf, ls, le, "WFVAR" as *u8) == 1 { if wf_has_rng(buf, ls, le, tok) == 1 { if wf_has_rng(buf, ls, le, kt) == 1 {
747 wf_val_line_end(buf, ls, le, " v=" as *u8, out, cap)
748 found = 1
749 } } }
750 ls = le + 1
751 }
752 return found
753}
754// substitute {key} placeholders from the run's variable plane; -1 LOUD on unknown or unclosed
755func wf_subst(cx: *i64, rid: *u8, src: *u8, out: *u8, cap: i64) -> i64 {
756 let kb: *u8 = sys_mmap(128)
757 let vb: *u8 = sys_mmap(512)
758 var i: i64 = 0
759 var o: i64 = 0
760 while src[i] != (0 as u8) {
761 if src[i] == (123 as u8) {
762 var j: i64 = i + 1
763 var kn: i64 = 0
764 var closed: i64 = 0
765 while src[j] != (0 as u8) {
766 if src[j] == (125 as u8) { closed = 1; break }
767 if kn < 127 { kb[kn] = src[j]; kn = kn + 1 }
768 j = j + 1
769 }
770 kb[kn] = 0 as u8
771 if closed == 0 { p("WFLOW subst ERROR unclosed placeholder -- fail loud\n" as *u8); return 0 - 1 }
772 if wf_var_get(cx, rid, kb, vb, 512) == 0 {
773 p("WFLOW subst ERROR unknown variable=" as *u8); p(kb); p(" -- fail loud\n" as *u8)
774 return 0 - 1
775 }
776 var t: i64 = 0
777 while vb[t] != (0 as u8) { if o < cap - 1 { out[o] = vb[t]; o = o + 1 } t = t + 1 }
778 i = j + 1
779 } else {
780 if o < cap - 1 { out[o] = src[i]; o = o + 1 }
781 i = i + 1
782 }
783 }
784 out[o] = 0 as u8
785 return o
786}
787// bind every k=v field of the trigger event as an initial run variable (segment 0 = kind, skipped)
788func wf_bind_event(cx: *i64, rid: *u8, ev: *u8) -> i64 {
789 let led: *u8 = cx[2] as *u8
790 let kb: *u8 = sys_mmap(160)
791 let vb: *u8 = sys_mmap(512)
792 var seg: i64 = 0
793 var i: i64 = 0
794 var bound: i64 = 0
795 while ev[i] != (0 as u8) {
796 if ev[i] == (126 as u8) {
797 seg = seg + 1
798 var j: i64 = i + 1
799 var kn: i64 = 0
800 var haseq: i64 = 0
801 while ev[j] != (0 as u8) {
802 if ev[j] == (126 as u8) { break }
803 if ev[j] == (61 as u8) { haseq = 1; j = j + 1; break }
804 if kn < 158 { kb[kn] = ev[j]; kn = kn + 1 }
805 j = j + 1
806 }
807 kb[kn] = 0 as u8
808 if haseq == 1 {
809 var vn: i64 = 0
810 while ev[j] != (0 as u8) { if ev[j] == (126 as u8) { break } if vn < 510 { vb[vn] = ev[j]; vn = vn + 1 } j = j + 1 }
811 vb[vn] = 0 as u8
812 if kn > 0 { wf_emit_var(led, rid, kb, vb); bound = bound + 1 }
813 }
814 }
815 i = i + 1
816 }
817 return bound
818}
819
820// ---- R12 versioning: @version N directive in the flows file; runs stamp WFVER; resume flags drift ----
821func wf_defs_version(path: *u8) -> i64 {
822 let szp: *i64 = sys_mmap(8) as *i64
823 let buf: *u8 = wf_readall(path, szp)
824 if (buf as i64) == 0 { return 0 }
825 let n: i64 = szp[0]
826 var ls: i64 = 0
827 while ls < n {
828 var le: i64 = ls
829 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
830 if wf_has_rng(buf, ls, le, "@version " as *u8) == 1 {
831 let v: i64 = wf_num_after(buf, ls, le, "@version " as *u8)
832 if v > 0 { return v }
833 }
834 ls = le + 1
835 }
836 return 0
837}
838func wf_emit_ver(led: *u8, rid: *u8, ver: i64) -> i64 {
839 let ln: *u8 = sys_mmap(512)
840 var o: i64 = 0
841 o = wf_cat(ln, o, "WFVER rid=" as *u8)
842 o = wf_cat(ln, o, rid)
843 o = wf_cat(ln, o, " ver=" as *u8)
844 o = wf_catn(ln, o, ver)
845 ln[o] = 10 as u8
846 ln[o+1] = 0 as u8
847 return wf_append(led, ln)
848}
849func wf_emit_drift(led: *u8, rid: *u8, ranv: i64, nowv: i64) -> i64 {
850 let ln: *u8 = sys_mmap(512)
851 var o: i64 = 0
852 o = wf_cat(ln, o, "WFVERDRIFT rid=" as *u8)
853 o = wf_cat(ln, o, rid)
854 o = wf_cat(ln, o, " ran=" as *u8)
855 o = wf_catn(ln, o, ranv)
856 o = wf_cat(ln, o, " now=" as *u8)
857 o = wf_catn(ln, o, nowv)
858 ln[o] = 10 as u8
859 ln[o+1] = 0 as u8
860 return wf_append(led, ln)
861}
862// the version a run was STARTED under (last WFVER line for the rid; 0 = unversioned)
863func wf_run_ver(buf: *u8, n: i64, ridtok: *u8) -> i64 {
864 var v: i64 = 0
865 var ls: i64 = 0
866 while ls < n {
867 var le: i64 = ls
868 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
869 if wf_has_rng(buf, ls, le, "WFVER " as *u8) == 1 { if wf_has_rng(buf, ls, le, ridtok) == 1 {
870 let x: i64 = wf_num_after(buf, ls, le, " ver=" as *u8)
871 if x > 0 { v = x }
872 } }
873 ls = le + 1
874 }
875 return v
876}
877
878// ---- R13 templates: gallery = <dir>/index (name|description rows) + <name>.flows/.steps pairs ----
879func wf_tpl_exists(path: *u8) -> i64 {
880 let fd: i64 = __syscall(257, 0 - 100, path as i64, 0, 0, 0, 0)
881 if fd < 0 { return 0 }
882 __syscall(3, fd, 0, 0, 0, 0, 0)
883 return 1
884}
885func wf_tpl_path(dir: *u8, name: *u8, ext: *u8, out: *u8) -> i64 {
886 var o: i64 = 0
887 o = wf_cat(out, o, dir)
888 out[o] = 47 as u8
889 o = o + 1
890 o = wf_cat(out, o, name)
891 o = wf_cat(out, o, ext)
892 out[o] = 0 as u8
893 return o
894}
895func wf_tpl_index_path(dir: *u8, out: *u8) -> i64 {
896 var o: i64 = 0
897 o = wf_cat(out, o, dir)
898 out[o] = 47 as u8
899 o = o + 1
900 o = wf_cat(out, o, "index" as *u8)
901 out[o] = 0 as u8
902 return o
903}
904// list the gallery: every index row printed with file-existence + version; returns USABLE count, -1 loud
905func wf_tpl_list(dir: *u8) -> i64 {
906 let ip: *u8 = sys_mmap(512)
907 wf_tpl_index_path(dir, ip)
908 let szp: *i64 = sys_mmap(8) as *i64
909 let buf: *u8 = wf_readall(ip, szp)
910 if (buf as i64) == 0 { return 0 - 1 }
911 let n: i64 = szp[0]
912 if n == 0 { p("WFLOW templates ERROR missing gallery index -- fail loud\n" as *u8); return 0 - 1 }
913 let nm: *u8 = sys_mmap(128)
914 let ds: *u8 = sys_mmap(512)
915 let fp2: *u8 = sys_mmap(512)
916 let sp2: *u8 = sys_mmap(512)
917 var cnt: i64 = 0
918 var ls: i64 = 0
919 while ls < n {
920 var le: i64 = ls
921 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
922 var keep: i64 = 0
923 if le > ls { keep = 1 }
924 if keep == 1 { if buf[ls] == (35 as u8) { keep = 0 } }
925 if keep == 1 {
926 let row: *u8 = sys_mmap(le - ls + 2)
927 var k: i64 = 0
928 while ls + k < le { row[k] = buf[ls+k]; k = k + 1 }
929 row[k] = 0 as u8
930 pipe_field(row, 0, nm, 128)
931 pipe_field(row, 1, ds, 512)
932 wf_tpl_path(dir, nm, ".flows" as *u8, fp2)
933 wf_tpl_path(dir, nm, ".steps" as *u8, sp2)
934 let fe: i64 = wf_tpl_exists(fp2)
935 let se: i64 = wf_tpl_exists(sp2)
936 let v: i64 = wf_defs_version(fp2)
937 p("WFTPL name=" as *u8); p(nm)
938 p(" ver=" as *u8); pn(v)
939 p(" flows=" as *u8); pn(fe)
940 p(" steps=" as *u8); pn(se)
941 p(" desc=" as *u8); p(ds)
942 p("\n" as *u8)
943 if fe == 1 { if se == 1 { cnt = cnt + 1 } }
944 }
945 ls = le + 1
946 }
947 p("WFTPL-TOTAL usable=" as *u8); pn(cnt); p("\n" as *u8)
948 return cnt
949}
950func wf_tpl_indexed(dir: *u8, name: *u8) -> i64 {
951 let ip: *u8 = sys_mmap(512)
952 wf_tpl_index_path(dir, ip)
953 let szp: *i64 = sys_mmap(8) as *i64
954 let buf: *u8 = wf_readall(ip, szp)
955 if (buf as i64) == 0 { return 0 - 1 }
956 let n: i64 = szp[0]
957 if n == 0 { return 0 - 1 }
958 let nm: *u8 = sys_mmap(128)
959 var ls: i64 = 0
960 while ls < n {
961 var le: i64 = ls
962 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
963 if le > ls { if buf[ls] != (35 as u8) {
964 let row: *u8 = sys_mmap(le - ls + 2)
965 var k: i64 = 0
966 while ls + k < le { row[k] = buf[ls+k]; k = k + 1 }
967 row[k] = 0 as u8
968 pipe_field(row, 0, nm, 128)
969 if seq(nm, name) == 1 { return 1 }
970 } }
971 ls = le + 1
972 }
973 return 0
974}
975// instantiate: fail-closed to INDEXED templates; byte-copy both files; whole-set validate the COPY
976// (a broken template never lands silently). Returns the template version (0 = unversioned), -1 loud.
977func wf_tpl_instantiate(dir: *u8, name: *u8, df: *u8, ds: *u8) -> i64 {
978 let ix: i64 = wf_tpl_indexed(dir, name)
979 if ix != 1 { p("WFLOW instantiate ERROR template not in gallery index -- fail closed\n" as *u8); return 0 - 1 }
980 let fp2: *u8 = sys_mmap(512)
981 let sp2: *u8 = sys_mmap(512)
982 wf_tpl_path(dir, name, ".flows" as *u8, fp2)
983 wf_tpl_path(dir, name, ".steps" as *u8, sp2)
984 let szp: *i64 = sys_mmap(8) as *i64
985 let fb: *u8 = wf_readall(fp2, szp)
986 if (fb as i64) == 0 { return 0 - 1 }
987 if szp[0] == 0 { p("WFLOW instantiate ERROR template flows file missing -- fail loud\n" as *u8); return 0 - 1 }
988 if wf_writeall(df, fb, szp[0]) != 0 { return 0 - 1 }
989 let sb: *u8 = wf_readall(sp2, szp)
990 if (sb as i64) == 0 { return 0 - 1 }
991 if szp[0] == 0 { p("WFLOW instantiate ERROR template steps file missing -- fail loud\n" as *u8); return 0 - 1 }
992 if wf_writeall(ds, sb, szp[0]) != 0 { return 0 - 1 }
993 let fl: *i64 = sys_mmap(8 * 64) as *i64
994 let st2: *i64 = sys_mmap(8 * 128) as *i64
995 let nf: i64 = wf_lines_load(df, fl, 64)
996 if nf < 1 { return 0 - 1 }
997 let nst: i64 = wf_lines_load(ds, st2, 128)
998 if nst < 1 { return 0 - 1 }
999 if wf_load_flows(fl, nf) < 0 { return 0 - 1 }
1000 if wf_load(st2, nst) < 0 { return 0 - 1 }
1001 let v: i64 = wf_defs_version(df)
1002 p("WFLOW-INSTANTIATE template=" as *u8); p(name); p(" version=" as *u8); pn(v); p(" flows=" as *u8); pn(nf); p(" steps=" as *u8); pn(nst); p("\n" as *u8)
1003 return v
1004}
1005
1006// R9 branching: per-step condition k=v against the run variable plane. Missing variable or
1007// mismatch = 0 (the step SKIPs, recorded in the ledger); match = 1. Equality only, v1.
1008func wf_cond_ok(cx: *i64, rid: *u8, cond: *u8) -> i64 {
1009 let kb: *u8 = sys_mmap(160)
1010 let want: *u8 = sys_mmap(512)
1011 var i: i64 = 0
1012 var kn: i64 = 0
1013 while cond[i] != (0 as u8) { if cond[i] == (61 as u8) { break } if kn < 158 { kb[kn] = cond[i]; kn = kn + 1 } i = i + 1 }
1014 kb[kn] = 0 as u8
1015 if cond[i] != (61 as u8) { return 0 }
1016 if kn == 0 { return 0 }
1017 var wn: i64 = 0
1018 var q: i64 = i + 1
1019 while cond[q] != (0 as u8) { if wn < 510 { want[wn] = cond[q]; wn = wn + 1 } q = q + 1 }
1020 want[wn] = 0 as u8
1021 let vb: *u8 = sys_mmap(512)
1022 if wf_var_get(cx, rid, kb, vb, 512) == 0 { return 0 }
1023 return seq(vb, want)
1024}
1025
1026// ---- the durable executor: run flow fid for run rid, starting AFTER step s0 ----
1027// cx bundle: cx[0]=steps ptr, cx[1]=nsteps, cx[2]=ledger path
1028func wf_run_from(cx: *i64, rid: *u8, fid: *u8, s0: i64) -> i64 {
1029 let steps: *i64 = cx[0] as *i64
1030 let ns: i64 = cx[1]
1031 let led: *u8 = cx[2] as *u8
1032 let f2: *u8 = sys_mmap(64)
1033 let ix: *u8 = sys_mmap(32)
1034 let a: *u8 = sys_mmap(64)
1035 let ch: *u8 = sys_mmap(64)
1036 let arg: *u8 = sys_mmap(512)
1037 let ma: *u8 = sys_mmap(32)
1038 var i: i64 = 0
1039 var done: i64 = 0
1040 var failed: i64 = 0
1041 var parked: i64 = 0
1042 var jumpto: i64 = 0
1043 while i < ns {
1044 let r: *u8 = steps[i] as *u8
1045 pipe_field(r, 0, f2, 64)
1046 var use: i64 = 0
1047 if seq(f2, fid) == 1 { use = 1 }
1048 var idx: i64 = 0
1049 if use == 1 { pipe_field(r, 1, ix, 32); idx = wf_atoi(ix); if idx <= s0 { use = 0 } }
1050 // R10: an active on-fail jump skips the remaining normal-path steps (recorded) until the target
1051 if use == 1 { if jumpto > 0 { if idx < jumpto { wf_emit(led, rid, fid, idx, "SKIP" as *u8, 0); use = 0 } else { jumpto = 0 } } }
1052 if use == 1 {
1053 pipe_field(r, 2, a, 64)
1054 pipe_field(r, 3, ch, 64)
1055 pipe_field(r, 4, arg, 512)
1056 pipe_field(r, 5, ma, 32)
1057 // R9: optional per-step condition (7th field, k=v) vs the variable plane -> SKIP on mismatch
1058 let cond: *u8 = sys_mmap(256)
1059 pipe_field(r, 6, cond, 256)
1060 let argx: *u8 = sys_mmap(WF_MAGIC_1024)
1061 var mode: i64 = 0
1062 var docond: i64 = 0
1063 if cond[0] != (0 as u8) { docond = 1 }
1064 if cond[0] == (45 as u8) { if cond[1] == (0 as u8) { docond = 0 } }
1065 if docond == 1 { if wf_cond_ok(cx, rid, cond) == 0 { mode = 4 } }
1066 // R8: resolve {key} placeholders from the run variable plane; unknown = LOUD step failure.
1067 // for-each bypasses step-level subst ({item} binds inside the loop) -- the loop substitutes
1068 // the list once and the per-item template each iteration.
1069 if mode == 0 {
1070 var rawarg: i64 = 0
1071 if seq(a, "for-each" as *u8) == 1 { rawarg = 1 }
1072 if rawarg == 1 {
1073 var cc: i64 = 0
1074 while arg[cc] != (0 as u8) { argx[cc] = arg[cc]; cc = cc + 1 }
1075 argx[cc] = 0 as u8
1076 } else {
1077 if wf_subst(cx, rid, arg, argx, WF_MAGIC_1024) < 0 { mode = 3 }
1078 }
1079 }
1080 if mode == 4 {
1081 p("WFLOW-SKIP run=" as *u8); p(rid); p(" step=" as *u8); pn(idx); p(" cond=" as *u8); p(cond); p("\n" as *u8)
1082 wf_emit(led, rid, fid, idx, "SKIP" as *u8, 0)
1083 }
1084 if mode == 0 { if seq(a, "set-var" as *u8) == 1 { mode = 2 } }
1085 if mode == 0 { if seq(a, "approve" as *u8) == 1 { mode = 1 } }
1086 if mode == 3 {
1087 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, 0)
1088 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0)
1089 failed = 1
1090 i = ns
1091 }
1092 if mode == 2 {
1093 let kb: *u8 = sys_mmap(160)
1094 let vb: *u8 = sys_mmap(512)
1095 var e: i64 = 0
1096 var kn: i64 = 0
1097 while argx[e] != (0 as u8) { if argx[e] == (61 as u8) { break } if kn < 158 { kb[kn] = argx[e]; kn = kn + 1 } e = e + 1 }
1098 kb[kn] = 0 as u8
1099 if argx[e] == (61 as u8) { if kn > 0 {
1100 var vn: i64 = 0
1101 var q2: i64 = e + 1
1102 while argx[q2] != (0 as u8) { if vn < 510 { vb[vn] = argx[q2]; vn = vn + 1 } q2 = q2 + 1 }
1103 vb[vn] = 0 as u8
1104 wf_emit_var(led, rid, kb, vb)
1105 p("WFLOW-SETVAR run=" as *u8); p(rid); p(" k=" as *u8); p(kb); p(" v=" as *u8); p(vb); p("\n" as *u8)
1106 wf_emit(led, rid, fid, idx, "OK" as *u8, 1)
1107 done = done + 1
1108 } else { mode = 3 } } else { mode = 3 }
1109 if mode == 3 {
1110 p("WFLOW set-var ERROR needs key=value -- fail loud\n" as *u8)
1111 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, 0)
1112 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0)
1113 failed = 1
1114 i = ns
1115 }
1116 }
1117 if mode == 1 {
1118 let dec: i64 = wf_decision(cx, rid, idx)
1119 if dec == 1 {
1120 p("WFLOW-APPROVAL granted run=" as *u8); p(rid); p(" step=" as *u8); pn(idx); p("\n" as *u8)
1121 wf_emit(led, rid, fid, idx, "OK" as *u8, 1)
1122 done = done + 1
1123 }
1124 if dec == (0 - 1) {
1125 p("WFLOW-APPROVAL DENIED run=" as *u8); p(rid); p(" step=" as *u8); pn(idx); p(" -- run fails (deny wins)\n" as *u8)
1126 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, 0)
1127 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0)
1128 failed = 1
1129 i = ns
1130 }
1131 if dec == 0 {
1132 p("WFLOW-APPROVAL pending run=" as *u8); p(rid); p(" step=" as *u8); pn(idx)
1133 p(" approver=" as *u8); p(ch); p(" what=" as *u8); p(argx); p(" -- run PARKED\n" as *u8)
1134 wf_emit(led, rid, fid, idx, "WAIT" as *u8, 0)
1135 parked = 1
1136 i = ns
1137 }
1138 }
1139 if mode == 0 {
1140 let maxa: i64 = wf_atoi(ma)
1141 var att: i64 = 1
1142 var okd: i64 = 0
1143 while att <= maxa {
1144 wf_emit(led, rid, fid, idx, "ATT" as *u8, att)
1145 let rc: i64 = wf_exec(cx, a, ch, argx, att, rid)
1146 if rc == 1 { wf_emit(led, rid, fid, idx, "OK" as *u8, att); okd = 1; att = maxa + 1 } else { att = att + 1 }
1147 }
1148 if okd == 0 {
1149 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, maxa)
1150 // R10: on-fail routing (8th field) -- exhausted ACTION failures only; the failure stays
1151 // recorded (FAILSTEP + ONFAIL), wf_error vars land on the plane, execution jumps forward
1152 let onf: *u8 = sys_mmap(32)
1153 pipe_field(r, 7, onf, 32)
1154 var tj: i64 = 0
1155 var haveof: i64 = 0
1156 if onf[0] != (0 as u8) { haveof = 1 }
1157 if onf[0] == (45 as u8) { if onf[1] == (0 as u8) { haveof = 0 } }
1158 if haveof == 1 { tj = wf_atoi(onf) }
1159 if tj > 0 {
1160 wf_emit(led, rid, fid, idx, "ONFAIL" as *u8, tj)
1161 wf_emit_var(led, rid, "wf_error" as *u8, "yes" as *u8)
1162 let sv: *u8 = sys_mmap(32)
1163 var so: i64 = 0
1164 so = wf_cat(sv, so, "step" as *u8)
1165 so = wf_catn(sv, so, idx)
1166 sv[so] = 0 as u8
1167 wf_emit_var(led, rid, "wf_error_step" as *u8, sv)
1168 p("WFLOW-ONFAIL run=" as *u8); p(rid); p(" step=" as *u8); pn(idx); p(" routes-to=" as *u8); pn(tj); p("\n" as *u8)
1169 jumpto = tj
1170 } else {
1171 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0)
1172 failed = 1
1173 i = ns
1174 }
1175 } else { done = done + 1 }
1176 }
1177 }
1178 i = i + 1
1179 }
1180 if failed == 1 { return 0 - 2 }
1181 if parked == 1 { return 0 - 3 }
1182 wf_emit(led, rid, fid, 0, "DONE" as *u8, 0)
1183 return done
1184}
1185
1186// ---- fire an event: durable ledger-deduped, one run per matched flow ----
1187func wf_fire(cx: *i64, flows: *i64, nf: i64, ev: *u8) -> i64 {
1188 let evid: *u8 = sys_mmap(128)
1189 if wf_evid(ev, evid, 128) == 0 { p("WFLOW fire ERROR event has no id -- fail loud\n" as *u8); return 0 - 1 }
1190 let led: *u8 = cx[2] as *u8
1191 let szp: *i64 = sys_mmap(8) as *i64
1192 let buf: *u8 = wf_readall(led, szp)
1193 if (buf as i64) == 0 { return 0 - 1 }
1194 let fid: *u8 = sys_mmap(64)
1195 let rid: *u8 = sys_mmap(200)
1196 let tok: *u8 = sys_mmap(256)
1197 var started: i64 = 0
1198 var matched: i64 = 0
1199 var f: i64 = 0
1200 while f < nf {
1201 let fh: *u8 = flows[f] as *u8
1202 if wf_match(fh, ev) == 1 {
1203 matched = matched + 1
1204 pipe_field(fh, 0, fid, 64)
1205 var o: i64 = 0
1206 o = wf_cat(rid, o, evid)
1207 rid[o] = 46 as u8
1208 o = o + 1
1209 o = wf_cat(rid, o, fid)
1210 rid[o] = 0 as u8
1211 var q: i64 = 0
1212 q = wf_cat(tok, q, " rid=" as *u8)
1213 q = wf_cat(tok, q, rid)
1214 tok[q] = 32 as u8
1215 tok[q+1] = 0 as u8
1216 if wf_lines_with2(buf, szp[0], tok, "status=START" as *u8) > 0 {
1217 p("WFLOW dedup: run " as *u8); p(rid); p(" already started (durable) -- skipped\n" as *u8)
1218 } else {
1219 wf_emit(led, rid, fid, 0, "START" as *u8, 0)
1220 if cx[4] > 0 { wf_emit_ver(led, rid, cx[4]) }
1221 wf_bind_event(cx, rid, ev)
1222 wf_run_from(cx, rid, fid, 0)
1223 started = started + 1
1224 }
1225 }
1226 f = f + 1
1227 }
1228 if matched > 0 { if started == 0 { return 0 - 100 } }
1229 return started
1230}
1231
1232// any step row for this flow?
1233func wf_flow_known(cx: *i64, fid: *u8) -> i64 {
1234 let steps: *i64 = cx[0] as *i64
1235 let ns: i64 = cx[1]
1236 let f2: *u8 = sys_mmap(64)
1237 var i: i64 = 0
1238 while i < ns {
1239 pipe_field(steps[i] as *u8, 0, f2, 64)
1240 if seq(f2, fid) == 1 { return 1 }
1241 i = i + 1
1242 }
1243 return 0
1244}
1245
1246// ---- crash recovery: complete every in-flight run from its first non-OK step ----
1247func wf_resume(cx: *i64) -> i64 {
1248 let led: *u8 = cx[2] as *u8
1249 let szp: *i64 = sys_mmap(8) as *i64
1250 let buf: *u8 = wf_readall(led, szp)
1251 if (buf as i64) == 0 { return 0 - 1 }
1252 let n: i64 = szp[0]
1253 let rid: *u8 = sys_mmap(200)
1254 let fid: *u8 = sys_mmap(64)
1255 let tok: *u8 = sys_mmap(256)
1256 var resumed: i64 = 0
1257 var ls: i64 = 0
1258 while ls < n {
1259 var le: i64 = ls
1260 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
1261 if wf_has_rng(buf, ls, le, "status=START" as *u8) == 1 {
1262 wf_val_after(buf, ls, le, " rid=" as *u8, rid, 200)
1263 wf_val_after(buf, ls, le, " flow=" as *u8, fid, 64)
1264 var q: i64 = 0
1265 q = wf_cat(tok, q, " rid=" as *u8)
1266 q = wf_cat(tok, q, rid)
1267 tok[q] = 32 as u8
1268 tok[q+1] = 0 as u8
1269 let dn: i64 = wf_lines_with2(buf, n, tok, "status=DONE" as *u8)
1270 let fl: i64 = wf_lines_with2(buf, n, tok, "status=FAILED" as *u8)
1271 if dn == 0 { if fl == 0 {
1272 if wf_flow_known(cx, fid) == 0 { p("WFLOW resume ERROR unknown flow=" as *u8); p(fid); p(" -- fail loud, nothing fabricated\n" as *u8); return 0 - 1 }
1273 // R12: flag version drift (run started under a different definitions version) -- resume
1274 // proceeds (additive semantics) but the drift is RECORDED and printed
1275 let curv: i64 = cx[4]
1276 if curv > 0 {
1277 let ranv: i64 = wf_run_ver(buf, n, tok)
1278 if ranv > 0 { if ranv != curv {
1279 wf_emit_drift(led, rid, ranv, curv)
1280 p("WFLOW-RESUME-WARN version drift run=" as *u8); p(rid); p(" ran=" as *u8); pn(ranv); p(" now=" as *u8); pn(curv); p("\n" as *u8)
1281 } }
1282 }
1283 let s0: i64 = wf_max_ok_step(buf, n, tok)
1284 p("WFLOW-RESUME run=" as *u8); p(rid); p(" from-step=" as *u8); pn(s0 + 1); p("\n" as *u8)
1285 wf_run_from(cx, rid, fid, s0)
1286 resumed = resumed + 1
1287 } }
1288 }
1289 ls = le + 1
1290 }
1291 return resumed
1292}
1293
1294// ---- definitions as DATA FILES: load rows, validate whole, fire/resume (the production entry path) ----
1295// skips comment (leading 35) and empty lines; -1 loud on missing/empty file or row-cap
1296func wf_lines_load(path: *u8, arr: *i64, cap: i64) -> i64 {
1297 let szp: *i64 = sys_mmap(8) as *i64
1298 let buf: *u8 = wf_readall(path, szp)
1299 if (buf as i64) == 0 { return 0 - 1 }
1300 let n: i64 = szp[0]
1301 if n == 0 { p("WFLOW files ERROR missing or empty definition file -- fail loud\n" as *u8); return 0 - 1 }
1302 var cnt: i64 = 0
1303 var ls: i64 = 0
1304 while ls < n {
1305 var le: i64 = ls
1306 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
1307 var keep: i64 = 0
1308 if le > ls { keep = 1 }
1309 if keep == 1 { if buf[ls] == (35 as u8) { keep = 0 } }
1310 if keep == 1 { if buf[ls] == (64 as u8) { keep = 0 } }
1311 if keep == 1 {
1312 if cnt >= cap { p("WFLOW files ERROR too many definition rows -- fail loud\n" as *u8); return 0 - 1 }
1313 let s: *u8 = sys_mmap(le - ls + 2)
1314 var k: i64 = 0
1315 while ls + k < le { s[k] = buf[ls+k]; k = k + 1 }
1316 s[k] = 0 as u8
1317 arr[cnt] = s as i64
1318 cnt = cnt + 1
1319 }
1320 ls = le + 1
1321 }
1322 return cnt
1323}
1324func wf_cx_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8, cx: *i64, flows: *i64) -> i64 {
1325 let steps: *i64 = sys_mmap(8 * 128) as *i64
1326 let nf: i64 = wf_lines_load(flowsp, flows, 64)
1327 if nf < 1 { return 0 - 1 }
1328 let nst: i64 = wf_lines_load(stepsp, steps, 128)
1329 if nst < 1 { return 0 - 1 }
1330 if wf_load_flows(flows, nf) < 0 { return 0 - 1 }
1331 if wf_load(steps, nst) < 0 { return 0 - 1 }
1332 cx[0] = steps as i64
1333 cx[1] = nst
1334 cx[2] = led as i64
1335 cx[3] = 0
1336 var usecat: i64 = 1
1337 if cat[0] == (45 as u8) { if cat[1] == (0 as u8) { usecat = 0 } }
1338 if usecat == 1 { cx[3] = cat as i64 }
1339 cx[4] = wf_defs_version(flowsp)
1340 return nf
1341}
1342func wf_fire_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8, ev: *u8) -> i64 {
1343 let cx: *i64 = sys_mmap(40) as *i64
1344 let flows: *i64 = sys_mmap(8 * 64) as *i64
1345 let nf: i64 = wf_cx_files(flowsp, stepsp, led, cat, cx, flows)
1346 if nf < 0 { return 0 - 1 }
1347 return wf_fire(cx, flows, nf, ev)
1348}
1349func wf_resume_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8) -> i64 {
1350 let cx: *i64 = sys_mmap(40) as *i64
1351 let flows: *i64 = sys_mmap(8 * 64) as *i64
1352 let nf: i64 = wf_cx_files(flowsp, stepsp, led, cat, cx, flows)
1353 if nf < 0 { return 0 - 1 }
1354 return wf_resume(cx)
1355}
1356
1357// ---- R6 run observability: board + 0-JS HTML derived ENTIRELY by ledger replay (nothing self-reported) ----
1358func wf_state_name(st: i64) -> *u8 {
1359 if st == 3 { return "DONE" as *u8 }
1360 if st == 2 { return "FAILED" as *u8 }
1361 if st == 1 { return "PARKED" as *u8 }
1362 return "RUNNING" as *u8
1363}
1364// derived state: DONE > FAILED > PARKED(waiting approval) > RUNNING(in-flight)
1365func wf_run_state(buf: *u8, n: i64, tok: *u8) -> i64 {
1366 if wf_lines_with2(buf, n, tok, "status=DONE" as *u8) > 0 { return 3 }
1367 if wf_lines_with2(buf, n, tok, "status=FAILED" as *u8) > 0 { return 2 }
1368 if wf_lines_with2(buf, n, tok, "status=WAIT" as *u8) > 0 { return 1 }
1369 return 0
1370}
1371// collect distinct run ids from START lines; returns count, -1 loud over cap
1372func wf_obs_runs(buf: *u8, n: i64, rids: *i64, cap: i64) -> i64 {
1373 var cnt: i64 = 0
1374 var ls: i64 = 0
1375 while ls < n {
1376 var le: i64 = ls
1377 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 }
1378 if wf_has_rng(buf, ls, le, "status=START" as *u8) == 1 {
1379 let r: *u8 = sys_mmap(200)
1380 wf_val_after(buf, ls, le, " rid=" as *u8, r, 200)
1381 var dup: i64 = 0
1382 var i: i64 = 0
1383 while i < cnt { if seq(rids[i] as *u8, r) == 1 { dup = 1; i = cnt } else { i = i + 1 } }
1384 if dup == 0 {
1385 if cnt >= cap { p("WFLOW obs ERROR run cap exceeded -- fail loud\n" as *u8); return 0 - 1 }
1386 rids[cnt] = r as i64
1387 cnt = cnt + 1
1388 }
1389 }
1390 ls = le + 1
1391 }
1392 return cnt
1393}
1394func wf_obs_tok(tok: *u8, rid: *u8) -> i64 {
1395 var q: i64 = 0
1396 q = wf_cat(tok, q, " rid=" as *u8)
1397 q = wf_cat(tok, q, rid)
1398 tok[q] = 32 as u8
1399 tok[q+1] = 0 as u8
1400 return q
1401}
1402func wf_obs_board(led: *u8) -> i64 {
1403 let szp: *i64 = sys_mmap(8) as *i64
1404 let buf: *u8 = wf_readall(led, szp)
1405 if (buf as i64) == 0 { return 0 - 1 }
1406 let n: i64 = szp[0]
1407 let rids: *i64 = sys_mmap(8 * 128) as *i64
1408 let cnt: i64 = wf_obs_runs(buf, n, rids, 128)
1409 if cnt < 0 { return 0 - 1 }
1410 let tok: *u8 = sys_mmap(256)
1411 var d: i64 = 0
1412 var f: i64 = 0
1413 var w: i64 = 0
1414 var ru: i64 = 0
1415 var i: i64 = 0
1416 while i < cnt {
1417 let rid: *u8 = rids[i] as *u8
1418 wf_obs_tok(tok, rid)
1419 let st: i64 = wf_run_state(buf, n, tok)
1420 let oks: i64 = wf_lines_with2(buf, n, tok, "status=OK" as *u8)
1421 let ats: i64 = wf_lines_with2(buf, n, tok, "status=ATT" as *u8)
1422 p("WFOBS run=" as *u8); p(rid)
1423 p(" state=" as *u8); p(wf_state_name(st))
1424 p(" oksteps=" as *u8); pn(oks)
1425 p(" attempts=" as *u8); pn(ats)
1426 p("\n" as *u8)
1427 if st == 3 { d = d + 1 }
1428 if st == 2 { f = f + 1 }
1429 if st == 1 { w = w + 1 }
1430 if st == 0 { ru = ru + 1 }
1431 i = i + 1
1432 }
1433 p("WFOBS-TOTALS runs=" as *u8); pn(cnt); p(" done=" as *u8); pn(d); p(" failed=" as *u8); pn(f); p(" parked=" as *u8); pn(w); p(" running=" as *u8); pn(ru); p("\n" as *u8)
1434 return cnt
1435}
1436// 0-JS sovereign run board page; unquoted HTML5 attrs + rgb() colors (no quote/hash/bang literals)
1437func wf_obs_html(led: *u8, out: *u8) -> i64 {
1438 let szp: *i64 = sys_mmap(8) as *i64
1439 let buf: *u8 = wf_readall(led, szp)
1440 if (buf as i64) == 0 { return 0 - 1 }
1441 let n: i64 = szp[0]
1442 let rids: *i64 = sys_mmap(8 * 128) as *i64
1443 let cnt: i64 = wf_obs_runs(buf, n, rids, 128)
1444 if cnt < 0 { return 0 - 1 }
1445 let hb: *u8 = sys_mmap(WF_MAGIC_65536)
1446 var o: i64 = 0
1447 o = wf_cat(hb, o, "<" as *u8)
1448 hb[o] = 33 as u8
1449 o = o + 1
1450 o = wf_cat(hb, o, "DOCTYPE html><html><head><meta charset=utf-8><title>Nishi Workflow Runs</title><style>body{font-family:monospace;background:rgb(16,17,22);color:rgb(222,224,230);margin:2em}table{border-collapse:collapse}td,th{border:1px solid rgb(60,62,72);padding:6px 12px}h1{color:rgb(140,190,255)}.DONE{color:rgb(90,210,130)}.FAILED{color:rgb(245,95,95)}.PARKED{color:rgb(235,195,95)}.RUNNING{color:rgb(120,175,245)}</style></head><body><h1>Nishi Workflow Runs</h1><p>Derived by REPLAY of the event-sourced run ledger, nothing self-reported.</p><table><tr><th>run</th><th>state</th><th>ok steps</th><th>attempts</th></tr>" as *u8)
1451 let tok: *u8 = sys_mmap(256)
1452 var d: i64 = 0
1453 var f: i64 = 0
1454 var w: i64 = 0
1455 var ru: i64 = 0
1456 var i: i64 = 0
1457 while i < cnt {
1458 if o > WF_MAGIC_60000 { p("WFLOW obs ERROR page over cap -- fail loud\n" as *u8); return 0 - 1 }
1459 let rid: *u8 = rids[i] as *u8
1460 wf_obs_tok(tok, rid)
1461 let st: i64 = wf_run_state(buf, n, tok)
1462 let oks: i64 = wf_lines_with2(buf, n, tok, "status=OK" as *u8)
1463 let ats: i64 = wf_lines_with2(buf, n, tok, "status=ATT" as *u8)
1464 o = wf_cat(hb, o, "<tr><td>" as *u8)
1465 o = wf_cat(hb, o, rid)
1466 o = wf_cat(hb, o, "</td><td class=" as *u8)
1467 o = wf_cat(hb, o, wf_state_name(st))
1468 o = wf_cat(hb, o, ">" as *u8)
1469 o = wf_cat(hb, o, wf_state_name(st))
1470 o = wf_cat(hb, o, "</td><td>" as *u8)
1471 o = wf_catn(hb, o, oks)
1472 o = wf_cat(hb, o, "</td><td>" as *u8)
1473 o = wf_catn(hb, o, ats)
1474 o = wf_cat(hb, o, "</td></tr>" as *u8)
1475 if st == 3 { d = d + 1 }
1476 if st == 2 { f = f + 1 }
1477 if st == 1 { w = w + 1 }
1478 if st == 0 { ru = ru + 1 }
1479 i = i + 1
1480 }
1481 o = wf_cat(hb, o, "</table><p>runs " as *u8)
1482 o = wf_catn(hb, o, cnt)
1483 o = wf_cat(hb, o, " · done " as *u8)
1484 o = wf_catn(hb, o, d)
1485 o = wf_cat(hb, o, " · failed " as *u8)
1486 o = wf_catn(hb, o, f)
1487 o = wf_cat(hb, o, " · parked " as *u8)
1488 o = wf_catn(hb, o, w)
1489 o = wf_cat(hb, o, " · running " as *u8)
1490 o = wf_catn(hb, o, ru)
1491 o = wf_cat(hb, o, "</p><p>Generated by nx_wflow_engine wf_obs_html (sovereign, 0-JS, ledger replay).</p></body></html>" as *u8)
1492 if wf_writeall(out, hb, o) != 0 { return 0 - 1 }
1493 return cnt
1494}
1495
1496func wf_selftest() -> i64 {
1497 p("=== NX-WFLOW-ENGINE SELFTEST (R1 keystone: multi-step durable runs, retry, ledger-dedup, crash-resume) ===\n" as *u8)
1498 var ok: i64 = 1
1499
1500 // fresh unique ledger under /tmp (getpid is broken on this backend -- use epoch seconds)
1501 let led: *u8 = sys_mmap(128)
1502 var lo: i64 = 0
1503 lo = wf_cat(led, lo, "/tmp/wflow_st_" as *u8)
1504 lo = wf_catn(led, lo, wf_now())
1505 lo = wf_cat(led, lo, ".log" as *u8)
1506 led[lo] = 0 as u8
1507
1508 let flows: *i64 = sys_mmap(8 * 8) as *i64
1509 flows[0] = "f1|deal-stage|to=Negotiation" as *u8 as i64
1510 flows[1] = "f2|giving-received|-" as *u8 as i64
1511 flows[2] = "f3|ops-check|-" as *u8 as i64
1512 let steps: *i64 = sys_mmap(8 * 16) as *i64
1513 steps[0] = "f1|1|create-task|-|Prepare contract and pricing|1" as *u8 as i64
1514 steps[1] = "f1|2|send|card|Congrats on reaching Negotiation|2" as *u8 as i64
1515 steps[2] = "f1|3|notify|-|Deal advanced|1" as *u8 as i64
1516 steps[3] = "f2|1|probe-fail|-|3|4" as *u8 as i64
1517 steps[4] = "f2|2|update-field|-|status=thanked|1" as *u8 as i64
1518 steps[5] = "f3|1|probe-fail|-|9|2" as *u8 as i64
1519 steps[6] = "f3|2|notify|-|never reached|1" as *u8 as i64
1520
1521 let cx: *i64 = sys_mmap(32) as *i64
1522 cx[0] = steps as i64
1523 cx[1] = 7
1524 cx[2] = led as i64
1525 cx[3] = 0
1526
1527 // T1 whole-set validation, loud refusals
1528 let l1: i64 = wf_load(steps, 7)
1529 let l2: i64 = wf_load_flows(flows, 3)
1530 p(" T1 load steps=" as *u8); pn(l1); p(" flows=" as *u8); pn(l2); p("\n" as *u8)
1531 if l1 != 7 { ok = 0 }
1532 if l2 != 3 { ok = 0 }
1533 let bad1: *i64 = sys_mmap(16) as *i64
1534 bad1[0] = "f9|1|explode|-|boom|1" as *u8 as i64
1535 if wf_load(bad1, 1) != (0 - 1) { ok = 0 }
1536 let bad2: *i64 = sys_mmap(16) as *i64
1537 bad2[0] = "f9|1|send|fax|x|1" as *u8 as i64
1538 if wf_load(bad2, 1) != (0 - 1) { ok = 0 }
1539 let bad3: *i64 = sys_mmap(16) as *i64
1540 bad3[0] = "f9|1|notify|-|x|0" as *u8 as i64
1541 if wf_load(bad3, 1) != (0 - 1) { ok = 0 }
1542
1543 // T2 fire -> 3-step run completes in order, DONE recorded
1544 let r1: i64 = wf_fire(cx, flows, 3, "deal-stage~deal=Riverside Contract~to=Negotiation~id=E1" as *u8)
1545 let szp: *i64 = sys_mmap(8) as *i64
1546 var buf: *u8 = wf_readall(led, szp)
1547 let okE1: i64 = wf_lines_with2(buf, szp[0], " rid=E1.f1 " as *u8, "status=OK" as *u8)
1548 let dnE1: i64 = wf_lines_with2(buf, szp[0], " rid=E1.f1 " as *u8, "status=DONE" as *u8)
1549 p(" T2 fire started=" as *u8); pn(r1); p(" okSteps=" as *u8); pn(okE1); p(" done=" as *u8); pn(dnE1); p("\n" as *u8)
1550 if r1 != 1 { ok = 0 }
1551 if okE1 != 3 { ok = 0 }
1552 if dnE1 != 1 { ok = 0 }
1553
1554 // T3 durable dedup: same event id never starts twice (survives across processes via the ledger)
1555 let r2: i64 = wf_fire(cx, flows, 3, "deal-stage~deal=Riverside Contract~to=Negotiation~id=E1" as *u8)
1556 buf = wf_readall(led, szp)
1557 let stE1: i64 = wf_lines_with2(buf, szp[0], " rid=E1.f1 " as *u8, "status=START" as *u8)
1558 p(" T3 refire rc=" as *u8); pn(r2); p(" starts=" as *u8); pn(stE1); p("\n" as *u8)
1559 if r2 != (0 - 100) { ok = 0 }
1560 if stE1 != 1 { ok = 0 }
1561
1562 // T4 per-step retry: probe-fail 3 with max 4 succeeds on attempt 3
1563 let r3: i64 = wf_fire(cx, flows, 3, "giving-received~person=Rose~amount=500~id=E2" as *u8)
1564 buf = wf_readall(led, szp)
1565 let ok3: i64 = wf_lines_with2(buf, szp[0], " rid=E2.f2 " as *u8, "status=OK att=3" as *u8)
1566 let atts: i64 = wf_lines_with2(buf, szp[0], " rid=E2.f2 " as *u8, "status=ATT" as *u8)
1567 let dnE2: i64 = wf_lines_with2(buf, szp[0], " rid=E2.f2 " as *u8, "status=DONE" as *u8)
1568 p(" T4 retry started=" as *u8); pn(r3); p(" okAtt3=" as *u8); pn(ok3); p(" attempts=" as *u8); pn(atts); p(" done=" as *u8); pn(dnE2); p("\n" as *u8)
1569 if r3 != 1 { ok = 0 }
1570 if ok3 != 1 { ok = 0 }
1571 if atts != 4 { ok = 0 }
1572 if dnE2 != 1 { ok = 0 }
1573
1574 // T5 exhausted retries -> FAILED at step, later steps never run
1575 let r4: i64 = wf_fire(cx, flows, 3, "ops-check~probe=edge~id=E3" as *u8)
1576 buf = wf_readall(led, szp)
1577 let fsE3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, "status=FAILSTEP" as *u8)
1578 let flE3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, "status=FAILED" as *u8)
1579 let s2E3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, " step=2 " as *u8)
1580 let dnE3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, "status=DONE" as *u8)
1581 p(" T5 exhaust rc=" as *u8); pn(r4); p(" failstep=" as *u8); pn(fsE3); p(" failed=" as *u8); pn(flE3); p(" step2lines=" as *u8); pn(s2E3); p(" done=" as *u8); pn(dnE3); p("\n" as *u8)
1582 if fsE3 != 1 { ok = 0 }
1583 if flE3 != 1 { ok = 0 }
1584 if s2E3 != 0 { ok = 0 }
1585 if dnE3 != 0 { ok = 0 }
1586
1587 // T6 THE CROWN -- durable crash-resume: synthetic run crashed after step 1; resume completes 2..3
1588 // WITHOUT re-executing step 1 (event-sourced replay, additive only)
1589 wf_emit(led, "E9.f1" as *u8, "f1" as *u8, 0, "START" as *u8, 0)
1590 wf_emit(led, "E9.f1" as *u8, "f1" as *u8, 1, "ATT" as *u8, 1)
1591 wf_emit(led, "E9.f1" as *u8, "f1" as *u8, 1, "OK" as *u8, 1)
1592 let rs: i64 = wf_resume(cx)
1593 buf = wf_readall(led, szp)
1594 let okE9: i64 = wf_lines_with2(buf, szp[0], " rid=E9.f1 " as *u8, "status=OK" as *u8)
1595 let ok1E9: i64 = wf_lines_with2(buf, szp[0], " rid=E9.f1 " as *u8, " step=1 status=OK" as *u8)
1596 let dnE9: i64 = wf_lines_with2(buf, szp[0], " rid=E9.f1 " as *u8, "status=DONE" as *u8)
1597 p(" T6 resume resumed=" as *u8); pn(rs); p(" okSteps=" as *u8); pn(okE9); p(" step1okOnce=" as *u8); pn(ok1E9); p(" done=" as *u8); pn(dnE9); p("\n" as *u8)
1598 if rs != 1 { ok = 0 }
1599 if okE9 != 3 { ok = 0 }
1600 if ok1E9 != 1 { ok = 0 }
1601 if dnE9 != 1 { ok = 0 }
1602
1603 // T7 liar-kill negatives: unknown-flow in-flight run refuses loudly (no fabricated DONE);
1604 // nonsense status never appears
1605 wf_emit(led, "EX.zz" as *u8, "zz" as *u8, 0, "START" as *u8, 0)
1606 let rz: i64 = wf_resume(cx)
1607 buf = wf_readall(led, szp)
1608 let dnZZ: i64 = wf_lines_with2(buf, szp[0], " rid=EX.zz " as *u8, "status=DONE" as *u8)
1609 let ban: i64 = wf_lines_with2(buf, szp[0], "WFRUN" as *u8, "status=BANANA" as *u8)
1610 p(" T7 negctl resumeRc=" as *u8); pn(rz); p(" zzDone=" as *u8); pn(dnZZ); p(" nonsense=" as *u8); pn(ban); p("\n" as *u8)
1611 if rz != (0 - 1) { ok = 0 }
1612 if dnZZ != 0 { ok = 0 }
1613 if ban != 0 { ok = 0 }
1614
1615 p("NX-WFLOW-ENGINE-SELFTEST ledger=" as *u8); p(led)
1616 p(" runs: done=3 failed=1 dedup=1 resumed=1 " as *u8)
1617 if ok == 1 { p("verdict=GREEN\n" as *u8); return 0 }
1618 p("verdict=RED\n" as *u8)
1619 return 1
1620}