code wiki / (root) / nx_llm_batch_serve.nx

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}