1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78
use serde::{Deserialize, Serialize}; use std::marker::PhantomData; use super::{ DistributedIteratorMulti, DistributedReducer, PushReducer, ReduceFactory, Reducer, ReducerA }; use crate::pool::ProcessSend; #[must_use] pub struct ForEach<I, F> { i: I, f: F, } impl<I, F> ForEach<I, F> { pub(super) fn new(i: I, f: F) -> Self { Self { i, f } } } impl<I: DistributedIteratorMulti<Source>, Source, F> DistributedReducer<I, Source, ()> for ForEach<I, F> where F: FnMut(I::Item) + Clone + ProcessSend, I::Item: 'static, { type ReduceAFactory = ForEachReducerFactory<I::Item, F>; type ReduceA = ForEachReducer<I::Item, F>; type ReduceB = PushReducer<()>; fn reducers(self) -> (I, Self::ReduceAFactory, Self::ReduceB) { ( self.i, ForEachReducerFactory(self.f, PhantomData), PushReducer((), PhantomData), ) } } pub struct ForEachReducerFactory<A, F>(F, PhantomData<fn(A)>); impl<A, F> ReduceFactory for ForEachReducerFactory<A, F> where F: FnMut(A) + Clone, { type Reducer = ForEachReducer<A, F>; fn make(&self) -> Self::Reducer { ForEachReducer(self.0.clone(), PhantomData) } } #[derive(Serialize, Deserialize)] #[serde( bound(serialize = "F: Serialize"), bound(deserialize = "F: Deserialize<'de>") )] pub struct ForEachReducer<A, F>(F, PhantomData<fn(A)>); impl<A, F> Reducer for ForEachReducer<A, F> where F: FnMut(A) + Clone, { type Item = A; type Output = (); #[inline(always)] fn push(&mut self, item: Self::Item) -> bool { self.0(item); true } fn ret(self) -> Self::Output {} } impl<A, F> ReducerA for ForEachReducer<A, F> where A: 'static, F: FnMut(A) + Clone + ProcessSend, { type Output = (); }