Skip to main content

moirai_parallel/
ops.rs

1//! Synchronous data-parallel operators over the unified scheduler.
2//!
3//! # Safety
4//!
5//! The mutable operators split a buffer across worker tasks through Melinoe's
6//! [`WriterShard`](melinoe::region::WriterShard) /
7//! [`ParChunks`](melinoe::region::ParChunks), which own the disjoint-partition
8//! contract and the range math once; each operator brands the caller's slice in
9//! place ([`MelinoeCell::from_mut_slice`](melinoe::MelinoeCell::from_mut_slice))
10//! and vends one partition per task. One
11//! executor contract makes every such access sound, and each per-site `SAFETY`
12//! comment appeals to it:
13//!
14//! - **Disjoint partition.** `global().for_each_indexed(count, f)` invokes `f`
15//!   with each index in `0..count` exactly once — the same guarantee
16//!   `ParChunks::get_unchecked_chunk` requires of its caller. The operators map
17//!   that index (or a `chunk_size`-strided range derived from it) to a
18//!   *disjoint* slice of the buffer, so no two concurrent tasks ever form
19//!   `&mut` to the same element. Multi-buffer operators additionally rely on the
20//!   caller's distinct `&mut [_]` arguments being non-aliasing (guaranteed by
21//!   the borrow checker at the call site).
22//! - **All-or-error collect.** The `map_collect_*` helpers build a
23//!   `Vec<MaybeUninit<R>>`, `set_len` it (sound — `MaybeUninit` needs no
24//!   initialization), fill every slot through the disjoint-partition contract,
25//!   then reinterpret it as `Vec<R>`. `for_each_indexed` returns `Ok` only after
26//!   writing every index, so the reinterpretation is reached only when all slots
27//!   are initialized; a task panic instead surfaces as `Err`, which `.expect()`
28//!   turns into a propagating panic that unwinds with the buffer still typed
29//!   `MaybeUninit<R>` (its contents are not dropped — a leak of the written
30//!   values on panic, never a use of uninitialized memory).
31
32use crate::policy::{ExecutionPolicy, Parallel};
33use melinoe::MelinoeCell;
34use melinoe::region::WriterShard;
35use moirai_core::error::{ExecutorError, ExecutorResult};
36use moirai_executor::{HybridExecutor, SchedulerScope, SyncTask, global};
37use std::sync::Mutex;
38
39mod chunks;
40mod shards;
41mod unit_tasks;
42pub use chunks::{
43    ChunkBuffersError, for_each_chunk_buffers_mut_enumerated_with,
44    for_each_chunk_mut_enumerated_with, for_each_chunk_mut_with, for_each_chunk_mut_with_state,
45    for_each_chunk_pair_mut_enumerated_with, for_each_chunk_quad_mut_enumerated_with,
46    for_each_chunk_triple_mut_enumerated_with,
47};
48pub use unit_tasks::{
49    UNIT_TASK_BYTES, for_each_unit_task_many_mut_with, for_each_unit_task_mut_with,
50    for_each_unit_task_pair_mut_with, for_each_unit_task_range_with,
51    for_each_unit_task_triple_mut_with, units_per_task,
52};
53
54/// State of the scheduled branch of a join.
55///
56/// The scheduler can refuse a job — while shutting down, or when a worker's
57/// bounded admission queue is full — and it drops the refused job before
58/// returning the error. A branch owned by that job would go with it, so the
59/// closure lives here instead and whichever lane reaches it first takes it.
60/// The caller can therefore still run a branch the scheduler never did.
61enum Branch<F, R> {
62    /// Nobody has claimed this branch yet.
63    Pending(F),
64    /// A lane claimed the branch and has not published a result. Observing
65    /// this once the scope has joined means that lane unwound.
66    Claimed,
67    /// Ran to completion.
68    Done(R),
69}
70
71impl<F, R> Branch<F, R>
72where
73    F: FnOnce() -> R,
74{
75    /// Take the closure if this lane is the one that gets to run it.
76    fn claim(&mut self) -> Option<F> {
77        match std::mem::replace(self, Self::Claimed) {
78            Self::Pending(branch) => Some(branch),
79            other => {
80                *self = other;
81                None
82            }
83        }
84    }
85
86    fn complete(&mut self, result: R) {
87        *self = Self::Done(result);
88    }
89
90    /// Run the branch on a lane that shares the slot, unless another lane
91    /// already claimed it.
92    ///
93    /// The lock is released before the closure runs, so a branch never holds it
94    /// across arbitrary caller code and a panicking branch cannot poison it.
95    fn run_shared(slot: &Mutex<Self>) {
96        let Some(branch) = lock(slot).claim() else {
97            return;
98        };
99        let result = branch();
100        lock(slot).complete(result);
101    }
102
103    /// Run a branch that never leaves this thread.
104    fn run_here(&mut self) {
105        let Some(branch) = self.claim() else {
106            return;
107        };
108        let result = branch();
109        self.complete(result);
110    }
111
112    /// Take the finished value.
113    fn into_result(self) -> R {
114        match self {
115            Self::Done(result) => result,
116            _ => panic!("invariant: a join branch neither ran nor reported failure"),
117        }
118    }
119}
120
121/// Lock without propagating poisoning: every path that touches a slot leaves
122/// it in a consistent state, and [`Branch::run_shared`] never holds the lock
123/// across the branch closure, so a poisoned flag carries no information here.
124fn lock<T>(slot: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
125    slot.lock()
126        .unwrap_or_else(std::sync::PoisonError::into_inner)
127}
128
129/// Reclaim a slot once every lane that shared it has finished.
130fn lock_owned<T>(slot: Mutex<T>) -> T {
131    slot.into_inner()
132        .unwrap_or_else(std::sync::PoisonError::into_inner)
133}
134
135/// Run two closures to completion and return both results.
136///
137/// This is the synchronous Rayon-style `join` shape. The policy is selected at
138/// compile time; [`Sequential`](crate::Sequential) runs both closures on the
139/// caller, [`Parallel`] schedules the left closure on the unified scheduler and
140/// runs the right closure on the caller lane, and [`crate::Adaptive`] currently
141/// stays sequential for a fixed two-branch join.
142///
143/// A branch the scheduler refuses runs on the caller instead, so a shutting-down
144/// or saturated executor makes the join sequential rather than losing a branch.
145///
146/// # Panics
147///
148/// Panics if a branch panicked, propagating the failure on the caller's thread
149/// as rayon does.
150pub fn join_with<P, A, B, RA, RB>(left: A, right: B) -> (RA, RB)
151where
152    P: ExecutionPolicy,
153    A: FnOnce() -> RA + Send,
154    B: FnOnce() -> RB,
155    RA: Send,
156{
157    if !P::parallelize_pair() {
158        return (left(), right());
159    }
160
161    join_on(global(), left, right)
162}
163
164/// [`join_with`]'s parallel path against a named executor.
165///
166/// Separate from the public entry so the refusal path can be exercised against
167/// a shut-down executor in tests.
168pub(crate) fn join_on<A, B, RA, RB>(executor: &HybridExecutor, left: A, right: B) -> (RA, RB)
169where
170    A: FnOnce() -> RA + Send,
171    B: FnOnce() -> RB,
172    RA: Send,
173{
174    // Only the left branch crosses a lane boundary, so only it needs a shared
175    // slot. The right branch stays on this thread and is reborrowed by the
176    // scope body, which leaves it runnable here if the body returns early.
177    let left_slot = Mutex::new(Branch::Pending(left));
178    let mut right_slot = Branch::Pending(right);
179
180    let forked = executor.scope::<SyncTask, _>(|scope| {
181        scope.spawn(|_| Branch::run_shared(&left_slot))?;
182        // Enter the scheduler before the caller takes its own branch, so the
183        // two overlap instead of running back to back.
184        scope.flush()?;
185        right_slot.run_here();
186        Ok(())
187    });
188
189    match forked {
190        Ok(()) => {}
191        // The scheduler refused the job and dropped it unexecuted, so neither
192        // branch is guaranteed to have run. Both claims are idempotent: a
193        // branch that did run is no longer `Pending`.
194        Err(ExecutorError::ShuttingDown | ExecutorError::ResourceExhausted(_)) => {
195            Branch::run_shared(&left_slot);
196            right_slot.run_here();
197        }
198        Err(error) => panic!("invariant: scheduled join branch failed ({error})"),
199    }
200
201    (
202        lock_owned(left_slot).into_result(),
203        right_slot.into_result(),
204    )
205}
206
207/// Adaptive Rayon-style two-closure join.
208///
209/// Use [`join_with`] to force a specific execution policy.
210pub fn join<A, B, RA, RB>(left: A, right: B) -> (RA, RB)
211where
212    A: FnOnce() -> RA + Send,
213    B: FnOnce() -> RB,
214    RA: Send,
215{
216    join_with::<crate::Adaptive, _, _, _, _>(left, right)
217}
218
219/// Borrowing scope for spawning parallel sub-tasks that may capture non-`'static`
220/// references.
221///
222/// Created by [`scope`]. Each [`Scope::spawn`] call registers a job on the unified
223/// scheduler; the scope blocks until every spawned job has completed before the
224/// body closure returns, so borrowed data cannot escape the scope.
225///
226/// This is the Rayon-style `scope` shape, adapted to Moirai's unified hybrid
227/// scheduler. Unlike [`join`], which forks exactly two branches, `scope` allows
228/// an arbitrary number of sub-tasks to be spawned and joined within a single
229/// region.
230pub struct Scope<'scope, 'env: 'scope> {
231    inner: &'scope SchedulerScope<'env, SyncTask>,
232}
233
234impl<'scope, 'env: 'scope> Scope<'scope, 'env> {
235    /// Spawn a parallel sub-task within this scope.
236    ///
237    /// The task may borrow values that outlive the scope call. The scope waits
238    /// for every spawned task before returning, so borrowed data cannot escape.
239    /// A task the scheduler's admission queue turns away runs on the calling
240    /// thread instead, so backpressure costs parallelism rather than the task.
241    ///
242    /// # Panics
243    ///
244    /// Panics if the underlying scheduler refuses to register the task, which
245    /// registration itself does not do — the failure surfaces from [`scope`]
246    /// when the scheduler is shutting down.
247    #[inline]
248    pub fn spawn<F>(&self, task: F)
249    where
250        F: FnOnce() + Send + 'env,
251    {
252        self.inner
253            .spawn(move |_| task())
254            .expect("moirai global executor: scope spawn");
255    }
256}
257
258/// Create a borrowing scope for parallel sub-tasks.
259///
260/// Within the body closure, [`Scope::spawn`] registers jobs on the unified
261/// scheduler. The scope blocks until every spawned job has completed before
262/// returning, so tasks may borrow non-`'static` data from the enclosing
263/// environment.
264///
265/// This is the Rayon-style `scope` shape, adapted to Moirai's unified hybrid
266/// scheduler.
267///
268/// # Examples
269///
270/// ```
271/// use moirai_parallel::scope;
272///
273/// let data: Vec<u64> = (0..1000).collect();
274/// use std::sync::atomic::{AtomicU64, Ordering};
275///
276/// let sum = AtomicU64::new(0);
277/// scope(|s| {
278///     s.spawn(|| {
279///         sum.fetch_add(data.iter().sum::<u64>(), Ordering::Relaxed);
280///     });
281///     s.spawn(|| {
282///         sum.fetch_add(data.len() as u64, Ordering::Relaxed);
283///     });
284/// });
285/// assert_eq!(sum.load(Ordering::Relaxed), data.iter().sum::<u64>() + 1000);
286/// ```
287///
288/// A task cannot borrow a value local to the body, which is dropped before the
289/// scheduler runs the buffered task:
290///
291/// ```compile_fail,E0597
292/// moirai_parallel::scope(|s| {
293///     let local = vec![7_u8; 8];
294///     let borrowed = &local;
295///     s.spawn(move || assert_eq!(borrowed[0], 7));
296/// });
297/// ```
298#[inline]
299pub fn scope<'env, F, R>(body: F) -> R
300where
301    F: for<'scope> FnOnce(&Scope<'scope, 'env>) -> R,
302    R: Send,
303{
304    let mut result = None;
305    global()
306        .scope::<SyncTask, _>(|inner| {
307            let scope = Scope { inner };
308            result = Some(body(&scope));
309            ExecutorResult::Ok(())
310        })
311        .expect("moirai global executor: scope");
312
313    result.expect("scoped body must complete")
314}
315
316/// Apply `f` to every element of `data`, scheduled by policy `P`.
317pub fn for_each_with<P, T, F>(data: &[T], f: F)
318where
319    P: ExecutionPolicy,
320    T: Sync,
321    F: Fn(&T) + Send + Sync,
322{
323    let n = data.len();
324    if n == 0 {
325        return;
326    }
327    if !P::parallelize(n) {
328        data.iter().for_each(f);
329        return;
330    }
331    let f = &f;
332    global()
333        .for_each_indexed::<SyncTask, _>(n, move |i| f(&data[i]))
334        .expect("moirai global executor: for_each_with");
335}
336
337/// Partition width for an indexed region: about one worker-sized chunk per task.
338///
339/// `for_each_indexed` already splits its domain into contiguous worker-sized
340/// ranges, so matching the partition to that width keeps the same element
341/// distribution while constructing one `WriterShard` per range instead of one
342/// per element. A width of `1` would put a shard build (and its range math) on
343/// every element's path for no benefit.
344fn worker_chunk_size(len: usize) -> usize {
345    let workers = moirai_core::executor::logical_parallelism().max(1);
346    len.div_ceil(workers).max(1)
347}
348
349/// Apply `f` to every element of `data` in place, scheduled by policy `P`.
350pub fn for_each_mut_with<P, T, F>(data: &mut [T], f: F)
351where
352    P: ExecutionPolicy,
353    T: Send,
354    F: Fn(&mut T) + Send + Sync,
355{
356    let n = data.len();
357    if n == 0 {
358        return;
359    }
360    if !P::parallelize(n) {
361        data.iter_mut().for_each(f);
362        return;
363    }
364    let partitions =
365        WriterShard::new(MelinoeCell::from_mut_slice(data)).par_chunks(worker_chunk_size(n));
366    let tasks = partitions.len();
367    let f = &f;
368    global()
369        .for_each_indexed::<SyncTask, _>(tasks, move |c| {
370            // SAFETY: the scheduler visits each index in `0..tasks` exactly once,
371            // which is the contract `get_unchecked_chunk` documents; distinct
372            // partitions name disjoint element ranges, so no two tasks form a
373            // `&mut` to the same element. `data` is branded for the whole
374            // joined call, so the view cannot outlive it.
375            let mut shard = unsafe { partitions.get_unchecked_chunk(c) };
376            for element in shard.iter_mut() {
377                f(element);
378            }
379        })
380        .expect("moirai global executor: for_each_mut_with");
381}
382
383/// Apply `f(index, &element)` to every element of `data`, scheduled by policy `P`.
384pub fn enumerate_with<P, T, F>(data: &[T], f: F)
385where
386    P: ExecutionPolicy,
387    T: Sync,
388    F: Fn(usize, &T) + Send + Sync,
389{
390    let n = data.len();
391    if n == 0 {
392        return;
393    }
394    if !P::parallelize(n) {
395        data.iter().enumerate().for_each(|(i, x)| f(i, x));
396        return;
397    }
398    let f = &f;
399    global()
400        .for_each_indexed::<SyncTask, _>(n, move |i| f(i, &data[i]))
401        .expect("moirai global executor: enumerate_with");
402}
403
404/// Apply `f(index, &mut element)` to every element of `data` in place,
405/// scheduled by policy `P`.
406pub fn enumerate_mut_with<P, T, F>(data: &mut [T], f: F)
407where
408    P: ExecutionPolicy,
409    T: Send,
410    F: Fn(usize, &mut T) + Send + Sync,
411{
412    let n = data.len();
413    if n == 0 {
414        return;
415    }
416    if !P::parallelize(n) {
417        data.iter_mut().enumerate().for_each(|(i, x)| f(i, x));
418        return;
419    }
420    let chunk_size = worker_chunk_size(n);
421    let partitions = WriterShard::new(MelinoeCell::from_mut_slice(data)).par_chunks(chunk_size);
422    let tasks = partitions.len();
423    let f = &f;
424    global()
425        .for_each_indexed::<SyncTask, _>(tasks, move |c| {
426            // Partition `c` covers the elements `[c * chunk_size, …)`.
427            let start = c * chunk_size;
428            // SAFETY: each index in `0..tasks` is visited exactly once; distinct
429            // partitions are disjoint element ranges, so the element at absolute
430            // index `start + offset` is touched by exactly one task. See
431            // `for_each_mut_with`.
432            let mut shard = unsafe { partitions.get_unchecked_chunk(c) };
433            for (offset, element) in shard.iter_mut().enumerate() {
434                f(start + offset, element);
435            }
436        })
437        .expect("moirai global executor: enumerate_mut_with");
438}
439
440/// Apply `f` to every index in `0..len` in parallel, scheduled by policy `P`.
441///
442/// Synchronous equivalent of rayon's `(0..len).into_par_iter().for_each(f)`. Use
443/// when the work is keyed by index and writes through external disjoint state
444/// (atomics, per-index channels) rather than returning a value.
445pub fn for_each_index_with<P, F>(len: usize, f: F)
446where
447    P: ExecutionPolicy,
448    F: Fn(usize) + Send + Sync,
449{
450    if len == 0 {
451        return;
452    }
453    if !P::parallelize(len) {
454        (0..len).for_each(f);
455        return;
456    }
457    let f = &f;
458    global()
459        .for_each_indexed::<SyncTask, _>(len, f)
460        .expect("moirai global executor: for_each_index_with");
461}
462
463/// Map each element of `data` with `f`, collecting into a `Vec<R>` in order,
464/// scheduled by policy `P`.
465pub fn map_collect_with<P, T, R, F>(data: &[T], f: F) -> Vec<R>
466where
467    P: ExecutionPolicy,
468    T: Sync,
469    R: Send,
470    F: Fn(&T) -> R + Send + Sync,
471{
472    let n = data.len();
473    if !P::parallelize(n) {
474        return data.iter().map(f).collect();
475    }
476    let mut out: Vec<core::mem::MaybeUninit<R>> = Vec::with_capacity(n);
477    // SAFETY: capacity is `n`; every slot is written exactly once below before
478    // being read, and `MaybeUninit` makes `set_len` sound without initialization.
479    unsafe {
480        out.set_len(n);
481    }
482    enumerate_mut_with::<Parallel, _, _>(&mut out, |i, slot| {
483        slot.write(f(&data[i]));
484    });
485    // SAFETY: every slot initialized above; `MaybeUninit<R>` shares `R`'s layout.
486    let mut out = core::mem::ManuallyDrop::new(out);
487    unsafe { Vec::from_raw_parts(out.as_mut_ptr().cast::<R>(), n, out.capacity()) }
488}
489
490/// Map-reduce over `data`, scheduled by policy `P`.
491///
492/// `reduce` must be associative and `identity` its neutral element, since chunk
493/// boundaries and combination order are unspecified.
494pub fn map_reduce_with<P, T, R, M, Rd>(data: &[T], identity: R, map: M, reduce: Rd) -> R
495where
496    P: ExecutionPolicy,
497    T: Sync,
498    R: Send + Sync + Clone,
499    M: Fn(&T) -> R + Send + Sync,
500    Rd: Fn(R, R) -> R + Send + Sync,
501{
502    let n = data.len();
503    if n == 0 || !P::parallelize(n) {
504        let mut acc = identity;
505        for item in data {
506            acc = reduce(acc, map(item));
507        }
508        return acc;
509    }
510    let map = &map;
511    let reduce = &reduce;
512    // The executor folds each worker chunk locally (seeded by `identity`) then
513    // combines chunk results, so `map` is per-element and `reduce` per-pair.
514    global()
515        .map_reduce_indexed::<SyncTask, _, _, _>(n, identity, move |i| map(&data[i]), reduce)
516        .expect("moirai global executor: map_reduce_with")
517}
518
519/// Parallel fold-reduce over the index domain `0..len`, scheduled by policy `P`.
520///
521/// Each worker chunk creates one accumulator with `init()`, folds its indices
522/// into it with `fold`, and the per-chunk accumulators are combined with
523/// `reduce`. Unlike [`reduce_index_with`], `fold` mutates a single accumulator
524/// per chunk (no per-element temporary), which is the efficient shape for
525/// accumulating into a collection — e.g. grouping entries into a `HashMap`.
526/// `reduce` must be associative; `init()` must yield its neutral element.
527pub fn fold_reduce_with<P, A, Init, Fold, Red>(len: usize, init: Init, fold: Fold, reduce: Red) -> A
528where
529    P: ExecutionPolicy,
530    A: Send,
531    Init: Fn() -> A + Send + Sync,
532    Fold: Fn(A, usize) -> A + Send + Sync,
533    Red: Fn(A, A) -> A,
534{
535    if len == 0 {
536        return init();
537    }
538    if !P::parallelize(len) {
539        let mut acc = init();
540        for i in 0..len {
541            acc = fold(acc, i);
542        }
543        return acc;
544    }
545    let workers = moirai_core::executor::logical_parallelism().max(1);
546    let chunks = workers.min(len).max(1);
547    let chunk = len.div_ceil(chunks);
548    let mut slots: Vec<Option<A>> = (0..chunks).map(|_| None).collect();
549    let partitions =
550        WriterShard::new(MelinoeCell::from_mut_slice(slots.as_mut_slice())).par_chunks(1);
551    let init_ref = &init;
552    let fold_ref = &fold;
553    global()
554        .for_each_indexed::<SyncTask, _>(chunks, move |ci| {
555            let start = ci * chunk;
556            if start >= len {
557                return;
558            }
559            let end = (start + chunk).min(len);
560            let mut acc = init_ref();
561            for i in start..end {
562                acc = fold_ref(acc, i);
563            }
564            // SAFETY: each `ci` in `0..chunks` is visited exactly once, and
565            // partition `ci` is exactly slot `ci` (one cell per partition), so
566            // every slot is written exactly once and no two tasks alias.
567            let mut shard = unsafe { partitions.get_unchecked_chunk(ci) };
568            if let Some(slot) = shard.get_mut(0) {
569                *slot = Some(acc);
570            }
571        })
572        .expect("moirai global executor: fold_reduce_with");
573    slots
574        .into_iter()
575        .flatten()
576        .reduce(reduce)
577        .unwrap_or_else(init)
578}
579
580/// Parallel map over the index domain `0..len`, collecting into a `Vec<R>` in
581/// order, scheduled by policy `P`.
582///
583/// `map(i)` produces the element at index `i`. Use this for index-aligned maps
584/// over multiple slices that [`map_collect_with`] cannot express — e.g. an
585/// elementwise product `map_collect_index_with::<Adaptive>(n, |i| a[i] * b[i])`.
586pub fn map_collect_index_with<P, R, Map>(len: usize, map: Map) -> Vec<R>
587where
588    P: ExecutionPolicy,
589    R: Send,
590    Map: Fn(usize) -> R + Send + Sync,
591{
592    if !P::parallelize(len) {
593        return (0..len).map(map).collect();
594    }
595    let mut out: Vec<core::mem::MaybeUninit<R>> = Vec::with_capacity(len);
596    // SAFETY: capacity is `len`; every slot is written exactly once below.
597    unsafe {
598        out.set_len(len);
599    }
600    enumerate_mut_with::<Parallel, _, _>(&mut out, |i, slot| {
601        slot.write(map(i));
602    });
603    // SAFETY: every slot initialized; `MaybeUninit<R>` shares `R`'s layout.
604    let mut out = core::mem::ManuallyDrop::new(out);
605    unsafe { Vec::from_raw_parts(out.as_mut_ptr().cast::<R>(), len, out.capacity()) }
606}
607
608/// Map each element of `data` in place with `f(index, &mut element)`, collecting
609/// each returned value into a `Vec<R>` in order, scheduled by policy `P`.
610///
611/// The synchronous equivalent of rayon's
612/// `data.par_iter_mut().enumerate().map(f).collect()`: each element is mutated
613/// and produces a result. Use for parallel solve-in-place-and-collect loops.
614pub fn map_collect_mut_with<P, T, R, F>(data: &mut [T], f: F) -> Vec<R>
615where
616    P: ExecutionPolicy,
617    T: Send,
618    R: Send,
619    F: Fn(usize, &mut T) -> R + Send + Sync,
620{
621    let n = data.len();
622    if !P::parallelize(n) {
623        return data.iter_mut().enumerate().map(|(i, x)| f(i, x)).collect();
624    }
625    let mut out: Vec<core::mem::MaybeUninit<R>> = Vec::with_capacity(n);
626    // SAFETY: capacity is `n`; every slot is written exactly once below.
627    unsafe {
628        out.set_len(n);
629    }
630    let chunk_size = worker_chunk_size(n);
631    let data_partitions =
632        WriterShard::new(MelinoeCell::from_mut_slice(data)).par_chunks(chunk_size);
633    let out_partitions =
634        WriterShard::new(MelinoeCell::from_mut_slice(out.as_mut_slice())).par_chunks(chunk_size);
635    let tasks = data_partitions.len();
636    let f = &f;
637    global()
638        .for_each_indexed::<SyncTask, _>(tasks, move |c| {
639            // Partition `c` covers the elements `[c * chunk_size, …)` in both
640            // regions; the two regions have equal length, so their partition
641            // counts agree.
642            let start = c * chunk_size;
643            // SAFETY: each index in `0..tasks` is visited exactly once, for both
644            // the input element and the output slot at `start + offset`. The two
645            // regions are distinct borrows, and distinct partitions are disjoint
646            // ranges in each, so neither the element nor the slot at an absolute
647            // index aliases another task's.
648            let mut data_shard = unsafe { data_partitions.get_unchecked_chunk(c) };
649            let mut out_shard = unsafe { out_partitions.get_unchecked_chunk(c) };
650            for (offset, (element, slot)) in
651                data_shard.iter_mut().zip(out_shard.iter_mut()).enumerate()
652            {
653                slot.write(f(start + offset, element));
654            }
655        })
656        .expect("moirai global executor: map_collect_mut_with");
657    // SAFETY: every slot initialized; `MaybeUninit<R>` shares `R`'s layout.
658    let mut out = core::mem::ManuallyDrop::new(out);
659    unsafe { Vec::from_raw_parts(out.as_mut_ptr().cast::<R>(), n, out.capacity()) }
660}
661
662/// Parallel reduction over the index domain `0..len`, scheduled by policy `P`.
663///
664/// `map(i)` produces a value for index `i`; results are folded within and across
665/// chunks with `reduce`, seeded by `identity` (which must be `reduce`'s neutral
666/// element). Use this for index-aligned reductions over multiple slices that
667/// [`map_reduce_with`] cannot express — e.g. a dot product
668/// `reduce_index_with::<Adaptive>(n, T::zero(), |i| a[i] * b[i], |x, y| x + y)`.
669pub fn reduce_index_with<P, R, Map, Red>(len: usize, identity: R, map: Map, reduce: Red) -> R
670where
671    P: ExecutionPolicy,
672    R: Send + Sync + Clone,
673    Map: Fn(usize) -> R + Send + Sync,
674    Red: Fn(R, R) -> R + Send + Sync,
675{
676    if len == 0 || !P::parallelize(len) {
677        let mut acc = identity;
678        for i in 0..len {
679            acc = reduce(acc, map(i));
680        }
681        return acc;
682    }
683    global()
684        .map_reduce_indexed::<SyncTask, _, _, _>(len, identity, map, reduce)
685        .expect("moirai global executor: reduce_index_with")
686}