code wiki / _hdl_build / nx_gen_multiworker_gate.nx

nx_gen_multiworker_gate.nx source

↩ module page · 188 lines · 12706 B

1// nx_gen_multiworker_gate.nx -- proves the orchestrator's SOVEREIGN multi-worker dispatch (worker list comes from 2// cfg = the supervisor's launch command, NO files/TSVs). Real loopback rig: fork the REAL orchestrator (go_serve_conn) 3// + a mock GPU worker, then POST /api/generate over TCP: 4// (1) availability health-check: a listening worker -> up, an unbound port -> down (non-blocking, no hang) 5// (2) availability selection WITH FAILOVER: a down primary + a down extra are skipped, the request reaches the 6// one reachable worker (cfg extra #2) and returns CIDs 7// (3) all workers down -> graceful 503 (never a hang, never a crash) 8// license_tier: ORIGINAL expect_exit: 0 9import "nx_gen_orchestrator.nx" // go_worker_up, go_serve_conn, go_handle_generate 10import "nx_syscalls.nx" 11 12import "nx_connect.nx" // bounded connect: a raw sys_connect hangs ~127s on a black-holed host 13func ig(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} sys_write(1,s,n); return 0 } 14func ig_n(v: i64) -> i64 { var m: i64=v; if m<0{ig("-" as *u8);m=0-m} let t:*u8=sys_mmap(24); var k:i64=0; if m==0{t[0]=48 as u8;k=1} while m>0{t[k]=(48+(m%10)) as u8;m=m/10;k=k+1} var i:i64=0; let o:*u8=sys_mmap(24); while i<k{o[i]=t[k-1-i];i=i+1} sys_write(1,o,k); return 0 } 15func ig_row(id: *u8, ok: i64, pass: *i64) -> i64 { ig(" " as *u8); ig(id); ig(": " as *u8); if ok==1 { ig("OK\n" as *u8); pass[0]=pass[0]+1 } else { ig("FAIL\n" as *u8) } return 0 } 16func ig_cat(dst: *u8, off: i64, s: *u8) -> i64 { var o: i64=off; var i: i64=0; while s[i]!=(0 as u8){dst[o]=s[i]; o=o+1; i=i+1} return o } 17func ig_catb(dst: *u8, off: i64, src: *u8, n: i64) -> i64 { var o: i64=off; var i: i64=0; while i<n {dst[o]=src[i]; o=o+1; i=i+1} return o } 18func ig_itoa(dst: *u8, off: i64, v: i64) -> i64 { let t: *u8=sys_mmap(28); var m: i64=v; var k: i64=0; if m==0{t[0]=48 as u8;k=1} while m>0{t[k]=(48+(m%10)) as u8;m=m/10;k=k+1} var o: i64=off; var q: i64=k-1; while q>=0{dst[o]=t[q];o=o+1;q=q-1} return o } 19func ig_find(buf: *u8, n: i64, needle: *u8) -> i64 { var nl: i64=0; while needle[nl]!=(0 as u8){nl=nl+1} if nl==0 {return 0} var i: i64=0; while i+nl<=n { var j: i64=0; var ok: i64=1; while j<nl { if (buf[i+j]&0xff)!=(needle[j]&0xff){ok=0;j=nl} else {j=j+1} } if ok==1 {return i} i=i+1 } return 0-1 } 20func ig_has(buf: *u8, n: i64, needle: *u8) -> i64 { if ig_find(buf,n,needle)>=0 {return 1} return 0 } 21func ig_nap(ms: i64) -> i64 { let ts: *i64 = sys_mmap(16) as *i64; ts[0]=0; ts[1]=ms*1000000; __syscall(35, ts as i64, 0, 0, 0, 0, 0); return 0 } 22func ig_freshns(prefix: *u8) -> i64 { let mfn: *u8 = sys_mmap(256); var mo: i64=0; while (prefix[mo]&0xff)!=0 { mfn[mo]=prefix[mo]; mo=mo+1 } let suf: *u8 = "manifest.txt" as *u8; var so: i64=0; while (suf[so]&0xff)!=0 { mfn[mo+so]=suf[so]; so=so+1 } mfn[mo+so]=0 as u8; let fd: i64 = sys_openat_wr(mfn, 0x1a4); if fd>=0 {sys_close(fd)} return 0 } 23 24func ig_basepng(b: *u8) -> i64 { 25 b[0]=137 as u8; b[1]=80 as u8; b[2]=78 as u8; b[3]=71 as u8 26 b[4]=13 as u8; b[5]=10 as u8; b[6]=26 as u8; b[7]=10 as u8 27 b[8]=0 as u8; b[9]=0 as u8; b[10]=0 as u8; b[11]=14 as u8 28 b[12]=73 as u8; b[13]=72 as u8; b[14]=68 as u8; b[15]=82 as u8 29 b[16]=0 as u8; b[17]=0 as u8; b[18]=0 as u8; b[19]=0 as u8; b[20]=1 as u8 30 b[21]=0 as u8; b[22]=0 as u8; b[23]=0 as u8; b[24]=1 as u8; b[25]=8 as u8 31 b[26]=6 as u8; b[27]=0 as u8; b[28]=0 as u8; b[29]=0 as u8 32 b[30]=115 as u8; b[31]=22 as u8; b[32]=206 as u8; b[33]=242 as u8 33 b[34]=0 as u8; b[35]=0 as u8; b[36]=0 as u8; b[37]=0 as u8 34 b[38]=73 as u8; b[39]=69 as u8; b[40]=78 as u8; b[41]=68 as u8 35 b[42]=174 as u8; b[43]=66 as u8; b[44]=96 as u8; b[45]=130 as u8 36 return 46 37} 38func ig_enc(v: i64) -> i64 { let x: i64=v&0x3f; if x<26 {return 65+x} if x<52 {return 97+(x-26)} if x<62 {return 48+(x-52)} if x==62 {return 43} return 47 } 39func ig_b64encode(inp: *u8, n: i64, out: *u8) -> i64 { 40 var i: i64 = 0; var o: i64 = 0 41 while i + 3 <= n { 42 let b0: i64=inp[i]&0xff; let b1: i64=inp[i+1]&0xff; let b2: i64=inp[i+2]&0xff 43 out[o]=ig_enc(b0>>2) as u8; out[o+1]=ig_enc(((b0<<4)|(b1>>4))&0x3f) as u8; out[o+2]=ig_enc(((b1<<2)|(b2>>6))&0x3f) as u8; out[o+3]=ig_enc(b2&0x3f) as u8 44 o=o+4; i=i+3 45 } 46 let rem: i64 = n - i 47 if rem == 1 { let b0: i64=inp[i]&0xff; out[o]=ig_enc(b0>>2) as u8; out[o+1]=ig_enc((b0<<4)&0x3f) as u8; out[o+2]=61 as u8; out[o+3]=61 as u8; o=o+4 } 48 if rem == 2 { let b0: i64=inp[i]&0xff; let b1: i64=inp[i+1]&0xff; out[o]=ig_enc(b0>>2) as u8; out[o+1]=ig_enc(((b0<<4)|(b1>>4))&0x3f) as u8; out[o+2]=ig_enc((b1<<2)&0x3f) as u8; out[o+3]=61 as u8; o=o+4 } 49 return o 50} 51func ig_listen(port: i64) -> i64 { 52 let lfd: i64 = sys_socket(2,1,0); if lfd<0 {return 0-1} 53 let optv: *u8 = sys_mmap(4); optv[0]=1 as u8; sys_setsockopt(lfd,1,2,optv,4) 54 let addr: *u8 = sys_mmap(16) 55 addr[0]=2 as u8; addr[1]=0 as u8; addr[2]=((port>>8)&0xff) as u8; addr[3]=(port&0xff) as u8 56 addr[4]=127 as u8; addr[5]=0 as u8; addr[6]=0 as u8; addr[7]=1 as u8 57 var zi: i64=8; while zi<16 {addr[zi]=0 as u8; zi=zi+1} 58 if sys_bind(lfd, addr, 16) < 0 { return 0-1 } 59 if sys_listen(lfd, 16) < 0 { return 0-1 } 60 return lfd 61} 62// connect 127.0.0.1:port (retry), send raw, read until close. 63func ig_send(port: i64, raw: *u8, rawlen: i64, resp: *u8, cap: i64) -> i64 { 64 var tries: i64 = 0 65 while tries < 300 { 66 let fd: i64 = sys_socket(2,1,0) 67 if fd >= 0 { 68 let tv: *u8 = sys_mmap(16); tv[0]=15 as u8; var tz: i64=1; while tz<16 {tv[tz]=0 as u8; tz=tz+1} 69 sys_setsockopt(fd, 1, 20, tv, 16) 70 let addr: *u8 = sys_mmap(16) 71 addr[0]=2 as u8; addr[1]=0 as u8; addr[2]=((port>>8)&0xff) as u8; addr[3]=(port&0xff) as u8 72 addr[4]=127 as u8; addr[5]=0 as u8; addr[6]=0 as u8; addr[7]=1 as u8 73 var z: i64=8; while z<16 {addr[z]=0 as u8; z=z+1} 74 if nx_connect_bounded(fd, addr, 16, NX_CONN_DEFAULT_MS) == 0 { 75 sys_write(fd, raw, rawlen) 76 var rn: i64 = 0; var g: i64 = 1 77 while g==1 { let r: i64 = sys_read(fd, (resp as i64+rn) as *u8, cap-rn); if r<=0 {g=0} else {rn=rn+r; if rn>=cap {g=0}} } 78 sys_close(fd) 79 return rn 80 } 81 sys_close(fd) 82 } 83 ig_nap(15) 84 tries = tries + 1 85 } 86 return 0 - 1 87} 88 89func main() -> i64 { 90 let pass: *i64 = sys_mmap(8) as *i64; pass[0]=0 91 ig("=== NX-GEN-MULTIWORKER GATE (sovereign cfg worker list -- NO files/TSVs; health + selection + failover) ===\n" as *u8) 92 93 let UP_PORT: i64 = 39501 94 let DOWN_PORT: i64 = 39502 // no listener -> down 95 let DOWN_PRIMARY: i64 = 39503 // no listener -> down (the primary worker) 96 let ORCH2: i64 = 39504 // orchestrator under test for the FAILOVER case 97 let ORCH3: i64 = 39505 // orchestrator under test for the ALL-DOWN case 98 99 let png: *u8 = sys_mmap(256); let plen: i64 = ig_basepng(png) 100 let b64: *u8 = sys_mmap(512); let b64len: i64 = ig_b64encode(png, plen, b64) 101 102 // cfg for the FAILOVER orchestrator: primary DOWN + extras = [DOWN_PORT, UP_PORT] (all from cfg, no file). 103 ig_freshns("/tmp/mwgate-" as *u8); ig_freshns("/tmp/mwgateblob-" as *u8) 104 let cfg2: *i64 = sys_mmap(8*32) as *i64 105 cfg2[0]=127; cfg2[1]=0; cfg2[2]=0; cfg2[3]=1; cfg2[4]=DOWN_PRIMARY 106 cfg2[5]="/tmp/mwgate-" as *u8 as i64; cfg2[6]="/tmp/mwgateblob-" as *u8 as i64; cfg2[7]="/tmp/mwgate_cids.tsv" as *u8 as i64; cfg2[8]="laptop" as *u8 as i64 107 cfg2[9]=2 108 cfg2[10]=127; cfg2[11]=0; cfg2[12]=0; cfg2[13]=1; cfg2[14]=DOWN_PORT 109 cfg2[15]=127; cfg2[16]=0; cfg2[17]=0; cfg2[18]=1; cfg2[19]=UP_PORT 110 // cfg for the ALL-DOWN orchestrator: primary DOWN + 1 down extra -> must 503. 111 let cfg3: *i64 = sys_mmap(8*32) as *i64 112 cfg3[0]=127; cfg3[1]=0; cfg3[2]=0; cfg3[3]=1; cfg3[4]=DOWN_PRIMARY 113 cfg3[5]="/tmp/mwgate-" as *u8 as i64; cfg3[6]="/tmp/mwgateblob-" as *u8 as i64; cfg3[7]="/tmp/mwgate_cids.tsv" as *u8 as i64; cfg3[8]="laptop" as *u8 as i64 114 cfg3[9]=1 115 cfg3[10]=127; cfg3[11]=0; cfg3[12]=0; cfg3[13]=1; cfg3[14]=DOWN_PORT 116 117 let upfd: i64 = ig_listen(UP_PORT); if upfd<0 { ig("FAIL up listen\n" as *u8); sys_exit(1); return 1 } 118 let ofd2: i64 = ig_listen(ORCH2); if ofd2<0 { ig("FAIL orch2 listen\n" as *u8); sys_exit(1); return 1 } 119 let ofd3: i64 = ig_listen(ORCH3); if ofd3<0 { ig("FAIL orch3 listen\n" as *u8); sys_exit(1); return 1 } 120 121 // mock GPU worker: tolerates empty health-check connects + serves the real /v1/images/generations POST. 122 let pidW: i64 = sys_fork() 123 if pidW == 0 { 124 sys_close(ofd2); sys_close(ofd3) 125 var k: i64 = 0 126 while k < 6 { 127 let c: i64 = sys_accept(upfd) 128 if c >= 0 { 129 let rb: *u8 = sys_mmap(16384); let r: i64 = sys_read(c, rb, 16384) 130 if r > 0 { 131 let bod: *u8 = sys_mmap(4096); var bo: i64 = 0 132 bo = ig_cat(bod, bo, "{\"data\":[{\"b64_json\":\"" as *u8); bo = ig_catb(bod, bo, b64, b64len); bo = ig_cat(bod, bo, "\"}]}" as *u8) 133 let resp: *u8 = sys_mmap(8192); var o: i64 = 0 134 o = ig_cat(resp, o, "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nConnection: close\r\nContent-Length: " as *u8) 135 o = ig_itoa(resp, o, bo); o = ig_cat(resp, o, "\r\n\r\n" as *u8); o = ig_catb(resp, o, bod, bo) 136 sys_write(c, resp, o) 137 } 138 sys_close(c) 139 } 140 k = k + 1 141 } 142 sys_exit(0) 143 } 144 // failover orchestrator (serve 1) 145 let pidO2: i64 = sys_fork() 146 if pidO2 == 0 { sys_close(upfd); sys_close(ofd3); let c: i64 = sys_accept(ofd2); if c>=0 { go_serve_conn(c, cfg2); sys_close(c) } sys_exit(0) } 147 // all-down orchestrator (serve 1) 148 let pidO3: i64 = sys_fork() 149 if pidO3 == 0 { sys_close(upfd); sys_close(ofd2); let c: i64 = sys_accept(ofd3); if c>=0 { go_serve_conn(c, cfg3); sys_close(c) } sys_exit(0) } 150 sys_close(upfd); sys_close(ofd2); sys_close(ofd3) 151 ig_nap(150) 152 153 // (1) availability health-check (non-blocking; a down port must not hang): UP -> 1, DOWN -> 0 154 let up_ok: i64 = go_worker_up(127, 0, 0, 1, UP_PORT) 155 let down_ok: i64 = go_worker_up(127, 0, 0, 1, DOWN_PORT) 156 let r1: i64 = ((up_ok==1) as i64) & ((down_ok==0) as i64) 157 ig(" health: up_listener="); ig_n(up_ok); ig(" down_port="); ig_n(down_ok); ig("\n" as *u8) 158 159 let raw: *u8 = sys_mmap(8192); let resp: *u8 = sys_mmap(262144) 160 let gb: *u8 = "{\"prompt\":\"a fox\",\"width\":256,\"height\":256,\"steps\":8,\"cfg\":\"1.0\",\"seed\":7,\"count\":1,\"sampler\":\"euler\"}" as *u8 161 var gbl: i64 = 0; while gb[gbl]!=(0 as u8){gbl=gbl+1} 162 163 // (2) FAILOVER: POST to ORCH2 (primary down, extra0 down, extra1 UP) -> reaches UP -> 200 + CIDs 164 var rl: i64 = ig_cat(raw, 0, "POST /api/generate HTTP/1.0\r\nContent-Type: application/json\r\nConnection: close\r\nContent-Length: " as *u8) 165 rl = ig_itoa(raw, rl, gbl); rl = ig_cat(raw, rl, "\r\n\r\n" as *u8); rl = ig_catb(raw, rl, gb, gbl) 166 let n2: i64 = ig_send(ORCH2, raw, rl, resp, 262144) 167 let r2: i64 = ((ig_has(resp, n2, "200 OK" as *u8)==1) as i64) & ((ig_has(resp, n2, "\"cids\":[\"" as *u8)==1) as i64) 168 ig(" failover resp bytes="); ig_n(n2); ig("\n" as *u8) 169 170 // (3) ALL DOWN: POST to ORCH3 (primary down + 1 down extra) -> graceful 503 171 rl = ig_cat(raw, 0, "POST /api/generate HTTP/1.0\r\nContent-Type: application/json\r\nConnection: close\r\nContent-Length: " as *u8) 172 rl = ig_itoa(raw, rl, gbl); rl = ig_cat(raw, rl, "\r\n\r\n" as *u8); rl = ig_catb(raw, rl, gb, gbl) 173 let n3: i64 = ig_send(ORCH3, raw, rl, resp, 262144) 174 let r3: i64 = (ig_has(resp, n3, "503" as *u8)==1) as i64 175 ig(" all-down resp bytes="); ig_n(n3); ig("\n" as *u8) 176 177 __syscall(129, pidW, 9, 0, 0, 0, 0) // rv64 kill=129 -- raw x86 62 is an RV64 KEY the backend translated to lseek(8), so these SIGKILLs never happened (debt idx 2277) 178 __syscall(129, pidO2, 9, 0, 0, 0, 0) 179 __syscall(129, pidO3, 9, 0, 0, 0, 0) 180 let st: *i64 = sys_mmap(8) as *i64; var rr: i64=0; while rr<20 { sys_wait4(0-1, st, 1); rr=rr+1 } 181 182 ig_row("availability health-check: up listener -> available, down port -> not (non-blocking, no hang)" as *u8, r1, pass) 183 ig_row("availability selection + FAILOVER: down primary + down extra skipped -> reaches UP -> CIDs" as *u8, r2, pass) 184 ig_row("all workers down -> graceful 503 (never a hang or crash)" as *u8, r3, pass) 185 ig("GEN-MULTIWORKER-GATE rows=3 pass="); ig_n(pass[0]) 186 if pass[0]==3 { ig(" verdict=GREEN\n" as *u8); sys_exit(0); return 0 } 187 ig(" verdict=RED\n" as *u8); sys_exit(1); return 1 188}