Skip to main content

moirai_iter/parallel/adapters/
side_effect.rs

1use super::super::{Consumer, InspectConsumer, ParallelIterator};
2
3/// Inspect adapter that observes items by shared reference without changing them.
4pub struct Inspect<I, F> {
5    base: I,
6    inspect_fn: F,
7}
8
9impl<I, F> Inspect<I, F> {
10    pub(in crate::parallel) fn new(base: I, inspect_fn: F) -> Self {
11        Self { base, inspect_fn }
12    }
13}
14
15impl<I, F> ParallelIterator for Inspect<I, F>
16where
17    I: ParallelIterator,
18    F: Fn(&I::Item) + Send + Sync + Clone,
19    I::Item: Sync + 'static,
20{
21    type Item = I::Item;
22
23    fn seq_items(self) -> Vec<Self::Item> {
24        let inspect_fn = self.inspect_fn;
25        let items = self.base.seq_items();
26        for item in &items {
27            inspect_fn(item);
28        }
29        items
30    }
31
32    fn drive<C, R>(self, consumer: C) -> R
33    where
34        C: Consumer<Self::Item, Result = R> + Send + Sync,
35        R: Send,
36    {
37        self.base
38            .drive(InspectConsumer::new(consumer, self.inspect_fn))
39    }
40}
41
42#[derive(Clone, Copy, Debug, Default)]
43struct PanicFusePolicy;
44
45/// Panic-fuse adapter for the non-indexed iterator subset.
46pub struct PanicFuse<I> {
47    base: I,
48    _policy: PanicFusePolicy,
49}
50
51impl<I> PanicFuse<I> {
52    pub(in crate::parallel) fn new(base: I) -> Self {
53        Self {
54            base,
55            _policy: PanicFusePolicy,
56        }
57    }
58}
59
60impl<I> ParallelIterator for PanicFuse<I>
61where
62    I: ParallelIterator,
63    I::Item: Sync,
64{
65    type Item = I::Item;
66
67    fn seq_items(self) -> Vec<Self::Item> {
68        self.base.seq_items()
69    }
70
71    fn drive<C, R>(self, consumer: C) -> R
72    where
73        C: Consumer<Self::Item, Result = R> + Send + Sync,
74        R: Send,
75    {
76        self.base.drive(consumer)
77    }
78}
79
80#[cfg(test)]
81mod tests {
82    use super::PanicFusePolicy;
83
84    #[test]
85    fn panic_fuse_policy_is_zero_sized() {
86        assert_eq!(std::mem::size_of::<PanicFusePolicy>(), 0);
87    }
88}