feat: added cleanscript
This commit is contained in:
710
ch21/ch21-03-graceful-shutdown-and-cleanup.html
Normal file
710
ch21/ch21-03-graceful-shutdown-and-cleanup.html
Normal file
@@ -0,0 +1,710 @@
|
||||
<!DOCTYPE html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="UTF-8">
|
||||
<title>Graceful Shutdown and Cleanup</title>
|
||||
</head>
|
||||
<body>
|
||||
<h2 id="graceful-shutdown-and-cleanup"><a class="header" href="#graceful-shutdown-and-cleanup">Graceful Shutdown and Cleanup</a></h2>
|
||||
<p>The code in Listing 21-20 is responding to requests asynchronously through the
|
||||
use of a thread pool, as we intended. We get some warnings about the <code>workers</code>,
|
||||
<code>id</code>, and <code>thread</code> fields that we’re not using in a direct way that reminds us
|
||||
we’re not cleaning up anything. When we use the less elegant
|
||||
<kbd>ctrl</kbd>-<kbd>C</kbd> method to halt the main thread, all other threads
|
||||
are stopped immediately as well, even if they’re in the middle of serving a
|
||||
request.</p>
|
||||
<p>Next, then, we’ll implement the <code>Drop</code> trait to call <code>join</code> on each of the
|
||||
threads in the pool so that they can finish the requests they’re working on
|
||||
before closing. Then, we’ll implement a way to tell the threads they should
|
||||
stop accepting new requests and shut down. To see this code in action, we’ll
|
||||
modify our server to accept only two requests before gracefully shutting down
|
||||
its thread pool.</p>
|
||||
<p>One thing to notice as we go: None of this affects the parts of the code that
|
||||
handle executing the closures, so everything here would be the same if we were
|
||||
using a thread pool for an async runtime.</p>
|
||||
<h3 id="implementing-the-drop-trait-on-threadpool"><a class="header" href="#implementing-the-drop-trait-on-threadpool">Implementing the <code>Drop</code> Trait on <code>ThreadPool</code></a></h3>
|
||||
<p>Let’s start with implementing <code>Drop</code> on our thread pool. When the pool is
|
||||
dropped, our threads should all join to make sure they finish their work.
|
||||
Listing 21-22 shows a first attempt at a <code>Drop</code> implementation; this code won’t
|
||||
quite work yet.</p>
|
||||
<figure class="listing" id="listing-21-22">
|
||||
<span class="file-name">Filename: src/lib.rs</span>
|
||||
<pre><code class="language-rust ignore does_not_compile"><span class="boring">use std::{
|
||||
</span><span class="boring"> sync::{Arc, Mutex, mpsc},
|
||||
</span><span class="boring"> thread,
|
||||
</span><span class="boring">};
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">pub struct ThreadPool {
|
||||
</span><span class="boring"> workers: Vec<Worker>,
|
||||
</span><span class="boring"> sender: mpsc::Sender<Job>,
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">type Job = Box<dyn FnOnce() + Send + 'static>;
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">impl ThreadPool {
|
||||
</span><span class="boring"> /// Create a new ThreadPool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The size is the number of threads in the pool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// # Panics
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The `new` function will panic if the size is zero.
|
||||
</span><span class="boring"> pub fn new(size: usize) -> ThreadPool {
|
||||
</span><span class="boring"> assert!(size > 0);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let (sender, receiver) = mpsc::channel();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let receiver = Arc::new(Mutex::new(receiver));
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let mut workers = Vec::with_capacity(size);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> for id in 0..size {
|
||||
</span><span class="boring"> workers.push(Worker::new(id, Arc::clone(&receiver)));
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> ThreadPool { workers, sender }
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> pub fn execute<F>(&self, f: F)
|
||||
</span><span class="boring"> where
|
||||
</span><span class="boring"> F: FnOnce() + Send + 'static,
|
||||
</span><span class="boring"> {
|
||||
</span><span class="boring"> let job = Box::new(f);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> self.sender.send(job).unwrap();
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span>impl Drop for ThreadPool {
|
||||
fn drop(&mut self) {
|
||||
for worker in &mut self.workers {
|
||||
println!("Shutting down worker {}", worker.id);
|
||||
|
||||
worker.thread.join().unwrap();
|
||||
}
|
||||
}
|
||||
}
|
||||
<span class="boring">
|
||||
</span><span class="boring">struct Worker {
|
||||
</span><span class="boring"> id: usize,
|
||||
</span><span class="boring"> thread: thread::JoinHandle<()>,
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">impl Worker {
|
||||
</span><span class="boring"> fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
|
||||
</span><span class="boring"> let thread = thread::spawn(move || {
|
||||
</span><span class="boring"> loop {
|
||||
</span><span class="boring"> let job = receiver.lock().unwrap().recv().unwrap();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> println!("Worker {id} got a job; executing.");
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> job();
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring"> });
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> Worker { id, thread }
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">}</span></code></pre>
|
||||
<figcaption><a href="#listing-21-22">Listing 21-22</a>: Joining each thread when the thread pool goes out of scope</figcaption>
|
||||
</figure>
|
||||
<p>First, we loop through each of the thread pool <code>workers</code>. We use <code>&mut</code> for this
|
||||
because <code>self</code> is a mutable reference, and we also need to be able to mutate
|
||||
<code>worker</code>. For each <code>worker</code>, we print a message saying that this particular
|
||||
<code>Worker</code> instance is shutting down, and then we call <code>join</code> on that <code>Worker</code>
|
||||
instance’s thread. If the call to <code>join</code> fails, we use <code>unwrap</code> to make Rust
|
||||
panic and go into an ungraceful shutdown.</p>
|
||||
<p>Here is the error we get when we compile this code:</p>
|
||||
<pre><code class="language-console">$ 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/1159e78c4747b02ef996e55082b704c09b970588/library/std/src/thread/mod.rs:1921:17
|
||||
|
||||
For more information about this error, try `rustc --explain E0507`.
|
||||
error: could not compile `hello` (lib) due to 1 previous error
|
||||
</code></pre>
|
||||
<p>The error tells us we can’t call <code>join</code> because we only have a mutable borrow
|
||||
of each <code>worker</code> and <code>join</code> takes ownership of its argument. To solve this
|
||||
issue, we need to move the thread out of the <code>Worker</code> instance that owns
|
||||
<code>thread</code> so that <code>join</code> can consume the thread. One way to do this is to take
|
||||
the same approach we took in Listing 18-15. If <code>Worker</code> held an
|
||||
<code>Option<thread::JoinHandle<()>></code>, we could call the <code>take</code> method on the
|
||||
<code>Option</code> to move the value out of the <code>Some</code> variant and leave a <code>None</code> variant
|
||||
in its place. In other words, a <code>Worker</code> that is running would have a <code>Some</code>
|
||||
variant in <code>thread</code>, and when we wanted to clean up a <code>Worker</code>, we’d replace
|
||||
<code>Some</code> with <code>None</code> so that the <code>Worker</code> wouldn’t have a thread to run.</p>
|
||||
<p>However, the <em>only</em> time this would come up would be when dropping the
|
||||
<code>Worker</code>. In exchange, we’d have to deal with an
|
||||
<code>Option<thread::JoinHandle<()>></code> anywhere we accessed <code>worker.thread</code>.
|
||||
Idiomatic Rust uses <code>Option</code> quite a bit, but when you find yourself wrapping
|
||||
something you know will always be present in an <code>Option</code> as a workaround like
|
||||
this, it’s a good idea to look for alternative approaches to make your code
|
||||
cleaner and less error-prone.</p>
|
||||
<p>In this case, a better alternative exists: the <code>Vec::drain</code> method. It accepts
|
||||
a range parameter to specify which items to remove from the vector and returns
|
||||
an iterator of those items. Passing the <code>..</code> range syntax will remove <em>every</em>
|
||||
value from the vector.</p>
|
||||
<p>So, we need to update the <code>ThreadPool</code> <code>drop</code> implementation like this:</p>
|
||||
<figure class="listing">
|
||||
<span class="file-name">Filename: src/lib.rs</span>
|
||||
<pre class="playground"><code class="language-rust edition2024"><span class="boring">#![allow(unused)]
|
||||
</span><span class="boring">fn main() {
|
||||
</span><span class="boring">use std::{
|
||||
</span><span class="boring"> sync::{Arc, Mutex, mpsc},
|
||||
</span><span class="boring"> thread,
|
||||
</span><span class="boring">};
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">pub struct ThreadPool {
|
||||
</span><span class="boring"> workers: Vec<Worker>,
|
||||
</span><span class="boring"> sender: mpsc::Sender<Job>,
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">type Job = Box<dyn FnOnce() + Send + 'static>;
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">impl ThreadPool {
|
||||
</span><span class="boring"> /// Create a new ThreadPool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The size is the number of threads in the pool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// # Panics
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The `new` function will panic if the size is zero.
|
||||
</span><span class="boring"> pub fn new(size: usize) -> ThreadPool {
|
||||
</span><span class="boring"> assert!(size > 0);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let (sender, receiver) = mpsc::channel();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let receiver = Arc::new(Mutex::new(receiver));
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let mut workers = Vec::with_capacity(size);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> for id in 0..size {
|
||||
</span><span class="boring"> workers.push(Worker::new(id, Arc::clone(&receiver)));
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> ThreadPool { workers, sender }
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> pub fn execute<F>(&self, f: F)
|
||||
</span><span class="boring"> where
|
||||
</span><span class="boring"> F: FnOnce() + Send + 'static,
|
||||
</span><span class="boring"> {
|
||||
</span><span class="boring"> let job = Box::new(f);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> self.sender.send(job).unwrap();
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span>impl Drop for ThreadPool {
|
||||
fn drop(&mut self) {
|
||||
for worker in self.workers.drain(..) {
|
||||
println!("Shutting down worker {}", worker.id);
|
||||
|
||||
worker.thread.join().unwrap();
|
||||
}
|
||||
}
|
||||
}
|
||||
<span class="boring">
|
||||
</span><span class="boring">struct Worker {
|
||||
</span><span class="boring"> id: usize,
|
||||
</span><span class="boring"> thread: thread::JoinHandle<()>,
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">impl Worker {
|
||||
</span><span class="boring"> fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
|
||||
</span><span class="boring"> let thread = thread::spawn(move || {
|
||||
</span><span class="boring"> loop {
|
||||
</span><span class="boring"> let job = receiver.lock().unwrap().recv().unwrap();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> println!("Worker {id} got a job; executing.");
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> job();
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring"> });
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> Worker { id, thread }
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">}</span></code></pre>
|
||||
</figure>
|
||||
<p>This resolves the compiler error and does not require any other changes to our
|
||||
code. Note that, because drop can be called when panicking, the unwrap
|
||||
could also panic and cause a double panic, which immediately crashes the
|
||||
program and ends any cleanup in progress. This is fine for an example program,
|
||||
but it isn’t recommended for production code.</p>
|
||||
<h3 id="signaling-to-the-threads-to-stop-listening-for-jobs"><a class="header" href="#signaling-to-the-threads-to-stop-listening-for-jobs">Signaling to the Threads to Stop Listening for Jobs</a></h3>
|
||||
<p>With all the changes we’ve made, our code compiles without any warnings.
|
||||
However, the bad news is that this code doesn’t function the way we want it to
|
||||
yet. The key is the logic in the closures run by the threads of the <code>Worker</code>
|
||||
instances: At the moment, we call <code>join</code>, but that won’t shut down the threads,
|
||||
because they <code>loop</code> forever looking for jobs. If we try to drop our
|
||||
<code>ThreadPool</code> with our current implementation of <code>drop</code>, the main thread will
|
||||
block forever, waiting for the first thread to finish.</p>
|
||||
<p>To fix this problem, we’ll need a change in the <code>ThreadPool</code> <code>drop</code>
|
||||
implementation and then a change in the <code>Worker</code> loop.</p>
|
||||
<p>First, we’ll change the <code>ThreadPool</code> <code>drop</code> implementation to explicitly drop
|
||||
the <code>sender</code> before waiting for the threads to finish. Listing 21-23 shows the
|
||||
changes to <code>ThreadPool</code> to explicitly drop <code>sender</code>. Unlike with the thread,
|
||||
here we <em>do</em> need to use an <code>Option</code> to be able to move <code>sender</code> out of
|
||||
<code>ThreadPool</code> with <code>Option::take</code>.</p>
|
||||
<figure class="listing" id="listing-21-23">
|
||||
<span class="file-name">Filename: src/lib.rs</span>
|
||||
<pre><code class="language-rust noplayground not_desired_behavior"><span class="boring">use std::{
|
||||
</span><span class="boring"> sync::{Arc, Mutex, mpsc},
|
||||
</span><span class="boring"> thread,
|
||||
</span><span class="boring">};
|
||||
</span><span class="boring">
|
||||
</span>pub struct ThreadPool {
|
||||
workers: Vec<Worker>,
|
||||
sender: Option<mpsc::Sender<Job>>,
|
||||
}
|
||||
// --snip--
|
||||
<span class="boring">
|
||||
</span><span class="boring">type Job = Box<dyn FnOnce() + Send + 'static>;
|
||||
</span><span class="boring">
|
||||
</span>impl ThreadPool {
|
||||
<span class="boring"> /// Create a new ThreadPool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The size is the number of threads in the pool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// # Panics
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The `new` function will panic if the size is zero.
|
||||
</span> pub fn new(size: usize) -> ThreadPool {
|
||||
// --snip--
|
||||
|
||||
<span class="boring"> assert!(size > 0);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let (sender, receiver) = mpsc::channel();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let receiver = Arc::new(Mutex::new(receiver));
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let mut workers = Vec::with_capacity(size);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> for id in 0..size {
|
||||
</span><span class="boring"> workers.push(Worker::new(id, Arc::clone(&receiver)));
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">
|
||||
</span> 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
<span class="boring">
|
||||
</span><span class="boring">struct Worker {
|
||||
</span><span class="boring"> id: usize,
|
||||
</span><span class="boring"> thread: thread::JoinHandle<()>,
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">impl Worker {
|
||||
</span><span class="boring"> fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
|
||||
</span><span class="boring"> let thread = thread::spawn(move || {
|
||||
</span><span class="boring"> loop {
|
||||
</span><span class="boring"> let job = receiver.lock().unwrap().recv().unwrap();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> println!("Worker {id} got a job; executing.");
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> job();
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring"> });
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> Worker { id, thread }
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">}</span></code></pre>
|
||||
<figcaption><a href="#listing-21-23">Listing 21-23</a>: Explicitly dropping <code>sender</code> before joining the <code>Worker</code> threads</figcaption>
|
||||
</figure>
|
||||
<p>Dropping <code>sender</code> closes the channel, which indicates no more messages will be
|
||||
sent. When that happens, all the calls to <code>recv</code> that the <code>Worker</code> instances do
|
||||
in the infinite loop will return an error. In Listing 21-24, we change the
|
||||
<code>Worker</code> loop to gracefully exit the loop in that case, which means the threads
|
||||
will finish when the <code>ThreadPool</code> <code>drop</code> implementation calls <code>join</code> on them.</p>
|
||||
<figure class="listing" id="listing-21-24">
|
||||
<span class="file-name">Filename: src/lib.rs</span>
|
||||
<pre><code class="language-rust noplayground"><span class="boring">use std::{
|
||||
</span><span class="boring"> sync::{Arc, Mutex, mpsc},
|
||||
</span><span class="boring"> thread,
|
||||
</span><span class="boring">};
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">pub struct ThreadPool {
|
||||
</span><span class="boring"> workers: Vec<Worker>,
|
||||
</span><span class="boring"> sender: Option<mpsc::Sender<Job>>,
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">type Job = Box<dyn FnOnce() + Send + 'static>;
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">impl ThreadPool {
|
||||
</span><span class="boring"> /// Create a new ThreadPool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The size is the number of threads in the pool.
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// # Panics
|
||||
</span><span class="boring"> ///
|
||||
</span><span class="boring"> /// The `new` function will panic if the size is zero.
|
||||
</span><span class="boring"> pub fn new(size: usize) -> ThreadPool {
|
||||
</span><span class="boring"> assert!(size > 0);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let (sender, receiver) = mpsc::channel();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let receiver = Arc::new(Mutex::new(receiver));
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let mut workers = Vec::with_capacity(size);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> for id in 0..size {
|
||||
</span><span class="boring"> workers.push(Worker::new(id, Arc::clone(&receiver)));
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> ThreadPool {
|
||||
</span><span class="boring"> workers,
|
||||
</span><span class="boring"> sender: Some(sender),
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> pub fn execute<F>(&self, f: F)
|
||||
</span><span class="boring"> where
|
||||
</span><span class="boring"> F: FnOnce() + Send + 'static,
|
||||
</span><span class="boring"> {
|
||||
</span><span class="boring"> let job = Box::new(f);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> self.sender.as_ref().unwrap().send(job).unwrap();
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">impl Drop for ThreadPool {
|
||||
</span><span class="boring"> fn drop(&mut self) {
|
||||
</span><span class="boring"> drop(self.sender.take());
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> for worker in self.workers.drain(..) {
|
||||
</span><span class="boring"> println!("Shutting down worker {}", worker.id);
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> worker.thread.join().unwrap();
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span><span class="boring">struct Worker {
|
||||
</span><span class="boring"> id: usize,
|
||||
</span><span class="boring"> thread: thread::JoinHandle<()>,
|
||||
</span><span class="boring">}
|
||||
</span><span class="boring">
|
||||
</span>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 }
|
||||
}
|
||||
}</code></pre>
|
||||
<figcaption><a href="#listing-21-24">Listing 21-24</a>: Explicitly breaking out of the loop when <code>recv</code> returns an error</figcaption>
|
||||
</figure>
|
||||
<p>To see this code in action, let’s modify <code>main</code> to accept only two requests
|
||||
before gracefully shutting down the server, as shown in Listing 21-25.</p>
|
||||
<figure class="listing" id="listing-21-25">
|
||||
<span class="file-name">Filename: src/main.rs</span>
|
||||
<pre><code class="language-rust ignore"><span class="boring">use hello::ThreadPool;
|
||||
</span><span class="boring">use std::{
|
||||
</span><span class="boring"> fs,
|
||||
</span><span class="boring"> io::{BufReader, prelude::*},
|
||||
</span><span class="boring"> net::{TcpListener, TcpStream},
|
||||
</span><span class="boring"> thread,
|
||||
</span><span class="boring"> time::Duration,
|
||||
</span><span class="boring">};
|
||||
</span><span class="boring">
|
||||
</span>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.");
|
||||
}
|
||||
<span class="boring">
|
||||
</span><span class="boring">fn handle_connection(mut stream: TcpStream) {
|
||||
</span><span class="boring"> let buf_reader = BufReader::new(&stream);
|
||||
</span><span class="boring"> let request_line = buf_reader.lines().next().unwrap().unwrap();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let (status_line, filename) = match &request_line[..] {
|
||||
</span><span class="boring"> "GET / HTTP/1.1" => ("HTTP/1.1 200 OK", "hello.html"),
|
||||
</span><span class="boring"> "GET /sleep HTTP/1.1" => {
|
||||
</span><span class="boring"> thread::sleep(Duration::from_secs(5));
|
||||
</span><span class="boring"> ("HTTP/1.1 200 OK", "hello.html")
|
||||
</span><span class="boring"> }
|
||||
</span><span class="boring"> _ => ("HTTP/1.1 404 NOT FOUND", "404.html"),
|
||||
</span><span class="boring"> };
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let contents = fs::read_to_string(filename).unwrap();
|
||||
</span><span class="boring"> let length = contents.len();
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> let response =
|
||||
</span><span class="boring"> format!("{status_line}\r\nContent-Length: {length}\r\n\r\n{contents}");
|
||||
</span><span class="boring">
|
||||
</span><span class="boring"> stream.write_all(response.as_bytes()).unwrap();
|
||||
</span><span class="boring">}</span></code></pre>
|
||||
<figcaption><a href="#listing-21-25">Listing 21-25</a>: Shutting down the server after serving two requests by exiting the loop</figcaption>
|
||||
</figure>
|
||||
<p>You wouldn’t want a real-world web server to shut down after serving only two
|
||||
requests. This code just demonstrates that the graceful shutdown and cleanup is
|
||||
in working order.</p>
|
||||
<p>The <code>take</code> method is defined in the <code>Iterator</code> trait and limits the iteration
|
||||
to the first two items at most. The <code>ThreadPool</code> will go out of scope at the
|
||||
end of <code>main</code>, and the <code>drop</code> implementation will run.</p>
|
||||
<p>Start the server with <code>cargo run</code> and make three requests. The third request
|
||||
should error, and in your terminal, you should see output similar to this:</p>
|
||||
<!-- manual-regeneration
|
||||
cd listings/ch21-web-server/listing-21-25
|
||||
cargo run
|
||||
curl http://127.0.0.1:7878
|
||||
curl http://127.0.0.1:7878
|
||||
curl http://127.0.0.1:7878
|
||||
third request will error because server will have shut down
|
||||
copy output below
|
||||
Can't automate because the output depends on making requests
|
||||
-->
|
||||
<pre><code class="language-console">$ 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
|
||||
</code></pre>
|
||||
<p>You might see a different ordering of <code>Worker</code> IDs and messages printed. We can
|
||||
see how this code works from the messages: <code>Worker</code> instances 0 and 3 got the
|
||||
first two requests. The server stopped accepting connections after the second
|
||||
connection, and the <code>Drop</code> implementation on <code>ThreadPool</code> starts executing
|
||||
before <code>Worker 3</code> even starts its job. Dropping the <code>sender</code> disconnects all the
|
||||
<code>Worker</code> instances and tells them to shut down. The <code>Worker</code> instances each
|
||||
print a message when they disconnect, and then the thread pool calls <code>join</code> to
|
||||
wait for each <code>Worker</code> thread to finish.</p>
|
||||
<p>Notice one interesting aspect of this particular execution: The <code>ThreadPool</code>
|
||||
dropped the <code>sender</code>, and before any <code>Worker</code> received an error, we tried to
|
||||
join <code>Worker 0</code>. <code>Worker 0</code> had not yet gotten an error from <code>recv</code>, so the main
|
||||
thread blocked, waiting for <code>Worker 0</code> to finish. In the meantime, <code>Worker 3</code>
|
||||
received a job and then all threads received an error. When <code>Worker 0</code> finished,
|
||||
the main thread waited for the rest of the <code>Worker</code> instances to finish. At that
|
||||
point, they had all exited their loops and stopped.</p>
|
||||
<p>Congrats! We’ve now completed our project; we have a basic web server that uses
|
||||
a thread pool to respond asynchronously. We’re able to perform a graceful
|
||||
shutdown of the server, which cleans up all the threads in the pool.</p>
|
||||
<p>Here’s the full code for reference:</p>
|
||||
<figure class="listing">
|
||||
<span class="file-name">Filename: src/main.rs</span>
|
||||
<pre><code class="language-rust ignore">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();
|
||||
}</code></pre>
|
||||
</figure>
|
||||
<figure class="listing">
|
||||
<span class="file-name">Filename: src/lib.rs</span>
|
||||
<pre><code class="language-rust noplayground">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),
|
||||
}
|
||||
}
|
||||
}</code></pre>
|
||||
</figure>
|
||||
<p>We could do more here! If you want to continue enhancing this project, here are
|
||||
some ideas:</p>
|
||||
<ul>
|
||||
<li>Add more documentation to <code>ThreadPool</code> and its public methods.</li>
|
||||
<li>Add tests of the library’s functionality.</li>
|
||||
<li>Change calls to <code>unwrap</code> to more robust error handling.</li>
|
||||
<li>Use <code>ThreadPool</code> to perform some task other than serving web requests.</li>
|
||||
<li>Find a thread pool crate on <a href="https://crates.io/">crates.io</a> and implement a
|
||||
similar web server using the crate instead. Then, compare its API and
|
||||
robustness to the thread pool we implemented.</li>
|
||||
</ul>
|
||||
<h2 id="summary"><a class="header" href="#summary">Summary</a></h2>
|
||||
<p>Well done! You’ve made it to the end of the book! We want to thank you for
|
||||
joining us on this tour of Rust. You’re now ready to implement your own Rust
|
||||
projects and help with other people’s projects. Keep in mind that there is a
|
||||
welcoming community of other Rustaceans who would love to help you with any
|
||||
challenges you encounter on your Rust journey.</p>
|
||||
</body>
|
||||
</html>
|
||||
Reference in New Issue
Block a user