code wiki / _hdl_build / nx_doctor_queue.nx
nx_doctor_queue.nx source
↩ module page · 134 lines · 7771 B
1// nx_doctor_queue.nx -- THE NISHI DOCTOR's serialized fix INTAKE (so a fix is literally "sent to the Doctor"
2// and applied one-at-a-time, no competition). Mirrors the Publisher's submit->queue->ship: workstreams SUBMIT
3// (append a PENDING fix request); the Doctor PUMPS (drains PENDING one-at-a-time under a pump-lock, applying
4// each via the gated graceful doc_apply_mem, flipping the queue line to APPLIED/CONFLICT/AMBIGUOUS, ledgering).
5// Idempotent (rule #10): only PENDING lines are processed; a re-run pumps 0. RACI: the Doctor APPLIES; the
6// Engineer re-verifies + the Warden admits (separate roles) -- this organ never blesses its own work.
7//
8// queue line (TSV): status<TAB>ts<TAB>requester<TAB>target<TAB>oldstr-path<TAB>newstr-path
9// usage: nx_doctor_queue submit <target> <oldstr-file> <newstr-file> <requester>
10// nx_doctor_queue pump
11// license_tier: ORIGINAL
12import "nx_syscalls.nx"
13import "nx_arbiter.nx" // fl_acquire / fl_release
14import "nx_doctor_apply.nx" // doc_apply_mem / da_read / DA_* / da_slen
15const DQ_MAGIC_4096: i64 = 4096
16const DQ_MAGIC_100000: i64 = 100000
17const DQ_MAGIC_1024: i64 = 1024
18const DQ_MAGIC_2048: i64 = 2048
19
20const DQ_DIR: *u8 = "knowledge/doctor"
21const DQ_QUEUE: *u8 = "knowledge/doctor/fix_queue.tsv"
22const DQ_LEDGER: *u8 = "knowledge/doctor/fix_ledger.tsv"
23const DQ_PUMPLOCK: *u8 = "doctor_pump"
24const DQ_CAP: i64 = 1048576
25
26func dq_w(s: *u8) -> i64 { return sys_write(1, s, da_slen(s)) }
27func dq_streq(a: *u8, b: *u8) -> i64 { var i: i64 = 0; while a[i] != (0 as u8) { if a[i] != b[i] { return 0 } i = i + 1 } if b[i] != (0 as u8) { return 0 } return 1 }
28func dq_cat(dst: *u8, off: i64, s: *u8) -> i64 { var o: i64 = off; var i: i64 = 0; while s[i] != (0 as u8) { dst[o] = s[i]; o = o + 1; i = i + 1 } return o }
29func dq_n(dst: *u8, off: i64, v: i64) -> i64 {
30 var o: i64 = off; var m: i64 = v; if m < 0 { dst[o]=45 as u8; o=o+1; m = 0 - m }
31 let t: *u8 = sys_mmap(24); var k: i64 = 0; if m == 0 { t[0]=48 as u8; k=1 }
32 while m > 0 { t[k]=(48+(m-(m/10)*10)) as u8; m=m/10; k=k+1 }
33 var i: i64 = 0; while i < k { dst[o+i]=t[k-1-i]; i=i+1 } return o + k
34}
35// idx-th tab/newline-delimited field of line[0..len) -> out (NUL-term). returns 1 if found.
36func dq_field(line: *u8, len: i64, idx: i64, out: *u8, outcap: i64) -> i64 {
37 var f: i64 = 0; var i: i64 = 0; var w: i64 = 0
38 while i < len {
39 let c: i64 = line[i] as i64
40 if c == 9 { if f == idx { out[w] = 0 as u8; return 1 } f = f + 1 } else { if f == idx { if w < outcap - 1 { out[w] = line[i]; w = w + 1 } } }
41 i = i + 1
42 }
43 if f == idx { out[w] = 0 as u8; return 1 }
44 return 0
45}
46
47func doc_submit(target: *u8, oldpath: *u8, newpath: *u8, requester: *u8) -> i64 {
48 sys_mkdir(DQ_DIR, 0x1ff)
49 let rec: *u8 = sys_mmap(DQ_MAGIC_4096); var o: i64 = 0
50 o = dq_cat(rec, o, "PENDING" as *u8); rec[o]=9 as u8; o=o+1
51 o = dq_n(rec, o, sys_now_realtime_sec()); rec[o]=9 as u8; o=o+1
52 o = dq_cat(rec, o, requester); rec[o]=9 as u8; o=o+1
53 o = dq_cat(rec, o, target); rec[o]=9 as u8; o=o+1
54 o = dq_cat(rec, o, oldpath); rec[o]=9 as u8; o=o+1
55 o = dq_cat(rec, o, newpath); rec[o]=10 as u8; o=o+1
56 let fd: i64 = sys_openat_append(DQ_QUEUE, 0x1a4); if fd < 0 { return 0 }
57 sys_write(fd, rec, o); sys_close(fd)
58 return 1
59}
60
61// drain PENDING -> apply -> flip status -> ledger. Returns count processed.
62func doc_pump(queue: *u8, ledger: *u8) -> i64 {
63 let lk: i64 = fl_acquire(DQ_PUMPLOCK, DQ_MAGIC_100000, 1)
64 let buf: *u8 = sys_mmap(DQ_CAP); let n: i64 = da_read(queue, buf, DQ_CAP)
65 if n <= 0 { fl_release(lk); return 0 }
66 let out: *u8 = sys_mmap(DQ_CAP + DQ_MAGIC_4096); var oo: i64 = 0
67 let st: *u8 = sys_mmap(64); let ts: *u8 = sys_mmap(64); let req: *u8 = sys_mmap(256)
68 let tgt: *u8 = sys_mmap(DQ_MAGIC_1024); let oldp: *u8 = sys_mmap(DQ_MAGIC_1024); let newp: *u8 = sys_mmap(DQ_MAGIC_1024)
69 let oldb: *u8 = sys_mmap(DA_CAP); let newb: *u8 = sys_mmap(DA_CAP)
70 var processed: i64 = 0
71 var ls: i64 = 0; var i: i64 = 0
72 while i <= n {
73 var nl: i64 = 0
74 if i >= n { nl = 1 } else { if buf[i] == (10 as u8) { nl = 1 } }
75 if nl == 1 {
76 let ll: i64 = i - ls
77 if ll > 0 {
78 let line: *u8 = ((buf as i64) + ls) as *u8
79 dq_field(line, ll, 0, st, 64); dq_field(line, ll, 1, ts, 64); dq_field(line, ll, 2, req, 256)
80 dq_field(line, ll, 3, tgt, DQ_MAGIC_1024); dq_field(line, ll, 4, oldp, DQ_MAGIC_1024); dq_field(line, ll, 5, newp, DQ_MAGIC_1024)
81 var newst: *u8 = st // default: pass through unchanged
82 if dq_streq(st, "PENDING" as *u8) == 1 {
83 let oldlen: i64 = da_read(oldp, oldb, DA_CAP)
84 let newlen: i64 = da_read(newp, newb, DA_CAP)
85 var code: i64 = DA_ERR
86 if oldlen > 0 { code = doc_apply_mem(tgt, oldb, oldlen, newb, newlen) }
87 if code == DA_APPLIED { newst = "APPLIED" as *u8 }
88 if code == DA_CONFLICT { newst = "CONFLICT" as *u8 }
89 if code == DA_AMBIG { newst = "AMBIGUOUS" as *u8 }
90 if code == DA_ERR { newst = "ERROR" as *u8 }
91 // ledger: result ts requester target
92 let lr: *u8 = sys_mmap(DQ_MAGIC_2048); var lo: i64 = 0
93 lo = dq_cat(lr, lo, newst); lr[lo]=9 as u8; lo=lo+1
94 lo = dq_cat(lr, lo, ts); lr[lo]=9 as u8; lo=lo+1
95 lo = dq_cat(lr, lo, req); lr[lo]=9 as u8; lo=lo+1
96 lo = dq_cat(lr, lo, tgt); lr[lo]=10 as u8; lo=lo+1
97 let lfd: i64 = sys_openat_append(ledger, 0x1a4); if lfd >= 0 { sys_write(lfd, lr, lo); sys_close(lfd) }
98 processed = processed + 1
99 }
100 // re-emit the line (only status may change)
101 oo = dq_cat(out, oo, newst); out[oo]=9 as u8; oo=oo+1
102 oo = dq_cat(out, oo, ts); out[oo]=9 as u8; oo=oo+1
103 oo = dq_cat(out, oo, req); out[oo]=9 as u8; oo=oo+1
104 oo = dq_cat(out, oo, tgt); out[oo]=9 as u8; oo=oo+1
105 oo = dq_cat(out, oo, oldp); out[oo]=9 as u8; oo=oo+1
106 oo = dq_cat(out, oo, newp); out[oo]=10 as u8; oo=oo+1
107 }
108 ls = i + 1
109 }
110 i = i + 1
111 }
112 // atomic queue rewrite
113 let tmp: *u8 = sys_mmap(DQ_MAGIC_1024); var t: i64 = 0; while queue[t] != (0 as u8) { tmp[t]=queue[t]; t=t+1 }
114 let sfx: *u8 = ".tmp"; var s: i64 = 0; while sfx[s] != (0 as u8) { tmp[t]=sfx[s]; t=t+1; s=s+1 } tmp[t]=0 as u8
115 let wfd: i64 = sys_openat_wr(tmp, 0x1a4); if wfd >= 0 { sys_write(wfd, out, oo); sys_close(wfd); sys_renameat(tmp, queue) }
116 fl_release(lk)
117 return processed
118}
119
120func main(argc: i64, argv: *i64) -> i64 {
121 if argc < 2 { dq_w("usage: nx_doctor_queue submit <target> <oldfile> <newfile> <requester> | pump\n" as *u8); return 2 }
122 let cmd: *u8 = argv[1] as *u8
123 if dq_streq(cmd, "submit" as *u8) == 1 {
124 if argc < 6 { dq_w("usage: nx_doctor_queue submit <target> <oldfile> <newfile> <requester>\n" as *u8); return 2 }
125 let rc: i64 = doc_submit(argv[2] as *u8, argv[3] as *u8, argv[4] as *u8, argv[5] as *u8)
126 if rc == 1 { dq_w("DOCTOR-SUBMIT OK (queued PENDING)\n" as *u8); return 0 }
127 dq_w("DOCTOR-SUBMIT FAIL\n" as *u8); return 1
128 }
129 if dq_streq(cmd, "pump" as *u8) == 1 {
130 let np: i64 = doc_pump(DQ_QUEUE, DQ_LEDGER)
131 dq_w("DOCTOR-PUMP processed " as *u8); let b: *u8=sys_mmap(24); let e: i64=dq_n(b,0,np); sys_write(1,b,e); dq_w(" PENDING fix(es)\n" as *u8); return 0
132 }
133 dq_w("unknown subcommand\n" as *u8); return 2
134}