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, AdaptiveWithThreshold, ExecutionPolicy, Parallel, Sequential,
50    ADAPTIVE_PARALLEL_THRESHOLD,
51};
52
53use core::marker::PhantomData;
54
55/// Pointer wrapper used to hand disjoint `&mut` sub-slices to worker tasks.
56///
57/// The `Send`/`Sync` impls are sound only because the `*_mut` operations assign
58/// each task a non-overlapping index range, so the pointer is never used to form
59/// aliasing references.
60pub(crate) struct DisjointMutPtr<T>(pub(crate) *mut T);
61
62// SAFETY: callers dereference pairwise-disjoint ranges only, so the pointer
63// never forms aliasing `&mut` references; `T: Send` permits moving element
64// access across worker threads.
65unsafe impl<T: Send> Send for DisjointMutPtr<T> {}
66unsafe impl<T: Send> Sync for DisjointMutPtr<T> {}
67
68impl<T> DisjointMutPtr<T> {
69    /// Return a `&mut` to element `i`.
70    ///
71    /// # Safety
72    /// `i` must be in bounds and visited at most once across all concurrent
73    /// tasks, so the returned reference never aliases another.
74    #[inline]
75    pub(crate) unsafe fn get_mut<'a>(&self, i: usize) -> &'a mut T {
76        // SAFETY: guaranteed by the caller's per-index-once contract.
77        unsafe { &mut *self.0.add(i) }
78    }
79
80    /// Return the wrapped base pointer. Taking `&self` forces a closure to
81    /// capture the whole (`Send`/`Sync`) wrapper rather than the bare `*mut T`
82    /// field under 2021 disjoint capture.
83    #[inline]
84    pub(crate) fn base(&self) -> *mut T {
85        self.0
86    }
87}
88/// Synchronous data-parallel operators and free functions.
89pub mod ops;
90pub use ops::{
91    enumerate_mut_with, enumerate_with, fold_reduce_with, for_each_chunk_mut_enumerated_with,
92    for_each_chunk_mut_with, for_each_chunk_mut_with_state,
93    for_each_chunk_pair_mut_enumerated_with, for_each_chunk_quad_mut_enumerated_with,
94    for_each_chunk_triple_mut_enumerated_with, for_each_index_with, for_each_mut_with,
95    for_each_with, join, join_with, map_collect_index_with, map_collect_mut_with, map_collect_with,
96    map_reduce_with, reduce_index_with, scope, Scope,
97};
98
99// ---------------------------------------------------------------------------
100// Extension traits: trait-based, type-selected parallel views over slices
101// ---------------------------------------------------------------------------
102
103/// A read-only parallel view of a slice bound to execution policy `P`.
104///
105/// Construct via [`ParallelSlice::par`]. Zero-sized beyond the borrowed slice.
106pub struct ParRef<'a, T, P> {
107    data: &'a [T],
108    _policy: PhantomData<P>,
109}
110
111impl<'a, T, P: ExecutionPolicy> ParRef<'a, T, P> {
112    /// Apply `f` to every element. See [`for_each_with`].
113    pub fn for_each<F: Fn(&T) + Send + Sync>(self, f: F)
114    where
115        T: Sync,
116    {
117        for_each_with::<P, _, _>(self.data, f);
118    }
119
120    /// Apply `f(index, &element)` to every element. See [`enumerate_with`].
121    pub fn enumerate<F: Fn(usize, &T) + Send + Sync>(self, f: F)
122    where
123        T: Sync,
124    {
125        enumerate_with::<P, _, _>(self.data, f);
126    }
127
128    /// Map then collect into a `Vec<R>` in order. See [`map_collect_with`].
129    pub fn map_collect<R: Send, F: Fn(&T) -> R + Send + Sync>(self, f: F) -> Vec<R>
130    where
131        T: Sync,
132    {
133        map_collect_with::<P, _, _, _>(self.data, f)
134    }
135
136    /// Map each `(index, &element)` pair, then collect into a `Vec<R>` in order.
137    ///
138    /// This is the indexed counterpart to [`Self::map_collect`]. It exposes the
139    /// logical slice index without requiring callers to derive it from pointers.
140    pub fn map_collect_index<R: Send, F: Fn(usize, &T) -> R + Send + Sync>(self, f: F) -> Vec<R>
141    where
142        T: Sync,
143    {
144        map_collect_index_with::<P, _, _>(self.data.len(), |i| f(i, &self.data[i]))
145    }
146
147    /// Map-reduce. See [`map_reduce_with`].
148    pub fn map_reduce<R, M, Rd>(self, identity: R, map: M, reduce: Rd) -> R
149    where
150        T: Sync,
151        R: Send + Sync + Clone,
152        M: Fn(&T) -> R + Send + Sync,
153        Rd: Fn(R, R) -> R + Send + Sync,
154    {
155        map_reduce_with::<P, _, _, _, _>(self.data, identity, map, reduce)
156    }
157}
158
159/// A mutable parallel view of a slice bound to execution policy `P`.
160pub struct ParMut<'a, T, P> {
161    data: &'a mut [T],
162    _policy: PhantomData<P>,
163}
164
165impl<'a, T, P: ExecutionPolicy> ParMut<'a, T, P> {
166    /// Apply `f` to every element in place. See [`for_each_mut_with`].
167    pub fn for_each<F: Fn(&mut T) + Send + Sync>(self, f: F)
168    where
169        T: Send,
170    {
171        for_each_mut_with::<P, _, _>(self.data, f);
172    }
173
174    /// Apply `f(index, &mut element)` to every element in place.
175    pub fn enumerate<F: Fn(usize, &mut T) + Send + Sync>(self, f: F)
176    where
177        T: Send,
178    {
179        enumerate_mut_with::<P, _, _>(self.data, f);
180    }
181}
182
183/// Extension trait providing an adaptive parallel view over `&[T]`.
184pub trait ParallelSlice<T> {
185    /// Adaptive, auto-routing parallel view (the everyday entry point).
186    fn par(&self) -> ParRef<'_, T, Adaptive>;
187}
188
189impl<T> ParallelSlice<T> for [T] {
190    #[inline]
191    fn par(&self) -> ParRef<'_, T, Adaptive> {
192        ParRef {
193            data: self,
194            _policy: PhantomData,
195        }
196    }
197}
198
199/// Extension trait providing an adaptive mutable parallel view over `&mut [T]`.
200pub trait ParallelSliceMut<T> {
201    /// Adaptive, auto-routing mutable parallel view (the everyday entry point).
202    fn par_mut(&mut self) -> ParMut<'_, T, Adaptive>;
203}
204
205impl<T> ParallelSliceMut<T> for [T] {
206    #[inline]
207    fn par_mut(&mut self) -> ParMut<'_, T, Adaptive> {
208        ParMut {
209            data: self,
210            _policy: PhantomData,
211        }
212    }
213}
214
215#[cfg(feature = "melinoe")]
216pub mod melinoe_ext;
217
218#[cfg(test)]
219#[path = "tests.rs"]
220mod tests;