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    /// # Why this stays sequential
50    ///
51    /// The retained prefix ends at the first item failing the predicate
52    /// anywhere in the stream, so a shard cannot decide its own items without
53    /// knowing whether an earlier shard already stopped — the prefix dependency
54    /// documented on [`WhileSome`](super::filter::WhileSome).
55    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    /// # Why this stays sequential
88    ///
89    /// The discarded prefix ends at the first item failing the predicate
90    /// anywhere in the stream, a prefix dependency per [`TakeAnyWhile`].
91    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}