nx_llm_batch_serve.nx source
↩ module page · 288 lines · 12980 B
1// nx_llm_batch_serve.nx -- the SOVEREIGN CONTINUOUS-BATCHING LLM SEAT.
2// Wraps nx_llm_sched (batch-invariance GATED by nx_llm_sched_gate) in the
3// proven nx_f32_llm_serve accept-loop pattern: binds 127.0.0.1:11435
4// (:11434 stays the production single-stream seat), accept WITH TIMEOUT
5// between decode rounds -- new POST /chat requests admit MID-FLIGHT and
6// join the running batch; each round = ONE batched forward for every
7// active request; finished requests get their OpenAI-chat JSON and close.
8//
9// Modes: (no args) = serve forever. `once <N>` = serve exactly N requests
10// then exit 0 (the testable smoke mode).
11// expect: LIVE on :11435 (blocks) / exit 0 in once-mode.
12// license_tier: ORIGINAL module: nishi-core.seat.llm-batch
13
14import "nx_syscalls.nx"
15import "nx_connect.nx" // bounded connect: a raw sys_connect hangs ~127s on a black-holed host
16import "nx_tier.nx"
17import "nx_le.nx"
18import "nx_bpe.nx"
19import "nx_gguf.nx"
20import "nx_gguf_load.nx"
21import "nx_gguf_meta.nx"
22import "nx_f32.nx"
23import "nx_f32_kv_cache.nx"
24import "nx_f32_lazy_weight.nx"
25import "nx_f32_llama_block.nx"
26import "nx_f32_llama_block_v4.nx"
27import "nx_f32_llama_stack_v4.nx"
28import "nx_f32_llama_layer_lazy_load.nx"
29import "nx_f32_llm.nx"
30import "nx_f32_llm_v4.nx"
31import "nx_f32_llm_read_dims.nx"
32import "nx_f32_bpe_load.nx"
33import "nx_f32_llm_special_tokens.nx"
34import "nx_f32_sampler.nx"
35import "nx_prng.nx"
36import "nx_reasoning.nx"
37import "nx_kvcache.nx"
38import "nx_f32_attn_paged.nx"
39import "nx_f32_llama_v4p.nx"
40import "nx_f32_llama_v4b.nx"
41import "nx_llm_sched.nx"
42import "nx_http_client.nx"
43
44// ONE integer port; sockaddr computed (shared helper), self-probed at boot.
45const NX_BSEAT_PORT: i64 = 11435
46
47func bs_puts(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} sys_write(1,s,n); return 0 }
48func bs_slen(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} return n }
49func bs_find(buf: *u8, n: i64, ndl: *u8, nl: i64) -> i64 {
50 if nl<=0 { return 0-1 }
51 var i: i64=0
52 while i+nl<=n { var j: i64=0; var hit: i64=1; while j<nl { if buf[i+j]!=ndl[j] { hit=0; j=nl } else { j=j+1 } } if hit==1 { return i } i=i+1 }
53 return 0-1
54}
55func bs_cat(dst: *u8, off: i64, s: *u8) -> i64 { var i: i64=0; while s[i]!=(0 as u8){ dst[off]=s[i]; off=off+1; i=i+1 } return off }
56func bs_write_all(fd: i64, buf: *u8, n: i64) -> i64 { var off: i64=0; while off<n { let w: i64=sys_write(fd, ((buf as i64)+off) as *u8, n-off); if w<=0 { off=n } else { off=off+w } } return 0 }
57func bs_extract_prompt(body: *u8, blen: i64, out: *u8, cap: i64) -> i64 {
58 let key: *u8 = "\"content\":\"" as *u8
59 let kl: i64 = bs_slen(key)
60 let pos: i64 = bs_find(body, blen, key, kl)
61 if pos<0 { return 0 }
62 var i: i64 = pos+kl
63 var o: i64 = 0
64 var go: i64 = 1
65 while go==1 {
66 if i>=blen { go=0 } else {
67 let c: i64 = body[i] as i64
68 if c==92 {
69 if i+1<blen { let d: i64=body[i+1] as i64
70 if d==110 { if o<cap { out[o]=10 as u8; o=o+1 } } else { if o<cap { out[o]=body[i+1]; o=o+1 } }
71 i=i+2
72 } else { i=i+1 }
73 } else { if c==34 { go=0 } else { if o<cap { out[o]=body[i]; o=o+1 } i=i+1 } }
74 }
75 }
76 return o
77}
78func bs_json_esc(src: *u8, sn: i64, dst: *u8, cap: i64) -> i64 {
79 var o: i64=0; var i: i64=0
80 while i<sn { if o+2>=cap { i=sn } else {
81 let c: i64=src[i] as i64
82 if c==34 { dst[o]=92 as u8; dst[o+1]=34 as u8; o=o+2 }
83 else { if c==92 { dst[o]=92 as u8; dst[o+1]=92 as u8; o=o+2 }
84 else { if c==10 { dst[o]=92 as u8; dst[o+1]=110 as u8; o=o+2 }
85 else { if c==13 { }
86 else { if c==9 { dst[o]=92 as u8; dst[o+1]=116 as u8; o=o+2 }
87 else { if c<32 { } else { dst[o]=src[i] as u8; o=o+1 } } } } } }
88 i=i+1
89 } }
90 return o
91}
92func bs_atoi(s: *u8) -> i64 {
93 var v: i64 = 0
94 var i: i64 = 0
95 while s[i] != (0 as u8) {
96 let d: i64 = s[i] as i64
97 if d < 48 { return v }
98 if d > 57 { return v }
99 v = v * 10 + (d - 48)
100 i = i + 1
101 }
102 return v
103}
104
105// respond one finished request: OpenAI-chat JSON, close, release.
106func bs_respond(S: *NxLlmSched, slot: nx_int, fd: i64,
107 esc: *u8, jbuf: *u8, resp: *u8) -> i64 {
108 let R: *NxLlmReq = nx_sched_req(S, slot)
109 let el: i64 = bs_json_esc(R.out, R.out_len, esc, 131072)
110 var jo: i64 = 0
111 jo=bs_cat(jbuf,jo,"{\"model\":\"nishi-qwen-0.5b-sovereign-batch\",\"choices\":[{\"message\":{\"role\":\"assistant\",\"content\":\"" as *u8)
112 var k: i64=0; while k<el { jbuf[jo]=esc[k]; jo=jo+1; k=k+1 }
113 jo=bs_cat(jbuf,jo,"\"}}]}" as *u8)
114 var ro: i64=0
115 ro=bs_cat(resp,ro,"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nConnection: close\r\nContent-Length: " as *u8)
116 let db: *u8=sys_mmap(24); var m: i64=jo; var dc: i64=0; if m==0 { db[0]=48 as u8; dc=1 } while m>0 { db[dc]=(48+(m%10)) as u8; m=m/10; dc=dc+1 }
117 var di: i64=dc-1; while di>=0 { resp[ro]=db[di]; ro=ro+1; di=di-1 }
118 ro=bs_cat(resp,ro,"\r\n\r\n" as *u8)
119 var jj: i64=0; while jj<jo { resp[ro]=jbuf[jj]; ro=ro+1; jj=jj+1 }
120 bs_write_all(fd, resp, ro)
121 sys_close(fd)
122 nx_sched_release(S, slot)
123 return 0
124}
125
126func main(argc: i64, argv: *i64) -> i64 {
127 var once_n: i64 = 0 - 1 // -1 = forever
128 if argc >= 3 { once_n = bs_atoi(argv[2] as *u8) }
129
130 bs_puts("nx_llm_batch_serve: loading /tmp/nx_real_model.gguf ...\n" as *u8)
131 let path: *u8 = "/tmp/nx_real_model.gguf" as *u8
132 let len_out: *i64 = sys_mmap(8) as *i64
133 let buf: *u8 = sys_read_file(path, len_out)
134 if buf==(0 as *u8) { bs_puts("batch-seat: no model\n" as *u8); sys_exit(10); return 10 }
135 let hdr: *NxGgufHeader = sys_mmap(NX_GGUF_HDR_BYTES) as *NxGgufHeader
136 if nx_gguf_parse(buf, len_out[0], hdr) != NX_GGUF_OK { sys_exit(20); return 20 }
137 let model: *NxF32LlamaModel = nx_f32_llama_model_alloc()
138 let out_err: *i64 = sys_mmap(8) as *i64
139 if nx_f32_llm_read_dims_from_gguf(buf, len_out[0], hdr, model, out_err) != NX_FLD_OK { sys_exit(30); return 30 }
140 if nx_f32_llm_load_weights_v4_from_gguf(buf, hdr, model, out_err) != NX_FLV4_OK { sys_exit(40); return 40 }
141 let vocab: *NxBpeVocab = nx_bpe_vocab_new(67108864, 262144, 524288)
142 let nt: *i64 = sys_mmap(8) as *i64
143 let nm: *i64 = sys_mmap(8) as *i64
144 if nx_f32_bpe_load_from_gguf(buf, len_out[0], hdr, vocab, nt, nm, out_err) != NX_FBL_OK { sys_exit(50); return 50 }
145 let eos: nx_int = nx_f32_llm_read_eos(buf, len_out[0], hdr)
146
147 let rc: *NxReasonCfg = nx_reason_cfg_alloc()
148 rc.model = model
149 rc.vocab = vocab
150 rc.cache = 0 as *NxF32KVCache
151 rc.max_new = 24
152 rc.inv_temp_f32 = 0x3FA00000
153 rc.top_k = 40
154 rc.eps = 0x358637BD
155 rc.attn_scale = 0x3E000000
156 rc.rope_log_base = 0x415D0EAB
157 rc.eos = eos
158 rc.im_start = 151644
159 rc.im_end = 151645
160
161 let kv_dim: nx_int = model.n_kv_heads * model.head_dim
162 let pool: *NxPagedPool = nx_pkv_pool_new(64, model.n_layers, kv_dim)
163 let S: *NxLlmSched = nx_sched_new(rc, pool)
164
165 let lfd: i64 = sys_socket(2, 1 | 0x80000, 0)
166 if lfd<0 { sys_exit(60); return 60 }
167 let optval: *i64 = sys_mmap(8) as *i64; optval[0]=1
168 sys_setsockopt(lfd, 1, 2, optval as *u8, 4)
169 // sockaddr COMPUTED from the integer port (the hand-encoded-bytes class
170 // put the :11434 seat on 11438 for its whole life -- never again).
171 let addr: *u8 = sys_mmap(16)
172 nx_http_client_sockaddr_ipv4(addr, 127, 0, 0, 1, NX_BSEAT_PORT)
173 if sys_bind(lfd, addr, 16) < 0 { bs_puts("batch-seat: bind fail (11435 in use?)\n" as *u8); sys_exit(70); return 70 }
174 sys_listen(lfd, 16)
175 sys_set_socket_timeout(lfd, 1) // accept-between-rounds
176 // SELF-PROBE before claiming LIVE (execution-level self-verification).
177 let pfd: i64 = sys_socket(2, 1, 0)
178 let paddr: *u8 = sys_mmap(16)
179 nx_http_client_sockaddr_ipv4(paddr, 127, 0, 0, 1, NX_BSEAT_PORT)
180 if nx_connect_bounded(pfd, paddr, 16, NX_CONN_DEFAULT_MS) < 0 {
181 bs_puts("batch-seat: SELF-PROBE FAILED -- not reachable on 127.0.0.1:11435; refusing to claim LIVE\n" as *u8)
182 sys_exit(71); return 71
183 }
184 sys_close(pfd)
185 let dfd: i64 = sys_accept(lfd) // drain the probe conn
186 if dfd >= 0 { sys_close(dfd) }
187 bs_puts("nx_llm_batch_serve: SELF-PROBE OK -- LIVE continuous-batching seat on 127.0.0.1:11435 (POST /chat)\n" as *u8)
188
189 let req: *u8 = sys_mmap(1048576)
190 let prompt: *u8 = sys_mmap(65536)
191 let toks: *i64 = sys_mmap(1024 * 8) as *i64
192 let esc: *u8 = sys_mmap(131072)
193 let resp: *u8 = sys_mmap(262144)
194 let jbuf: *u8 = sys_mmap(200000)
195 let fds: *i64 = sys_mmap(NX_SCHED_MAX_REQS * 8) as *i64
196 let outs: *u8 = sys_mmap(NX_SCHED_MAX_REQS * 4096)
197 var served: i64 = 0
198 var seed_ctr: i64 = 20260709
199
200 var alive: i64 = 1
201 while alive==1 {
202 // 1) admission: accept (times out ~1s when nothing arrives).
203 let fd: i64 = sys_accept(lfd)
204 if fd>=0 {
205 sys_set_socket_timeout(fd, 1)
206 // read until headers AND Content-Length body are complete (the
207 // request can arrive in multiple segments -- measured live:
208 // first read returned just "POST /chat HTTP/1.1\r\n", 21 bytes).
209 var rn: i64 = 0
210 var tries: i64 = 0
211 while tries < 8 {
212 let r1: i64 = sys_read(fd, ((req as i64)+rn) as *u8, 1048576-rn)
213 if r1 <= 0 { tries = 8 } else {
214 rn = rn + r1
215 let hb0: i64 = bs_find(req, rn, "\r\n\r\n" as *u8, 4)
216 if hb0 >= 0 {
217 var cl: i64 = 0
218 let cp: i64 = bs_find(req, rn, "Content-Length: " as *u8, 16)
219 if cp >= 0 { cl = bs_atoi(((req as i64)+cp+16) as *u8) }
220 if rn >= hb0 + 4 + cl { tries = 8 } else { tries = tries + 1 }
221 } else { tries = tries + 1 }
222 }
223 }
224 var handled: i64 = 0
225 if rn>0 {
226 if bs_find(req, rn, "POST /chat" as *u8, 10) >= 0 {
227 let hb: i64 = bs_find(req, rn, "\r\n\r\n" as *u8, 4)
228 var bstart: i64 = 0; if hb>=0 { bstart=hb+4 }
229 let pl: i64 = bs_extract_prompt(((req as i64)+bstart) as *u8, rn-bstart, prompt, 65534)
230 if pl > 0 {
231 let ntk: nx_int = nx_reason_build_chat_toks(rc, prompt, pl as nx_int, toks)
232 seed_ctr = seed_ctr + 7919
233 let ob: *u8 = ((outs as i64) + 0) as *u8 // slot buffer chosen below
234 let slot: nx_int = nx_sched_admit(S, toks, ntk, seed_ctr, rc.max_new,
235 ((outs as i64) + (S.n_reqs as i64) * 4096) as *u8, 4096)
236 if slot >= 0 { fds[slot] = fd; handled = 1 }
237 }
238 }
239 }
240 if handled == 0 {
241 // diagnostic: why did this request miss the /chat path?
242 bs_puts("fallback rn=" as *u8)
243 let dn2: *u8 = sys_mmap(24)
244 var mm: i64 = rn
245 if mm < 0 { bs_puts("-" as *u8); mm = 0 - mm }
246 var dk: i64 = 0
247 if mm == 0 { dn2[0]=48 as u8; dk=1 }
248 while mm > 0 { dn2[dk]=(48+(mm%10)) as u8; mm=mm/10; dk=dk+1 }
249 var dj: i64 = dk - 1
250 let ob2: *u8 = sys_mmap(24)
251 var oo: i64 = 0
252 while dj >= 0 { ob2[oo]=dn2[dj]; oo=oo+1; dj=dj-1 }
253 sys_write(1, ob2, oo)
254 bs_puts(" head=[" as *u8)
255 var hh: i64 = 0
256 let hs: *u8 = sys_mmap(16)
257 while hh < 12 { var cc: i64 = req[hh] as i64; if cc < 32 { cc = 46 } hs[hh]=cc as u8; hh=hh+1 }
258 sys_write(1, hs, 12)
259 bs_puts("]\n" as *u8)
260 let hm: *u8 = "HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nConnection: close\r\nContent-Length: 24\r\n\r\nnishi-llm-batch-seat ok\n" as *u8
261 bs_write_all(fd, hm, bs_slen(hm))
262 sys_close(fd)
263 }
264 }
265
266 // 2) one decode round for every active request.
267 nx_sched_round(S)
268
269 // 3) respond + release finished requests.
270 var i: nx_int = 0
271 while i < S.n_reqs {
272 let R: *NxLlmReq = nx_sched_req(S, i)
273 if R.in_use == 1 {
274 var fin: i64 = 0
275 if R.done == 1 { fin = 1 }
276 if R.emitted >= (R.max_new as i64) { fin = 1 }
277 if fin == 1 {
278 bs_respond(S, i, fds[i], esc, jbuf, resp)
279 served = served + 1
280 }
281 }
282 i = i + 1
283 }
284 if once_n >= 0 { if served >= once_n { alive = 0 } }
285 }
286 bs_puts("nx_llm_batch_serve: once-mode complete\n" as *u8)
287 sys_exit(0); return 0
288}