Skip to main content

moirai_iter/parallel/adapters/
stride.rs

1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3#[repr(transparent)]
4#[derive(Clone, Copy, Debug, Eq, PartialEq)]
5struct StepSize(usize);
6
7impl StepSize {
8    fn new(value: usize) -> Self {
9        assert!(value != 0, "step size must be non-zero");
10        Self(value)
11    }
12
13    const fn get(self) -> usize {
14        self.0
15    }
16}
17
18/// Indexed step-by adapter with exact-size source semantics.
19pub struct StepBy<I> {
20    pub(in crate::parallel) base: I,
21    step: StepSize,
22}
23
24impl<I> StepBy<I> {
25    pub(in crate::parallel) fn new(base: I, step: usize) -> Self {
26        Self {
27            base,
28            step: StepSize::new(step),
29        }
30    }
31
32    pub(in crate::parallel) fn step(&self) -> usize {
33        self.step.get()
34    }
35}
36
37impl<I> ParallelIterator for StepBy<I>
38where
39    I: ParallelIterator,
40    I::Item: Sync + 'static,
41{
42    type Item = I::Item;
43
44    fn seq_items(self) -> Vec<Self::Item> {
45        self.base
46            .seq_items()
47            .into_iter()
48            .step_by(self.step.get())
49            .collect()
50    }
51
52    /// # Why this stays sequential
53    ///
54    /// Which of a shard's items survive depends on the shard's logical start
55    /// index modulo `step`, the same absent offset documented on
56    /// [`Enumerate`](super::ref_ops::Enumerate).
57    fn drive<C, R>(self, consumer: C) -> R
58    where
59        C: Consumer<Self::Item, Result = R> + Send + Sync,
60        R: Send,
61    {
62        consumer.consume(VecParIter::new(self.seq_items()))
63    }
64}