Arrêt propre et nettoyage
Le code de l’encart 21-20 répond aux requêtes de manière asynchrone grâce à l’utilisation du groupe de fils d’exécution, comme nous l’espérions. Nous avons quelques avertissements sur les champs operateurs, id et tache que nous n’utilisons pas directement et qui nous rappellent que nous ne nettoyons rien. Lorsque nous arrêtons brutalement la tâche principale en appuyant sur ctrl-C, tous les autres fils d’exécution sont également immédiatement stoppées, même si ils sont en train de servir une requête.
Après cela, nous allons donc implémenter le trait Drop afin d’appeler join sur chacun des fils d’exécution du groupe afin qu’ils puissent finir les requêtes qu’ils sont en train de traiter avant de s’arrêter. Ensuite, nous allons implémenter un moyen de demander aux fils d’exécution d’arrêter d’accepter de nouvelles requêtes et de s’arrêter. Pour voir ce code en action, nous allons modifier notre serveur pour n’accepter que deux requêtes avant d’arrêter proprement son groupe de fils d’exécution.
Une remarque au passage : rien de ceci n’affecte les portions du code qui se charge d’exécuter les fermetures, donc tout serait ici identique si nous utilisions un groupe de fils pour un moteur d’exécution asynchrone.
Implémentation du trait Drop sur GroupeTaches
Commençons par implémenter Drop sur notre groupe de tâches. Lorsque le groupe est nettoyé, nos tâches doivent toutes faire appel à join pour s’assurer qu’elles finissent leur travail. L’encart 21-22 montre une première tentative d’implémentation de Drop ; ce code ne fonctionne pas encore tout à fait.
use std::{
sync::{Arc, Mutex, mpsc},
thread,
};
pub struct GroupeTaches {
operateurs: Vec<Operateur>,
envoi: mpsc::Sender<Mission>,
}
type Mission = Box<dyn FnOnce() + Send + 'static>;
impl GroupeTaches {
/// Crée un nouveau GroupeTaches.
///
/// La taille est le nombre de tâches présentes dans le groupe.
///
/// # Panics
///
/// La fonction `new` devrait paniquer si la taille vaut zéro.
pub fn new(taille: usize) -> GroupeTaches {
assert!(taille > 0);
let (envoi, reception) = mpsc::channel();
let reception = Arc::new(Mutex::new(reception));
let mut operateurs = Vec::with_capacity(taille);
for id in 0..taille {
operateurs.push(Operateur::new(id, Arc::clone(&reception)));
}
GroupeTaches { operateurs, envoi }
}
pub fn executer<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
let mission = Box::new(f);
self.envoi.send(mission).unwrap();
}
}
impl Drop for GroupeTaches {
fn drop(&mut self) {
for operateur in &mut self.operateurs {
println!("Arrêt de l'opérateur {}", operateur.id);
operateur.tache.join().unwrap();
}
}
}
struct Operateur {
id: usize,
tache: thread::JoinHandle<()>,
}
impl Operateur {
fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
let tache = thread::spawn(move || {
loop {
let mission = reception.lock().unwrap().recv().unwrap();
println!("L'opérateur {id} a reçu une mission ; il l'exécute.");
mission();
}
});
Operateur { id, tache }
}
}
D’abord, nous faisons une boucle sur tous les operateurs du groupe de fils d’exécution. Pour ce faire, nous utilisons &mut car self n’est qu’une référence mutable du groupe de fils d’exécution, et nous aurons également besoin de pouvoir modifier operateur. Pour chaque operateur, nous affichons un message qui indique que cette instance de operateur est en train de s’arrêter, puis nous faisons appel à join sur le fil d’exécution de cette instance de operateur. Si l’appel à join échoue, nous utilisons unwrap pour faire paniquer Rust et ainsi procéder à un arrêt brutal.
Voici l’erreur que nous obtenons lorsque nous compilons ce code :
$ cargo check
Checking salutations v0.1.0 (file:///projects/salutations)
error[E0507]: cannot move out of `operateur.tache` which is behind a mutable reference
--> src/lib.rs:52:13
|
52 | operateur.tache.join().unwrap();
| ^^^^^^^^^^^^^^^ ------ `operateur.tache` moved due to this method call
| |
| move occurs because `operateur.tache` has type `JoinHandle<()>`, which does not implement the `Copy` trait
|
note: `JoinHandle::<T>::join` takes ownership of the receiver `self`, which moves `operateur.tache`
--> /rustc/2d8144b7880597b6e6d3dfd63a9a9efae3f533d3/library/std/src/thread/join_handle.rs:149:16
For more information about this error, try `rustc --explain E0507`.
error: could not compile `salutations` (lib) due to 1 previous error
L’erreur nous informe que nous ne pouvons pas faire appel à join car nous avons seulement fait un emprunt mutable pour chacun des operateur alors que join prend possession de son argument. Pour résoudre ce problème, nous devons sortir la tache de l’instance de Operateur qui la possède afin que join puisse la consommer. Une manière de faire ceci est d’avoir la même approche que dans l’encart 18-15. Si Operateur contenait un Option<thread::JoinHandle<()>>, nous pourrions utiliser la méthode take sur Option pour sortir la valeur de la variante Some et y mettre à la place une variante None. Autrement dit, un Operateur qui est en cours d’exécution aurait une variante Some dans tache, et lorsque nous souhaiterions nettoyer un Operateur, nous remplacerions Some par None afin que Operateur n’ait pas de tâche à exécuter.
Toutefois, la seule fois où ce cas se présenterait serait lors de la suppression de Operateur. En contrepartie, nous devrions gérer un Option<thread::JoinHandle<()>> à chaque endroit où nous accèderions à operateur.tache. La manière classique de coder en Rust utilise beaucoup Option, mais quand vous vous retrouvez à encapsuler dans une Option quelque chose dont vous savez pertinemment qu’il sera toujours présent, il est judicieux de rechercher d’autres approches afin de rendre votre code plus propre et moins sujet aux erreurs.
Dans notre cas, il existe une meilleure solution : la méthode Vec::drain. Elle accepte un paramètre de plage permettant de préciser quels éléments sont à enlever du vecteur et renvoie un itérateur sur ces éléments. En passant la syntaxe de plage .., chacune des valeurs du vecteur seront supprimées.
Il nous faut donc mettre à jour l’implémentation de drop de GroupeTaches de la manière suivante :
#![allow(unused)]
fn main() {
use std::{
sync::{Arc, Mutex, mpsc},
thread,
};
pub struct GroupeTaches {
operateurs: Vec<Operateur>,
envoi: mpsc::Sender<Mission>,
}
type Mission = Box<dyn FnOnce() + Send + 'static>;
impl GroupeTaches {
/// Crée un nouveau GroupeTaches.
///
/// La taille est le nombre de tâches présentes dans le groupe.
///
/// # Panics
///
/// La fonction `new` devrait paniquer si la taille vaut zéro.
pub fn new(taille: usize) -> GroupeTaches {
assert!(taille > 0);
let (envoi, reception) = mpsc::channel();
let reception = Arc::new(Mutex::new(reception));
let mut operateurs = Vec::with_capacity(taille);
for id in 0..taille {
operateurs.push(Operateur::new(id, Arc::clone(&reception)));
}
GroupeTaches { operateurs, envoi }
}
pub fn executer<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
let mission = Box::new(f);
self.envoi.send(mission).unwrap();
}
}
impl Drop for GroupeTaches {
fn drop(&mut self) {
for operateur in self.operateurs.drain(..) {
println!("Arrêt de l'opérateur {}", operateur.id);
operateur.tache.join().unwrap();
}
}
}
struct Operateur {
id: usize,
tache: thread::JoinHandle<()>,
}
impl Operateur {
fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
let tache = thread::spawn(move || {
loop {
let mission = reception.lock().unwrap().recv().unwrap();
println!("L'opérateur {id} a reçu une mission ; il l'exécute.");
mission();
}
});
Operateur { id, tache }
}
}
}
Cela résout l’erreur de compilation et ne nécessite pas d’autres changements dans notre code. Notez bien que, dans la mesure où drop peut être appelé lors d’une panique, unwrap pourrait paniquer à son tour et provoquer ainsi une double panique ; ceci ferait immédiatement planter le programme et mettrait fin à tout processus de nettoyage en cours. Cela est convenable pour un programme d’exemple, mais ce n’est pas recommandé pour du code de production.
Demander aux fils d’exécution d’arrêter d’attendre des missions
Avec tous ces changements, notre code se compile désormais sans aucun avertissement. Toutefois, la mauvaise nouvelle est que, pour l’instant, ce code ne fonctionne toujours pas comme nous le souhaitons. La cause se situe dans la logique des fermetures qui sont exécutées par les fils d’exécution des instances de Operateur : pour le moment, nous faisons appel à join, mais cela ne va pas arrêter les fils d’exécution car ils font une boucle infinie avec loop pour attendre des missions. Si nous essayons de nettoyer notre GroupeTaches avec l’implémentation actuelle de drop, la tâche principale va se bloquer pour toujours en attendant en vain que la première tâche se termine.
Pour résoudre ce problème, nous devons modifier l’implémentation de drop de GroupeTaches et ensuite modifier la boucle dans Operateur.
Tout d’abord, nous allons changer l’implémentation drop de GroupeTaches pour supprimer de manière explicite envoi avant d’attendre que tous les fils se soient terminés. L’encart 21-23 montre les changements à GroupeTaches pour supprimer envoi de manière explicite. Au contraire du fil d’exécution, nous devons ici utiliser une Option pour pouvoir déplacer envoi hors de GroupeTaches avec Option::take.
use std::{
sync::{Arc, Mutex, mpsc},
thread,
};
pub struct GroupeTaches {
operateurs: Vec<Operateur>,
envoi: Option<mpsc::Sender<Mission>>,
}
// -- partie masquée ici --
type Mission = Box<dyn FnOnce() + Send + 'static>;
impl GroupeTaches {
/// Crée un nouveau GroupeTaches.
///
/// La taille est le nombre de tâches présentes dans le groupe.
///
/// # Panics
///
/// La fonction `new` devrait paniquer si la taille vaut zéro.
pub fn new(taille: usize) -> GroupeTaches {
// -- partie masquée ici --
assert!(taille > 0);
let (envoi, reception) = mpsc::channel();
let reception = Arc::new(Mutex::new(reception));
let mut operateurs = Vec::with_capacity(taille);
for id in 0..taille {
operateurs.push(Operateur::new(id, Arc::clone(&reception)));
}
GroupeTaches {
operateurs,
envoi: Some(envoi),
}
}
pub fn executer<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
let mission = Box::new(f);
self.envoi.as_ref().unwrap().send(mission).unwrap();
}
}
impl Drop for GroupeTaches {
fn drop(&mut self) {
drop(self.envoi.take());
for operateur in self.operateurs.drain(..) {
println!("Arrêt de l'opérateur {}", operateur.id);
operateur.tache.join().unwrap();
}
}
}
struct Operateur {
id: usize,
tache: thread::JoinHandle<()>,
}
impl Operateur {
fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
let tache = thread::spawn(move || {
loop {
let mission = reception.lock().unwrap().recv().unwrap();
println!("L'opérateur {id} a reçu une mission ; il l'exécute.");
mission();
}
});
Operateur { id, tache }
}
}
sender before joining the Worker threadsLa suppression de envoi referme le canal, ce qui indique que plus aucun message ne sera envoyé. Quand ceci survient, tous les appels à recv que les instances d’Operateur font dans la boucle infinie vont renvoyer une erreur. Dans l’encart 21-24, nous modifions la boucle de Operateur afin de quitter la boucle proprement dans ce cas, ce qui implique que les fils d’exécution vont se terminer quand l’implémentation drop de GroupeTaches appelle join sur eux.
use std::{
sync::{Arc, Mutex, mpsc},
thread,
};
pub struct GroupeTaches {
operateurs: Vec<Operateur>,
envoi: Option<mpsc::Sender<Mission>>,
}
type Mission = Box<dyn FnOnce() + Send + 'static>;
impl GroupeTaches {
/// Crée un nouveau GroupeTaches.
///
/// La taille est le nombre de tâches présentes dans le groupe.
///
/// # Panics
///
/// La fonction `new` devrait paniquer si la taille vaut zéro.
pub fn new(taille: usize) -> GroupeTaches {
assert!(taille > 0);
let (envoi, reception) = mpsc::channel();
let reception = Arc::new(Mutex::new(reception));
let mut operateurs = Vec::with_capacity(taille);
for id in 0..taille {
operateurs.push(Operateur::new(id, Arc::clone(&reception)));
}
GroupeTaches {
operateurs,
envoi: Some(envoi),
}
}
pub fn executer<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
let mission = Box::new(f);
self.envoi.as_ref().unwrap().send(mission).unwrap();
}
}
impl Drop for GroupeTaches {
fn drop(&mut self) {
drop(self.envoi.take());
for operateur in self.operateurs.drain(..) {
println!("Arrêt de l'opérateur {}", operateur.id);
operateur.tache.join().unwrap();
}
}
}
struct Operateur {
id: usize,
tache: thread::JoinHandle<()>,
}
impl Operateur {
fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
let tache = thread::spawn(move || {
loop {
let message = reception.lock().unwrap().recv();
match message {
Ok(mission) => {
println!("L'opérateur {id} a reçu une mission ; il l'exécute.");
mission();
}
Err(_) => {
println!("L'opérateur {id} a reçu l'instruction d'arrêt.");
break;
}
}
}
});
Operateur { id, tache }
}
}
recv returns an errorPour observer ce code en action, modifions notre main pour accepter uniquement deux requêtes avant d’arrêter proprement le serveur, comme dans l’encart 21-25.
use salutations::GroupeTaches;
use std::{
fs,
io::{BufReader, prelude::*},
net::{TcpListener, TcpStream},
thread,
time::Duration,
};
fn main() {
let ecouteur = TcpListener::bind("127.0.0.1:7878").unwrap();
let groupe = GroupeTaches::new(4);
for flux in ecouteur.incoming().take(2) {
let flux = flux.unwrap();
groupe.executer(|| {
gestion_connexion(flux);
});
}
println!("Arrêt complet.");
}
fn gestion_connexion(mut flux: TcpStream) {
let tampon = BufReader::new(&flux);
let ligne_requete = tampon.lines().next().unwrap().unwrap();
let (ligne_statut, nom_fichier) = match &ligne_requete[..] {
"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 NON TROUVÉ", "404.html"),
};
let contenu = fs::read_to_string(nom_fichier).unwrap();
let longueur = contenu.len();
let reponse =
format!("{ligne_statut}\r\nContent-Length: {longueur}\r\n\r\n{contenu}");
flux.write_all(reponse.as_bytes()).unwrap();
}
Dans la réalité, on ne voudrait pas qu’un serveur web s’arrête après avoir servi seulement deux requêtes. Ce code sert uniquement à montrer que l’arrêt et le nettoyage s’effectuent bien proprement.
La méthode take est définie dans le trait Iterator et limite l’itération aux deux premiers éléments au maximum. Le GroupeTaches va sortir de la portée à la fin du main et l’implémentation de drop va s’exécuter.
Démarrez le serveur avec cargo run et faites trois requêtes. La troisième requête devrait renvoyer une erreur tandis que dans votre terminal vous devriez avoir une sortie similaire à ceci :
$ cargo run
Compiling salutations v0.1.0 (file:///projects/salutations)
Finished `dev` profile [unoptimized + debuginfo] target(s) in 0.41s
Running `target/debug/salutations`
L'opérateur 0 a reçu une mission ; il l'exécute.
Arrêt complet.
Arrêt de l'opérateur 0
L'opérateur 3 a reçu une mission ; il l'exécute.
L'opérateur 1 a reçu l'instruction d'arrêt.
L'opérateur 2 a reçu l'instruction d'arrêt.
L'opérateur 3 a reçu l'instruction d'arrêt.
L'opérateur 0 a reçu l'instruction d'arrêt.
Arrêt de l'opérateur 1
Arrêt de l'opérateur 2
Arrêt de l'opérateur 3
Vous pourriez avoir un ordre différent entre les identifiants d’Operateur et les messages affichés. Nous pouvons constater la façon dont ce code fonctionne grâce aux messages : les instances d’Operateur 0 et 3 ont reçu les deux premières requêtes. Le serveur a arrêté d’accepter des connexions après la deuxième connexion, et l’implémentation de Drop sur GroupeTaches commence à s’exécuter avant même qu’Operateur 3 ne commence son travail. La suppression de envoi entraîne la déconnexion de toutes les instances d’Operateur et leur demande de s’arrêter. Les instances d’Operateur affichent chacune un message lors de leur déconnexion, puis le groupe de fils d’exécution appelle join pour attendre la fin de chaque fil d’exécution Operateur.
Remarquez un aspect intéressant spécifique à cette exécution : le GroupeTaches a supprimé le envoi et, avant qu’aucun Operateur n’aie reçu d’erreur, nous avons essayé de joindre l’Operateur 0. Cet Operateur 0 n’avait alors pas encore reçu d’erreur de recv, donc la tâche principale s’est bloquée dans l’attente de la fin de Operateur 0. Pendant ce temps, Operateur 3 a reçu une mission, et ensuite tous les fils d’exécution ont reçu une erreur. Quand Operateur 0 a terminé, le fil d’exécution principal a attendu que le reste des instances d’Operateur se terminent. À ce stade, elles étaient toutes sorties de leurs boucles et s’étaient arrêtées.
Félicitations ! Nous avons maintenant terminé notre projet ; nous avons un serveur web basique qui utilise un groupe de tâches pour répondre de manière asynchrone. Nous pouvons demander un arrêt propre du serveur qui va nettoyer toutes les tâches du groupe.
Voici le code complet afin que vous puissiez vous y référer :
use salutations::GroupeTaches;
use std::{
fs,
io::{BufReader, prelude::*},
net::{TcpListener, TcpStream},
thread,
time::Duration,
};
fn main() {
let ecouteur = TcpListener::bind("127.0.0.1:7878").unwrap();
let groupe = GroupeTaches::new(4);
for flux in ecouteur.incoming().take(2) {
let flux = flux.unwrap();
groupe.executer(|| {
gestion_connexion(flux);
});
}
println!("Arrêt complet.");
}
fn gestion_connexion(mut flux: TcpStream) {
let tampon = BufReader::new(&flux);
let ligne_requete = tampon.lines().next().unwrap().unwrap();
let (ligne_statut, nom_fichier) = match &ligne_requete[..] {
"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 NON TROUVÉ", "404.html"),
};
let contenu = fs::read_to_string(nom_fichier).unwrap();
let longueur = contenu.len();
let reponse =
format!("{ligne_statut}\r\nContent-Length: {longueur}\r\n\r\n{contenu}");
flux.write_all(reponse.as_bytes()).unwrap();
}
use std::{
sync::{Arc, Mutex, mpsc},
thread,
};
pub struct GroupeTaches {
operateurs: Vec<Operateur>,
envoi: Option<mpsc::Sender<Mission>>,
}
type Mission = Box<dyn FnOnce() + Send + 'static>;
impl GroupeTaches {
/// Crée un nouveau GroupeTaches.
///
/// La taille est le nombre de tâches présentes dans le groupe.
///
/// # Panics
///
/// La fonction `new` devrait paniquer si la taille vaut zéro.
pub fn new(taille: usize) -> GroupeTaches {
assert!(taille > 0);
let (envoi, reception) = mpsc::channel();
let reception = Arc::new(Mutex::new(reception));
let mut operateurs = Vec::with_capacity(taille);
for id in 0..taille {
operateurs.push(Operateur::new(id, Arc::clone(&reception)));
}
GroupeTaches {
operateurs,
envoi: Some(envoi),
}
}
pub fn executer<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
let mission = Box::new(f);
self.envoi.as_ref().unwrap().send(mission).unwrap();
}
}
impl Drop for GroupeTaches {
fn drop(&mut self) {
drop(self.envoi.take());
for operateur in &mut self.operateurs {
println!("Arrêt de l'opérateur {}", operateur.id);
if let Some(tache) = operateur.tache.take() {
tache.join().unwrap();
}
}
}
}
struct Operateur {
id: usize,
tache: Option<thread::JoinHandle<()>>,
}
impl Operateur {
fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
let tache = thread::spawn(move || {
loop {
let message = reception.lock().unwrap().recv();
match message {
Ok(mission) => {
println!("L'opérateur {id} a reçu une mission ; il l'exécute.");
mission();
}
Err(_) => {
println!("L'opérateur {id} a reçu l'instruction d'arrêt.");
break;
}
}
}
});
Operateur {
id,
tache: Some(tache),
}
}
}
Nous aurions pu faire bien plus ! Si vous souhaitez continuer à améliorer ce projet, voici quelques idées :
- ajouter de la documentation à
GroupeTacheset aux méthodes publiques ; - ajouter des tests sur les fonctionnalités de la bibliothèque ;
- remplacer les appels à
unwrappour fournir une meilleure gestion des erreurs ; - utiliser
GroupeTachespour exécuter d’autres tâches que de répondre à des requêtes web ; - trouver une crate de groupe de fils d’exécution (NdT : thread pool) sur crates.io et implémenter un serveur web similaire en l’utilisant. Comparer ensuite son API et sa robustesse au groupe de fils d’exécution que nous avons implémenté.
Résumé
Bravo ! Vous êtes arrivé à la fin du livre ! Nous tenons à vous remercier chaleureusement de nous avoir accompagné pendant cette présentation de Rust. Vous êtes maintenant fin prêt(e) à créer vos propres projets Rust et aider les projets des autres développeurs. Rappelez-vous qu’il existe une communauté accueillante de Rustacés qui adorerait vous aider à relever tous les défis que vous rencontrerez dans votre aventure avec Rust.