code wiki / _hdl_build / nx_qserver.nx

nx_qserver.nx source

↩ module page · 166 lines · 8565 B

1// nx_qserver.nx -- MULTI-SLOT sovereign job queue: one persistent process that scans a directory for *.job 2// files (sys_getdents64, the same syscall nx_claude_bridge's ls uses) and runs each, so the agent loop can 3// have MANY jobs in flight (vs nx_jobserver's single fixed slot). Per job <base>.job: 4// read it (line1 elf, lines2.. args VERBATIM) -> fork/exec directly with a WATCHDOG (kill -9 at <watchdog>s) 5// -> capture stdout+stderr -> <base>.out -> write exit code -> <base>.res -> delete <base>.job (consume). 6// Collects names FIRST, then processes (so unlinking never disturbs the getdents scan). Stop via <dir>/nxq.stop. 7// usage: nx_qserver <queuedir> <pollbudget> <watchdog-sec> (budget=0 = forever). No sh, no quoting. 8// Sovereign: nx_syscalls only. license_tier: ORIGINAL 9import "nx_syscalls.nx" 10import "nx_itoa_lib.nx" // shared MSB-first emitter (zero-alloc) 11const K_MAGIC_1000000: i64 = 1000000 12const K_MAGIC_65536: i64 = 65536 13const K_MAGIC_1024: i64 = 1024 14 15func qs_w(s: *u8) -> i64 { var n: i64 = 0; while s[n] != (0 as u8) { n = n + 1 } sys_write(1, s, n); return 0 } 16// MIGRATED to the shared emitter (debt 1785563586). The old body mmapped a scratch buffer 17// per call and never freed it. At PAGE granularity that is 4096B leaked PER CALL -- the 18// defect that took 28.5GB of a 36GB host in nx_ts_lumadiff (2MB input, ~3.66M calls). 19// nxi_* is MSB-first, allocates NOTHING, and emits identical bytes including the sign. 20func qs_wn(v: i64) -> i64 { nxi_out(v); return 0 } 21func qs_slen(s: *u8) -> i64 { var n: i64 = 0; while s[n] != (0 as u8) { n = n + 1 } return n } 22func qs_atoi(s: *u8) -> i64 { var v: i64 = 0; var i: i64 = 0; while s[i] != (0 as u8) { let c: i64 = s[i] as i64; if c >= 48 { if c <= 57 { v = v * 10 + (c - 48) } } i = i + 1 } return v } 23func qs_unlink(path: *u8) -> i64 { __syscall(263, AT_FDCWD, path, 0, 0, 0, 0) return 0 } 24func qs_sleep_ms(ms: i64) -> i64 { let ts: *i64 = sys_mmap(16) as *i64; ts[0] = ms / 1000; ts[1] = (ms % 1000) * K_MAGIC_1000000; return __syscall(SYS_CLOCK_NANOSLEEP, 0, 0, ts, 0, 0, 0) } 25// MIGRATED to the shared emitter (debt 1785563586). The old body mmapped a scratch buffer 26// per call and never freed it. At PAGE granularity that is 4096B leaked PER CALL -- the 27// defect that took 28.5GB of a 36GB host in nx_ts_lumadiff (2MB input, ~3.66M calls). 28// nxi_* is MSB-first, allocates NOTHING, and emits identical bytes including the sign. 29func qs_wn_fd(fd: i64, v: i64) -> i64 { nxi_fd(fd, v); return 0 } 30func qs_path(dir: *u8, name: *u8, out: *u8) -> i64 { 31 var o: i64 = 0; var i: i64 = 0 32 while dir[i] != (0 as u8) { out[o] = dir[i]; o = o + 1; i = i + 1 } 33 out[o] = 47 as u8; o = o + 1 34 i = 0; while name[i] != (0 as u8) { out[o] = name[i]; o = o + 1; i = i + 1 } 35 out[o] = 0 as u8 36 return o 37} 38// dir + "/" + name[0..namelen) + suffix 39func qs_rpath(dir: *u8, name: *u8, namelen: i64, suffix: *u8, out: *u8) -> i64 { 40 var o: i64 = 0; var i: i64 = 0 41 while dir[i] != (0 as u8) { out[o] = dir[i]; o = o + 1; i = i + 1 } 42 out[o] = 47 as u8; o = o + 1 43 i = 0; while i < namelen { out[o] = name[i]; o = o + 1; i = i + 1 } 44 i = 0; while suffix[i] != (0 as u8) { out[o] = suffix[i]; o = o + 1; i = i + 1 } 45 out[o] = 0 as u8 46 return o 47} 48func qs_ends_job(name: *u8) -> i64 { 49 let n: i64 = qs_slen(name) 50 if n < 5 { return 0 } 51 if (name[n - 4] as i64) != 46 { return 0 } 52 if (name[n - 3] as i64) != 106 { return 0 } 53 if (name[n - 2] as i64) != 111 { return 0 } 54 if (name[n - 1] as i64) != 98 { return 0 } 55 return 1 56} 57// parse a job buffer + fork/exec with output->out_path + watchdog. returns exit code. 58func qs_run_job(reqbuf: *u8, n: i64, out_path: *u8, wdog_sec: i64) -> i64 { 59 let buf: *u8 = sys_mmap(n + 8) 60 var ci: i64 = 0; while ci < n { buf[ci] = reqbuf[ci]; ci = ci + 1 } buf[n] = 10 as u8 61 let cargv: *i64 = sys_mmap(8 * 128) as *i64 62 var ac: i64 = 0; var i: i64 = 0; var start: i64 = 0 63 while i <= n { 64 if (buf[i] as i64) == 10 { 65 var end: i64 = i 66 if end > start { if (buf[end - 1] as i64) == 13 { end = end - 1 } } 67 if end > start { buf[end] = 0 as u8; if ac < 127 { cargv[ac] = ((buf as i64) + start) as i64; ac = ac + 1 } } 68 start = i + 1 69 } 70 i = i + 1 71 } 72 cargv[ac] = 0 73 if ac == 0 { return 125 } 74 let out_fd: i64 = sys_openat_wr(out_path, 0x180) 75 if out_fd < 0 { return 126 } 76 let pid: i64 = sys_fork() 77 if pid == 0 { 78 sys_dup3(out_fd, 1, 0); sys_dup3(out_fd, 2, 0) 79 let envp: *i64 = sys_mmap(16) as *i64 80 envp[0] = "PATH=/usr/bin:/bin" as *u8 as i64; envp[1] = 0 81 sys_execve(cargv[0] as *u8, cargv, envp) 82 sys_exit(127) 83 } 84 let st: *i64 = sys_mmap(16) as *i64 85 var ticks: i64 = 0; let deadline: i64 = wdog_sec * 5; var rc: i64 = 124; var waiting: i64 = 1 86 while waiting == 1 { 87 let w: i64 = sys_wait4(pid, st, 1) 88 if w == pid { waiting = 0; if (st[0] % 128) != 0 { rc = 128 + (st[0] % 128) } else { rc = (st[0] >> 8) & 0xff } } 89 else { ticks = ticks + 1; if ticks >= deadline { waiting = 0; nx_kill(pid, 9); sys_wait4(pid, st, 0); rc = 124 } else { qs_sleep_ms(200) } } 90 } 91 sys_close(out_fd) 92 return rc 93} 94// scan dir for *.job, run each. returns count run this pass. 95func qs_scan(dir: *u8, wdog: i64) -> i64 { 96 let fd: i64 = sys_openat_rd(dir) 97 if fd < 0 { return 0 } 98 let dbuf: *u8 = sys_mmap(K_MAGIC_65536) 99 let names: *u8 = sys_mmap(128 * 256) 100 var ncount: i64 = 0 101 var keepd: i64 = 1 102 while keepd == 1 { 103 let nread: i64 = sys_getdents64(fd, dbuf, K_MAGIC_65536) 104 if nread <= 0 { keepd = 0 } else { 105 var pos: i64 = 0 106 while pos < nread { 107 let reclen: i64 = (dbuf[pos + 16] as i64) + ((dbuf[pos + 17] as i64) << 8) 108 if reclen <= 0 { pos = nread } else { 109 let name: *u8 = ((dbuf as i64) + pos + 19) as *u8 110 if qs_ends_job(name) == 1 { 111 if ncount < 128 { 112 let dst: *u8 = ((names as i64) + ncount * 256) as *u8 113 var ci: i64 = 0; while name[ci] != (0 as u8) { dst[ci] = name[ci]; ci = ci + 1 } dst[ci] = 0 as u8 114 ncount = ncount + 1 115 } 116 } 117 pos = pos + reclen 118 } 119 } 120 } 121 } 122 sys_close(fd) 123 var ran: i64 = 0 124 var k: i64 = 0 125 while k < ncount { 126 let nm: *u8 = ((names as i64) + k * 256) as *u8 127 let namelen: i64 = qs_slen(nm) 128 let jobpath: *u8 = sys_mmap(K_MAGIC_1024); qs_path(dir, nm, jobpath) 129 let outpath: *u8 = sys_mmap(K_MAGIC_1024); qs_rpath(dir, nm, namelen - 4, ".out\x00" as *u8, outpath) 130 let respath: *u8 = sys_mmap(K_MAGIC_1024); qs_rpath(dir, nm, namelen - 4, ".res\x00" as *u8, respath) 131 let lbox: *i64 = sys_mmap(16) as *i64; lbox[0] = 0 132 let jb: *u8 = sys_read_file(jobpath, lbox) 133 if (jb as i64) != 0 { 134 if lbox[0] > 0 { 135 let code: i64 = qs_run_job(jb, lbox[0], outpath, wdog) 136 let rfd: i64 = sys_openat_wr(respath, 0x180) 137 if rfd >= 0 { qs_wn_fd(rfd, code); sys_write(rfd, "\n" as *u8, 1); sys_close(rfd) } 138 qs_unlink(jobpath) 139 ran = ran + 1 140 qs_w(" JOB " as *u8); qs_w(nm); qs_w(" exit=" as *u8); qs_wn(code); qs_w("\n" as *u8) 141 } 142 } 143 k = k + 1 144 } 145 return ran 146} 147 148func main(argc: i64, argv: *i64) -> i64 { 149 if argc < 4 { sys_write(2, "usage: nx_qserver <queuedir> <pollbudget> <watchdog-sec>\n" as *u8, 56); return 1 } 150 let dir: *u8 = argv[1] as *u8 151 let budget: i64 = qs_atoi(argv[2] as *u8) 152 let wdog: i64 = qs_atoi(argv[3] as *u8) 153 let stop_path: *u8 = sys_mmap(K_MAGIC_1024); qs_path(dir, "nxq.stop" as *u8, stop_path) 154 qs_w("nx_qserver: watching " as *u8); qs_w(dir); qs_w(" for *.job (multi-slot; stop via nxq.stop)\n" as *u8) 155 var polls: i64 = 0; var keep: i64 = 1; var total: i64 = 0 156 while keep == 1 { 157 total = total + qs_scan(dir, wdog) 158 let sbox: *i64 = sys_mmap(16) as *i64; sbox[0] = 0 159 let sf: *u8 = sys_read_file(stop_path, sbox) 160 if (sf as i64) != 0 { keep = 0; qs_unlink(stop_path) } 161 if keep == 1 { if budget > 0 { polls = polls + 1; if polls >= budget { keep = 0 } } } 162 if keep == 1 { qs_sleep_ms(200) } 163 } 164 qs_w("nx_qserver: exit (jobs run=" as *u8); qs_wn(total); qs_w(")\n" as *u8) 165 return 0 166}