ludic/runtime/native/jobs.ludic
Orkuncakilkaya b0b0b62bce feat(lang): L7 memory is safe unless it says unsafe
The typed buffers are slices: words/floats/fixeds/doubles/pointers(n) make
zeroed, bounds-checked []int/[]float/... and the type names mean them. buffer(n)
is a []byte, with text_of, Fs.read_bytes/write_bytes and view(xs, start, n).
bytes(), indexing a raw pointer or bytes, free, resize, Memory.*, raw file calls,
data_of and C externs are refused outside unsafe { } / unsafe function, and a
project's own files may write unsafe only with --unsafe; the runtime and packages
are the platform. A slice passed to an extern goes as its data.

What the change found: Sync's atomics on a slice header, words(n) uninitialised,
input's fixed axes in ints, truetype's fixed outlines as ints, skin matrices
typed int, gl_shader's source table made from raw bytes. render3d gets safe
entry points (safe_api.ludic). Rendering is byte-identical; a frame costs the same.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-24 12:53:27 +03:00

580 lines
20 KiB
Text

# ============================================================================
# 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
var jb_ready: bool = false
var jb_state: words = null # J_*
var jb_kind: words = null # JK_*
var jb_result: words = null # the success value
var jb_error: words = null # the failure code
var jb_arg: words = null # compute input n
var jb_i: words = null # compute progress counter
var jb_acc: words = null # compute accumulator
var jb_acc2: words = null # compute second accumulator (Fibonacci)
var jb_nmem: words = null # group member count
var jb_mem: words = null # flat [JOB_SLOTS * JOB_MAXMEM] of member handles
function jb_init() -> void {
if jb_ready { return }
jb_state = words(JOB_SLOTS); fill(jb_state, 0, JOB_SLOTS * 4)
jb_kind = words(JOB_SLOTS); fill(jb_kind, 0, JOB_SLOTS * 4)
jb_result = words(JOB_SLOTS); fill(jb_result, 0, JOB_SLOTS * 4)
jb_error = words(JOB_SLOTS); fill(jb_error, 0, JOB_SLOTS * 4)
jb_arg = words(JOB_SLOTS); fill(jb_arg, 0, JOB_SLOTS * 4)
jb_i = words(JOB_SLOTS); fill(jb_i, 0, JOB_SLOTS * 4)
jb_acc = words(JOB_SLOTS); fill(jb_acc, 0, JOB_SLOTS * 4)
jb_acc2 = words(JOB_SLOTS); fill(jb_acc2, 0, JOB_SLOTS * 4)
jb_nmem = words(JOB_SLOTS); fill(jb_nmem, 0, JOB_SLOTS * 4)
jb_mem = words(JOB_SLOTS * JOB_MAXMEM); fill(jb_mem, 0, JOB_SLOTS * JOB_MAXMEM * 4)
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(kind: int) -> int {
jb_init()
var i = 0
while i < JOB_SLOTS {
if jb_state[i] == J_FREE {
jb_state[i] = J_PENDING
jb_kind[i] = kind
jb_result[i] = 0; jb_error[i] = 0
jb_arg[i] = 0; jb_i[i] = 0; jb_acc[i] = 0; jb_acc2[i] = 0
jb_nmem[i] = 0
return i + 1
}
i += 1
}
return 0
}
function jb_valid(h: int) -> bool {
jb_init()
if (h < 1) or (h > JOB_SLOTS) { return false }
return 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() -> int { return jb_alloc(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(kind: int, arg: int) -> int {
let h = jb_alloc(kind)
if h == 0 { return 0 }
let s = h - 1
jb_arg[s] = arg
if kind == JK_FIB { jb_acc[s] = 0; 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(h: int, value: int) -> void {
if not jb_valid(h) { return }
let s = h - 1
if jb_state[s] != J_PENDING { return }
jb_state[s] = J_DONE
jb_result[s] = value
}
# Resolve a pending job as failed with error code `err` (no-op once resolved).
function job_fail(h: int, err: int) -> void {
if not jb_valid(h) { return }
let s = h - 1
if jb_state[s] != J_PENDING { return }
jb_state[s] = J_FAILED
jb_error[s] = err
}
# Cancel a pending job (no-op if it already resolved).
function job_cancel(h: int) -> void {
if not jb_valid(h) { return }
let s = h - 1
if jb_state[s] == J_PENDING { 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(s: int) -> void {
if jb_state[s] != J_PENDING { return }
let k = jb_kind[s]
if (k != JK_ALL) and (k != JK_RACE) { return }
let n = 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 = jb_mem[base + i]
if jb_valid(mh) {
let ms = mh - 1
let mst = 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 { jb_state[s] = J_DONE; jb_result[s] = n }
else { if settled == n { jb_state[s] = J_FAILED; jb_error[s] = n - ok } }
} else {
if first_ok != 0 { jb_state[s] = J_DONE; jb_result[s] = first_ok }
else { if settled == n { jb_state[s] = J_FAILED; jb_error[s] = n } }
}
}
# resolved in any terminal state?
function job_done(h: int) -> bool {
if not jb_valid(h) { return false }
jb_refresh_group(h - 1)
return jb_state[h - 1] != J_PENDING
}
function job_ok(h: int) -> bool {
if not jb_valid(h) { return false }
jb_refresh_group(h - 1)
return jb_state[h - 1] == J_DONE
}
function job_failed(h: int) -> bool {
if not jb_valid(h) { return false }
jb_refresh_group(h - 1)
return jb_state[h - 1] == J_FAILED
}
function job_cancelled(h: int) -> bool {
if not jb_valid(h) { return false }
return jb_state[h - 1] == J_CANCELLED
}
# the success value (0 unless the job is done-ok)
function job_result(h: int) -> int {
if not jb_valid(h) { return 0 }
jb_refresh_group(h - 1)
if jb_state[h - 1] != J_DONE { return 0 }
return jb_result[h - 1]
}
# the failure code (0 unless the job failed)
function job_error(h: int) -> int {
if not jb_valid(h) { return 0 }
jb_refresh_group(h - 1)
if jb_state[h - 1] != J_FAILED { return 0 }
return jb_error[h - 1]
}
# how many jobs are still pending (a ready-made loading-screen denominator).
function job_pending() -> int {
jb_init()
var n = 0
var i = 0
while i < JOB_SLOTS {
if jb_state[i] == J_PENDING { n += 1 }
i += 1
}
return n
}
# release a slot back to the pool.
function job_free(h: int) -> void {
if not jb_valid(h) { return }
jb_state[h - 1] = J_FREE
}
# advance one compute job by a single step; returns 1 if it just finished.
function jb_step(s: int) -> int {
let k = jb_kind[s]
let n = jb_arg[s]
var i = jb_i[s]
if k == JK_SUM {
jb_acc[s] = jb_acc[s] + (i + 1)
i += 1
jb_i[s] = i
if i >= n { jb_state[s] = J_DONE; jb_result[s] = jb_acc[s]; return 1 }
return 0
}
if k == JK_FIB {
if i >= n { jb_state[s] = J_DONE; jb_result[s] = jb_acc[s]; return 1 }
let t = jb_acc[s] + jb_acc2[s]
jb_acc[s] = jb_acc2[s]
jb_acc2[s] = t
i += 1
jb_i[s] = i
if i >= n { jb_state[s] = J_DONE; jb_result[s] = jb_acc[s]; return 1 }
return 0
}
if k == JK_PRIMES {
if jb_is_prime(i) { jb_acc[s] += 1 }
i += 1
jb_i[s] = i
if i > n { jb_state[s] = J_DONE; jb_result[s] = 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(budget: int) -> int {
jb_init()
var completed = 0
var spent = 0
var s = 0
while s < JOB_SLOTS {
let k = jb_kind[s]
let compute = (k == JK_SUM) or (k == JK_FIB) or (k == JK_PRIMES)
while (jb_state[s] == J_PENDING) and compute {
if (budget > 0) and (spent >= budget) { s = JOB_SLOTS + 1; break }
let fin = jb_step(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 jb_state[s] == J_PENDING {
let k = jb_kind[s]
if (k == JK_ALL) or (k == JK_RACE) {
jb_refresh_group(s)
if 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(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 { jb_mem[base + i] = handles[i]; i += 1 }
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(handles: []int) -> int {
let h = jb_alloc(JK_ALL)
if h == 0 { return 0 }
jb_set_members(h - 1, handles)
jb_refresh_group(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(handles: []int) -> int {
let h = jb_alloc(JK_RACE)
if h == 0 { return 0 }
jb_set_members(h - 1, handles)
jb_refresh_group(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(handles: []int) -> int {
var n = 0
var i = 0
while i < len(handles) {
if job_done(handles[i]) { n += 1 }
i += 1
}
return n
}
function prom_all_done(handles: []int) -> bool {
var i = 0
while i < len(handles) {
if not job_done(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"
var sy_ready: bool = false
var mx_used: words = null
var mx_obj: pointers = null # the native mutex behind each handle
var at_used: words = null
var at_cell: pointers = null # the int each atomic handle names (a words(1) of its own)
var ch_used: words = null
var ch_head: words = null
var ch_count: words = null
var ch_buf: words = null # flat [SYNC_CHAN * CHAN_CAP]
var ch_lock: pointers = null # a mutex per channel
function sy_init() -> void {
if sy_ready { return }
mx_used = words(SYNC_MUTEX); fill(mx_used, 0, SYNC_MUTEX * 4)
mx_obj = pointers(SYNC_MUTEX); fill(mx_obj, 0, SYNC_MUTEX * 8)
at_used = words(SYNC_ATOMIC); fill(at_used, 0, SYNC_ATOMIC * 4)
at_cell = pointers(SYNC_ATOMIC); fill(at_cell, 0, SYNC_ATOMIC * 8)
ch_used = words(SYNC_CHAN); fill(ch_used, 0, SYNC_CHAN * 4)
ch_head = words(SYNC_CHAN); fill(ch_head, 0, SYNC_CHAN * 4)
ch_count = words(SYNC_CHAN); fill(ch_count, 0, SYNC_CHAN * 4)
ch_buf = words(SYNC_CHAN * CHAN_CAP); fill(ch_buf, 0, SYNC_CHAN * CHAN_CAP * 4)
ch_lock = pointers(SYNC_CHAN); fill(ch_lock, 0, SYNC_CHAN * 8)
sy_ready = true
}
# ---- mutex -----------------------------------------------------------------
function sync_mutex() -> int {
sy_init()
var i = 0
while i < SYNC_MUTEX {
if mx_used[i] == 0 {
mx_used[i] = 1
if mx_obj[i] == null { mx_obj[i] = thr_mutex_new() }
return i + 1
}
i += 1
}
return 0
}
function sync_lock(m: int) -> void {
sy_init()
if (m < 1) or (m > SYNC_MUTEX) { return }
thr_lock(mx_obj[m - 1])
}
function sync_unlock(m: int) -> void {
sy_init()
if (m < 1) or (m > SYNC_MUTEX) { return }
thr_unlock(mx_obj[m - 1])
}
# take the lock only if it is free; returns whether it was taken.
function sync_try_lock(m: int) -> bool {
sy_init()
if (m < 1) or (m > SYNC_MUTEX) { return false }
return thr_trylock(mx_obj[m - 1]) == 1
}
# ---- atomic counter --------------------------------------------------------
function sync_atomic() -> int {
sy_init()
var i = 0
while i < SYNC_ATOMIC {
if at_used[i] == 0 {
at_used[i] = 1
if at_cell[i] == null { at_cell[i] = data_of(words(1)) } # the atomics take the cell's address, not its slice
thr_store(at_cell[i], 0)
return i + 1
}
i += 1
}
return 0
}
function sync_get(a: int) -> int {
sy_init()
if (a < 1) or (a > SYNC_ATOMIC) { return 0 }
return thr_load(at_cell[a - 1])
}
function sync_set(a: int, v: int) -> void {
sy_init()
if (a < 1) or (a > SYNC_ATOMIC) { return }
thr_store(at_cell[a - 1], v)
}
# add `delta` and return the new value.
function sync_add(a: int, delta: int) -> int {
sy_init()
if (a < 1) or (a > SYNC_ATOMIC) { return 0 }
return thr_atomic_add(at_cell[a - 1], delta)
}
# compare-and-set: if the value equals `expect`, store `next` and return true.
function sync_cas(a: int, expect: int, next: int) -> bool {
sy_init()
if (a < 1) or (a > SYNC_ATOMIC) { return false }
return thr_cas(at_cell[a - 1], expect, next) == 1
}
# ---- channel (a bounded int FIFO, behind its own mutex) ---------------------
function sync_channel() -> int {
sy_init()
var i = 0
while i < SYNC_CHAN {
if ch_used[i] == 0 {
ch_used[i] = 1; ch_head[i] = 0; ch_count[i] = 0
if ch_lock[i] == null { 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(c: int, v: int) -> bool {
sy_init()
if (c < 1) or (c > SYNC_CHAN) { return false }
let s = c - 1
thr_lock(ch_lock[s])
if ch_count[s] >= CHAN_CAP { thr_unlock(ch_lock[s]); return false }
let pos = ch_head[s] + ch_count[s]
var idx = pos
if idx >= CHAN_CAP { idx -= CHAN_CAP }
ch_buf[s * CHAN_CAP + idx] = v
ch_count[s] += 1
thr_unlock(ch_lock[s])
return true
}
# dequeue the oldest value; returns 0 on an empty channel (guard with can_recv).
function sync_recv(c: int) -> int {
sy_init()
if (c < 1) or (c > SYNC_CHAN) { return 0 }
let s = c - 1
thr_lock(ch_lock[s])
if ch_count[s] == 0 { thr_unlock(ch_lock[s]); return 0 }
let v = ch_buf[s * CHAN_CAP + ch_head[s]]
var nh = ch_head[s] + 1
if nh >= CHAN_CAP { nh = 0 }
ch_head[s] = nh
ch_count[s] -= 1
thr_unlock(ch_lock[s])
return v
}
function sync_can_recv(c: int) -> bool {
sy_init()
if (c < 1) or (c > SYNC_CHAN) { return false }
thr_lock(ch_lock[c - 1])
let has = ch_count[c - 1] > 0
thr_unlock(ch_lock[c - 1])
return has
}
function sync_len(c: int) -> int {
sy_init()
if (c < 1) or (c > SYNC_CHAN) { return 0 }
thr_lock(ch_lock[c - 1])
let n = ch_count[c - 1]
thr_unlock(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 }