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}