Skip to main content

vole_document/
parallel.rs

1//! Bounded CPU parallelism for the naturally independent ingest work
2//! (Phase 15.4).
3//!
4//! The pool is compiled **only** with the non-default `parallel` feature (which
5//! implies `field`). Without that feature [`WorkerPool`] is a zero-sized marker
6//! that is never constructed, and every ingest site takes its serial path; the
7//! exact core therefore keeps no new dependency and no behavior change.
8//!
9//! ## The one architectural rule
10//!
11//! A worker may only run a **pure** map over immutable inputs. It must never
12//! touch a [`FieldStore`](crate::field::FieldStore), its `Rc`-based
13//! [`IoCounters`](crate::store::IoCounters), a derived-cache memo, or the disk.
14//! The store's single-threaded accounting is the fuse that keeps parallelism
15//! honest: the parallel phase computes owned canonical bytes and [`NodeId`]s,
16//! and a *serial* fold merges them in physical order.
17//!
18//! Determinism is therefore structural, not a tuning property. Indexed
19//! `par_iter().collect()` preserves physical order, and every order-dependent
20//! gate (the running `MAX_TOTAL_DECODED` / `MAX_OBJSTM*` caps, the existence
21//! probes and counters) is replayed in that same serial fold — so the parallel
22//! path yields byte-identical nodes, ids, index entries, and limits enforcement
23//! to the serial path.
24//!
25//! [`NodeId`]: crate::store::NodeId
26
27#[cfg(feature = "parallel")]
28use crate::error::{Error, Result};
29
30/// A bounded worker pool. `workers == 1` is a valid (serial) pool, though the
31/// ingest sites take their serial branch directly rather than paying to build it.
32#[cfg(feature = "parallel")]
33pub struct WorkerPool {
34    pool: rayon::ThreadPool,
35    workers: usize,
36}
37
38#[cfg(feature = "parallel")]
39impl WorkerPool {
40    /// Build a pool of exactly `workers` threads, or return a typed error.
41    ///
42    /// `workers == 0` is rejected here: the "auto" policy (`available_parallelism`)
43    /// is resolved explicitly by the CLI, never silently by the library.
44    pub fn new(workers: usize) -> Result<Self> {
45        if workers == 0 {
46            return Err(Error::usage("a worker pool needs at least one worker"));
47        }
48        let pool = rayon::ThreadPoolBuilder::new()
49            .num_threads(workers)
50            .thread_name(|i| format!("vole-worker-{i}"))
51            .build()
52            .map_err(|e| Error::internal_invariant(format!("worker pool build failed: {e}")))?;
53        Ok(WorkerPool { pool, workers })
54    }
55
56    /// The number of threads in the pool.
57    pub fn workers(&self) -> usize {
58        self.workers
59    }
60
61    /// Run `f` on the pool and block until it returns. The rayon thread-local
62    /// context is scoped to `f`.
63    pub fn install<R: Send>(&self, f: impl FnOnce() -> R + Send) -> R {
64        self.pool.install(f)
65    }
66}
67
68/// Without the `parallel` feature this is a zero-sized marker that is never
69/// constructed; the ingest sites keep their serial path (`pool` is always `None`).
70#[cfg(not(feature = "parallel"))]
71pub struct WorkerPool;