pub mod error;
mod message;
mod worker;
#[cfg(feature = "crossbeam")]
pub use crossbeam_channel;
#[cfg(feature = "crossbeam")]
use crossbeam_channel::{unbounded, Sender};
#[cfg(feature = "mpsc")]
use std::sync::mpsc::{channel, Sender};
#[cfg(feature = "mpsc")]
use std::sync::{Arc, Mutex};
use error::{FailedToSendJob, FailedToSpawnThread};
use message::Message;
use worker::Worker;
type Job = Box<dyn FnOnce() + Send + 'static>;
#[derive(Debug)]
pub struct ThreadPool {
sender: Sender<Message>,
workers: Vec<Worker>,
}
impl ThreadPool {
#[cfg(feature = "crossbeam")]
pub fn new(worker: usize) -> Result<ThreadPool, FailedToSpawnThread> {
let workers = Vec::with_capacity(worker);
let (sender, receiver) = unbounded();
let mut threadpool = ThreadPool { workers, sender };
for _ in 0..worker {
let thread_builder = std::thread::Builder::new();
let worker = Worker::new(receiver.clone(), thread_builder)
.or_else(|_| Err(FailedToSpawnThread))?;
threadpool.workers.push(worker);
}
Ok(threadpool)
}
#[cfg(feature = "mpsc")]
pub fn new(worker: usize) -> Result<ThreadPool, FailedToSpawnThread> {
let workers = Vec::with_capacity(worker);
let (sender, receiver) = channel();
let receiver = Arc::new(Mutex::new(receiver));
let mut threadpool = ThreadPool { sender, workers };
for _ in 0..worker {
let thread_builder = std::thread::Builder::new();
let worker = Worker::new(Arc::clone(&receiver), thread_builder)
.or_else(|_| Err(FailedToSpawnThread))?;
threadpool.workers.push(worker);
}
Ok(threadpool)
}
pub fn execute<F>(&self, job: F) -> Result<(), FailedToSendJob>
where
F: FnOnce() + Send + 'static,
{
self.sender
.send(Message::NewJob(Box::new(job)))
.or_else(|_| Err(FailedToSendJob))?;
Ok(())
}
}
impl Drop for ThreadPool {
fn drop(&mut self) {
for _ in &self.workers {
self.sender.send(Message::Terminate).unwrap();
}
for worker in &mut self.workers {
if let Some(thread) = worker.take_thread() {
thread.join().unwrap();
}
}
}
}