moirai_iter/parallel/adapters/
side_effect.rs1use super::super::{Consumer, InspectConsumer, ParallelIterator};
2
3pub 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
45pub 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}