Skip to main content

moirai_parallel/
policy.rs

1//! Compile-time execution policies for data-parallel operations.
2//!
3//! A policy decides *whether* a data-parallel operation runs in parallel. The
4//! types are zero-sized **type-level markers** — the decision is an associated
5//! function, so generic code over `P: ExecutionPolicy` monomorphizes to one
6//! concrete path with no value passed and no dynamic dispatch:
7//!
8//! - [`Sequential`] / [`Parallel`] return a constant, so the unused branch is
9//!   eliminated entirely at compile time.
10//! - [`Adaptive`] parallelizes only at or above [`ADAPTIVE_PARALLEL_THRESHOLD`],
11//!   a cheap inlined runtime check that routes per workload size (and thus across
12//!   the worker threads only when worthwhile).
13//!
14//! Select a policy by type via the [`ParallelSlice`](crate::ParallelSlice) /
15//! [`ParallelSliceMut`](crate::ParallelSliceMut) extension traits
16//! (`slice.par_with::<Parallel>()`) or the `*_with::<P>` functions; the `par_*`
17//! helpers and `slice.par()` use [`Adaptive`] as the unset default.
18
19/// Element count at or above which [`Adaptive`] chooses parallel execution.
20///
21/// # This value encodes an assumption about per-element cost
22///
23/// An element count cannot decide this on its own. Parallel wins once
24/// `n * per_element_cost` exceeds the fixed dispatch cost, so the crossover
25/// moves with how expensive the body is, and any single count is right for one
26/// body weight and wrong for the others.
27///
28/// Measured on this workstation (best of 30 blocks, `map_reduce`, parallel
29/// against the same fold run serially; the dispatch floor is ~11.9 us, which
30/// is one task spawned per worker chunk plus the joins):
31///
32/// ```text
33///   body weight                     crossover   parallel/serial at n = 1024
34///   one multiply                    ~21K-32K    20.6x worse
35///   sqrt + ln_1p                    ~8K          4.1x worse
36///   24 chained fused multiply-adds  ~512-1024    0.25x  (parallel wins)
37/// ```
38///
39/// So 1024 is tuned for an expensive body. A caller folding a cheap expression
40/// over 1K-16K elements pays 1.3x to 20.6x for the parallel choice, and the
41/// earlier claim here — that below this value dispatch overhead "typically
42/// exceeds the benefit" — is the opposite of what happens above it for such a
43/// body.
44///
45/// It is left at 1024 deliberately rather than raised: the stack's own heavy
46/// consumers (spherical-harmonic mode loops, for one) fold expensive bodies
47/// over exactly the 1K-16K range where raising it would serialize them. The
48/// two ways out are re-deriving it against a stated body weight, or shrinking
49/// the dispatch floor so the choice matters less; both are tracked rather than
50/// guessed at here.
51///
52/// A caller who knows its body is cheap should select [`Sequential`], and one
53/// who knows it is expensive should select [`Parallel`]. `Adaptive` is for
54/// callers who know neither, and it cannot be right for both.
55pub const ADAPTIVE_PARALLEL_THRESHOLD: usize = 1024;
56
57/// Compile-time strategy selector for the data-parallel operations in this crate.
58///
59/// Implemented by zero-sized marker types; used purely as a type parameter so
60/// each operation monomorphizes to a single concrete path.
61pub trait ExecutionPolicy: Send + Sync + 'static {
62    /// Return `true` if an operation over `len` elements should run in parallel.
63    fn parallelize(len: usize) -> bool;
64
65    /// Return `true` if an operation over `len` elements, partitioned into
66    /// `chunks` logical chunks, should run in parallel.
67    ///
68    /// The default preserves policies expressed only in terms of element
69    /// count. Policies for chunked operations may override this method when
70    /// chunk geometry also determines whether scheduling is worthwhile. An
71    /// operator may coalesce logical chunks into fewer scheduled worker tasks.
72    #[inline(always)]
73    fn parallelize_chunks(len: usize, chunks: usize) -> bool {
74        let _ = chunks;
75        Self::parallelize(len)
76    }
77
78    /// Return `true` if an operation over `len` elements, scheduled as `chunks`
79    /// tasks and moving `bytes` in total, should run in parallel.
80    ///
81    /// The default keeps policies expressed in element and chunk counts. A
82    /// policy whose crossover is a quantity of data rather than a count of
83    /// elements overrides this: a pass over wide elements, or over an output
84    /// beside a wider input, moves more than its element count says (ADR 0059).
85    #[inline(always)]
86    fn parallelize_work(len: usize, chunks: usize, bytes: usize) -> bool {
87        let _ = bytes;
88        Self::parallelize_chunks(len, chunks)
89    }
90
91    /// Return `true` if a fixed two-branch operation should run in parallel.
92    #[inline(always)]
93    fn parallelize_pair() -> bool {
94        Self::parallelize(2)
95    }
96}
97
98/// Always run sequentially (single thread, no scheduling).
99#[derive(Debug, Clone, Copy, Default)]
100pub struct Sequential;
101
102/// Always run in parallel on the shared work-stealing pool.
103#[derive(Debug, Clone, Copy, Default)]
104pub struct Parallel;
105
106/// Run in parallel only for inputs at or above [`ADAPTIVE_PARALLEL_THRESHOLD`].
107#[derive(Debug, Clone, Copy, Default)]
108pub struct Adaptive;
109
110impl ExecutionPolicy for Sequential {
111    #[inline(always)]
112    fn parallelize(_len: usize) -> bool {
113        false
114    }
115}
116
117impl ExecutionPolicy for Parallel {
118    #[inline(always)]
119    fn parallelize(_len: usize) -> bool {
120        true
121    }
122
123    #[inline(always)]
124    fn parallelize_pair() -> bool {
125        true
126    }
127}
128
129impl ExecutionPolicy for Adaptive {
130    #[inline(always)]
131    fn parallelize(len: usize) -> bool {
132        len >= ADAPTIVE_PARALLEL_THRESHOLD
133    }
134}
135
136/// Run in parallel only for inputs at or above the custom threshold `N`.
137#[derive(Debug, Clone, Copy, Default)]
138pub struct AdaptiveWithThreshold<const N: usize>;
139
140impl<const N: usize> ExecutionPolicy for AdaptiveWithThreshold<N> {
141    #[inline(always)]
142    fn parallelize(len: usize) -> bool {
143        len >= N
144    }
145}
146
147/// Run in parallel only for an operation that moves at least `N` bytes.
148///
149/// Byte-reporting operators such as
150/// [`for_each_unit_task_mut_with`](crate::for_each_unit_task_mut_with) decide
151/// through [`ExecutionPolicy::parallelize_work`]. Entry points that report only
152/// an element count are treated as moving one byte per element, a lower bound
153/// for any element type, so they lean serial rather than over-schedule.
154#[derive(Debug, Clone, Copy, Default)]
155pub struct WorkBytes<const N: usize>;
156
157impl<const N: usize> ExecutionPolicy for WorkBytes<N> {
158    #[inline(always)]
159    fn parallelize(len: usize) -> bool {
160        len >= N
161    }
162
163    #[inline(always)]
164    fn parallelize_work(_len: usize, _chunks: usize, bytes: usize) -> bool {
165        bytes >= N
166    }
167}