use std::marker::PhantomData;
use crate::task_queue::TaskQueue;
use crate::utils;
use crate::worker_thread::WorkerThreads;
use super::iterators::iterator::*;
use super::collector::*;
pub trait ParallelMapIter<I, V,F,T>
where Self: AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) -> T + Send + Sync,
V: Send + Sync,
T:Send + Sync
{
fn map(self,f:F) -> ParallelMap<V,F,T,Self>{
ParallelMap::new(self,f)
}
}
impl<I,V,F,T> ParallelMapIter<I, V,F,T> for I
where I: AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) -> T + Send + Sync,
V: Send + Sync,
T:Send + Sync {}
pub struct ParallelMap<V,F,T,I>
where I: AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) -> T + Send + Sync,
V: Send,
T:Send
{
pub iter: TaskQueue<I,V>,
pub f:F,
pub num_threads:usize,
pub v: PhantomData<V>,
pub t: PhantomData<T>,
}
#[allow(dead_code)]
impl<I,V,F,T> ParallelMap<V,F,T,I>
where I:AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) -> T + Send + Sync,
V: Send + Sync,
T:Send + Sync
{
pub fn new(iter:I,f:F) -> Self
{
Self {
iter: TaskQueue { iter },
f,
num_threads: Self::max_threads(),
v: PhantomData,
t: PhantomData
}
}
fn max_threads() -> usize {
utils::max_threads()
}
pub fn threads(mut self, nthreads:usize) -> Self {
self.num_threads = usize::min(Self::max_threads(),nthreads);
self
}
pub fn collect<C>(self) -> C
where C: Collector<T>
{
let num_threads = self.num_threads;
WorkerThreads { nthreads: num_threads }
.collect(self)
}
}