Skip to main content

moirai_executor/schedule/reduce/
mod.rs

1//! Result slots for scoped indexed reductions.
2
3use std::{
4    cell::UnsafeCell,
5    mem::MaybeUninit,
6    sync::atomic::{AtomicBool, Ordering},
7};
8
9pub(crate) struct ReduceSlots<T> {
10    slots: Box<[ReduceSlot<T>]>,
11}
12
13struct ReduceSlot<T> {
14    value: UnsafeCell<MaybeUninit<T>>,
15    initialized: AtomicBool,
16}
17
18// Safety: each slot is written by exactly one scheduled chunk and read by the
19// parent after the scoped completion counter reaches zero. `T: Send` is
20// required because values cross worker-thread boundaries.
21unsafe impl<T: Send> Sync for ReduceSlot<T> {}
22
23impl<T> ReduceSlots<T> {
24    pub(crate) fn new(len: usize) -> Self {
25        Self {
26            slots: (0..len)
27                .map(|_| ReduceSlot::new())
28                .collect::<Vec<_>>()
29                .into_boxed_slice(),
30        }
31    }
32
33    pub(crate) fn write(&self, index: usize, value: T) {
34        self.slots[index].write(value);
35    }
36
37    pub(crate) fn reduce<F>(&self, identity: T, reduce: F) -> T
38    where
39        F: Fn(T, T) -> T,
40    {
41        self.slots
42            .iter()
43            .fold(identity, |accumulator, slot| match slot.take() {
44                Some(value) => reduce(accumulator, value),
45                _ => accumulator,
46            })
47    }
48}
49
50impl<T> Drop for ReduceSlots<T> {
51    fn drop(&mut self) {
52        for slot in &self.slots {
53            let _ = slot.take();
54        }
55    }
56}
57
58impl<T> ReduceSlot<T> {
59    fn new() -> Self {
60        Self {
61            value: UnsafeCell::new(MaybeUninit::uninit()),
62            initialized: AtomicBool::new(false),
63        }
64    }
65
66    fn write(&self, value: T) {
67        // Safety: a chunk owns exactly one slot index. No other worker writes
68        // the same slot, and the parent reads only after scoped completion.
69        unsafe {
70            (*self.value.get()).write(value);
71        }
72        self.initialized.store(true, Ordering::Release);
73    }
74
75    fn take(&self) -> Option<T> {
76        if self.initialized.swap(false, Ordering::AcqRel) {
77            // Safety: initialized=true is stored only after `write` initializes
78            // the slot. The swap gives the caller unique read ownership.
79            Some(unsafe { (*self.value.get()).assume_init_read() })
80        } else {
81            None
82        }
83    }
84}