code wiki / (root) / nx_queries_merge_lib.nx

nx_queries_merge_lib.nx source

↩ module page · 144 lines · 6671 B

1// nx_queries_merge_lib.nx -- the row-aligned union of K query files (see nx_queries_merge.nx). LIB (no main) so the 2// gate drives qm_merge in-process on fixtures. license_tier: ORIGINAL No hw writes (Rule 26). 3import "nx_syscalls.nx" 4 5const QM_I64: i64 = 8 6const QM_ST_SLOTS: i64 = 8 // [0]=rows [1]=bytes_out [2]=bad_row (1-based, 0 = none) [3]=bad_file (0-based) [4]=rows_in_bad_file 7const QM_ROW_SLOTS: i64 = 4 // [0]=line start [1]=line end (CR trimmed) [2]=first TAB or line end [3]=NL index or n 8const QM_EXIT_OK: i64 = 0 9const QM_EXIT_USAGE: i64 = 2 10const QM_EXIT_REFUSED: i64 = 3 11const QM_EXIT_IO: i64 = 4 12const QM_TAB: i64 = 9 13const QM_NL: i64 = 10 14const QM_CR: i64 = 13 15const QM_SPACE: i64 = 32 16const QM_FILE_MODE: i64 = 420 17const QM_TMP_SUFFIX: *u8 = ".tmp" as *u8 18const QM_PATH_CAP: i64 = 1024 19const QM_STDOUT: i64 = 1 20 21func qm_slen(s: *u8) -> i64 { var n: i64 = 0; while s[n] != (0 as u8) { n = n + 1 } return n } 22func qm_puts(s: *u8) -> i64 { let n: i64 = qm_slen(s); if n > 0 { sys_write(QM_STDOUT, s, n) } return 0 } 23func qm_putn(v: i64) -> i64 { 24 let t: *u8 = sys_mmap(32) 25 var m: i64 = v 26 var k: i64 = 0 27 if m < 0 { sys_write(QM_STDOUT, "-" as *u8, 1); m = 0 - m } 28 if m == 0 { t[0] = 48 as u8; k = 1 } 29 while m > 0 { t[k] = (48 + (m % 10)) as u8; m = m / 10; k = k + 1 } 30 let b: *u8 = sys_mmap(32) 31 var i: i64 = 0 32 while i < k { b[i] = t[k - 1 - i]; i = i + 1 } 33 sys_write(QM_STDOUT, b, k) 34 return 0 35} 36func qm_cat(dst: *u8, off: i64, s: *u8) -> i64 { var i: i64 = 0; var o: i64 = off; while s[i] != (0 as u8) { dst[o] = s[i]; o = o + 1; i = i + 1 } dst[o] = 0 as u8; return o } 37func qm_count_rows(buf: *u8, n: i64) -> i64 { 38 var c: i64 = 0 39 var i: i64 = 0 40 while i < n { if buf[i] == (QM_NL as u8) { c = c + 1 } i = i + 1 } 41 if n > 0 { if buf[n - 1] != (QM_NL as u8) { c = c + 1 } } 42 return c 43} 44// the row starting at byte p: o4 = [line start, line end (CR trimmed), first TAB or line end, NL index or n]; 1 found, 0 past the end 45func qm_row_at(buf: *u8, n: i64, p: i64, o4: *i64) -> i64 { 46 if p >= n { return 0 } 47 var e: i64 = p 48 var found: i64 = 0 49 while found == 0 { if e >= n { found = 1 } else { if buf[e] == (QM_NL as u8) { found = 1 } else { e = e + 1 } } } 50 var le: i64 = e 51 if le > p { if buf[le - 1] == (QM_CR as u8) { le = le - 1 } } 52 var tb: i64 = p 53 var ft: i64 = 0 54 while ft == 0 { if tb >= le { ft = 1 } else { if buf[tb] == (QM_TAB as u8) { ft = 1 } else { tb = tb + 1 } } } 55 o4[0] = p; o4[1] = le; o4[2] = tb; o4[3] = e 56 return 1 57} 58func qm_ids_equal(a: *u8, a0: i64, a1: i64, b: *u8, b0: i64, b1: i64) -> i64 { 59 if a1 - a0 != b1 - b0 { return 0 } 60 var i: i64 = 0 61 while i < a1 - a0 { if a[a0 + i] != b[b0 + i] { return 0 } i = i + 1 } 62 return 1 63} 64func qm_write_all(fd: i64, buf: *u8, n: i64) -> i64 { 65 var off: i64 = 0 66 while off < n { let w: i64 = sys_write(fd, buf + off, n - off); if w <= 0 { return 0 - 1 } off = off + w } 67 return 0 68} 69// the merge. Returns QM_EXIT_*; st as documented on QM_ST_SLOTS. 70func qm_merge(out: *u8, ins: *i64, n_in: i64, st: *i64) -> i64 { 71 var z: i64 = 0 72 while z < QM_ST_SLOTS { st[z] = 0; z = z + 1 } 73 let bufs: *i64 = sys_mmap(n_in * QM_I64) as *i64 74 let lens: *i64 = sys_mmap(n_in * QM_I64) as *i64 75 let rows: *i64 = sys_mmap(n_in * QM_I64) as *i64 76 let cur: *i64 = sys_mmap(n_in * QM_I64) as *i64 // per file: byte cursor at the start of the next row 77 var f: i64 = 0 78 var total: i64 = 0 79 while f < n_in { 80 let ln: *i64 = sys_mmap(2 * QM_I64) as *i64 81 let b: *u8 = sys_read_file(ins[f] as *u8, ln) 82 if (b as i64) == 0 { st[3] = f; return QM_EXIT_IO } 83 bufs[f] = b as i64 84 lens[f] = ln[0] 85 rows[f] = qm_count_rows(b, ln[0]) 86 total = total + ln[0] 87 cur[f] = 0 88 f = f + 1 89 } 90 // every file must carry the same number of rows 91 f = 1 92 while f < n_in { if rows[f] != rows[0] { st[3] = f; st[4] = rows[f]; st[0] = rows[0]; return QM_EXIT_REFUSED } f = f + 1 } 93 let nrows: i64 = rows[0] 94 let o: *u8 = sys_mmap(total + nrows * (n_in + 2) + 1) 95 var oo: i64 = 0 96 let a4: *i64 = sys_mmap(QM_ROW_SLOTS * QM_I64) as *i64 97 let b4: *i64 = sys_mmap(QM_ROW_SLOTS * QM_I64) as *i64 98 var r: i64 = 0 99 while r < nrows { 100 let ab: *u8 = bufs[0] as *u8 101 if qm_row_at(ab, lens[0], cur[0], a4) == 0 { st[2] = r + 1; st[3] = 0; return QM_EXIT_REFUSED } 102 let id0: i64 = a4[0] 103 let id1: i64 = a4[2] 104 var i: i64 = id0 105 while i < id1 { o[oo] = ab[i]; oo = oo + 1; i = i + 1 } 106 o[oo] = QM_TAB as u8; oo = oo + 1 107 f = 0 108 while f < n_in { 109 let fb: *u8 = bufs[f] as *u8 110 if qm_row_at(fb, lens[f], cur[f], b4) == 0 { st[2] = r + 1; st[3] = f; return QM_EXIT_REFUSED } 111 if qm_ids_equal(ab, id0, id1, fb, b4[0], b4[2]) == 0 { st[2] = r + 1; st[3] = f; return QM_EXIT_REFUSED } 112 var t0: i64 = b4[2] 113 if t0 < b4[1] { t0 = t0 + 1 } // skip the TAB when present 114 if f > 0 { o[oo] = QM_SPACE as u8; oo = oo + 1 } 115 i = t0 116 while i < b4[1] { o[oo] = fb[i]; oo = oo + 1; i = i + 1 } 117 cur[f] = b4[3] + 1 // past the NL (or past the end) 118 f = f + 1 119 } 120 o[oo] = QM_NL as u8; oo = oo + 1 121 r = r + 1 122 } 123 let tmp: *u8 = sys_mmap(QM_PATH_CAP) 124 var tl: i64 = qm_cat(tmp, 0, out) 125 tl = qm_cat(tmp, tl, QM_TMP_SUFFIX) 126 let fd: i64 = sys_openat_wr(tmp, QM_FILE_MODE) 127 if fd < 0 { return QM_EXIT_IO } 128 if qm_write_all(fd, o, oo) != 0 { sys_close(fd); return QM_EXIT_IO } 129 sys_close(fd) 130 if sys_renameat(tmp, out) < 0 { return QM_EXIT_IO } 131 st[0] = nrows 132 st[1] = oo 133 return QM_EXIT_OK 134} 135func qm_report(out: *u8, n_in: i64, st: *i64, rc: i64) -> i64 { 136 if rc == QM_EXIT_OK { qm_puts("MERGED rows=" as *u8); qm_putn(st[0]); qm_puts(" files=" as *u8); qm_putn(n_in); qm_puts(" bytes=" as *u8); qm_putn(st[1]); qm_puts(" -> " as *u8); qm_puts(out); qm_puts("\n" as *u8); return 0 } 137 if rc == QM_EXIT_REFUSED { 138 if st[2] > 0 { qm_puts("REFUSED: query id mismatch at row " as *u8); qm_putn(st[2]); qm_puts(" in input " as *u8); qm_putn(st[3]); qm_puts(" -- a misaligned merge would score one query's rewrite against another's judgments\n" as *u8) } 139 else { qm_puts("REFUSED: row-count mismatch: input " as *u8); qm_putn(st[3]); qm_puts(" has " as *u8); qm_putn(st[4]); qm_puts(" rows, input 0 has " as *u8); qm_putn(st[0]); qm_puts("\n" as *u8) } 140 return 0 141 } 142 qm_puts("IO-FAIL input=" as *u8); qm_putn(st[3]); qm_puts("\n" as *u8) 143 return 0 144}