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}