711 lines
33 KiB
HTML
711 lines
33 KiB
HTML
<!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>
|