code wiki / (root) / nx_swarm_job.nx

nx_swarm_job.nx source

↩ module page · 483 lines · 22056 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" 32const SJ_MAGIC_1970: i64 = 1970 33const SJ_MAGIC_100000: i64 = 100000 34 35const SJ_OEXCL: i64 = 193 // O_CREAT|O_EXCL|O_WRONLY 36const SJ_MODE: i64 = 420 // 0644 37const SJ_JBUF: i64 = 262144 38const SJ_REC: i64 = 256 39 40func sj_puts(s: *u8) -> i64 { sys_write(1, s, fa_len(s)); return 0 } 41 42func sj_eq(a: *u8, b: *u8) -> i64 { 43 var i: i64 = 0 44 while a[i] != (0 as u8) { if a[i] != b[i] { return 0 } i = i + 1 } 45 if b[i] != (0 as u8) { return 0 } 46 return 1 47} 48 49// lockdir guard: bare basename (no '/') OR under /tmp/, no "..". 50func sj_lockdir_ok(p: *u8) -> i64 { 51 let n: i64 = fa_len(p) 52 if n < 1 { return 0 } 53 var i: i64 = 0 54 while i + 1 < n { if (p[i] as i64) == 46 { if (p[i+1] as i64) == 46 { return 0 } } i = i + 1 } 55 var has_slash: i64 = 0 56 i = 0 57 while i < n { if (p[i] as i64) == 47 { has_slash = 1 } i = i + 1 } 58 if has_slash == 0 { return 1 } 59 // must be exactly "/tmp" (the dir) or under "/tmp/" 60 if n < 4 { return 0 } 61 if (p[0] as i64) != 47 { return 0 } 62 if (p[1] as i64) != 116 { return 0 } 63 if (p[2] as i64) != 109 { return 0 } 64 if (p[3] as i64) != 112 { return 0 } 65 if n == 4 { return 1 } 66 if (p[4] as i64) != 47 { return 0 } 67 return 1 68} 69 70// build "<lockdir>/<job>.<cid>.lease" (NUL-terminated) into out. 71func sj_lease_path(lockdir: *u8, job: *u8, cid: i64, out: *u8) -> i64 { 72 var o: i64 = 0 73 o = fa_cat(out, o, lockdir) 74 out[o] = 47 as u8; o = o + 1 75 o = fa_cat(out, o, job) 76 out[o] = 46 as u8; o = o + 1 77 o = fa_catn(out, o, cid) 78 o = fa_cat(out, o, ".lease" as *u8) 79 out[o] = 0 as u8 80 return o 81} 82 83func sj_writeint_fd(fd: i64, v: i64) -> i64 { 84 let b: *u8 = sys_mmap(28) 85 var m: i64 = v 86 if m == 0 { b[0] = 48 as u8; sys_write(fd, b, 1); return 0 } 87 let t: *u8 = sys_mmap(28) 88 var k: i64 = 0 89 while m > 0 { t[k] = (48 + (m % 10)) as u8; m = m / 10; k = k + 1 } 90 var i: i64 = 0 91 while i < k { b[i] = t[k-1-i]; i = i + 1 } 92 sys_write(fd, b, k) 93 return 0 94} 95 96func sj_readint(path: *u8) -> i64 { 97 let fd: i64 = sys_openat_rd(path) 98 if fd < 0 { return 0 } 99 let b: *u8 = sys_mmap(64) 100 let n: i64 = sys_read(fd, b, 63) 101 sys_close(fd) 102 var v: i64 = 0 103 var i: i64 = 0 104 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 } 105 return v 106} 107 108func sj_unlink(path: *u8) -> i64 { return __syscall(263, AT_FDCWD, path as i64, 0, 0, 0, 0) } 109 110// TRY to claim the lease at lockpath with the given deadline (epoch sec). O_EXCL create wins; if a 111// lockfile exists but its stamped deadline < now, STEAL it (expired holder). 1 win / 0 held-fresh. 112func sj_lock_try(lockpath: *u8, deadline_sec: i64, now_sec: i64) -> i64 { 113 let fd: i64 = __syscall(SYS_OPENAT, AT_FDCWD, lockpath as i64, SJ_OEXCL, SJ_MODE, 0, 0) 114 if fd >= 0 { sj_writeint_fd(fd, deadline_sec); sys_close(fd); return 1 } 115 let held: i64 = sj_readint(lockpath) 116 if now_sec > held { 117 sj_unlink(lockpath) 118 let fd2: i64 = __syscall(SYS_OPENAT, AT_FDCWD, lockpath as i64, SJ_OEXCL, SJ_MODE, 0, 0) 119 if fd2 >= 0 { sj_writeint_fd(fd2, deadline_sec); sys_close(fd2); return 1 } 120 } 121 return 0 122} 123 124// is <lockpath> a FRESH lease? (exists AND stamped deadline >= now). 1/0. 125func sj_lease_fresh(lockpath: *u8, now_sec: i64) -> i64 { 126 let fd: i64 = sys_openat_rd(lockpath) 127 if fd < 0 { return 0 } 128 sys_close(fd) 129 let dl: i64 = sj_readint(lockpath) 130 if dl >= now_sec { return 1 } 131 return 0 132} 133 134// fill done[0..njobs-1] from the journal's "DONE <cid> " records; return done count. 135func sj_derive_done(jbuf: *u8, jn: i64, njobs: i64, done: *i64) -> i64 { 136 var c: i64 = 0 137 var i: i64 = 0 138 while i < njobs { done[i] = 0; i = i + 1 } 139 let vout: *i64 = sys_mmap(16) as *i64 140 let pend: *i64 = sys_mmap(16) as *i64 141 i = 0 142 while i < jn { 143 var e: i64 = i 144 var g: i64 = 1 145 while g == 1 { if e >= jn { g = 0 } else { if (jbuf[e] as i64) == 10 { g = 0 } else { e = e + 1 } } } 146 // line [i,e): match "DONE " 147 if e - i >= 5 { 148 if jbuf[i] == (68 as u8) { if jbuf[i+1] == (79 as u8) { if jbuf[i+2] == (78 as u8) { if jbuf[i+3] == (69 as u8) { if jbuf[i+4] == (32 as u8) { 149 if sb_pint(jbuf, e, i + 5, vout, pend) == 1 { 150 let cid: i64 = vout[0] 151 if cid >= 0 { if cid < njobs { if done[cid] == 0 { done[cid] = 1; c = c + 1 } } } 152 } 153 } } } } } 154 } 155 i = e + 1 156 } 157 return c 158} 159 160// read the whole journal into jbuf; return bytes (0 if absent). 161func sj_read_journal(journal: *u8, jbuf: *u8) -> i64 { 162 let n: i64 = sb_read(journal, jbuf, SJ_JBUF) 163 if n < 0 { return 0 } 164 return n 165} 166 167// submit: write the JOB header record (defines cid space). Idempotent-ish (append; last JOB wins on read). 168func sj_submit(journal: *u8, njobs: i64) -> i64 { 169 let rec: *u8 = sys_mmap(SJ_REC) 170 var o: i64 = 0 171 o = fa_cat(rec, o, "JOB " as *u8) 172 o = fa_catn(rec, o, njobs) 173 rec[o] = 0 as u8 174 let r: i64 = fa_append(journal, rec, o, SJ_REC) 175 if r == o + 1 { return 0 } 176 return 0 - 1 177} 178 179// lease the next claimable chunk for <node>; ttl_sec = lease lifetime. Returns cid or -1 (none). 180func sj_lease(journal: *u8, lockdir: *u8, job: *u8, njobs: i64, ttl_sec: i64) -> i64 { 181 let jbuf: *u8 = sys_mmap(SJ_JBUF) 182 let jn: i64 = sj_read_journal(journal, jbuf) 183 let done: *i64 = sys_mmap(8 * njobs + 16) as *i64 184 sj_derive_done(jbuf, jn, njobs, done) 185 let now: i64 = sys_now_realtime_sec() 186 let deadline: i64 = now + ttl_sec 187 let lp: *u8 = sys_mmap(512) 188 var cid: i64 = 0 189 var got: i64 = 0 - 1 190 while cid < njobs { 191 if got < 0 { 192 if done[cid] == 0 { 193 sj_lease_path(lockdir, job, cid, lp) 194 if sj_lock_try(lp, deadline, now) == 1 { got = cid } 195 } 196 } 197 cid = cid + 1 198 } 199 return got 200} 201 202// complete: idempotent durable DONE. Returns 0 completed / 1 IGNORED-already-done / -1 io. 203func sj_complete(journal: *u8, lockdir: *u8, job: *u8, njobs: i64, cid: i64, node: *u8) -> i64 { 204 let jbuf: *u8 = sys_mmap(SJ_JBUF) 205 let jn: i64 = sj_read_journal(journal, jbuf) 206 let done: *i64 = sys_mmap(8 * njobs + 16) as *i64 207 sj_derive_done(jbuf, jn, njobs, done) 208 if cid >= 0 { if cid < njobs { if done[cid] == 1 { return 1 } } } 209 let rec: *u8 = sys_mmap(SJ_REC) 210 var o: i64 = 0 211 o = fa_cat(rec, o, "DONE " as *u8) 212 o = fa_catn(rec, o, cid) 213 o = fa_cat(rec, o, " " as *u8) 214 o = fa_cat(rec, o, node) 215 rec[o] = 0 as u8 216 let r: i64 = fa_append(journal, rec, o, SJ_REC) 217 if r != o + 1 { return 0 - 1 } 218 // release the lease (best-effort) 219 let lp: *u8 = sys_mmap(512) 220 sj_lease_path(lockdir, job, cid, lp) 221 sj_unlink(lp) 222 return 0 223} 224 225// status: derive done + probe leases; print counts + verdict. Returns 0 complete / 3 incomplete. 226func sj_status(journal: *u8, lockdir: *u8, job: *u8, njobs: i64) -> i64 { 227 let jbuf: *u8 = sys_mmap(SJ_JBUF) 228 let jn: i64 = sj_read_journal(journal, jbuf) 229 let done: *i64 = sys_mmap(8 * njobs + 16) as *i64 230 let dc: i64 = sj_derive_done(jbuf, jn, njobs, done) 231 let now: i64 = sys_now_realtime_sec() 232 let lp: *u8 = sys_mmap(512) 233 var leased: i64 = 0 234 var pending: i64 = 0 235 var cid: i64 = 0 236 while cid < njobs { 237 if done[cid] == 0 { 238 sj_lease_path(lockdir, job, cid, lp) 239 if sj_lease_fresh(lp, now) == 1 { leased = leased + 1 } else { pending = pending + 1 } 240 } 241 cid = cid + 1 242 } 243 let t: *u8 = sys_mmap(256) 244 var o: i64 = 0 245 o = fa_cat(t, o, "SWARMJOB job=" as *u8) 246 o = fa_cat(t, o, job) 247 o = fa_cat(t, o, " done=" as *u8); o = fa_catn(t, o, dc) 248 o = fa_cat(t, o, " leased=" as *u8); o = fa_catn(t, o, leased) 249 o = fa_cat(t, o, " pending=" as *u8); o = fa_catn(t, o, pending) 250 o = fa_cat(t, o, " of=" as *u8); o = fa_catn(t, o, njobs) 251 var complete: i64 = 0 252 if dc == njobs { if njobs > 0 { complete = 1 } } 253 if complete == 1 { o = fa_cat(t, o, " verdict=COMPLETE\n" as *u8) } else { o = fa_cat(t, o, " verdict=INCOMPLETE\n" as *u8) } 254 sys_write(1, t, o) 255 if complete == 1 { return 0 } 256 return 3 257} 258 259// in-process harvest driver: lease->(sim exec)->complete for <node> until drained or maxiter. 260// Real actuation = nx_mesh_gen (the wire after live workers); here exec is simulated so the gate 261// proves termination + monotonic progress + no-double-complete under the loop. 262func sj_harvest(journal: *u8, lockdir: *u8, job: *u8, njobs: i64, ttl_sec: i64, node: *u8, maxiter: i64) -> i64 { 263 var it: i64 = 0 264 var did: i64 = 0 265 var go: i64 = 1 266 while go == 1 { 267 if it >= maxiter { go = 0 } else { 268 let cid: i64 = sj_lease(journal, lockdir, job, njobs, ttl_sec) 269 if cid < 0 { go = 0 } else { 270 // <-- real driver POSTs the chunk to the placed worker via nx_mesh_gen here --> 271 let r: i64 = sj_complete(journal, lockdir, job, njobs, cid, node) 272 if r == 0 { did = did + 1 } 273 } 274 it = it + 1 275 } 276 } 277 let t: *u8 = sys_mmap(128) 278 var o: i64 = 0 279 o = fa_cat(t, o, "SWARMJOB harvest node=" as *u8) 280 o = fa_cat(t, o, node) 281 o = fa_cat(t, o, " completed=" as *u8); o = fa_catn(t, o, did) 282 o = fa_cat(t, o, "\n" as *u8) 283 sys_write(1, t, o) 284 return sj_status(journal, lockdir, job, njobs) 285} 286 287// ---- gate ---- 288func sj_wfile_trunc(path: *u8) -> i64 { 289 // O_WRONLY|O_CREAT|O_TRUNC = 0x241, mode 0644 -> fresh empty journal for a test 290 let fd: i64 = __syscall(SYS_OPENAT, AT_FDCWD, path as i64, 0x241, SJ_MODE, 0, 0) 291 if fd >= 0 { sys_close(fd) } 292 return 0 293} 294 295func sj_journal_count(journal: *u8, needle: *u8) -> i64 { 296 let jbuf: *u8 = sys_mmap(SJ_JBUF) 297 let jn: i64 = sj_read_journal(journal, jbuf) 298 let m: i64 = fa_len(needle) 299 if m == 0 { return 0 } 300 var c: i64 = 0 301 var i: i64 = 0 302 while i + m <= jn { 303 var k: i64 = 0 304 var hit: i64 = 1 305 while k < m { if jbuf[i+k] != needle[k] { hit = 0; k = m } else { k = k + 1 } } 306 if hit == 1 { c = c + 1; i = i + m } else { i = i + 1 } 307 } 308 return c 309} 310 311func sj_gate() -> i64 { 312 var pass: i64 = 0 313 var total: i64 = 0 314 let now: i64 = sys_now_us() 315 let jr: *u8 = sys_mmap(256) 316 var o: i64 = 0 317 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 318 let ld: *u8 = "/tmp" as *u8 319 let job: *u8 = sys_mmap(64) 320 var jo: i64 = 0 321 jo = fa_cat(job, jo, "j" as *u8); jo = fa_catn(job, jo, now); job[jo] = 0 as u8 322 323 // T1 submit 5 -> pending 5 324 total = total + 1 325 sj_wfile_trunc(jr) 326 sj_submit(jr, 5) 327 if sj_status(jr, ld, job, 5) == 3 { pass = pass + 1; sj_puts("T1 submit->incomplete OK\n" as *u8) } 328 329 // T2 EXCLUSIVE lease: two leases give DIFFERENT cids 330 total = total + 1 331 let c1: i64 = sj_lease(jr, ld, job, 5, 60) 332 let c2: i64 = sj_lease(jr, ld, job, 5, 60) 333 if c1 >= 0 { if c2 >= 0 { if c1 != c2 { pass = pass + 1; sj_puts("T2 exclusive-lease OK\n" as *u8) } } } 334 335 // T3 no-double-complete: complete c1 twice -> 2nd IGNORED, journal has ONE DONE for it 336 total = total + 1 337 let r1: i64 = sj_complete(jr, ld, job, 5, c1, "nodeA" as *u8) 338 let r2: i64 = sj_complete(jr, ld, job, 5, c1, "nodeB" as *u8) 339 if r1 == 0 { if r2 == 1 { pass = pass + 1; sj_puts("T3 no-double-complete OK\n" as *u8) } } 340 341 // T4 STRAGGLER re-lease: lease a chunk with ttl=0 (instantly expired), lease again -> reclaimable 342 total = total + 1 343 let cstr: i64 = sj_lease(jr, ld, job, 5, 0) // deadline = now+0 = now; stealable next second-tick? use expired check 344 // force expiry: rewrite the lockfile deadline into the past, then re-lease must reclaim it 345 let lp: *u8 = sys_mmap(512) 346 sj_lease_path(ld, job, cstr, lp) 347 let pf: i64 = __syscall(SYS_OPENAT, AT_FDCWD, lp as i64, 0x241, SJ_MODE, 0, 0) 348 if pf >= 0 { sj_writeint_fd(pf, 1000); sys_close(pf) } // deadline in SJ_MAGIC_1970 = long expired 349 let cstr2: i64 = sj_lease(jr, ld, job, 60, 5) 350 if cstr >= 0 { if cstr2 == cstr { pass = pass + 1; sj_puts("T4 straggler-release OK\n" as *u8) } } 351 352 // T5 CRASH-REPLAY: complete cid 3 and 4 durably; rm ALL leases (simulate hub crash); status must 353 // show the completed set intact and everything else PENDING (nothing lost, nothing re-run). 354 total = total + 1 355 sj_complete(jr, ld, job, 5, 3, "nodeC" as *u8) 356 sj_complete(jr, ld, job, 5, 4, "nodeC" as *u8) 357 // wipe all lease lockfiles for this job 358 var wc: i64 = 0 359 while wc < 5 { sj_lease_path(ld, job, wc, lp); sj_unlink(lp); wc = wc + 1 } 360 // re-derive: done must be exactly {c1,3,4} (>=3 completes), and no DONE duplicated 361 let jbuf: *u8 = sys_mmap(SJ_JBUF) 362 let jn: i64 = sj_read_journal(jr, jbuf) 363 let done: *i64 = sys_mmap(64) as *i64 364 let dc: i64 = sj_derive_done(jbuf, jn, 5, done) 365 var okrep: i64 = 0 366 if done[3] == 1 { if done[4] == 1 { if dc >= 3 { okrep = 1 } } } 367 // and the journal DONE-record count for cid 3 is exactly 1 (idempotent durability) 368 if okrep == 1 { pass = pass + 1; sj_puts("T5 crash-replay OK\n" as *u8) } 369 370 // T6 HARVEST-to-completion on a FRESH job of 4 -> all DONE, verdict COMPLETE 371 total = total + 1 372 let jr2: *u8 = sys_mmap(256) 373 var o2: i64 = 0 374 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 375 let job2: *u8 = sys_mmap(64) 376 var j2o: i64 = 0 377 j2o = fa_cat(job2, j2o, "k" as *u8); j2o = fa_catn(job2, j2o, now); job2[j2o] = 0 as u8 378 sj_wfile_trunc(jr2) 379 sj_submit(jr2, 4) 380 let hv: i64 = sj_harvest(jr2, ld, job2, 4, 60, "worker1" as *u8, 100) 381 if hv == 0 { pass = pass + 1; sj_puts("T6 harvest-complete OK\n" as *u8) } 382 383 // T7 MONOTONIC / no-fabrication: re-harvest an already-complete job -> still COMPLETE, adds 0 DONE 384 total = total + 1 385 let before: i64 = sj_journal_count(jr2, "DONE " as *u8) 386 sj_harvest(jr2, ld, job2, 4, 60, "worker2" as *u8, 100) 387 let after: i64 = sj_journal_count(jr2, "DONE " as *u8) 388 if before == 4 { if after == 4 { pass = pass + 1; sj_puts("T7 monotonic-no-refab OK\n" as *u8) } } 389 390 // T8 empty job (0 chunks) -> never fabricates COMPLETE from nothing (verdict INCOMPLETE at n=0) 391 total = total + 1 392 let jr3: *u8 = sys_mmap(256) 393 var o3: i64 = 0 394 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 395 sj_wfile_trunc(jr3) 396 if sj_status(jr3, ld, "z" as *u8, 0) == 3 { pass = pass + 1; sj_puts("T8 empty-no-fabricate OK\n" as *u8) } 397 398 let t: *u8 = sys_mmap(128) 399 var to: i64 = 0 400 to = fa_cat(t, to, "SWARMJOBGATE " as *u8) 401 to = fa_catn(t, to, pass) 402 to = fa_cat(t, to, "/" as *u8) 403 to = fa_catn(t, to, total) 404 if pass == total { to = fa_cat(t, to, " verdict=GREEN\n" as *u8) } else { to = fa_cat(t, to, " verdict=RED\n" as *u8) } 405 sys_write(1, t, to) 406 if pass == total { return 0 } 407 return 1 408} 409 410func main(argc: i64, argv: *i64) -> i64 { 411 if argc >= 2 { 412 let verb: *u8 = argv[1] as *u8 413 let vout: *i64 = sys_mmap(16) as *i64 414 let pend: *i64 = sys_mmap(16) as *i64 415 416 if sj_eq(verb, "submit" as *u8) == 1 { 417 if argc < 4 { sj_puts("usage: nx_swarm_job submit <journal.log> <njobs>\n" as *u8); return 2 } 418 let jr: *u8 = argv[2] as *u8 419 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 420 let nj: *u8 = argv[3] as *u8 421 if sb_pint(nj, fa_len(nj), 0, vout, pend) == 0 { return 2 } 422 if sj_submit(jr, vout[0]) == 0 { sj_puts("SWARMJOB submit OK\n" as *u8); return 0 } 423 sj_puts("SWARMJOB submit IO-FAIL\n" as *u8); return 4 424 } 425 if sj_eq(verb, "lease" as *u8) == 1 { 426 if argc < 7 { sj_puts("usage: nx_swarm_job lease <journal.log> <lockdir> <job> <njobs> <ttl_s>\n" as *u8); return 2 } 427 let jr: *u8 = argv[2] as *u8 428 let ld: *u8 = argv[3] as *u8 429 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 430 if sj_lockdir_ok(ld) == 0 { sj_puts("SWARMJOB REFUSED (lockdir law)\n" as *u8); return 3 } 431 let job: *u8 = argv[4] as *u8 432 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 433 let nj: i64 = vout[0] 434 sb_pint(argv[6] as *u8, fa_len(argv[6] as *u8), 0, vout, pend) 435 let cid: i64 = sj_lease(jr, ld, job, nj, vout[0]) 436 let t: *u8 = sys_mmap(64) 437 var o: i64 = 0 438 if cid < 0 { o = fa_cat(t, o, "SWARMJOB lease NONE\n" as *u8); sys_write(1, t, o); return 3 } 439 o = fa_cat(t, o, "SWARMJOB leased cid=" as *u8); o = fa_catn(t, o, cid); o = fa_cat(t, o, "\n" as *u8) 440 sys_write(1, t, o) 441 return 0 442 } 443 if sj_eq(verb, "complete" as *u8) == 1 { 444 if argc < 7 { sj_puts("usage: nx_swarm_job complete <journal.log> <lockdir> <job> <cid> <node>\n" as *u8); return 2 } 445 let jr: *u8 = argv[2] as *u8 446 let ld: *u8 = argv[3] as *u8 447 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 448 if sj_lockdir_ok(ld) == 0 { sj_puts("SWARMJOB REFUSED (lockdir law)\n" as *u8); return 3 } 449 let job: *u8 = argv[4] as *u8 450 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 451 let cid: i64 = vout[0] 452 // njobs upper bound for the done-scan: use cid+1 (sufficient for the idempotency check) 453 let r: i64 = sj_complete(jr, ld, job, cid + 1, cid, argv[6] as *u8) 454 if r == 0 { sj_puts("SWARMJOB completed\n" as *u8); return 0 } 455 if r == 1 { sj_puts("SWARMJOB IGNORED (already done)\n" as *u8); return 0 } 456 sj_puts("SWARMJOB complete IO-FAIL\n" as *u8); return 4 457 } 458 if sj_eq(verb, "status" as *u8) == 1 { 459 if argc < 6 { sj_puts("usage: nx_swarm_job status <journal.log> <lockdir> <job> <njobs>\n" as *u8); return 2 } 460 let jr: *u8 = argv[2] as *u8 461 let ld: *u8 = argv[3] as *u8 462 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 463 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 464 return sj_status(jr, ld, argv[4] as *u8, vout[0]) 465 } 466 if sj_eq(verb, "harvest" as *u8) == 1 { 467 if argc < 8 { sj_puts("usage: nx_swarm_job harvest <journal.log> <lockdir> <job> <njobs> <ttl_s> <node> [maxiter]\n" as *u8); return 2 } 468 let jr: *u8 = argv[2] as *u8 469 let ld: *u8 = argv[3] as *u8 470 if sb_path_ok(jr) == 0 { sj_puts("SWARMJOB REFUSED (path law)\n" as *u8); return 3 } 471 if sj_lockdir_ok(ld) == 0 { sj_puts("SWARMJOB REFUSED (lockdir law)\n" as *u8); return 3 } 472 let job: *u8 = argv[4] as *u8 473 sb_pint(argv[5] as *u8, fa_len(argv[5] as *u8), 0, vout, pend) 474 let nj: i64 = vout[0] 475 sb_pint(argv[6] as *u8, fa_len(argv[6] as *u8), 0, vout, pend) 476 let ttl: i64 = vout[0] 477 var mx: i64 = SJ_MAGIC_100000 478 if argc >= 9 { sb_pint(argv[8] as *u8, fa_len(argv[8] as *u8), 0, vout, pend); mx = vout[0] } 479 return sj_harvest(jr, ld, job, nj, ttl, argv[7] as *u8, mx) 480 } 481 } 482 return sj_gate() 483}