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}