use crate::infallible::recursive::utils;
use crate::{Par, ParDrain, ParUse, Params, ThreadPool, infallible::Xap, runner::ParRunner};
use alloc::vec::Vec;
pub fn reduce<R, C, X, F, I, E>(
mut runner: R,
params: Params,
iter: C,
xap: X,
extend: E,
f: F,
) -> Option<X::O>
where
R: ParRunner,
C: IntoIterator,
X: Xap<I = C::Item>,
I: IntoIterator<Item = X::I>,
E: Fn(&X::I) -> I + Send + Copy,
F: Fn(X::O, X::O) -> X::O + Send + Copy,
X::O: Send,
X::I: Send,
{
let max_threads: usize = runner.pool().max_num_threads().into();
let mut data: Vec<_> = (0..max_threads).map(|_| Vec::<X::I>::new()).collect();
let mut inputs: Vec<_> = iter.into_iter().collect();
let par = inputs.par_drain(..).runner(&mut runner);
let par = params.apply(par).use_slice(&mut data);
let mut result = par
.flat_map(move |u, i| {
u.extend(extend(&i));
xap.xap(i)
})
.reduce(move |_, a, b| f(a, b));
utils::into_outer_par(&mut inputs, &mut data, |x| x, &mut runner);
while !inputs.is_empty() {
let par = inputs.par_drain(..).runner(&mut runner);
let par = params.apply(par).use_slice(&mut data);
let result_wave = par
.flat_map(move |u, i| {
u.extend(extend(&i));
xap.xap(i)
})
.reduce(move |_, a, b| f(a, b));
result = match (result, result_wave) {
(Some(a), Some(b)) => Some(f(a, b)),
(Some(a), None) => Some(a),
(None, Some(b)) => Some(b),
(None, None) => None,
};
utils::into_outer_par(&mut inputs, &mut data, |x| x, &mut runner);
}
result
}