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}