code wiki / (root) / nx_swarm_tend.nx

nx_swarm_tend.nx source

↩ module page · 139 lines · 5627 B

1// nx_swarm_tend.nx -- ACTIVE QUEUE TENDER (operator 2026-07-16: "whoever in the nishi RACI is part 2// of the queue -- use it WELL, not fire-and-forget letting things block and get traffic jams; SOTA"). 3// THE GAP: nx_swarm_queue reaps dead holders PASSIVELY (only when a waiter polls sq_try_grant). If a 4// holder dies with NO waiters polling, its slot LEAKS forever; starvation + jams are invisible. My own 5// 1.5B job queued 45min and timed out unnoticed -- the fire-and-forget failure. THE FIX: an active 6// tender the CONDUCTOR (RACI schedule_resources|R = Responsible-for-draining) runs on a loop / on 7// demand: 8// * REAP: proactively DONE any HOLD whose pid is DEAD -> slots free even with no waiter polling. 9// * DETECT: a waiter starved past a jam threshold = JAMMED (the traffic-jam alarm), the oldest 10// wait-age reported so the PM sees pressure building BEFORE a timeout. 11// * VERDICT: OK / DRAINING (waiters progressing) / JAMMED (starvation) -- health, not silence. 12// tend reap + scan + report + verdict-exit (0 OK / 1 DRAINING / 2 JAMMED) 13// Sovereign store (knowledge/store/swarm_sched) is the SSOT; PURE tend_verdict is gate-testable. 14// license_tier: ORIGINAL 15import "nx_swarm_queue.nx" 16 17const TEND_JAM_MS: i64 = 300000 // a waiter silent > 5min = JAMMED (data-driven default) 18const TEND_OK: i64 = 0 19const TEND_DRAINING: i64 = 1 20const TEND_JAMMED: i64 = 2 21 22// PURE health verdict: no waiters -> OK; oldest waiter past jam -> JAMMED; else DRAINING. gate-core. 23func tend_verdict(waiters: i64, oldest_wait_ms: i64, jam_ms: i64) -> i64 { 24 if waiters <= 0 { return TEND_OK } 25 if oldest_wait_ms > jam_ms { return TEND_JAMMED } 26 return TEND_DRAINING 27} 28 29// proactively reap dead HOLDers: for each q:<pid> that is HOLD and whose pid is dead, write DONE so 30// the slot frees WITHOUT waiting for a waiter to poll. returns count reaped. 31func tend_reap() -> i64 { 32 let ids: *u8 = sys_mmap(262144) as *u8 33 let idn: i64 = sov_get_copy(SQ_STORE, "q:ids" as *u8, ids, 262143) 34 if idn <= 0 { return 0 } 35 let val: *u8 = sys_mmap(256) as *u8 36 let key: *u8 = sys_mmap(32) as *u8 37 let row: *u8 = sys_mmap(256) as *u8 38 var reaped: i64 = 0 39 var i: i64 = 0 40 while i < idn { 41 var e: i64 = i 42 var sc: i64 = 1 43 while sc == 1 { 44 if e >= idn { sc = 0 } 45 if sc == 1 { if ids[e] == (10 as u8) { sc = 0 } else { e = e + 1 } } 46 } 47 let kl: i64 = e - i 48 if kl > 0 && kl < 28 { 49 key[0] = 113 as u8 50 key[1] = 58 as u8 51 var m: i64 = 0 52 while m < kl { key[2 + m] = ids[i + m]; m = m + 1 } 53 key[2 + kl] = 0 as u8 54 let vn: i64 = sov_get_copy(SQ_STORE, key, val, 255) 55 if vn > 0 { 56 val[vn] = 0 as u8 57 // parse pid (field 1) + type (first char). HOLD = 'H' (72). 58 let pid: i64 = sq_ifield(val, 0, vn, 1) 59 if val[0] == (72 as u8) { 60 if pid > 0 { 61 if ml_pid_alive(pid) == 0 { 62 sq_mkrow("DONE" as *u8, pid, 0, 0, sys_now_us(), row) 63 sq_put_state(pid, row) 64 reaped = reaped + 1 65 } 66 } 67 } 68 } 69 } 70 i = e + 1 71 } 72 return reaped 73} 74 75// scan current health: fills holders/waiters + oldest-waiter age(ms). returns waiters. 76func tend_health(holders_out: *i64, oldest_ms_out: *i64) -> i64 { 77 let buf: *u8 = sys_mmap(262144) as *u8 78 let n: i64 = sq_materialize(buf, 262143) 79 let wp: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 80 let wpr: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 81 let we: *i64 = sys_mmap(SQ_MAXROW * 8) as *i64 82 let nwp: *i64 = sys_mmap(8) as *i64 83 var holders: i64 = 0 84 if n > 0 { holders = sq_scan(buf, n, wp, wpr, we, nwp) } 85 let nw: i64 = nwp[0] 86 holders_out[0] = holders 87 var oldest_enq: i64 = 0 88 var have: i64 = 0 89 var j: i64 = 0 90 while j < nw { 91 if have == 0 { oldest_enq = we[j]; have = 1 } 92 if have == 1 { if we[j] < oldest_enq { oldest_enq = we[j] } } 93 j = j + 1 94 } 95 var age_ms: i64 = 0 96 if have == 1 { age_ms = (sys_now_us() - oldest_enq) / 1000 } 97 if age_ms < 0 { age_ms = 0 } 98 oldest_ms_out[0] = age_ms 99 return nw 100} 101 102func tend_run() -> i64 { 103 let reaped: i64 = tend_reap() 104 let holders: *i64 = sys_mmap(8) as *i64 105 let oldest: *i64 = sys_mmap(8) as *i64 106 let waiters: i64 = tend_health(holders, oldest) 107 let v: i64 = tend_verdict(waiters, oldest[0], TEND_JAM_MS) 108 std_puts("SWARM-TEND reaped=" as *u8) 109 std_pdec(reaped) 110 std_puts(" holders=" as *u8) 111 std_pdec(holders[0]) 112 std_puts(" waiters=" as *u8) 113 std_pdec(waiters) 114 std_puts(" oldest_wait_ms=" as *u8) 115 std_pdec(oldest[0]) 116 std_puts(" -> " as *u8) 117 if v == TEND_OK { std_putln("OK (queue clear)" as *u8) } 118 if v == TEND_DRAINING { std_putln("DRAINING (waiters progressing)" as *u8) } 119 if v == TEND_JAMMED { std_putln("JAMMED (a waiter starved past threshold -- PM: add capacity or preempt)" as *u8) } 120 return v 121} 122 123func main(argc: i64, argv: *i64) -> i64 { 124 if argc < 2 { 125 std_putln("usage: nx_swarm_tend tend (conductor drains: reap dead holders + report jams)" as *u8) 126 sys_exit(2) 127 return 2 128 } 129 let a1: i64 = argv[1] 130 let verb: *u8 = a1 as *u8 131 if std_streq(verb, "tend" as *u8) == 1 { 132 let v: i64 = tend_run() 133 sys_exit(v) 134 return v 135 } 136 std_putln("unknown verb (only: tend)" as *u8) 137 sys_exit(2) 138 return 2 139}