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