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}