moirai_iter/parallel/adapters/
stride.rs1use 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
18pub 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 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}