code wiki / (root) / nx_ingest_batch.nx

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}