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}