moirai_iter/parallel/adapters/
slice_ops.rs1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3pub 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
61pub 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
98pub 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
147pub 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
199pub 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}