code wiki / (root) / nx_wflow_engine.nx

nx_wflow_engine.nx source

↩ module page · 1618 lines · 68064 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 sys_wait4(pid, st, 0) 515 let ec: i64 = (st[0] >> 8) & 0xff 516 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) 517 if ec != 0 { return 0 } 518 if docap == 1 { 519 if total > 480 { p("WFLOW capture ERROR output over 480 bytes -- fail loud\n" as *u8); return 0 } 520 var t: i64 = 0 521 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 } 522 while total > 0 { if cbuf[total-1] == (32 as u8) { total = total - 1 } else { break } } 523 cbuf[total] = 0 as u8 524 let ledc: *u8 = cx[2] as *u8 525 wf_emit_var(ledc, rid, capv, cbuf) 526 p("WFLOW-CAPTURE run=" as *u8); p(rid); p(" var=" as *u8); p(capv); p(" v=" as *u8); p(cbuf); p("\n" as *u8) 527 } 528 return 1 529} 530 531// ---- R11 loops: apply-to-each over the connector plane ---- 532func wf_emit_iter(led: *u8, rid: *u8, n: i64, val: *u8) -> i64 { 533 let ln: *u8 = sys_mmap(768) 534 var o: i64 = 0 535 o = wf_cat(ln, o, "WFITER rid=" as *u8) 536 o = wf_cat(ln, o, rid) 537 o = wf_cat(ln, o, " n=" as *u8) 538 o = wf_catn(ln, o, n) 539 o = wf_cat(ln, o, " item=" as *u8) 540 o = wf_cat(ln, o, val) 541 ln[o] = 10 as u8 542 ln[o+1] = 0 as u8 543 return wf_append(led, ln) 544} 545// for-each: arg = <listexpr>~<per-item template> (template default {item}). The list substitutes ONCE 546// (comma-split, empties skipped, 32-item loud cap); per item: bind {item} as a run var, emit WFITER, 547// substitute the template, fork the connector. Any failed item fails the STEP (retry/on-fail apply). 548// Empty list = zero iterations, step OK (apply-to-each on an empty collection is a no-op). 549func wf_for_each(cx: *i64, chspec: *u8, argraw: *u8, rid: *u8) -> i64 { 550 let led: *u8 = cx[2] as *u8 551 let lb: *u8 = sys_mmap(WF_MAGIC_1024) 552 let tpl: *u8 = sys_mmap(512) 553 var i: i64 = 0 554 var sep: i64 = 0 - 1 555 while argraw[i] != (0 as u8) { if sep < 0 { if argraw[i] == (126 as u8) { sep = i } } i = i + 1 } 556 let alen: i64 = i 557 var lend: i64 = alen 558 if sep >= 0 { lend = sep } 559 var li: i64 = 0 560 while li < lend { lb[li] = argraw[li]; li = li + 1 } 561 lb[li] = 0 as u8 562 if sep >= 0 { 563 var ti: i64 = 0 564 var q: i64 = sep + 1 565 while q < alen { if ti < 510 { tpl[ti] = argraw[q]; ti = ti + 1 } q = q + 1 } 566 tpl[ti] = 0 as u8 567 } else { 568 var to: i64 = 0 569 to = wf_cat(tpl, to, "{item}" as *u8) 570 tpl[to] = 0 as u8 571 } 572 let lbx: *u8 = sys_mmap(WF_MAGIC_1024) 573 if wf_subst(cx, rid, lb, lbx, WF_MAGIC_1024) < 0 { return 0 } 574 var ln2: i64 = 0 575 while lbx[ln2] != (0 as u8) { ln2 = ln2 + 1 } 576 let iv: *i64 = sys_mmap(8 * 33) as *i64 577 var ni: i64 = 0 578 var st0: i64 = 0 579 var j: i64 = 0 580 var over: i64 = 0 581 while j <= ln2 { 582 var cut: i64 = 0 583 if j == ln2 { cut = 1 } 584 if cut == 0 { if lbx[j] == (44 as u8) { cut = 1 } } 585 if cut == 1 { 586 lbx[j] = 0 as u8 587 if lbx[st0] != (0 as u8) { 588 if ni >= 32 { over = 1 } else { iv[ni] = (lbx as i64) + st0; ni = ni + 1 } 589 } 590 st0 = j + 1 591 } 592 j = j + 1 593 } 594 if over == 1 { p("WFLOW for-each ERROR over 32 items -- fail loud\n" as *u8); return 0 } 595 if ni == 0 { p("WFLOW-FOREACH empty list, zero iterations\n" as *u8); return 1 } 596 let ibuf: *u8 = sys_mmap(WF_MAGIC_1024) 597 var k: i64 = 0 598 while k < ni { 599 let itv: *u8 = iv[k] as *u8 600 wf_emit_var(led, rid, "item" as *u8, itv) 601 wf_emit_iter(led, rid, k + 1, itv) 602 if wf_subst(cx, rid, tpl, ibuf, WF_MAGIC_1024) < 0 { return 0 } 603 p("WFLOW-FOREACH n=" as *u8); pn(k + 1); p(" item=" as *u8); p(itv); p("\n" as *u8) 604 if wf_exec_organ(cx, chspec, ibuf, rid) == 0 { p("WFLOW for-each item FAILED -- step fails\n" as *u8); return 0 } 605 k = k + 1 606 } 607 return 1 608} 609 610// ---- action execution (in-process v1 + exec-organ; send-channel semantics from nx_send) ---- 611func wf_exec(cx: *i64, action: *u8, ch: *u8, arg: *u8, att: i64, rid: *u8) -> i64 { 612 p("WFLOW-ACTION action=" as *u8); p(action) 613 if seq(action, "send" as *u8) == 1 { p(" channel=" as *u8); p(ch); p(" route=" as *u8); p(sr_route(ch)) } 614 p(" arg=" as *u8); p(arg) 615 p(" att=" as *u8); pn(att) 616 p("\n" as *u8) 617 if seq(action, "probe-fail" as *u8) == 1 { 618 let k: i64 = wf_atoi(arg) 619 if att < k { return 0 } 620 return 1 621 } 622 if seq(action, "exec-organ" as *u8) == 1 { return wf_exec_organ(cx, ch, arg, rid) } 623 if seq(action, "for-each" as *u8) == 1 { return wf_for_each(cx, ch, arg, rid) } 624 return 1 625} 626 627// ---- R4 human-in-the-loop approvals: WFDEC decision records on the same ledger ---- 628// WFDEC rid=<rid> step=<n> decision=APPROVED|DENIED by=<who>. DENY WINS if both exist (fail-safe, 629// deny-by-default). No decision = the approve step PARKS the run (status=WAIT, in-flight); wf_resume 630// re-examines it and completes once a decision lands. 631func wf_decision(cx: *i64, rid: *u8, idx: i64) -> i64 { 632 let led: *u8 = cx[2] as *u8 633 let szp: *i64 = sys_mmap(8) as *i64 634 let buf: *u8 = wf_readall(led, szp) 635 if (buf as i64) == 0 { return 0 } 636 let tok: *u8 = sys_mmap(256) 637 var q: i64 = 0 638 q = wf_cat(tok, q, " rid=" as *u8) 639 q = wf_cat(tok, q, rid) 640 tok[q] = 32 as u8 641 tok[q+1] = 0 as u8 642 let nb: *u8 = sys_mmap(128) 643 var o: i64 = 0 644 o = wf_cat(nb, o, " step=" as *u8) 645 o = wf_catn(nb, o, idx) 646 o = wf_cat(nb, o, " decision=DENIED" as *u8) 647 nb[o] = 0 as u8 648 if wf_lines_with2(buf, szp[0], tok, nb) > 0 { return 0 - 1 } 649 var o2: i64 = 0 650 o2 = wf_cat(nb, o2, " step=" as *u8) 651 o2 = wf_catn(nb, o2, idx) 652 o2 = wf_cat(nb, o2, " decision=APPROVED" as *u8) 653 nb[o2] = 0 as u8 654 if wf_lines_with2(buf, szp[0], tok, nb) > 0 { return 1 } 655 return 0 656} 657// the operator surface: append a decision record (only APPROVED or DENIED accepted, loud else) 658func wf_decide(cx: *i64, rid: *u8, idx: i64, decision: *u8, who: *u8) -> i64 { 659 var okd: i64 = 0 660 if seq(decision, "APPROVED" as *u8) == 1 { okd = 1 } 661 if seq(decision, "DENIED" as *u8) == 1 { okd = 1 } 662 if okd == 0 { p("WFLOW decide ERROR decision must be APPROVED or DENIED -- fail loud\n" as *u8); return 0 - 1 } 663 let led: *u8 = cx[2] as *u8 664 let ln: *u8 = sys_mmap(512) 665 var o: i64 = 0 666 o = wf_cat(ln, o, "WFDEC rid=" as *u8) 667 o = wf_cat(ln, o, rid) 668 o = wf_cat(ln, o, " step=" as *u8) 669 o = wf_catn(ln, o, idx) 670 o = wf_cat(ln, o, " decision=" as *u8) 671 o = wf_cat(ln, o, decision) 672 o = wf_cat(ln, o, " by=" as *u8) 673 o = wf_cat(ln, o, who) 674 ln[o] = 10 as u8 675 ln[o+1] = 0 as u8 676 return wf_append(led, ln) 677} 678 679// ---- R8 data passing: run-scoped variables on the SAME event-sourced ledger ---- 680// WFVAR rid=<rid> k=<key> v=<value-to-end-of-line>. Trigger-event fields auto-bind at fire; a set-var 681// step writes derived vars; {key} placeholders in any step arg substitute at execution; LAST write wins; 682// unknown placeholder = LOUD step failure (the send_merge law). Values must not contain ~ or newline. 683func wf_emit_var(led: *u8, rid: *u8, k: *u8, v: *u8) -> i64 { 684 let ln: *u8 = sys_mmap(768) 685 var o: i64 = 0 686 o = wf_cat(ln, o, "WFVAR rid=" as *u8) 687 o = wf_cat(ln, o, rid) 688 o = wf_cat(ln, o, " k=" as *u8) 689 o = wf_cat(ln, o, k) 690 o = wf_cat(ln, o, " v=" as *u8) 691 o = wf_cat(ln, o, v) 692 ln[o] = 10 as u8 693 ln[o+1] = 0 as u8 694 return wf_append(led, ln) 695} 696// copy value after key to END OF LINE (values may contain spaces) 697func wf_val_line_end(buf: *u8, ls: i64, le: i64, key: *u8, out: *u8, cap: i64) -> i64 { 698 let m: i64 = slen(key) 699 var i: i64 = ls 700 while i + m <= le { 701 var k: i64 = 0 702 var hit: i64 = 1 703 while k < m { if buf[i+k] != key[k] { hit = 0; k = m } else { k = k + 1 } } 704 if hit == 1 { 705 var q: i64 = i + m 706 var t: i64 = 0 707 while q < le { if t < cap - 1 { out[t] = buf[q]; t = t + 1 } q = q + 1 } 708 out[t] = 0 as u8 709 return 1 710 } 711 i = i + 1 712 } 713 out[0] = 0 as u8 714 return 0 715} 716func wf_obs_tok_fwd(tok: *u8, rid: *u8) -> i64 { 717 var q: i64 = 0 718 q = wf_cat(tok, q, " rid=" as *u8) 719 q = wf_cat(tok, q, rid) 720 tok[q] = 32 as u8 721 tok[q+1] = 0 as u8 722 return q 723} 724// latest WFVAR value for (rid,key); 1 found / 0 not 725func wf_var_get(cx: *i64, rid: *u8, key: *u8, out: *u8, cap: i64) -> i64 { 726 let led: *u8 = cx[2] as *u8 727 let szp: *i64 = sys_mmap(8) as *i64 728 let buf: *u8 = wf_readall(led, szp) 729 if (buf as i64) == 0 { return 0 } 730 let n: i64 = szp[0] 731 let tok: *u8 = sys_mmap(256) 732 wf_obs_tok_fwd(tok, rid) 733 let kt: *u8 = sys_mmap(160) 734 var q: i64 = 0 735 q = wf_cat(kt, q, " k=" as *u8) 736 q = wf_cat(kt, q, key) 737 kt[q] = 32 as u8 738 kt[q+1] = 0 as u8 739 var found: i64 = 0 740 var ls: i64 = 0 741 while ls < n { 742 var le: i64 = ls 743 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 744 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 { 745 wf_val_line_end(buf, ls, le, " v=" as *u8, out, cap) 746 found = 1 747 } } } 748 ls = le + 1 749 } 750 return found 751} 752// substitute {key} placeholders from the run's variable plane; -1 LOUD on unknown or unclosed 753func wf_subst(cx: *i64, rid: *u8, src: *u8, out: *u8, cap: i64) -> i64 { 754 let kb: *u8 = sys_mmap(128) 755 let vb: *u8 = sys_mmap(512) 756 var i: i64 = 0 757 var o: i64 = 0 758 while src[i] != (0 as u8) { 759 if src[i] == (123 as u8) { 760 var j: i64 = i + 1 761 var kn: i64 = 0 762 var closed: i64 = 0 763 while src[j] != (0 as u8) { 764 if src[j] == (125 as u8) { closed = 1; break } 765 if kn < 127 { kb[kn] = src[j]; kn = kn + 1 } 766 j = j + 1 767 } 768 kb[kn] = 0 as u8 769 if closed == 0 { p("WFLOW subst ERROR unclosed placeholder -- fail loud\n" as *u8); return 0 - 1 } 770 if wf_var_get(cx, rid, kb, vb, 512) == 0 { 771 p("WFLOW subst ERROR unknown variable=" as *u8); p(kb); p(" -- fail loud\n" as *u8) 772 return 0 - 1 773 } 774 var t: i64 = 0 775 while vb[t] != (0 as u8) { if o < cap - 1 { out[o] = vb[t]; o = o + 1 } t = t + 1 } 776 i = j + 1 777 } else { 778 if o < cap - 1 { out[o] = src[i]; o = o + 1 } 779 i = i + 1 780 } 781 } 782 out[o] = 0 as u8 783 return o 784} 785// bind every k=v field of the trigger event as an initial run variable (segment 0 = kind, skipped) 786func wf_bind_event(cx: *i64, rid: *u8, ev: *u8) -> i64 { 787 let led: *u8 = cx[2] as *u8 788 let kb: *u8 = sys_mmap(160) 789 let vb: *u8 = sys_mmap(512) 790 var seg: i64 = 0 791 var i: i64 = 0 792 var bound: i64 = 0 793 while ev[i] != (0 as u8) { 794 if ev[i] == (126 as u8) { 795 seg = seg + 1 796 var j: i64 = i + 1 797 var kn: i64 = 0 798 var haseq: i64 = 0 799 while ev[j] != (0 as u8) { 800 if ev[j] == (126 as u8) { break } 801 if ev[j] == (61 as u8) { haseq = 1; j = j + 1; break } 802 if kn < 158 { kb[kn] = ev[j]; kn = kn + 1 } 803 j = j + 1 804 } 805 kb[kn] = 0 as u8 806 if haseq == 1 { 807 var vn: i64 = 0 808 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 } 809 vb[vn] = 0 as u8 810 if kn > 0 { wf_emit_var(led, rid, kb, vb); bound = bound + 1 } 811 } 812 } 813 i = i + 1 814 } 815 return bound 816} 817 818// ---- R12 versioning: @version N directive in the flows file; runs stamp WFVER; resume flags drift ---- 819func wf_defs_version(path: *u8) -> i64 { 820 let szp: *i64 = sys_mmap(8) as *i64 821 let buf: *u8 = wf_readall(path, szp) 822 if (buf as i64) == 0 { return 0 } 823 let n: i64 = szp[0] 824 var ls: i64 = 0 825 while ls < n { 826 var le: i64 = ls 827 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 828 if wf_has_rng(buf, ls, le, "@version " as *u8) == 1 { 829 let v: i64 = wf_num_after(buf, ls, le, "@version " as *u8) 830 if v > 0 { return v } 831 } 832 ls = le + 1 833 } 834 return 0 835} 836func wf_emit_ver(led: *u8, rid: *u8, ver: i64) -> i64 { 837 let ln: *u8 = sys_mmap(512) 838 var o: i64 = 0 839 o = wf_cat(ln, o, "WFVER rid=" as *u8) 840 o = wf_cat(ln, o, rid) 841 o = wf_cat(ln, o, " ver=" as *u8) 842 o = wf_catn(ln, o, ver) 843 ln[o] = 10 as u8 844 ln[o+1] = 0 as u8 845 return wf_append(led, ln) 846} 847func wf_emit_drift(led: *u8, rid: *u8, ranv: i64, nowv: i64) -> i64 { 848 let ln: *u8 = sys_mmap(512) 849 var o: i64 = 0 850 o = wf_cat(ln, o, "WFVERDRIFT rid=" as *u8) 851 o = wf_cat(ln, o, rid) 852 o = wf_cat(ln, o, " ran=" as *u8) 853 o = wf_catn(ln, o, ranv) 854 o = wf_cat(ln, o, " now=" as *u8) 855 o = wf_catn(ln, o, nowv) 856 ln[o] = 10 as u8 857 ln[o+1] = 0 as u8 858 return wf_append(led, ln) 859} 860// the version a run was STARTED under (last WFVER line for the rid; 0 = unversioned) 861func wf_run_ver(buf: *u8, n: i64, ridtok: *u8) -> i64 { 862 var v: i64 = 0 863 var ls: i64 = 0 864 while ls < n { 865 var le: i64 = ls 866 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 867 if wf_has_rng(buf, ls, le, "WFVER " as *u8) == 1 { if wf_has_rng(buf, ls, le, ridtok) == 1 { 868 let x: i64 = wf_num_after(buf, ls, le, " ver=" as *u8) 869 if x > 0 { v = x } 870 } } 871 ls = le + 1 872 } 873 return v 874} 875 876// ---- R13 templates: gallery = <dir>/index (name|description rows) + <name>.flows/.steps pairs ---- 877func wf_tpl_exists(path: *u8) -> i64 { 878 let fd: i64 = __syscall(257, 0 - 100, path as i64, 0, 0, 0, 0) 879 if fd < 0 { return 0 } 880 __syscall(3, fd, 0, 0, 0, 0, 0) 881 return 1 882} 883func wf_tpl_path(dir: *u8, name: *u8, ext: *u8, out: *u8) -> i64 { 884 var o: i64 = 0 885 o = wf_cat(out, o, dir) 886 out[o] = 47 as u8 887 o = o + 1 888 o = wf_cat(out, o, name) 889 o = wf_cat(out, o, ext) 890 out[o] = 0 as u8 891 return o 892} 893func wf_tpl_index_path(dir: *u8, out: *u8) -> i64 { 894 var o: i64 = 0 895 o = wf_cat(out, o, dir) 896 out[o] = 47 as u8 897 o = o + 1 898 o = wf_cat(out, o, "index" as *u8) 899 out[o] = 0 as u8 900 return o 901} 902// list the gallery: every index row printed with file-existence + version; returns USABLE count, -1 loud 903func wf_tpl_list(dir: *u8) -> i64 { 904 let ip: *u8 = sys_mmap(512) 905 wf_tpl_index_path(dir, ip) 906 let szp: *i64 = sys_mmap(8) as *i64 907 let buf: *u8 = wf_readall(ip, szp) 908 if (buf as i64) == 0 { return 0 - 1 } 909 let n: i64 = szp[0] 910 if n == 0 { p("WFLOW templates ERROR missing gallery index -- fail loud\n" as *u8); return 0 - 1 } 911 let nm: *u8 = sys_mmap(128) 912 let ds: *u8 = sys_mmap(512) 913 let fp2: *u8 = sys_mmap(512) 914 let sp2: *u8 = sys_mmap(512) 915 var cnt: i64 = 0 916 var ls: i64 = 0 917 while ls < n { 918 var le: i64 = ls 919 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 920 var keep: i64 = 0 921 if le > ls { keep = 1 } 922 if keep == 1 { if buf[ls] == (35 as u8) { keep = 0 } } 923 if keep == 1 { 924 let row: *u8 = sys_mmap(le - ls + 2) 925 var k: i64 = 0 926 while ls + k < le { row[k] = buf[ls+k]; k = k + 1 } 927 row[k] = 0 as u8 928 pipe_field(row, 0, nm, 128) 929 pipe_field(row, 1, ds, 512) 930 wf_tpl_path(dir, nm, ".flows" as *u8, fp2) 931 wf_tpl_path(dir, nm, ".steps" as *u8, sp2) 932 let fe: i64 = wf_tpl_exists(fp2) 933 let se: i64 = wf_tpl_exists(sp2) 934 let v: i64 = wf_defs_version(fp2) 935 p("WFTPL name=" as *u8); p(nm) 936 p(" ver=" as *u8); pn(v) 937 p(" flows=" as *u8); pn(fe) 938 p(" steps=" as *u8); pn(se) 939 p(" desc=" as *u8); p(ds) 940 p("\n" as *u8) 941 if fe == 1 { if se == 1 { cnt = cnt + 1 } } 942 } 943 ls = le + 1 944 } 945 p("WFTPL-TOTAL usable=" as *u8); pn(cnt); p("\n" as *u8) 946 return cnt 947} 948func wf_tpl_indexed(dir: *u8, name: *u8) -> i64 { 949 let ip: *u8 = sys_mmap(512) 950 wf_tpl_index_path(dir, ip) 951 let szp: *i64 = sys_mmap(8) as *i64 952 let buf: *u8 = wf_readall(ip, szp) 953 if (buf as i64) == 0 { return 0 - 1 } 954 let n: i64 = szp[0] 955 if n == 0 { return 0 - 1 } 956 let nm: *u8 = sys_mmap(128) 957 var ls: i64 = 0 958 while ls < n { 959 var le: i64 = ls 960 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 961 if le > ls { if buf[ls] != (35 as u8) { 962 let row: *u8 = sys_mmap(le - ls + 2) 963 var k: i64 = 0 964 while ls + k < le { row[k] = buf[ls+k]; k = k + 1 } 965 row[k] = 0 as u8 966 pipe_field(row, 0, nm, 128) 967 if seq(nm, name) == 1 { return 1 } 968 } } 969 ls = le + 1 970 } 971 return 0 972} 973// instantiate: fail-closed to INDEXED templates; byte-copy both files; whole-set validate the COPY 974// (a broken template never lands silently). Returns the template version (0 = unversioned), -1 loud. 975func wf_tpl_instantiate(dir: *u8, name: *u8, df: *u8, ds: *u8) -> i64 { 976 let ix: i64 = wf_tpl_indexed(dir, name) 977 if ix != 1 { p("WFLOW instantiate ERROR template not in gallery index -- fail closed\n" as *u8); return 0 - 1 } 978 let fp2: *u8 = sys_mmap(512) 979 let sp2: *u8 = sys_mmap(512) 980 wf_tpl_path(dir, name, ".flows" as *u8, fp2) 981 wf_tpl_path(dir, name, ".steps" as *u8, sp2) 982 let szp: *i64 = sys_mmap(8) as *i64 983 let fb: *u8 = wf_readall(fp2, szp) 984 if (fb as i64) == 0 { return 0 - 1 } 985 if szp[0] == 0 { p("WFLOW instantiate ERROR template flows file missing -- fail loud\n" as *u8); return 0 - 1 } 986 if wf_writeall(df, fb, szp[0]) != 0 { return 0 - 1 } 987 let sb: *u8 = wf_readall(sp2, szp) 988 if (sb as i64) == 0 { return 0 - 1 } 989 if szp[0] == 0 { p("WFLOW instantiate ERROR template steps file missing -- fail loud\n" as *u8); return 0 - 1 } 990 if wf_writeall(ds, sb, szp[0]) != 0 { return 0 - 1 } 991 let fl: *i64 = sys_mmap(8 * 64) as *i64 992 let st2: *i64 = sys_mmap(8 * 128) as *i64 993 let nf: i64 = wf_lines_load(df, fl, 64) 994 if nf < 1 { return 0 - 1 } 995 let nst: i64 = wf_lines_load(ds, st2, 128) 996 if nst < 1 { return 0 - 1 } 997 if wf_load_flows(fl, nf) < 0 { return 0 - 1 } 998 if wf_load(st2, nst) < 0 { return 0 - 1 } 999 let v: i64 = wf_defs_version(df) 1000 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) 1001 return v 1002} 1003 1004// R9 branching: per-step condition k=v against the run variable plane. Missing variable or 1005// mismatch = 0 (the step SKIPs, recorded in the ledger); match = 1. Equality only, v1. 1006func wf_cond_ok(cx: *i64, rid: *u8, cond: *u8) -> i64 { 1007 let kb: *u8 = sys_mmap(160) 1008 let want: *u8 = sys_mmap(512) 1009 var i: i64 = 0 1010 var kn: i64 = 0 1011 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 } 1012 kb[kn] = 0 as u8 1013 if cond[i] != (61 as u8) { return 0 } 1014 if kn == 0 { return 0 } 1015 var wn: i64 = 0 1016 var q: i64 = i + 1 1017 while cond[q] != (0 as u8) { if wn < 510 { want[wn] = cond[q]; wn = wn + 1 } q = q + 1 } 1018 want[wn] = 0 as u8 1019 let vb: *u8 = sys_mmap(512) 1020 if wf_var_get(cx, rid, kb, vb, 512) == 0 { return 0 } 1021 return seq(vb, want) 1022} 1023 1024// ---- the durable executor: run flow fid for run rid, starting AFTER step s0 ---- 1025// cx bundle: cx[0]=steps ptr, cx[1]=nsteps, cx[2]=ledger path 1026func wf_run_from(cx: *i64, rid: *u8, fid: *u8, s0: i64) -> i64 { 1027 let steps: *i64 = cx[0] as *i64 1028 let ns: i64 = cx[1] 1029 let led: *u8 = cx[2] as *u8 1030 let f2: *u8 = sys_mmap(64) 1031 let ix: *u8 = sys_mmap(32) 1032 let a: *u8 = sys_mmap(64) 1033 let ch: *u8 = sys_mmap(64) 1034 let arg: *u8 = sys_mmap(512) 1035 let ma: *u8 = sys_mmap(32) 1036 var i: i64 = 0 1037 var done: i64 = 0 1038 var failed: i64 = 0 1039 var parked: i64 = 0 1040 var jumpto: i64 = 0 1041 while i < ns { 1042 let r: *u8 = steps[i] as *u8 1043 pipe_field(r, 0, f2, 64) 1044 var use: i64 = 0 1045 if seq(f2, fid) == 1 { use = 1 } 1046 var idx: i64 = 0 1047 if use == 1 { pipe_field(r, 1, ix, 32); idx = wf_atoi(ix); if idx <= s0 { use = 0 } } 1048 // R10: an active on-fail jump skips the remaining normal-path steps (recorded) until the target 1049 if use == 1 { if jumpto > 0 { if idx < jumpto { wf_emit(led, rid, fid, idx, "SKIP" as *u8, 0); use = 0 } else { jumpto = 0 } } } 1050 if use == 1 { 1051 pipe_field(r, 2, a, 64) 1052 pipe_field(r, 3, ch, 64) 1053 pipe_field(r, 4, arg, 512) 1054 pipe_field(r, 5, ma, 32) 1055 // R9: optional per-step condition (7th field, k=v) vs the variable plane -> SKIP on mismatch 1056 let cond: *u8 = sys_mmap(256) 1057 pipe_field(r, 6, cond, 256) 1058 let argx: *u8 = sys_mmap(WF_MAGIC_1024) 1059 var mode: i64 = 0 1060 var docond: i64 = 0 1061 if cond[0] != (0 as u8) { docond = 1 } 1062 if cond[0] == (45 as u8) { if cond[1] == (0 as u8) { docond = 0 } } 1063 if docond == 1 { if wf_cond_ok(cx, rid, cond) == 0 { mode = 4 } } 1064 // R8: resolve {key} placeholders from the run variable plane; unknown = LOUD step failure. 1065 // for-each bypasses step-level subst ({item} binds inside the loop) -- the loop substitutes 1066 // the list once and the per-item template each iteration. 1067 if mode == 0 { 1068 var rawarg: i64 = 0 1069 if seq(a, "for-each" as *u8) == 1 { rawarg = 1 } 1070 if rawarg == 1 { 1071 var cc: i64 = 0 1072 while arg[cc] != (0 as u8) { argx[cc] = arg[cc]; cc = cc + 1 } 1073 argx[cc] = 0 as u8 1074 } else { 1075 if wf_subst(cx, rid, arg, argx, WF_MAGIC_1024) < 0 { mode = 3 } 1076 } 1077 } 1078 if mode == 4 { 1079 p("WFLOW-SKIP run=" as *u8); p(rid); p(" step=" as *u8); pn(idx); p(" cond=" as *u8); p(cond); p("\n" as *u8) 1080 wf_emit(led, rid, fid, idx, "SKIP" as *u8, 0) 1081 } 1082 if mode == 0 { if seq(a, "set-var" as *u8) == 1 { mode = 2 } } 1083 if mode == 0 { if seq(a, "approve" as *u8) == 1 { mode = 1 } } 1084 if mode == 3 { 1085 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, 0) 1086 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0) 1087 failed = 1 1088 i = ns 1089 } 1090 if mode == 2 { 1091 let kb: *u8 = sys_mmap(160) 1092 let vb: *u8 = sys_mmap(512) 1093 var e: i64 = 0 1094 var kn: i64 = 0 1095 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 } 1096 kb[kn] = 0 as u8 1097 if argx[e] == (61 as u8) { if kn > 0 { 1098 var vn: i64 = 0 1099 var q2: i64 = e + 1 1100 while argx[q2] != (0 as u8) { if vn < 510 { vb[vn] = argx[q2]; vn = vn + 1 } q2 = q2 + 1 } 1101 vb[vn] = 0 as u8 1102 wf_emit_var(led, rid, kb, vb) 1103 p("WFLOW-SETVAR run=" as *u8); p(rid); p(" k=" as *u8); p(kb); p(" v=" as *u8); p(vb); p("\n" as *u8) 1104 wf_emit(led, rid, fid, idx, "OK" as *u8, 1) 1105 done = done + 1 1106 } else { mode = 3 } } else { mode = 3 } 1107 if mode == 3 { 1108 p("WFLOW set-var ERROR needs key=value -- fail loud\n" as *u8) 1109 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, 0) 1110 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0) 1111 failed = 1 1112 i = ns 1113 } 1114 } 1115 if mode == 1 { 1116 let dec: i64 = wf_decision(cx, rid, idx) 1117 if dec == 1 { 1118 p("WFLOW-APPROVAL granted run=" as *u8); p(rid); p(" step=" as *u8); pn(idx); p("\n" as *u8) 1119 wf_emit(led, rid, fid, idx, "OK" as *u8, 1) 1120 done = done + 1 1121 } 1122 if dec == (0 - 1) { 1123 p("WFLOW-APPROVAL DENIED run=" as *u8); p(rid); p(" step=" as *u8); pn(idx); p(" -- run fails (deny wins)\n" as *u8) 1124 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, 0) 1125 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0) 1126 failed = 1 1127 i = ns 1128 } 1129 if dec == 0 { 1130 p("WFLOW-APPROVAL pending run=" as *u8); p(rid); p(" step=" as *u8); pn(idx) 1131 p(" approver=" as *u8); p(ch); p(" what=" as *u8); p(argx); p(" -- run PARKED\n" as *u8) 1132 wf_emit(led, rid, fid, idx, "WAIT" as *u8, 0) 1133 parked = 1 1134 i = ns 1135 } 1136 } 1137 if mode == 0 { 1138 let maxa: i64 = wf_atoi(ma) 1139 var att: i64 = 1 1140 var okd: i64 = 0 1141 while att <= maxa { 1142 wf_emit(led, rid, fid, idx, "ATT" as *u8, att) 1143 let rc: i64 = wf_exec(cx, a, ch, argx, att, rid) 1144 if rc == 1 { wf_emit(led, rid, fid, idx, "OK" as *u8, att); okd = 1; att = maxa + 1 } else { att = att + 1 } 1145 } 1146 if okd == 0 { 1147 wf_emit(led, rid, fid, idx, "FAILSTEP" as *u8, maxa) 1148 // R10: on-fail routing (8th field) -- exhausted ACTION failures only; the failure stays 1149 // recorded (FAILSTEP + ONFAIL), wf_error vars land on the plane, execution jumps forward 1150 let onf: *u8 = sys_mmap(32) 1151 pipe_field(r, 7, onf, 32) 1152 var tj: i64 = 0 1153 var haveof: i64 = 0 1154 if onf[0] != (0 as u8) { haveof = 1 } 1155 if onf[0] == (45 as u8) { if onf[1] == (0 as u8) { haveof = 0 } } 1156 if haveof == 1 { tj = wf_atoi(onf) } 1157 if tj > 0 { 1158 wf_emit(led, rid, fid, idx, "ONFAIL" as *u8, tj) 1159 wf_emit_var(led, rid, "wf_error" as *u8, "yes" as *u8) 1160 let sv: *u8 = sys_mmap(32) 1161 var so: i64 = 0 1162 so = wf_cat(sv, so, "step" as *u8) 1163 so = wf_catn(sv, so, idx) 1164 sv[so] = 0 as u8 1165 wf_emit_var(led, rid, "wf_error_step" as *u8, sv) 1166 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) 1167 jumpto = tj 1168 } else { 1169 wf_emit(led, rid, fid, idx, "FAILED" as *u8, 0) 1170 failed = 1 1171 i = ns 1172 } 1173 } else { done = done + 1 } 1174 } 1175 } 1176 i = i + 1 1177 } 1178 if failed == 1 { return 0 - 2 } 1179 if parked == 1 { return 0 - 3 } 1180 wf_emit(led, rid, fid, 0, "DONE" as *u8, 0) 1181 return done 1182} 1183 1184// ---- fire an event: durable ledger-deduped, one run per matched flow ---- 1185func wf_fire(cx: *i64, flows: *i64, nf: i64, ev: *u8) -> i64 { 1186 let evid: *u8 = sys_mmap(128) 1187 if wf_evid(ev, evid, 128) == 0 { p("WFLOW fire ERROR event has no id -- fail loud\n" as *u8); return 0 - 1 } 1188 let led: *u8 = cx[2] as *u8 1189 let szp: *i64 = sys_mmap(8) as *i64 1190 let buf: *u8 = wf_readall(led, szp) 1191 if (buf as i64) == 0 { return 0 - 1 } 1192 let fid: *u8 = sys_mmap(64) 1193 let rid: *u8 = sys_mmap(200) 1194 let tok: *u8 = sys_mmap(256) 1195 var started: i64 = 0 1196 var matched: i64 = 0 1197 var f: i64 = 0 1198 while f < nf { 1199 let fh: *u8 = flows[f] as *u8 1200 if wf_match(fh, ev) == 1 { 1201 matched = matched + 1 1202 pipe_field(fh, 0, fid, 64) 1203 var o: i64 = 0 1204 o = wf_cat(rid, o, evid) 1205 rid[o] = 46 as u8 1206 o = o + 1 1207 o = wf_cat(rid, o, fid) 1208 rid[o] = 0 as u8 1209 var q: i64 = 0 1210 q = wf_cat(tok, q, " rid=" as *u8) 1211 q = wf_cat(tok, q, rid) 1212 tok[q] = 32 as u8 1213 tok[q+1] = 0 as u8 1214 if wf_lines_with2(buf, szp[0], tok, "status=START" as *u8) > 0 { 1215 p("WFLOW dedup: run " as *u8); p(rid); p(" already started (durable) -- skipped\n" as *u8) 1216 } else { 1217 wf_emit(led, rid, fid, 0, "START" as *u8, 0) 1218 if cx[4] > 0 { wf_emit_ver(led, rid, cx[4]) } 1219 wf_bind_event(cx, rid, ev) 1220 wf_run_from(cx, rid, fid, 0) 1221 started = started + 1 1222 } 1223 } 1224 f = f + 1 1225 } 1226 if matched > 0 { if started == 0 { return 0 - 100 } } 1227 return started 1228} 1229 1230// any step row for this flow? 1231func wf_flow_known(cx: *i64, fid: *u8) -> i64 { 1232 let steps: *i64 = cx[0] as *i64 1233 let ns: i64 = cx[1] 1234 let f2: *u8 = sys_mmap(64) 1235 var i: i64 = 0 1236 while i < ns { 1237 pipe_field(steps[i] as *u8, 0, f2, 64) 1238 if seq(f2, fid) == 1 { return 1 } 1239 i = i + 1 1240 } 1241 return 0 1242} 1243 1244// ---- crash recovery: complete every in-flight run from its first non-OK step ---- 1245func wf_resume(cx: *i64) -> i64 { 1246 let led: *u8 = cx[2] as *u8 1247 let szp: *i64 = sys_mmap(8) as *i64 1248 let buf: *u8 = wf_readall(led, szp) 1249 if (buf as i64) == 0 { return 0 - 1 } 1250 let n: i64 = szp[0] 1251 let rid: *u8 = sys_mmap(200) 1252 let fid: *u8 = sys_mmap(64) 1253 let tok: *u8 = sys_mmap(256) 1254 var resumed: i64 = 0 1255 var ls: i64 = 0 1256 while ls < n { 1257 var le: i64 = ls 1258 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 1259 if wf_has_rng(buf, ls, le, "status=START" as *u8) == 1 { 1260 wf_val_after(buf, ls, le, " rid=" as *u8, rid, 200) 1261 wf_val_after(buf, ls, le, " flow=" as *u8, fid, 64) 1262 var q: i64 = 0 1263 q = wf_cat(tok, q, " rid=" as *u8) 1264 q = wf_cat(tok, q, rid) 1265 tok[q] = 32 as u8 1266 tok[q+1] = 0 as u8 1267 let dn: i64 = wf_lines_with2(buf, n, tok, "status=DONE" as *u8) 1268 let fl: i64 = wf_lines_with2(buf, n, tok, "status=FAILED" as *u8) 1269 if dn == 0 { if fl == 0 { 1270 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 } 1271 // R12: flag version drift (run started under a different definitions version) -- resume 1272 // proceeds (additive semantics) but the drift is RECORDED and printed 1273 let curv: i64 = cx[4] 1274 if curv > 0 { 1275 let ranv: i64 = wf_run_ver(buf, n, tok) 1276 if ranv > 0 { if ranv != curv { 1277 wf_emit_drift(led, rid, ranv, curv) 1278 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) 1279 } } 1280 } 1281 let s0: i64 = wf_max_ok_step(buf, n, tok) 1282 p("WFLOW-RESUME run=" as *u8); p(rid); p(" from-step=" as *u8); pn(s0 + 1); p("\n" as *u8) 1283 wf_run_from(cx, rid, fid, s0) 1284 resumed = resumed + 1 1285 } } 1286 } 1287 ls = le + 1 1288 } 1289 return resumed 1290} 1291 1292// ---- definitions as DATA FILES: load rows, validate whole, fire/resume (the production entry path) ---- 1293// skips comment (leading 35) and empty lines; -1 loud on missing/empty file or row-cap 1294func wf_lines_load(path: *u8, arr: *i64, cap: i64) -> i64 { 1295 let szp: *i64 = sys_mmap(8) as *i64 1296 let buf: *u8 = wf_readall(path, szp) 1297 if (buf as i64) == 0 { return 0 - 1 } 1298 let n: i64 = szp[0] 1299 if n == 0 { p("WFLOW files ERROR missing or empty definition file -- fail loud\n" as *u8); return 0 - 1 } 1300 var cnt: i64 = 0 1301 var ls: i64 = 0 1302 while ls < n { 1303 var le: i64 = ls 1304 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 1305 var keep: i64 = 0 1306 if le > ls { keep = 1 } 1307 if keep == 1 { if buf[ls] == (35 as u8) { keep = 0 } } 1308 if keep == 1 { if buf[ls] == (64 as u8) { keep = 0 } } 1309 if keep == 1 { 1310 if cnt >= cap { p("WFLOW files ERROR too many definition rows -- fail loud\n" as *u8); return 0 - 1 } 1311 let s: *u8 = sys_mmap(le - ls + 2) 1312 var k: i64 = 0 1313 while ls + k < le { s[k] = buf[ls+k]; k = k + 1 } 1314 s[k] = 0 as u8 1315 arr[cnt] = s as i64 1316 cnt = cnt + 1 1317 } 1318 ls = le + 1 1319 } 1320 return cnt 1321} 1322func wf_cx_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8, cx: *i64, flows: *i64) -> i64 { 1323 let steps: *i64 = sys_mmap(8 * 128) as *i64 1324 let nf: i64 = wf_lines_load(flowsp, flows, 64) 1325 if nf < 1 { return 0 - 1 } 1326 let nst: i64 = wf_lines_load(stepsp, steps, 128) 1327 if nst < 1 { return 0 - 1 } 1328 if wf_load_flows(flows, nf) < 0 { return 0 - 1 } 1329 if wf_load(steps, nst) < 0 { return 0 - 1 } 1330 cx[0] = steps as i64 1331 cx[1] = nst 1332 cx[2] = led as i64 1333 cx[3] = 0 1334 var usecat: i64 = 1 1335 if cat[0] == (45 as u8) { if cat[1] == (0 as u8) { usecat = 0 } } 1336 if usecat == 1 { cx[3] = cat as i64 } 1337 cx[4] = wf_defs_version(flowsp) 1338 return nf 1339} 1340func wf_fire_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8, ev: *u8) -> i64 { 1341 let cx: *i64 = sys_mmap(40) as *i64 1342 let flows: *i64 = sys_mmap(8 * 64) as *i64 1343 let nf: i64 = wf_cx_files(flowsp, stepsp, led, cat, cx, flows) 1344 if nf < 0 { return 0 - 1 } 1345 return wf_fire(cx, flows, nf, ev) 1346} 1347func wf_resume_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8) -> i64 { 1348 let cx: *i64 = sys_mmap(40) as *i64 1349 let flows: *i64 = sys_mmap(8 * 64) as *i64 1350 let nf: i64 = wf_cx_files(flowsp, stepsp, led, cat, cx, flows) 1351 if nf < 0 { return 0 - 1 } 1352 return wf_resume(cx) 1353} 1354 1355// ---- R6 run observability: board + 0-JS HTML derived ENTIRELY by ledger replay (nothing self-reported) ---- 1356func wf_state_name(st: i64) -> *u8 { 1357 if st == 3 { return "DONE" as *u8 } 1358 if st == 2 { return "FAILED" as *u8 } 1359 if st == 1 { return "PARKED" as *u8 } 1360 return "RUNNING" as *u8 1361} 1362// derived state: DONE > FAILED > PARKED(waiting approval) > RUNNING(in-flight) 1363func wf_run_state(buf: *u8, n: i64, tok: *u8) -> i64 { 1364 if wf_lines_with2(buf, n, tok, "status=DONE" as *u8) > 0 { return 3 } 1365 if wf_lines_with2(buf, n, tok, "status=FAILED" as *u8) > 0 { return 2 } 1366 if wf_lines_with2(buf, n, tok, "status=WAIT" as *u8) > 0 { return 1 } 1367 return 0 1368} 1369// collect distinct run ids from START lines; returns count, -1 loud over cap 1370func wf_obs_runs(buf: *u8, n: i64, rids: *i64, cap: i64) -> i64 { 1371 var cnt: i64 = 0 1372 var ls: i64 = 0 1373 while ls < n { 1374 var le: i64 = ls 1375 while le < n { if buf[le] == (10 as u8) { break } le = le + 1 } 1376 if wf_has_rng(buf, ls, le, "status=START" as *u8) == 1 { 1377 let r: *u8 = sys_mmap(200) 1378 wf_val_after(buf, ls, le, " rid=" as *u8, r, 200) 1379 var dup: i64 = 0 1380 var i: i64 = 0 1381 while i < cnt { if seq(rids[i] as *u8, r) == 1 { dup = 1; i = cnt } else { i = i + 1 } } 1382 if dup == 0 { 1383 if cnt >= cap { p("WFLOW obs ERROR run cap exceeded -- fail loud\n" as *u8); return 0 - 1 } 1384 rids[cnt] = r as i64 1385 cnt = cnt + 1 1386 } 1387 } 1388 ls = le + 1 1389 } 1390 return cnt 1391} 1392func wf_obs_tok(tok: *u8, rid: *u8) -> i64 { 1393 var q: i64 = 0 1394 q = wf_cat(tok, q, " rid=" as *u8) 1395 q = wf_cat(tok, q, rid) 1396 tok[q] = 32 as u8 1397 tok[q+1] = 0 as u8 1398 return q 1399} 1400func wf_obs_board(led: *u8) -> i64 { 1401 let szp: *i64 = sys_mmap(8) as *i64 1402 let buf: *u8 = wf_readall(led, szp) 1403 if (buf as i64) == 0 { return 0 - 1 } 1404 let n: i64 = szp[0] 1405 let rids: *i64 = sys_mmap(8 * 128) as *i64 1406 let cnt: i64 = wf_obs_runs(buf, n, rids, 128) 1407 if cnt < 0 { return 0 - 1 } 1408 let tok: *u8 = sys_mmap(256) 1409 var d: i64 = 0 1410 var f: i64 = 0 1411 var w: i64 = 0 1412 var ru: i64 = 0 1413 var i: i64 = 0 1414 while i < cnt { 1415 let rid: *u8 = rids[i] as *u8 1416 wf_obs_tok(tok, rid) 1417 let st: i64 = wf_run_state(buf, n, tok) 1418 let oks: i64 = wf_lines_with2(buf, n, tok, "status=OK" as *u8) 1419 let ats: i64 = wf_lines_with2(buf, n, tok, "status=ATT" as *u8) 1420 p("WFOBS run=" as *u8); p(rid) 1421 p(" state=" as *u8); p(wf_state_name(st)) 1422 p(" oksteps=" as *u8); pn(oks) 1423 p(" attempts=" as *u8); pn(ats) 1424 p("\n" as *u8) 1425 if st == 3 { d = d + 1 } 1426 if st == 2 { f = f + 1 } 1427 if st == 1 { w = w + 1 } 1428 if st == 0 { ru = ru + 1 } 1429 i = i + 1 1430 } 1431 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) 1432 return cnt 1433} 1434// 0-JS sovereign run board page; unquoted HTML5 attrs + rgb() colors (no quote/hash/bang literals) 1435func wf_obs_html(led: *u8, out: *u8) -> i64 { 1436 let szp: *i64 = sys_mmap(8) as *i64 1437 let buf: *u8 = wf_readall(led, szp) 1438 if (buf as i64) == 0 { return 0 - 1 } 1439 let n: i64 = szp[0] 1440 let rids: *i64 = sys_mmap(8 * 128) as *i64 1441 let cnt: i64 = wf_obs_runs(buf, n, rids, 128) 1442 if cnt < 0 { return 0 - 1 } 1443 let hb: *u8 = sys_mmap(WF_MAGIC_65536) 1444 var o: i64 = 0 1445 o = wf_cat(hb, o, "<" as *u8) 1446 hb[o] = 33 as u8 1447 o = o + 1 1448 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) 1449 let tok: *u8 = sys_mmap(256) 1450 var d: i64 = 0 1451 var f: i64 = 0 1452 var w: i64 = 0 1453 var ru: i64 = 0 1454 var i: i64 = 0 1455 while i < cnt { 1456 if o > WF_MAGIC_60000 { p("WFLOW obs ERROR page over cap -- fail loud\n" as *u8); return 0 - 1 } 1457 let rid: *u8 = rids[i] as *u8 1458 wf_obs_tok(tok, rid) 1459 let st: i64 = wf_run_state(buf, n, tok) 1460 let oks: i64 = wf_lines_with2(buf, n, tok, "status=OK" as *u8) 1461 let ats: i64 = wf_lines_with2(buf, n, tok, "status=ATT" as *u8) 1462 o = wf_cat(hb, o, "<tr><td>" as *u8) 1463 o = wf_cat(hb, o, rid) 1464 o = wf_cat(hb, o, "</td><td class=" as *u8) 1465 o = wf_cat(hb, o, wf_state_name(st)) 1466 o = wf_cat(hb, o, ">" as *u8) 1467 o = wf_cat(hb, o, wf_state_name(st)) 1468 o = wf_cat(hb, o, "</td><td>" as *u8) 1469 o = wf_catn(hb, o, oks) 1470 o = wf_cat(hb, o, "</td><td>" as *u8) 1471 o = wf_catn(hb, o, ats) 1472 o = wf_cat(hb, o, "</td></tr>" as *u8) 1473 if st == 3 { d = d + 1 } 1474 if st == 2 { f = f + 1 } 1475 if st == 1 { w = w + 1 } 1476 if st == 0 { ru = ru + 1 } 1477 i = i + 1 1478 } 1479 o = wf_cat(hb, o, "</table><p>runs " as *u8) 1480 o = wf_catn(hb, o, cnt) 1481 o = wf_cat(hb, o, " &middot; done " as *u8) 1482 o = wf_catn(hb, o, d) 1483 o = wf_cat(hb, o, " &middot; failed " as *u8) 1484 o = wf_catn(hb, o, f) 1485 o = wf_cat(hb, o, " &middot; parked " as *u8) 1486 o = wf_catn(hb, o, w) 1487 o = wf_cat(hb, o, " &middot; running " as *u8) 1488 o = wf_catn(hb, o, ru) 1489 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) 1490 if wf_writeall(out, hb, o) != 0 { return 0 - 1 } 1491 return cnt 1492} 1493 1494func wf_selftest() -> i64 { 1495 p("=== NX-WFLOW-ENGINE SELFTEST (R1 keystone: multi-step durable runs, retry, ledger-dedup, crash-resume) ===\n" as *u8) 1496 var ok: i64 = 1 1497 1498 // fresh unique ledger under /tmp (getpid is broken on this backend -- use epoch seconds) 1499 let led: *u8 = sys_mmap(128) 1500 var lo: i64 = 0 1501 lo = wf_cat(led, lo, "/tmp/wflow_st_" as *u8) 1502 lo = wf_catn(led, lo, wf_now()) 1503 lo = wf_cat(led, lo, ".log" as *u8) 1504 led[lo] = 0 as u8 1505 1506 let flows: *i64 = sys_mmap(8 * 8) as *i64 1507 flows[0] = "f1|deal-stage|to=Negotiation" as *u8 as i64 1508 flows[1] = "f2|giving-received|-" as *u8 as i64 1509 flows[2] = "f3|ops-check|-" as *u8 as i64 1510 let steps: *i64 = sys_mmap(8 * 16) as *i64 1511 steps[0] = "f1|1|create-task|-|Prepare contract and pricing|1" as *u8 as i64 1512 steps[1] = "f1|2|send|card|Congrats on reaching Negotiation|2" as *u8 as i64 1513 steps[2] = "f1|3|notify|-|Deal advanced|1" as *u8 as i64 1514 steps[3] = "f2|1|probe-fail|-|3|4" as *u8 as i64 1515 steps[4] = "f2|2|update-field|-|status=thanked|1" as *u8 as i64 1516 steps[5] = "f3|1|probe-fail|-|9|2" as *u8 as i64 1517 steps[6] = "f3|2|notify|-|never reached|1" as *u8 as i64 1518 1519 let cx: *i64 = sys_mmap(32) as *i64 1520 cx[0] = steps as i64 1521 cx[1] = 7 1522 cx[2] = led as i64 1523 cx[3] = 0 1524 1525 // T1 whole-set validation, loud refusals 1526 let l1: i64 = wf_load(steps, 7) 1527 let l2: i64 = wf_load_flows(flows, 3) 1528 p(" T1 load steps=" as *u8); pn(l1); p(" flows=" as *u8); pn(l2); p("\n" as *u8) 1529 if l1 != 7 { ok = 0 } 1530 if l2 != 3 { ok = 0 } 1531 let bad1: *i64 = sys_mmap(16) as *i64 1532 bad1[0] = "f9|1|explode|-|boom|1" as *u8 as i64 1533 if wf_load(bad1, 1) != (0 - 1) { ok = 0 } 1534 let bad2: *i64 = sys_mmap(16) as *i64 1535 bad2[0] = "f9|1|send|fax|x|1" as *u8 as i64 1536 if wf_load(bad2, 1) != (0 - 1) { ok = 0 } 1537 let bad3: *i64 = sys_mmap(16) as *i64 1538 bad3[0] = "f9|1|notify|-|x|0" as *u8 as i64 1539 if wf_load(bad3, 1) != (0 - 1) { ok = 0 } 1540 1541 // T2 fire -> 3-step run completes in order, DONE recorded 1542 let r1: i64 = wf_fire(cx, flows, 3, "deal-stage~deal=Riverside Contract~to=Negotiation~id=E1" as *u8) 1543 let szp: *i64 = sys_mmap(8) as *i64 1544 var buf: *u8 = wf_readall(led, szp) 1545 let okE1: i64 = wf_lines_with2(buf, szp[0], " rid=E1.f1 " as *u8, "status=OK" as *u8) 1546 let dnE1: i64 = wf_lines_with2(buf, szp[0], " rid=E1.f1 " as *u8, "status=DONE" as *u8) 1547 p(" T2 fire started=" as *u8); pn(r1); p(" okSteps=" as *u8); pn(okE1); p(" done=" as *u8); pn(dnE1); p("\n" as *u8) 1548 if r1 != 1 { ok = 0 } 1549 if okE1 != 3 { ok = 0 } 1550 if dnE1 != 1 { ok = 0 } 1551 1552 // T3 durable dedup: same event id never starts twice (survives across processes via the ledger) 1553 let r2: i64 = wf_fire(cx, flows, 3, "deal-stage~deal=Riverside Contract~to=Negotiation~id=E1" as *u8) 1554 buf = wf_readall(led, szp) 1555 let stE1: i64 = wf_lines_with2(buf, szp[0], " rid=E1.f1 " as *u8, "status=START" as *u8) 1556 p(" T3 refire rc=" as *u8); pn(r2); p(" starts=" as *u8); pn(stE1); p("\n" as *u8) 1557 if r2 != (0 - 100) { ok = 0 } 1558 if stE1 != 1 { ok = 0 } 1559 1560 // T4 per-step retry: probe-fail 3 with max 4 succeeds on attempt 3 1561 let r3: i64 = wf_fire(cx, flows, 3, "giving-received~person=Rose~amount=500~id=E2" as *u8) 1562 buf = wf_readall(led, szp) 1563 let ok3: i64 = wf_lines_with2(buf, szp[0], " rid=E2.f2 " as *u8, "status=OK att=3" as *u8) 1564 let atts: i64 = wf_lines_with2(buf, szp[0], " rid=E2.f2 " as *u8, "status=ATT" as *u8) 1565 let dnE2: i64 = wf_lines_with2(buf, szp[0], " rid=E2.f2 " as *u8, "status=DONE" as *u8) 1566 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) 1567 if r3 != 1 { ok = 0 } 1568 if ok3 != 1 { ok = 0 } 1569 if atts != 4 { ok = 0 } 1570 if dnE2 != 1 { ok = 0 } 1571 1572 // T5 exhausted retries -> FAILED at step, later steps never run 1573 let r4: i64 = wf_fire(cx, flows, 3, "ops-check~probe=edge~id=E3" as *u8) 1574 buf = wf_readall(led, szp) 1575 let fsE3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, "status=FAILSTEP" as *u8) 1576 let flE3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, "status=FAILED" as *u8) 1577 let s2E3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, " step=2 " as *u8) 1578 let dnE3: i64 = wf_lines_with2(buf, szp[0], " rid=E3.f3 " as *u8, "status=DONE" as *u8) 1579 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) 1580 if fsE3 != 1 { ok = 0 } 1581 if flE3 != 1 { ok = 0 } 1582 if s2E3 != 0 { ok = 0 } 1583 if dnE3 != 0 { ok = 0 } 1584 1585 // T6 THE CROWN -- durable crash-resume: synthetic run crashed after step 1; resume completes 2..3 1586 // WITHOUT re-executing step 1 (event-sourced replay, additive only) 1587 wf_emit(led, "E9.f1" as *u8, "f1" as *u8, 0, "START" as *u8, 0) 1588 wf_emit(led, "E9.f1" as *u8, "f1" as *u8, 1, "ATT" as *u8, 1) 1589 wf_emit(led, "E9.f1" as *u8, "f1" as *u8, 1, "OK" as *u8, 1) 1590 let rs: i64 = wf_resume(cx) 1591 buf = wf_readall(led, szp) 1592 let okE9: i64 = wf_lines_with2(buf, szp[0], " rid=E9.f1 " as *u8, "status=OK" as *u8) 1593 let ok1E9: i64 = wf_lines_with2(buf, szp[0], " rid=E9.f1 " as *u8, " step=1 status=OK" as *u8) 1594 let dnE9: i64 = wf_lines_with2(buf, szp[0], " rid=E9.f1 " as *u8, "status=DONE" as *u8) 1595 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) 1596 if rs != 1 { ok = 0 } 1597 if okE9 != 3 { ok = 0 } 1598 if ok1E9 != 1 { ok = 0 } 1599 if dnE9 != 1 { ok = 0 } 1600 1601 // T7 liar-kill negatives: unknown-flow in-flight run refuses loudly (no fabricated DONE); 1602 // nonsense status never appears 1603 wf_emit(led, "EX.zz" as *u8, "zz" as *u8, 0, "START" as *u8, 0) 1604 let rz: i64 = wf_resume(cx) 1605 buf = wf_readall(led, szp) 1606 let dnZZ: i64 = wf_lines_with2(buf, szp[0], " rid=EX.zz " as *u8, "status=DONE" as *u8) 1607 let ban: i64 = wf_lines_with2(buf, szp[0], "WFRUN" as *u8, "status=BANANA" as *u8) 1608 p(" T7 negctl resumeRc=" as *u8); pn(rz); p(" zzDone=" as *u8); pn(dnZZ); p(" nonsense=" as *u8); pn(ban); p("\n" as *u8) 1609 if rz != (0 - 1) { ok = 0 } 1610 if dnZZ != 0 { ok = 0 } 1611 if ban != 0 { ok = 0 } 1612 1613 p("NX-WFLOW-ENGINE-SELFTEST ledger=" as *u8); p(led) 1614 p(" runs: done=3 failed=1 dedup=1 resumed=1 " as *u8) 1615 if ok == 1 { p("verdict=GREEN\n" as *u8); return 0 } 1616 p("verdict=RED\n" as *u8) 1617 return 1 1618}