orx-parallel 4.0.0

Performant parallel computations with an expressive iterator API.
Documentation
use crate::infallible::recursive::utils;
use crate::{Par, ParDrain, ParUse, Params, ThreadPool, infallible::Xap, runner::ParRunner};
use alloc::vec::Vec;

struct Local<I, B> {
    input: Vec<I>,
    fold: B,
}

impl<I, B> Local<I, B> {
    fn new(init: B) -> Self {
        Self {
            input: Default::default(),
            fold: init,
        }
    }
}

pub fn fold<R, C, X, B, Id, F, I, E>(
    mut runner: R,
    params: Params,
    iter: C,
    xap: X,
    extend: E,
    init: Id,
    f: F,
) -> Vec<B>
where
    R: ParRunner,
    C: IntoIterator,
    X: Xap<I = C::Item>,
    I: IntoIterator<Item = X::I>,
    E: Fn(&X::I) -> I + Send + Copy,
    B: Send,
    Id: Fn() -> B,
    F: Fn(&mut B, X::O) + Copy + Send,
    X::I: Send,
{
    let max_threads: usize = runner.pool().max_num_threads().into();

    let mut data: Vec<_> = (0..max_threads)
        .map(|_| Local::<X::I, B>::new(init()))
        .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);

    par.for_each(move |u, i| {
        u.input.extend(extend(&i));

        for i in xap.xap(i) {
            f(&mut u.fold, i);
        }
    });
    utils::into_outer_par(&mut inputs, &mut data, |x| &mut x.input, &mut runner);

    while !inputs.is_empty() {
        let par = inputs.par_drain(..).runner(&mut runner);
        let par = params.apply(par).use_slice(&mut data);

        par.for_each(move |u, i| {
            u.input.extend(extend(&i));

            for i in xap.xap(i) {
                f(&mut u.fold, i);
            }
        });
        utils::into_outer_par(&mut inputs, &mut data, |x| &mut x.input, &mut runner);
    }

    data.into_iter().map(|x| x.fold).collect()
}