use crate::ParExtend;
use crate::infallible::Xap;
use crate::infallible::recursive::execution;
use crate::infallible::recursive::par::ParRec;
use crate::infallible::recursive::par_core::ParRecCore;
use crate::infallible::xap::{FilMapOf, FilOf, FlatMapOf, FlattenOf, InsOf, MapOf};
use crate::parameters::{ChunkSize, IterationOrder, NumThreads, Params};
use crate::runner::{DefaultRunner, ParRunner};
use alloc::vec::Vec;
pub struct ParRecIter<I, X, Ix, Ex, R = DefaultRunner>
where
I: IntoIterator,
X: Xap<I = I::Item>,
R: ParRunner,
Ix: IntoIterator<Item = X::I>,
Ex: Fn(&I::Item) -> Ix + Send + Copy,
{
iter: I,
xap: X,
exe: R,
params: Params,
extend: Ex,
}
impl<I, X, Ix, Ex, R> ParRecIter<I, X, Ix, Ex, R>
where
I: IntoIterator,
X: Xap<I = I::Item>,
R: ParRunner,
Ix: IntoIterator<Item = X::I>,
Ex: Fn(&I::Item) -> Ix + Send + Copy,
{
pub(crate) fn new(iter: I, xap: X, exe: R, params: Params, extend: Ex) -> Self {
Self {
iter,
xap,
exe,
params,
extend,
}
}
pub(super) fn with_xap<Y: Xap<I = I::Item>>(self, xap: Y) -> ParRecIter<I, Y, Ix, Ex, R> {
ParRecIter::new(self.iter, xap, self.exe, self.params, self.extend)
}
fn destruct_x(self) -> (I, X, R, Params, Ex) {
(self.iter, self.xap, self.exe, self.params, self.extend)
}
}
impl<I, X, Ix, Ex, R> ParRecCore for ParRecIter<I, X, Ix, Ex, R>
where
I: IntoIterator,
X: Xap<I = I::Item>,
R: ParRunner,
Ix: IntoIterator<Item = X::I>,
Ex: Fn(&I::Item) -> Ix + Send + Copy,
{
type Item = X::O;
type Runner = R;
type Input = I;
type Xap = X;
fn destruct(self) -> (Self::Input, Self::Xap, Self::Runner, Params) {
(self.iter, self.xap, self.exe, self.params)
}
}
impl<I, X, Ix, Ex, R> ParRec for ParRecIter<I, X, Ix, Ex, R>
where
I: IntoIterator,
X: Xap<I = I::Item>,
R: ParRunner,
Ix: IntoIterator<Item = X::I>,
Ex: Fn(&I::Item) -> Ix + Send + Copy,
{
fn runner<Q: ParRunner>(
self,
runner: Q,
) -> impl ParRec<Item = Self::Item, Xap = Self::Xap, Input = Self::Input> {
let (iter, xap, _, params, extend) = self.destruct_x();
ParRecIter::new(iter, xap, runner, params, extend)
}
#[cfg(feature = "std")]
fn runner_with_diagnostics(
self,
) -> impl ParRec<Item = Self::Item, Xap = Self::Xap, Input = Self::Input> {
let (iter, xap, exe, params, extend) = self.destruct_x();
ParRecIter::new(iter, xap, exe.with_diagnostics(), params, extend)
}
fn num_threads(mut self, num_threads: impl Into<NumThreads>) -> Self {
self.params = self.params.with_num_threads(num_threads);
self
}
fn chunk_size(mut self, chunk_size: impl Into<ChunkSize>) -> Self {
self.params = self.params.with_chunk_size(chunk_size);
self
}
fn iteration_order(mut self, collect: IterationOrder) -> Self {
self.params = self.params.with_collect_ordering(collect);
self
}
fn map<Q, H>(
self,
h: H,
) -> impl ParRec<Item = Q, Xap = MapOf<Self::Xap, Q, H>, Input = Self::Input>
where
H: Fn(Self::Item) -> Q + Copy + Send,
{
let xap = self.xap.map(h);
self.with_xap(xap)
}
fn inspect<H>(
self,
h: H,
) -> impl ParRec<Item = Self::Item, Xap = InsOf<Self::Xap, H>, Input = Self::Input>
where
H: Fn(&Self::Item) + Copy + Send,
{
let xap = self.xap.inspect(h);
self.with_xap(xap)
}
fn filter<H>(
self,
h: H,
) -> impl ParRec<Item = Self::Item, Xap = FilOf<Self::Xap, H>, Input = Self::Input>
where
H: Fn(&Self::Item) -> bool + Copy + Send,
{
let xap = self.xap.filter(h);
self.with_xap(xap)
}
fn filter_map<Q, H>(
self,
h: H,
) -> impl ParRec<Item = Q, Xap = FilMapOf<Self::Xap, Q, H>, Input = Self::Input>
where
H: Fn(Self::Item) -> Option<Q> + Copy + Send,
{
let xap = self.xap.filter_map(h);
self.with_xap(xap)
}
fn flat_map<V, H>(
self,
h: H,
) -> impl ParRec<Item = V::Item, Xap = FlatMapOf<Self::Xap, V, H>, Input = Self::Input>
where
V: IntoIterator,
H: Fn(Self::Item) -> V + Copy + Send,
{
let xap = self.xap.flat_map(h);
self.with_xap(xap)
}
fn flatten(
self,
) -> impl ParRec<
Item = <Self::Item as IntoIterator>::Item,
Xap = FlattenOf<Self::Xap>,
Input = Self::Input,
>
where
Self::Item: IntoIterator,
{
let xap = self.xap.flatten();
self.with_xap(xap)
}
fn first(self) -> Option<Self::Item>
where
Self::Item: Send,
<Self::Input as IntoIterator>::Item: Send,
{
let (iter, x, exe, params, extend) = self.destruct_x();
match params.iteration_order {
IterationOrder::Ordered => execution::next(exe, params, iter, x, extend),
IterationOrder::Arbitrary => execution::next_any(exe, params, iter, x, extend),
}
}
fn reduce<F>(self, f: F) -> Option<Self::Item>
where
F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Copy,
Self::Item: Send,
<Self::Input as IntoIterator>::Item: Send,
{
let (iter, x, exe, params, extend) = self.destruct_x();
execution::reduce(exe, params, iter, x, extend, f)
}
fn fold<B, Id, F>(self, init: Id, f: F) -> Vec<B>
where
B: Send,
Id: Fn() -> B,
F: Fn(&mut B, Self::Item) + Copy + Send,
<Self::Input as IntoIterator>::Item: Send,
{
let (iter, x, exe, params, extend) = self.destruct_x();
execution::fold(exe, params, iter, x, extend, init, f)
}
fn collect_into<C>(self, dst: &mut C)
where
C: ParExtend<Self::Item>,
Self::Item: Send,
<Self::Input as IntoIterator>::Item: Send,
{
let (iter, x, exe, params, extend) = self.destruct_x();
match params.iteration_order {
IterationOrder::Ordered => execution::collect(exe, params, iter, x, extend, dst),
IterationOrder::Arbitrary => execution::collect_arb(exe, params, iter, x, extend, dst),
}
}
}