Skip to content

Thread pool and tlparallelfor (optional C)

Modulert.03 · side · C · Pass 6 · 3 to 4 h
You buildc/src/runtime/pool.c: tl_pool_create, tl_parallel_for, tl_pool_threads, tl_pool_destroy, and the helpers run_ranges, worker_main, shutdown_pool (the struct is given)
Contractcourse/contracts/c/include/tinyllm/pool.h · rules: c/ABI.md (pthreads only; batch invariance, rule 10)
Testscourse/tests/rt.03/test_pool.c (C, built twice: ASan and UBSan with the counting allocator, then ThreadSanitizer) (what they check: section 4)
Needsrt.02 runtime support: tl_alloc, tl_free, tl_set_last_error (or --ref-deps). Reading: lang.03 C
Used byL9.1, L9.3, L9.4, and L9.5 use the pool in their optional standalone C implementations
MilestoneMS-L9, the optional C module group
Optional depthHerlihy and Shavit, The Art of Multiprocessor Programming, ch. 16 (work distribution); Butenhof, Programming with POSIX Threads, ch. 3 and 7; the OpenMP specification, schedule(dynamic, chunk)
  • The partition is part of the contract: ranges are always [kg,min⁡((k+1)g,n))[kg, \min((k+1)g, n)), whoever runs them, so a kernel that reduces only inside a range gives the same bits with 1 or 8 threads (pooled_matmul_is_bitwise_serial).
  • One atomic counter hands out ranges: atomic_fetch_add gives each range to exactly one thread; a read followed by a write gives some ranges twice (each_index_runs_exactly_once, under ThreadSanitizer).
  • “Done” means every thread that woke has finished, not that the counter ran out; returning early lets the next job overwrite one still being read (each_index_runs_exactly_once).
  • A nested call on the same pool is TL_EBUSY, not a deadlock, and calls from two threads are serialized (nested_call_is_busy_not_a_deadlock, two_callers_are_serialized).
  • Destroy waits for a running call and joins every thread; a failed create leaks nothing (destroy_waits_for_running_work, create_failure_leaks_nothing).
Terminal window
ol start rt.03 # stubs pool.c into your repo (the struct is given)
ol tests rt.03 # read the test catalog first
ol check rt.03 # exit code is the verdict (ASan build, then TSan build)
OL_TSAN=0 ol check rt.03 # skip the ThreadSanitizer build while iterating
ol diff rt.03 # after passing: your code against the reference

Every test starts a 10-second watchdog: a deadlock reports FAIL (no progress in 10 s: a deadlock?) and the test’s name instead of hanging until the harness timeout.


Your C matmul (M03.1) runs on one core. The optional tiled kernel of L9.1 can split independent columns across workers; attention (L9.3) splits heads, and other optional C exercises parallelize rows or sequences. Each implementation still owns a fixed reduction order, so scheduling changes which worker computes a result but does not change the result’s bits. Starting threads per call costs tens of microseconds, so these standalone C exercises keep a pool that sleeps between calls. The production Rust engine owns its own Rust thread pool and independent kernels.

SymbolMeaningType / shape
nnthe number of indices to process, n≥0n \ge 0int64_t
ggthe grain: indices per range, g≥1g \ge 1int64_t
R=⌈n/g⌉R = \lceil n / g \rceilthe number of rangesint64_t
kka range index, 0≤k<R0 \le k < Rint64_t
TTthe number of workers, counting the callerint
wwa worker id, 0≤w<T0 \le w < T; the caller is 0int

2.1 Threads, mutexes, condition variables, atomics

Section titled “2.1 Threads, mutexes, condition variables, atomics”

A thread runs a function concurrently with the others in the same address space (pthread_create, pthread_join; c/ABI.md rule 8: pthreads only, since <threads.h> is missing on macOS). Two threads touching the same memory, at least one writing, with nothing ordering the two accesses, is a data race, which is undefined behavior in C11. Three tools order accesses: a mutex (pthread_mutex_lock/unlock: one holder at a time, and everything written before an unlock is visible after the next lock), a condition variable (pthread_cond_wait releases the mutex and sleeps until another thread signals; always re-check the condition in a loop, because wakeups can be spurious), and an atomic (<stdatomic.h>: atomic_fetch_add reads, adds, and writes as one indivisible step). ThreadSanitizer (-fsanitize=thread) instruments every access and reports any pair not ordered by these tools.

tl_parallel_for(p, n, g, fn, ctx) calls fn(ctx, lo, hi, w) once for each range

[ kg, min⁡((k+1)g, n) )for k=0,1,…,R−1,R=⌈n/g⌉=⌊(n−1)/g⌋+1 (n>0).[\,k g,\ \min((k + 1) g,\ n)\,) \quad \text{for } k = 0, 1, \dots, R - 1,\quad R = \lceil n/g \rceil = \lfloor (n - 1)/g \rfloor + 1 \ (n > 0).

The second form of RR avoids overflow in n+g−1n + g - 1. The last range is short when gg does not divide nn. Nothing about the ranges depends on TT: only which worker runs which range, and in which order, does. A kernel that writes each output element inside one range, from a reduction in a fixed order, is therefore bitwise identical for every TT. The obvious alternative, “split nn into TT equal chunks”, changes the ranges with TT and breaks this.

The pool has TT workers: the caller is worker 0 and T−1T - 1 pthreads are workers 1…T−11 \dots T - 1. A job publishes (n,g,R,fn,ctx)(n, g, R, \mathit{fn}, \mathit{ctx}) and resets an atomic counter next to 0. Every participant then loops:

for (;;) { int64_t k = atomic_fetch_add(&next, 1); if (k >= R) break; run range k; }

atomic_fetch_add returns each value of the counter to exactly one caller, so each range runs exactly once, and fast workers simply take more ranges (dynamic scheduling, which balances uneven ranges). Writing it as k = next; next = k + 1; lets two threads read the same kk: a race, and a range run twice (ThreadSanitizer reports it; the per-index counters of each_index_runs_exactly_once see a 2).

2.4 Waking, finishing, and reusing the job

Section titled “2.4 Waking, finishing, and reusing the job”

Workers sleep on a condition variable wake and wait for a generation number to change. tl_parallel_for locks the mutex, writes the job, sets active to the number of pthreads, bumps the generation, broadcasts, unlocks, and runs ranges itself. Each pthread that woke runs ranges until the counter passes RR, then, under the mutex, decrements active; the one that reaches 0 signals idle. The caller returns only after waiting for active == 0.

Why wait for active, when the counter alone says no range is left? Because “no range left to hand out” is not “no range still running”, and a worker that took the last range may still be inside fn. Worse, a worker that is between waking and reading the job would read the next job’s fields if the caller returned and was called again. Waiting for every woken thread to check in is what makes the job struct safe to overwrite.

2.5 Nested calls, concurrent callers, destroy

Section titled “2.5 Nested calls, concurrent callers, destroy”

Nested. If fn calls tl_parallel_for on the same pool, the inner call would wait for workers that are busy running the outer call’s ranges, forever. A thread-local pointer records which pool’s ranges the current thread is running; a call on that pool returns TL_EBUSY. A 1-thread pool has no pthreads but must refuse too, because the contract does not depend on TT. A call with p = NULL (serial) inside a range is fine.

Concurrent callers. Two threads calling on one pool would overwrite each other’s job. A second mutex, call_lock, held for the whole call, makes the second wait for the first.

Destroy. tl_pool_destroy first takes and releases call_lock, which waits for a call running on another thread to return; then it sets stop, broadcasts, joins every pthread, destroys the mutexes and condition variables, and frees through the hook. Joining the workers alone is not enough: the caller of the running call (worker 0) still reads the pool after the workers have run out of ranges.

Create. All allocations (the pool, the thread table, the thread arguments) happen on the creating thread, so the allocator hook never needs to be thread-safe. If a pthread_create fails, the threads already started are stopped and joined, everything is freed, and the result is TL_EIO; a failed allocation gives TL_ENOMEM. n_threads == 0 means the hardware concurrency, sysconf(_SC_NPROCESSORS_ONLN).

n=10n = 10, g=4g = 4: R=⌊9/4⌋+1=3R = \lfloor 9 / 4 \rfloor + 1 = 3.

kkrangelength
0[0,4)[0, 4)4
1[4,8)[4, 8)4
2[8,min⁡(12,10))=[8,10)[8, \min(12, 10)) = [8, 10)2 (the short last range)

With p = NULL the three calls run in this order on the caller with worker = 0. With a 4-thread pool, one possible run: the caller takes k=0k = 0, worker 2 takes k=1k = 1, worker 1 takes k=2k = 2, worker 3 finds the counter at 3 ≥R\ge R and takes nothing. Another run gives other workers other ranges, but the set of three ranges is always the table above, which is what hand_example checks (sorted by lo) on 4 threads after checking the serial order exactly.

For the bitwise claim: a 37×21137 \times 211 by 211×29211 \times 29 product where range [i0,i1)[i_0, i_1) computes rows i0…i1−1i_0 \dots i_1 - 1, each element summed over kk from 0 to 210 in order. Whether row 5 is computed by worker 0 or worker 3, the same 211 products are added in the same order, so pooled_matmul_is_bitwise_serial compares with memcmp.

typedef struct tl_pool tl_pool;
typedef void (*tl_range_fn)(void *ctx, int64_t lo, int64_t hi, int worker);
tl_status tl_pool_create(int n_threads, tl_pool **out); /* 0 = hardware concurrency; TL_EINVAL, TL_ENOMEM, TL_EIO */
tl_status tl_parallel_for(tl_pool *p, int64_t n, int64_t grain, tl_range_fn fn, void *ctx);
/* blocks; p == NULL is serial; TL_EINVAL, TL_EBUSY (nested) */
int tl_pool_threads(const tl_pool *p); /* T; 1 for NULL */
void tl_pool_destroy(tl_pool *p); /* waits for a running call; NULL is a no-op */

worker is stable during one call of fn and is below tl_pool_threads(p), so a kernel can index per-worker scratch with it (one rt.02 arena per worker, for instance).

TestKINDChecksWhy it matters downstream
hand_exampleunit, smokesection 3: serial order and worker 0 with p = NULL; the same three ranges on 4 threadsyou and the test agree on the partition
grain_edges_and_empty_rangeboundaryg=1g = 1, g=ng = n, g>ng > n, a short last range, n=0n = 0 (no call), on 1 and 3 threadsevery kernel’s edge shapes
bad_argumentsboundaryn<0n < 0, g<1g < 1, fn == NULL, T<0T < 0, out == NULL; no range runsdefined behavior at the boundary
each_index_runs_exactly_onceproperty200 seeded (n,g)(n, g) on 8 threads: an atomic counter per index ends at exactly 1no lost or doubled rows
pooled_matmul_is_bitwise_serialdifferentiala row-range matmul on 1, 2, 3, 8 threads and grains 1 and 5 equals the serial one under memcmpbatch invariance (c/ABI.md rule 10)
all_workers_run_concurrentlyunit4 ranges on 4 threads each wait for all 4 to start; worker ids 0 to 3 each seen oncereal parallelism, correct worker ids
nested_call_is_busy_not_a_deadlockfaulta range calling on its own pool gets TL_EBUSY (1 and 3 threads); p = NULL inside worksno deadlock in composed kernels
two_callers_are_serializedfaulttwo threads, 50 calls each of 1000 indices: the total is 100000; TSan cleanthe engine and a benchmark sharing a pool
destroy_waits_for_running_workfaultdestroy during another thread’s call (whose caller is still inside a 100 ms range) returns after all 16 rangesclean engine shutdown
create_failure_leaks_nothingfaultthe n-th allocation fails for each n: TL_ENOMEM, out == NULL, no leak; success still worksevery constructor’s failure path
hardware_concurrency_defaultunitT=0T = 0 resolves to at least 1 and runs the exact partition[engine].threads unset
PitfallSymptomCaught by
1. taking a range with a plain read and write of the countersome ranges run twice, others never; TSan reports the raceeach_index_runs_exactly_once (mutant s01)
2. splitting nn into TT chunks instead of by graindifferent sums with different thread counts; the edge partitions failhand_example, grain_edges_and_empty_range (mutant s02)
3. R=⌊n/g⌋R = \lfloor n/g \rfloorthe short last range is droppedgrain_edges_and_empty_range (mutant s03)
4. returning when the counter runs out, not when the workers finisha range still running after return; the next job is read half-writteneach_index_runs_exactly_once (mutant s05)
5. not resetting the counter for each jobthe second call runs nothinggrain_edges_and_empty_range (mutant s06)
6. the caller doing all the work, or every worker reporting id 0no speedup; per-worker scratch shared by allall_workers_run_concurrently (mutants s07, s08)
7. no nested-call check, or none for the 1-thread poolthe inner call waits forevernested_call_is_busy_not_a_deadlock (mutants s09, s10)
8. no lock between two callersone caller’s job overwrites the other’stwo_callers_are_serialized (mutant s11)
9. destroy that only joins the workers, or detaches themthe pool freed under a running call (ASan, TSan)destroy_waits_for_running_work (mutants s12, s13)
10. a failed create that frees only part of what it allocateda leak per failed createcreate_failure_leaks_nothing (mutant s14)
DirectionModuleHow it uses this
Backrt.02tl_alloc/tl_free for the pool and its tables (counted in tests), tl_set_last_error for failures
ForwardL9.1tl_matmul_f32(..., tp) splits rows of C into ranges; each output element is reduced in one range
ForwardL9.3FlashAttention gives each (batch, head) pair to a range; worker ids index per-worker scratch arenas (rt.02)
ForwardL9.4paged decode attention splits the batch’s sequences and heads into ranges
ForwardL9.5the int4 and int8 matmuls split output rows into ranges
ForwardL9.1, L9.3, L9.4, L9.5Optional standalone C kernels use the pool to schedule independent output regions

If you skip this module, ol check L9.1 stops with BLOCKED ... needs rt.03: build it, or pass --ref-deps.

Your pieceProduction equivalentWhat it addsWhere to look
tl_parallel_forOpenMP parallel for schedule(dynamic, g), TBB parallel_forwork stealing between per-thread deques, nested parallelism without deadlockoneTBB include/oneapi/tbb/parallel_for.h
fixed partitionggml’s threadpool and ggml_compute_forward_mul_matsplits rows by ith/nth per thread, with a barrier between graph nodesggml/src/ggml-cpu/ggml-cpu.c
sleeping workersRust rayonwork-stealing scheduler, join and parallel iterators, a global poolthe rayon rayon-core crate
the counterPyTorch at::parallel_fora grain size to avoid overhead on small loops, thread-local intra-op poolsaten/src/ATen/Parallel.h