use std::marker::PhantomData;
use crate::task_queue::TaskQueue;
use crate::utils;
use crate::worker_thread::WorkerThreads;
use super::iterators::iterator::AtomicIterator;
pub trait ParallelForEachIter<I, V,F>
where Self: AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) + Send + Sync,
V: Send + Sync,
{
fn for_each(self,f:F) {
ParallelForEach::new(self,f).run()
}
}
impl<I,V,F> ParallelForEachIter<I, V,F> for I
where I: AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) + Send + Sync,
V: Send + Sync,
{}
pub struct ParallelForEach<V,F,I>
where I: AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) + Send + Sync,
V: Send
{
pub iter: TaskQueue<I,V>,
pub f:F,
pub num_threads:usize,
pub v: PhantomData<V>,
}
#[allow(dead_code)]
impl<I,V,F> ParallelForEach<V,F,I>
where I:AtomicIterator<AtomicItem = V> + Send + Sized,
F: Fn(V) + Send + Sync,
V: Send + Sync,
{
pub fn new(iter:I,f:F) -> Self
{
Self {
iter: TaskQueue { iter },
f,
num_threads: Self::max_threads(),
v: 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 run(self)
{
let num_threads = self.num_threads;
WorkerThreads { nthreads: num_threads }
.run(self)
}
}