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

D’un serveur monotâche à un serveur multitâches

Pour le moment, le serveur va traiter chaque requête l’une après l’autre, ce qui signifie qu’il ne traitera pas une deuxième connexion tant que la première connexion n’aura pas fini d’être traitée. Si le serveur reçoit encore plus de requêtes, cette exécution en série sera de moins en moins adaptée. Si le serveur reçoit une requête qui prend longtemps à traiter, les demandes suivantes devront attendre que la longue requête à traiter soit terminée, même si les nouvelles requêtes peuvent être traitées rapidement. Nous devons corriger cela, mais d’abord, observons le problème se produire pour de vrai.

Simulation d’une requête longue

Nous allons voir comment une requête longue à traiter peut affecter le traitement des autres requêtes avec l’implémentation actuelle de notre serveur. L’encart 21-10 rajoute le traitement d’une requête pour /pause qui va simuler une longue réponse qui va faire en sorte que le serveur soit en pause pendant cinq secondes avant de pouvoir répondre à nouveau.

Filename: src/main.rs
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};
// -- partie masquée ici --

fn main() {
    let ecouteur = TcpListener::bind("127.0.0.1:7878").unwrap();

    for flux in ecouteur.incoming() {
        let flux = flux.unwrap();

        gestion_connexion(flux);
    }
}

fn gestion_connexion(mut flux: TcpStream) {
    // -- partie masquée ici --

    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"),
    };

    // -- partie masquée ici --

    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();
}
Listing 21-10: Simulating a slow request by sleeping for five seconds

Nous sommes passés du if au match, maintenant que nous avons trois cas de figure. Nous devons explicitement faire une correspondance sur une slice de ligne_requete pour comparer les valeurs littérales de chaînes de caractères ; match ne fait pas de déréférencement automatique, contrairement à la méthode d’égalité.

La première branche est la même que le bloc if block de l’encart 21-9. La deuxième branche correspond à une requête vers /sleep. Lorsque cette requête est reçue, le serveur va se mettre en pause pendant 5 secondes avant de générer la page HTML de succès. La troisième branche est la même que le bloc else de l’encart 21-9.

Vous pouvez constater à quel point notre serveur est primitif : une bibliothèque digne de ce nom devrait gérer la détection de différents types de requêtes de manière bien moins verbeuse !

Démarrez le serveur en utilisant cargo run. Ouvrez ensuite deux fenêtres de navigateur web : une pour http://127.0.0.1:7878/ et l’autre pour http://127.0.0.1:7878/pause. Si vous demandez l’URI / plusieurs fois, comme vous l’avez fait précédemment, vous constaterez que le serveur répond rapidement. Mais lorsque vous saisirez /pause et que vous chargerez ensuite /, vous constaterez que / attend que pause ait fini sa pause de cinq secondes avant de se charger.

Nous pourrions utiliser plusieurs techniques permettant d’éviter d’accumuler des requêtes après une requête dont le traitement est long, dont l’utilisation de async comme dans le chapitre 17 ; celle que nous allons implémenter est un groupe de fils d’exécution.

Amélioration du débit avec un groupe de tâches

Un groupe de fils est un groupe constitué de fils d’exécution qui ont été créées au préalable, qui sont prêts et qui attendent de prendre en charge des missions. Lorsque le programme reçoit une nouvelle mission, il assigne un des fils d’exécution du groupe pour cette mission, et ce fil va traiter la mission. Les fils restants dans le groupe restent disponibles pour traiter d’autres missions qui peuvent arriver pendant que le premier fil est en cours de traitement. Lorsque le premier fil en a fini avec sa mission, il retourne dans le groupe des fils inactifs, prêt à gérer une nouvelle tâche. Un groupe de fils vous permet de traiter plusieurs connexions en simultané, ce qui augmente le débit de votre serveur.

Nous allons limiter le nombre de tâches dans le groupe à un petit nombre pour nous protéger d’attaques par déni de service (Denial of Service, DoS) ; si notre programme créait une nouvelle tâche à chaque requête qu’il reçoit, quelqu’un qui ferait 10 millions de requêtes à notre serveur pourrait faire des ravages en utilisant toutes les ressources de notre serveur et bloquer ainsi le traitement de toute nouvelle requête.

Plutôt que de générer des fils d’exécution en quantité illimitée, nous allons donc faire en sorte qu’il y ait un nombre fixe de fils qui seront en attente dans le groupe. Les requêtes entrantes sont envoyées au groupe pour être traitées. Le groupe gèrera une file d’attente pour les requêtes entrantes. Chaque fil dans le groupe va récupérer une requête dans cette liste d’attente, la traiter puis demander une autre requête à la file d’attente. Avec ce fonctionnement, nous pouvons traiter jusqu’à N requêtes en concurrence, où N est le nombre de fils d’exécution. Si chaque fil répond à une requête longue à traiter, les requêtes suivantes vont s’accumuler dans la file d’attente, mais nous aurons quand même augmenté le nombre de requêtes longues que nous pouvons traiter avant d’en arriver là.

Cette technique n’est qu’une des nombreuses manières d’améliorer le débit d’un serveur web. D’autres options que vous devriez envisager sont le modèle fork/join, le modèle d’entrée-sortie asynchrone monotâche, et le modèle d’entrée-sortie asynchrone multifils. Si vous êtes intéressés par ce sujet, vous pouvez aussi en apprendre plus sur ces autres solutions et essayer de les implémenter en Rust ; avec un langage bas niveau comme Rust, toutes les options restent possibles.

Avant que nous ne commencions l’implémentation du groupe de tâches, parlons de l’utilisation du groupe. Lorsque vous essayez de concevoir du code, commencer par écrire l’interface client peut vous aider à vous guider dans la conception. Ecrivez l’API du code afin qu’il soit structuré de la manière dont vous souhaitez l’appeler ; puis implémentez ensuite la fonctionnalité au sein de cette structure, plutôt que d’implémenter la fonctionnalité puis de concevoir l’API publique.

De la même manière que nous avons utilisé le développement piloté par les tests dans le projet du chapitre 12, nous allons utiliser ici le développement orienté par le compilateur. Nous allons écrire le code qui appelle les fonctions que nous souhaitons, et ensuite nous analyserons les erreurs du compilateur pour déterminer ce qu’il faut ensuite corriger pour que le code fonctionne. Cependant, avant de faire cela, en guise de point de départ, nous allons explorer la technique que nous n’allons pas mettre en œuvre.

Création d’un fil d’exécution par requête

Pour commencer, voyons à quoi ressemblerait notre code s’il créait un nouveau fil pour chaque connexion. Comme nous l’avons évoqué précédemment, cela ne sera pas notre solution finale à cause des problèmes liés à la création potentielle d’un nombre illimité de fils, mais c’est un début pour ensuite aboutir à un serveur multifils opérationnel. Par la suite, nous ajouterons le groupe de fils d’exécution comme une amélioration, de sorte que la comparaison entre les deux solutions sera plus aisée.

L’encart 21-11 montre les changements à apporter au main pour créer un nouveau fil d’exécution pour gérer chaque flux avec une boucle for.

Filename: src/main.rs
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};

fn main() {
    let ecouteur = TcpListener::bind("127.0.0.1:7878").unwrap();

    for flux in ecouteur.incoming() {
        let flux = flux.unwrap();

        thread::spawn(|| {
            gestion_connexion(flux);
        });
    }
}

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();
}
Listing 21-11: Spawning a new thread for each stream

Comme vous l’avez appris au chapitre 16, thread::spawn va créer un nouveau fil d’exécution puis exécuter dans ce nouveau fil le code présent dans la fermeture. Si vous exécutez ce code et chargez /pause dans votre navigateur, et que vous ouvrez / dans deux nouveaux onglets, vous constaterez en effet que les requêtes vers / n’auront pas à attendre que /pause se finisse. Toutefois, comme nous l’avons mentionné, cela peut potentiellement surcharger le système si vous créez de nouveaux fils d’exécution sans aucune limite.

Peut-être vous souvenez-vous du chapitre 17, qu’il s’agit précisément du genre de situation où async et await brillent ! Gardez cela en tête pendant que nous construisons le groupe de fils et pensez à la manière dont les choses seraient différentes ou similaires avec async.

Création d’un nombre fini de fils d’exécution

Nous souhaitons faire en sorte que notre groupe de fils fonctionne de la même manière, donc passer des fils à un groupe de fils ne devrait pas nécessiter de gros changements au code qui utilise notre API. L’encart 21-12 montre une interface possible pour une structure GroupeTaches que nous souhaitons utiliser à la place de thread::spawn.

Filename: src/main.rs
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() {
        let flux = flux.unwrap();

        groupe.executer(|| {
            gestion_connexion(flux);
        });
    }
}

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();
}
Listing 21-12: Our ideal ThreadPool interface

Nous avons utilisé GroupeTaches::new pour créer un nouveau groupe de fils avec un nombre configurable de fils, dans notre cas, quatre. Ensuite, dans la boucle for, groupe.executer a une interface similaire à thread::spawn qui prend une fermeture que le groupe devra exécuter pour chaque flux. Nous devons implémenter groupe.executer pour qu’il prenne la fermeture et la donne à un fil dans le groupe pour qu’il l’exécute. Ce code ne se compile pas encore, mais nous allons faire comme si c’était le cas pour que le compilateur puisse nous guider dans la résolution des problèmes.

Construction de la structure GroupeTaches en utilisant le développement orienté par le compilateur

Faites les changements de l’encart 21-12 dans votre src/main.rs, et utilisez ensuite les erreurs du compilateur lors du cargo check pour orienter votre développement. Voici la première erreur que nous obtenons :

$ cargo check
    Checking salutations v0.1.0 (file:///projects/salutations)
error[E0433]: cannot find type `ThreadPool` in this scope
  --> src/main.rs:10:16
   |
11 |     let groupe = GroupeTaches::new(4);
   |                  ^^^^^^^^^^^^ use of undeclared type `GroupeTaches`

For more information about this error, try `rustc --explain E0433`.
error: could not compile `salutations` (bin "salutations") due to 1 previous error

Bien ! Cette erreur nous informe que nous avons besoin d’un type ou d’un module qui s’appelle GroupeTaches, donc nous allons le créer. Notre implémentation de GroupeTaches sera indépendante du type de travail qu’accomplira notre serveur web. Donc, transformons la crate binaire salutations en crate de bibliothèque pour y implémenter notre GroupeTaches. Après l’avoir changé en crate de bibliothèque, nous pourrons utiliser ensuite cette bibliothèque de groupe de fils dans n’importe quel projet où nous aurons besoin d’un groupe de fils, et pas seulement pour servir des requêtes web.

Créez un fichier src/lib.rs qui contient ce qui suit, à savoir la définition la plus simple d’une structure GroupeTaches que nous pouvons avoir pour le moment :

Filename: src/lib.rs
pub struct GroupeTaches;

Ensuite, éditez le fichier main.rs pour importer GroupeTaches dans la portée à partir de la bibliothèque de crate en ajoutant le code suivant au début de src/main.rs :

Filename: src/main.rs
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() {
        let flux = flux.unwrap();

        groupe.executer(|| {
            gestion_connexion(flux);
        });
    }
}

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();
}

Ce code ne fonctionne toujours pas, mais vérifions-le à nouveau pour obtenir l’erreur que nous devons maintenant résoudre :

$ cargo check
    Checking salutations v0.1.0 (file:///projects/salutations)
error[E0599]: no associated function or constant named `new` found for struct `GroupeTaches` in the current scope
  --> src/main.rs:12:28
   |
12 |     let groupe = GroupeTaches::new(4);
   |                                ^^^ associated function or constant not found in `GroupeTaches`

For more information about this error, try `rustc --explain E0599`.
error: could not compile `salutations` (bin "salutations") due to 1 previous error

Cette erreur indique que nous devons ensuite créer une fonction associée new pour GroupeTaches. Nous savons aussi que new nécessite d’avoir un paramètre qui peut accepter 4 comme argument et doit retourner une instance de GroupeTaches. Implémentons la fonction new la plus simple possible qui aura ces caractéristiques :

Filename: src/lib.rs
pub struct GroupeTaches;

impl GroupeTaches {
    pub fn new(taille: usize) -> GroupeTaches {
        GroupeTaches
    }
}

Nous avons choisi usize comme type du paramètre taille, car nous savons qu’un nombre négatif de fils d’exécution n’a pas de sens. Nous savons également que nous allons utiliser ce 4 comme étant le nombre d’éléments dans une collection de fils d’exécution, ce qui est à quoi sert le type usize, comme nous l’avons vu dans la section “Types de nombres entiers” du chapitre 3.

Vérifions à nouveau le code :

$ cargo check
    Checking salutations v0.1.0 (file:///projects/salutations)
error[E0599]: no method named `executer` found for struct `GroupeTaches` in the current scope
  --> src/main.rs:17:14
   |
17 |         groupe.executer(|| {
   |                ^^^^^^^^ method not found in `GroupeTaches`

For more information about this error, try `rustc --explain E0599`.
error: could not compile `salutations` (bin "salutations") due to 1 previous error

Désormais, nous obtenons une erreur car nous n’avons pas implémenté la méthode executer sur GroupeTaches. Souvenez-vous que nous avions décidé dans la section “Création d’un nombre fini de fils d’exécution” que notre groupe de fils d’exécution devrait avoir une interface similaire à thread::spawn. C’est pourquoi nous allons implémenter la fonction executer pour qu’elle prenne en argument la fermeture qu’on lui donne et qu’elle la passe à un fil d’exécution inactif du groupe pour qu’elle l’exécute.

Nous allons définir la méthode executer sur GroupeTaches pour prendre en paramètre une fermeture. Souvenez-vous que nous avions vu dans la section “Déplacement de valeurs capturées en-dehors des fermetures du chapitre 13 que nous pouvions prendre en paramètre les fermetures avec trois types de traits différents : Fn, FnMut, et FnOnce. Nous devons décider quel genre de fermeture nous allons utiliser ici. Nous savons que nous allons faire quelque chose de sensiblement identique à l’implémentation du thread::spawn de la bibliothèque standard, donc nous pouvons nous inspirer de ce qui lie la signature de thread::spawn à son paramètre. La documentation nous donne ceci :

pub fn spawn<F, T>(f: F) -> JoinHandle<T>
    where
        F: FnOnce() -> T,
        F: Send + 'static,
        T: Send + 'static,

Le paramètre de type F est celui qui nous intéresse ici ; le paramètre de type T est lié à la valeur de retour, et ceci ne nous intéresse pas ici. Nous pouvons constater que spawn utilise le trait FnOnce lié à F. C’est probablement ce dont nous avons besoin, parce que nous allons sûrement passer cet argument dans le execute de spawn. Nous pouvons aussi être sûr que FnOnce est le trait dont nous avons besoin car le fil qui va traiter une requête ne va le faire qu’une seule fois, ce qui correspond à la partie Once dans FnOnce.

Le paramètre de type F a aussi le trait lié Send et la durée de vie liée 'static, qui sont utiles dans notre situation : nous avons besoin de Send pour transférer la fermeture d’un fil d’exécution vers une autre et de 'static car nous ne connaissons pas la durée d’exécution du fil. Créons donc une méthode executer sur GroupeTaches qui va utiliser un paramètre générique de type F avec les liens suivants :

Filename: src/lib.rs
pub struct GroupeTaches;

impl GroupeTaches {
    // -- partie masquée ici --
    pub fn new(taille: usize) -> GroupeTaches {
        GroupeTaches
    }

    pub fn executer<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
    }
}

Nous utilisons toujours le () après FnOne car ce FnOnce représente une fermeture qui ne prend pas de paramètres et retourne le type unité (). Exactement comme les définitions de fonctions, le type de retour peut être omis de la signature, mais même si elle ne contient pas de paramètre, nous avons tout de même besoin des parenthèses.

À nouveau, c’est l’implémentation la plus simpliste de la méthode executer : elle ne fait rien, mais nous essayons seulement de faire en sorte que notre code se compile. Vérifions-le à nouveau :

$ cargo check
    Checking salutations v0.1.0 (file:///projects/salutations)
    Finished `dev` profile [unoptimized + debuginfo] target(s) in 0.24s

Cela se compile ! Mais remarquez que si vous lancez cargo run et faites la requête dans votre navigateur web, vous verrez l’erreur dans le navigateur que nous avions tout au début du chapitre. Notre bibliothèque n’exécute pas encore la fermeture envoyée à executer !

Remarque : un dicton que vous avez probablement déjà entendu à propos des compilateurs stricts, comme Haskell et Rust, est que “si le code se compile, il fonctionne”. Mais ce dicton n’est pas toujours vrai. Notre projet se compile, mais il ne fait absolument rien ! Si nous construisions un vrai projet, complexe, il serait bon de commencer à écrire des tests unitaires pour vérifier que ce code compile et qu’il suit le comportement que nous souhaitons.

Songeons-y : qu’est-ce qui pourrait être différent ici si nous exécutions une future à la place d’une fermeture ?

Validation du nombre de fils d’exécutions envoyé à new

Nous ne faisons rien avec les paramètres passés à new et executer. Implémentons le corps de ces fonctions avec le comportement que nous souhaitons. Pour commencer, réfléchissons à new. Précédemment, nous avions choisi un type sans signe pour le paramètre taille, car un groupe avec un nombre négatif de fils n’a pas de sens. Cependant, un groupe avec aucun fil n’a pas non plus de sens, alors que zéro est une valeur parfaitement valide pour usize. Nous allons ajouter du code pour vérifier que taille est plus grand que zéro avant de retourner une instance de GroupeTaches et nous allons faire en sorte que le programme panique s’il reçoit un zéro, en utilisant la macro assert! comme dans l’encart 21-13.

Filename: src/lib.rs
pub struct GroupeTaches;

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);

        GroupeTaches
    }

    // -- partie masquée ici --
    pub fn executer<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
    }
}
Listing 21-13: Implementing ThreadPool::new to panic if size is zero

Nous avons aussi ajouté un peu de documentation pour notre GroupeTaches avec des commentaires de documentation. Remarquez que nous avons suivi les pratiques de bonne documentation en ajoutant une section qui liste les situations pour lesquelles notre fonction peut paniquer, comme nous l’avons vu dans le chapitre 14. Essayez de lancer cargo doc --open et de cliquer sur la structure GroupeTaches pour voir à quoi ressemble la documentation générée pour new !

Au lieu d’ajouter la macro assert! comme nous venons de le faire, nous aurions pu changer new en build et retourner un Result, comme nous l’avions fait avec Config::build dans le projet d’entrée/sortie dans l’encart 12-9. Mais nous avons décidé que dans le cas présent, la création d’un groupe de fils d’exécution sans aucun fil devait être une erreur irrécupérable. Si vous en sentez l’envie, essayez d’écrire une fonction nommée build avec la signature suivante, pour comparer avec la fonction new :

pub fn build(taille: usize) -> Result<GroupeTaches, ErreurGroupeTaches> {

Création d’espace de rangement des fils d’exécution

Maintenant que nous avons une manière de savoir si nous avons un nombre valide de fils d’exécution à stocker dans le groupe, nous pouvons créer ces fils et les stocker dans la structure GroupeTaches avant de retourner la structure. Mais comment “stocker” un fil d’exécution ? Regardons à nouveau la signature de thread::spawn :

pub fn spawn<F, T>(f: F) -> JoinHandle<T>
    where
        F: FnOnce() -> T,
        F: Send + 'static,
        T: Send + 'static,

La fonction spawn retourne un JoinHandle<T>, où T est le type que retourne notre fermeture. Essayons d’utiliser nous aussi JoinHandle pour voir ce qu’il va se passer. Dans notre cas, les fermetures que nous passons dans le groupe de tâches vont traiter les connexions mais ne vont rien retourner, donc T sera le type unité, ().

Le code de l’encart 21-14 va se compiler mais ne va pas encore créer de tâches pour le moment. Nous avons changé la définition de GroupeTaches pour qu’elle possède un vecteur d’instances thread::JoinHandle<()>, nous avons initialisé le vecteur avec une capacité de la valeur de taille, mis en place une boucle for qui va exécuter du code pour créer les tâches puis nous avons retourné une instance de GroupeTaches qui les contient.

Filename: src/lib.rs
use std::thread;

pub struct GroupeTaches {
    taches: Vec<thread::JoinHandle<()>>,
}

impl GroupeTaches {
    // -- partie masquée ici --
    /// 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 mut taches = Vec::with_capacity(taille);

        for _ in 0..taille {
            // on crée quelques tâches ici et on les stocke dans le vecteur
        }

        GroupeTaches { taches }
    }
    // -- partie masquée ici --

    pub fn executer<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
    }
}
Listing 21-14: Creating a vector for ThreadPool to hold the threads

Nous avons importé std::thread dans la portée de la crate de bibliothèque, car nous utilisons thread::JoinHandle comme étant le type des éléments du vecteur dans GroupeTaches.

Une fois qu’une taille valide est reçue, notre GroupeTaches crée un nouveau vecteur qui peut stocker taille éléments. La fonction with_capacity fait la même chose que Vec::new mais avec une grosse différence : elle pré-alloue l’espace dans le vecteur. Comme nous savons que nous avons besoin de stocker taille éléments dans le vecteur, faire cette allocation en amont est bien plus efficace que d’utiliser Vec::new qui va se redimensionner lorsque des éléments lui seront ajoutés.

Lorsque vous lancez à nouveau cargo check, cela devrait être un succès.

Envoi de code depuis GroupeTaches vers un fil d’exécution

Nous avions laissé un commentaire dans la boucle for dans l’encart 21-14 qui concernait la création des fils d’exécution. Maintenant, nous allons voir comment créer ces fils. La bibliothèque standard fournit un moyen de créer des fils d’exécution avec thread::spawn à qui il faut passer le code que le fil doit exécuter dès qu’il est créé. Cependant, dans notre cas, nous souhaitons créer des fils et faire en sorte qu’elles attendent du code que nous leur enverrons plus tard. L’implémentation des fils d’exécution de la bibliothèque standard n’offre aucun moyen de faire ceci ; nous devons donc implémenter cela nous-même.

Nous allons implémenter ce comportement en introduisant une nouvelle structure de données entre le GroupeTaches et les fils qui va gérer ce nouveau comportement. Nous allons appeler cette structure Operateur, nom qui lui est traditionnellement donné avec Worker dans les implémentations de groupe de fils d’exécution. Le Worker récupère le code qui doit être exécuté et l’exécute dans son fil d’exécution.

Imaginez des personnes qui travaillent dans la cuisine d’un restaurant : les opérateurs attendent les commandes des clients puis sont chargés de prendre en charge ces commandes et d’y répondre.

Au lieu de stocker un vecteur d’instances JoinHandle<()> dans le groupe de fils d’exécution, nous allons stocker des instances de structure Operateur. Chaque Operateur va stocker une seule instance de JoinHandle<()>. Ensuite nous implémenterons une méthode sur Operateur qui va prendre en argument une fermeture de code à exécuter et l’envoyer au fil qui fonctionne déjà pour exécution. Nous allons aussi donner à chaque Operateur un identifiant id afin que nous puissions distinguer les différentes instances de Operateur du groupe dans les journaux ou lors de débogages.

Voici le nouveau traitement qui va avoir lieu quand nous créons un GroupeTaches. Nous allons implémenter le code qui envoie la fermeture au fil d’exécution après avoir configuré Operateur de cette manière :

  1. définir une structure Operateur qui possède un id et un JoinHandle<()> ;
  2. modifier le GroupeTaches afin qu’il possède un vecteur d’instances de Operateur ;
  3. définir une fonction Operateur::new qui prend en argument un numéro d’id et retourne une instance de Operateur qui contient l’ id et un fil d’exécution créé avec une fermeture vide ;
  4. dans GroupeTaches::new, utiliser le compteur de la boucle for pour générer un id, créer un nouveau Operateur avec cet id et stocker l’opérateur dans le vecteur.

Si vous vous sentez prêt à relever le défi, essayez de faire ces changements de votre côté avant de regarder le code de l’encart 21-15.

Vous êtes prêt(e) ? Voici l’encart 21-15 qui propose une solution pour procéder aux changements listés précédemment.

Filename: src/lib.rs
use std::thread;

pub struct GroupeTaches {
    operateurs: Vec<Operateur>,
}

impl GroupeTaches {
    // -- partie masquée ici --
    /// 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 mut operateurs = Vec::with_capacity(taille);

        for id in 0..taille {
            operateurs.push(Operateur::new(id));
        }

        GroupeTaches { operateurs }
    }
    // -- partie masquée ici --

    pub fn executer<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
    }
}

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

impl Operateur {
    fn new(id: usize) -> Operateur {
        let tache = thread::spawn(|| {});

        Operateur { id, tache }
    }
}
Listing 21-15: Modifying ThreadPool to hold Worker instances instead of holding threads directly

Nous avons changé le nom du champ taches de GroupeTaches en operateurs car il stocke maintenant des instances de Operateur plutôt que des instances de JoinHandle<()>. Nous utilisons le compteur de la boucle for comme argument de Operateur::new et nous stockons chacun des nouveaux Operateur dans le vecteur operateurs.

Le code externe (comme celui de notre serveur dans src/bin/main.rs) n’a pas besoin de connaître les détails de l’implémentation qui utilise une structure Operateur dans GroupeTaches, donc nous faisons en sorte que la structure Operateur et sa fonction new soient privées. La fonction Operateur::new utilise l’ id que nous lui donnons et stocke une instance de JoinHandle<()> qui est créée en instanciant un nouveau fil d’exécution utilisant une fermeture vide.

Remarque : si le système d’exploitation ne peut pas créer de fil d’exécution faute de ressources systèmes, thread::spawn va paniquer. Ceci va provoquer la panique de l’ensemble de notre serveur, bien que la création de certains fils d’exécution puisse réussir. Par souci de simplicité, ce comportement est acceptable, mais dans une implémentation d’un groupe de fils d’exécution dans le monde réel, vous utiliseriez certainement std::thread::Builder et sa méthode spawn qui, à la place, renvoie Result.

Ce code va se compiler et stocker le nombre d’instances de Operateur que nous avons renseigné en argument de GroupeTaches::new. Mais nous n’exécutons toujours pas la fermeture que nous obtenons de executer. Voyons maintenant comment faire cela.

Envoyer des requêtes à des fils d’exécution via des canaux

Le problème suivant auquel nous allons nous attaquer est le fait que les fermetures passées à thread::spawn ne font absolument rien. Actuellement, nous obtenons la fermeture que nous souhaitons exécuter dans la méthode executer. Mais nous avons besoin de donner une fermeture à thread::spawn à exécuter lorsque nous créons chaque Operateur lors de la création de GroupeTaches.

Nous souhaitons que les structures Operateur que nous venons de créer récupèrent le code à exécuter dans une liste d’attente présente dans le GroupeTaches et renvoient ce code à leur fil d’exécution pour l’exécuter.

Les canaux que nous avons vus dans le chapitre 16 (une manière simple de communiquer entre deux fils d’exécution) seront parfaits pour ce cas d’emploi. Nous allons utiliser un canal pour les fonctions pour créer la liste d’attente des missions, et executer devrait envoyer une mission de GroupeTaches vers les instances Operateur, qui vont passer la mission à leurs fils d’exécution. Voici le plan :

  1. le GroupeTaches va créer un canal et se connecter à l’émetteur ;
  2. chaque Operateur va se connecter au récepteur ;
  3. nous allons créer une nouvelle structure Mission qui va stocker les fermetures que nous souhaitons envoyer dans le canal ;
  4. la méthode executer va envoyer la mission qu’elle souhaite executer en passant par l’émetteur ;
  5. dans son propre fil d’exécution, l’Operateur va vérifier en permanence le récepteur et exécuter les fermetures de toutes les missions qu’il va recevoir.

Commençons par créer un canal dans GroupeTaches::new et stocker l’émetteur dans l’instance de GroupeTaches, comme dans l’encart 21-16. La structure Mission ne contient rien pour le moment mais sera le type d’éléments que nous enverrons dans le canal.

Filename: src/lib.rs
use std::{sync::mpsc, thread};

pub struct GroupeTaches {
    operateurs: Vec<Operateur>,
    envoi: mpsc::Sender<Mission>,
}

struct Mission;

impl GroupeTaches {
    // -- partie masquée ici --
    /// 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 mut operateurs = Vec::with_capacity(taille);

        for id in 0..taille {
            operateurs.push(Operateur::new(id));
        }

        GroupeTaches { operateurs, envoi }
    }
    // -- partie masquée ici --

    pub fn executer<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
    }
}

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

impl Operateur {
    fn new(id: usize) -> Operateur {
        let tache = thread::spawn(|| {});

        Operateur { id, tache }
    }
}
Listing 21-16: Modifying ThreadPool to store the sender of a channel that transmits Job instances

Dans GroupeTaches::new, nous créons notre nouveau canal et faisons en sorte que le groupe stocke l’émetteur. Cela devrait pouvoir se compiler.

Essayons de passer un récepteur du canal à chaque Operateur lorsque le groupe de fils d’exécution crée le canal. Nous savons que nous voulons utiliser le récepteur dans le fil que l’instance de Operateurcrée, donc nous allons créer une référence vers le paramètre reception dans la fermeture. Le code de l’encart 21-17 ne se compile pas encore.

Filename: src/lib.rs
use std::{sync::mpsc, thread};

pub struct GroupeTaches {
    operateurs: Vec<Operateur>,
    envoi: mpsc::Sender<Mission>,
}

struct Mission;

impl GroupeTaches {
    // -- partie masquée ici --
    /// 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 mut operateurs = Vec::with_capacity(taille);

        for id in 0..taille {
            operateurs.push(Operateur::new(id, reception));
        }

        GroupeTaches { operateurs, envoi }
    }
    // -- partie masquée ici --

    pub fn executer<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
    }
}

// -- partie masquée ici --


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

impl Operateur {
    fn new(id: usize, reception: mpsc::Receiver<Mission>) -> Operateur {
        let tache = thread::spawn(|| {
            reception;
        });

        Operateur { id, tache }
    }
}
Listing 21-17: Passing the receiver to each Worker

Nous avons juste fait de petites modifications simples : nous envoyons le récepteur dans Operateur::new puis nous l’utilisons dans la fermeture.

Lorsque nous essayons de vérifier ce code, nous obtenons cette erreur :

$ cargo check
    Checking salutations v0.1.0 (file:///projects/salutations)
error[E0382]: use of moved value: `reception`
  --> src/lib.rs:26:42
   |
22 |         let (envoi, reception) = mpsc::channel();
   |                     --------- move occurs because `reception` has type `std::sync::mpsc::Receiver<Mission>`, which does not implement the `Copy` trait
...
26 |         for id in 0..taille {
   |         ------------------- inside of this loop
27 |             operateurs.push(Operateur::new(id, reception));
   |                                                ^^^^^^^^^ value moved here, in previous iteration of loop
   |
note: consider changing this parameter type in method `new` to borrow instead if owning the value isn't necessary
  --> src/lib.rs:47:33
   |
50 |     fn new(id: usize, reception: mpsc::Receiver<Mission>) -> Operateur {
   |        --- in this method        ^^^^^^^^^^^^^^^^^^^^^^^ this parameter takes ownership of the value

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

Le code essaye d’envoyer reception dans plusieurs instances de Operateur. Ceci ne fonctionne pas, comme vous l’avez appris au chapitre 16 : l’implémentation du canal que fournit Rust est du type plusieurs producteurs, un seul consommateur. Cela signifie que nous ne pouvons pas simplement cloner la partie réceptrice du canal pour corriger ce code. Nous ne voulons pas non plus envoyer un message plusieurs fois à plusieurs consommateurs ; nous voulons une liste de messages avec plusieurs instances de Operateur de manière à ce que chaque message soit traité une fois.

De plus, obtenir une mission de la file d’attente du canal implique de modifier la reception, donc les tâches ont besoin d’une méthode sécurisée pour partager et modifier reception ; autrement, nous risquons de nous trouver dans des situations de concurrence (comme nous l’avons vu dans le chapitre 16).

Souvenez-vous des pointeurs intelligents conçus pour les échanges entre les tâches que nous avons vus au chapitre 16 : pour partager la possession entre plusieurs tâches et permettre aux tâches de modifier la valeur, nous avons besoin d’utiliser Arc<Mutex<T>>. Le type Arc va permettre à plusieurs instances de Operateur de posséder la réception tandis que Mutex va s’assurer qu’un Operateur obtienne une mission de la réception à un moment donné. L’encart 21-18 montre les changements que nous devons apporter.

Filename: src/lib.rs
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};
// -- partie masquée ici --

pub struct GroupeTaches {
    operateurs: Vec<Operateur>,
    envoi: mpsc::Sender<Mission>,
}

struct Mission;

impl GroupeTaches {
    // -- partie masquée ici --
    /// 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 }
    }

    // -- partie masquée ici --

    pub fn executer<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
    }
}

// -- partie masquée ici --

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

impl Operateur {
    fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
        // -- partie masquée ici --
        let tache = thread::spawn(|| {
            reception;
        });

        Operateur { id, tache }
    }
}
Listing 21-18: Sharing the receiver among the Worker instances using Arc and Mutex

Dans GroupeTaches::new, nous installons le récepteur dans un Arc et un Mutex. Pour chaque nouveau Operateur, nous clonons le Arc pour augmenter le compteur de références afin que les instances de Operateur puissent se partager la possession du récepteur.

Grâce à ces changements, le code se compile ! Nous touchons au but !

Implémentation de la méthode executer

Finissons en implémentant la méthode executer de GroupeTaches. Nous allons également modifier la structure Mission pour la transformer en un alias de type pour un objet trait qui contiendra le type de fermeture que executer recevra. Comme nous l’avons vu dans la section “Synonymes de types et alias de types” du chapitre 20, les alias de type nous permettent de raccourcir les types un peu trop longs. Voyez cela dans l’encart 21-19.

Filename: src/lib.rs
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct GroupeTaches {
    operateurs: Vec<Operateur>,
    envoi: mpsc::Sender<Mission>,
}

// -- partie masquée ici --

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

impl GroupeTaches {
    // -- partie masquée ici --
    /// 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();
    }
}

// -- partie masquée ici --

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

impl Operateur {
    fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
        let tache = thread::spawn(|| {
            reception;
        });

        Operateur { id, tache }
    }
}
Listing 21-19: Creating a Job type alias for a Box that holds each closure and then sending the job down the channel

Après avoir créé une nouvelle instance Mission en utilisant la fermeture que nous obtenons dans executer, nous envoyons cette mission dans le canal via la partie émettrice. Nous utilisons unwrap sur send pour les cas où l’envoi échoue. Cela peut arriver si, par exemple, nous stoppons l’exécution de tous les fils d’exécution, ce qui signifiera que les parties réceptrices auront fini de recevoir des nouveaux messages. Pour le moment, nous ne pouvons pas stopper l’exécution de nos fils d’exécution : nos fils continuerons à s’exécuter aussi longtemps que le groupe existe. La raison pour laquelle nous utilisons unwrap est que nous savons que le cas d’échec ne va pas se produire, mais le compilateur ne le sait pas.

Mais nous n’avons pas encore tout à fait fini ! Dans Operateur, notre fermeture envoyée à thread::spawn ne fait que référencer la partie réception du canal. Au lieu de ça, nous avons besoin que la fermeture boucle à l’infini, demandant une mission à la partie réceptrice du canal et l’exécutant quand elle en obtient une. Appliquons les changements montrés dans l’encart 21-20 à Operateur::new.

Filename: src/lib.rs
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();
    }
}

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

// -- partie masquée ici --

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 }
    }
}
Listing 21-20: Receiving and executing the jobs in the Worker instance’s thread

Ici, nous faisons d’abord appel à lock sur reception pour obtenir le mutex, puis nous faisons appel à unwrap pour paniquer dès qu’il y a une erreur. L’acquisition d’un verrou peut échouer si le mutex est dans un état empoisonné, ce qui peut arriver si d’autres fils d’exécution ont paniqué pendant qu’elles avaient le verrou au lieu de le rendre. Dans cette situation, l’appel à unwrap fera paniquer le fil, ce qui est la bonne chose à faire. Vous pouvez aussi changer ce unwrap en un expect avec un message d’erreur qui sera plus explicite pour vous.

Si nous obtenons le verrou du mutex, nous faisons appel à recv pour recevoir une Mission provenant du canal. Un unwrap final s’occupe lui aussi des cas d’erreurs qui peuvent se produire si le fil d’exécution qui est connecté à l’émetteur se termine, de la même manière que la méthode send enverrait Err si le récepteur se fermerait.

L’appel à recv bloque l’exécution, donc s’il n’y a pas encore de mission, le fil courant va attendre jusqu’à ce qu’une mission soit disponible. Le Mutex<T> s’assure qu’un seul fil d’exécution d’Operateur essaie d’obtenir une mission à un instant donné.

Notre groupe de fils d’exécution est désormais en état de fonctionner ! Faites un cargo run et faites quelques requêtes :

$ cargo run
   Compiling salutations v0.1.0 (file:///projects/salutations)
warning: field `operateurs` is never read
 --> src/lib.rs:7:5
  |
6 | pub struct GroupeTaches {
  |            ------------ field in this struct
7 |     operateurs: Vec<Operateur>,
  |     ^^^^^^^^^^
  |
  = note: `#[warn(dead_code)]` on by default

warning: fields `id` and `tache` are never read
  --> src/lib.rs:48:5
   |
47 | struct Operateur {
   |        --------- fields in this struct
48 |     id: usize,
   |     ^^
49 |     tache: thread::JoinHandle<()>,
   |     ^^^^^

warning: `salutations` (lib) generated 2 warnings
    Finished `dev` profile [unoptimized + debuginfo] target(s) in 4.91s
     Running `target/debug/salutations`
L'opérateur 0 a reçu une mission ; il l'exécute.
L'opérateur 2 a reçu une mission ; il l'exécute.
L'opérateur 1 a reçu une mission ; il l'exécute.
L'opérateur 3 a reçu une mission ; il l'exécute.
L'opérateur 0 a reçu une mission ; il l'exécute.
L'opérateur 2 a reçu une mission ; il l'exécute.
L'opérateur 1 a reçu une mission ; il l'exécute.
L'opérateur 3 a reçu une mission ; il l'exécute.
L'opérateur 0 a reçu une mission ; il l'exécute.
L'opérateur 2 a reçu une mission ; il l'exécute.

Parfait ! Nous avons maintenant un groupe de fils qui exécute des connexions de manière asynchrone. Il n’y a jamais plus de quatre fils qui sont créées, donc notre système ne sera pas surchargé si le serveur reçoit beaucoup de requêtes. Si nous faisons une requête vers /pause, le serveur sera toujours capable de servir les autres requêtes grâce aux autres fils qui pourront les exécuter.

Remarque : si vous ouvrez /pause dans plusieurs fenêtres de navigation en simultané, elles peuvent parfois être chargées une par une avec cinq secondes d’intervalle. Certains navigateurs web exécutent plusieurs instances de la même requête de manière séquentielle pour des raisons de mise en cache. Cette limitation n’est pas imputable à notre serveur web.

Le moment est bien choisi pour faire une pause et voir comment le code des encarts 21-18, 21-19 et 21-20 pourraient être différents si nous avions utilisé des futures à la place d’une fermeture pour accomplir le travail. Quels types changeraient ? Dans quelle mesure les signatures de méthodes seraient différentes, et le seraient-elles vraiment ? Quelles portions du code resteraient identiques ?

Ayant appris la boucle while let dans les chapitres 17 et 18, vous pourriez vous demander pourquoi nous n’avons pas écrit le code des fils d’exécution de Operateur comme dans l’encart 21-21.

Filename: src/lib.rs
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();
    }
}

struct Operateur {
    id: usize,
    tache: thread::JoinHandle<()>,
}
// -- partie masquée ici --

impl Operateur {
    fn new(id: usize, reception: Arc<Mutex<mpsc::Receiver<Mission>>>) -> Operateur {
        let tache = thread::spawn(move || {
            while let Ok(mission) = reception.lock().unwrap().recv() {
                println!("L'opérateur {id} a reçu une mission ; il l'exécute.");

                mission();
            }
        });

        Operateur { id, tache }
    }
}
Listing 21-21: An alternative implementation of Worker::new using while let

Ce code se compile et s’exécute mais ne se produit pas le comportement des tâches que nous souhaitons : une requête lente à traiter va continuer à mettre en attente de traitement les autres requêtes. La raison à cela est subtile : la structure Mutex n’a pas de méthode publique unlock car la propriété du verrou se base sur la durée de vie du MutexGuard<T> au sein du LockResult<MutexGuard<T>> que retourne la méthode lock. À la compilation, le vérificateur d’emprunt peut ensuite vérifier la règle qui dit qu’une ressource gardée par un Mutex ne peut être accessible que si nous avons ce verrou. Cependant, cette implémentation peut aussi conduire à ce que nous gardions le verrou plus longtemps que prévu si nous ne prenons pas compte la durée de vie du MutexGuard<T>.

Le code de l’encart 21-20 qui utilise let mission = reception.lock().unwrap().recv().unwrap(); fonctionne, car avec let, toute valeur temporaire utilisée dans la partie droite du signe égal est libérée immédiatement lorsque l’instruction let se termine. Cependant, while let (ainsi que if let et match) ne libèrent pas les valeurs temporaires avant la fin du bloc associé. Dans l’encart 21-21, le verrou continue à être maintenu pendant toute la durée de l’appel à mission(), ce qui veut dire que les autres instances de Operateur ne peuvent pas recevoir de tâches.