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;