let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
let pool = ThreadPool::new(4);
- for stream in listener.incoming() {
+ for stream in listener.incoming() /* for testing cleanup: .take(2) */ {
let stream = stream.unwrap();
pool.execute(|| handle_connection(stream));
}
type Job = Box<dyn FnOnce() + Send + 'static>;
+enum Message {
+ NewJob(Job),
+ Terminate,
+}
+
struct Worker {
id: usize,
thread: Option<thread::JoinHandle<()>>,
}
impl Worker {
- fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
+ fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Message>>>) -> Worker {
let thread = Some(thread::spawn(move || loop {
- let job = receiver.lock().unwrap().recv().unwrap();
- println!("Worker {} got a job, executing", id);
- job();
+ let message = receiver.lock().unwrap().recv().unwrap();
+
+ match message {
+ Message::NewJob(job) => {
+ println!("Worker {} got a job, executing", id);
+ job();
+ },
+
+ Message::Terminate => {
+ println!("Worker {} got terminated", id);
+ break;
+ }
+ }
}));
Worker { id, thread }
}
pub struct ThreadPool {
workers: Vec<Worker>,
- sender: mpsc::Sender<Job>,
+ sender: mpsc::Sender<Message>,
}
impl ThreadPool {
pub fn execute<F>(&self, f: F)
where F: FnOnce() + Send + 'static
{
- self.sender.send(Box::new(f)).unwrap();
+ self.sender.send(Message::NewJob(Box::new(f))).unwrap();
}
}
impl Drop for ThreadPool {
fn drop(&mut self) {
+ for _ in &self.workers {
+ self.sender.send(Message::Terminate).unwrap();
+ }
+
for worker in &mut self.workers {
println!("Shutting down worker {}", worker.id);