Skip to main content

moirai_iter/parallel/adapters/
blocks.rs

1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
4struct ExponentialBlockPolicy;
5
6#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
7struct UniformBlockPolicy;
8
9#[repr(transparent)]
10#[derive(Clone, Copy, Debug, Eq, PartialEq)]
11struct BlockSize(usize);
12
13impl BlockSize {
14    fn new(value: usize) -> Self {
15        assert!(value != 0, "block size must be non-zero");
16        Self(value)
17    }
18}
19
20/// Indexed source adapter for exponential logical block scheduling.
21pub struct ExponentialBlocks<I> {
22    base: I,
23    policy: ExponentialBlockPolicy,
24}
25
26impl<I> ExponentialBlocks<I> {
27    pub(in crate::parallel) fn new(base: I) -> Self {
28        Self {
29            base,
30            policy: ExponentialBlockPolicy,
31        }
32    }
33}
34
35impl<I> ParallelIterator for ExponentialBlocks<I>
36where
37    I: ParallelIterator,
38    I::Item: Sync + 'static,
39{
40    type Item = I::Item;
41
42    fn seq_items(self) -> Vec<Self::Item> {
43        let _policy = self.policy;
44        self.base.seq_items()
45    }
46
47    fn drive<C, R>(self, consumer: C) -> R
48    where
49        C: Consumer<Self::Item, Result = R> + Send + Sync,
50        R: Send,
51    {
52        consumer.consume(VecParIter::new(self.seq_items()))
53    }
54}
55
56/// Indexed source adapter for uniform logical block scheduling.
57pub struct UniformBlocks<I> {
58    base: I,
59    block_size: BlockSize,
60    policy: UniformBlockPolicy,
61}
62
63impl<I> UniformBlocks<I> {
64    pub(in crate::parallel) fn new(base: I, block_size: usize) -> Self {
65        Self {
66            base,
67            block_size: BlockSize::new(block_size),
68            policy: UniformBlockPolicy,
69        }
70    }
71}
72
73impl<I> ParallelIterator for UniformBlocks<I>
74where
75    I: ParallelIterator,
76    I::Item: Sync + 'static,
77{
78    type Item = I::Item;
79
80    fn seq_items(self) -> Vec<Self::Item> {
81        let _block_size = self.block_size;
82        let _policy = self.policy;
83        self.base.seq_items()
84    }
85
86    fn drive<C, R>(self, consumer: C) -> R
87    where
88        C: Consumer<Self::Item, Result = R> + Send + Sync,
89        R: Send,
90    {
91        consumer.consume(VecParIter::new(self.seq_items()))
92    }
93}
94
95#[cfg(test)]
96mod tests {
97    use super::*;
98
99    #[test]
100    fn block_policy_markers_are_zero_sized() {
101        assert_eq!(std::mem::size_of::<ExponentialBlockPolicy>(), 0);
102        assert_eq!(std::mem::size_of::<UniformBlockPolicy>(), 0);
103    }
104}