use crossbeam_channel::{Receiver as CbReceiver, Sender as CbSender, unbounded};
use std::sync::mpsc::{Receiver, Sender, channel};
use crate::ecs::plugin::Plugin;
type Job = Box<dyn FnOnce() + Send + 'static>;
pub struct BackgroundTasks {
job_tx: CbSender<Job>,
}
impl BackgroundTasks {
pub fn new(worker_count: usize) -> Self {
let (job_tx, job_rx): (CbSender<Job>, CbReceiver<Job>) = unbounded();
for _ in 0..worker_count.max(1) {
let job_rx = job_rx.clone(); std::thread::spawn(move || {
while let Ok(job) = job_rx.recv() {
job();
}
});
}
Self { job_tx }
}
pub fn spawn<T: Send + 'static>(
&self,
work: impl FnOnce() -> T + Send + 'static,
) -> TaskHandle<T> {
let (result_tx, result_rx) = channel::<T>();
let job: Job = Box::new(move || {
let result = work();
let _ = result_tx.send(result); });
let _ = self.job_tx.send(job);
TaskHandle { rx: result_rx }
}
}
pub struct TaskHandle<T> {
rx: Receiver<T>,
}
impl<T> TaskHandle<T> {
pub fn try_recv(&mut self) -> Option<T> {
self.rx.try_recv().ok()
}
}
pub struct BackgroundTasksPlugin {
worker_count: usize,
}
impl BackgroundTasksPlugin {
pub fn new(worker_count: usize) -> Self {
Self { worker_count }
}
}
impl Plugin for BackgroundTasksPlugin {
fn build(&self, app: &mut crate::prelude::App) {
app.add_resource(BackgroundTasks::new(self.worker_count));
}
}