code wiki / (root) / nx_parquet.nx

nx_parquet.nx source

↩ module page · 357 lines · 14649 B

1// nx_parquet.nx -- THE PARQUET CLI over nx_parquet_lib (search R0d be_bright: BRIGHT ships as parquet). 2// 3// usage: nx_parquet meta <file> 4// nx_parquet count <file> 5// nx_parquet show <file> <leaf name or index> <n> 6// nx_parquet tsv <file> <leaf a> <leaf b> <out.tsv> 7// meta prints the footer as data (schema tree with max levels, row groups, every column chunk's codec, encodings, 8// counts and offsets) and then WALKS every BYTE_ARRAY column chunk's pages, so what it prints is what the 9// reader decoded, never what the header claims; the last line is verdict=OK or verdict=REFUSED with a count. 10// count the decoded entry, present-value and row counts per leaf across all row groups beside the declared rows: 11// the known-answer test for a dataset card (BRIGHT declares num_examples per split; this must reproduce it). 12// show the first n entries of one leaf: row, definition level, repetition level, value (a list leaf shows its 13// elements one per entry with the row they belong to). 14// tsv two flat leaves as <a>TAB<b> rows (tabs, newlines and carriage returns inside a cell become spaces), written 15// to <out>.tmp and renamed into place, one row group in memory at a time; rows where either side is null are 16// COUNTED and skipped, never written as empty strings. 17// Exit 0 OK, 1 REFUSED (the refusal named on the last line), 2 usage. 18// license_tier: ORIGINAL No hw writes (Rule 26). 19import "nx_syscalls.nx" 20import "nx_parquet_lib.nx" 21 22const NP_OK: i64 = 0 23const NP_REFUSED: i64 = 1 24const NP_USAGE: i64 = 2 25const NP_STDOUT: i64 = 1 26const NP_PATH_CAP: i64 = 1024 27const NP_TMP_SUFFIX: *u8 = ".tmp" 28 29func np_slen(s: *u8) -> i64 { 30 var n: i64 = 0 31 while (s[n] & PQ_BYTE) as i64 != 0 { n = n + 1 } 32 return n 33} 34func np_atoi(s: *u8) -> i64 { 35 var v: i64 = 0 36 var i: i64 = 0 37 if (s[0] & PQ_BYTE) as i64 == 0 { return PQ_NONE } 38 while (s[i] & PQ_BYTE) as i64 != 0 { 39 let ch: i64 = (s[i] & PQ_BYTE) as i64 40 if ch < PQ_DIGIT0 || ch > PQ_DIGIT0 + PQ_DECIMAL - 1 { return PQ_NONE } 41 v = v * PQ_DECIMAL + (ch - PQ_DIGIT0) 42 i = i + 1 43 } 44 return v 45} 46func np_streq(a: *u8, b: *u8) -> i64 { 47 var i: i64 = 0 48 while (a[i] & PQ_BYTE) as i64 != 0 || (b[i] & PQ_BYTE) as i64 != 0 { 49 if (a[i] & PQ_BYTE) as i64 != (b[i] & PQ_BYTE) as i64 { return 0 } 50 i = i + 1 51 } 52 return 1 53} 54// a leaf by name or by index; PQ_NONE when neither resolves 55func np_leaf_arg(pq: *i64, s: *u8) -> i64 { 56 let byname: i64 = pq_leaf_named(pq, s) 57 if byname >= 0 { return byname } 58 let ix: i64 = np_atoi(s) 59 if ix >= 0 && ix < pq[PQ_NLEAF] { return ix } 60 return PQ_NONE 61} 62func np_refusal(w: *i64, pq: *i64, rc: i64) -> i64 { 63 wb_puts(w, "REFUSED " as *u8) 64 wb_puts(w, pq_errname(rc)) 65 wb_byte(w, PQ_SPACE) 66 wb_kv(w, "code" as *u8, rc) 67 wb_kv(w, "what" as *u8, pq[PQ_WHAT]) 68 return 0 69} 70func np_name(w: *i64, pq: *i64, se: *i64) -> i64 { 71 wb_bytes(w, pq[PQ_BUF] as *u8, se[SE_NAMEOFF], se[SE_NAMELEN]) 72 return 0 73} 74 75func np_meta(pq: *i64, w: *i64, path: *u8) -> i64 { 76 wb_puts(w, "NX-PARQUET meta path=" as *u8) 77 wb_puts(w, path) 78 wb_byte(w, PQ_SPACE) 79 wb_kv(w, "bytes" as *u8, pq[PQ_LEN]) 80 wb_kv(w, "footer" as *u8, pq[PQ_META_LEN]) 81 wb_kv(w, "version" as *u8, pq[PQ_VERSION]) 82 wb_kv(w, "rows" as *u8, pq[PQ_NROWS]) 83 wb_kv(w, "schema" as *u8, pq[PQ_NSE]) 84 wb_kv(w, "leaves" as *u8, pq[PQ_NLEAF]) 85 wb_kv(w, "row_groups" as *u8, pq[PQ_NRG]) 86 wb_puts(w, "created_by=" as *u8) 87 wb_cell(w, pq[PQ_BUF] as *u8, pq[PQ_CREATED_OFF], pq[PQ_CREATED_LEN]) 88 wb_byte(w, PQ_NL) 89 var i: i64 = 0 90 while i < pq[PQ_NSE] { 91 let se: *i64 = pq_se(pq, i) 92 wb_puts(w, "schema|" as *u8) 93 wb_putn(w, i) 94 wb_byte(w, PQ_PIPE) 95 np_name(w, pq, se) 96 wb_byte(w, PQ_PIPE) 97 wb_kv(w, "type" as *u8, se[SE_TYPE]) 98 wb_kv(w, "rep" as *u8, se[SE_REP]) 99 wb_kv(w, "children" as *u8, se[SE_NCHILD]) 100 wb_kv(w, "converted" as *u8, se[SE_CONV]) 101 wb_kv(w, "parent" as *u8, se[SE_PARENT]) 102 wb_kv(w, "maxdef" as *u8, se[SE_MAXDEF]) 103 wb_kv(w, "maxrep" as *u8, se[SE_MAXREP]) 104 wb_kv(w, "leaf" as *u8, se[SE_LEAF]) 105 wb_byte(w, PQ_NL) 106 i = i + 1 107 } 108 var r: i64 = 0 109 while r < pq[PQ_NRG] { 110 let rg: *i64 = pq_rg(pq, r) 111 wb_puts(w, "rg|" as *u8) 112 wb_putn(w, r) 113 wb_byte(w, PQ_PIPE) 114 wb_kv(w, "rows" as *u8, rg[RG_NROWS]) 115 wb_kv(w, "total_bytes" as *u8, rg[RG_TOTAL]) 116 wb_byte(w, PQ_NL) 117 var j: i64 = 0 118 while j < pq[PQ_NLEAF] { 119 let cc: *i64 = pq_cc(pq, r, j) 120 wb_puts(w, "col|" as *u8) 121 wb_putn(w, r) 122 wb_byte(w, PQ_PIPE) 123 wb_putn(w, j) 124 wb_byte(w, PQ_PIPE) 125 if cc[CC_PATHLEN] > 0 { wb_bytes(w, pq[PQ_BUF] as *u8, cc[CC_PATHOFF], cc[CC_PATHLEN]) } 126 wb_byte(w, PQ_PIPE) 127 wb_kv(w, "type" as *u8, cc[CC_TYPE]) 128 wb_kv(w, "codec" as *u8, cc[CC_CODEC]) 129 wb_kv(w, "encmask" as *u8, cc[CC_ENCMASK]) 130 wb_kv(w, "nvals" as *u8, cc[CC_NVALS]) 131 wb_kv(w, "uncompressed" as *u8, cc[CC_UNCOMP]) 132 wb_kv(w, "compressed" as *u8, cc[CC_COMP]) 133 wb_kv(w, "data_off" as *u8, cc[CC_DATAOFF]) 134 wb_kv(w, "dict_off" as *u8, cc[CC_DICTOFF]) 135 wb_byte(w, PQ_NL) 136 j = j + 1 137 } 138 r = r + 1 139 } 140 var refused: i64 = 0 141 var walked: i64 = 0 142 var skipped: i64 = 0 143 r = 0 144 while r < pq[PQ_NRG] { 145 var j2: i64 = 0 146 while j2 < pq[PQ_NLEAF] { 147 let cc2: *i64 = pq_cc(pq, r, j2) 148 wb_puts(w, "pages|" as *u8) 149 wb_putn(w, r) 150 wb_byte(w, PQ_PIPE) 151 wb_putn(w, j2) 152 wb_byte(w, PQ_PIPE) 153 if cc2[CC_TYPE] == PQ_T_BYTE_ARRAY { 154 let pr: *i64 = pr_alloc(pq, r, j2) 155 let rc: i64 = pq_read_column(pq, r, j2, pr) 156 if rc < 0 { np_refusal(w, pq, rc); refused = refused + 1 } else { wb_puts(w, "OK " as *u8); walked = walked + 1 } 157 wb_kv(w, "pages" as *u8, pr[PR_PAGES]) 158 wb_kv(w, "page_types_mask" as *u8, pr[PR_PTYPES]) 159 wb_kv(w, "value_encodings_mask" as *u8, pr[PR_ENCMASK]) 160 wb_kv(w, "entries" as *u8, pr[PR_N]) 161 wb_kv(w, "present" as *u8, pr[PR_NONNULL]) 162 wb_kv(w, "rows_seen" as *u8, pr[PR_ROWS]) 163 wb_kv(w, "dictionary" as *u8, pr[PR_DICTN]) 164 wb_kv(w, "arena_used" as *u8, pr[PR_AUSED]) 165 wb_kv(w, "arena_cap" as *u8, pr[PR_ACAP]) 166 pr_free(pr) 167 } else { 168 wb_puts(w, "SKIPPED-TYPE " as *u8) 169 wb_kv(w, "type" as *u8, cc2[CC_TYPE]) 170 skipped = skipped + 1 171 } 172 wb_byte(w, PQ_NL) 173 j2 = j2 + 1 174 } 175 r = r + 1 176 } 177 wb_kv(w, "walked" as *u8, walked) 178 wb_kv(w, "skipped" as *u8, skipped) 179 wb_kv(w, "refused" as *u8, refused) 180 if refused == 0 { wb_puts(w, "verdict=OK" as *u8) } else { wb_puts(w, "verdict=REFUSED" as *u8) } 181 wb_byte(w, PQ_NL) 182 if refused == 0 { return NP_OK } 183 return NP_REFUSED 184} 185 186func np_count(pq: *i64, w: *i64, path: *u8) -> i64 { 187 wb_puts(w, "NX-PARQUET count path=" as *u8) 188 wb_puts(w, path) 189 wb_byte(w, PQ_SPACE) 190 wb_kv(w, "declared_rows" as *u8, pq[PQ_NROWS]) 191 wb_kv(w, "row_groups" as *u8, pq[PQ_NRG]) 192 wb_kv(w, "leaves" as *u8, pq[PQ_NLEAF]) 193 wb_byte(w, PQ_NL) 194 var rgrows: i64 = 0 195 var r0: i64 = 0 196 while r0 < pq[PQ_NRG] { let rg: *i64 = pq_rg(pq, r0); if rg[RG_NROWS] > 0 { rgrows = rgrows + rg[RG_NROWS] } r0 = r0 + 1 } 197 var refused: i64 = 0 198 var disagree: i64 = 0 199 var j: i64 = 0 200 while j < pq[PQ_NLEAF] { 201 let sei: i64 = pq_leaf_index(pq, j) 202 let se: *i64 = pq_se(pq, sei) 203 var entries: i64 = 0 204 var present: i64 = 0 205 var rows: i64 = 0 206 var declared: i64 = 0 207 var rc_worst: i64 = 0 208 var r: i64 = 0 209 while r < pq[PQ_NRG] { 210 let cc: *i64 = pq_cc(pq, r, j) 211 if cc[CC_NVALS] > 0 { declared = declared + cc[CC_NVALS] } 212 if cc[CC_TYPE] == PQ_T_BYTE_ARRAY { 213 let pr: *i64 = pr_alloc(pq, r, j) 214 let rc: i64 = pq_read_column(pq, r, j, pr) 215 if rc < 0 { rc_worst = rc } else { entries = entries + pr[PR_N]; present = present + pr[PR_NONNULL]; rows = rows + pr[PR_ROWS] } 216 pr_free(pr) 217 } else { rc_worst = PQ_E_UNSUPPORTED; pq[PQ_WHAT] = cc[CC_TYPE] } 218 r = r + 1 219 } 220 wb_puts(w, "leaf|" as *u8) 221 wb_putn(w, j) 222 wb_byte(w, PQ_PIPE) 223 np_name(w, pq, se) 224 wb_byte(w, PQ_PIPE) 225 wb_kv(w, "maxdef" as *u8, se[SE_MAXDEF]) 226 wb_kv(w, "maxrep" as *u8, se[SE_MAXREP]) 227 wb_kv(w, "declared_values" as *u8, declared) 228 wb_kv(w, "entries" as *u8, entries) 229 wb_kv(w, "present" as *u8, present) 230 wb_kv(w, "rows" as *u8, rows) 231 wb_kv(w, "rowgroup_rows" as *u8, rgrows) 232 if rc_worst < 0 { np_refusal(w, pq, rc_worst); refused = refused + 1 } else { 233 if rows != rgrows { wb_puts(w, "ROWS-DISAGREE " as *u8); disagree = disagree + 1 } else { wb_puts(w, "ROWS-AGREE " as *u8) } 234 } 235 wb_byte(w, PQ_NL) 236 j = j + 1 237 } 238 wb_kv(w, "rowgroup_rows" as *u8, rgrows) 239 wb_kv(w, "declared_rows" as *u8, pq[PQ_NROWS]) 240 wb_kv(w, "refused" as *u8, refused) 241 wb_kv(w, "disagree" as *u8, disagree) 242 if refused == 0 && disagree == 0 && rgrows == pq[PQ_NROWS] { wb_puts(w, "verdict=OK" as *u8) } else { wb_puts(w, "verdict=REFUSED" as *u8) } 243 wb_byte(w, PQ_NL) 244 if refused == 0 && disagree == 0 && rgrows == pq[PQ_NROWS] { return NP_OK } 245 return NP_REFUSED 246} 247 248func np_show(pq: *i64, w: *i64, leaf: i64, n: i64) -> i64 { 249 var shown: i64 = 0 250 var r: i64 = 0 251 var rowbase: i64 = 0 252 while r < pq[PQ_NRG] && shown < n { 253 let pr: *i64 = pr_alloc(pq, r, leaf) 254 let rc: i64 = pq_read_column(pq, r, leaf, pr) 255 if rc < 0 { np_refusal(w, pq, rc); wb_byte(w, PQ_NL); pr_free(pr); wb_puts(w, "verdict=REFUSED" as *u8); wb_byte(w, PQ_NL); return NP_REFUSED } 256 let off: *i64 = pr[PR_OFF] as *i64 257 let len: *i64 = pr[PR_LEN] as *i64 258 let def: *i64 = pr[PR_DEF] as *i64 259 let rep: *i64 = pr[PR_REP] as *i64 260 let row: *i64 = pr[PR_ROW] as *i64 261 var i: i64 = 0 262 while i < pr[PR_N] && shown < n { 263 wb_puts(w, "entry|" as *u8) 264 wb_putn(w, shown) 265 wb_byte(w, PQ_PIPE) 266 wb_kv(w, "row" as *u8, rowbase + row[i]) 267 wb_kv(w, "def" as *u8, def[i]) 268 wb_kv(w, "rep" as *u8, rep[i]) 269 wb_kv(w, "len" as *u8, len[i]) 270 wb_byte(w, PQ_PIPE) 271 if off[i] < 0 { wb_puts(w, "NULL" as *u8) } else { wb_cell(w, pr[PR_ARENA] as *u8, off[i], len[i]) } 272 wb_byte(w, PQ_NL) 273 shown = shown + 1 274 i = i + 1 275 } 276 rowbase = rowbase + pr[PR_ROWS] 277 pr_free(pr) 278 r = r + 1 279 } 280 wb_kv(w, "shown" as *u8, shown) 281 wb_puts(w, "verdict=OK" as *u8) 282 wb_byte(w, PQ_NL) 283 return NP_OK 284} 285 286func np_tsv(pq: *i64, w: *i64, la: i64, lb: i64, out: *u8) -> i64 { 287 let stats: *i64 = sys_mmap(PT_SLOTS * PQ_I64) as *i64 288 let rc: i64 = pq_tsv(pq, la, lb, out, stats) 289 wb_puts(w, "NX-PARQUET tsv " as *u8) 290 wb_kv(w, "rows" as *u8, stats[PT_ROWS]) 291 wb_kv(w, "written" as *u8, stats[PT_WRITTEN]) 292 wb_kv(w, "null_rows" as *u8, stats[PT_NULLROWS]) 293 wb_kv(w, "bytes" as *u8, stats[PT_BYTES]) 294 wb_puts(w, "out=" as *u8) 295 wb_puts(w, out) 296 wb_byte(w, PQ_SPACE) 297 if rc < 0 { 298 np_refusal(w, pq, rc) 299 if rc == PQ_E_ARG { wb_puts(w, "(tsv takes two flat leaves; a list leaf needs show or a converter that joins by row) " as *u8) } 300 wb_puts(w, "verdict=REFUSED" as *u8) 301 wb_byte(w, PQ_NL) 302 return NP_REFUSED 303 } 304 wb_puts(w, "verdict=OK" as *u8) 305 wb_byte(w, PQ_NL) 306 return NP_OK 307} 308 309func main(argc: i64, argv: *i64) -> i64 { 310 let w: *i64 = wb_new(NP_STDOUT) 311 if argc < 3 { 312 wb_puts(w, "usage: nx_parquet meta <file> | count <file> | show <file> <leaf> <n> | tsv <file> <leaf-a> <leaf-b> <out.tsv>" as *u8) 313 wb_byte(w, PQ_NL) 314 wb_flush(w) 315 sys_exit(NP_USAGE) 316 return NP_USAGE 317 } 318 let verb: *u8 = argv[1] as *u8 319 let path: *u8 = argv[2] as *u8 320 let pq: *i64 = sys_mmap(PQ_SLOTS * PQ_I64) as *i64 321 let rc: i64 = pq_open(pq, path) 322 if rc < 0 { 323 wb_puts(w, "NX-PARQUET path=" as *u8) 324 wb_puts(w, path) 325 wb_byte(w, PQ_SPACE) 326 np_refusal(w, pq, rc) 327 wb_puts(w, "verdict=REFUSED" as *u8) 328 wb_byte(w, PQ_NL) 329 wb_flush(w) 330 sys_exit(NP_REFUSED) 331 return NP_REFUSED 332 } 333 var code: i64 = NP_USAGE 334 var known: i64 = 0 335 if np_streq(verb, "meta" as *u8) == 1 { code = np_meta(pq, w, path); known = 1 } 336 if np_streq(verb, "count" as *u8) == 1 { code = np_count(pq, w, path); known = 1 } 337 if np_streq(verb, "show" as *u8) == 1 { 338 known = 1 339 if argc < 5 { wb_puts(w, "usage: nx_parquet show <file> <leaf> <n>" as *u8); wb_byte(w, PQ_NL); code = NP_USAGE } else { 340 let leaf: i64 = np_leaf_arg(pq, argv[3] as *u8) 341 let n: i64 = np_atoi(argv[4] as *u8) 342 if leaf < 0 || n < 0 { wb_puts(w, "REFUSED no such leaf or bad count " as *u8); wb_puts(w, argv[3] as *u8); wb_byte(w, PQ_NL); wb_puts(w, "verdict=REFUSED" as *u8); wb_byte(w, PQ_NL); code = NP_REFUSED } else { code = np_show(pq, w, leaf, n) } 343 } 344 } 345 if np_streq(verb, "tsv" as *u8) == 1 { 346 known = 1 347 if argc < 6 { wb_puts(w, "usage: nx_parquet tsv <file> <leaf-a> <leaf-b> <out.tsv>" as *u8); wb_byte(w, PQ_NL); code = NP_USAGE } else { 348 let la: i64 = np_leaf_arg(pq, argv[3] as *u8) 349 let lb: i64 = np_leaf_arg(pq, argv[4] as *u8) 350 if la < 0 || lb < 0 { wb_puts(w, "REFUSED no such leaf " as *u8); wb_puts(w, argv[3] as *u8); wb_byte(w, PQ_SPACE); wb_puts(w, argv[4] as *u8); wb_byte(w, PQ_NL); wb_puts(w, "verdict=REFUSED" as *u8); wb_byte(w, PQ_NL); code = NP_REFUSED } else { code = np_tsv(pq, w, la, lb, argv[5] as *u8) } 351 } 352 } 353 if known == 0 { wb_puts(w, "usage: nx_parquet meta <file> | count <file> | show <file> <leaf> <n> | tsv <file> <leaf-a> <leaf-b> <out.tsv>" as *u8); wb_byte(w, PQ_NL) } 354 wb_flush(w) 355 sys_exit(code) 356 return code 357}