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 = ↦
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}