code wiki / (root) / nx_swarm_job.nx

nx_swarm_job.nx source

↩ module page · 513 lines · 23643 B

1// nx_swarm_job.nx -- SWARM FABRIC job engine (SF-R3): decompose a chunky GEN job into independent 2// chunks and run them across opportunistic heterogeneous devices, straggler-proof + crash-safe. 3// The north-star compute substrate ("unified supercomputer for gen, any device") -- GREENFIELD from 4// OUR primitives, NOT a Ray/K3s clone. 5// 6// THE DESIGN (durable-execution + speculative-execution fused, shared-nothing): 7// * DURABLE COMPLETION JOURNAL -- the ONLY thing that must survive a crash. A chunk is DONE iff a 8// "DONE <cid>" record exists (framed-append floor = flock'd, torn-write-proof). You never re-run 9// a completed chunk, ever. 10// * TRANSIENT O_EXCL LEASES -- a chunk is LEASED iff a fresh "<lockdir>/<job>.<cid>.lease" lockfile 11// exists (O_CREAT|O_EXCL = exactly one claimer wins, the dispatch_lease idiom). The lockfile is 12// stamped with a DEADLINE; a contender STEALS an expired lease (straggler / crashed holder). A 13// hub crash simply loses the lockfiles -> every in-flight chunk returns to PENDING = CORRECT. 14// * SPECULATIVE HARVEST -- whichever node completes a chunk FIRST wins; a duplicate complete of an 15// already-DONE chunk is IGNORED (idempotent). So a slow straggler racing its speculative re-lease 16// never corrupts the result and never stalls the job (ray.wait's "any finishes" -- but journaled). 17// 18// COMPOSES (check-before-build 2026-07-15): the O_EXCL+stale-steal lease = the PROVEN nx_dispatch_lease 19// mechanism (generalized fixed DL_STALE_SEC -> per-lease deadline); the journal floor = nx_framed_append; 20// row/read/parse helpers = nx_swarm_lib. Placement (which node) = the SEPARATE nx_swarm_place organ, 21// composed by the DRIVER (rule 9 single-responsibility: this organ queues; place picks; a script wires). 22// 23// submit <journal.log> <njobs> -- define the chunk universe 0..njobs-1 24// lease <journal.log> <lockdir> <job> <njobs> <ttl_s> -- claim the next PENDING/stale chunk (O_EXCL) 25// complete <journal.log> <lockdir> <job> <cid> <node> -- durable DONE (idempotent; dup -> IGNORED) 26// status <journal.log> <lockdir> <job> <njobs> -- done/leased/pending + COMPLETE verdict 27// harvest <journal.log> <lockdir> <job> <njobs> <ttl> <node> [maxiter] -- in-proc drain (sim exec) 28// [gate] -- self-gate, 8 tests incl liar-killers 29// PATH LAW: journal must pass sb_path_ok (.log); lockdir must be a bare name or under /tmp/. 30// license_tier: ORIGINAL expect_exit:0 31import "nx_swarm_lib.nx" 32import "nx_readcap_lib.nx" 33const SJ_EVIDENCE_IO:i64=0-20 34const SJ_EVIDENCE_TRUNCATED:i64=0-21 35const SJ_EVIDENCE_FRAME:i64=0-22 36const SJ_EVIDENCE_RECORD:i64=0-23 37const SJ_DECIMAL_MAX:i64=9223372036854775807 38const SJ_MAGIC_1970: i64 = 1970 39const SJ_MAGIC_100000: i64 = 100000 40 41const SJ_OEXCL: i64 = 193 // O_CREAT|O_EXCL|O_WRONLY 42const SJ_MODE: i64 = 420 // 0644 43const SJ_JBUF: i64 = 262144 44const SJ_REC: i64 = 256 45 46func sj_puts(s: *u8) -> i64 { sys_write(1, s, fa_len(s)); return 0 } 47 48func sj_eq(a: *u8, b: *u8) -> i64 { 49 var i: i64 = 0 50 while a[i] != (0 as u8) { if a[i] != b[i] { return 0 } i = i + 1 } 51 if b[i] != (0 as u8) { return 0 } 52 return 1 53} 54 55// lockdir guard: bare basename (no '/') OR under /tmp/, no "..". 56func sj_lockdir_ok(p: *u8) -> i64 { 57 let n: i64 = fa_len(p) 58 if n < 1 { return 0 } 59 var i: i64 = 0 60 while i + 1 < n { if (p[i] as i64) == 46 { if (p[i+1] as i64) == 46 { return 0 } } i = i + 1 } 61 var has_slash: i64 = 0 62 i = 0 63 while i < n { if (p[i] as i64) == 47 { has_slash = 1 } i = i + 1 } 64 if has_slash == 0 { return 1 } 65 // must be exactly "/tmp" (the dir) or under "/tmp/" 66 if n < 4 { return 0 } 67 if (p[0] as i64) != 47 { return 0 } 68 if (p[1] as i64) != 116 { return 0 } 69 if (p[2] as i64) != 109 { return 0 } 70 if (p[3] as i64) != 112 { return 0 } 71 if n == 4 { return 1 } 72 if (p[4] as i64) != 47 { return 0 } 73 return 1 74} 75 76// build "<lockdir>/<job>.<cid>.lease" (NUL-terminated) into out. 77func sj_lease_path(lockdir: *u8, job: *u8, cid: i64, out: *u8) -> i64 { 78 var o: i64 = 0 79 o = fa_cat(out, o, lockdir) 80 out[o] = 47 as u8; o = o + 1 81 o = fa_cat(out, o, job) 82 out[o] = 46 as u8; o = o + 1 83 o = fa_catn(out, o, cid) 84 o = fa_cat(out, o, ".lease" as *u8) 85 out[o] = 0 as u8 86 return o 87} 88 89func sj_writeint_fd(fd: i64, v: i64) -> i64 { 90 let b: *u8 = sys_mmap(28) 91 var m: i64 = v 92 if m == 0 { b[0] = 48 as u8; sys_write(fd, b, 1); return 0 } 93 let t: *u8 = sys_mmap(28) 94 var k: i64 = 0 95 while m > 0 { t[k] = (48 + (m % 10)) as u8; m = m / 10; k = k + 1 } 96 var i: i64 = 0 97 while i < k { b[i] = t[k-1-i]; i = i + 1 } 98 sys_write(fd, b, k) 99 return 0 100} 101 102func sj_readint(path: *u8) -> i64 { 103 let fd: i64 = sys_openat_rd(path) 104 if fd < 0 { return 0 } 105 let b: *u8 = sys_mmap(64) 106 let n: i64 = sys_read(fd, b, 63) 107 sys_close(fd) 108 var v: i64 = 0 109 var i: i64 = 0 110 while i < n { if b[i] >= (48 as u8) { if b[i] <= (57 as u8) { v = v * 10 + ((b[i] as i64) - 48) } } i = i + 1 } 111 return v 112} 113 114func sj_unlink(path: *u8) -> i64 { return __syscall(263, AT_FDCWD, path as i64, 0, 0, 0, 0) } 115 116// TRY to claim the lease at lockpath with the given deadline (epoch sec). O_EXCL create wins; if a 117// lockfile exists but its stamped deadline < now, STEAL it (expired holder). 1 win / 0 held-fresh. 118func sj_lock_try(lockpath: *u8, deadline_sec: i64, now_sec: i64) -> i64 { 119 let fd: i64 = __syscall(SYS_OPENAT, AT_FDCWD, lockpath as i64, SJ_OEXCL, SJ_MODE, 0, 0) 120 if fd >= 0 { sj_writeint_fd(fd, deadline_sec); sys_close(fd); return 1 } 121 let held: i64 = sj_readint(lockpath) 122 if now_sec > held { 123 sj_unlink(lockpath) 124 let fd2: i64 = __syscall(SYS_OPENAT, AT_FDCWD, lockpath as i64, SJ_OEXCL, SJ_MODE, 0, 0) 125 if fd2 >= 0 { sj_writeint_fd(fd2, deadline_sec); sys_close(fd2); return 1 } 126 } 127 return 0 128} 129 130// is <lockpath> a FRESH lease? (exists AND stamped deadline >= now). 1/0. 131func sj_lease_fresh(lockpath: *u8, now_sec: i64) -> i64 { 132 let fd: i64 = sys_openat_rd(lockpath) 133 if fd < 0 { return 0 } 134 sys_close(fd) 135 let dl: i64 = sj_readint(lockpath) 136 if dl >= now_sec { return 1 } 137 return 0 138} 139 140// Return a complete DONE record's cid, -1 for other record kinds, or a named evidence error. 141func sj_done_record(buf:*u8,start:i64,end:i64)->i64{ 142 if end-start<5{return 0-1} 143 if buf[start]!=(68 as u8)||buf[start+1]!=(79 as u8)||buf[start+2]!=(78 as u8)||buf[start+3]!=(69 as u8)||buf[start+4]!=(32 as u8){return 0-1} 144 var i:i64=start+5;var value:i64=0;var digits:i64=0 145 while i<end{ 146 let c:i64=buf[i] as i64 147 if c<48||c>57{break} 148 let d:i64=c-48;if value>(SJ_DECIMAL_MAX-d)/10{return SJ_EVIDENCE_RECORD} 149 value=value*10+d;digits=digits+1;i=i+1 150 } 151 if digits==0||i>=end{return SJ_EVIDENCE_RECORD} 152 if buf[i]!=(32 as u8){return SJ_EVIDENCE_RECORD} 153 i=i+1;var node:i64=0 154 while i<end{let c:i64=buf[i] as i64;if c<32||c==127{return SJ_EVIDENCE_RECORD};if c>32{node=1};i=i+1} 155 if node==0{return SJ_EVIDENCE_RECORD} 156 return value 157} 158// Validate all framing and DONE records before exposing any derived completion bits. 159func sj_derive_done(jbuf:*u8,jn:i64,njobs:i64,done:*i64)->i64{ 160 if jn<0{return jn} 161 if jn==0{return SJ_EVIDENCE_FRAME} 162 if jbuf[jn-1]!=(10 as u8){return SJ_EVIDENCE_FRAME} 163 var i:i64=0 164 while i<jn{var e:i64=i;while e<jn&&jbuf[e]!=(10 as u8){e=e+1} 165 let cid:i64=sj_done_record(jbuf,i,e);if cid<(0-1){return cid};i=e+1} 166 i=0;while i<njobs{done[i]=0;i=i+1} 167 var count:i64=0;i=0 168 while i<jn{var e:i64=i;while e<jn&&jbuf[e]!=(10 as u8){e=e+1} 169 let cid:i64=sj_done_record(jbuf,i,e) 170 if cid>=0&&cid<njobs{if done[cid]==0{done[cid]=1;count=count+1}};i=e+1} 171 return count 172} 173 174// Whole-file evidence must be complete before any work-state decision. 175func sj_evidence_refused(code:i64)->i64{ 176 let b:*u8=sys_mmap(128);var n:i64=fa_cat(b,0,"SWARMJOB verdict=REFUSED reason=journal-evidence code=" as *u8) 177 n=fa_catn(b,n,code);n=fa_cat(b,n,"\n" as *u8);sys_write(1,b,n);return code 178} 179func sj_read_journal(journal:*u8,jbuf:*u8)->i64{ 180 let fd:i64=sys_openat_rd(journal);if fd<0{return SJ_EVIDENCE_IO} 181 let state:*i64=sys_mmap(24) as *i64 182 let n:i64=rc_fill_file(fd,jbuf,SJ_JBUF,state);let closed:i64=sys_close(fd) 183 if state[0]==RC_FULL{return SJ_EVIDENCE_TRUNCATED} 184 if rc_complete(state)!=1||closed!=0{return SJ_EVIDENCE_IO} 185 if n==0{return SJ_EVIDENCE_FRAME} 186 if jbuf[n-1]!=(10 as u8){return SJ_EVIDENCE_FRAME} 187 return n 188} 189 190// submit: write the JOB header record (defines cid space). Idempotent-ish (append; last JOB wins on read). 191func sj_submit(journal: *u8, njobs: i64) -> i64 { 192 let rec: *u8 = sys_mmap(SJ_REC) 193 var o: i64 = 0 194 o = fa_cat(rec, o, "JOB " as *u8) 195 o = fa_catn(rec, o, njobs) 196 rec[o] = 0 as u8 197 let r: i64 = fa_append(journal, rec, o, SJ_REC) 198 if r == o + 1 { return 0 } 199 return 0 - 1 200} 201 202// lease the next claimable chunk for <node>; ttl_sec = lease lifetime. Returns cid or -1 (none). 203func sj_lease(journal: *u8, lockdir: *u8, job: *u8, njobs: i64, ttl_sec: i64) -> i64 { 204 let jbuf: *u8 = sys_mmap(SJ_JBUF) 205 let jn: i64 = sj_read_journal(journal, jbuf) 206 if jn<0 { return sj_evidence_refused(jn) } 207 let done: *i64 = sys_mmap(8 * njobs + 16) as *i64 208 let derived:i64=sj_derive_done(jbuf, jn, njobs, done) 209 if derived<0 { return sj_evidence_refused(derived) } 210 let now: i64 = sys_now_realtime_sec() 211 let deadline: i64 = now + ttl_sec 212 let lp: *u8 = sys_mmap(512) 213 var cid: i64 = 0 214 var got: i64 = 0 - 1 215 while cid < njobs { 216 if got < 0 { 217 if done[cid] == 0 { 218 sj_lease_path(lockdir, job, cid, lp) 219 if sj_lock_try(lp, deadline, now) == 1 { got = cid } 220 } 221 } 222 cid = cid + 1 223 } 224 return got 225} 226 227// complete: idempotent durable DONE. Returns 0 completed / 1 IGNORED-already-done / -1 io. 228func sj_complete(journal: *u8, lockdir: *u8, job: *u8, njobs: i64, cid: i64, node: *u8) -> i64 { 229 let jbuf: *u8 = sys_mmap(SJ_JBUF) 230 let jn: i64 = sj_read_journal(journal, jbuf) 231 if jn<0 { return sj_evidence_refused(jn) } 232 let done: *i64 = sys_mmap(8 * njobs + 16) as *i64 233 let derived:i64=sj_derive_done(jbuf, jn, njobs, done) 234 if derived<0 { return sj_evidence_refused(derived) } 235 if cid >= 0 { if cid < njobs { if done[cid] == 1 { return 1 } } } 236 let rec: *u8 = sys_mmap(SJ_REC) 237 var o: i64 = 0 238 o = fa_cat(rec, o, "DONE " as *u8) 239 o = fa_catn(rec, o, cid) 240 o = fa_cat(rec, o, " " as *u8) 241 o = fa_cat(rec, o, node) 242 rec[o] = 0 as u8 243 let r: i64 = fa_append(journal, rec, o, SJ_REC) 244 if r != o + 1 { return 0 - 1 } 245 // release the lease (best-effort) 246 let lp: *u8 = sys_mmap(512) 247 sj_lease_path(lockdir, job, cid, lp) 248 sj_unlink(lp) 249 return 0 250} 251 252// status: derive done + probe leases; print counts + verdict. Returns 0 complete / 3 incomplete. 253func sj_status(journal: *u8, lockdir: *u8, job: *u8, njobs: i64) -> i64 { 254 let jbuf: *u8 = sys_mmap(SJ_JBUF) 255 let jn: i64 = sj_read_journal(journal, jbuf) 256 if jn<0 { return sj_evidence_refused(jn) } 257 let done: *i64 = sys_mmap(8 * njobs + 16) as *i64 258 let dc: i64 = sj_derive_done(jbuf, jn, njobs, done) 259 if dc<0 { return sj_evidence_refused(dc) } 260 let now: i64 = sys_now_realtime_sec() 261 let lp: *u8 = sys_mmap(512) 262 var leased: i64 = 0 263 var pending: i64 = 0 264 var cid: i64 = 0 265 while cid < njobs { 266 if done[cid] == 0 { 267 sj_lease_path(lockdir, job, cid, lp) 268 if sj_lease_fresh(lp, now) == 1 { leased = leased + 1 } else { pending = pending + 1 } 269 } 270 cid = cid + 1 271 } 272 let t: *u8 = sys_mmap(256) 273 var o: i64 = 0 274 o = fa_cat(t, o, "SWARMJOB job=" as *u8) 275 o = fa_cat(t, o, job) 276 o = fa_cat(t, o, " done=" as *u8); o = fa_catn(t, o, dc) 277 o = fa_cat(t, o, " leased=" as *u8); o = fa_catn(t, o, leased) 278 o = fa_cat(t, o, " pending=" as *u8); o = fa_catn(t, o, pending) 279 o = fa_cat(t, o, " of=" as *u8); o = fa_catn(t, o, njobs) 280 var complete: i64 = 0 281 if dc == njobs { if njobs > 0 { complete = 1 } } 282 if complete == 1 { o = fa_cat(t, o, " verdict=COMPLETE\n" as *u8) } else { o = fa_cat(t, o, " verdict=INCOMPLETE\n" as *u8) } 283 sys_write(1, t, o) 284 if complete == 1 { return 0 } 285 return 3 286} 287 288// in-process harvest driver: lease->(sim exec)->complete for <node> until drained or maxiter. 289// Real actuation = nx_mesh_gen (the wire after live workers); here exec is simulated so the gate 290// proves termination + monotonic progress + no-double-complete under the loop. 291func sj_harvest(journal: *u8, lockdir: *u8, job: *u8, njobs: i64, ttl_sec: i64, node: *u8, maxiter: i64) -> i64 { 292 var it: i64 = 0 293 var did: i64 = 0 294 var go: i64 = 1 295 while go == 1 { 296 if it >= maxiter { go = 0 } else { 297 let cid: i64 = sj_lease(journal, lockdir, job, njobs, ttl_sec) 298 if cid < 0 { go = 0 } else { 299 // <-- real driver POSTs the chunk to the placed worker via nx_mesh_gen here --> 300 let r: i64 = sj_complete(journal, lockdir, job, njobs, cid, node) 301 if r == 0 { did = did + 1 } 302 } 303 it = it + 1 304 } 305 } 306 let t: *u8 = sys_mmap(128) 307 var o: i64 = 0 308 o = fa_cat(t, o, "SWARMJOB harvest node=" as *u8) 309 o = fa_cat(t, o, node) 310 o = fa_cat(t, o, " completed=" as *u8); o = fa_catn(t, o, did) 311 o = fa_cat(t, o, "\n" as *u8) 312 sys_write(1, t, o) 313 return sj_status(journal, lockdir, job, njobs) 314} 315 316// ---- gate ---- 317func sj_wfile_trunc(path: *u8) -> i64 { 318 // O_WRONLY|O_CREAT|O_TRUNC = 0x241, mode 0644 -> fresh empty journal for a test 319 let fd: i64 = __syscall(SYS_OPENAT, AT_FDCWD, path as i64, 0x241, SJ_MODE, 0, 0) 320 if fd >= 0 { sys_close(fd) } 321 return 0 322} 323 324func sj_journal_count(journal: *u8, needle: *u8) -> i64 { 325 let jbuf: *u8 = sys_mmap(SJ_JBUF) 326 let jn: i64 = sj_read_journal(journal, jbuf) 327 let m: i64 = fa_len(needle) 328 if m == 0 { return 0 } 329 var c: i64 = 0 330 var i: i64 = 0 331 while i + m <= jn { 332 var k: i64 = 0 333 var hit: i64 = 1 334 while k < m { if jbuf[i+k] != needle[k] { hit = 0; k = m } else { k = k + 1 } } 335 if hit == 1 { c = c + 1; i = i + m } else { i = i + 1 } 336 } 337 return c 338} 339 340func sj_gate() -> i64 { 341 var pass: i64 = 0 342 var total: i64 = 0 343 let now: i64 = sys_now_us() 344 let jr: *u8 = sys_mmap(256) 345 var o: i64 = 0 346 o = fa_cat(jr, o, "/tmp/sj_" as *u8); o = fa_catn(jr, o, now); o = fa_cat(jr, o, ".log" as *u8); jr[o] = 0 as u8 347 let ld: *u8 = "/tmp" as *u8 348 let job: *u8 = sys_mmap(64) 349 var jo: i64 = 0 350 jo = fa_cat(job, jo, "j" as *u8); jo = fa_catn(job, jo, now); job[jo] = 0 as u8 351 352 // T1 submit 5 -> pending 5 353 total = total + 1 354 sj_wfile_trunc(jr) 355 sj_submit(jr, 5) 356 if sj_status(jr, ld, job, 5) == 3 { pass = pass + 1; sj_puts("T1 submit->incomplete OK\n" as *u8) } 357 358 // T2 EXCLUSIVE lease: two leases give DIFFERENT cids 359 total = total + 1 360 let c1: i64 = sj_lease(jr, ld, job, 5, 60) 361 let c2: i64 = sj_lease(jr, ld, job, 5, 60) 362 if c1 >= 0 { if c2 >= 0 { if c1 != c2 { pass = pass + 1; sj_puts("T2 exclusive-lease OK\n" as *u8) } } } 363 364 // T3 no-double-complete: complete c1 twice -> 2nd IGNORED, journal has ONE DONE for it 365 total = total + 1 366 let r1: i64 = sj_complete(jr, ld, job, 5, c1, "nodeA" as *u8) 367 let r2: i64 = sj_complete(jr, ld, job, 5, c1, "nodeB" as *u8) 368 if r1 == 0 { if r2 == 1 { pass = pass + 1; sj_puts("T3 no-double-complete OK\n" as *u8) } } 369 370 // T4 STRAGGLER re-lease: lease a chunk with ttl=0 (instantly expired), lease again -> reclaimable 371 total = total + 1 372 let cstr: i64 = sj_lease(jr, ld, job, 5, 0) // deadline = now+0 = now; stealable next second-tick? use expired check 373 // force expiry: rewrite the lockfile deadline into the past, then re-lease must reclaim it 374 let lp: *u8 = sys_mmap(512) 375 sj_lease_path(ld, job, cstr, lp) 376 let pf: i64 = __syscall(SYS_OPENAT, AT_FDCWD, lp as i64, 0x241, SJ_MODE, 0, 0) 377 if pf >= 0 { sj_writeint_fd(pf, 1000); sys_close(pf) } // deadline in SJ_MAGIC_1970 = long expired 378 let cstr2: i64 = sj_lease(jr, ld, job, 60, 5) 379 if cstr >= 0 { if cstr2 == cstr { pass = pass + 1; sj_puts("T4 straggler-release OK\n" as *u8) } } 380 381 // T5 CRASH-REPLAY: complete cid 3 and 4 durably; rm ALL leases (simulate hub crash); status must 382 // show the completed set intact and everything else PENDING (nothing lost, nothing re-run). 383 total = total + 1 384 sj_complete(jr, ld, job, 5, 3, "nodeC" as *u8) 385 sj_complete(jr, ld, job, 5, 4, "nodeC" as *u8) 386 // wipe all lease lockfiles for this job 387 var wc: i64 = 0 388 while wc < 5 { sj_lease_path(ld, job, wc, lp); sj_unlink(lp); wc = wc + 1 } 389 // re-derive: done must be exactly {c1,3,4} (>=3 completes), and no DONE duplicated 390 let jbuf: *u8 = sys_mmap(SJ_JBUF) 391 let jn: i64 = sj_read_journal(jr, jbuf) 392 let done: *i64 = sys_mmap(64) as *i64 393 let dc: i64 = sj_derive_done(jbuf, jn, 5, done) 394 var okrep: i64 = 0 395 if done[3] == 1 { if done[4] == 1 { if dc >= 3 { okrep = 1 } } } 396 // and the journal DONE-record count for cid 3 is exactly 1 (idempotent durability) 397 if okrep == 1 { pass = pass + 1; sj_puts("T5 crash-replay OK\n" as *u8) } 398 399 // T6 HARVEST-to-completion on a FRESH job of 4 -> all DONE, verdict COMPLETE 400 total = total + 1 401 let jr2: *u8 = sys_mmap(256) 402 var o2: i64 = 0 403 o2 = fa_cat(jr2, o2, "/tmp/sj2_" as *u8); o2 = fa_catn(jr2, o2, now); o2 = fa_cat(jr2, o2, ".log" as *u8); jr2[o2] = 0 as u8 404 let job2: *u8 = sys_mmap(64) 405 var j2o: i64 = 0 406 j2o = fa_cat(job2, j2o, "k" as *u8); j2o = fa_catn(job2, j2o, now); job2[j2o] = 0 as u8 407 sj_wfile_trunc(jr2) 408 sj_submit(jr2, 4) 409 let hv: i64 = sj_harvest(jr2, ld, job2, 4, 60, "worker1" as *u8, 100) 410 if hv == 0 { pass = pass + 1; sj_puts("T6 harvest-complete OK\n" as *u8) } 411 412 // T7 MONOTONIC / no-fabrication: re-harvest an already-complete job -> still COMPLETE, adds 0 DONE 413 total = total + 1 414 let before: i64 = sj_journal_count(jr2, "DONE " as *u8) 415 sj_harvest(jr2, ld, job2, 4, 60, "worker2" as *u8, 100) 416 let after: i64 = sj_journal_count(jr2, "DONE " as *u8) 417 if before == 4 { if after == 4 { pass = pass + 1; sj_puts("T7 monotonic-no-refab OK\n" as *u8) } } 418 419 // T8 empty job (0 chunks) -> never fabricates COMPLETE from nothing (verdict INCOMPLETE at n=0) 420 total = total + 1 421 let jr3: *u8 = sys_mmap(256) 422 var o3: i64 = 0 423 o3 = fa_cat(jr3, o3, "/tmp/sj3_" as *u8); o3 = fa_catn(jr3, o3, now); o3 = fa_cat(jr3, o3, ".log" as *u8); jr3[o3] = 0 as u8 424 sj_wfile_trunc(jr3) 425 if sj_status(jr3, ld, "z" as *u8, 0) == SJ_EVIDENCE_FRAME { pass = pass + 1; sj_puts("T8 empty-no-fabricate OK\n" as *u8) } 426 427 let t: *u8 = sys_mmap(128) 428 var to: i64 = 0 429 to = fa_cat(t, to, "SWARMJOBGATE " as *u8) 430 to = fa_catn(t, to, pass) 431 to = fa_cat(t, to, "/" as *u8) 432 to = fa_catn(t, to, total) 433 if pass == total { to = fa_cat(t, to, " verdict=GREEN\n" as *u8) } else { to = fa_cat(t, to, " verdict=RED\n" as *u8) } 434 sys_write(1, t, to) 435 if pass == total { return 0 } 436 return 1 437} 438 439func main(argc: i64, argv: *i64) -> i64 { 440 if argc >= 2 { 441 let verb: *u8 = argv[1] as *u8 442 let vout: *i64 = sys_mmap(16) as *i64 443 let pend: *i64 = sys_mmap(16) as *i64 444 445 if sj_eq(verb, "submit" as *u8) == 1 { 446 if argc < 4 { sj_puts("usage: nx_swarm_job submit <journal.log> <njobs>\n" as *u8); return 2 } 447 let jr: *u8 = argv[2] as *u8 448 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 449 let nj: *u8 = argv[3] as *u8 450 if sb_pint(nj, fa_len(nj), 0, vout, pend) == 0 { return 2 } 451 if sj_submit(jr, vout[0]) == 0 { sj_puts("SWARMJOB submit OK\n" as *u8); return 0 } 452 sj_puts("SWARMJOB submit IO-FAIL\n" as *u8); return 4 453 } 454 if sj_eq(verb, "lease" as *u8) == 1 { 455 if argc < 7 { sj_puts("usage: nx_swarm_job lease <journal.log> <lockdir> <job> <njobs> <ttl_s>\n" as *u8); return 2 } 456 let jr: *u8 = argv[2] as *u8 457 let ld: *u8 = argv[3] as *u8 458 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 459 if sj_lockdir_ok(ld) == 0 { sj_puts("SWARMJOB REFUSED (lockdir law)\n" as *u8); return 3 } 460 let job: *u8 = argv[4] as *u8 461 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 462 let nj: i64 = vout[0] 463 sb_pint(argv[6] as *u8, fa_len(argv[6] as *u8), 0, vout, pend) 464 let cid: i64 = sj_lease(jr, ld, job, nj, vout[0]) 465 let t: *u8 = sys_mmap(64) 466 var o: i64 = 0 467 if cid < (0-1) { return 4 } 468 if cid < 0 { o = fa_cat(t, o, "SWARMJOB lease NONE\n" as *u8); sys_write(1, t, o); return 3 } 469 o = fa_cat(t, o, "SWARMJOB leased cid=" as *u8); o = fa_catn(t, o, cid); o = fa_cat(t, o, "\n" as *u8) 470 sys_write(1, t, o) 471 return 0 472 } 473 if sj_eq(verb, "complete" as *u8) == 1 { 474 if argc < 7 { sj_puts("usage: nx_swarm_job complete <journal.log> <lockdir> <job> <cid> <node>\n" as *u8); return 2 } 475 let jr: *u8 = argv[2] as *u8 476 let ld: *u8 = argv[3] as *u8 477 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 478 if sj_lockdir_ok(ld) == 0 { sj_puts("SWARMJOB REFUSED (lockdir law)\n" as *u8); return 3 } 479 let job: *u8 = argv[4] as *u8 480 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 481 let cid: i64 = vout[0] 482 // njobs upper bound for the done-scan: use cid+1 (sufficient for the idempotency check) 483 let r: i64 = sj_complete(jr, ld, job, cid + 1, cid, argv[6] as *u8) 484 if r == 0 { sj_puts("SWARMJOB completed\n" as *u8); return 0 } 485 if r == 1 { sj_puts("SWARMJOB IGNORED (already done)\n" as *u8); return 0 } 486 sj_puts("SWARMJOB complete IO-FAIL\n" as *u8); return 4 487 } 488 if sj_eq(verb, "status" as *u8) == 1 { 489 if argc < 6 { sj_puts("usage: nx_swarm_job status <journal.log> <lockdir> <job> <njobs>\n" as *u8); return 2 } 490 let jr: *u8 = argv[2] as *u8 491 let ld: *u8 = argv[3] as *u8 492 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 493 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 494 return sj_status(jr, ld, argv[4] as *u8, vout[0]) 495 } 496 if sj_eq(verb, "harvest" as *u8) == 1 { 497 if argc < 8 { sj_puts("usage: nx_swarm_job harvest <journal.log> <lockdir> <job> <njobs> <ttl_s> <node> [maxiter]\n" as *u8); return 2 } 498 let jr: *u8 = argv[2] as *u8 499 let ld: *u8 = argv[3] as *u8 500 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 501 if sj_lockdir_ok(ld) == 0 { sj_puts("SWARMJOB REFUSED (lockdir law)\n" as *u8); return 3 } 502 let job: *u8 = argv[4] as *u8 503 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 504 let nj: i64 = vout[0] 505 sb_pint(argv[6] as *u8, fa_len(argv[6] as *u8), 0, vout, pend) 506 let ttl: i64 = vout[0] 507 var mx: i64 = SJ_MAGIC_100000 508 if argc >= 9 { sb_pint(argv[8] as *u8, fa_len(argv[8] as *u8), 0, vout, pend); mx = vout[0] } 509 return sj_harvest(jr, ld, job, nj, ttl, argv[7] as *u8, mx) 510 } 511 } 512 return sj_gate() 513}