Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Graceful Shutdown and Cleanup

โค้ดในโค้ดตัวอย่างที่ 21-20 กำลังตอบกลับคำร้องขอแบบไม่พร้อมกัน (asynchronously) ผ่านการใช้พูลของเธรดตามที่เราตั้งใจ เราได้รับการแจ้งเตือนบางอย่างเกี่ยวกับฟิลด์ workers, id, และ thread ที่เราไม่ได้ใช้โดยตรง ซึ่งเตือนเราว่าเราไม่ได้ล้างข้อมูล (clean up) อะไรเลย เมื่อเราใช้วิธี ctrl-C ที่สวยงามน้อยกว่าในการหยุดเธรดหลัก เธรดอื่นทั้งหมดก็ถูกหยุดลงทันทีเช่นกัน แม้ว่าพวกมันกำลังอยู่ในระหว่างการให้บริการคำร้องขออยู่ก็ตาม

ถัดไป เราจะอิมพลีเมนต์ Drop trait เพื่อเรียก join บนแต่ละเธรดในพูล เพื่อให้พวกมันสามารถทำคำร้องขอที่กำลังทำอยู่ให้เสร็จสิ้นก่อนที่จะปิดลง จากนั้น เราจะอิมพลีเมนต์วิธีที่จะบอกเธรดว่าพวกมันควรหยุดรับคำร้องขอใหม่และปิดตัวลง เพื่อให้เห็นโค้ดนี้ทำงานจริง เราจะปรับเปลี่ยนเซิร์ฟเวอร์ของเราให้รับคำร้องขอเพียงสองรายการก่อนที่จะปิดตัวพูลของเธรดอย่างนุ่มนวล (graceful shutdown)

สิ่งหนึ่งที่ต้องสังเกตขณะเราดำเนินการ: ไม่มีสิ่งใดในนี้ส่งผลกระทบต่อส่วนของโค้ดที่จัดการการรันโคลเชอร์ ดังนั้นทุกอย่างในที่นี้จะเหมือนเดิมแม้ว่าเราจะใช้พูลของเธรดสำหรับ async runtime ก็ตาม

Implementing the Drop Trait on ThreadPool

เรามาเริ่มต้นด้วยการอิมพลีเมนต์ Drop บนพูลของเธรดของเรา เมื่อพูลถูกดรอป เธรดของเราทั้งหมดควรเรียก join เพื่อให้แน่ใจว่าพวกมันทำงานเสร็จสิ้น โค้ดตัวอย่างที่ 21-22 แสดงความพยายามครั้งแรกในการอิมพลีเมนต์ Drop; โค้ดนี้จะยังไม่คอมไพล์ผ่าน

use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: mpsc::Sender<Job>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool { workers, sender }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        for worker in &mut self.workers {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let job = receiver.lock().unwrap().recv().unwrap();

                println!("Worker {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}

แรกเริ่ม เราวนลูปผ่านแต่ละ workers ของพูลของเธรด เราใช้ &mut สำหรับสิ่งนี้เนื่องจาก self เป็นการอ้างอิงที่เปลี่ยนแปลงค่าได้ (mutable reference) และเรายังจำเป็นต้องสามารถเปลี่ยนแปลงค่า worker ได้ด้วย สำหรับแต่ละ worker เราพิมพ์ข้อความระบุว่าอินสแตนซ์ Worker เฉพาะนี้กำลังปิดตัวลง จากนั้นเราเรียก join บนเธรดของอินสแตนซ์ Worker นั้น หากการเรียก join ล้มเหลว เราใช้ unwrap เพื่อทำให้ Rust เกิด panic และขยับไปสู่การปิดตัวที่ไม่นุ่มนวล

นี่คือข้อผิดพลาดที่เราได้รับเมื่อคอมไพล์โค้ดนี้:

$ cargo check
    Checking hello v0.1.0 (file:///projects/hello)
error[E0507]: cannot move out of `worker.thread` which is behind a mutable reference
  --> src/lib.rs:52:13
   |
52 |             worker.thread.join().unwrap();
   |             ^^^^^^^^^^^^^ ------ `worker.thread` moved due to this method call
   |             |
   |             move occurs because `worker.thread` has type `JoinHandle<()>`, which does not implement the `Copy` trait
   |
note: `JoinHandle::<T>::join` takes ownership of the receiver `self`, which moves `worker.thread`
  --> /rustc/ac68faa20c58cbccd01ee7208bf3b6e93a7d7f96/library/std/src/thread/join_handle.rs:149:16

For more information about this error, try `rustc --explain E0507`.
error: could not compile `hello` (lib) due to 1 previous error

ข้อผิดพลาดบอกเราว่าเราไม่สามารถเรียก join ได้เนื่องจากเรามีเพียงยืมแบบเปลี่ยนแปลงค่าได้ (mutable borrow) ของแต่ละ worker และ join รับความเป็นเจ้าของอาร์กิวเมนต์ของมัน เพื่อแก้ปัญหานี้ เราจำเป็นต้องย้ายเธรดออกจากอินสแตนซ์ Worker ที่เป็นเจ้าของ thread เพื่อให้ join สามารถบริโภคเธรดได้ วิธีหนึ่งในการทำเช่นนี้คือการใช้แนวทางเดียวกับที่เราใช้ในโค้ดตัวอย่างที่ 18-15 หาก Worker ถือครอง Option<thread::JoinHandle<()>> เราจะสามารถเรียกเมธอด take บน Option เพื่อย้ายค่าออกจากตัวแปรผัน Some และเหลือตัวแปรผัน None ไว้แทนที่ กล่าวอีกนัยหนึ่ง Worker ที่กำลังรันจะมีตัวแปรผัน Some ใน thread และเมื่อเราต้องการล้างข้อมูล Worker เราก็จะแทนที่ Some ด้วย None เพื่อให้ Worker ไม่มีเธรดให้รัน

อย่างไรก็ตาม เวลา_เดียว_ที่จะเกิดสิ่งนี้ขึ้นคือเมื่อทำการดรอป Worker ในทางกลับกัน เราจะต้องจัดการกับ Option<thread::JoinHandle<()>> ทุกๆ ที่ที่เราเข้าถึง worker.thread Rust แบบสำนวนนิยม (Idiomatic Rust) ใช้ Option ค่อนข้างมาก แต่เมื่อคุณพบว่าตัวเองหุ้มสิ่งที่คุณรู้ว่าจะมีอยู่เสมอไว้ใน Option เพื่อเป็นทางแก้ปัญหาชั่วคราวเช่นนี้ จึงเป็นความคิดที่ดีที่จะมองหาแนวทางอื่นเพื่อให้โค้ดของคุณสะอาดขึ้นและเกิดข้อผิดพลาดได้น้อยลง

ในกรณีนี้ มีทางเลือกที่ดีกว่า: เมธอด Vec::drain มันรับพารามิเตอร์ช่วง (range parameter) เพื่อระบุไอเทมที่จะลบออกจากเวกเตอร์ และคืนค่าตัววนซ้ำของไอเทมเหล่านั้น การส่งไวยากรณ์ช่วง .. จะลบ ทุก ค่าออกจากเวกเตอร์

ดังนั้น เราจำเป็นต้องอัปเดตการอิมพลีเมนต์ drop ของ ThreadPool ดังนี้:

#![allow(unused)]
fn main() {
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: mpsc::Sender<Job>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool { workers, sender }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let job = receiver.lock().unwrap().recv().unwrap();

                println!("Worker {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}
}

สิ่งนี้จะแก้ข้อผิดพลาดคอมไพเลอร์และไม่ต้องการการเปลี่ยนแปลงอื่นใดในโค้ดของเรา โปรดทราบว่า เนื่องจาก drop สามารถถูกเรียกใช้เมื่อเกิด panic ได้ unwrap จึงอาจเกิด panic และทำให้เกิด double panic ซึ่งจะทำให้โปรแกรมค้าง/แครชทันทีและยุติการล้างข้อมูลที่กำลังดำเนินการอยู่ สิ่งนี้ใช้ได้สำหรับโปรแกรมตัวอย่าง แต่ไม่แนะนำสำหรับโค้ดสำหรับการทำงานจริง

Signaling to the Threads to Stop Listening for Jobs

ด้วยการเปลี่ยนแปลงทั้งหมดที่เราได้ทำ โค้ดของเราคอมไพล์ผ่านโดยไม่มีคำเตือนใดๆ อย่างไรก็ตาม ข่าวร้ายคือโค้ดนี้ยังไม่ได้ทำงานในแบบที่เราต้องการ สิ่งสำคัญคือตรรกะในโคลเชอร์ที่รันโดยเธรดของอินสแตนซ์ Worker: ในขณะนี้ เราเรียก join แต่นั่นจะไม่ปิดเธรดลง เนื่องจากพวกมันวน loop ตลอดไปเพื่อมองหางาน หากเราลองดรอป ThreadPool ด้วยการอิมพลีเมนต์ drop ปัจจุบันของเรา เธรดหลักจะบล็อกตลอดไป โดยรอให้เธรดแรกทำงานเสร็จ

เพื่อแก้ไขปัญหานี้ เราจำเป็นต้องมีการเปลี่ยนแปลงในการอิมพลีเมนต์ drop ของ ThreadPool จากนั้นจึงเปลี่ยนลูปใน Worker

ประการแรก เราจะเปลี่ยนการอิมพลีเมนต์ drop ของ ThreadPool ให้ดรอป sender อย่างชัดเจนก่อนที่จะรอให้เธรดทำงานเสร็จ โค้ดตัวอย่างที่ 21-23 แสดงการเปลี่ยนแปลงกับ ThreadPool เพื่อดรอป sender อย่างชัดเจน ไม่เหมือนกับเธรด ในที่นี้เรา จำเป็นต้อง ใช้ Option เพื่อให้สามารถย้าย sender ออกจาก ThreadPool ด้วย Option::take ได้

use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}
// --snip--

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        // --snip--

        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool {
            workers,
            sender: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let job = receiver.lock().unwrap().recv().unwrap();

                println!("Worker {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}

การดรอป sender จะปิดช่องทางสื่อสาร ซึ่งบ่งบอกว่าจะไม่มีการส่งข้อความเพิ่มเติมอีก เมื่อสิ่งนั้นเกิดขึ้น การเรียก recv ทั้งหมดที่อินสแตนซ์ Worker ทำในลูปอนันต์จะคืนค่าเป็นข้อผิดพลาด ในโค้ดตัวอย่างที่ 21-24 เราเปลี่ยนลูปใน Worker ให้ออกจากลูปอย่างนุ่มนวลในกรณีนั้น ซึ่งหมายความว่าเธรดจะเสร็จสิ้นเมื่อการอิมพลีเมนต์ drop ของ ThreadPool เรียก join บนพวกมัน

use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool {
            workers,
            sender: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let message = receiver.lock().unwrap().recv();

                match message {
                    Ok(job) => {
                        println!("Worker {id} got a job; executing.");

                        job();
                    }
                    Err(_) => {
                        println!("Worker {id} disconnected; shutting down.");
                        break;
                    }
                }
            }
        });

        Worker { id, thread }
    }
}

เพื่อให้เห็นโค้ดนี้ทำงานจริง เรามาร่วมกันปรับเปลี่ยน main ให้รับคำร้องขอเพียงสองรายการก่อนที่จะปิดเซิร์ฟเวอร์อย่างนุ่มนวล ดังแสดงในโค้ดตัวอย่างที่ 21-25

use hello::ThreadPool;
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};

fn main() {
    let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
    let pool = ThreadPool::new(4);

    for stream in listener.incoming().take(2) {
        let stream = stream.unwrap();

        pool.execute(|| {
            handle_connection(stream);
        });
    }

    println!("Shutting down.");
}

fn handle_connection(mut stream: TcpStream) {
    let buf_reader = BufReader::new(&stream);
    let request_line = buf_reader.lines().next().unwrap().unwrap();

    let (status_line, filename) = match &request_line[..] {
        "GET / HTTP/1.1" => ("HTTP/1.1 200 OK", "hello.html"),
        "GET /sleep HTTP/1.1" => {
            thread::sleep(Duration::from_secs(5));
            ("HTTP/1.1 200 OK", "hello.html")
        }
        _ => ("HTTP/1.1 404 NOT FOUND", "404.html"),
    };

    let contents = fs::read_to_string(filename).unwrap();
    let length = contents.len();

    let response =
        format!("{status_line}\r\nContent-Length: {length}\r\n\r\n{contents}");

    stream.write_all(response.as_bytes()).unwrap();
}

คุณคงไม่อยากให้เว็บเซิร์ฟเวอร์ในโลกจริงปิดตัวลงหลังจากให้บริการคำร้องขอเพียงสองรายการ โค้ดนี้เพียงสาธิตให้เห็นว่าการปิดตัวลงและการล้างข้อมูลอย่างนุ่มนวลนั้นทำงานได้เรียบร้อยแล้ว

เมธอด take ถูกนิยามไว้ใน Iterator trait และจำกัดการวนซ้ำไว้ที่สองไอเทมแรกเป็นอย่างมาก ThreadPool จะหลุดออกจากขอบเขตเมื่อสิ้นสุด main และการอิมพลีเมนต์ drop จะรัน

เริ่มเซิร์ฟเวอร์ด้วย cargo run แล้วทำการร้องขอสามรายการ คำร้องขอที่สามควรเกิดข้อผิดพลาด และในเทอร์มินัลของคุณ คุณควรเห็นเอาต์พุตคล้ายดังนี้:

$ cargo run
   Compiling hello v0.1.0 (file:///projects/hello)
    Finished `dev` profile [unoptimized + debuginfo] target(s) in 0.41s
      Running `target/debug/hello`
Worker 0 got a job; executing.
Shutting down.
Shutting down worker 0
Worker 3 got a job; executing.
Worker 1 disconnected; shutting down.
Worker 2 disconnected; shutting down.
Worker 3 disconnected; shutting down.
Worker 0 disconnected; shutting down.
Shutting down worker 1
Shutting down worker 2
Shutting down worker 3

คุณอาจเห็นลำดับของ IDs ของ Worker และข้อความที่พิมพ์ออกมาแตกต่างกัน เราสามารถดูวิธีที่โค้ดนี้ทำงานได้จากข้อความ: อินสแตนซ์ Worker 0 และ 3 ได้รับคำร้องขอสองรายการแรก เซิร์ฟเวอร์หยุดรับการเชื่อมต่อหลังจากเรารับการเชื่อมต่อที่สอง และการอิมพลีเมนต์ Drop บน ThreadPool เริ่มทำงานก่อนที่ Worker 3 จะเริ่มงานของมันเสียอีก การดรอป sender จะตัดการเชื่อมต่ออินสแตนซ์ Worker ทั้งหมดและบอกให้พวกมันปิดตัวลง อินสแตนซ์ Worker แต่ละอันจะพิมพ์ข้อความเมื่อพวกมันตัดการเชื่อมต่อ จากนั้นพูลของเธรดจะเรียก join เพื่อรอให้เธรด Worker แต่ละอันทำงานเสร็จสิ้น

สังเกตแง่มุมที่น่าสนใจประการหนึ่งของการทำงานเฉพาะนี้: ThreadPool ดรอป sender และก่อนที่ Worker ใดๆ จะได้รับข้อผิดพลาด เราลองเรียก join บน Worker 0 โดย Worker 0 ยังไม่ได้รับข้อผิดพลาดจาก recv ดังนั้นเธรดหลักจึงบล็อก โดยรอให้ Worker 0 ทำงานเสร็จสิ้น ในระหว่างนั้น Worker 3 ได้รับงาน จากนั้นเธรดทั้งหมดก็ได้รับข้อผิดพลาด เมื่อ Worker 0 ทำงานเสร็จ เธรดหลักจึงรอให้อินสแตนซ์ Worker ที่เหลือทำงานเสร็จสิ้น ณ จุดนั้น พวกมันทั้งหมดได้ออกจากลูปและหยุดทำงานลงแล้ว

ยินดีด้วย! ตอนนี้เราทำโปรเจกต์ของเราเสร็จสมบูรณ์แล้ว; เรามีเว็บเซิร์ฟเวอร์พื้นฐานที่ใช้พูลของเธรดในการตอบกลับแบบไม่พร้อมกัน เราสามารถทำการปิดตัวเซิร์ฟเวอร์อย่างนุ่มนวล ซึ่งจะล้างข้อมูลเธรดทั้งหมดในพูล

นี่คือโค้ดฉบับเต็มสำหรับการอ้างอิง:

use hello::ThreadPool;
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};

fn main() {
    let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
    let pool = ThreadPool::new(4);

    for stream in listener.incoming().take(2) {
        let stream = stream.unwrap();

        pool.execute(|| {
            handle_connection(stream);
        });
    }

    println!("Shutting down.");
}

fn handle_connection(mut stream: TcpStream) {
    let buf_reader = BufReader::new(&stream);
    let request_line = buf_reader.lines().next().unwrap().unwrap();

    let (status_line, filename) = match &request_line[..] {
        "GET / HTTP/1.1" => ("HTTP/1.1 200 OK", "hello.html"),
        "GET /sleep HTTP/1.1" => {
            thread::sleep(Duration::from_secs(5));
            ("HTTP/1.1 200 OK", "hello.html")
        }
        _ => ("HTTP/1.1 404 NOT FOUND", "404.html"),
    };

    let contents = fs::read_to_string(filename).unwrap();
    let length = contents.len();

    let response =
        format!("{status_line}\r\nContent-Length: {length}\r\n\r\n{contents}");

    stream.write_all(response.as_bytes()).unwrap();
}
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool {
            workers,
            sender: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in &mut self.workers {
            println!("Shutting down worker {}", worker.id);

            if let Some(thread) = worker.thread.take() {
                thread.join().unwrap();
            }
        }
    }
}

struct Worker {
    id: usize,
    thread: Option<thread::JoinHandle<()>>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let message = receiver.lock().unwrap().recv();

                match message {
                    Ok(job) => {
                        println!("Worker {id} got a job; executing.");

                        job();
                    }
                    Err(_) => {
                        println!("Worker {id} disconnected; shutting down.");
                        break;
                    }
                }
            }
        });

        Worker {
            id,
            thread: Some(thread),
        }
    }
}

เราสามารถทำอะไรเพิ่มเติมที่นี่ได้อีก! หากคุณต้องการเพิ่มประสิทธิภาพโปรเจกต์นี้ต่อไป นี่คือไอเดียบางประการ:

  • เพิ่มเอกสารประกอบให้กับ ThreadPool และเมธอดสาธารณะของมัน
  • เพิ่มการทดสอบ (tests) สำหรับฟังก์ชันการทำงานของไลบรารี
  • เปลี่ยนการเรียกใช้ unwrap ให้เป็นการจัดการข้อผิดพลาดที่แข็งแกร่งยิ่งขึ้น
  • ใช้ ThreadPool เพื่อทำภารกิจอื่นนอกเหนือจากการให้บริการคำร้องขอของเว็บ
  • ค้นหาเครตพูลของเธรดบน crates.io แล้วอิมพลีเมนต์เว็บเซิร์ฟเวอร์ที่คล้ายกันโดยใช้เครตนั้นแทน จากนั้นเปรียบเทียบ API และความแข็งแกร่งของมันกับพูลของเธรดที่เราอิมพลีเมนต์

Summary

ทำได้ดีมาก! คุณทำจนถึงตอนจบของหนังสือเล่มนี้แล้ว! เราขอขอบคุณที่คุณมาร่วมทัวร์ Rust กับเรา ตอนนี้คุณพร้อมแล้วที่จะอิมพลีเมนต์โปรเจกต์ Rust ของคุณเอง และช่วยเหลือโปรเจกต์ของผู้อื่น โปรดจำไว้ว่ายังมีชุมชน Rustaceans ที่เปิดรับต้อนรับซึ่งยินดีที่จะช่วยเหลือคุณเกี่ยวกับความท้าทายต่างๆ ที่คุณพบเจอในการเดินทางของ Rust