code wiki / _hdl_build / nx_ws_kickoff_sync.nx

nx_ws_kickoff_sync.nx source

↩ module page · 303 lines · 13945 B

1// nx_ws_kickoff_sync.nx -- CR-R1b MECHANICAL WORKSTREAM CAPTURE: the "sync from kickoff, survive 2// interruptions" rung of the orchestration north-star (Pillar 7). A workstream KICKS OFF, BEATS its 3// state continuously, DONEs on completion; the append-only journal IS the resume surface -- a crash 4// costs nothing but reading it. Each frame is ONE O_APPEND write (atomic for <4KB lines on one host) 5// so PARALLEL sessions converge without clobber (the coindex append-and-derive property). 6// license_tier: ORIGINAL expect_exit: 0 7// kickoff <journal> <ws> <actor> <note> 8// beat <journal> <ws> <actor> <note> 9// done <journal> <ws> <actor> <note> 10// resume <journal> -> in-flight (KICKOFF w/o DONE) + last checkpoint = resume list 11// board <journal> [window_sec] -> per-ws latest state + liveness (ACTIVE/STALE/DONE) 12// selftest <journal> -> gate T1..T6 (caller pre-cleans the path); exit 0 iff all pass 13import "nx_syscalls.nx" 14import "nx_itoa_lib.nx" // shared MSB-first emitter (zero-alloc) 15const K_MAGIC_4096: i64 = 4096 16const K_MAGIC_262144: i64 = 262144 17const K_MAGIC_262140: i64 = 262140 18const K_MAGIC_3600: i64 = 3600 19 20func ks_puts(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} sys_write(1,s,n); return 0 } 21// MIGRATED to the shared emitter (debt 1785563586). The old body mmapped a scratch buffer 22// per call and never freed it. At PAGE granularity that is 4096B leaked PER CALL -- the 23// defect that took 28.5GB of a 36GB host in nx_ts_lumadiff (2MB input, ~3.66M calls). 24// nxi_* is MSB-first, allocates NOTHING, and emits identical bytes including the sign. 25func ks_putn(v: i64) -> i64 { nxi_out(v); return 0 } 26func ks_cat(d: *u8, o: i64, s: *u8) -> i64 { var i: i64=0; while s[i]!=(0 as u8){ d[o]=s[i]; o=o+1; i=i+1 } return o } 27func ks_catn(d: *u8, o: i64, v: i64) -> i64 { let t: *u8=sys_mmap(28); var m: i64=v; if m<0{m=0} var k: i64=0; if m==0{t[0]=48 as u8;k=1} while m>0{t[k]=(48+(m%10)) as u8;m=m/10;k=k+1} var i: i64=0; while i<k{d[o]=t[k-1-i];o=o+1;i=i+1} return o } 28func ks_vlen(s: *u8) -> i64 { var n: i64=0; while s[n]!=(0 as u8){n=n+1} return n } 29func ks_atoi_z(s: *u8) -> i64 { var v: i64=0; var i: i64=0; while s[i]!=(0 as u8){ let c: i64=s[i] as i64; if c>=48 { if c<=57 { v=v*10+(c-48) } } i=i+1 } return v } 30 31// append one frame: <ts>\t<verb>\t<ws>\t<actor>\t<note>\n (ts<0 => real clock) 32func ks_append(journal: *u8, ts: i64, verb: *u8, ws: *u8, actor: *u8, note: *u8) -> i64 { 33 var t: i64 = ts 34 if t < 0 { t = sys_now_realtime_sec() } 35 let ln: *u8 = sys_mmap(K_MAGIC_4096) 36 var o: i64 = 0 37 o = ks_catn(ln, o, t) 38 ln[o]=9 as u8; o=o+1 39 o = ks_cat(ln, o, verb) 40 ln[o]=9 as u8; o=o+1 41 o = ks_cat(ln, o, ws) 42 ln[o]=9 as u8; o=o+1 43 o = ks_cat(ln, o, actor) 44 ln[o]=9 as u8; o=o+1 45 o = ks_cat(ln, o, note) 46 ln[o]=10 as u8; o=o+1 47 let fd: i64 = sys_openat_append(journal, 0x1a4) 48 if fd < 0 { return -1 } 49 sys_write(fd, ln, o) 50 sys_close(fd) 51 return 0 52} 53func ks_read(path: *u8, buf: *u8, cap: i64) -> i64 { 54 let fd: i64 = sys_openat_rd(path) 55 if fd < 0 { return 0 } 56 var n: i64 = 0; var go: i64 = 1 57 while go == 1 { let r: i64 = sys_read(fd, ((buf as i64)+n) as *u8, cap-n); if r <= 0 { go = 0 } else { n = n + r } if n >= cap { go = 0 } } 58 sys_close(fd) 59 return n 60} 61// span [out0,out1) of column c in line [ls,le); 1 if found 62func ks_col(q: *u8, ls: i64, le: i64, c: i64, out: *i64) -> i64 { 63 var col: i64 = 0 64 var p: i64 = ls 65 while col < c { 66 var s: i64 = 1 67 while s == 1 { if p >= le { return 0 } if q[p]==(9 as u8) { s = 0 } else { p = p+1 } } 68 p = p + 1 69 col = col + 1 70 } 71 var e: i64 = p 72 var s2: i64 = 1 73 while s2 == 1 { if e >= le { s2 = 0 } else { if q[e]==(9 as u8) { s2 = 0 } else { e = e+1 } } } 74 out[0] = p; out[1] = e 75 return 1 76} 77func ks_lit_eq(q: *u8, s: i64, e: i64, lit: *u8) -> i64 { 78 var i: i64 = 0 79 while s+i < e { if lit[i]==(0 as u8) { return 0 } if q[s+i]!=lit[i] { return 0 } i=i+1 } 80 if lit[i]!=(0 as u8) { return 0 } 81 return 1 82} 83func ks_span_eq(q: *u8, s1: i64, e1: i64, s2: i64, e2: i64) -> i64 { 84 if e1-s1 != e2-s2 { return 0 } 85 var i: i64 = 0 86 while s1+i < e1 { if q[s1+i]!=q[s2+i] { return 0 } i=i+1 } 87 return 1 88} 89func ks_atoi(q: *u8, s: i64, e: i64) -> i64 { 90 var v: i64 = 0; var i: i64 = s 91 while i < e { let c: i64 = q[i] as i64; if c>=48 { if c<=57 { v = v*10 + (c-48) } } i=i+1 } 92 return v 93} 94// line end from i 95func ks_le(q: *u8, i: i64, n: i64) -> i64 { 96 var le: i64 = i; var s: i64 = 1 97 while s==1 { if le>=n { s=0 } else { if q[le]==(10 as u8){s=0} else {le=le+1} } } 98 return le 99} 100// any frame verb==VERB (literal) and col2 span == [ws_s,ws_e) ? 101func ks_has(q: *u8, n: i64, verb: *u8, ws_s: i64, ws_e: i64) -> i64 { 102 let cv: *i64 = sys_mmap(16) as *i64 103 let cw: *i64 = sys_mmap(16) as *i64 104 var i: i64 = 0 105 while i < n { 106 let le: i64 = ks_le(q,i,n) 107 if ks_col(q,i,le,1,cv)==1 { if ks_lit_eq(q,cv[0],cv[1],verb)==1 { 108 if ks_col(q,i,le,2,cw)==1 { if ks_span_eq(q,cw[0],cw[1],ws_s,ws_e)==1 { return 1 } } 109 } } 110 i = le + 1 111 } 112 return 0 113} 114// any frame verb==VERB (literal) and col2 == ws (literal) ? (for selftest) 115func ks_has_lit(q: *u8, n: i64, verb: *u8, ws: *u8) -> i64 { 116 let cv: *i64 = sys_mmap(16) as *i64 117 let cw: *i64 = sys_mmap(16) as *i64 118 var i: i64 = 0 119 while i < n { 120 let le: i64 = ks_le(q,i,n) 121 if ks_col(q,i,le,1,cv)==1 { if ks_lit_eq(q,cv[0],cv[1],verb)==1 { 122 if ks_col(q,i,le,2,cw)==1 { if ks_lit_eq(q,cw[0],cw[1],ws)==1 { return 1 } } 123 } } 124 i = le + 1 125 } 126 return 0 127} 128// first KICKOFF for ws span at line-start upto? (no earlier KICKOFF same ws in [0,upto)) 129func ks_first_kick(q: *u8, upto: i64, ws_s: i64, ws_e: i64) -> i64 { 130 let cv: *i64 = sys_mmap(16) as *i64 131 let cw: *i64 = sys_mmap(16) as *i64 132 var i: i64 = 0 133 while i < upto { 134 let le: i64 = ks_le(q,i,upto) 135 if ks_col(q,i,le,1,cv)==1 { if ks_lit_eq(q,cv[0],cv[1],"KICKOFF" as *u8)==1 { 136 if ks_col(q,i,le,2,cw)==1 { if ks_span_eq(q,cw[0],cw[1],ws_s,ws_e)==1 { return 0 } } 137 } } 138 i = le + 1 139 } 140 return 1 141} 142// ts of the LAST (file-order) frame for ws span; -1 if none 143func ks_last_ts(q: *u8, n: i64, ws_s: i64, ws_e: i64) -> i64 { 144 let cw: *i64 = sys_mmap(16) as *i64 145 let c0: *i64 = sys_mmap(16) as *i64 146 var last: i64 = -1 147 var i: i64 = 0 148 while i < n { 149 let le: i64 = ks_le(q,i,n) 150 if ks_col(q,i,le,2,cw)==1 { if ks_span_eq(q,cw[0],cw[1],ws_s,ws_e)==1 { 151 if ks_col(q,i,le,0,c0)==1 { last = ks_atoi(q,c0[0],c0[1]) } 152 } } 153 i = le + 1 154 } 155 return last 156} 157// print the last checkpoint (note col4) of the last frame for ws span 158func ks_print_last_note(q: *u8, n: i64, ws_s: i64, ws_e: i64) -> i64 { 159 let cw: *i64 = sys_mmap(16) as *i64 160 let cn: *i64 = sys_mmap(16) as *i64 161 var ns: i64 = -1; var ne: i64 = -1 162 var i: i64 = 0 163 while i < n { 164 let le: i64 = ks_le(q,i,n) 165 if ks_col(q,i,le,2,cw)==1 { if ks_span_eq(q,cw[0],cw[1],ws_s,ws_e)==1 { 166 if ks_col(q,i,le,4,cn)==1 { ns=cn[0]; ne=cn[1] } 167 } } 168 i = le + 1 169 } 170 if ns>=0 { sys_write(1,(q as i64 + ns) as *u8, ne-ns) } 171 return 0 172} 173// RESUME: print in-flight workstreams (KICKOFF w/o DONE) + last checkpoint; returns count 174func ks_resume(q: *u8, n: i64) -> i64 { 175 let cv: *i64 = sys_mmap(16) as *i64 176 let cw: *i64 = sys_mmap(16) as *i64 177 var inflight: i64 = 0 178 var i: i64 = 0 179 while i < n { 180 let le: i64 = ks_le(q,i,n) 181 if ks_col(q,i,le,1,cv)==1 { if ks_lit_eq(q,cv[0],cv[1],"KICKOFF" as *u8)==1 { 182 if ks_col(q,i,le,2,cw)==1 { if ks_first_kick(q,i,cw[0],cw[1])==1 { 183 if ks_has(q,n,"DONE" as *u8,cw[0],cw[1])==0 { 184 inflight = inflight + 1 185 ks_puts(" IN-FLIGHT ws=" as *u8); sys_write(1,(q as i64 + cw[0]) as *u8, cw[1]-cw[0]) 186 ks_puts(" @" as *u8); ks_putn(ks_last_ts(q,n,cw[0],cw[1])) 187 ks_puts(" last=" as *u8); ks_print_last_note(q,n,cw[0],cw[1]) 188 ks_puts("\n" as *u8) 189 } 190 } } 191 } } 192 i = le + 1 193 } 194 return inflight 195} 196// BOARD: distinct ws -> state + liveness 197func ks_board(q: *u8, n: i64, window: i64) -> i64 { 198 let cv: *i64 = sys_mmap(16) as *i64 199 let cw: *i64 = sys_mmap(16) as *i64 200 let now: i64 = sys_now_realtime_sec() 201 var i: i64 = 0 202 while i < n { 203 let le: i64 = ks_le(q,i,n) 204 if ks_col(q,i,le,1,cv)==1 { if ks_lit_eq(q,cv[0],cv[1],"KICKOFF" as *u8)==1 { 205 if ks_col(q,i,le,2,cw)==1 { if ks_first_kick(q,i,cw[0],cw[1])==1 { 206 ks_puts(" ws=" as *u8); sys_write(1,(q as i64 + cw[0]) as *u8, cw[1]-cw[0]) 207 if ks_has(q,n,"DONE" as *u8,cw[0],cw[1])==1 { ks_puts(" DONE" as *u8) } 208 else { 209 let lt: i64 = ks_last_ts(q,n,cw[0],cw[1]) 210 if now - lt <= window { ks_puts(" ACTIVE" as *u8) } else { ks_puts(" STALE" as *u8) } 211 ks_puts(" last=" as *u8); ks_print_last_note(q,n,cw[0],cw[1]) 212 } 213 ks_puts("\n" as *u8) 214 } } 215 } } 216 i = le + 1 217 } 218 return 0 219} 220func ks_assert(name: *u8, got: i64, want: i64, fails: *i64) -> i64 { 221 if got==want { ks_puts(" PASS " as *u8) } else { ks_puts(" FAIL " as *u8); fails[0]=fails[0]+1 } 222 ks_puts(name); ks_puts(" got=" as *u8); ks_putn(got); ks_puts(" want=" as *u8); ks_putn(want); ks_puts("\n" as *u8) 223 return 0 224} 225// inflight-count via a fresh read 226func ks_inflight(journal: *u8) -> i64 { 227 let q: *u8 = sys_mmap(K_MAGIC_262144) 228 let n: i64 = ks_read(journal, q, K_MAGIC_262140) 229 let cv: *i64 = sys_mmap(16) as *i64 230 let cw: *i64 = sys_mmap(16) as *i64 231 var c: i64 = 0 232 var i: i64 = 0 233 while i < n { 234 let le: i64 = ks_le(q,i,n) 235 if ks_col(q,i,le,1,cv)==1 { if ks_lit_eq(q,cv[0],cv[1],"KICKOFF" as *u8)==1 { 236 if ks_col(q,i,le,2,cw)==1 { if ks_first_kick(q,i,cw[0],cw[1])==1 { 237 if ks_has(q,n,"DONE" as *u8,cw[0],cw[1])==0 { c=c+1 } 238 } } 239 } } 240 i = le + 1 241 } 242 return c 243} 244func ks_selftest(journal: *u8) -> i64 { 245 let fails: *i64 = sys_mmap(16) as *i64 246 fails[0]=0 247 ks_puts("== nx_ws_kickoff_sync selftest ==\n" as *u8) 248 // T6 neg: empty journal (caller pre-cleaned) -> 0 in-flight, no fabrication 249 ks_assert("T6-empty-no-fabrication" as *u8, ks_inflight(journal), 0, fails) 250 // T1 kickoff A -> in-flight 1 251 ks_append(journal, -1, "KICKOFF" as *u8, "A" as *u8, "s1" as *u8, "start-A" as *u8) 252 ks_assert("T1-kickoff-tracked" as *u8, ks_inflight(journal), 1, fails) 253 // T2 beat A -> still 1, last note advances (checked via has) 254 ks_append(journal, -1, "BEAT" as *u8, "A" as *u8, "s1" as *u8, "checkpoint-A1" as *u8) 255 ks_assert("T2-beat-keeps-inflight" as *u8, ks_inflight(journal), 1, fails) 256 // T5 kickoff B -> in-flight 2 (two concurrent, no clobber) 257 ks_append(journal, -1, "KICKOFF" as *u8, "B" as *u8, "s2" as *u8, "start-B" as *u8) 258 ks_assert("T5-two-concurrent" as *u8, ks_inflight(journal), 2, fails) 259 // T3/T4 done A -> in-flight 1 (B remains = the CRASH-RESUME property: B kicked-off, never done, still resumable) 260 ks_append(journal, -1, "DONE" as *u8, "A" as *u8, "s1" as *u8, "finished-A" as *u8) 261 ks_assert("T3-done-clears" as *u8, ks_inflight(journal), 1, fails) 262 // load-bearing: A must be gone, B must remain 263 let q: *u8 = sys_mmap(K_MAGIC_262144) 264 let n: i64 = ks_read(journal, q, K_MAGIC_262140) 265 var bfail: i64 = 0 266 if ks_has_lit(q,n,"DONE" as *u8,"A" as *u8)==0 { bfail=1 } 267 if ks_has_lit(q,n,"KICKOFF" as *u8,"B" as *u8)==0 { bfail=1 } 268 if ks_has_lit(q,n,"DONE" as *u8,"B" as *u8)==1 { bfail=1 } 269 ks_assert("T4-crash-resume-B-survives" as *u8, bfail, 0, fails) 270 ks_puts("== RESUME view ==\n" as *u8) 271 ks_resume(q, n) 272 if fails[0]==0 { ks_puts("VERDICT=GREEN (6/6)\n" as *u8); return 0 } 273 ks_puts("VERDICT=RED fails=" as *u8); ks_putn(fails[0]); ks_puts("\n" as *u8) 274 return 1 275} 276func main(argc: i64, argv: *i64) -> i64 { 277 if argc < 3 { ks_puts("usage: nx_ws_kickoff_sync {kickoff|beat|done <journal> <ws> <actor> <note> | resume|board <journal> [window] | selftest <journal>}\n" as *u8); sys_exit(2); return 2 } 278 let verb: *u8 = argv[1] as *u8 279 let journal: *u8 = argv[2] as *u8 280 let vv: *u8 = verb 281 if ks_lit_eq(vv,0,ks_vlen(vv),"selftest" as *u8)==1 { sys_exit(ks_selftest(journal)); return 0 } 282 if ks_lit_eq(vv,0,ks_vlen(vv),"resume" as *u8)==1 { 283 let q: *u8 = sys_mmap(K_MAGIC_262144); let n: i64 = ks_read(journal,q,K_MAGIC_262140) 284 let c: i64 = ks_resume(q,n); ks_puts(" RESUME in-flight=" as *u8); ks_putn(c); ks_puts("\n" as *u8); sys_exit(0); return 0 285 } 286 if ks_lit_eq(vv,0,ks_vlen(vv),"board" as *u8)==1 { 287 var w: i64 = K_MAGIC_3600 288 if argc>=4 { w = ks_atoi_z(argv[3] as *u8) } 289 let q: *u8 = sys_mmap(K_MAGIC_262144); let n: i64 = ks_read(journal,q,K_MAGIC_262140) 290 ks_board(q,n,w); sys_exit(0); return 0 291 } 292 // append verbs 293 if argc < 6 { ks_puts("kickoff|beat|done need <journal> <ws> <actor> <note>\n" as *u8); sys_exit(2); return 2 } 294 let ws: *u8 = argv[3] as *u8 295 let actor: *u8 = argv[4] as *u8 296 let note: *u8 = argv[5] as *u8 297 var vtag: *u8 = "BEAT" as *u8 298 if ks_lit_eq(vv,0,ks_vlen(vv),"kickoff" as *u8)==1 { vtag = "KICKOFF" as *u8 } 299 if ks_lit_eq(vv,0,ks_vlen(vv),"done" as *u8)==1 { vtag = "DONE" as *u8 } 300 let rc: i64 = ks_append(journal, -1, vtag, ws, actor, note) 301 if rc==0 { ks_puts(" SYNC " as *u8); ks_puts(vtag); ks_puts(" ws=" as *u8); ks_puts(ws); ks_puts("\n" as *u8); sys_exit(0); return 0 } 302 ks_puts(" append FAILED\n" as *u8); sys_exit(1); return 1 303}