Creates a shared-memory task pool and spawns its worker processes. Per-submitter injection rings feed the workers, and results are published through result slots. A submission is one SHM ring write plus at most one directed wake: no dispatcher process is in the loop. Each worker owns a work-stealing deque. Idle workers steal from busy peers and consume the injection rings. A fairness tick bounds the latency of external submissions on a saturated pool.
Usage
mizu_pool(
n_workers = 1L,
max_workers = n_workers,
max_submitters = 8L,
injection_cap = 1024L,
per_worker_cap = 1024L,
result_slots = 4096L,
slot_size = 512L,
launcher = mizu_launcher(),
startup_timeout = 30
)Arguments
- n_workers
number of worker processes to spawn, at most
max_workers.- max_workers
worker registry capacity (at most 64).
- max_submitters
submitter registry capacity (at most 64). Each submitter owns its own injection ring and an equal share of
result_slots.- injection_cap
entries per submitter injection ring. A power of two.
- per_worker_cap
entries per worker work-stealing deque. A power of two.
- result_slots
total result slots, partitioned equally across the submitter slots and rounded up to a multiple of
max_submitters. This bounds the outstanding (uncollected) tasks of each submitter.- slot_size
bytes per queue entry and result slot. A power of two between 128 and 2^20. A payload (task or result) that serializes past the inline budget travels in a fresh region per payload. This is an order-of-magnitude latency cliff, surfaced per submitter as
mizu_pool_stats()$submitters$spills. The default512Lkeeps typical expression-plus-arguments tasks inline. Pools that move only scalar payloads can drop to256L.- launcher
a
function(token, slot)that arranges for an R process to callmizu:::worker_main(token, slot). The defaultmizu_launcher()spawnsRscriptand propagates the.libPaths()of the host. Itsstdoutandstderrarguments direct the worker output. A custom launcher must arrange the library paths itself.- startup_timeout
seconds to wait for all workers to join. On expiry, mizu destroys the pool and raises
mizu_error_startup(see mizu_error).
Value
A pool handle (class "mizu_pool") holding submitter slot 0.
Handles are process-private and do not survive fork().
Details
The lifetime of the pool is bound to the creating process, which holds
submitter slot 0 of the returned handle. Use this handle directly with
mizu_submit() and mizu_collect(). Other processes join as submitters
through mizu_pool_attach(). Dropping the handle (or exiting R) shuts
the pool down as mizu_pool_stop() does, but without the wait.
Payload contents interoperate transparently with mori. A mori::share()d
object anywhere inside the arguments of a task or its result serializes
to its short identifier wire form through the mori hooks. It maps
zero-copy on the other side.
The liveness lock files of the pool (the death-detection verdict) live
in a per-platform directory chosen at create time. This is /dev/shm on
Linux, and the per-user temporary directory on macOS and Windows. The
chosen path is recorded in the region, so every participant uses the
same files. The environment variable MIZU_LIVENESS_DIR, read in the
creating process, overrides the default.
Examples
p <- mizu_pool()
t <- mizu_submit(p, x + y, x = 1, y = 2)
mizu_collect(t)
#> [1] 3
mizu_pool_stop(p)