# ============================================================================ # jobs.ludic — Jobs, Promises & opt-in Sync, in Ludic (#14). # # Concurrency that keeps the *simple* thing simple. Ludic's simulation is # single-threaded and deterministic on purpose — the ECS schedule, lockstep # networking and replays all depend on it — so this library is layered: # # 1. Job.* / Promise.* (the safe default): kick off work, keep the frame # moving, and collect the result *on the main thread* at a point you choose # (Job.pump). No shared mutable state, no locks in game code. # 2. Sync.* (advanced, opt-in — "here be dragons"): mutex, atomic # counters and channels for engine-level systems that pass messages around. # # The whole thing is a deterministic *cooperative* scheduler: a background Job # runs a little each Job.pump(budget) and finishes after enough frames, so a # heavy computation (terrain, pathfinding pre-bake, a checksum) is spread out # instead of hitching one frame. Same jobs + same total budget over time -> # byte-identical results and completion order, every run and every target # (native and wasm alike). That determinism is why it is cooperative rather than # a preemptive OS-thread pool: a real-thread backend can slot in behind this same # API later without changing a line of game code, and the safe tier keeps its # hard rule either way — collect results on the main thread; never let a Job # reach into the ECS world itself. # # Ludic has no first-class functions, so a Job carries a small compute *kind* + # an integer argument (or you drive a hand-made future with Job.defer + # Job.fulfill) rather than a closure, and Promise progress is polled # (Promise.count_done) rather than chained through a callback. ludicc splices # this file into any program that mentions Job.* / Promise.* / Sync.*; it is # self-contained (only compiler intrinsics), so a plain tool works like a game. # The namespaces (emit_call.ludic) alias each method to a function below. # ============================================================================ const JOB_SLOTS: int = 128 # concurrent jobs (handles are 1-based slot ids) const JOB_MAXMEM: int = 64 # members a single Promise group can hold # job state (the terminal states are DONE / FAILED / CANCELLED) const J_FREE: int = 0 # slot not allocated const J_PENDING: int = 1 # allocated, not yet resolved const J_DONE: int = 2 # resolved successfully -> result const J_FAILED: int = 3 # resolved with an error -> error const J_CANCELLED: int = 4 # cancelled before it finished # job kind (what a PENDING slot is) const JK_DEFER: int = 0 # a hand-driven future (Job.fulfill / Job.fail) const JK_SUM: int = 1 # compute: 1 + 2 + ... + arg const JK_FIB: int = 2 # compute: the arg-th Fibonacci number const JK_PRIMES: int = 3 # compute: how many primes are <= arg const JK_ALL: int = 4 # Promise group: succeeds when every member does const JK_RACE: int = 5 # Promise group: succeeds when the first member does export state RtJobsState { jb_ready: bool = false jb_state: words = null # J_* jb_kind: words = null # JK_* jb_result: words = null # the success value jb_error: words = null # the failure code jb_arg: words = null # compute input n jb_i: words = null # compute progress counter jb_acc: words = null # compute accumulator jb_acc2: words = null # compute second accumulator (Fibonacci) jb_nmem: words = null # group member count jb_mem: words = null # flat [JOB_SLOTS * JOB_MAXMEM] of member handles sy_ready: bool = false mx_used: words = null mx_obj: pointers = null # the native mutex behind each handle at_used: words = null at_cell: pointers = null # the int each atomic handle names (a words(1) of its own) ch_used: words = null ch_head: words = null ch_count: words = null ch_buf: words = null # flat [SYNC_CHAN * CHAN_CAP] ch_lock: pointers = null # a mutex per channel } function jb_init(rt_jobs_st: mut RtJobsState) -> void { if rt_jobs_st.jb_ready { return } rt_jobs_st.jb_state = words(JOB_SLOTS); fill(rt_jobs_st.jb_state, 0, JOB_SLOTS * 4) rt_jobs_st.jb_kind = words(JOB_SLOTS); fill(rt_jobs_st.jb_kind, 0, JOB_SLOTS * 4) rt_jobs_st.jb_result = words(JOB_SLOTS); fill(rt_jobs_st.jb_result, 0, JOB_SLOTS * 4) rt_jobs_st.jb_error = words(JOB_SLOTS); fill(rt_jobs_st.jb_error, 0, JOB_SLOTS * 4) rt_jobs_st.jb_arg = words(JOB_SLOTS); fill(rt_jobs_st.jb_arg, 0, JOB_SLOTS * 4) rt_jobs_st.jb_i = words(JOB_SLOTS); fill(rt_jobs_st.jb_i, 0, JOB_SLOTS * 4) rt_jobs_st.jb_acc = words(JOB_SLOTS); fill(rt_jobs_st.jb_acc, 0, JOB_SLOTS * 4) rt_jobs_st.jb_acc2 = words(JOB_SLOTS); fill(rt_jobs_st.jb_acc2, 0, JOB_SLOTS * 4) rt_jobs_st.jb_nmem = words(JOB_SLOTS); fill(rt_jobs_st.jb_nmem, 0, JOB_SLOTS * 4) rt_jobs_st.jb_mem = words(JOB_SLOTS * JOB_MAXMEM); fill(rt_jobs_st.jb_mem, 0, JOB_SLOTS * JOB_MAXMEM * 4) rt_jobs_st.jb_ready = true } # claim a free slot as PENDING with the given kind; returns a 1-based handle, or # 0 if the table is full. function jb_alloc(rt_jobs_st: mut RtJobsState, kind: int) -> int { jb_init(rt_jobs_st) var i = 0 while i < JOB_SLOTS { if rt_jobs_st.jb_state[i] == J_FREE { rt_jobs_st.jb_state[i] = J_PENDING rt_jobs_st.jb_kind[i] = kind rt_jobs_st.jb_result[i] = 0; rt_jobs_st.jb_error[i] = 0 rt_jobs_st.jb_arg[i] = 0; rt_jobs_st.jb_i[i] = 0; rt_jobs_st.jb_acc[i] = 0; rt_jobs_st.jb_acc2[i] = 0 rt_jobs_st.jb_nmem[i] = 0 return i + 1 } i += 1 } return 0 } function jb_valid(rt_jobs_st: mut RtJobsState, h: int) -> bool { jb_init(rt_jobs_st) if (h < 1) or (h > JOB_SLOTS) { return false } return rt_jobs_st.jb_state[h - 1] != J_FREE } # ---- the safe tier: futures ------------------------------------------------ # A hand-driven future: PENDING until you call Job.fulfill / Job.fail on it. function job_defer(rt_jobs_st: mut RtJobsState) -> int { return jb_alloc(rt_jobs_st, JK_DEFER) } # Kick off a background compute job (kind = JK_SUM / JK_FIB / JK_PRIMES). It runs # a little each Job.pump and resolves when it finishes. `arg` is its input. function job_run(rt_jobs_st: mut RtJobsState, kind: int, arg: int) -> int { let h = jb_alloc(rt_jobs_st, kind) if h == 0 { return 0 } let s = h - 1 rt_jobs_st.jb_arg[s] = arg if kind == JK_FIB { rt_jobs_st.jb_acc[s] = 0; rt_jobs_st.jb_acc2[s] = 1 } # fib(0)=0, fib(1)=1 return h } # Resolve a pending job successfully with `value` (no-op once resolved). function job_fulfill(rt_jobs_st: mut RtJobsState, h: int, value: int) -> void { if not jb_valid(rt_jobs_st, h) { return } let s = h - 1 if rt_jobs_st.jb_state[s] != J_PENDING { return } rt_jobs_st.jb_state[s] = J_DONE rt_jobs_st.jb_result[s] = value } # Resolve a pending job as failed with error code `err` (no-op once resolved). function job_fail(rt_jobs_st: mut RtJobsState, h: int, err: int) -> void { if not jb_valid(rt_jobs_st, h) { return } let s = h - 1 if rt_jobs_st.jb_state[s] != J_PENDING { return } rt_jobs_st.jb_state[s] = J_FAILED rt_jobs_st.jb_error[s] = err } # Cancel a pending job (no-op if it already resolved). function job_cancel(rt_jobs_st: mut RtJobsState, h: int) -> void { if not jb_valid(rt_jobs_st, h) { return } let s = h - 1 if rt_jobs_st.jb_state[s] == J_PENDING { rt_jobs_st.jb_state[s] = J_CANCELLED } } # Recompute a group job (JK_ALL / JK_RACE) from its members. A no-op unless the # slot is a still-PENDING group. This is what "resolve on the main thread" means: # a Promise settles only when you look at it (done/ok/...) or pump. function jb_refresh_group(rt_jobs_st: mut RtJobsState, s: int) -> void { if rt_jobs_st.jb_state[s] != J_PENDING { return } let k = rt_jobs_st.jb_kind[s] if (k != JK_ALL) and (k != JK_RACE) { return } let n = rt_jobs_st.jb_nmem[s] let base = s * JOB_MAXMEM var i = 0 var settled = 0 # members in a terminal state var ok = 0 # members that succeeded var first_ok = 0 # winning handle for RACE while i < n { let mh = rt_jobs_st.jb_mem[base + i] if jb_valid(rt_jobs_st, mh) { let ms = mh - 1 let mst = rt_jobs_st.jb_state[ms] if mst != J_PENDING { settled += 1 if mst == J_DONE { ok += 1 if first_ok == 0 { first_ok = mh } } } } else { settled += 1 # a freed/invalid member counts as settled-failed } i += 1 } if k == JK_ALL { if ok == n { rt_jobs_st.jb_state[s] = J_DONE; rt_jobs_st.jb_result[s] = n } else { if settled == n { rt_jobs_st.jb_state[s] = J_FAILED; rt_jobs_st.jb_error[s] = n - ok } } } else { if first_ok != 0 { rt_jobs_st.jb_state[s] = J_DONE; rt_jobs_st.jb_result[s] = first_ok } else { if settled == n { rt_jobs_st.jb_state[s] = J_FAILED; rt_jobs_st.jb_error[s] = n } } } } # resolved in any terminal state? function job_done(rt_jobs_st: mut RtJobsState, h: int) -> bool { if not jb_valid(rt_jobs_st, h) { return false } jb_refresh_group(rt_jobs_st, h - 1) return rt_jobs_st.jb_state[h - 1] != J_PENDING } function job_ok(rt_jobs_st: mut RtJobsState, h: int) -> bool { if not jb_valid(rt_jobs_st, h) { return false } jb_refresh_group(rt_jobs_st, h - 1) return rt_jobs_st.jb_state[h - 1] == J_DONE } function job_failed(rt_jobs_st: mut RtJobsState, h: int) -> bool { if not jb_valid(rt_jobs_st, h) { return false } jb_refresh_group(rt_jobs_st, h - 1) return rt_jobs_st.jb_state[h - 1] == J_FAILED } function job_cancelled(rt_jobs_st: mut RtJobsState, h: int) -> bool { if not jb_valid(rt_jobs_st, h) { return false } return rt_jobs_st.jb_state[h - 1] == J_CANCELLED } # the success value (0 unless the job is done-ok) function job_result(rt_jobs_st: mut RtJobsState, h: int) -> int { if not jb_valid(rt_jobs_st, h) { return 0 } jb_refresh_group(rt_jobs_st, h - 1) if rt_jobs_st.jb_state[h - 1] != J_DONE { return 0 } return rt_jobs_st.jb_result[h - 1] } # the failure code (0 unless the job failed) function job_error(rt_jobs_st: mut RtJobsState, h: int) -> int { if not jb_valid(rt_jobs_st, h) { return 0 } jb_refresh_group(rt_jobs_st, h - 1) if rt_jobs_st.jb_state[h - 1] != J_FAILED { return 0 } return rt_jobs_st.jb_error[h - 1] } # how many jobs are still pending (a ready-made loading-screen denominator). function job_pending(rt_jobs_st: mut RtJobsState) -> int { jb_init(rt_jobs_st) var n = 0 var i = 0 while i < JOB_SLOTS { if rt_jobs_st.jb_state[i] == J_PENDING { n += 1 } i += 1 } return n } # release a slot back to the pool. function job_free(rt_jobs_st: mut RtJobsState, h: int) -> void { if not jb_valid(rt_jobs_st, h) { return } rt_jobs_st.jb_state[h - 1] = J_FREE } # advance one compute job by a single step; returns 1 if it just finished. function jb_step(rt_jobs_st: mut RtJobsState, s: int) -> int { let k = rt_jobs_st.jb_kind[s] let n = rt_jobs_st.jb_arg[s] var i = rt_jobs_st.jb_i[s] if k == JK_SUM { rt_jobs_st.jb_acc[s] = rt_jobs_st.jb_acc[s] + (i + 1) i += 1 rt_jobs_st.jb_i[s] = i if i >= n { rt_jobs_st.jb_state[s] = J_DONE; rt_jobs_st.jb_result[s] = rt_jobs_st.jb_acc[s]; return 1 } return 0 } if k == JK_FIB { if i >= n { rt_jobs_st.jb_state[s] = J_DONE; rt_jobs_st.jb_result[s] = rt_jobs_st.jb_acc[s]; return 1 } let t = rt_jobs_st.jb_acc[s] + rt_jobs_st.jb_acc2[s] rt_jobs_st.jb_acc[s] = rt_jobs_st.jb_acc2[s] rt_jobs_st.jb_acc2[s] = t i += 1 rt_jobs_st.jb_i[s] = i if i >= n { rt_jobs_st.jb_state[s] = J_DONE; rt_jobs_st.jb_result[s] = rt_jobs_st.jb_acc[s]; return 1 } return 0 } if k == JK_PRIMES { if jb_is_prime(i) { rt_jobs_st.jb_acc[s] += 1 } i += 1 rt_jobs_st.jb_i[s] = i if i > n { rt_jobs_st.jb_state[s] = J_DONE; rt_jobs_st.jb_result[s] = rt_jobs_st.jb_acc[s]; return 1 } return 0 } return 0 } function jb_is_prime(v: int) -> bool { if v < 2 { return false } var d = 2 while d * d <= v { if v - (v / d) * d == 0 { return false } d += 1 } return true } # Advance every pending compute job, spending up to `budget` steps in total, and # resolve any Promise groups. Call it once per frame (or wherever you want the # results to land). Returns how many jobs finished during this call. `budget` <= # 0 means "run every compute job to completion right now". function job_pump(rt_jobs_st: mut RtJobsState, budget: int) -> int { jb_init(rt_jobs_st) var completed = 0 var spent = 0 var s = 0 while s < JOB_SLOTS { let k = rt_jobs_st.jb_kind[s] let compute = (k == JK_SUM) or (k == JK_FIB) or (k == JK_PRIMES) while (rt_jobs_st.jb_state[s] == J_PENDING) and compute { if (budget > 0) and (spent >= budget) { s = JOB_SLOTS + 1; break } let fin = jb_step(rt_jobs_st, s) spent += 1 if fin == 1 { completed += 1 } } s += 1 } # settle groups after the compute jobs advanced this frame. s = 0 while s < JOB_SLOTS { if rt_jobs_st.jb_state[s] == J_PENDING { let k = rt_jobs_st.jb_kind[s] if (k == JK_ALL) or (k == JK_RACE) { jb_refresh_group(rt_jobs_st, s) if rt_jobs_st.jb_state[s] != J_PENDING { completed += 1 } } } s += 1 } return completed } # ---- Promise combinators (over a []int of job handles) --------------------- # store up to JOB_MAXMEM handles as the members of group slot `s`. function jb_set_members(rt_jobs_st: mut RtJobsState, s: int, handles: []int) -> void { var n = len(handles) if n > JOB_MAXMEM { n = JOB_MAXMEM } let base = s * JOB_MAXMEM var i = 0 while i < n { rt_jobs_st.jb_mem[base + i] = handles[i]; i += 1 } rt_jobs_st.jb_nmem[s] = n } # A promise that succeeds once every member has succeeded, and fails as soon as # the whole set has settled with at least one non-success. Returns a job handle. function prom_all(rt_jobs_st: mut RtJobsState, handles: []int) -> int { let h = jb_alloc(rt_jobs_st, JK_ALL) if h == 0 { return 0 } jb_set_members(rt_jobs_st, h - 1, handles) jb_refresh_group(rt_jobs_st, h - 1) return h } # A promise that succeeds as soon as the first member succeeds (its handle is the # result), and fails only if every member settles without success. function prom_race(rt_jobs_st: mut RtJobsState, handles: []int) -> int { let h = jb_alloc(rt_jobs_st, JK_RACE) if h == 0 { return 0 } jb_set_members(rt_jobs_st, h - 1, handles) jb_refresh_group(rt_jobs_st, h - 1) return h } # how many of `handles` have resolved (any terminal state) — a loading bar's # numerator; pair with len(handles) for the denominator. function prom_count_done(rt_jobs_st: mut RtJobsState, handles: []int) -> int { var n = 0 var i = 0 while i < len(handles) { if job_done(rt_jobs_st, handles[i]) { n += 1 } i += 1 } return n } function prom_all_done(rt_jobs_st: mut RtJobsState, handles: []int) -> bool { var i = 0 while i < len(handles) { if not job_done(rt_jobs_st, handles[i]) { return false } i += 1 } return true } # ============================================================================ # Sync.* — the advanced, opt-in tier. HERE BE DRAGONS. # # These are the raw building blocks — a lock, an atomic counter, a channel — for # engine-level systems that hand data between producers and consumers. On today's # single-threaded, deterministic runtime they are cooperative: correct, ordered # and replayable, and impossible to deadlock (there is one thread). They exist so # a message-passing system reads the same in game code now as it will when a # preemptive OS-thread backend lands behind this same API. Beginners never need # to touch this — reach for Job.* / Promise.* instead. # ============================================================================ const SYNC_MUTEX: int = 32 const SYNC_ATOMIC: int = 64 const SYNC_CHAN: int = 32 const CHAN_CAP: int = 64 # capacity of each channel's ring buffer # The OS side (threads.ll / threads_win.ll, linked with this file). Sync handles stay small ints; # behind each is a real mutex, or an int read and written with atomic instructions, so the calls # are safe from Job.parallel_for workers. Make the objects on the main thread before starting work. extern function thr_cpu_count() -> int = "thr_cpu_count" extern function thr_is_worker() -> int = "thr_is_worker" extern function thr_parallel_for(count: int, work: pointer, ctx: pointer) = "thr_parallel_for" extern function thr_mutex_new() -> pointer = "thr_mutex_new" extern function thr_lock(m: pointer) = "thr_lock" extern function thr_unlock(m: pointer) = "thr_unlock" extern function thr_trylock(m: pointer) -> int = "thr_trylock" extern function thr_atomic_add(p: pointer, delta: int) -> int = "thr_atomic_add" extern function thr_cas(p: pointer, expect: int, next: int) -> int = "thr_cas" extern function thr_load(p: pointer) -> int = "thr_load" extern function thr_store(p: pointer, v: int) = "thr_store" function sy_init(rt_jobs_st: mut RtJobsState) -> void { if rt_jobs_st.sy_ready { return } rt_jobs_st.mx_used = words(SYNC_MUTEX); fill(rt_jobs_st.mx_used, 0, SYNC_MUTEX * 4) rt_jobs_st.mx_obj = pointers(SYNC_MUTEX); fill(rt_jobs_st.mx_obj, 0, SYNC_MUTEX * 8) rt_jobs_st.at_used = words(SYNC_ATOMIC); fill(rt_jobs_st.at_used, 0, SYNC_ATOMIC * 4) rt_jobs_st.at_cell = pointers(SYNC_ATOMIC); fill(rt_jobs_st.at_cell, 0, SYNC_ATOMIC * 8) rt_jobs_st.ch_used = words(SYNC_CHAN); fill(rt_jobs_st.ch_used, 0, SYNC_CHAN * 4) rt_jobs_st.ch_head = words(SYNC_CHAN); fill(rt_jobs_st.ch_head, 0, SYNC_CHAN * 4) rt_jobs_st.ch_count = words(SYNC_CHAN); fill(rt_jobs_st.ch_count, 0, SYNC_CHAN * 4) rt_jobs_st.ch_buf = words(SYNC_CHAN * CHAN_CAP); fill(rt_jobs_st.ch_buf, 0, SYNC_CHAN * CHAN_CAP * 4) rt_jobs_st.ch_lock = pointers(SYNC_CHAN); fill(rt_jobs_st.ch_lock, 0, SYNC_CHAN * 8) rt_jobs_st.sy_ready = true } # ---- mutex ----------------------------------------------------------------- function sync_mutex(rt_jobs_st: mut RtJobsState) -> int { sy_init(rt_jobs_st) var i = 0 while i < SYNC_MUTEX { if rt_jobs_st.mx_used[i] == 0 { rt_jobs_st.mx_used[i] = 1 if rt_jobs_st.mx_obj[i] == null { rt_jobs_st.mx_obj[i] = thr_mutex_new() } return i + 1 } i += 1 } return 0 } function sync_lock(rt_jobs_st: mut RtJobsState, m: int) -> void { sy_init(rt_jobs_st) if (m < 1) or (m > SYNC_MUTEX) { return } thr_lock(rt_jobs_st.mx_obj[m - 1]) } function sync_unlock(rt_jobs_st: mut RtJobsState, m: int) -> void { sy_init(rt_jobs_st) if (m < 1) or (m > SYNC_MUTEX) { return } thr_unlock(rt_jobs_st.mx_obj[m - 1]) } # take the lock only if it is free; returns whether it was taken. function sync_try_lock(rt_jobs_st: mut RtJobsState, m: int) -> bool { sy_init(rt_jobs_st) if (m < 1) or (m > SYNC_MUTEX) { return false } return thr_trylock(rt_jobs_st.mx_obj[m - 1]) == 1 } # ---- atomic counter -------------------------------------------------------- function sync_atomic(rt_jobs_st: mut RtJobsState) -> int { sy_init(rt_jobs_st) var i = 0 while i < SYNC_ATOMIC { if rt_jobs_st.at_used[i] == 0 { rt_jobs_st.at_used[i] = 1 if rt_jobs_st.at_cell[i] == null { rt_jobs_st.at_cell[i] = data_of(words(1)) } # the atomics take the cell's address, not its slice thr_store(rt_jobs_st.at_cell[i], 0) return i + 1 } i += 1 } return 0 } function sync_get(rt_jobs_st: mut RtJobsState, a: int) -> int { sy_init(rt_jobs_st) if (a < 1) or (a > SYNC_ATOMIC) { return 0 } return thr_load(rt_jobs_st.at_cell[a - 1]) } function sync_set(rt_jobs_st: mut RtJobsState, a: int, v: int) -> void { sy_init(rt_jobs_st) if (a < 1) or (a > SYNC_ATOMIC) { return } thr_store(rt_jobs_st.at_cell[a - 1], v) } # add `delta` and return the new value. function sync_add(rt_jobs_st: mut RtJobsState, a: int, delta: int) -> int { sy_init(rt_jobs_st) if (a < 1) or (a > SYNC_ATOMIC) { return 0 } return thr_atomic_add(rt_jobs_st.at_cell[a - 1], delta) } # compare-and-set: if the value equals `expect`, store `next` and return true. function sync_cas(rt_jobs_st: mut RtJobsState, a: int, expect: int, next: int) -> bool { sy_init(rt_jobs_st) if (a < 1) or (a > SYNC_ATOMIC) { return false } return thr_cas(rt_jobs_st.at_cell[a - 1], expect, next) == 1 } # ---- channel (a bounded int FIFO, behind its own mutex) --------------------- function sync_channel(rt_jobs_st: mut RtJobsState) -> int { sy_init(rt_jobs_st) var i = 0 while i < SYNC_CHAN { if rt_jobs_st.ch_used[i] == 0 { rt_jobs_st.ch_used[i] = 1; rt_jobs_st.ch_head[i] = 0; rt_jobs_st.ch_count[i] = 0 if rt_jobs_st.ch_lock[i] == null { rt_jobs_st.ch_lock[i] = thr_mutex_new() } return i + 1 } i += 1 } return 0 } # enqueue `v`; returns false if the channel is full. function sync_send(rt_jobs_st: mut RtJobsState, c: int, v: int) -> bool { sy_init(rt_jobs_st) if (c < 1) or (c > SYNC_CHAN) { return false } let s = c - 1 thr_lock(rt_jobs_st.ch_lock[s]) if rt_jobs_st.ch_count[s] >= CHAN_CAP { thr_unlock(rt_jobs_st.ch_lock[s]); return false } let pos = rt_jobs_st.ch_head[s] + rt_jobs_st.ch_count[s] var idx = pos if idx >= CHAN_CAP { idx -= CHAN_CAP } rt_jobs_st.ch_buf[s * CHAN_CAP + idx] = v rt_jobs_st.ch_count[s] += 1 thr_unlock(rt_jobs_st.ch_lock[s]) return true } # dequeue the oldest value; returns 0 on an empty channel (guard with can_recv). function sync_recv(rt_jobs_st: mut RtJobsState, c: int) -> int { sy_init(rt_jobs_st) if (c < 1) or (c > SYNC_CHAN) { return 0 } let s = c - 1 thr_lock(rt_jobs_st.ch_lock[s]) if rt_jobs_st.ch_count[s] == 0 { thr_unlock(rt_jobs_st.ch_lock[s]); return 0 } let v = rt_jobs_st.ch_buf[s * CHAN_CAP + rt_jobs_st.ch_head[s]] var nh = rt_jobs_st.ch_head[s] + 1 if nh >= CHAN_CAP { nh = 0 } rt_jobs_st.ch_head[s] = nh rt_jobs_st.ch_count[s] -= 1 thr_unlock(rt_jobs_st.ch_lock[s]) return v } function sync_can_recv(rt_jobs_st: mut RtJobsState, c: int) -> bool { sy_init(rt_jobs_st) if (c < 1) or (c > SYNC_CHAN) { return false } thr_lock(rt_jobs_st.ch_lock[c - 1]) let has = rt_jobs_st.ch_count[c - 1] > 0 thr_unlock(rt_jobs_st.ch_lock[c - 1]) return has } function sync_len(rt_jobs_st: mut RtJobsState, c: int) -> int { sy_init(rt_jobs_st) if (c < 1) or (c > SYNC_CHAN) { return 0 } thr_lock(rt_jobs_st.ch_lock[c - 1]) let n = rt_jobs_st.ch_count[c - 1] thr_unlock(rt_jobs_st.ch_lock[c - 1]) return n } # the machine's logical cores: how many threads Job.parallel_for spreads work across function sync_cpu_count() -> int { return thr_cpu_count() } # ---- Job.parallel_for (real threads) ----------------------------------------- # `work(i, ctx)` for every i in [0, count), across one worker per core but one and the calling # thread; returns when every call has returned. `work` is a function reference (`fn name`) taking # (int, pointer). The rule: a worker computes on what `ctx` points at and writes its results there # - it never spawns, despawns, pushes onto a list another thread can see, or touches the world. function job_parallel_for(count: int, work: pointer, ctx: pointer) -> void { thr_parallel_for(count, work, ctx) } # true on a Job.parallel_for worker thread function job_is_worker() -> bool { return thr_is_worker() == 1 }