Skip to main content

moirai_iter/parallel/adapters/
slice_ops.rs

1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3/// Take adapter with prefix-bounded value semantics.
4pub struct Take<I> {
5    pub(super) base: I,
6    pub(super) count: usize,
7}
8
9impl<I> Take<I> {
10    pub(crate) fn new(base: I, count: usize) -> Self {
11        Self { base, count }
12    }
13}
14
15impl<I> ParallelIterator for Take<I>
16where
17    I: ParallelIterator,
18    I::Item: Sync + 'static,
19{
20    type Item = I::Item;
21
22    fn seq_items(self) -> Vec<Self::Item> {
23        self.base.seq_items_window(0, Some(self.count))
24    }
25
26    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
27        if skip >= self.count {
28            return Vec::new();
29        }
30
31        let remaining = self.count - skip;
32        let count = take.map_or(remaining, |count| count.min(remaining));
33        self.base.seq_items_window(skip, Some(count))
34    }
35
36    fn seq_items_reversed(self) -> Vec<Self::Item> {
37        self.base
38            .seq_items_window(0, Some(self.count))
39            .into_iter()
40            .rev()
41            .collect()
42    }
43
44    fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
45        let mut items = self.base.seq_items_window(0, Some(self.count));
46        let keep = count.min(items.len());
47        items.drain(..items.len().saturating_sub(keep));
48        items.reverse();
49        items
50    }
51
52    /// # Why this stays sequential
53    ///
54    /// The retained prefix is defined by a count over the whole logical stream,
55    /// so a shard cannot tell how many items precede it — the absent offset
56    /// documented on [`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}
65
66/// Skip adapter with prefix-discarding value semantics.
67pub struct Skip<I> {
68    pub(super) base: I,
69    pub(super) count: usize,
70}
71
72impl<I> Skip<I> {
73    pub(crate) fn new(base: I, count: usize) -> Self {
74        Self { base, count }
75    }
76}
77
78impl<I> ParallelIterator for Skip<I>
79where
80    I: ParallelIterator,
81    I::Item: Sync + 'static,
82{
83    type Item = I::Item;
84
85    fn seq_items(self) -> Vec<Self::Item> {
86        self.base.seq_items_window(self.count, None)
87    }
88
89    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
90        self.base
91            .seq_items_window(self.count.saturating_add(skip), take)
92    }
93
94    /// # Why this stays sequential
95    ///
96    /// The discarded prefix is a logical count, per [`Take`].
97    fn drive<C, R>(self, consumer: C) -> R
98    where
99        C: Consumer<Self::Item, Result = R> + Send + Sync,
100        R: Send,
101    {
102        consumer.consume(VecParIter::new(self.seq_items()))
103    }
104}
105
106/// Chain adapter with left-to-right concatenation semantics.
107pub struct Chain<I, J> {
108    pub(super) left: I,
109    pub(super) right: J,
110}
111
112impl<I, J> Chain<I, J> {
113    pub(crate) fn new(left: I, right: J) -> Self {
114        Self { left, right }
115    }
116}
117
118impl<I, J> ParallelIterator for Chain<I, J>
119where
120    I: ParallelIterator,
121    J: ParallelIterator<Item = I::Item>,
122    I::Item: Sync + 'static,
123{
124    type Item = I::Item;
125
126    fn seq_items(self) -> Vec<Self::Item> {
127        let mut left = self.left.seq_items();
128        left.extend(self.right.seq_items());
129        left
130    }
131
132    fn seq_items_reversed(self) -> Vec<Self::Item> {
133        let mut items = self.right.seq_items_reversed();
134        items.extend(self.left.seq_items_reversed());
135        items
136    }
137
138    fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
139        let mut items = self.right.seq_items_reversed_prefix(count);
140        if items.len() < count {
141            items.extend(self.left.seq_items_reversed_prefix(count - items.len()));
142        }
143        items
144    }
145
146    /// # Why this stays sequential
147    ///
148    /// Both branches could drive their own consumer and combine in order, but
149    /// `Consumer::split_at` takes the left branch's logical length, which this
150    /// adapter cannot know without consuming the left input first — the very
151    /// materialisation the conversion exists to remove. Concatenation belongs
152    /// at the same indexed producer boundary as the other length-bearing
153    /// adapters.
154    fn drive<C, R>(self, consumer: C) -> R
155    where
156        C: Consumer<Self::Item, Result = R> + Send + Sync,
157        R: Send,
158    {
159        consumer.consume(VecParIter::new(self.seq_items()))
160    }
161}
162
163/// Intersperse adapter with separator insertion between adjacent items.
164pub struct Intersperse<I>
165where
166    I: ParallelIterator,
167{
168    pub(super) base: I,
169    pub(super) separator: I::Item,
170}
171
172impl<I> Intersperse<I>
173where
174    I: ParallelIterator,
175{
176    pub(crate) fn new(base: I, separator: I::Item) -> Self {
177        Self { base, separator }
178    }
179}
180
181impl<I> ParallelIterator for Intersperse<I>
182where
183    I: ParallelIterator,
184    I::Item: Clone + Sync + 'static,
185{
186    type Item = I::Item;
187
188    fn seq_items(self) -> Vec<Self::Item> {
189        let items = self.base.seq_items();
190        if items.len() <= 1 {
191            return items;
192        }
193
194        let mut output = Vec::with_capacity(items.len().saturating_mul(2).saturating_sub(1));
195        let mut iter = items.into_iter();
196        if let Some(first) = iter.next() {
197            output.push(first);
198        }
199        for item in iter {
200            output.push(self.separator.clone());
201            output.push(item);
202        }
203        output
204    }
205
206    /// # Why this stays sequential
207    ///
208    /// The separator goes *between* adjacent items, so the first item of every
209    /// shard but the logically first needs a separator ahead of it, and the
210    /// shard cannot tell which one it is. Combining shards that each interspersed
211    /// their own range would drop exactly one separator per shard boundary.
212    fn drive<C, R>(self, consumer: C) -> R
213    where
214        C: Consumer<Self::Item, Result = R> + Send + Sync,
215        R: Send,
216    {
217        consumer.consume(VecParIter::new(self.seq_items()))
218    }
219}
220
221/// Reverse adapter with logical-order reversal semantics.
222pub struct Rev<I> {
223    pub(super) base: I,
224}
225
226impl<I> Rev<I> {
227    pub(crate) fn new(base: I) -> Self {
228        Self { base }
229    }
230}
231
232impl<I> ParallelIterator for Rev<I>
233where
234    I: ParallelIterator,
235    I::Item: Sync + 'static,
236{
237    type Item = I::Item;
238
239    fn seq_items(self) -> Vec<Self::Item> {
240        self.base.seq_items_reversed()
241    }
242
243    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
244        let count = take.unwrap_or(usize::MAX);
245        let prefix = skip.saturating_add(count);
246        let mut items = self.base.seq_items_reversed_prefix(prefix);
247        if skip >= items.len() {
248            return Vec::new();
249        }
250        items.drain(..skip);
251        if let Some(count) = take {
252            items.truncate(count);
253        }
254        items
255    }
256
257    fn seq_items_reversed(self) -> Vec<Self::Item> {
258        self.base.seq_items()
259    }
260
261    fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
262        self.base.seq_items_window(0, Some(count))
263    }
264
265    /// # Why this stays sequential
266    ///
267    /// Reversal needs both each shard's own stream reversed and the shard order
268    /// flipped, but `Consumer::combine(left, right)` fixes the merge order as
269    /// logically-earlier-first. Expressing this needs an order-reversing
270    /// consumer wrapper whose base combine is inverted, which is a protocol
271    /// change rather than an adapter push.
272    fn drive<C, R>(self, consumer: C) -> R
273    where
274        C: Consumer<Self::Item, Result = R> + Send + Sync,
275        R: Send,
276    {
277        consumer.consume(VecParIter::new(self.seq_items()))
278    }
279}