code wiki / (root) / nx_writejob_sweep.nx

nx_writejob_sweep.nx source

↩ module page · 173 lines · 7400 B

1// nx_writejob_sweep.nx -- drain the writer's durable job queue (operator 2026-08-05: the hub 2// stores requests and runs them later; nothing pushed from a spoke may take the hub down). 3// Walks knowledge/staging/writejobs.queue (append-only ledger of job ids), and for each id that 4// has a .req ticket but no finished .txt: re-runs ./nx_write.elf with the stored mode+prompt. 5// ADMISSION COMPOSES FOR FREE: nx_write itself yields (NX-WRITE-DEFERRED) while a game holds the 6// GPU, so a sweep during play leaves everything queued and costs one probe. Batch per sweep is 7// conf-driven (writejob_sweep_batch). A consumed ticket is renamed .req.done (never deleted). 8// license_tier: ORIGINAL module: nishi-core.write.sweep 9import "nx_syscalls.nx" 10import "nx_lane_conf.nx" 11 12const SW_BUF: i64 = 262144 13const SW_CAP: i64 = 200000 14const SW_LINE: i64 = 4096 15 16func sw_len(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} return n } 17func sw_w(s: *u8) -> i64 { sys_write(1, s, sw_len(s)); return 0 } 18func sw_cat(d: *u8, o: i64, s: *u8) -> i64 { var i: i64=0; while s[i]!=(0 as u8){d[o+i]=s[i];i=i+1} return o+i } 19func sw_u(dst: *u8, off: i64, v: i64) -> i64 { 20 var m: i64=v; if m==0 { dst[off]=48 as u8; return off+1 } 21 let t: *u8=sys_mmap(28); var k: i64=0 22 while m>0 { t[k]=(48+(m-(m/10)*10)) as u8; m=m/10; k=k+1 } 23 var o: i64=off; var i: i64=0 24 while i<k { dst[o]=t[k-1-i]; o=o+1; i=i+1 } 25 return o 26} 27func sw_read(path: *u8, buf: *u8, cap: i64) -> i64 { 28 let fd: i64=sys_openat_rd(path) 29 if fd<0 { return 0-1 } 30 var total: i64=0 31 var r: i64=1 32 while r>0 { r=sys_read(fd, ((buf as i64)+total) as *u8, cap-1-total); if r>0 { total=total+r } if total>=cap-1 { r=0 } } 33 sys_close(fd) 34 buf[total]=0 as u8 35 return total 36} 37func sw_writeall(fd: i64, buf: *u8, n: i64) -> i64 { 38 var off: i64=0 39 while off<n { let w: i64=sys_write(fd, ((buf as i64)+off) as *u8, n-off); if w<=0 { return 0 } off=off+w } 40 return 0 41} 42// fork ./nx_write.elf <mode> <prompt>, capture stdout (cap). Returns bytes. (compact copy of the 43// writer daemon's helper -- third consumer makes this a rule-15 lib extraction next touch) 44func sw_fork(mode: *u8, prompt: *u8, ob: *u8, cap: i64) -> i64 { 45 let pf: *i64=sys_mmap(32) as *i64 46 if sys_pipe2(pf, 0)<0 { return 0-1 } 47 let packed: i64=pf[0] 48 let rfd: i64=packed & 0xFFFFFFFF 49 let wfd: i64=(packed>>32) & 0xFFFFFFFF 50 let pid: i64=sys_fork() 51 if pid==0 { 52 sys_close(rfd) 53 sys_dup3(wfd, 1, 0) 54 let av: *i64=sys_mmap(40) as *i64 55 av[0]="./nx_write.elf" as i64 56 av[1]=mode as i64 57 av[2]=prompt as i64 58 av[3]=0 59 sys_execve("./nx_write.elf" as *u8, av as *i64, 0 as *i64) 60 sys_exit(127) 61 } 62 if pid<0 { sys_close(rfd); sys_close(wfd); return 0-1 } 63 sys_close(wfd) 64 var got: i64=0 65 var rd: i64=1 66 while rd==1 { 67 let n: i64=sys_read(rfd, ((ob as i64)+got) as *u8, cap-got) 68 if n<=0 { rd=0 } else { got=got+n; if got>=cap { rd=0 } } 69 } 70 sys_close(rfd) 71 let stw: *i64=sys_mmap(16) as *i64 72 stw[0]=0 73 sys_wait4(pid, stw, 0) 74 return got 75} 76 77func main(argc: i64, argv: *i64) -> i64 { 78 let q: *u8=sys_mmap(SW_BUF) 79 let qn: i64=sw_read("knowledge/staging/writejobs.queue" as *u8, q, SW_BUF) 80 if qn<=0 { sw_w("WRITEJOB-SWEEP queue empty\n" as *u8); return 0 } 81 let batch: i64=lc_write_lane("writejob_sweep_batch" as *u8, 4) 82 var ran: i64=0 83 var done: i64=0 84 var pending: i64=0 85 var i: i64=0 86 while i<qn { 87 // one id per line 88 let jid: *u8=sys_mmap(64) 89 var jo: i64=0 90 var lrun: i64=1 91 while lrun==1 { 92 if i>=qn { lrun=0 } else { 93 if q[i]==(10 as u8) { lrun=0; i=i+1 } else { 94 if jo<48 { jid[jo]=q[i]; jo=jo+1 } 95 i=i+1 96 } 97 } 98 } 99 jid[jo]=0 as u8 100 if jo>0 { 101 let fin: *u8=sys_mmap(SW_LINE) 102 var fo: i64=sw_cat(fin, 0, "knowledge/staging/writejob_" as *u8) 103 fo=sw_cat(fin, fo, jid) 104 fo=sw_cat(fin, fo, ".txt" as *u8) 105 fin[fo]=0 as u8 106 let ffd: i64=sys_openat_rd(fin) 107 if ffd>=0 { sys_close(ffd); done=done+1 } else { 108 let req: *u8=sys_mmap(SW_LINE) 109 var ro: i64=sw_cat(req, 0, "knowledge/staging/writejob_" as *u8) 110 ro=sw_cat(req, ro, jid) 111 ro=sw_cat(req, ro, ".req" as *u8) 112 req[ro]=0 as u8 113 let rb: *u8=sys_mmap(SW_LINE*4) 114 let rn: i64=sw_read(req, rb, SW_LINE*4) 115 if rn>0 { 116 if ran<batch { 117 // line1=mode, rest=prompt 118 let mode: *u8=sys_mmap(64) 119 var mo: i64=0 120 var pi: i64=0 121 var mrun: i64=1 122 while mrun==1 { 123 if pi>=rn { mrun=0 } else { 124 if rb[pi]==(10 as u8) { mrun=0; pi=pi+1 } else { 125 if mo<48 { mode[mo]=rb[pi]; mo=mo+1 } 126 pi=pi+1 127 } 128 } 129 } 130 mode[mo]=0 as u8 131 let prompt: *u8=((rb as i64)+pi) as *u8 132 let ob: *u8=sys_mmap(SW_BUF) 133 let got: i64=sw_fork(mode, prompt, ob, SW_CAP) 134 var real: i64=1 135 if got<=0 { real=0 } 136 if real==1 { if ob[0]==(110 as u8) { if ob[1]==(120 as u8) { if ob[2]==(95 as u8) { real=0 } } } } 137 if real==1 { if ob[0]==(78 as u8) { if ob[1]==(88 as u8) { if ob[2]==(45 as u8) { real=0 } } } } 138 if real==1 { 139 let part: *u8=sys_mmap(SW_LINE) 140 var po: i64=sw_cat(part, 0, "knowledge/staging/writejob_" as *u8) 141 po=sw_cat(part, po, jid) 142 po=sw_cat(part, po, ".txt.part" as *u8) 143 part[po]=0 as u8 144 let wfd: i64=sys_openat_wr(part, 420) 145 if wfd>=0 { 146 sw_writeall(wfd, ob, got) 147 sys_close(wfd) 148 sys_renameat(part, fin) 149 let dn: *u8=sys_mmap(SW_LINE) 150 var dno: i64=sw_cat(dn, 0, "knowledge/staging/writejob_" as *u8) 151 dno=sw_cat(dn, dno, jid) 152 dno=sw_cat(dn, dno, ".req.done" as *u8) 153 dn[dno]=0 as u8 154 sys_renameat(req, dn) 155 ran=ran+1 156 } 157 } else { pending=pending+1 } 158 } else { pending=pending+1 } 159 } 160 } 161 } 162 } 163 let out: *u8=sys_mmap(256) 164 var o: i64=sw_cat(out, 0, "WRITEJOB-SWEEP ran=" as *u8) 165 o=sw_u(out, o, ran) 166 o=sw_cat(out, o, " already_done=" as *u8) 167 o=sw_u(out, o, done) 168 o=sw_cat(out, o, " still_pending=" as *u8) 169 o=sw_u(out, o, pending) 170 o=sw_cat(out, o, "\n" as *u8) 171 sys_write(1, out, o) 172 return 0 173}