feat(jobs): real OS threads - Job.parallel_for, fn name, thread-safe Sync
- `fn name` names a top-level function as a value (E_FNREF, lowers to @fn_<name>); the worker entry point for Job.parallel_for, which checks it takes (int, pointer-like) and returns void. - runtime/native/threads.ll (pthreads) and threads_win.ll (Win32 SRWLOCK/CONDITION_VARIABLE): a pool of one worker per core but one, parked between batches; every thread claims chunks by compare-and-swap. Linked only into programs that use Job/Promise/Sync, by `ludicc -o`, `ludic build` and the test suite's build helper. - Sync.* is real: native mutexes, atomics as cmpxchg retry loops (neither clang takes atomicrw, the PC's rejects seq_consistent), mutex-guarded channels, Sync.cpu_count from the OS. - spawn/despawn on a pool thread stop the program with a located panic. - examples/library/threads.ludic and its test; docs for fn, Job.parallel_for, Job.is_worker. - Reseeded (bootstrap-cfree: out.ll == seed.ll). 141/141 on macOS; jobs, threads and the guard pass on Windows from the reseeded Windows seed. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
parent
f5d1a62ccf
commit
c10abd9f9f
22 changed files with 57464 additions and 55794 deletions
|
|
@ -386,35 +386,56 @@ 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_held: words = null
|
||||
var mx_obj: pointers = null # the native mutex behind each handle
|
||||
var at_used: words = null
|
||||
var at_val: 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_held = words(SYNC_MUTEX); fill(mx_held, 0, SYNC_MUTEX * 4)
|
||||
mx_obj = bytes(SYNC_MUTEX * 8); fill(mx_obj, 0, SYNC_MUTEX * 8)
|
||||
at_used = words(SYNC_ATOMIC); fill(at_used, 0, SYNC_ATOMIC * 4)
|
||||
at_val = words(SYNC_ATOMIC); fill(at_val, 0, SYNC_ATOMIC * 4)
|
||||
at_cell = bytes(SYNC_ATOMIC * 8); 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 = bytes(SYNC_CHAN * 8); fill(ch_lock, 0, SYNC_CHAN * 8)
|
||||
sy_ready = true
|
||||
}
|
||||
|
||||
# ---- mutex (a cooperative lock) --------------------------------------------
|
||||
# ---- mutex -----------------------------------------------------------------
|
||||
function sync_mutex() -> int {
|
||||
sy_init()
|
||||
var i = 0
|
||||
while i < SYNC_MUTEX {
|
||||
if mx_used[i] == 0 { mx_used[i] = 1; mx_held[i] = 0; return i + 1 }
|
||||
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
|
||||
|
|
@ -423,22 +444,20 @@ function sync_mutex() -> int {
|
|||
function sync_lock(m: int) -> void {
|
||||
sy_init()
|
||||
if (m < 1) or (m > SYNC_MUTEX) { return }
|
||||
mx_held[m - 1] = 1
|
||||
thr_lock(mx_obj[m - 1])
|
||||
}
|
||||
|
||||
function sync_unlock(m: int) -> void {
|
||||
sy_init()
|
||||
if (m < 1) or (m > SYNC_MUTEX) { return }
|
||||
mx_held[m - 1] = 0
|
||||
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 }
|
||||
if mx_held[m - 1] != 0 { return false }
|
||||
mx_held[m - 1] = 1
|
||||
return true
|
||||
return thr_trylock(mx_obj[m - 1]) == 1
|
||||
}
|
||||
|
||||
# ---- atomic counter --------------------------------------------------------
|
||||
|
|
@ -446,7 +465,12 @@ function sync_atomic() -> int {
|
|||
sy_init()
|
||||
var i = 0
|
||||
while i < SYNC_ATOMIC {
|
||||
if at_used[i] == 0 { at_used[i] = 1; at_val[i] = 0; return i + 1 }
|
||||
if at_used[i] == 0 {
|
||||
at_used[i] = 1
|
||||
if at_cell[i] == null { at_cell[i] = words(1) }
|
||||
thr_store(at_cell[i], 0)
|
||||
return i + 1
|
||||
}
|
||||
i += 1
|
||||
}
|
||||
return 0
|
||||
|
|
@ -455,38 +479,39 @@ function sync_atomic() -> int {
|
|||
function sync_get(a: int) -> int {
|
||||
sy_init()
|
||||
if (a < 1) or (a > SYNC_ATOMIC) { return 0 }
|
||||
return at_val[a - 1]
|
||||
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 }
|
||||
at_val[a - 1] = v
|
||||
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 }
|
||||
at_val[a - 1] += delta
|
||||
return at_val[a - 1]
|
||||
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 }
|
||||
if at_val[a - 1] != expect { return false }
|
||||
at_val[a - 1] = next
|
||||
return true
|
||||
return thr_cas(at_cell[a - 1], expect, next) == 1
|
||||
}
|
||||
|
||||
# ---- channel (a bounded int FIFO) ------------------------------------------
|
||||
# ---- 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; return i + 1 }
|
||||
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
|
||||
|
|
@ -497,12 +522,14 @@ function sync_send(c: int, v: int) -> bool {
|
|||
sy_init()
|
||||
if (c < 1) or (c > SYNC_CHAN) { return false }
|
||||
let s = c - 1
|
||||
if ch_count[s] >= CHAN_CAP { return false }
|
||||
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
|
||||
}
|
||||
|
||||
|
|
@ -511,27 +538,43 @@ function sync_recv(c: int) -> int {
|
|||
sy_init()
|
||||
if (c < 1) or (c > SYNC_CHAN) { return 0 }
|
||||
let s = c - 1
|
||||
if ch_count[s] == 0 { return 0 }
|
||||
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 }
|
||||
return ch_count[c - 1] > 0
|
||||
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 }
|
||||
return ch_count[c - 1]
|
||||
thr_lock(ch_lock[c - 1])
|
||||
let n = ch_count[c - 1]
|
||||
thr_unlock(ch_lock[c - 1])
|
||||
return n
|
||||
}
|
||||
|
||||
# worker lanes available to the scheduler. One today (the deterministic main
|
||||
# thread); a future OS-thread backend would report the real core count here.
|
||||
function sync_cpu_count() -> int { return 1 }
|
||||
# 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 }
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue