Skip to main content

moirai_iter/parallel/adapters/
window.rs

1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3/// Deterministic predicate-window adapter for the audited `take_any_while` subset.
4pub struct TakeAnyWhile<I, F> {
5    base: I,
6    predicate: F,
7}
8
9/// Deterministic predicate-window adapter for the audited `skip_any_while` subset.
10pub 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}