nx_ship_fleet.nx source
↩ module page · 234 lines · 12244 B
1// nx_ship_fleet.nx -- INTELLIGENT (adaptive-concurrency) FLEET SHIPPER. Drains a list of targets by running
2// as many concurrent ships as the box can take RIGHT NOW -- never serially (too slow to drain a backlog),
3// never blindly (3 blind concurrent ships stormed the box on 2026-09-02). It composes the estate's existing
4// pieces into the DRIVER that was missing:
5// SENSE nx_sysload (sl_ncpu / sl_loadavg_milli / sl_freemem_mb / sl_active_conns_8443)
6// GOVERN nx_resource_governor rg_worker_budget -> the polite headroom budget, used here as a concurrency
7// WIDTH (the governed drivers computed this budget and then fired ONE task per beat -- the budget
8// was never used as a width, which is exactly why concurrent shipping was ungoverned)
9// ADMIT fork _offc/nx_build_admit.elf -> its exit code is the DIRECT I/O-storm signal (D-state witness),
10// which the CPU-load budget alone cannot see (our storms are I/O on a degraded RAID, not CPU)
11// AIMD nx_shipfleet_lib -- +1 on progress, halve on a storm refusal, ceilinged by the polite budget
12// POOL fork _offc/nx_organ_ship.elf per target; reap with wait4(WNOHANG); requeue nothing blindly
13//
14// Research (Sept 2026): AIMD adaptive concurrency limits (Netflix concurrency-limits, Envoy
15// adaptive_concurrency filter, ThomWright/congestion-limiter) + HPA-style headroom budgeting. The field
16// infers congestion from latency; our exceed is a measured congestion signal (build_admit reads /proc).
17//
18// usage: nx_ship_fleet <targets-file> [plan]
19// <targets-file> one target name per line; '#' and blank lines skipped
20// plan SENSE + print the width the controller WOULD use, fork NOTHING (safe live demo)
21// license_tier: ORIGINAL
22
23import "nx_syscalls.nx"
24import "nx_sysload.nx"
25import "nx_resource_governor.nx"
26import "nx_shipfleet_lib.nx"
27
28const SF_MEM_FLOOR_MB: i64 = 256
29const SF_CEIL_MILLI: i64 = 800 // back off at 0.8 core/cpu (rg default; shipping is background, not latency-sacred)
30const SF_BASE_BACKOFF_MS:i64 = 2000 // AIMD storm back-off base (rg_backoff_ms scales it up with overload)
31const SF_TICK_MS: i64 = 250 // settle time between launch ticks (avoid busy-spin on /proc)
32const SF_MAX_STORM_TICKS:i64 = 120 // ~ bounded: consecutive storm/busy ticks before giving up (killable, deferred reported)
33const SF_MAXTARGETS: i64 = 4096
34const SF_NAMECAP: i64 = 128
35const SF_INFLIGHT_CAP: i64 = 256
36const SF_READCAP: i64 = 1048576
37const SF_SHIP_ELF: *u8 = "_offc/nx_organ_ship.elf" as *u8
38const SF_ADMIT_ELF: *u8 = "_offc/nx_build_admit.elf" as *u8
39const SF_DEVNULL: *u8 = "/dev/null\x00" as *u8
40
41func sf_slen(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} return n }
42func sf_w(s: *u8) -> i64 { sys_write(1, s, sf_slen(s)); return 0 }
43func sf_n(v: i64) -> i64 {
44 if v==0 { sys_write(1,"0" as *u8,1); return 0 }
45 var m: i64=v; if m<0 { sys_write(1,"-" as *u8,1); m=0-m }
46 let t: *u8=sys_mmap(24); var k: i64=0
47 while m>0 { t[k]=(48+(m%10)) as u8; m=m/10; k=k+1 }
48 let o: *u8=sys_mmap(24); var w: i64=0; var q: i64=k-1
49 while q>=0 { o[w]=t[q]; w=w+1; q=q-1 }
50 sys_write(1,o,w); return 0
51}
52
53// read the whole targets file into buf (NUL-terminated); returns byte count, 0 on open fail.
54func sf_read(path: *u8, buf: *u8, cap: i64) -> i64 {
55 let fd: i64 = sys_openat_rd(path); if fd<0 { return 0 }
56 var n: i64=0; var r: i64=sys_read(fd, buf, cap-1)
57 while r>0 { n=n+r; if n>=cap-1 { r=0 } else { r=sys_read(fd, buf+n, cap-1-n) } }
58 sys_close(fd); buf[n]=0 as u8; return n
59}
60
61// the DIRECT congestion signal: fork build_admit 'check', exit 0 = GRANT (headroom), else REFUSE (storm/busy).
62// cannot-fork or a missing admit binary is treated as NO headroom (fail-safe: not knowing is not permission).
63func sf_admit_grant() -> i64 {
64 let pid: i64 = sys_fork()
65 if pid < 0 { return 0 }
66 if pid == 0 {
67 let dn: i64 = sys_openat_wr(SF_DEVNULL, 420)
68 if dn>=0 { sys_dup3(dn,1,0); sys_dup3(dn,2,0) }
69 let argv: *i64 = sys_mmap(64) as *i64
70 argv[0]=SF_ADMIT_ELF as i64; argv[1]="check" as *u8 as i64; argv[2]=0
71 let envp: *i64 = sys_mmap(16) as *i64; envp[0]="PATH=/usr/bin:/bin" as *u8 as i64; envp[1]=0
72 sys_execve(SF_ADMIT_ELF, argv, envp); sys_exit(127)
73 }
74 let st: *i64 = sys_mmap(16) as *i64; sys_wait4(pid, st, 0)
75 let ex: i64 = (st[0]>>8)&0xff
76 if ex==0 { return 1 }
77 return 0
78}
79
80// fork one ship for a target (detached -- parent tracks the pid); returns pid (>0) or -1 on fork fail.
81func sf_launch(target: *u8) -> i64 {
82 let pid: i64 = sys_fork()
83 if pid < 0 { return 0-1 }
84 if pid == 0 {
85 let dn: i64 = sys_openat_wr(SF_DEVNULL, 420)
86 if dn>=0 { sys_dup3(dn,1,0); sys_dup3(dn,2,0) }
87 let argv: *i64 = sys_mmap(64) as *i64
88 argv[0]=SF_SHIP_ELF as i64; argv[1]=target as i64; argv[2]=0
89 let envp: *i64 = sys_mmap(16) as *i64; envp[0]="PATH=/usr/bin:/bin" as *u8 as i64; envp[1]=0
90 sys_execve(SF_SHIP_ELF, argv, envp); sys_exit(127)
91 }
92 return pid
93}
94
95// the polite governed WIDTH ceiling for this instant (0 if memory is below floor).
96func sf_rg_budget(ncpu: i64, load: i64, freemb: i64) -> i64 {
97 if rg_memory_ok(freemb, SF_MEM_FLOOR_MB)==0 { return 0 }
98 let reserve: i64 = rg_hw_class_reserve(ncpu) // sensor->supercomputer class-aware reserve
99 return rg_worker_budget(RG_POLITE, ncpu, load, reserve, SF_CEIL_MILLI)
100}
101
102func main(argc: i64, argv: *i64) -> i64 {
103 if argc < 2 {
104 sf_w("usage: nx_ship_fleet <targets-file> [plan]\n" as *u8)
105 return 1
106 }
107 let path: *u8 = argv[1] as *u8
108 var plan_only: i64 = 0
109 if argc >= 3 { let a2: *u8 = argv[2] as *u8; if a2[0]==(112 as u8) { plan_only=1 } } // 'p'lan
110
111 // parse the target list.
112 let buf: *u8 = sys_mmap(SF_READCAP)
113 let n: i64 = sf_read(path, buf, SF_READCAP)
114 if n <= 0 { sf_w("SHIP-FLEET no-targets file=" as *u8); sf_w(path); sf_w("\n" as *u8); return 1 }
115 let names: *u8 = sys_mmap(SF_MAXTARGETS * SF_NAMECAP)
116 var nt: i64 = 0
117 var ls: i64 = 0; var i: i64 = 0
118 while i <= n {
119 var eol: i64 = 0
120 if i==n { eol=1 } else { if buf[i]==(10 as u8) { eol=1 } }
121 if eol==1 {
122 if i>ls { if buf[ls]!=(35 as u8) { // skip '#'
123 if nt < SF_MAXTARGETS {
124 let dst: *u8 = (names + nt*SF_NAMECAP) as *u8
125 var k: i64 = 0; var p: i64 = ls
126 while p<i { if buf[p]!=(13 as u8) { if k < SF_NAMECAP-1 { dst[k]=buf[p]; k=k+1 } } p=p+1 }
127 dst[k]=0 as u8
128 if k>0 { nt=nt+1 }
129 }
130 } }
131 ls=i+1
132 }
133 i=i+1
134 }
135 sf_w("=== nx_ship_fleet: adaptive-concurrency drain of " as *u8); sf_n(nt); sf_w(" targets from " as *u8); sf_w(path); sf_w(" ===\n" as *u8)
136 if nt==0 { sf_w("SHIP-FLEET nothing-to-ship verdict=GREEN\n" as *u8); return 0 }
137
138 // SENSE once for the header / plan mode.
139 let ncpu0: i64=sl_ncpu(); let load0: i64=sl_loadavg_milli(); let free0: i64=sl_freemem_mb(); let conns0: i64=sl_active_conns_8443()
140 let rgb0: i64 = sf_rg_budget(ncpu0, load0, free0)
141 let admit0: i64 = sf_admit_grant()
142 let w0: i64 = sf_effective_width(rgb0, rgb0, admit0, rg_serve_first(conns0))
143 sf_w(" SENSE ncpu=" as *u8); sf_n(ncpu0); sf_w(" load=" as *u8); sf_n(load0); sf_w("milli free=" as *u8); sf_n(free0)
144 sf_w("MB conns=" as *u8); sf_n(conns0); sf_w(" -> rg_budget=" as *u8); sf_n(rgb0)
145 sf_w(" admit=" as *u8); if admit0==1 { sf_w("GRANT" as *u8) } else { sf_w("REFUSE" as *u8) }
146 sf_w(" -> start width=" as *u8); sf_n(w0); sf_w("\n" as *u8)
147 if plan_only==1 {
148 sf_w("SHIP-FLEET plan-only (forked nothing) would_start_width=" as *u8); sf_n(w0)
149 sf_w(" targets=" as *u8); sf_n(nt); sf_w(" verdict=GREEN\n" as *u8)
150 return 0
151 }
152
153 // the adaptive concurrency loop.
154 let pids: *i64 = sys_mmap(SF_INFLIGHT_CAP*8) as *i64
155 let ptgt: *i64 = sys_mmap(SF_INFLIGHT_CAP*8) as *i64
156 var next: i64=0; var inflight: i64=0; var done: i64=0; var failed: i64=0
157 var maxseen: i64=0; var cap: i64=1; var storm_ticks: i64=0
158 let st: *i64 = sys_mmap(16) as *i64
159
160 while (next < nt) | (inflight > 0) {
161 // reap every finished child (non-blocking).
162 var reaping: i64 = 1
163 while reaping==1 {
164 let r: i64 = sys_wait4(0-1, st, 1) // WNOHANG
165 if r > 0 {
166 let ex: i64 = (st[0]>>8)&0xff
167 var j: i64=0; var found: i64=0-1
168 while j<inflight { if pids[j]==r { found=j; j=inflight } else { j=j+1 } }
169 if found>=0 {
170 let ti: i64 = ptgt[found]
171 if ex==0 { done=done+1; sf_w(" SHIPPED " as *u8) } else { failed=failed+1; sf_w(" FAILED " as *u8) }
172 sf_w((names + ti*SF_NAMECAP) as *u8); sf_w(" exit=" as *u8); sf_n(ex)
173 sf_w(" (done=" as *u8); sf_n(done); sf_w(" failed=" as *u8); sf_n(failed); sf_w(")\n" as *u8)
174 let last: i64 = inflight-1
175 pids[found]=pids[last]; ptgt[found]=ptgt[last]
176 inflight=last
177 }
178 } else { reaping=0 }
179 }
180 if (next >= nt) & (inflight==0) { reaping=0; next=next } // done; fall through to exit
181
182 if (next < nt) | (inflight > 0) {
183 // SENSE + GOVERN + ADMIT.
184 let ncpu: i64=sl_ncpu(); let load: i64=sl_loadavg_milli(); let freemb: i64=sl_freemem_mb(); let conns: i64=sl_active_conns_8443()
185 let rgb: i64 = sf_rg_budget(ncpu, load, freemb)
186 let admit: i64 = sf_admit_grant()
187 let width: i64 = sf_effective_width(rgb, cap, admit, rg_serve_first(conns))
188 if width==0 {
189 cap = sf_aimd_next(cap, rgb, SF_STORM)
190 if inflight==0 {
191 storm_ticks = storm_ticks+1
192 if storm_ticks > SF_MAX_STORM_TICKS {
193 sf_w(" STORM-GIVEUP after " as *u8); sf_n(storm_ticks); sf_w(" idle storm ticks -- deferring the rest (killable, re-run to resume)\n" as *u8)
194 next = nt // stop launching; nothing in flight -> loop exits
195 }
196 }
197 var bo: i64 = rg_backoff_ms(RG_POLITE, ncpu, load, SF_CEIL_MILLI, SF_BASE_BACKOFF_MS)
198 if bo < SF_BASE_BACKOFF_MS { bo = SF_BASE_BACKOFF_MS }
199 sys_sleep_ms(bo)
200 } else {
201 storm_ticks = 0
202 let slots: i64 = sf_launch_slots(width, inflight)
203 var launched: i64=0; var s: i64=0
204 while (s < slots) & (next < nt) {
205 if inflight < SF_INFLIGHT_CAP {
206 let pid: i64 = sf_launch((names + next*SF_NAMECAP) as *u8)
207 if pid > 0 {
208 pids[inflight]=pid; ptgt[inflight]=next; inflight=inflight+1; next=next+1; launched=launched+1
209 } else { s = slots } // fork failed -> stop launching this tick
210 } else { s = slots }
211 s=s+1
212 }
213 if inflight > maxseen { maxseen=inflight }
214 if launched>0 {
215 cap = sf_aimd_next(cap, rgb, SF_OK)
216 sf_w(" tick launch=" as *u8); sf_n(launched); sf_w(" inflight=" as *u8); sf_n(inflight)
217 sf_w(" width=" as *u8); sf_n(width); sf_w(" cap=" as *u8); sf_n(cap); sf_w(" load=" as *u8); sf_n(load); sf_w("milli\n" as *u8)
218 } else {
219 cap = sf_aimd_next(cap, rgb, SF_HOLD)
220 }
221 sys_sleep_ms(SF_TICK_MS)
222 }
223 }
224 }
225
226 let deferred: i64 = nt - done - failed
227 sf_w("SHIP-FLEET targets=" as *u8); sf_n(nt); sf_w(" shipped=" as *u8); sf_n(done)
228 sf_w(" failed=" as *u8); sf_n(failed); sf_w(" deferred=" as *u8); sf_n(deferred)
229 sf_w(" max_concurrency=" as *u8); sf_n(maxseen); sf_w(" final_cap=" as *u8); sf_n(cap)
230 sf_w(" verdict=" as *u8)
231 if (failed==0) & (deferred==0) { sf_w("GREEN\n" as *u8); return 0 }
232 sf_w("AMBER\n" as *u8)
233 return 0
234}