nx_ingest_batch.nx source
↩ module page · 437 lines · 17481 B
1// nx_ingest_batch.nx -- the actual ingestion SYSTEM.
2//
3// module: nishi-core.ingest.batch
4// depends: nishi-core.io.syscalls, nishi-core.io.dir,
5// nishi-core.data.value, nishi-core.data.value_parse_json
6// disk_kb: 6
7// capability: CORE_IO
8// wired_status: FULLY_WIRED
9//
10// license_tier: PUBLIC_NISHI_SUBSTRATE
11// genealogy_id: nishi_build_the_system_cardinal_2026 +
12// unix_pipeline_composition_tradition +
13// wiremock_2011_http_stub_pattern
14//
15// Brick #3 of the bits-up ingestion stack. No more hand-coded
16// per-fixture paths or curl loops. Composition:
17//
18// nx_dir_list(dir_path, ...) -- enumerate files
19// -> for each NxDirRow with name ending in suffix:
20// sys_read_file(<dir>/<name>) -- load bytes
21// -> nx_value_parse_json -- bytes -> NxValue tree
22// -> caller-supplied extract callback emits one
23// JSONL row to the output file descriptor
24// -> aggregate verdict + count
25//
26// "Caller-supplied" means: the batch driver itself is generic over
27// upstream-source shape. Today's smoke ships an extractor for
28// GBIF /v1/species/match responses, but any caller can drop in a
29// different field-extractor for a different upstream shape and the
30// batch substrate composes identically.
31
32import "syscalls.nx"
33import "nx_dir.nx"
34import "nx_value.nx"
35import "nx_value_parse_json.nx"
36import "nx_value_parse_csv.nx"
37const NX_MAGIC_1024: i64 = 1024
38const NX_MAGIC_1048576: i64 = 1048576
39
40// ===== Verdict ====================================================
41
42const NX_INGEST_OK: i64 = 1
43const NX_INGEST_DIR_OPEN_FAILED: i64 = 2
44const NX_INGEST_OUT_OPEN_FAILED: i64 = 3
45const NX_INGEST_BAD_ARGS: i64 = 4
46const NX_INGEST_NO_FILES_MATCHED: i64 = 5
47
48func nx_ingest_verdict_name(v: i64) -> *u8 {
49 if v == NX_INGEST_OK { return "OK" }
50 if v == NX_INGEST_DIR_OPEN_FAILED { return "DIR_OPEN_FAILED" }
51 if v == NX_INGEST_OUT_OPEN_FAILED { return "OUT_OPEN_FAILED" }
52 if v == NX_INGEST_BAD_ARGS { return "BAD_ARGS" }
53 if v == NX_INGEST_NO_FILES_MATCHED { return "NO_FILES_MATCHED" }
54 return "UNKNOWN"
55}
56
57// ===== Per-batch counters =========================================
58
59struct NxIngestReport {
60 report_hk: i64,
61 files_scanned: i64,
62 files_matched: i64,
63 parses_ok: i64,
64 parses_failed: i64,
65 rows_emitted: i64,
66 last_parse_verdict: i64,
67 bytes_total_read: i64,
68 bytes_total_written: i64,
69 verdict: i64,
70}
71
72const NX_INGEST_REPORT_BYTES: i64 = 80 // 10 fields * 8 bytes
73
74// ===== Tiny write helpers (no JSONL formatter dependency) =========
75
76func nx_ingest_w(fd: i64, s: *u8, n: i64) -> i64 {
77 return sys_write(fd, s, n)
78}
79
80func nx_ingest_w_int(fd: i64, value: i64) -> i64 {
81 let buf: *u8 = sys_mmap(32)
82 var v: i64 = value
83 var sign: i64 = 1
84 if v < 0 { sign = -1; v = -v }
85 var n: i64 = 0
86 var iter: i64 = 0
87 var verdict: i64 = 0
88 while verdict == 0 && iter < 24 {
89 let d: i64 = v % 10
90 buf[n] = (0x30 + d) as u8
91 n = n + 1
92 v = v / 10
93 if v == 0 { verdict = 1 }
94 iter = iter + 1
95 }
96 if sign == -1 { buf[n] = 0x2D as u8; n = n + 1 }
97 var lo: i64 = 0
98 var hi: i64 = n - 1
99 var r_iter: i64 = 0
100 var r_verdict: i64 = 0
101 while r_verdict == 0 && r_iter < 32 {
102 if lo >= hi { r_verdict = 1 }
103 if r_verdict == 0 {
104 let t: i64 = buf[lo] as i64
105 buf[lo] = buf[hi]
106 buf[hi] = t as u8
107 lo = lo + 1
108 hi = hi - 1
109 r_iter = r_iter + 1
110 }
111 }
112 return sys_write(fd, buf, n)
113}
114
115// ===== Generic JSONL emit: emit ALL string + int keys verbatim ===
116//
117// Walks an NxValue object and emits every string-valued and
118// integer-valued key as a JSONL field. Skips nested arrays/objects
119// (they need source-specific handling at a higher layer).
120//
121// This is THE shape-free emit: no upstream-source-specific knowledge.
122
123func nx_ingest_emit_jsonl_flat(fd: i64, obj: *NxValue) -> i64 {
124 if fd < 0 { return 0 }
125 if obj == 0 as *NxValue { return 0 }
126 if obj.kind != NX_VAL_OBJECT { return 0 }
127
128 let lbrace: *u8 = sys_mmap(2); lbrace[0] = 0x7B
129 let rbrace: *u8 = sys_mmap(2); rbrace[0] = 0x7D
130 let comma: *u8 = sys_mmap(2); comma[0] = 0x2C
131 let colon: *u8 = sys_mmap(2); colon[0] = 0x3A
132 let quote: *u8 = sys_mmap(2); quote[0] = 0x22
133 let nl: *u8 = sys_mmap(2); nl[0] = 0x0A
134
135 var written: i64 = 0
136 sys_write(fd, lbrace, 1)
137 written = written + 1
138
139 var i: i64 = 0
140 var iter: i64 = 0
141 var verdict: i64 = 0
142 var emitted: i64 = 0
143 while verdict == 0 && iter < NX_MAGIC_1024 {
144 if i >= obj.n_items { verdict = 1 }
145 if verdict == 0 {
146 let kp: *u8 = obj.keys_ptr[i]
147 let kl: i64 = obj.key_lens_ptr[i]
148 let val: *NxValue = obj.items_ptr[i]
149 var should_emit: i64 = 0
150 if val.kind == NX_VAL_STRING { should_emit = 1 }
151 if val.kind == NX_VAL_INT { should_emit = 1 }
152
153 if should_emit == 1 {
154 if emitted > 0 { sys_write(fd, comma, 1); written = written + 1 }
155 sys_write(fd, quote, 1); written = written + 1
156 sys_write(fd, kp, kl); written = written + kl
157 sys_write(fd, quote, 1); written = written + 1
158 sys_write(fd, colon, 1); written = written + 1
159 if val.kind == NX_VAL_STRING {
160 sys_write(fd, quote, 1); written = written + 1
161 sys_write(fd, val.str_ptr, val.str_len); written = written + val.str_len
162 sys_write(fd, quote, 1); written = written + 1
163 }
164 if val.kind == NX_VAL_INT {
165 written = written + nx_ingest_w_int(fd, val.int_value)
166 }
167 emitted = emitted + 1
168 }
169 i = i + 1
170 iter = iter + 1
171 }
172 }
173 sys_write(fd, rbrace, 1); written = written + 1
174 sys_write(fd, nl, 1); written = written + 1
175 return written
176}
177
178// ===== Build absolute path: dir + "/" + name ======================
179
180const NX_INGEST_MAX_PATH: i64 = 1024
181
182func nx_ingest_join_path(
183 dir: *u8,
184 dir_len: i64,
185 name: *u8,
186 name_len: i64,
187 out: *u8
188) -> i64 {
189 var i: i64 = 0
190 var iter: i64 = 0
191 var verdict: i64 = 0
192 while verdict == 0 && iter < NX_INGEST_MAX_PATH {
193 if i >= dir_len { verdict = 1 }
194 if verdict == 0 {
195 out[i] = dir[i]
196 i = i + 1
197 iter = iter + 1
198 }
199 }
200 out[i] = 0x2F // '/'
201 i = i + 1
202 var j: i64 = 0
203 var j_iter: i64 = 0
204 var j_verdict: i64 = 0
205 while j_verdict == 0 && j_iter < NX_INGEST_MAX_PATH {
206 if j >= name_len { j_verdict = 1 }
207 if j_verdict == 0 {
208 out[i] = name[j]
209 i = i + 1
210 j = j + 1
211 j_iter = j_iter + 1
212 }
213 }
214 out[i] = 0 // NUL
215 return i
216}
217
218// ===== Batch driver ===============================================
219//
220// Enumerates dir_path, for each regular file whose name ends with
221// suffix, reads bytes, parses JSON, emits flat-JSONL row. Caller
222// owns dir_path + suffix + out_path lifetimes.
223
224const NX_INGEST_MAX_ROWS_CAP: i64 = 8192
225const NX_INGEST_NAME_ARENA_BYTES: i64 = 1048576 // 1 MiB
226
227func nx_ingest_batch_run(
228 dir_path: *u8,
229 dir_len: i64,
230 suffix: *u8,
231 suffix_len: i64,
232 out_path: *u8,
233 report: *NxIngestReport,
234 now_unix: i64
235) -> i64 {
236 if report == 0 as *NxIngestReport { return NX_INGEST_BAD_ARGS }
237 report.report_hk = 0
238 report.files_scanned = 0
239 report.files_matched = 0
240 report.parses_ok = 0
241 report.parses_failed = 0
242 report.rows_emitted = 0
243 report.last_parse_verdict = 0
244 report.bytes_total_read = 0
245 report.bytes_total_written = 0
246 report.verdict = NX_INGEST_OK
247
248 if dir_path == 0 as *u8 { report.verdict = NX_INGEST_BAD_ARGS; return NX_INGEST_BAD_ARGS }
249 if out_path == 0 as *u8 { report.verdict = NX_INGEST_BAD_ARGS; return NX_INGEST_BAD_ARGS }
250
251 // Enumerate directory via nx_dir
252 let rows_raw: *u8 = sys_mmap(NX_INGEST_MAX_ROWS_CAP * NX_DIR_ROW_BYTES)
253 let rows: *NxDirRow = rows_raw as *NxDirRow
254 let arena: *u8 = sys_mmap(NX_INGEST_NAME_ARENA_BYTES)
255 let res_raw: *u8 = sys_mmap(NX_DIR_RESULT_BYTES)
256 let res: *NxDirResult = res_raw as *NxDirResult
257
258 let dir_v: i64 = nx_dir_list(
259 dir_path, rows, NX_INGEST_MAX_ROWS_CAP,
260 arena, NX_INGEST_NAME_ARENA_BYTES, 0, res)
261 if dir_v != NX_DIR_OK {
262 if dir_v != NX_DIR_TRUNCATED {
263 report.verdict = NX_INGEST_DIR_OPEN_FAILED
264 return NX_INGEST_DIR_OPEN_FAILED
265 }
266 }
267 report.files_scanned = res.n_filled
268
269 // Open output file (truncate)
270 let fd: i64 = sys_openat_wr(out_path, 0x1A4) // 0644
271 if fd < 0 {
272 report.verdict = NX_INGEST_OUT_OPEN_FAILED
273 return NX_INGEST_OUT_OPEN_FAILED
274 }
275
276 let lenbox_raw: *u8 = sys_mmap(8)
277 let path_buf: *u8 = sys_mmap(NX_INGEST_MAX_PATH)
278
279 var i: i64 = 0
280 var i_iter: i64 = 0
281 var i_verdict: i64 = 0
282 while i_verdict == 0 && i_iter < NX_INGEST_MAX_ROWS_CAP {
283 if i >= res.n_filled { i_verdict = 1 }
284 if i_verdict == 0 {
285 let row: *NxDirRow = nx_dir_row_at(rows, i)
286 // Filter: only regular files matching suffix.
287 let is_reg: i64 = nx_dir_row_is_regular_file(row)
288 var keep: i64 = 0
289 if is_reg == 1 {
290 if nx_dir_name_ends_with(row.name_ptr, row.name_len, suffix, suffix_len) == 1 {
291 keep = 1
292 }
293 }
294 // Always skip .jsonl files -- they are our OWN output
295 // shape (substrate-emitted JSONL); ingesting them would
296 // be self-cannibalization.
297 if keep == 1 {
298 let ext_jsonl: *u8 = sys_mmap(8)
299 ext_jsonl[0] = 0x2E; ext_jsonl[1] = 0x6A; ext_jsonl[2] = 0x73
300 ext_jsonl[3] = 0x6F; ext_jsonl[4] = 0x6E; ext_jsonl[5] = 0x6C // ".jsonl"
301 if nx_dir_name_ends_with(row.name_ptr, row.name_len, ext_jsonl, 6) == 1 {
302 keep = 0
303 }
304 }
305 if keep == 1 {
306 report.files_matched = report.files_matched + 1
307 // Build full path
308 nx_ingest_join_path(dir_path, dir_len, row.name_ptr, row.name_len, path_buf)
309 let lp: *i64 = lenbox_raw as *i64
310 let bytes: *u8 = sys_read_file(path_buf, lp)
311 if bytes == 0 as *u8 {
312 report.parses_failed = report.parses_failed + 1
313 }
314 if bytes != 0 as *u8 {
315 report.bytes_total_read = report.bytes_total_read + *lp
316 // Dispatch by suffix: .csv -> csv parser, .tsv ->
317 // tsv parser, .raw / .json / anything-else -> json
318 // parser.
319 var is_csv: i64 = 0
320 var is_tsv: i64 = 0
321 let ext_csv: *u8 = sys_mmap(8)
322 ext_csv[0] = 0x2E; ext_csv[1] = 0x63; ext_csv[2] = 0x73; ext_csv[3] = 0x76 // ".csv"
323 let ext_tsv: *u8 = sys_mmap(8)
324 ext_tsv[0] = 0x2E; ext_tsv[1] = 0x74; ext_tsv[2] = 0x73; ext_tsv[3] = 0x76 // ".tsv"
325 if nx_dir_name_ends_with(row.name_ptr, row.name_len, ext_csv, 4) == 1 { is_csv = 1 }
326 if nx_dir_name_ends_with(row.name_ptr, row.name_len, ext_tsv, 4) == 1 { is_tsv = 1 }
327
328 var pv: i64 = 0
329 var tree: *NxValue = 0 as *NxValue
330 if is_csv == 0 {
331 if is_tsv == 0 { tree = nx_value_parse_json(bytes, *lp, &pv) }
332 }
333 if is_csv == 1 {
334 tree = nx_value_parse_csv(bytes, *lp, 0 as *NxCsvOpts, &pv)
335 }
336 if is_tsv == 1 {
337 let topts_raw: *u8 = sys_mmap(NX_CSV_OPTS_BYTES)
338 let topts: *NxCsvOpts = topts_raw as *NxCsvOpts
339 nx_csv_opts_tsv(topts)
340 tree = nx_value_parse_csv(bytes, *lp, topts, &pv)
341 }
342 report.last_parse_verdict = pv
343 var parse_ok: i64 = 0
344 if is_csv == 1 {
345 if pv == NX_CSV_OK { parse_ok = 1 }
346 }
347 if is_tsv == 1 {
348 if pv == NX_CSV_OK { parse_ok = 1 }
349 }
350 if is_csv == 0 {
351 if is_tsv == 0 {
352 if pv == NX_VAL_PARSE_OK { parse_ok = 1 }
353 }
354 }
355 if parse_ok == 0 {
356 report.parses_failed = report.parses_failed + 1
357 }
358 if parse_ok == 1 {
359 report.parses_ok = report.parses_ok + 1
360 // ENVELOPE UNWRAP: if tree is OBJECT with key
361 // "results" or "data" or "items" whose value
362 // is ARRAY-of-OBJECT, treat that ARRAY as the
363 // emit set. GBIF / Elasticsearch / generic
364 // pagination shape: {"offset":N,"results":[...]}.
365 var emit_tree: *NxValue = tree
366 if tree.kind == NX_VAL_OBJECT {
367 let k_results: *u8 = sys_mmap(16)
368 k_results[0] = 0x72; k_results[1] = 0x65; k_results[2] = 0x73
369 k_results[3] = 0x75; k_results[4] = 0x6C; k_results[5] = 0x74
370 k_results[6] = 0x73 // "results"
371 let env: *NxValue = nx_value_object_get(tree, k_results, 7)
372 if env != 0 as *NxValue {
373 if env.kind == NX_VAL_ARRAY { emit_tree = env }
374 }
375 if emit_tree == tree {
376 let k_data: *u8 = sys_mmap(8)
377 k_data[0] = 0x64; k_data[1] = 0x61; k_data[2] = 0x74; k_data[3] = 0x61 // "data"
378 let env2: *NxValue = nx_value_object_get(tree, k_data, 4)
379 if env2 != 0 as *NxValue {
380 if env2.kind == NX_VAL_ARRAY { emit_tree = env2 }
381 }
382 }
383 if emit_tree == tree {
384 let k_items: *u8 = sys_mmap(8)
385 k_items[0] = 0x69; k_items[1] = 0x74; k_items[2] = 0x65
386 k_items[3] = 0x6D; k_items[4] = 0x73 // "items"
387 let env3: *NxValue = nx_value_object_get(tree, k_items, 5)
388 if env3 != 0 as *NxValue {
389 if env3.kind == NX_VAL_ARRAY { emit_tree = env3 }
390 }
391 }
392 }
393 // After unwrap: emit either single OBJECT or
394 // each inner OBJECT of an ARRAY.
395 if emit_tree.kind == NX_VAL_OBJECT {
396 let w: i64 = nx_ingest_emit_jsonl_flat(fd, emit_tree)
397 if w > 0 {
398 report.rows_emitted = report.rows_emitted + 1
399 report.bytes_total_written = report.bytes_total_written + w
400 }
401 }
402 if emit_tree.kind == NX_VAL_ARRAY {
403 var ai: i64 = 0
404 var ai_iter: i64 = 0
405 var ai_verdict: i64 = 0
406 while ai_verdict == 0 && ai_iter < NX_MAGIC_1048576 {
407 if ai >= emit_tree.n_items { ai_verdict = 1 }
408 if ai_verdict == 0 {
409 let inner: *NxValue = nx_value_array_get(emit_tree, ai)
410 if inner.kind == NX_VAL_OBJECT {
411 let iw: i64 = nx_ingest_emit_jsonl_flat(fd, inner)
412 if iw > 0 {
413 report.rows_emitted = report.rows_emitted + 1
414 report.bytes_total_written = report.bytes_total_written + iw
415 }
416 }
417 ai = ai + 1
418 ai_iter = ai_iter + 1
419 }
420 }
421 }
422 }
423 }
424 }
425 i = i + 1
426 i_iter = i_iter + 1
427 }
428 }
429 sys_close(fd)
430
431 if report.files_matched == 0 {
432 report.verdict = NX_INGEST_NO_FILES_MATCHED
433 return NX_INGEST_NO_FILES_MATCHED
434 }
435 report.verdict = NX_INGEST_OK
436 return NX_INGEST_OK
437}