nx_par_pool.nx source
↩ module page · 147 lines · 6792 B
1// nx_par_pool.nx -- SOVEREIGN persistent worker POOL via PIPES (robust fix; busy-spin deadlocked on LICM,
2// ledger #21). Forks PP_NW workers ONCE; dispatches MANY matmuls via per-worker go/done PIPES -- kernel-
3// BLOCKING read/write (no busy-spin, no LICM hazard, no idle CPU burn). Data (A/B/C) in sys_mmap_shared (set
4// up BEFORE fork). Dispatch = write 1 byte to a worker's go-pipe; the worker computes its fixed range then
5// writes 1 byte to its done-pipe; shutdown = close the go-pipe (read returns 0 -> worker exits). Pure syscalls
6// (pipe2/fork/read/write/close/wait4) -- NO libvulkan/CUDA/pthread. Proves: fork-once + blocking dispatch =>
7// N-core speedup that HOLDS across many calls without per-call fork. license_tier: ORIGINAL
8import "nx_syscalls.nx"
9import "nx_itoa_lib.nx" // shared MSB-first emitter (zero-alloc)
10import "nx_f32.nx"
11import "nx_f32_hw.nx"
12
13const PP_M: i64 = 256
14const PP_K: i64 = 256
15const PP_N: i64 = 256
16const PP_NW: i64 = 8
17const PP_ITERS: i64 = 20
18
19func pp_puts(s: *u8) -> i64 { var n: i64 = 0; while s[n] != (0 as u8) { n = n + 1 } sys_write(1, s, n); return 0 }
20// MIGRATED to the shared emitter (debt 1785563586). The old body mmapped a scratch buffer
21// per call and never freed it. At PAGE granularity that is 4096B leaked PER CALL -- the
22// defect that took 28.5GB of a 36GB host in nx_ts_lumadiff (2MB input, ~3.66M calls).
23// nxi_* is MSB-first, allocates NOTHING, and emits identical bytes including the sign.
24func pp_putn(v: i64) -> i64 { nxi_out(v); return 0 }
25
26func pp_mm_range(A: *i64, B: *i64, C: *i64, K: i64, N: i64, lo: i64, hi: i64) -> i64 {
27 var idx: i64 = lo
28 while idx < hi {
29 let i: i64 = idx / N
30 let j: i64 = idx % N
31 var sum: i64 = 0
32 var kk: i64 = 0
33 while kk < K { sum = __f32_add(sum, __f32_mul(A[i * K + kk], B[kk * N + j])); kk = kk + 1 }
34 C[idx] = sum
35 idx = idx + 1
36 }
37 return 0
38}
39
40// worker: block on go-pipe read; on a byte, compute fixed range; signal done; EOF (close) -> exit
41func pp_worker(w: i64, go_rd: i64, done_wr: i64, A: *i64, B: *i64, C: *i64, M: i64, K: i64, N: i64) -> i64 {
42 let total: i64 = M * N
43 let lo: i64 = (w * total) / PP_NW
44 let hi: i64 = ((w + 1) * total) / PP_NW
45 let buf: *u8 = sys_mmap(8)
46 var live: i64 = 1
47 while live == 1 {
48 let got: i64 = sys_read(go_rd, buf, 1) // BLOCKS until dispatched
49 if got <= 0 { live = 0 } else {
50 pp_mm_range(A, B, C, K, N, lo, hi)
51 sys_write(done_wr, buf, 1) // signal done
52 }
53 }
54 return 0
55}
56
57func main() -> i64 {
58 let A: *i64 = sys_mmap_shared(PP_M * PP_K * 8) as *i64
59 let B: *i64 = sys_mmap_shared(PP_K * PP_N * 8) as *i64
60 let Cser: *i64 = sys_mmap(PP_M * PP_N * 8) as *i64
61 let Cpool: *i64 = sys_mmap_shared(PP_M * PP_N * 8) as *i64
62 let go_rd: *i64 = sys_mmap(PP_NW * 8) as *i64
63 let go_wr: *i64 = sys_mmap(PP_NW * 8) as *i64
64 let done_rd: *i64 = sys_mmap(PP_NW * 8) as *i64
65 let done_wr: *i64 = sys_mmap(PP_NW * 8) as *i64
66 var i: i64 = 0
67 while i < PP_M * PP_K { A[i] = f32_of((i % 5) - 2); i = i + 1 }
68 i = 0
69 while i < PP_K * PP_N { B[i] = f32_of((i % 7) - 3); i = i + 1 }
70
71 // create per-worker go + done pipes (fds[0]=read, fds[1]=write)
72 var w: i64 = 0
73 while w < PP_NW {
74 let fg: *i64 = sys_mmap(32) as *i64
75 if sys_pipe2(fg, 0) < 0 { pp_puts("pipe fail\n"); return 1 }
76 go_rd[w] = fg[0]; go_wr[w] = fg[1]
77 let fd: *i64 = sys_mmap(32) as *i64
78 if sys_pipe2(fd, 0) < 0 { pp_puts("pipe fail\n"); return 1 }
79 done_rd[w] = fd[0]; done_wr[w] = fd[1]
80 w = w + 1
81 }
82
83 // SERIAL baseline
84 let t0: i64 = sys_now_ms()
85 var it: i64 = 0
86 while it < PP_ITERS { pp_mm_range(A, B, Cser, PP_K, PP_N, 0, PP_M * PP_N); it = it + 1 }
87 let serial_ms: i64 = sys_now_ms() - t0
88
89 // fork the pool ONCE
90 w = 0
91 while w < PP_NW {
92 let pid: i64 = sys_fork()
93 if pid == 0 {
94 // child closes EVERY fd except its own go_rd[w] + done_wr[w], so the main's shutdown
95 // close(go_wr[w]) is the LAST writer -> read returns 0 (EOF) -> worker exits cleanly.
96 var v: i64 = 0
97 while v < PP_NW {
98 sys_close(go_wr[v]); sys_close(done_rd[v])
99 if v != w { sys_close(go_rd[v]); sys_close(done_wr[v]) }
100 v = v + 1
101 }
102 pp_worker(w, go_rd[w], done_wr[w], A, B, Cpool, PP_M, PP_K, PP_N); sys_exit(0)
103 }
104 w = w + 1
105 }
106 // parent closes the ends it never uses (it writes go_wr, reads done_rd)
107 w = 0
108 while w < PP_NW { sys_close(go_rd[w]); sys_close(done_wr[w]); w = w + 1 }
109
110 // POOL: PP_ITERS dispatches (no per-call fork; blocking pipe handshake)
111 let buf: *u8 = sys_mmap(8)
112 let t1: i64 = sys_now_ms()
113 it = 0
114 while it < PP_ITERS {
115 w = 0
116 while w < PP_NW { sys_write(go_wr[w], buf, 1); w = w + 1 } // dispatch all
117 w = 0
118 while w < PP_NW { sys_read(done_rd[w], buf, 1); w = w + 1 } // join all (blocks)
119 it = it + 1
120 }
121 let pool_ms: i64 = sys_now_ms() - t1
122
123 // shutdown (close go-pipes -> workers' read returns 0 -> exit) + reap
124 w = 0
125 while w < PP_NW { sys_close(go_wr[w]); w = w + 1 }
126 let stp: *i64 = sys_mmap(8) as *i64
127 var d: i64 = 0
128 while d < PP_NW { sys_wait4(0 - 1, stp, 0); d = d + 1 }
129
130 var mism: i64 = 0
131 i = 0
132 while i < PP_M * PP_N { if Cpool[i] != Cser[i] { mism = mism + 1 } i = i + 1 }
133
134 pp_puts("=== nx_par_pool: SOVEREIGN worker pool via PIPES (fork ONCE, blocking dispatch) ===\n")
135 pp_puts(" "); pp_putn(PP_M); pp_puts("x"); pp_putn(PP_K); pp_puts("x"); pp_putn(PP_N)
136 pp_puts(" matmul x "); pp_putn(PP_ITERS); pp_puts(" iters, "); pp_putn(PP_NW); pp_puts(" pooled workers\n")
137 pp_puts(" serial = "); pp_putn(serial_ms); pp_puts(" ms\n")
138 pp_puts(" pool = "); pp_putn(pool_ms); pp_puts(" ms (fork ONCE before the loop)\n")
139 if pool_ms > 0 { pp_puts(" speedup= "); pp_putn((serial_ms * 100) / pool_ms); pp_puts(" /100x ("); pp_putn(serial_ms / pool_ms); pp_puts("x)\n") }
140 pp_puts(" mismatches = "); pp_putn(mism); pp_puts(" (0 = bit-identical)\n")
141 var pass: i64 = 0
142 if mism == 0 { pp_puts(" T1 pool result bit-identical: PASS\n"); pass = pass + 1 } else { pp_puts(" T1: FAIL\n") }
143 if pool_ms < serial_ms { pp_puts(" T2 pool faster across many calls (fork-once + blocking dispatch): PASS\n"); pass = pass + 1 } else { pp_puts(" T2: FAIL\n") }
144 pp_puts("\n PASS="); pp_putn(pass); pp_puts("/2 ")
145 if pass == 2 { pp_puts("VERDICT=GREEN (sovereign pipe worker pool: the fix for per-call fork overhead)\n"); sys_exit(0); return 0 }
146 pp_puts("VERDICT=RED\n"); sys_exit(1); return 1
147}