1use super::{fallible, split, TryStreamItem};
2use super::{
3 Chain, Chunks, Cloned, CollectConsumer, Copied, Enumerate, Filter, FilterMap, FindConsumer,
4 FlatMap, Flatten, Inspect, Intersperse, Map, MapInit, MapWith, NullConsumer, PanicFuse,
5 Positions, ReduceConsumer, Reduction, Rev, SequentialAdapter, Skip, SkipAnyWhile, Take,
6 TakeAnyWhile, Update, WhileSome, Zip, ZipEq,
7};
8
9pub trait ParallelIterator: Sized + Send {
11 type Item: Send;
13
14 fn drive<C, R>(self, consumer: C) -> R
30 where
31 C: Consumer<Self::Item, Result = R> + Send + Sync,
32 R: Send;
33
34 fn seq_items(self) -> Vec<Self::Item>;
36
37 fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
39 let iter = self.seq_items().into_iter().skip(skip);
40 match take {
41 Some(count) => iter.take(count).collect(),
42 None => iter.collect(),
43 }
44 }
45
46 fn seq_items_reversed(self) -> Vec<Self::Item> {
48 let mut items = self.seq_items();
49 items.reverse();
50 items
51 }
52
53 fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
55 self.seq_items_reversed().into_iter().take(count).collect()
56 }
57
58 fn map<F, R>(self, map_fn: F) -> Map<Self, F>
60 where
61 F: Fn(Self::Item) -> R + Send + Sync + Clone,
62 R: Send,
63 {
64 Map::new(self, map_fn)
65 }
66
67 fn map_with<T, F, R>(self, init: T, map_fn: F) -> MapWith<Self, T, F>
69 where
70 T: Send + Clone,
71 F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone,
72 R: Send + Sync + 'static,
73 {
74 MapWith::new(self, init, map_fn)
75 }
76
77 fn map_init<Init, T, F, R>(self, init: Init, map_fn: F) -> MapInit<Self, Init, F>
79 where
80 Init: Fn() -> T + Send + Sync + Clone,
81 T: Send,
82 F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone,
83 R: Send + Sync + 'static,
84 {
85 MapInit::new(self, init, map_fn)
86 }
87
88 fn update<F>(self, update_fn: F) -> Update<Self, F>
90 where
91 F: Fn(&mut Self::Item) + Send + Sync + Clone,
92 Self::Item: Sync + 'static,
93 {
94 Update::new(self, update_fn)
95 }
96
97 fn filter<F>(self, filter_fn: F) -> Filter<Self, F>
99 where
100 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
101 {
102 Filter::new(self, filter_fn)
103 }
104
105 fn inspect<F>(self, inspect_fn: F) -> Inspect<Self, F>
107 where
108 F: Fn(&Self::Item) + Send + Sync + Clone,
109 Self::Item: Sync,
110 {
111 Inspect::new(self, inspect_fn)
112 }
113
114 fn panic_fuse(self) -> PanicFuse<Self>
116 where
117 Self::Item: Sync,
118 {
119 PanicFuse::new(self)
120 }
121
122 fn filter_map<F, R>(self, filter_map_fn: F) -> FilterMap<Self, F>
124 where
125 F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
126 R: Send + Sync + 'static,
127 {
128 FilterMap::new(self, filter_map_fn)
129 }
130
131 fn while_some<T>(self) -> WhileSome<Self>
133 where
134 Self: ParallelIterator<Item = Option<T>>,
135 T: Send + Sync + 'static,
136 {
137 WhileSome::new(self)
138 }
139
140 fn flat_map<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F>
142 where
143 F: Fn(Self::Item) -> U + Send + Sync + Clone,
144 U: IntoIterator,
145 U::Item: Send + Sync + 'static,
146 {
147 FlatMap::new(self, flat_map_fn)
148 }
149
150 fn flat_map_iter<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F>
152 where
153 F: Fn(Self::Item) -> U + Send + Sync + Clone,
154 U: IntoIterator,
155 U::Item: Send + Sync + 'static,
156 {
157 FlatMap::new(self, flat_map_fn)
158 }
159
160 fn flatten(self) -> Flatten<Self>
162 where
163 Self::Item: IntoIterator,
164 <Self::Item as IntoIterator>::Item: Send + Sync + 'static,
165 {
166 Flatten::new(self)
167 }
168
169 fn flatten_iter(self) -> Flatten<Self>
171 where
172 Self::Item: IntoIterator,
173 <Self::Item as IntoIterator>::Item: Send + Sync + 'static,
174 {
175 Flatten::new(self)
176 }
177
178 fn enumerate(self) -> Enumerate<Self>
180 where
181 Self::Item: Sync + 'static,
182 {
183 Enumerate::new(self)
184 }
185
186 fn zip<J>(self, other: J) -> Zip<Self, J>
188 where
189 J: ParallelIterator,
190 Self::Item: Sync + 'static,
191 J::Item: Sync + 'static,
192 {
193 Zip::new(self, other)
194 }
195
196 fn zip_eq<J>(self, other: J) -> ZipEq<Self, J>
198 where
199 J: ParallelIterator,
200 Self::Item: Sync + 'static,
201 J::Item: Sync + 'static,
202 {
203 ZipEq::new(self, other)
204 }
205
206 fn take(self, count: usize) -> Take<Self>
208 where
209 Self::Item: Sync + 'static,
210 {
211 Take::new(self, count)
212 }
213
214 fn take_any(self, count: usize) -> Take<Self>
216 where
217 Self::Item: Sync + 'static,
218 {
219 Take::new(self, count)
220 }
221
222 fn skip(self, count: usize) -> Skip<Self>
224 where
225 Self::Item: Sync + 'static,
226 {
227 Skip::new(self, count)
228 }
229
230 fn skip_any(self, count: usize) -> Skip<Self>
232 where
233 Self::Item: Sync + 'static,
234 {
235 Skip::new(self, count)
236 }
237
238 fn take_any_while<F>(self, predicate: F) -> TakeAnyWhile<Self, F>
240 where
241 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
242 Self::Item: Sync + 'static,
243 {
244 TakeAnyWhile::new(self, predicate)
245 }
246
247 fn skip_any_while<F>(self, predicate: F) -> SkipAnyWhile<Self, F>
249 where
250 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
251 Self::Item: Sync + 'static,
252 {
253 SkipAnyWhile::new(self, predicate)
254 }
255
256 fn chain<J>(self, other: J) -> Chain<Self, J>
258 where
259 J: ParallelIterator<Item = Self::Item>,
260 Self::Item: Sync + 'static,
261 {
262 Chain::new(self, other)
263 }
264
265 fn intersperse(self, separator: Self::Item) -> Intersperse<Self>
267 where
268 Self::Item: Clone + Sync + 'static,
269 {
270 Intersperse::new(self, separator)
271 }
272
273 fn rev(self) -> Rev<Self>
275 where
276 Self::Item: Sync + 'static,
277 {
278 Rev::new(self)
279 }
280
281 fn chunks(self, chunk_size: usize) -> Chunks<Self>
283 where
284 Self::Item: Sync + 'static,
285 {
286 Chunks::new(self, chunk_size)
287 }
288
289 fn copied<'data, T>(self) -> Copied<Self>
291 where
292 Self: ParallelIterator<Item = &'data T>,
293 T: Copy + Send + Sync + 'data + 'static,
294 {
295 Copied::new(self)
296 }
297
298 fn cloned<'data, T>(self) -> Cloned<Self>
300 where
301 Self: ParallelIterator<Item = &'data T>,
302 T: Clone + Send + Sync + 'data + 'static,
303 {
304 Cloned::new(self)
305 }
306
307 fn reduce<F>(self, reduce_fn: F) -> Option<Self::Item>
309 where
310 F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone,
311 Self::Item: Clone + Sync,
312 {
313 let reduction: Reduction<Self::Item, F> = self.drive(ReduceConsumer::new(reduce_fn));
314 reduction.into_value()
315 }
316
317 fn fold<T, F>(self, init: T, fold_fn: F) -> T
319 where
320 T: Send + Sync + Clone,
321 F: Fn(T, Self::Item) -> T + Send + Sync + Clone,
322 Self::Item: Sync,
323 {
324 self.drive(CollectConsumer::new())
328 .into_iter()
329 .fold(init, fold_fn)
330 }
331
332 fn collect<C>(self) -> C
334 where
335 C: ParallelExtend<Self::Item> + Default + Send,
336 {
337 let mut collection = C::default();
338 collection.par_extend(self);
339 collection
340 }
341
342 fn collect_vec_list(self) -> std::collections::LinkedList<Vec<Self::Item>> {
349 let items = self.seq_items();
350 let mut list = std::collections::LinkedList::new();
351 if !items.is_empty() {
352 list.push_back(items);
353 }
354 list
355 }
356
357 fn partition<C, F>(self, predicate: F) -> (C, C)
359 where
360 C: FromIterator<Self::Item> + Send,
361 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
362 Self::Item: Sync + 'static,
363 {
364 let (left_items, right_items): (Vec<Self::Item>, Vec<Self::Item>) = self
365 .seq_items()
366 .into_iter()
367 .partition(|item| predicate(item));
368
369 (
370 left_items.into_iter().collect(),
371 right_items.into_iter().collect(),
372 )
373 }
374
375 fn partition_map<A, B, P, L, R>(self, predicate: P) -> (A, B)
377 where
378 A: Default + Extend<L> + Send,
379 B: Default + Extend<R> + Send,
380 P: Fn(Self::Item) -> split::Either<L, R> + Send + Sync + Clone,
381 L: Send,
382 R: Send,
383 {
384 split::partition_map(self, predicate)
385 }
386
387 fn unzip<A, B, FromA, FromB>(self) -> (FromA, FromB)
389 where
390 Self: ParallelIterator<Item = (A, B)>,
391 FromA: Default + Extend<A> + Send,
392 FromB: Default + Extend<B> + Send,
393 A: Send,
394 B: Send,
395 {
396 self.seq_items().into_iter().unzip()
397 }
398
399 fn sequential(self) -> SequentialAdapter<Self> {
401 SequentialAdapter::new(self)
402 }
403
404 fn count(self) -> usize
406 where
407 Self::Item: Sync,
408 {
409 self.drive(CollectConsumer::new()).len()
410 }
411
412 fn find_first<F>(self, predicate: F) -> Option<Self::Item>
414 where
415 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
416 Self::Item: Sync,
417 {
418 self.find_any(predicate)
419 }
420
421 fn find_last<F>(self, predicate: F) -> Option<Self::Item>
423 where
424 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
425 {
426 self.seq_items().into_iter().rev().find(predicate)
427 }
428
429 fn position_first<F>(self, predicate: F) -> Option<usize>
431 where
432 F: Fn(Self::Item) -> bool + Send + Sync + Clone,
433 {
434 self.seq_items().into_iter().position(predicate)
435 }
436
437 fn position_any<F>(self, predicate: F) -> Option<usize>
439 where
440 F: Fn(Self::Item) -> bool + Send + Sync + Clone,
441 {
442 self.position_first(predicate)
443 }
444
445 fn position_last<F>(self, predicate: F) -> Option<usize>
447 where
448 F: Fn(Self::Item) -> bool + Send + Sync + Clone,
449 {
450 self.seq_items().into_iter().rposition(predicate)
451 }
452
453 fn positions<F>(self, predicate: F) -> Positions<Self, F>
455 where
456 F: Fn(Self::Item) -> bool + Send + Sync + Clone,
457 {
458 Positions::new(self, predicate)
459 }
460
461 fn find_map_first<F, R>(self, map_fn: F) -> Option<R>
463 where
464 F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
465 R: Send,
466 {
467 self.seq_items().into_iter().find_map(map_fn)
468 }
469
470 fn find_map_any<F, R>(self, map_fn: F) -> Option<R>
472 where
473 F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
474 R: Send,
475 {
476 self.find_map_first(map_fn)
477 }
478
479 fn find_map_last<F, R>(self, map_fn: F) -> Option<R>
481 where
482 F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
483 R: Send,
484 {
485 self.seq_items().into_iter().rev().find_map(map_fn)
486 }
487
488 fn any<F>(self, predicate: F) -> bool
490 where
491 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
492 Self::Item: Sync,
493 {
494 self.find_any(predicate).is_some()
495 }
496
497 fn all<F>(self, predicate: F) -> bool
499 where
500 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
501 Self::Item: Sync,
502 {
503 self.find_any(move |item| !predicate(item)).is_none()
504 }
505
506 fn for_each<F>(self, op: F)
508 where
509 F: Fn(Self::Item) + Send + Sync + Clone,
510 {
511 self.map(op).drive(NullConsumer::new())
512 }
513
514 fn for_each_with<T, F>(self, init: T, op: F)
516 where
517 T: Send + Clone,
518 F: Fn(&mut T, Self::Item) + Send + Sync + Clone,
519 {
520 let mut state = init;
521 for item in self.seq_items() {
522 op(&mut state, item);
523 }
524 }
525
526 fn for_each_init<Init, T, F>(self, init: Init, op: F)
528 where
529 Init: Fn() -> T + Send + Sync + Clone,
530 T: Send,
531 F: Fn(&mut T, Self::Item) + Send + Sync + Clone,
532 {
533 let mut state = init();
534 for item in self.seq_items() {
535 op(&mut state, item);
536 }
537 }
538
539 fn try_for_each<F, E>(self, op: F) -> Result<(), E>
541 where
542 F: Fn(Self::Item) -> Result<(), E> + Send + Sync + Clone,
543 E: Send,
544 {
545 for item in self.seq_items() {
546 op(item)?;
547 }
548 Ok(())
549 }
550
551 fn try_for_each_with<T, F, E>(self, init: T, op: F) -> Result<(), E>
553 where
554 T: Send + Clone,
555 F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone,
556 E: Send,
557 {
558 let mut state = init;
559 for item in self.seq_items() {
560 op(&mut state, item)?;
561 }
562 Ok(())
563 }
564
565 fn try_for_each_init<Init, T, F, E>(self, init: Init, op: F) -> Result<(), E>
567 where
568 Init: Fn() -> T + Send + Sync + Clone,
569 T: Send,
570 F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone,
571 E: Send,
572 {
573 let mut state = init();
574 for item in self.seq_items() {
575 op(&mut state, item)?;
576 }
577 Ok(())
578 }
579
580 fn reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item>
582 where
583 F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone,
584 Self::Item: Sync + Clone,
585 {
586 let reduction: Reduction<Self::Item, F> = self.drive(ReduceConsumer::new(reduce_fn));
587 reduction.into_value()
588 }
589
590 fn try_reduce<Identity, F, T, E>(self, identity: Identity, reduce_fn: F) -> Result<T, E>
592 where
593 Self::Item: Into<Result<T, E>>,
594 Identity: Fn() -> T + Send + Sync + Clone,
595 F: Fn(T, T) -> Result<T, E> + Send + Sync + Clone,
596 T: Send,
597 E: Send,
598 {
599 let mut accumulator = identity();
600 for item in self.seq_items() {
601 accumulator = reduce_fn(accumulator, item.into()?)?;
602 }
603 Ok(accumulator)
604 }
605
606 fn try_reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item>
608 where
609 Self::Item: TryStreamItem,
610 F: Fn(
611 <Self::Item as TryStreamItem>::Output,
612 <Self::Item as TryStreamItem>::Output,
613 ) -> Self::Item
614 + Send
615 + Sync
616 + Clone,
617 {
618 fallible::try_reduce_with(self, reduce_fn)
619 }
620
621 fn sum<S>(self) -> S
623 where
624 S: std::iter::Sum<Self::Item> + Send,
625 {
626 self.seq_items().into_iter().sum()
627 }
628
629 fn product<P>(self) -> P
631 where
632 P: std::iter::Product<Self::Item> + Send,
633 {
634 self.seq_items().into_iter().product()
635 }
636
637 fn min(self) -> Option<Self::Item>
639 where
640 Self::Item: Ord,
641 {
642 self.seq_items().into_iter().min()
643 }
644
645 fn max(self) -> Option<Self::Item>
647 where
648 Self::Item: Ord,
649 {
650 self.seq_items().into_iter().max()
651 }
652
653 fn min_by<F>(self, compare: F) -> Option<Self::Item>
655 where
656 F: Fn(&Self::Item, &Self::Item) -> std::cmp::Ordering + Send + Sync + Clone,
657 {
658 self.seq_items().into_iter().min_by(compare)
659 }
660
661 fn max_by<F>(self, compare: F) -> Option<Self::Item>
663 where
664 F: Fn(&Self::Item, &Self::Item) -> std::cmp::Ordering + Send + Sync + Clone,
665 {
666 self.seq_items().into_iter().max_by(compare)
667 }
668
669 fn min_by_key<K, F>(self, key_fn: F) -> Option<Self::Item>
671 where
672 K: Ord,
673 F: Fn(&Self::Item) -> K + Send + Sync + Clone,
674 {
675 self.seq_items().into_iter().min_by_key(key_fn)
676 }
677
678 fn max_by_key<K, F>(self, key_fn: F) -> Option<Self::Item>
680 where
681 K: Ord,
682 F: Fn(&Self::Item) -> K + Send + Sync + Clone,
683 {
684 self.seq_items().into_iter().max_by_key(key_fn)
685 }
686
687 fn find_any<F>(self, predicate: F) -> Option<Self::Item>
689 where
690 F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
691 Self::Item: Sync,
692 {
693 self.drive(FindConsumer::new(predicate))
694 }
695}
696
697pub trait Consumer<T>: Send + Sync {
699 type Result: Send;
700
701 fn consume<I>(self, iter: I) -> Self::Result
703 where
704 I: ParallelIterator<Item = T>;
705
706 fn split_at(self, index: usize) -> (Self, Self)
708 where
709 Self: Sized;
710
711 fn combine(left: Self::Result, right: Self::Result) -> Self::Result;
713}
714
715pub trait ParallelExtend<T>: Send {
717 fn par_extend<I>(&mut self, par_iter: I)
719 where
720 I: ParallelIterator<Item = T>;
721}
722
723pub trait IntoParallelIterator {
725 type Item: Send;
726 type Iter: ParallelIterator<Item = Self::Item>;
727
728 fn into_par_iter(self) -> Self::Iter;
729}
730
731pub trait IntoParallelRefIterator<'data> {
733 type Item: Send + Sync + 'data;
734 type Iter: ParallelIterator<Item = Self::Item>;
735
736 fn par_iter(&'data self) -> Self::Iter;
737}