code wiki / _hdl_build / nx_role_queue.nx

nx_role_queue.nx source

↩ module page · 115 lines · 7527 B

1// nx_role_queue.nx -- the ROLE-PARAMETERIZED serialized work intake (generalizes nx_doctor_queue so the SAME 2// mechanism serves EVERY change-applying role: Doctor=fix, Engineer=optimize, ... each with its OWN isolated 3// queue/ledger/lock). A role SUBMITs a change request; the role PUMPs its own queue one-at-a-time under its 4// per-role lock, applying each via the gated graceful doc_apply_mem (lock + exact-context-match + atomic), then 5// flips the line to APPLIED/CONFLICT/AMBIGUOUS + ledgers. Idempotent. Per-role isolation = roles never compete. 6// 7// queue line: status<TAB>ts<TAB>role<TAB>requester<TAB>target<TAB>oldstr-path<TAB>newstr-path 8// usage: nx_role_queue <role> submit <target> <oldfile> <newfile> <requester> 9// nx_role_queue <role> pump 10// license_tier: ORIGINAL 11import "nx_syscalls.nx" 12import "nx_arbiter.nx" 13import "nx_doctor_apply.nx" // doc_apply_mem / da_read / DA_* / da_slen 14import "nx_role_store.nx" // rs_append / rs_row_set / rs_get_seq / rs_count / rs_chan / rs_field -- the SOVEREIGN seg-store (NO TSV) 15const RQ_MAGIC_4096: i64 = 4096 16const RQ_MAGIC_100000: i64 = 100000 17const RQ_MAGIC_1024: i64 = 1024 18const RQ_MAGIC_2048: i64 = 2048 19 20const RQ_DIR: *u8 = "knowledge/roles" 21 22func rq_w(s: *u8) -> i64 { return sys_write(1, s, da_slen(s)) } 23func rq_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 } 24func rq_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 } 25func rq_n(dst: *u8, off: i64, v: i64) -> i64 { var o: i64=off; var m: i64=v; if m<0{dst[o]=45 as u8;o=o+1;m=0-m} let t: *u8=sys_mmap(24); var k: i64=0; if m==0{t[0]=48 as u8;k=1} while m>0{t[k]=(48+(m-(m/10)*10)) as u8;m=m/10;k=k+1} var i: i64=0; while i<k{dst[o+i]=t[k-1-i];i=i+1} return o+k } 26func rq_field(line: *u8, len: i64, idx: i64, out: *u8, outcap: i64) -> i64 { 27 var f: i64=0; var i: i64=0; var w: i64=0 28 while i<len { let c: i64=line[i] as i64; 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 } } } i=i+1 } 29 if f==idx { out[w]=0 as u8; return 1 } return 0 30} 31// build "knowledge/roles/<role><suffix>" into out 32func rq_path(role: *u8, suffix: *u8, out: *u8) -> i64 { 33 var o: i64 = rq_cat(out, 0, RQ_DIR); out[o]=47 as u8; o=o+1 // '/' 34 o = rq_cat(out, o, role); o = rq_cat(out, o, suffix); out[o]=0 as u8; return o 35} 36 37// SUBMIT: append a PENDING row to the role's queue CHANNEL in the sovereign seg-store (knowledge/store/roles-*). 38// Row schema UNCHANGED (status<TAB>ts<TAB>role<TAB>requester<TAB>target<TAB>oldp<TAB>newp); the value is the same 39// bytes the .tsv line held, minus the trailing newline -- only the HOME changed (seg-store record, not a .tsv file). 40func rq_submit(role: *u8, target: *u8, oldp: *u8, newp: *u8, requester: *u8) -> i64 { 41 let rec: *u8 = sys_mmap(RQ_MAGIC_4096); var o: i64 = 0 42 o = rq_cat(rec, o, "PENDING" as *u8); rec[o]=9 as u8; o=o+1 43 o = rq_n(rec, o, sys_now_realtime_sec()); rec[o]=9 as u8; o=o+1 44 o = rq_cat(rec, o, role); rec[o]=9 as u8; o=o+1 45 o = rq_cat(rec, o, requester); rec[o]=9 as u8; o=o+1 46 o = rq_cat(rec, o, target); rec[o]=9 as u8; o=o+1 47 o = rq_cat(rec, o, oldp); rec[o]=9 as u8; o=o+1 48 o = rq_cat(rec, o, newp) 49 let ch: *u8 = sys_mmap(96); rs_chan(role, 113, ch) // '<role>q' 50 if rs_append(ch, rec, o) < 0 { return 0 } 51 return 1 52} 53 54// PUMP: drain the role's queue CHANNEL in the sovereign seg-store. Walk every row; for each PENDING one, apply 55// via the gated doc_apply_mem, append a row to the role's LEDGER channel, and flip the row's status by writing a 56// new version of its <role>q:<seq> record (seg-store last-wins -- no file rewrite). Same semantics as the old 57// .tsv pump, store-backed. Held under the per-role pump lock (fl_acquire) so a role drains one-at-a-time. 58func rq_pump(role: *u8) -> i64 { 59 let qch: *u8 = sys_mmap(96); rs_chan(role, 113, qch) // '<role>q' 60 let lch: *u8 = sys_mmap(96); rs_chan(role, 108, lch) // '<role>l' 61 let lock: *u8 = sys_mmap(128); var lo: i64 = rq_cat(lock, 0, "rq_" as *u8); lo = rq_cat(lock, lo, role); lock[lo]=0 as u8 62 let lk: i64 = fl_acquire(lock, RQ_MAGIC_100000, 1) 63 let n: i64 = rs_count(qch) 64 let pq: *i64 = sys_mmap(16) as *i64; let lq: *i64 = sys_mmap(16) as *i64 65 let st: *u8=sys_mmap(64); let ts: *u8=sys_mmap(64); let rl: *u8=sys_mmap(64); let req: *u8=sys_mmap(256) 66 let tgt: *u8=sys_mmap(RQ_MAGIC_1024); let oldp: *u8=sys_mmap(RQ_MAGIC_1024); let newp: *u8=sys_mmap(RQ_MAGIC_1024) 67 let oldb: *u8=sys_mmap(DA_CAP); let newb: *u8=sys_mmap(DA_CAP) 68 var processed: i64 = 0; var seq: i64 = 0 69 while seq < n { 70 if rs_get_seq(qch, seq, pq, lq) == 1 { 71 let row: *u8 = pq[0] as *u8; let rlen: i64 = lq[0] 72 rs_field(row,rlen,0,st); rs_field(row,rlen,1,ts); rs_field(row,rlen,2,rl); rs_field(row,rlen,3,req) 73 rs_field(row,rlen,4,tgt); rs_field(row,rlen,5,oldp); rs_field(row,rlen,6,newp) 74 if rq_streq(st, "PENDING" as *u8) == 1 { 75 let oldlen: i64 = da_read(oldp, oldb, DA_CAP); let newlen: i64 = da_read(newp, newb, DA_CAP) 76 var code: i64 = DA_ERR 77 if oldlen > 0 { code = doc_apply_mem(tgt, oldb, oldlen, newb, newlen) } 78 var newst: *u8 = "ERROR" as *u8 79 if code==DA_APPLIED { newst="APPLIED" as *u8 } 80 if code==DA_CONFLICT { newst="CONFLICT" as *u8 } 81 if code==DA_AMBIG { newst="AMBIGUOUS" as *u8 } 82 let lr: *u8=sys_mmap(RQ_MAGIC_2048); var lz: i64=0 83 lz=rq_cat(lr,lz,newst); lr[lz]=9 as u8; lz=lz+1; lz=rq_cat(lr,lz,rl); lr[lz]=9 as u8; lz=lz+1 84 lz=rq_cat(lr,lz,req); lr[lz]=9 as u8; lz=lz+1; lz=rq_cat(lr,lz,tgt) 85 rs_append(lch, lr, lz) 86 let nr: *u8=sys_mmap(RQ_MAGIC_4096); var no: i64=0 87 no=rq_cat(nr,no,newst); nr[no]=9 as u8; no=no+1; no=rq_cat(nr,no,ts); nr[no]=9 as u8; no=no+1 88 no=rq_cat(nr,no,rl); nr[no]=9 as u8; no=no+1; no=rq_cat(nr,no,req); nr[no]=9 as u8; no=no+1 89 no=rq_cat(nr,no,tgt); nr[no]=9 as u8; no=no+1; no=rq_cat(nr,no,oldp); nr[no]=9 as u8; no=no+1 90 no=rq_cat(nr,no,newp) 91 rs_row_set(qch, seq, nr, no) 92 processed = processed + 1 93 } 94 } 95 seq = seq + 1 96 } 97 fl_release(lk) 98 return processed 99} 100 101func main(argc: i64, argv: *i64) -> i64 { 102 if argc < 3 { rq_w("usage: nx_role_queue <role> submit <target> <oldfile> <newfile> <requester> | <role> pump\n" as *u8); return 2 } 103 let role: *u8 = argv[1] as *u8; let cmd: *u8 = argv[2] as *u8 104 if rq_streq(cmd, "submit" as *u8) == 1 { 105 if argc < 7 { rq_w("usage: nx_role_queue <role> submit <target> <oldfile> <newfile> <requester>\n" as *u8); return 2 } 106 let rc: i64 = rq_submit(role, argv[3] as *u8, argv[4] as *u8, argv[5] as *u8, argv[6] as *u8) 107 if rc==1 { rq_w("ROLE-SUBMIT OK role=" as *u8); rq_w(role); rq_w(" (queued PENDING)\n" as *u8); return 0 } 108 rq_w("ROLE-SUBMIT FAIL\n" as *u8); return 1 109 } 110 if rq_streq(cmd, "pump" as *u8) == 1 { 111 let np: i64 = rq_pump(role) 112 rq_w("ROLE-PUMP role=" as *u8); rq_w(role); rq_w(" processed " as *u8); let b: *u8=sys_mmap(24); let e: i64=rq_n(b,0,np); sys_write(1,b,e); rq_w("\n" as *u8); return 0 113 } 114 rq_w("unknown subcommand\n" as *u8); return 2 115}