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    fn drive<C, R>(self, consumer: C) -> R
53    where
54        C: Consumer<Self::Item, Result = R> + Send + Sync,
55        R: Send,
56    {
57        consumer.consume(VecParIter::new(self.seq_items()))
58    }
59}
60
61/// Skip adapter with prefix-discarding value semantics.
62pub struct Skip<I> {
63    pub(super) base: I,
64    pub(super) count: usize,
65}
66
67impl<I> Skip<I> {
68    pub(crate) fn new(base: I, count: usize) -> Self {
69        Self { base, count }
70    }
71}
72
73impl<I> ParallelIterator for Skip<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        self.base.seq_items_window(self.count, None)
82    }
83
84    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
85        self.base
86            .seq_items_window(self.count.saturating_add(skip), take)
87    }
88
89    fn drive<C, R>(self, consumer: C) -> R
90    where
91        C: Consumer<Self::Item, Result = R> + Send + Sync,
92        R: Send,
93    {
94        consumer.consume(VecParIter::new(self.seq_items()))
95    }
96}
97
98/// Chain adapter with left-to-right concatenation semantics.
99pub struct Chain<I, J> {
100    pub(super) left: I,
101    pub(super) right: J,
102}
103
104impl<I, J> Chain<I, J> {
105    pub(crate) fn new(left: I, right: J) -> Self {
106        Self { left, right }
107    }
108}
109
110impl<I, J> ParallelIterator for Chain<I, J>
111where
112    I: ParallelIterator,
113    J: ParallelIterator<Item = I::Item>,
114    I::Item: Sync + 'static,
115{
116    type Item = I::Item;
117
118    fn seq_items(self) -> Vec<Self::Item> {
119        let mut left = self.left.seq_items();
120        left.extend(self.right.seq_items());
121        left
122    }
123
124    fn seq_items_reversed(self) -> Vec<Self::Item> {
125        let mut items = self.right.seq_items_reversed();
126        items.extend(self.left.seq_items_reversed());
127        items
128    }
129
130    fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
131        let mut items = self.right.seq_items_reversed_prefix(count);
132        if items.len() < count {
133            items.extend(self.left.seq_items_reversed_prefix(count - items.len()));
134        }
135        items
136    }
137
138    fn drive<C, R>(self, consumer: C) -> R
139    where
140        C: Consumer<Self::Item, Result = R> + Send + Sync,
141        R: Send,
142    {
143        consumer.consume(VecParIter::new(self.seq_items()))
144    }
145}
146
147/// Intersperse adapter with separator insertion between adjacent items.
148pub struct Intersperse<I>
149where
150    I: ParallelIterator,
151{
152    pub(super) base: I,
153    pub(super) separator: I::Item,
154}
155
156impl<I> Intersperse<I>
157where
158    I: ParallelIterator,
159{
160    pub(crate) fn new(base: I, separator: I::Item) -> Self {
161        Self { base, separator }
162    }
163}
164
165impl<I> ParallelIterator for Intersperse<I>
166where
167    I: ParallelIterator,
168    I::Item: Clone + Sync + 'static,
169{
170    type Item = I::Item;
171
172    fn seq_items(self) -> Vec<Self::Item> {
173        let items = self.base.seq_items();
174        if items.len() <= 1 {
175            return items;
176        }
177
178        let mut output = Vec::with_capacity(items.len().saturating_mul(2).saturating_sub(1));
179        let mut iter = items.into_iter();
180        if let Some(first) = iter.next() {
181            output.push(first);
182        }
183        for item in iter {
184            output.push(self.separator.clone());
185            output.push(item);
186        }
187        output
188    }
189
190    fn drive<C, R>(self, consumer: C) -> R
191    where
192        C: Consumer<Self::Item, Result = R> + Send + Sync,
193        R: Send,
194    {
195        consumer.consume(VecParIter::new(self.seq_items()))
196    }
197}
198
199/// Reverse adapter with logical-order reversal semantics.
200pub struct Rev<I> {
201    pub(super) base: I,
202}
203
204impl<I> Rev<I> {
205    pub(crate) fn new(base: I) -> Self {
206        Self { base }
207    }
208}
209
210impl<I> ParallelIterator for Rev<I>
211where
212    I: ParallelIterator,
213    I::Item: Sync + 'static,
214{
215    type Item = I::Item;
216
217    fn seq_items(self) -> Vec<Self::Item> {
218        self.base.seq_items_reversed()
219    }
220
221    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
222        let count = take.unwrap_or(usize::MAX);
223        let prefix = skip.saturating_add(count);
224        let mut items = self.base.seq_items_reversed_prefix(prefix);
225        if skip >= items.len() {
226            return Vec::new();
227        }
228        items.drain(..skip);
229        if let Some(count) = take {
230            items.truncate(count);
231        }
232        items
233    }
234
235    fn seq_items_reversed(self) -> Vec<Self::Item> {
236        self.base.seq_items()
237    }
238
239    fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
240        self.base.seq_items_window(0, Some(count))
241    }
242
243    fn drive<C, R>(self, consumer: C) -> R
244    where
245        C: Consumer<Self::Item, Result = R> + Send + Sync,
246        R: Send,
247    {
248        consumer.consume(VecParIter::new(self.seq_items()))
249    }
250}