code wiki / (root) / nx_par_pool.nx

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}