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
50 where
51 C: Consumer<Self::Item, Result = R> + Send + Sync,
52 R: Send,
53 {
54 consumer.consume(VecParIter::new(self.seq_items()))
55 }
56}
57
58impl<I, F> ParallelIterator for SkipAnyWhile<I, F>
59where
60 I: ParallelIterator,
61 F: Fn(&I::Item) -> bool + Send + Sync + Clone,
62 I::Item: Sync + 'static,
63{
64 type Item = I::Item;
65
66 fn seq_items(self) -> Vec<Self::Item> {
67 let predicate = self.predicate;
68 let mut items = self.base.seq_items().into_iter();
69
70 for item in items.by_ref() {
71 if !predicate(&item) {
72 let mut retained = vec![item];
73 retained.extend(items);
74 return retained;
75 }
76 }
77
78 Vec::new()
79 }
80
81 fn drive<C, R>(self, consumer: C) -> R
82 where
83 C: Consumer<Self::Item, Result = R> + Send + Sync,
84 R: Send,
85 {
86 consumer.consume(VecParIter::new(self.seq_items()))
87 }
88}