Canais (channels) são uma primitiva de concorrência que permitem a comunicação entre threads por meio do envio e recebimento de mensagens. Em Rust, o módulo std::sync::mpsc (Multiple Producer, Single Consumer) implementa canais com múltiplos produtores e um único consumidor. Essa abordagem é inspirada no modelo de atores e no CSP (Communicating Sequential Processes), promovendo um design seguro e livre de condições de corrida, já que os dados são transferidos entre threads sem compartilhamento de estado.

Nesta aula, exploraremos como criar canais, enviar e receber mensagens, lidar com múltiplos produtores e aplicar padrões comuns, como distribuição de tarefas e coleta de resultados.

mpsc

O módulo std::sync::mpsc fornece canais assíncronos e síncronos. A função channel() retorna um par (Sender, Receiver). O Sender é clonável para múltiplos produtores, enquanto o Receiver não pode ser clonado, garantindo um único consumidor. O canal é ilimitado por padrão (assíncrono), mas há também sync_channel com capacidade limitada.

Exemplo básico de criação de um canal:

use std::sync::mpsc;
use std::thread;

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        tx.send(String::from("oi")).unwrap();
    });

    let received = rx.recv().unwrap();
    println!("Recebido: {}", received);
}

O Sender é movido para a thread filha via move. O método send retorna Result; se o receptor foi descartado, retorna um erro. O recv bloqueia até receber uma mensagem ou até que todos os produtores sejam descartados, retornando erro se não houver mais mensagens.

send e recv

O método send envia um valor para o canal, consumindo-o. O método recv bloqueia a thread atual até que uma mensagem seja recebida. Há também try_recv para tentativa não bloqueante, retornando Result com Ok ou Err (vazio ou desconectado).

Exemplo com múltiplas mensagens e uso de try_recv:

use std::sync::mpsc;
use std::thread;
use std::time::Duration;

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        let vals = vec!["a", "b", "c"];
        for val in vals {
            tx.send(val).unwrap();
            thread::sleep(Duration::from_secs(1));
        }
    });

    for received in rx {
        println!("Recebido: {}", received);
    }
}

O Receiver implementa Iterator, permitindo iterar sobre as mensagens até que o canal seja fechado. O loop termina quando todos os Sender são descartados.

Múltiplos produtores

Para ter múltiplos produtores, clonamos o Sender usando clone(). Cada clone pode ser enviado para uma thread diferente. O Receiver coleta todas as mensagens em ordem de chegada (não há garantia de ordem de envio entre produtores).

Exemplo com dois produtores:

use std::sync::mpsc;
use std::thread;

fn main() {
    let (tx, rx) = mpsc::channel();
    let tx1 = tx.clone();

    thread::spawn(move || {
        tx.send("mensagem do produtor 1").unwrap();
    });

    thread::spawn(move || {
        tx1.send("mensagem do produtor 2").unwrap();
    });

    for received in rx {
        println!("Recebido: {}", received);
    }
}

Cuidado: o Sender original deve ser movido para uma thread; o clone para a outra. O Receiver não é Sync, portanto só pode ser usado em uma thread.

Padrões

Padrões comuns incluem: pool de workers (várias threads produtoras enviam tarefas para uma thread consumidora), coleta de resultados (threads enviam resultados de volta), e canais síncronos com sync_channel para limitar o buffer.

Exemplo de pool de workers com canais:

use std::sync::mpsc;
use std::thread;

fn main() {
    let (tx, rx) = mpsc::channel();
    let num_workers = 4;

    for id in 0..num_workers {
        let tx_clone = tx.clone();
        thread::spawn(move || {
            tx_clone.send(format!("worker {} finalizou", id)).unwrap();
        });
    }

    drop(tx); // descarta o sender original

    for received in rx {
        println!("{}", received);
    }
}

Outro padrão é usar canais para comunicação bidirecional com múltiplos canais ou usando mpsc com um canal de resposta.

Boas práticas

Sempre trate erros de send e recv (use unwrap apenas em exemplos). Prefira for received in rx em vez de loop com recv para iterar até o fechamento. Lembre-se de descartar Senders que não serão usados para evitar deadlocks. Para mensagens complexas, use tipos que implementam Send.

Referências

Exercícios

  1. Crie um programa que use um canal para enviar 5 números inteiros de uma thread para a thread principal e os imprima.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    
    fn main() {
        let (tx, rx) = mpsc::channel();
    
        thread::spawn(move || {
            for i in 1..=5 {
                tx.send(i).unwrap();
            }
        });
    
        for received in rx {
            println!("Recebido: {}", received);
        }
    }
  2. Implemente um programa com dois produtores (threads) que enviam strings para um único consumidor. Use clone do sender.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    
    fn main() {
        let (tx, rx) = mpsc::channel();
        let tx1 = tx.clone();
    
        thread::spawn(move || {
            tx.send("Produtor 1").unwrap();
        });
    
        thread::spawn(move || {
            tx1.send("Produtor 2").unwrap();
        });
    
        for received in rx {
            println!("{}", received);
        }
    }
  3. Use try_recv para verificar se há mensagens sem bloquear. Crie um loop que tenta receber 3 vezes com intervalo de 1 segundo.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    use std::time::Duration;
    
    fn main() {
        let (tx, rx) = mpsc::channel();
    
        thread::spawn(move || {
            thread::sleep(Duration::from_secs(2));
            tx.send("mensagem").unwrap();
        });
    
        for _ in 0..3 {
            match rx.try_recv() {
                Ok(msg) => println!("Recebido: {}", msg),
                Err(_) => println!("Nenhuma mensagem"),
            }
            thread::sleep(Duration::from_secs(1));
        }
    }
  4. Crie um canal síncrono com capacidade 2. Envie 3 mensagens de uma thread e veja o comportamento.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    use std::time::Duration;
    
    fn main() {
        let (tx, rx) = mpsc::sync_channel(2);
    
        thread::spawn(move || {
            for i in 1..=3 {
                println!("Enviando {}", i);
                tx.send(i).unwrap();
                println!("Enviado {}", i);
            }
        });
    
        thread::sleep(Duration::from_secs(1));
        for received in rx {
            println!("Recebido: {}", received);
            thread::sleep(Duration::from_secs(1));
        }
    }
  5. Implemente um padrão de múltiplos produtores onde cada produtor envia seu ID e um número aleatório. O consumidor imprime a soma total.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    use rand::Rng;
    
    fn main() {
        let (tx, rx) = mpsc::channel();
        let num_produtores = 3;
    
        for id in 0..num_produtores {
            let tx_clone = tx.clone();
            thread::spawn(move || {
                let mut rng = rand::thread_rng();
                let valor: i32 = rng.gen_range(1..100);
                tx_clone.send((id, valor)).unwrap();
            });
        }
    
        drop(tx);
        let mut soma = 0;
        for (id, valor) in rx {
            println!("Produtor {} enviou {}", id, valor);
            soma += valor;
        }
        println!("Soma total: {}", soma);
    }