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}