Skip to main content

moirai_parallel/
lib.rs

1//! Synchronous data-parallel primitives — Moirai's rayon-replacement surface.
2//!
3//! This crate is the **parallel** domain (throughput over data), distinct from
4//! the **concurrent** domain (`moirai-async`, async tasks/IO). All operations
5//! here are fully synchronous (no `async`, no `.await`), so they are safe inside
6//! pure compute kernels without async contagion, and operate on borrowed slices
7//! with in-place mutation (zero-copy).
8//!
9//! # Selecting an execution strategy
10//!
11//! Strategy is a zero-sized [`ExecutionPolicy`] type ([`Sequential`],
12//! [`Parallel`], [`Adaptive`]) chosen at compile time, so every form below
13//! monomorphizes with no dynamic dispatch:
14//!
15//! - **Extension traits** (the surface) — `slice.par()` / `slice.par_mut()`
16//!   return [`Adaptive`] handles, then `for_each` / `enumerate` / `map_collect` /
17//!   `map_reduce`:
18//!   ```
19//!   use moirai_parallel::{ParallelSlice, ParallelSliceMut};
20//!   let v: Vec<u64> = (0..1000).collect();
21//!   let sum = v.par().map_reduce(0, |&x| x, |a, b| a + b);   // auto-routes
22//!   let mut m = v.clone();
23//!   m.par_mut().for_each(|x| *x += 1);
24//!   ```
25//! - **`*_with::<P>` free functions** — a low-level override that pins the policy
26//!   via turbofish (`for_each_with::<Sequential>(&data, f)`), for the rare case
27//!   that needs to force sequential (determinism / nested regions) or parallel.
28//!   Most code should just use `.par()`.
29//!
30//! Because [`Adaptive`] is itself a zero-sized policy, `.par()` is a fully
31//! monomorphized, zero-cost abstraction that parallelizes only at or above
32//! [`ADAPTIVE_PARALLEL_THRESHOLD`] and runs sequentially below it — the
33//! parallel/sequential decision is automatic, with nothing to designate.
34//!
35//! These data-parallel ops are synchronous (they return values, not futures),
36//! but they run on the **same unified hybrid scheduler** as async work
37//! ([`moirai_executor::global`]) — not a separate pool. A `.par()` worker task
38//! can therefore spawn or drive async work (`moirai::global().spawn_async`/
39//! `block_on`) on that same runtime, so parallel processing and asynchronous
40//! tasks compose within one process. The sync return shape here is a property of
41//! the *operation*, not an isolation boundary.
42
43#![deny(missing_docs)]
44#![deny(unsafe_op_in_unsafe_fn)]
45
46mod policy;
47
48pub use policy::{
49    ADAPTIVE_PARALLEL_THRESHOLD, Adaptive, AdaptiveWithThreshold, ExecutionPolicy, Parallel,
50    Sequential, WorkBytes,
51};
52
53use core::marker::PhantomData;
54
55/// Synchronous data-parallel operators and free functions.
56pub mod ops;
57pub use ops::{
58    ChunkBuffersError, Scope, UNIT_TASK_BYTES, enumerate_mut_with, enumerate_with,
59    fold_reduce_with, for_each_chunk_buffers_mut_enumerated_with,
60    for_each_chunk_mut_enumerated_with, for_each_chunk_mut_with, for_each_chunk_mut_with_state,
61    for_each_chunk_pair_mut_enumerated_with, for_each_chunk_quad_mut_enumerated_with,
62    for_each_chunk_triple_mut_enumerated_with, for_each_index_with, for_each_mut_with,
63    for_each_unit_task_many_mut_with, for_each_unit_task_mut_with,
64    for_each_unit_task_pair_mut_with, for_each_unit_task_range_with,
65    for_each_unit_task_triple_mut_with, for_each_with, join, join_with, map_collect_index_with,
66    map_collect_mut_with, map_collect_with, map_reduce_with, reduce_index_with, scope,
67    units_per_task,
68};
69
70// ---------------------------------------------------------------------------
71// Extension traits: trait-based, type-selected parallel views over slices
72// ---------------------------------------------------------------------------
73
74/// A read-only parallel view of a slice bound to execution policy `P`.
75///
76/// Construct via [`ParallelSlice::par`]. Zero-sized beyond the borrowed slice.
77pub struct ParRef<'a, T, P> {
78    data: &'a [T],
79    _policy: PhantomData<P>,
80}
81
82impl<'a, T, P: ExecutionPolicy> ParRef<'a, T, P> {
83    /// Apply `f` to every element. See [`for_each_with`].
84    pub fn for_each<F: Fn(&T) + Send + Sync>(self, f: F)
85    where
86        T: Sync,
87    {
88        for_each_with::<P, _, _>(self.data, f);
89    }
90
91    /// Apply `f(index, &element)` to every element. See [`enumerate_with`].
92    pub fn enumerate<F: Fn(usize, &T) + Send + Sync>(self, f: F)
93    where
94        T: Sync,
95    {
96        enumerate_with::<P, _, _>(self.data, f);
97    }
98
99    /// Map then collect into a `Vec<R>` in order. See [`map_collect_with`].
100    pub fn map_collect<R: Send, F: Fn(&T) -> R + Send + Sync>(self, f: F) -> Vec<R>
101    where
102        T: Sync,
103    {
104        map_collect_with::<P, _, _, _>(self.data, f)
105    }
106
107    /// Map each `(index, &element)` pair, then collect into a `Vec<R>` in order.
108    ///
109    /// This is the indexed counterpart to [`Self::map_collect`]. It exposes the
110    /// logical slice index without requiring callers to derive it from pointers.
111    pub fn map_collect_index<R: Send, F: Fn(usize, &T) -> R + Send + Sync>(self, f: F) -> Vec<R>
112    where
113        T: Sync,
114    {
115        map_collect_index_with::<P, _, _>(self.data.len(), |i| f(i, &self.data[i]))
116    }
117
118    /// Map-reduce. See [`map_reduce_with`].
119    pub fn map_reduce<R, M, Rd>(self, identity: R, map: M, reduce: Rd) -> R
120    where
121        T: Sync,
122        R: Send + Sync + Clone,
123        M: Fn(&T) -> R + Send + Sync,
124        Rd: Fn(R, R) -> R + Send + Sync,
125    {
126        map_reduce_with::<P, _, _, _, _>(self.data, identity, map, reduce)
127    }
128}
129
130/// A mutable parallel view of a slice bound to execution policy `P`.
131pub struct ParMut<'a, T, P> {
132    data: &'a mut [T],
133    _policy: PhantomData<P>,
134}
135
136impl<'a, T, P: ExecutionPolicy> ParMut<'a, T, P> {
137    /// Apply `f` to every element in place. See [`for_each_mut_with`].
138    pub fn for_each<F: Fn(&mut T) + Send + Sync>(self, f: F)
139    where
140        T: Send,
141    {
142        for_each_mut_with::<P, _, _>(self.data, f);
143    }
144
145    /// Apply `f(index, &mut element)` to every element in place.
146    pub fn enumerate<F: Fn(usize, &mut T) + Send + Sync>(self, f: F)
147    where
148        T: Send,
149    {
150        enumerate_mut_with::<P, _, _>(self.data, f);
151    }
152}
153
154/// Extension trait providing an adaptive parallel view over `&[T]`.
155pub trait ParallelSlice<T> {
156    /// Adaptive, auto-routing parallel view (the everyday entry point).
157    fn par(&self) -> ParRef<'_, T, Adaptive>;
158}
159
160impl<T> ParallelSlice<T> for [T] {
161    #[inline]
162    fn par(&self) -> ParRef<'_, T, Adaptive> {
163        ParRef {
164            data: self,
165            _policy: PhantomData,
166        }
167    }
168}
169
170/// Extension trait providing an adaptive mutable parallel view over `&mut [T]`.
171pub trait ParallelSliceMut<T> {
172    /// Adaptive, auto-routing mutable parallel view (the everyday entry point).
173    fn par_mut(&mut self) -> ParMut<'_, T, Adaptive>;
174}
175
176impl<T> ParallelSliceMut<T> for [T] {
177    #[inline]
178    fn par_mut(&mut self) -> ParMut<'_, T, Adaptive> {
179        ParMut {
180            data: self,
181            _policy: PhantomData,
182        }
183    }
184}
185
186pub mod melinoe_ext;
187
188#[cfg(test)]
189#[path = "tests.rs"]
190mod tests;