code wiki / (root) / nx_ship_fleet.nx

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}