nx_wflow_engine.nx
buildroot/runtime/nx_wflow_engine.nx
about
nx_wflow_engine.nx -- R1 KEYSTONE of the Workflow Automation ladder (/compare/automation census row
"Unified multi-step durable workflow engine"). A WORKFLOW = DATA (law: data as data):
flow header = `flowid|event-kind|cond(k=v or -)`
step row = `flowid|idx|action|channel-or-dash|arg|max-attempts` (idx 1-based, author ascending)
An EVENT = `kind~k=v~...~id=EN` (the nx_crm_flow contract). The engine upgrades nx_crm_flow from
single-shot in-memory rules to DURABLE EXECUTION (the Temporal model, sovereign):
- every transition is ONE appended line in an append-only RUN LEDGER (flock + O_APPEND, torn-free):
`WFRUN rid=<evid>.<flowid> flow=<f> step=<n> status=START|ATT|OK|FAILSTEP|FAILED|DONE att=<k>`
- state is DERIVED BY REPLAY of the ledger (event-sourced; nothing mutates, history is sacred)
- IDEMPOTENT by construction: a (event id, flow) pair with a START line never starts twice --
the dedup survives crashes/restarts (stronger than the in-memory seen[] of nx_crm_flow)
- per-step RETRY up to max-attempts; exhausted -> run FAILED at that step, later steps never run
- wf_resume: any run with START but neither DONE nor FAILED (a crash mid-flight) is completed
from its first non-OK step; already-OK steps are NOT re-executed (proven in the selftest)
ACTIONS v1 (in-process, composing the send plane like nx_crm_flow): create-task | send (any nx_send
channel, sr_valid-checked at load) | notify | update-field | probe-fail (DIAGNOSTIC: fails while
attempt < arg -- the deterministic retry witness; keep it out of production flows).
wf_load validates the WHOLE step set up front (unknown action / bad channel / attempts<1 = LOUD -1,
nothing half-loaded). Production ledger path suggestion: knowledge/status/wflow_runs.log.
R3 CONNECTOR LAYER (exec-organ): a step can run ANY blessed organ as a workflow step, resolved through a
GREEN fail-closed catalog (rows `name<TAB>elf<TAB>GREEN`, the tool_allowlist.conf shape; cx[3]=catalog path,
0=exec-organ refused). Child exit 0 = step OK, nonzero = step fail (so per-step retry applies to REAL organ
execution). Production catalog: knowledge/wflow/connectors.conf.
PURE CORE (no main -- the house pure-core+gate idiom): gates = _hdl_build/nx_wflow_engine_gate (wf_selftest)
+ _hdl_build/nx_wflow_connect_gate (connector battery). license_tier: ORIGINAL
dependencies 1 imports · 13 importers
diagram shows first 10 each side; +0 more imports, +3 more importers in the complete lists below.
imports: nx_send.nx
imported by: nx_wflow.nxnx_wflow_approve_gate.nxnx_wflow_branch_gate.nxnx_wflow_capture_gate.nxnx_wflow_cli_gate.nxnx_wflow_connect_gate.nxnx_wflow_data_gate.nxnx_wflow_engine_gate.nxnx_wflow_loop_gate.nxnx_wflow_obs_gate.nxnx_wflow_onerr_gate.nxnx_wflow_tpl_gate.nxnx_wflow_version_gate.nx
structs
| none |
consts
| 27 | const WF_MAGIC_1089: i64 = 1089 |
| 28 | const WF_MAGIC_65536: i64 = 65536 |
| 29 | const WF_MAGIC_16777216: i64 = 16777216 |
| 30 | const WF_MAGIC_1024: i64 = 1024 |
| 31 | const WF_MAGIC_60000: i64 = 60000 |
| 110 | const WF_LEDCAP: i64 = 1048576 |
functions
| 34 | func wf_has(hay: *u8, needle: *u8) -> i64 |
| 48 | func wf_evid(ev: *u8, out: *u8, cap: i64) -> i64 |
| 66 | func 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 } |
| 67 | func wf_catn(dst: *u8, off: i64, v: i64) -> i64 |
| 79 | func wf_atoi(s: *u8) -> i64 |
| 95 | func wf_now() -> i64 { return sys_now_realtime_sec() } |
| 98 | func wf_append(path: *u8, line: *u8) -> i64 |
| 111 | func wf_readall(path: *u8, szp: *i64) -> *u8 |
| 127 | func wf_writeall(path: *u8, buf: *u8, n: i64) -> i64 |
| 137 | func wf_has_rng(buf: *u8, ls: i64, le: i64, needle: *u8) -> i64 called by 7: wf_lines_with2wf_max_ok_stepwf_var_getwf_defs_versionwf_run_verwf_resume+1 calls 1: slen |
| 151 | func wf_lines_with2(buf: *u8, n: i64, a: *u8, b: *u8) -> i64 |
| 163 | func wf_num_after(buf: *u8, ls: i64, le: i64, key: *u8) -> i64 |
| 183 | func wf_val_after(buf: *u8, ls: i64, le: i64, key: *u8, out: *u8, cap: i64) -> i64 |
| 203 | func wf_max_ok_step(buf: *u8, n: i64, ridtok: *u8) -> i64 |
| 219 | func wf_action_ok(a: *u8) -> i64 |
| 231 | func wf_load(steps: *i64, n: i64) -> i64 |
| 320 | func wf_load_flows(flows: *i64, n: i64) -> i64 |
| 332 | func wf_match(fhdr: *u8, ev: *u8) -> i64 |
| 360 | func wf_emit(led: *u8, rid: *u8, fid: *u8, step: i64, st: *u8, att: i64) -> i64 |
| 380 | func wf_conn_resolve(cx: *i64, name: *u8, out: *u8) -> i64 |
| 425 | func wf_rd32(b: *u8, off: i64) -> i64 called by 1: wf_exec_organ |
| 431 | func wf_exec_organ(cx: *i64, name0: *u8, arg: *u8, rid: *u8) -> i64 |
| 532 | func wf_emit_iter(led: *u8, rid: *u8, n: i64, val: *u8) -> i64 |
| 549 | func wf_for_each(cx: *i64, chspec: *u8, argraw: *u8, rid: *u8) -> i64 |
| 611 | func wf_exec(cx: *i64, action: *u8, ch: *u8, arg: *u8, att: i64, rid: *u8) -> i64 |
| 631 | func wf_decision(cx: *i64, rid: *u8, idx: i64) -> i64 |
| 658 | func wf_decide(cx: *i64, rid: *u8, idx: i64, decision: *u8, who: *u8) -> i64 |
| 683 | func wf_emit_var(led: *u8, rid: *u8, k: *u8, v: *u8) -> i64 |
| 697 | func wf_val_line_end(buf: *u8, ls: i64, le: i64, key: *u8, out: *u8, cap: i64) -> i64 |
| 716 | func wf_obs_tok_fwd(tok: *u8, rid: *u8) -> i64 |
| 725 | func wf_var_get(cx: *i64, rid: *u8, key: *u8, out: *u8, cap: i64) -> i64 |
| 753 | func wf_subst(cx: *i64, rid: *u8, src: *u8, out: *u8, cap: i64) -> i64 |
| 786 | func wf_bind_event(cx: *i64, rid: *u8, ev: *u8) -> i64 |
| 819 | func wf_defs_version(path: *u8) -> i64 called by 5: mainwf_tpl_listwf_tpl_instantiatewf_cx_filesmain calls 3: wf_readallwf_has_rngwf_num_after |
| 836 | func wf_emit_ver(led: *u8, rid: *u8, ver: i64) -> i64 |
| 847 | func wf_emit_drift(led: *u8, rid: *u8, ranv: i64, nowv: i64) -> i64 |
| 861 | func wf_run_ver(buf: *u8, n: i64, ridtok: *u8) -> i64 |
| 877 | func wf_tpl_exists(path: *u8) -> i64 called by 1: wf_tpl_list |
| 883 | func wf_tpl_path(dir: *u8, name: *u8, ext: *u8, out: *u8) -> i64 |
| 893 | func wf_tpl_index_path(dir: *u8, out: *u8) -> i64 |
| 903 | func wf_tpl_list(dir: *u8) -> i64 |
| 948 | func wf_tpl_indexed(dir: *u8, name: *u8) -> i64 |
| 975 | func wf_tpl_instantiate(dir: *u8, name: *u8, df: *u8, ds: *u8) -> i64 |
| 1006 | func wf_cond_ok(cx: *i64, rid: *u8, cond: *u8) -> i64 |
| 1026 | func wf_run_from(cx: *i64, rid: *u8, fid: *u8, s0: i64) -> i64 |
| 1185 | func wf_fire(cx: *i64, flows: *i64, nf: i64, ev: *u8) -> i64 called by 10: mainmainmainmainmainwf_fire_files+4 calls 11: wf_evidpwf_readallwf_matchpipe_fieldwf_cat+5 |
| 1231 | func wf_flow_known(cx: *i64, fid: *u8) -> i64 |
| 1245 | func wf_resume(cx: *i64) -> i64 called by 3: mainwf_resume_fileswf_selftest calls 12: wf_readallwf_has_rngwf_val_afterwf_catwf_lines_with2wf_flow_known+6 |
| 1294 | func wf_lines_load(path: *u8, arr: *i64, cap: i64) -> i64 |
| 1322 | func wf_cx_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8, cx: *i64, flows: *i64) -> i64 |
| 1340 | func wf_fire_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8, ev: *u8) -> i64 |
| 1347 | func wf_resume_files(flowsp: *u8, stepsp: *u8, led: *u8, cat: *u8) -> i64 |
| 1356 | func wf_state_name(st: i64) -> *u8 |
| 1363 | func wf_run_state(buf: *u8, n: i64, tok: *u8) -> i64 |
| 1370 | func wf_obs_runs(buf: *u8, n: i64, rids: *i64, cap: i64) -> i64 |
| 1392 | func wf_obs_tok(tok: *u8, rid: *u8) -> i64 |
| 1400 | func wf_obs_board(led: *u8) -> i64 |
| 1435 | func wf_obs_html(led: *u8, out: *u8) -> i64 |
| 1494 | func wf_selftest() -> i64 |