2024-09-29 17:21:11 +02:00
|
|
|
use std::convert::identity;
|
2024-09-26 18:43:23 +02:00
|
|
|
use std::sync::Arc;
|
|
|
|
|
2024-08-28 18:45:16 +02:00
|
|
|
use crossbeam_channel::{Receiver, Sender, TryRecvError};
|
2024-09-26 17:46:58 +02:00
|
|
|
use rayon::iter::{MapInit, ParallelIterator};
|
|
|
|
|
|
|
|
pub trait ParallelIteratorExt: ParallelIterator {
|
2024-09-29 16:46:58 +02:00
|
|
|
/// Maps items based on the init function.
|
|
|
|
///
|
|
|
|
/// The init function is ran only as necessary which is basically once by thread.
|
2024-09-26 17:46:58 +02:00
|
|
|
fn try_map_try_init<F, INIT, T, E, R>(
|
|
|
|
self,
|
|
|
|
init: INIT,
|
|
|
|
map_op: F,
|
|
|
|
) -> MapInit<
|
|
|
|
Self,
|
2024-09-26 18:43:23 +02:00
|
|
|
impl Fn() -> Result<T, Arc<E>> + Sync + Send + Clone,
|
|
|
|
impl Fn(&mut Result<T, Arc<E>>, Self::Item) -> Result<R, Arc<E>> + Sync + Send + Clone,
|
2024-09-26 17:46:58 +02:00
|
|
|
>
|
|
|
|
where
|
2024-09-26 18:43:23 +02:00
|
|
|
E: Send + Sync,
|
2024-09-26 17:46:58 +02:00
|
|
|
F: Fn(&mut T, Self::Item) -> Result<R, E> + Sync + Send + Clone,
|
|
|
|
INIT: Fn() -> Result<T, E> + Sync + Send + Clone,
|
|
|
|
R: Send,
|
|
|
|
{
|
|
|
|
self.map_init(
|
|
|
|
move || match init() {
|
|
|
|
Ok(t) => Ok(t),
|
2024-09-26 18:43:23 +02:00
|
|
|
Err(err) => Err(Arc::new(err)),
|
2024-09-26 17:46:58 +02:00
|
|
|
},
|
2024-09-29 17:21:11 +02:00
|
|
|
move |result, item| match result {
|
2024-09-26 18:43:23 +02:00
|
|
|
Ok(t) => map_op(t, item).map_err(Arc::new),
|
2024-09-29 17:21:11 +02:00
|
|
|
Err(err) => Err(err.clone()),
|
2024-09-26 17:46:58 +02:00
|
|
|
},
|
|
|
|
)
|
|
|
|
}
|
2024-09-29 17:21:11 +02:00
|
|
|
|
|
|
|
/// A method to run a closure of all the items and return an owned error.
|
|
|
|
///
|
|
|
|
/// The init function is ran only as necessary which is basically once by thread.
|
|
|
|
fn try_for_each_try_init<F, INIT, T, E>(self, init: INIT, op: F) -> Result<(), E>
|
|
|
|
where
|
|
|
|
E: Send + Sync,
|
|
|
|
F: Fn(&mut T, Self::Item) -> Result<(), Arc<E>> + Sync + Send + Clone,
|
|
|
|
INIT: Fn() -> Result<T, E> + Sync + Send + Clone,
|
|
|
|
{
|
|
|
|
let result = self.try_for_each_init(
|
|
|
|
move || match init() {
|
|
|
|
Ok(t) => Ok(t),
|
|
|
|
Err(err) => Err(Arc::new(err)),
|
|
|
|
},
|
|
|
|
move |result, item| match result {
|
|
|
|
Ok(t) => op(t, item),
|
|
|
|
Err(err) => Err(err.clone()),
|
|
|
|
},
|
|
|
|
);
|
|
|
|
|
|
|
|
match result {
|
|
|
|
Ok(()) => Ok(()),
|
|
|
|
Err(err) => Err(Arc::into_inner(err).expect("the error must be only owned by us")),
|
|
|
|
}
|
|
|
|
}
|
2024-09-26 17:46:58 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
impl<T: ParallelIterator> ParallelIteratorExt for T {}
|
2024-08-28 18:45:16 +02:00
|
|
|
|
|
|
|
/// A pool of items that can be pull and generated on demand.
|
|
|
|
pub struct ItemsPool<F, T, E>
|
|
|
|
where
|
|
|
|
F: Fn() -> Result<T, E>,
|
|
|
|
{
|
|
|
|
init: F,
|
|
|
|
sender: Sender<T>,
|
|
|
|
receiver: Receiver<T>,
|
|
|
|
}
|
|
|
|
|
|
|
|
impl<F, T, E> ItemsPool<F, T, E>
|
|
|
|
where
|
|
|
|
F: Fn() -> Result<T, E>,
|
|
|
|
{
|
|
|
|
/// Create a new unbounded items pool with the specified function
|
|
|
|
/// to generate items when needed.
|
|
|
|
///
|
|
|
|
/// The `init` function will be invoked whenever a call to `with` requires new items.
|
|
|
|
pub fn new(init: F) -> Self {
|
|
|
|
let (sender, receiver) = crossbeam_channel::unbounded();
|
|
|
|
ItemsPool { init, sender, receiver }
|
|
|
|
}
|
|
|
|
|
|
|
|
/// Consumes the pool to retrieve all remaining items.
|
|
|
|
///
|
|
|
|
/// This method is useful for cleaning up and managing the items once they are no longer needed.
|
|
|
|
pub fn into_items(self) -> crossbeam_channel::IntoIter<T> {
|
|
|
|
self.receiver.into_iter()
|
|
|
|
}
|
|
|
|
|
|
|
|
/// Allows running a function on an item from the pool,
|
|
|
|
/// potentially generating a new item if the pool is empty.
|
|
|
|
pub fn with<G, R>(&self, f: G) -> Result<R, E>
|
|
|
|
where
|
|
|
|
G: FnOnce(&mut T) -> Result<R, E>,
|
|
|
|
{
|
|
|
|
let mut item = match self.receiver.try_recv() {
|
|
|
|
Ok(item) => item,
|
|
|
|
Err(TryRecvError::Empty) => (self.init)()?,
|
|
|
|
Err(TryRecvError::Disconnected) => unreachable!(),
|
|
|
|
};
|
|
|
|
|
|
|
|
// Run the user's closure with the retrieved item
|
|
|
|
let result = f(&mut item);
|
|
|
|
|
|
|
|
if let Err(e) = self.sender.send(item) {
|
|
|
|
unreachable!("error when sending into channel {e}");
|
|
|
|
}
|
|
|
|
|
|
|
|
result
|
|
|
|
}
|
|
|
|
}
|