code wiki / (root) / nx_wflow_mapping_oldpair_fixture_t278.nx

nx_wflow_mapping_oldpair_fixture_t278.nx source

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