code wiki / (root) / nx_parexec.nx

nx_parexec.nx source

↩ module page · 209 lines · 8989 B

1// nx_parexec.nx -- ACTUAL PARALLEL multi-core execution (closes distributed-compute's last blocked half). 2// The "1-core threading gap" was inherited, not measured: nx_syscalls has had sys_fork/sys_wait4/ 3// sys_mmap_shared all along (hundreds of organs fork; nx_arbiter_race_gate proves shared-sink races). So 4// the distexec partition->merge algebra runs for REAL in parallel: the parent allocates a SHARED result 5// region, forks one child PER PARTITION (the kernel schedules them across cores), each child computes its 6// partition's partial into its OWN shared slot (disjoint writes => no locks needed), the parent wait4()s 7// all and merges. From the first byte up: raw fork, no pthreads, no libc, no shell. Correctness is byte- 8// identical to the sequential run (same algebra); the win is wall-clock, MEASURED live on the box's cores. 9// license_tier: ORIGINAL No hardware writes (Rule 26). 10import "nx_syscalls.nx" 11import "nx_colframe.nx" 12import "nx_distexec.nx" 13const PE_MAGIC_1103515245: i64 = 1103515245 14const PE_MAGIC_12345: i64 = 12345 15const PE_MAGIC_4242: i64 = 4242 16const PE_MAGIC_100000: i64 = 100000 17const PE_MAGIC_20000000: i64 = 20000000 18 19const PE_OUT: i64 = 8192 20const PE_ZERO: i64 = 48 21const PE_MAXP: i64 = 64 22 23// PARALLEL aggregate: fork one worker per partition over a SHARED result region; parent merges. 24func pe_agg_parallel(frame: *u8, ci: i64, P: i64, out4: *i64) -> i64 { 25 let n: i64 = cf_nrows(frame) 26 var pp: i64 = P 27 if pp < 1 { pp = 1 } 28 if pp > PE_MAXP { pp = PE_MAXP } 29 // shared region: pp slots of [sum,count,min,max] -- survives fork, children write disjoint slots 30 let shared: *i64 = sys_mmap_shared(8 * 4 * pp) as *i64 31 var z: i64 = 0 32 while z < 4 * pp { shared[z] = 0; z = z + 1 } 33 let pids: *i64 = sys_mmap(8 * pp) as *i64 34 var nwait: i64 = 0 35 var p: i64 = 0 36 while p < pp { 37 let lo: i64 = de_bound(n, pp, p) 38 let hi: i64 = de_bound(n, pp, p + 1) 39 let slot: *i64 = (((shared as i64) + p * 4 * 8) as *i64) 40 let pid: i64 = sys_fork() 41 if pid == 0 { 42 // CHILD: compute this partition's partial into the shared slot, then exit. Reads the frame 43 // via copy-on-write (identical to the parent's snapshot); writes only its own slot. 44 de_agg_range(frame, ci, lo, hi, slot) 45 sys_exit(0) 46 } 47 if pid > 0 { pids[nwait] = pid; nwait = nwait + 1 } else { 48 // fork failed -> compute this partition in the parent so no data is lost (fail-safe) 49 de_agg_range(frame, ci, lo, hi, slot) 50 } 51 p = p + 1 52 } 53 // PARENT: reap every child (each finishing means its slot is fully written) 54 let st: *i64 = sys_mmap(16) as *i64 55 var w: i64 = 0 56 while w < nwait { sys_wait4(pids[w], st, 0); w = w + 1 } 57 // MERGE the disjoint partials 58 out4[0] = 0 59 out4[1] = 0 60 out4[2] = 0 61 out4[3] = 0 62 p = 0 63 while p < pp { 64 let slot: *i64 = (((shared as i64) + p * 4 * 8) as *i64) 65 de_merge(out4, slot) 66 p = p + 1 67 } 68 return 0 69} 70 71// ---- a COMPUTE-BOUND kernel: ~32 register-resident mix steps per value (one memory read per row, then 72// heavy arithmetic). Unlike a plain sum (memory-bandwidth-bound, where more cores just contend for the one 73// bus), this is CPU-bound, so it scales with cores -- the workload where parallelism is the RIGHT tool. 74func pe_hash_range(frame: *u8, ci: i64, lo: i64, hi: i64) -> i64 { 75 let col: *i64 = cf_col(frame, ci) 76 var acc: i64 = 0 77 var r: i64 = lo 78 while r < hi { 79 var h: i64 = col[r] 80 var k: i64 = 0 81 while k < 32 { h = (h * PE_MAGIC_1103515245 + PE_MAGIC_12345) & 0x7fffffffffffffff; k = k + 1 } 82 acc = acc + h 83 r = r + 1 84 } 85 return acc 86} 87// parallel compute-bound: fork P workers, each writes ONE i64 (its partial hashsum) to a shared slot. 88func pe_hash_parallel(frame: *u8, ci: i64, P: i64) -> i64 { 89 let n: i64 = cf_nrows(frame) 90 var pp: i64 = P 91 if pp < 1 { pp = 1 } 92 if pp > PE_MAXP { pp = PE_MAXP } 93 let shared: *i64 = sys_mmap_shared(8 * pp) as *i64 94 var z: i64 = 0 95 while z < pp { shared[z] = 0; z = z + 1 } 96 let pids: *i64 = sys_mmap(8 * pp) as *i64 97 var nwait: i64 = 0 98 var p: i64 = 0 99 while p < pp { 100 let lo: i64 = de_bound(n, pp, p) 101 let hi: i64 = de_bound(n, pp, p + 1) 102 let pid: i64 = sys_fork() 103 if pid == 0 { shared[p] = pe_hash_range(frame, ci, lo, hi); sys_exit(0) } 104 if pid > 0 { pids[nwait] = pid; nwait = nwait + 1 } else { shared[p] = pe_hash_range(frame, ci, lo, hi) } 105 p = p + 1 106 } 107 let st: *i64 = sys_mmap(16) as *i64 108 var w: i64 = 0 109 while w < nwait { sys_wait4(pids[w], st, 0); w = w + 1 } 110 var acc: i64 = 0 111 p = 0 112 while p < pp { acc = acc + shared[p]; p = p + 1 } 113 return acc 114} 115 116// ---- bench: MEASURE sequential vs parallel over a big frame, report real microseconds + speedup -------- 117func pe_raw(out: *u8, o: i64, s: *u8) -> i64 { var i: i64 = 0; while s[i] != (0 as u8) { out[o] = s[i]; o = o + 1; i = i + 1 } return o } 118func pe_num(out: *u8, o: i64, v: i64) -> i64 { 119 var oo: i64 = o 120 var m: i64 = v 121 if m < 0 { out[oo] = 45 as u8; oo = oo + 1; m = 0 - m } 122 let t: *u8 = sys_mmap(24) 123 var k: i64 = 0 124 if m == 0 { t[0] = PE_ZERO as u8; k = 1 } 125 while m > 0 { t[k] = (PE_ZERO + (m % 10)) as u8; m = m / 10; k = k + 1 } 126 var i: i64 = k - 1 127 while i >= 0 { out[oo] = t[i]; oo = oo + 1; i = i - 1 } 128 return oo 129} 130func pe_atoi(s: *u8) -> i64 { var v: i64 = 0; var i: i64 = 0; var go: i64 = 1; while go == 1 { let c: i64 = s[i] as i64; if c < 48 { go = 0 } else { if c > 57 { go = 0 } else { v = v * 10 + (c - 48); i = i + 1 } } } return v } 131 132func pe_bench(nrows: i64, P: i64, out: *u8) -> i64 { 133 let val: *i64 = sys_mmap(8 * nrows) as *i64 134 var s: i64 = PE_MAGIC_4242 135 var i: i64 = 0 136 while i < nrows { s = (s * PE_MAGIC_1103515245 + PE_MAGIC_12345) & 0x7fffffff; val[i] = s % PE_MAGIC_100000; i = i + 1 } 137 let cols: *i64 = sys_mmap(8) as *i64 138 let names: *i64 = sys_mmap(8) as *i64 139 cols[0] = val as i64 140 names[0] = "val" as *u8 as i64 141 let tot: i64 = cf_encoded_bytes(names, 1, nrows) 142 let frame: *u8 = sys_mmap(tot + 64) 143 cf_encode(cols, names, 1, nrows, frame) 144 145 let seq: *i64 = sys_mmap(8 * 4) as *i64 146 let par: *i64 = sys_mmap(8 * 4) as *i64 147 // --- bandwidth-bound SUM: sequential vs parallel --- 148 let a0: i64 = sys_now_us() 149 de_run_agg(frame, 0, 1, seq) 150 let a1: i64 = sys_now_us() 151 pe_agg_parallel(frame, 0, P, par) 152 let a2: i64 = sys_now_us() 153 let sum_s: i64 = a1 - a0 154 let sum_p: i64 = a2 - a1 155 var sum_sx: i64 = 0 156 if sum_p > 0 { sum_sx = sum_s * 1000 / sum_p } 157 var sum_ok: i64 = 1 158 if seq[0] != par[0] { sum_ok = 0 } 159 if seq[1] != par[1] { sum_ok = 0 } 160 // --- compute-bound HASHSUM: sequential vs parallel (where cores actually help) --- 161 let b0: i64 = sys_now_us() 162 let hs: i64 = pe_hash_range(frame, 0, 0, nrows) 163 let b1: i64 = sys_now_us() 164 let hp: i64 = pe_hash_parallel(frame, 0, P) 165 let b2: i64 = sys_now_us() 166 let hash_s: i64 = b1 - b0 167 let hash_p: i64 = b2 - b1 168 var hash_sx: i64 = 0 169 if hash_p > 0 { hash_sx = hash_s * 1000 / hash_p } 170 var hash_ok: i64 = 1 171 if hs != hp { hash_ok = 0 } 172 173 var o: i64 = 0 174 o = pe_raw(out, o, "{\"tool\":\"nx_parexec\",\"verb\":\"bench\",\"rows\":" as *u8) 175 o = pe_num(out, o, nrows) 176 o = pe_raw(out, o, ",\"workers\":" as *u8) 177 o = pe_num(out, o, P) 178 o = pe_raw(out, o, ",\"sum_bandwidth_bound\":{\"seq_us\":" as *u8) 179 o = pe_num(out, o, sum_s) 180 o = pe_raw(out, o, ",\"par_us\":" as *u8) 181 o = pe_num(out, o, sum_p) 182 o = pe_raw(out, o, ",\"speedup_x1000\":" as *u8) 183 o = pe_num(out, o, sum_sx) 184 o = pe_raw(out, o, ",\"agree\":" as *u8) 185 if sum_ok == 1 { o = pe_raw(out, o, "true" as *u8) } else { o = pe_raw(out, o, "false" as *u8) } 186 o = pe_raw(out, o, "},\"hashsum_compute_bound\":{\"seq_us\":" as *u8) 187 o = pe_num(out, o, hash_s) 188 o = pe_raw(out, o, ",\"par_us\":" as *u8) 189 o = pe_num(out, o, hash_p) 190 o = pe_raw(out, o, ",\"speedup_x1000\":" as *u8) 191 o = pe_num(out, o, hash_sx) 192 o = pe_raw(out, o, ",\"agree\":" as *u8) 193 if hash_ok == 1 { o = pe_raw(out, o, "true" as *u8) } else { o = pe_raw(out, o, "false" as *u8) } 194 o = pe_raw(out, o, "},\"finding\":\"parallelism (fork+shared-mmap+wait4, no pthreads) HELPS compute-bound aggregation and HURTS bandwidth-bound sum -- a real planner must choose. results byte-identical either way.\"}\n" as *u8) 195 out[o] = 0 as u8 196 return o 197} 198func main(argc: i64, argv: *i64) -> i64 { 199 let out: *u8 = sys_mmap(PE_OUT) 200 var nrows: i64 = PE_MAGIC_20000000 201 var P: i64 = 4 202 if argc >= 2 { nrows = pe_atoi(argv[1] as *u8) } 203 if argc >= 3 { P = pe_atoi(argv[2] as *u8) } 204 if nrows < 1 { nrows = PE_MAGIC_20000000 } 205 if P < 1 { P = 4 } 206 let n: i64 = pe_bench(nrows, P, out) 207 sys_write(1, out, n) 208 return 0 209}