code wiki / (root) / nx_swarm_queue.nx

nx_swarm_queue.nx source

↩ module page · 449 lines · 17507 B

1// nx_swarm_queue.nx -- INTELLIGENT QUEUING + COORDINATION for the swarm admission plane (operator 2// 2026-07-16: "dont just hard deny it should be intelligent queing and coordination"). Upgrades 3// nx_swarm_admit's terminal DENY into a K3s/Ray-class PENDING->SCHEDULED queue: a heavy launch that 4// can't run NOW ENQUEUES and WAITS; the scheduler admits it the moment resources free, ordered by 5// PRIORITY then FIFO (no monopoly, no starvation), telemetry-gated (never oversubscribes). All 6// sessions coordinate through ONE SOVEREIGN STORE (seg_store q:<pid> latest-wins; sys_flock 7// serializes the grant transition + q:ids index) so cross-session heavy launches SEQUENCE instead 8// of colliding. Per-pid state (knowledge/store/swarm_sched, key q:<pid>): W|pid|prio|est|us · 9// H|pid|prio|est|us · D|pid|0|0|us -- latest-wins gives latest-state-per-pid natively (no log replay). 10// PURE gate-core: sq_rank (waiters ahead of me) + sq_should_admit (rank+holders<cap AND telem_ok). 11// wait-admit <pid> <prio> <est_mb> <max_conc> <max_wait_ms> · release <pid> · status 12// license_tier: ORIGINAL 13import "nx_model_lane_core.nx" 14import "nx_sov_ledger.nx" 15 16// SOVEREIGN store (no-TSV law, 2026-07-16): each pid's latest state = one key q:<pid> (seg_store 17// latest-wins = latest-state-per-pid natively, no log replay). q:ids = the pid index (maintained 18// under SQ_LOCK flock so a concurrent enqueue never loses an id). The gate-proven PURE parsers 19// (sq_scan/sq_rank/sq_heavy_holders/sq_admit_v2) are UNCHANGED -- sq_materialize rebuilds the exact 20// text buffer they expect from the store, so only the storage substrate moved flat->sovereign. 21const SQ_STORE: *u8 = "knowledge/store/swarm_sched" 22const SQ_LOCK: *u8 = "knowledge/status/swarm_sched.lock" 23const SQ_FLOOR_MB: i64 = 2048 24const SQ_MAXROW: i64 = 256 25const SQ_BACKFILL_MAX: i64 = 6000 // jobs <=6GB backfill idle capacity while a bigger job waits (data-driven default) 26 27// own pid WITHOUT the getpid syscall: parse the leading field of /proc/self/stat (wrapper-based, 28// pid-namespace-correct). getpid has NO working CONST __syscall form today: const numbers go through 29// the RV64->x86 translator where 39 COLLIDES (mistranslated -> -ENOTTY) and 172 is ABSENT from the 30// table (passes raw -> x86 iopl -> -ENOSYS). The runtime path (nr in a register) is NOT translated, 31// which is why nx_pid_diag's param-based getpid(39) "worked" -- opposite conventions per path. 32// Witnessed: nx_getpid_const_probe (const_fn=-25 var_fn=<pid> same binary) + nx_rv64_const_probe. 33func sq_getpid() -> i64 { 34 let fd: i64 = sys_openat_rd("/proc/self/stat" as *u8) 35 if fd < 0 { return 0 - 1 } 36 let b: *u8 = sys_mmap(128) as *u8 37 let n: i64 = sys_read(fd, b, 127) 38 sys_close(fd) 39 if n <= 0 { return 0 - 1 } 40 var v: i64 = 0 41 var i: i64 = 0 42 var seen: i64 = 0 43 while i < n { 44 let c: i64 = b[i] as i64 45 if c >= 48 && c <= 57 { v = v * 10 + (c - 48); seen = 1; i = i + 1 } else { i = n } 46 } 47 if seen == 0 { return 0 - 1 } 48 return v 49} 50 51// build "q:<pid>" into out (NUL-term). 52func sq_pidkey(pid: i64, out: *u8) -> i64 { 53 out[0] = 113 as u8 54 out[1] = 58 as u8 55 let l: i64 = std_itoa(pid, out + 2) 56 out[2 + l] = 0 as u8 57 return 2 + l 58} 59 60// write a pid's latest state line ("T|pid|prio|est|enq") -> q:<pid> (own key; no lock; safe to call 61// from WITHIN the grant flock -- avoids the self-deadlock of re-locking SQ_LOCK we already hold). 62func sq_put_state(pid: i64, line: *u8) -> i64 { 63 let key: *u8 = sys_mmap(32) as *u8 64 sq_pidkey(pid, key) 65 sov_put_str(SQ_STORE, key, line) 66 return 0 67} 68 69// ensure pid in q:ids under the shared flock (no lost-update on concurrent enqueue). MUST NOT be 70// called while already holding SQ_LOCK (enqueue calls it pre-grant-loop; try_grant/release do not). 71func sq_add_id(pid: i64) -> i64 { 72 let lf: i64 = sys_openat_append(SQ_LOCK, 0x1a4) 73 if lf >= 0 { sys_flock(lf, SYS_LOCK_EX) } 74 let ids: *u8 = sys_mmap(262144) as *u8 75 let idn: i64 = sov_get_copy(SQ_STORE, "q:ids" as *u8, ids, 262143) 76 var have: i64 = 0 77 let pk: *u8 = sys_mmap(32) as *u8 78 var pkl: i64 = std_itoa(pid, pk) 79 if idn > 0 { 80 var i: i64 = 0 81 while i < idn { 82 var e: i64 = i 83 var sc: i64 = 1 84 while sc == 1 { 85 if e >= idn { sc = 0 } 86 if sc == 1 { if ids[e] == (10 as u8) { sc = 0 } else { e = e + 1 } } 87 } 88 if e - i == pkl { 89 var m: i64 = 0 90 var eq: i64 = 1 91 while m < pkl { if ids[i + m] != pk[m] { eq = 0; m = pkl } else { m = m + 1 } } 92 if eq == 1 { have = 1 } 93 } 94 i = e + 1 95 } 96 } 97 if have == 0 { 98 var o: i64 = idn 99 if o < 0 { o = 0 } 100 var m: i64 = 0 101 while m < pkl { ids[o] = pk[m]; o = o + 1; m = m + 1 } 102 ids[o] = 10 as u8 103 o = o + 1 104 sov_put(SQ_STORE, "q:ids" as *u8, ids, o) 105 } 106 if lf >= 0 { sys_flock(lf, SYS_LOCK_UN); sys_close(lf) } 107 return 0 108} 109 110// rebuild the text buffer the PURE parsers expect: one "<line>\n" per pid in q:ids (from the store). 111func sq_materialize(buf: *u8, cap: i64) -> i64 { 112 let ids: *u8 = sys_mmap(262144) as *u8 113 let idn: i64 = sov_get_copy(SQ_STORE, "q:ids" as *u8, ids, 262143) 114 if idn <= 0 { buf[0] = 0 as u8; return 0 } 115 let key: *u8 = sys_mmap(32) as *u8 116 var o: i64 = 0 117 var i: i64 = 0 118 while i < idn { 119 var e: i64 = i 120 var sc: i64 = 1 121 while sc == 1 { 122 if e >= idn { sc = 0 } 123 if sc == 1 { if ids[e] == (10 as u8) { sc = 0 } else { e = e + 1 } } 124 } 125 let kl: i64 = e - i 126 if kl > 0 && kl < 28 { 127 key[0] = 113 as u8 128 key[1] = 58 as u8 129 var m: i64 = 0 130 while m < kl { key[2 + m] = ids[i + m]; m = m + 1 } 131 key[2 + kl] = 0 as u8 132 let q: *u8 = buf + o 133 let vn: i64 = sov_get_copy(SQ_STORE, key, q, cap - 1 - o) 134 if vn > 0 { o = o + vn; buf[o] = 10 as u8; o = o + 1 } 135 } 136 i = e + 1 137 } 138 buf[o] = 0 as u8 139 return o 140} 141 142// integer value of |-field k (0-based) in line [ls,le); TYPE field 0 is non-numeric -> use sq_type. 143func sq_ifield(buf: *u8, ls: i64, le: i64, k: i64) -> i64 { 144 var field: i64 = 0 145 var v: i64 = 0 146 var seen: i64 = 0 147 var i: i64 = ls 148 while i < le { 149 let c: i64 = buf[i] as i64 150 if c == 124 { field = field + 1; if field > k { i = le } } 151 if c != 124 { 152 if field == k { 153 if c >= 48 && c <= 57 { v = v * 10 + (c - 48); seen = 1 } 154 } 155 } 156 if i < le { i = i + 1 } 157 } 158 if seen == 0 { return 0 - 1 } 159 return v 160} 161 162// 0=WAIT 1=HOLD 2=DONE -1=other (unique first chars) 163func sq_type(buf: *u8, ls: i64) -> i64 { 164 let c: i64 = buf[ls] as i64 165 if c == 87 { return 0 } 166 if c == 72 { return 1 } 167 if c == 68 { return 2 } 168 return 0 - 1 169} 170 171// replay the log -> latest state per pid. Fills live-waiter arrays (wp/wpr/we) + returns holder count. 172// tables sized SQ_MAXROW. A pid alive & last-state WAIT = live waiter; last-state HOLD = holder. 173func sq_scan(buf: *u8, n: i64, wp: *i64, wpr: *i64, we: *i64, out_nw: *i64) -> i64 { 174 let tpid: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 175 let ttype: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 176 let tprio: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 177 let tenq: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 178 var tn: i64 = 0 179 var i: i64 = 0 180 while i < n { 181 var e: i64 = i 182 var sc: i64 = 1 183 while sc == 1 { 184 if e >= n { sc = 0 } 185 if sc == 1 { if buf[e] == (10 as u8) { sc = 0 } else { e = e + 1 } } 186 } 187 let ty: i64 = sq_type(buf, i) 188 if ty >= 0 { 189 let pid: i64 = sq_ifield(buf, i, e, 1) 190 if pid > 0 { 191 // find/insert pid in table 192 var idx: i64 = 0 - 1 193 var j: i64 = 0 194 while j < tn { if tpid[j] == pid { idx = j; j = tn } else { j = j + 1 } } 195 if idx < 0 { if tn < SQ_MAXROW { idx = tn; tpid[tn] = pid; tprio[tn] = 0; tenq[tn] = 0; tn = tn + 1 } } 196 if idx >= 0 { 197 ttype[idx] = ty 198 if ty == 0 { tprio[idx] = sq_ifield(buf, i, e, 2); tenq[idx] = sq_ifield(buf, i, e, 4) } 199 } 200 } 201 } 202 i = e + 1 203 } 204 // materialize live waiters + count holders 205 var holders: i64 = 0 206 var nw: i64 = 0 207 var k: i64 = 0 208 while k < tn { 209 let alive: i64 = ml_pid_alive(tpid[k]) 210 if alive == 1 { 211 if ttype[k] == 1 { holders = holders + 1 } 212 if ttype[k] == 0 { 213 if nw < SQ_MAXROW { wp[nw] = tpid[k]; wpr[nw] = tprio[k]; we[nw] = tenq[k]; nw = nw + 1 } 214 } 215 } 216 k = k + 1 217 } 218 out_nw[0] = nw 219 return holders 220} 221 222// count of live waiters strictly AHEAD of me (higher prio, or equal prio + earlier enqueue). PURE. 223func sq_rank(wp: *i64, wpr: *i64, we: *i64, nw: i64, my_pid: i64, my_prio: i64, my_enq: i64) -> i64 { 224 var ahead: i64 = 0 225 var j: i64 = 0 226 while j < nw { 227 if wp[j] != my_pid { 228 var isahead: i64 = 0 229 if wpr[j] > my_prio { isahead = 1 } 230 if wpr[j] == my_prio { if we[j] < my_enq { isahead = 1 } } 231 if isahead == 1 { ahead = ahead + 1 } 232 } 233 j = j + 1 234 } 235 return ahead 236} 237 238// PURE admission decision: within capacity for my rank AND telemetry allows. gate-core. 239func sq_should_admit(rank: i64, holders: i64, max_conc: i64, telem_ok: i64) -> i64 { 240 if telem_ok == 0 { return 0 } 241 if rank + holders < max_conc { return 1 } 242 return 0 243} 244 245// count LIVE holders whose est > threshold ("heavy" jobs; small backfill holders don't count 246// against the heavy-job capacity). Standalone scan (no churn on sq_scan's signature). 247func sq_heavy_holders(buf: *u8, n: i64, threshold: i64) -> i64 { 248 let tpid: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 249 let ttype: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 250 let test: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 251 var tn: i64 = 0 252 var i: i64 = 0 253 while i < n { 254 var e: i64 = i 255 var sc: i64 = 1 256 while sc == 1 { 257 if e >= n { sc = 0 } 258 if sc == 1 { if buf[e] == (10 as u8) { sc = 0 } else { e = e + 1 } } 259 } 260 let ty: i64 = sq_type(buf, i) 261 if ty >= 0 { 262 let pid: i64 = sq_ifield(buf, i, e, 1) 263 if pid > 0 { 264 var idx: i64 = 0 - 1 265 var j: i64 = 0 266 while j < tn { if tpid[j] == pid { idx = j; j = tn } else { j = j + 1 } } 267 if idx < 0 { if tn < SQ_MAXROW { idx = tn; tpid[tn] = pid; test[tn] = 0; tn = tn + 1 } } 268 if idx >= 0 { 269 ttype[idx] = ty 270 if ty == 1 { test[idx] = sq_ifield(buf, i, e, 3) } 271 } 272 } 273 } 274 i = e + 1 275 } 276 var heavy: i64 = 0 277 var k: i64 = 0 278 while k < tn { 279 if ttype[k] == 1 { 280 if ml_pid_alive(tpid[k]) == 1 { 281 if test[k] > threshold { heavy = heavy + 1 } 282 } 283 } 284 k = k + 1 285 } 286 return heavy 287} 288 289// PURE admission v2 with BACKFILL (operator 2026-07-16: "use the waittime intelligently... active 290// growth, don't just sit in line"). Never oversubscribe (telem gate). Front job admitted in order 291// (rank + HEAVY holders < cap). Else a SMALL job (est <= backfill_max) BACKFILLS the idle capacity 292// a big waiter can't use yet -- bounded footprint so it can't grow into or starve the big job. 293// 0 QUEUE · 1 GRANT-front (ordered) · 2 GRANT-backfill (small fills the gap). 294func sq_admit_v2(rank: i64, heavy_holders: i64, max_conc: i64, my_est: i64, backfill_max: i64, telem_ok: i64) -> i64 { 295 if telem_ok == 0 { return 0 } 296 if rank + heavy_holders < max_conc { return 1 } 297 if my_est <= backfill_max { return 2 } 298 return 0 299} 300 301// live telemetry gate (fail-closed): load below cores AND RAM headroom for est. 302func sq_telem_ok(est_mb: i64) -> i64 { 303 let l1: i64 = ml_load1() 304 let nc: i64 = ml_ncores() 305 let av: i64 = ml_meminfo_avail_mb() 306 if av < 0 { return 0 } 307 if nc > 0 { if l1 >= nc { return 0 } } 308 if av < est_mb + SQ_FLOOR_MB { return 0 } 309 return 1 310} 311 312func sq_mkrow(ty: *u8, pid: i64, prio: i64, est: i64, us: i64, out: *u8) -> i64 { 313 var o: i64 = 0 314 var i: i64 = 0 315 while ty[i] != (0 as u8) { out[o] = ty[i]; o = o + 1; i = i + 1 } 316 out[o] = 124 as u8; o = o + 1 317 o = o + std_itoa(pid, out + o) 318 out[o] = 124 as u8; o = o + 1 319 o = o + std_itoa(prio, out + o) 320 out[o] = 124 as u8; o = o + 1 321 o = o + std_itoa(est, out + o) 322 out[o] = 124 as u8; o = o + 1 323 o = o + std_itoa(us, out + o) 324 out[o] = 0 as u8 325 return o 326} 327 328func sq_enqueue(pid: i64, prio: i64, est: i64) -> i64 { 329 let now: i64 = sys_now_us() 330 let row: *u8 = sys_mmap(256) as *u8 331 sq_mkrow("WAIT" as *u8, pid, prio, est, now, row) 332 sq_put_state(pid, row) 333 sq_add_id(pid) 334 return now 335} 336 337func sq_release(pid: i64) -> i64 { 338 let now: i64 = sys_now_us() 339 let row: *u8 = sys_mmap(256) as *u8 340 sq_mkrow("DONE" as *u8, pid, 0, 0, now, row) 341 sq_put_state(pid, row) 342 return 0 343} 344 345// the grant transition, serialized under flock: read store -> decide -> on admit write HOLD. 0 GRANT / 1 QUEUE. 346func sq_try_grant(pid: i64, prio: i64, est: i64, max_conc: i64, my_enq: i64) -> i64 { 347 let lf: i64 = sys_openat_append(SQ_LOCK, 0x1a4) 348 if lf >= 0 { sys_flock(lf, SYS_LOCK_EX) } 349 let buf: *u8 = sys_mmap(262144) as *u8 350 let n: i64 = sq_materialize(buf, 262143) 351 let wp: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 352 let wpr: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 353 let we: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 354 let nwp: *i64 = sys_mmap(8) as *i64 355 var holders: i64 = 0 356 if n > 0 { holders = sq_scan(buf, n, wp, wpr, we, nwp) } 357 let nw: i64 = nwp[0] 358 let rank: i64 = sq_rank(wp, wpr, we, nw, pid, prio, my_enq) 359 let tok: i64 = sq_telem_ok(est) 360 var heavy: i64 = 0 361 if n > 0 { heavy = sq_heavy_holders(buf, n, SQ_BACKFILL_MAX) } 362 let dec: i64 = sq_admit_v2(rank, heavy, max_conc, est, SQ_BACKFILL_MAX, tok) 363 var verdict: i64 = 1 364 if dec >= 1 { 365 let row: *u8 = sys_mmap(256) as *u8 366 let now: i64 = sys_now_us() 367 sq_mkrow("HOLD" as *u8, pid, prio, est, now, row) 368 sq_put_state(pid, row) 369 verdict = 0 370 if dec == 2 { std_putln("SWARM-QUEUE (backfill: small job fills idle capacity while a bigger job waits)" as *u8) } 371 } 372 if lf >= 0 { sys_flock(lf, SYS_LOCK_UN); sys_close(lf) } 373 return verdict 374} 375 376func sq_wait_admit(pid: i64, prio: i64, est: i64, max_conc: i64, max_wait_ms: i64, poll_ms: i64) -> i64 { 377 let my_enq: i64 = sq_enqueue(pid, prio, est) 378 var waited: i64 = 0 379 var go: i64 = 1 380 while go == 1 { 381 let v: i64 = sq_try_grant(pid, prio, est, max_conc, my_enq) 382 if v == 0 { return 0 } 383 if waited >= max_wait_ms { go = 0 } 384 if go == 1 { sys_sleep_ms(poll_ms); waited = waited + poll_ms } 385 } 386 return 1 387} 388 389func main(argc: i64, argv: *i64) -> i64 { 390 if argc < 2 { 391 std_putln("usage: nx_swarm_queue wait-admit <pid> <prio> <est_mb> <max_conc> <max_wait_ms> | release <pid> | status" as *u8) 392 sys_exit(2) 393 return 2 394 } 395 let a1: i64 = argv[1] 396 let verb: *u8 = a1 as *u8 397 if std_streq(verb, "release" as *u8) == 1 { 398 if argc < 3 { std_putln("release needs <pid>" as *u8); sys_exit(2); return 2 } 399 let b2: i64 = argv[2] 400 let pb: *u8 = b2 as *u8 401 sq_release(std_atoi(pb)) 402 std_putln("SWARM-QUEUE released" as *u8) 403 sys_exit(0) 404 return 0 405 } 406 if std_streq(verb, "status" as *u8) == 1 { 407 let buf: *u8 = sys_mmap(262144) as *u8 408 let n: i64 = sq_materialize(buf, 262143) 409 let wp: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 410 let wpr: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 411 let we: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 412 let nwp: *i64 = sys_mmap(8) as *i64 413 var holders: i64 = 0 414 if n > 0 { holders = sq_scan(buf, n, wp, wpr, we, nwp) } 415 std_puts("SWARM-QUEUE holders=" as *u8) 416 std_pdec(holders) 417 std_puts(" waiters=" as *u8) 418 std_pdec(nwp[0]) 419 std_puts("\n" as *u8) 420 sys_exit(0) 421 return 0 422 } 423 if std_streq(verb, "wait-admit" as *u8) == 1 { 424 if argc < 7 { std_putln("wait-admit needs <pid> <prio> <est_mb> <max_conc> <max_wait_ms>" as *u8); sys_exit(2); return 2 } 425 let c2: i64 = argv[2] 426 let p2: *u8 = c2 as *u8 427 let pid: i64 = std_atoi(p2) 428 let c3: i64 = argv[3] 429 let p3: *u8 = c3 as *u8 430 let prio: i64 = std_atoi(p3) 431 let c4: i64 = argv[4] 432 let p4: *u8 = c4 as *u8 433 let est: i64 = std_atoi(p4) 434 let c5: i64 = argv[5] 435 let p5: *u8 = c5 as *u8 436 let maxc: i64 = std_atoi(p5) 437 let c6: i64 = argv[6] 438 let p6: *u8 = c6 as *u8 439 let maxw: i64 = std_atoi(p6) 440 let r: i64 = sq_wait_admit(pid, prio, est, maxc, maxw, 2000) 441 if r == 0 { std_putln("SWARM-QUEUE GRANT (admitted; run then: nx_swarm_queue release <pid>)" as *u8); sys_exit(0); return 0 } 442 std_putln("SWARM-QUEUE TIMEOUT (still queued; box stayed saturated past max_wait)" as *u8) 443 sys_exit(1) 444 return 1 445 } 446 std_putln("unknown verb" as *u8) 447 sys_exit(2) 448 return 2 449}