moirai_iter/parallel/adapters/
window.rs1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3pub struct TakeAnyWhile<I, F> {
5 base: I,
6 predicate: F,
7}
8
9pub struct SkipAnyWhile<I, F> {
11 base: I,
12 predicate: F,
13}
14
15impl<I, F> TakeAnyWhile<I, F> {
16 pub(in crate::parallel) fn new(base: I, predicate: F) -> Self {
17 Self { base, predicate }
18 }
19}
20
21impl<I, F> SkipAnyWhile<I, F> {
22 pub(in crate::parallel) fn new(base: I, predicate: F) -> Self {
23 Self { base, predicate }
24 }
25}
26
27impl<I, F> ParallelIterator for TakeAnyWhile<I, F>
28where
29 I: ParallelIterator,
30 F: Fn(&I::Item) -> bool + Send + Sync + Clone,
31 I::Item: Sync + 'static,
32{
33 type Item = I::Item;
34
35 fn seq_items(self) -> Vec<Self::Item> {
36 let predicate = self.predicate;
37 let mut retained = Vec::new();
38
39 for item in self.base.seq_items() {
40 if !predicate(&item) {
41 break;
42 }
43 retained.push(item);
44 }
45
46 retained
47 }
48
49 fn drive<C, R>(self, consumer: C) -> R
56 where
57 C: Consumer<Self::Item, Result = R> + Send + Sync,
58 R: Send,
59 {
60 consumer.consume(VecParIter::new(self.seq_items()))
61 }
62}
63
64impl<I, F> ParallelIterator for SkipAnyWhile<I, F>
65where
66 I: ParallelIterator,
67 F: Fn(&I::Item) -> bool + Send + Sync + Clone,
68 I::Item: Sync + 'static,
69{
70 type Item = I::Item;
71
72 fn seq_items(self) -> Vec<Self::Item> {
73 let predicate = self.predicate;
74 let mut items = self.base.seq_items().into_iter();
75
76 for item in items.by_ref() {
77 if !predicate(&item) {
78 let mut retained = vec![item];
79 retained.extend(items);
80 return retained;
81 }
82 }
83
84 Vec::new()
85 }
86
87 fn drive<C, R>(self, consumer: C) -> R
92 where
93 C: Consumer<Self::Item, Result = R> + Send + Sync,
94 R: Send,
95 {
96 consumer.consume(VecParIter::new(self.seq_items()))
97 }
98}