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}