Skip to main content

moirai_iter/parallel/adapters/
chunks.rs

1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3#[repr(transparent)]
4#[derive(Clone, Copy, Debug, Eq, PartialEq)]
5struct ChunkSize(usize);
6
7impl ChunkSize {
8    fn new(value: usize) -> Self {
9        assert!(value != 0, "chunk size must be non-zero");
10        Self(value)
11    }
12
13    const fn get(self) -> usize {
14        self.0
15    }
16}
17
18/// Chunk adapter with Rayon-style non-empty chunk-size semantics.
19pub struct Chunks<I> {
20    base: I,
21    chunk_size: ChunkSize,
22}
23
24impl<I> Chunks<I> {
25    pub(in crate::parallel) fn new(base: I, chunk_size: usize) -> Self {
26        Self {
27            base,
28            chunk_size: ChunkSize::new(chunk_size),
29        }
30    }
31
32    pub(in crate::parallel) fn into_parts(self) -> (I, usize) {
33        (self.base, self.chunk_size.get())
34    }
35}
36
37impl<I> ParallelIterator for Chunks<I>
38where
39    I: ParallelIterator,
40    I::Item: Sync + 'static,
41{
42    type Item = Vec<I::Item>;
43
44    fn seq_items(self) -> Vec<Self::Item> {
45        let (base, chunk_size) = self.into_parts();
46        let mut items = base.seq_items();
47        let mut chunks = Vec::with_capacity(items.len().div_ceil(chunk_size));
48
49        let tail_len = items.len() % chunk_size;
50        let tail = if tail_len == 0 {
51            None
52        } else {
53            Some(items.split_off(items.len() - tail_len))
54        };
55
56        while !items.is_empty() {
57            chunks.push(items.split_off(items.len() - chunk_size));
58        }
59        chunks.reverse();
60
61        if let Some(tail) = tail {
62            chunks.push(tail);
63        }
64
65        chunks
66    }
67
68    /// # Why this stays sequential (chunk boundaries are logical positions)
69    ///
70    /// A chunk is defined by its position in the logical stream, and the source
71    /// splits at its own midpoint, which is not in general a multiple of
72    /// `chunk_size`. A shard chunking its own range alone would emit a short
73    /// chunk at every internal shard boundary, so a stream whose only short
74    /// chunk should be the tail would gain one per split. Aligning splits to
75    /// chunk boundaries is a producer-side decision, not something this adapter
76    /// can express by pushing into a consumer.
77    fn drive<C, R>(self, consumer: C) -> R
78    where
79        C: Consumer<Self::Item, Result = R> + Send + Sync,
80        R: Send,
81    {
82        consumer.consume(VecParIter::new(self.seq_items()))
83    }
84}