code wiki / (root) / nx_xform.nx

nx_xform.nx source

↩ module page · 311 lines · 13423 B

1// nx_xform.nx -- TRANSFORM-DAG (F1008), the dbt-class core: MODELS AS DATA + dependency ordering + 2// FAIL-CLOSED data-tests. A model names a stored column, a metric to compute over it, and a test its 3// result must pass; models declare deps on other models and the engine builds them in topological order, 4// recording the dependency edges (its own lineage) and STOPPING the pipeline the moment a data-test fails. 5// 6// SCOPE, honest: this is the ORCHESTRATION + VALIDATION half of dbt -- DAG order, cycle refusal, data 7// tests, fail-fast, self-reported lineage -- with each model computing a scalar metric via the sovereign 8// analytics (ad_profile). Multi-model DATA FLOW (a model consuming an upstream materialised table), the 9// semantic layer, and generated docs are the named finer gaps. Models are DATA (transform_models.conf), 10// never code (rule 11). license_tier: ORIGINAL No hardware writes (Rule 26). 11import "nx_syscalls.nx" 12import "nx_estate_path.nx" // ep_anchor: the CWD must not decide this organ's verdict 13import "_hdl_build/nx_analyst_store.nx" 14const XF_MAGIC_100000: i64 = 100000 15 16const XF_MAXM: i64 = 32 17const XF_CONF_CAP: i64 = 8192 18const XF_OUT: i64 = 32768 19const XF_TAB: i64 = 9 20const XF_NL: i64 = 10 21const XF_CR: i64 = 13 22const XF_HASH: i64 = 35 23const XF_ZERO: i64 = 48 24 25func xf_raw(out: *u8, o: i64, s: *u8) -> i64 { var i: i64 = 0; while s[i] != (0 as u8) { out[o] = s[i]; o = o + 1; i = i + 1 } return o } 26func xf_num(out: *u8, o: i64, v: i64) -> i64 { 27 var oo: i64 = o 28 var m: i64 = v 29 if m < 0 { out[oo] = 45 as u8; oo = oo + 1; m = 0 - m } 30 let t: *u8 = sys_mmap(24) 31 var k: i64 = 0 32 if m == 0 { t[0] = XF_ZERO as u8; k = 1 } 33 while m > 0 { t[k] = (XF_ZERO + (m % 10)) as u8; m = m / 10; k = k + 1 } 34 var i: i64 = k - 1 35 while i >= 0 { out[oo] = t[i]; oo = oo + 1; i = i - 1 } 36 return oo 37} 38func xf_qstr(out: *u8, o: i64, s: *u8) -> i64 { 39 out[o] = 34 as u8 40 var oo: i64 = o + 1 41 var i: i64 = 0 42 while s[i] != (0 as u8) { 43 let c: i64 = s[i] as i64 44 if c == 34 { out[oo] = 92 as u8; oo = oo + 1 } else { if c == 92 { out[oo] = 92 as u8; oo = oo + 1 } } 45 out[oo] = s[i]; oo = oo + 1; i = i + 1 46 } 47 out[oo] = 34 as u8 48 return oo + 1 49} 50func xf_streq(a: *u8, b: *u8) -> i64 { 51 var i: i64 = 0 52 while a[i] != (0 as u8) { if a[i] != b[i] { return 0 } i = i + 1 } 53 if b[i] != (0 as u8) { return 0 } 54 return 1 55} 56func xf_atoi(s: *u8) -> i64 { 57 var i: i64 = 0 58 var neg: i64 = 0 59 if s[0] == (45 as u8) { neg = 1; i = 1 } 60 var v: i64 = 0 61 var go: i64 = 1 62 while go == 1 { let c: i64 = s[i] as i64; if c < 48 { go = 0 } else { if c > 57 { go = 0 } else { v = v * 10 + (c - 48); i = i + 1 } } } 63 if neg == 1 { return 0 - v } 64 return v 65} 66func xf_dup(buf: *u8, a: i64, b: i64) -> *u8 { 67 let n: i64 = b - a 68 let o: *u8 = sys_mmap(n + 1) 69 var i: i64 = 0 70 while i < n { o[i] = buf[a+i]; i = i + 1 } 71 o[n] = 0 as u8 72 return o 73} 74// does the comma-separated dep list contain tok as a whole element? a lone "-" means no deps. 75func xf_csv_has(csv: *u8, tok: *u8) -> i64 { 76 if csv[0] == (45 as u8) { if csv[1] == (0 as u8) { return 0 } } 77 var start: i64 = 0 78 var i: i64 = 0 79 var scanning: i64 = 1 80 while scanning == 1 { 81 let c: i64 = csv[i] as i64 82 var atend: i64 = 0 83 if c == 0 { atend = 1 } else { if c == 44 { atend = 1 } } 84 if atend == 1 { 85 // element is csv[start..i) 86 var full: i64 = 1 87 var k: i64 = 0 88 while tok[k] != (0 as u8) { if start + k >= i { full = 0 } else { if csv[start+k] != tok[k] { full = 0 } } k = k + 1 } 89 if full == 1 { if start + k == i { return 1 } } 90 start = i + 1 91 } 92 if c == 0 { scanning = 0 } 93 i = i + 1 94 } 95 return 0 96} 97func xf_dep_count(csv: *u8) -> i64 { 98 if csv[0] == (45 as u8) { if csv[1] == (0 as u8) { return 0 } } 99 if csv[0] == (0 as u8) { return 0 } 100 var n: i64 = 1 101 var i: i64 = 0 102 while csv[i] != (0 as u8) { if csv[i] == (44 as u8) { n = n + 1 } i = i + 1 } 103 return n 104} 105func xf_load(ids: *i64, deps: *i64, prefs: *i64, keys: *i64, fidx: *i64, ops: *i64, tops: *i64, targs: *i64) -> i64 { 106 let fd: i64 = sys_openat_rd("knowledge/registry/transform_models.conf" as *u8) 107 if fd < 0 { return 0 } 108 let cb: *u8 = sys_mmap(XF_CONF_CAP) 109 var n: i64 = 0 110 var r: i64 = 1 111 while r > 0 { if n >= XF_CONF_CAP { r = 0 } else { r = sys_read(fd, (((cb as i64) + n) as *u8), XF_CONF_CAP - n); if r > 0 { n = n + r } } } 112 sys_close(fd) 113 var cnt: i64 = 0 114 var ls: i64 = 0 115 while ls < n { 116 var le: i64 = ls 117 var e: i64 = 0 118 while e == 0 { if le >= n { e = 1 } else { if cb[le] == (XF_NL as u8) { e = 1 } else { le = le + 1 } } } 119 var te: i64 = le 120 if te > ls { if cb[te-1] == (XF_CR as u8) { te = te - 1 } } 121 if te > ls { if cb[ls] != (XF_HASH as u8) { if cnt < XF_MAXM { 122 let fo: *i64 = sys_mmap(8 * 18) as *i64 123 var nf: i64 = 0 124 var fstart: i64 = ls 125 var i: i64 = ls 126 while i <= te { 127 var cut: i64 = 0 128 if i == te { cut = 1 } else { if cb[i] == (XF_TAB as u8) { cut = 1 } } 129 if cut == 1 { if nf < 9 { fo[nf * 2] = fstart; fo[nf * 2 + 1] = i; nf = nf + 1 } fstart = i + 1 } 130 i = i + 1 131 } 132 if nf >= 8 { 133 ids[cnt] = xf_dup(cb, fo[0], fo[1]) as i64 134 deps[cnt] = xf_dup(cb, fo[2], fo[3]) as i64 135 prefs[cnt] = xf_dup(cb, fo[4], fo[5]) as i64 136 keys[cnt] = xf_dup(cb, fo[6], fo[7]) as i64 137 fidx[cnt] = xf_atoi(xf_dup(cb, fo[8], fo[9])) 138 ops[cnt] = xf_dup(cb, fo[10], fo[11]) as i64 139 tops[cnt] = xf_dup(cb, fo[12], fo[13]) as i64 140 targs[cnt] = xf_atoi(xf_dup(cb, fo[14], fo[15])) 141 cnt = cnt + 1 142 } 143 } } } 144 ls = le + 1 145 } 146 return cnt 147} 148// topological order into order[]; returns count emitted (< m signals a CYCLE) 149func xf_topo(ids: *i64, deps: *i64, m: i64, order: *i64) -> i64 { 150 let built: *i64 = sys_mmap(8 * XF_MAXM) as *i64 151 var i: i64 = 0 152 while i < m { built[i] = 0; i = i + 1 } 153 var emitted: i64 = 0 154 var progress: i64 = 1 155 while progress == 1 { 156 progress = 0 157 var k: i64 = 0 158 while k < m { 159 if built[k] == 0 { 160 var ready: i64 = 1 161 var d: i64 = 0 162 while d < m { 163 if d != k { if built[d] == 0 { if xf_csv_has(deps[k] as *u8, ids[d] as *u8) == 1 { ready = 0 } } } 164 d = d + 1 165 } 166 if ready == 1 { built[k] = 1; order[emitted] = k; emitted = emitted + 1; progress = 1 } 167 } 168 k = k + 1 169 } 170 } 171 return emitted 172} 173func xf_metric(pref: *u8, key: *u8, fidx: i64, op: *u8) -> i64 { 174 let col: *i64 = sys_mmap(8 * XF_MAGIC_100000) as *i64 175 let fl: *i64 = sys_mmap(8) as *i64 176 fl[0] = 0 177 let rows: i64 = asr_load_field_sharded(pref, key, fidx, col, XF_MAGIC_100000, 64, fl) 178 if rows <= 0 { return 0 } 179 let prof: *i64 = sys_mmap(8 * 16) as *i64 180 ad_profile(col, rows, 10, prof) 181 if xf_streq(op, "count" as *u8) == 1 { return prof[0] } 182 if xf_streq(op, "min" as *u8) == 1 { return prof[2] } 183 if xf_streq(op, "max" as *u8) == 1 { return prof[3] } 184 if xf_streq(op, "mean" as *u8) == 1 { return prof[4] } 185 if xf_streq(op, "median" as *u8) == 1 { return prof[5] } 186 if xf_streq(op, "stddev" as *u8) == 1 { return prof[6] } 187 return prof[0] 188} 189func xf_test(result: i64, top: *u8, targ: i64) -> i64 { 190 if xf_streq(top, "gt" as *u8) == 1 { if result > targ { return 1 } return 0 } 191 if xf_streq(top, "lt" as *u8) == 1 { if result < targ { return 1 } return 0 } 192 if xf_streq(top, "nonzero" as *u8) == 1 { if result != 0 { return 1 } return 0 } 193 return 0 194} 195func xf_order_json(out: *u8) -> i64 { 196 let ids: *i64 = sys_mmap(8 * XF_MAXM) as *i64 197 let deps: *i64 = sys_mmap(8 * XF_MAXM) as *i64 198 let prefs: *i64 = sys_mmap(8 * XF_MAXM) as *i64 199 let keys: *i64 = sys_mmap(8 * XF_MAXM) as *i64 200 let fidx: *i64 = sys_mmap(8 * XF_MAXM) as *i64 201 let ops: *i64 = sys_mmap(8 * XF_MAXM) as *i64 202 let tops: *i64 = sys_mmap(8 * XF_MAXM) as *i64 203 let targs: *i64 = sys_mmap(8 * XF_MAXM) as *i64 204 let m: i64 = xf_load(ids, deps, prefs, keys, fidx, ops, tops, targs) 205 let order: *i64 = sys_mmap(8 * XF_MAXM) as *i64 206 let emitted: i64 = xf_topo(ids, deps, m, order) 207 var o: i64 = 0 208 o = xf_raw(out, o, "{\"tool\":\"nx_xform\",\"verb\":\"order\",\"models\":" as *u8) 209 o = xf_num(out, o, m) 210 if emitted < m { 211 o = xf_raw(out, o, ",\"ok\":false,\"reason\":\"CYCLE -- dependency graph is not a DAG; refusing\",\"emitted\":" as *u8) 212 o = xf_num(out, o, emitted) 213 o = xf_raw(out, o, "}\n" as *u8) 214 out[o] = 0 as u8 215 return o 216 } 217 o = xf_raw(out, o, ",\"ok\":true,\"order\":[" as *u8) 218 var i: i64 = 0 219 while i < emitted { 220 if i > 0 { o = xf_raw(out, o, "," as *u8) } 221 o = xf_qstr(out, o, ids[order[i]] as *u8) 222 i = i + 1 223 } 224 o = xf_raw(out, o, "]}\n" as *u8) 225 out[o] = 0 as u8 226 return o 227} 228func xf_run_json(out: *u8) -> i64 { 229 let ids: *i64 = sys_mmap(8 * XF_MAXM) as *i64 230 let deps: *i64 = sys_mmap(8 * XF_MAXM) as *i64 231 let prefs: *i64 = sys_mmap(8 * XF_MAXM) as *i64 232 let keys: *i64 = sys_mmap(8 * XF_MAXM) as *i64 233 let fidx: *i64 = sys_mmap(8 * XF_MAXM) as *i64 234 let ops: *i64 = sys_mmap(8 * XF_MAXM) as *i64 235 let tops: *i64 = sys_mmap(8 * XF_MAXM) as *i64 236 let targs: *i64 = sys_mmap(8 * XF_MAXM) as *i64 237 let m: i64 = xf_load(ids, deps, prefs, keys, fidx, ops, tops, targs) 238 let order: *i64 = sys_mmap(8 * XF_MAXM) as *i64 239 let emitted: i64 = xf_topo(ids, deps, m, order) 240 var o: i64 = 0 241 o = xf_raw(out, o, "{\"tool\":\"nx_xform\",\"verb\":\"run\",\"models\":" as *u8) 242 o = xf_num(out, o, m) 243 if emitted < m { 244 o = xf_raw(out, o, ",\"ok\":false,\"reason\":\"CYCLE -- refusing to run a non-DAG\"}\n" as *u8) 245 out[o] = 0 as u8 246 return o 247 } 248 o = xf_raw(out, o, ",\"ok\":true,\"built\":[" as *u8) 249 // per-model status: 0 built+passed, 1 failed its own test, 2 skipped (a dep failed/was skipped). 250 // DESCENDANT-ONLY skip (the correct dbt semantic): a failed model does NOT halt unrelated models, 251 // only the ones that (transitively) depend on it -- and transitivity is automatic because we walk in 252 // topological order, so a dep's status is always finalised before its dependent is reached. 253 let status: *i64 = sys_mmap(8 * XF_MAXM) as *i64 254 var all_ok: i64 = 1 255 var i: i64 = 0 256 while i < emitted { 257 let k: i64 = order[i] 258 if i > 0 { o = xf_raw(out, o, "," as *u8) } 259 o = xf_raw(out, o, "{\"id\":" as *u8) 260 o = xf_qstr(out, o, ids[k] as *u8) 261 o = xf_raw(out, o, ",\"deps\":" as *u8) 262 o = xf_qstr(out, o, deps[k] as *u8) 263 // any dep failed or skipped? 264 var blocked: i64 = 0 265 var d: i64 = 0 266 while d < m { 267 if d != k { if xf_csv_has(deps[k] as *u8, ids[d] as *u8) == 1 { if status[d] != 0 { blocked = 1 } } } 268 d = d + 1 269 } 270 if blocked == 1 { 271 status[k] = 2 272 all_ok = 0 273 o = xf_raw(out, o, ",\"state\":\"SKIPPED (a dependency failed -- fail-closed)\"}" as *u8) 274 } else { 275 let res: i64 = xf_metric(prefs[k] as *u8, keys[k] as *u8, fidx[k], ops[k] as *u8) 276 let pass: i64 = xf_test(res, tops[k] as *u8, targs[k]) 277 o = xf_raw(out, o, ",\"op\":" as *u8) 278 o = xf_qstr(out, o, ops[k] as *u8) 279 o = xf_raw(out, o, ",\"result\":" as *u8) 280 o = xf_num(out, o, res) 281 o = xf_raw(out, o, ",\"test\":\"" as *u8) 282 o = xf_raw(out, o, tops[k] as *u8) 283 o = xf_raw(out, o, " " as *u8) 284 o = xf_num(out, o, targs[k]) 285 o = xf_raw(out, o, "\",\"pass\":" as *u8) 286 if pass == 1 { o = xf_raw(out, o, "true" as *u8); status[k] = 0 } else { o = xf_raw(out, o, "false" as *u8); status[k] = 1; all_ok = 0 } 287 o = xf_raw(out, o, ",\"state\":\"built\"}" as *u8) 288 } 289 i = i + 1 290 } 291 o = xf_raw(out, o, "],\"pipeline_pass\":" as *u8) 292 if all_ok == 1 { o = xf_raw(out, o, "true" as *u8) } else { o = xf_raw(out, o, "false" as *u8) } 293 o = xf_raw(out, o, "}\n" as *u8) 294 out[o] = 0 as u8 295 return o 296} 297func main(argc: i64, argv: *i64) -> i64 { 298 // ANCHOR FIRST (2026-08-04, nx_cwdguard finding): this organ reads a RELATIVE 299 // knowledge/ path, so its answer depended on where it was launched. No-op when 300 // already at the estate root, so the cron/MCP context is unchanged. 301 ep_anchor() 302 let out: *u8 = sys_mmap(XF_OUT) 303 if argc < 2 { let n: i64 = xf_raw(out, 0, "usage: nx_xform {order|run}\n" as *u8); sys_write(2, out, n); sys_exit(2); return 2 } 304 let verb: *u8 = argv[1] as *u8 305 if xf_streq(verb, "order" as *u8) == 1 { let n: i64 = xf_order_json(out); sys_write(1, out, n); return 0 } 306 if xf_streq(verb, "run" as *u8) == 1 { let n: i64 = xf_run_json(out); sys_write(1, out, n); return 0 } 307 let n: i64 = xf_raw(out, 0, "nx_xform: unknown verb\n" as *u8) 308 sys_write(2, out, n) 309 sys_exit(2) 310 return 2 311}