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}