Nesta aula, vamos mergulhar nos canais (channels) em Rust, uma ferramenta essencial para comunicação segura entre threads. Canais seguem o modelo de passagem de mensagens, onde um produtor envia dados e um consumidor os recebe, evitando problemas de concorrência como data races. O módulo std::sync::mpsc implementa canais com múltiplos produtores e único consumidor (multiple producer, single consumer), sendo a base para a maioria dos casos de uso.

Entender canais é fundamental para escrever programas concorrentes em Rust. Eles permitem que threads troquem informações de forma controlada, sem compartilhamento de estado, alinhando-se com a filosofia de segurança da linguagem. Veremos desde o básico até padrões mais avançados.

mpsc

O módulo std::sync::mpsc fornece canais assíncronos (com buffer) e síncronos (sem buffer, via SyncSender). A função channel() retorna um par (Sender, Receiver). O Sender pode ser clonado para múltiplos produtores, enquanto Receiver é único. O canal usa um buffer interno (atualmente com capacidade ilimitada) para armazenar mensagens até que sejam lidas.

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

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

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

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

O método send retorna Result<(), SendError>; se o receptor foi descartado, retorna erro. recv bloqueia até receber uma mensagem ou o canal fechar (todos os produtores descartados), retornando Err se não houver mais produtores.

send e recv

send envia um valor pelo canal. Se o receptor já foi dropado, retorna SendError contendo o valor de volta. recv bloqueia a thread atual até que uma mensagem esteja disponível ou o canal seja fechado. Há também try_recv, que não bloqueia e retorna TryRecvError se vazio ou fechado.

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(1));
        tx.send("mensagem").unwrap();
    });

    match rx.recv_timeout(Duration::from_secs(2)) {
        Ok(msg) => println!("Recebido: {}", msg),
        Err(e) => println!("Erro: {:?}", e),
    }
}

O método recv_timeout bloqueia por um tempo máximo. Também é possível iterar sobre o receptor com um loop for, que termina quando o canal é fechado.

Múltiplos produtores

Para ter múltiplos produtores, clonamos o Sender. Cada clone pode ser movido para uma thread diferente. O Receiver permanece único e coleta todas as mensagens na ordem em que foram enviadas (não necessariamente FIFO estrito devido ao buffer, mas geralmente FIFO).

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

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

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

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

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

Note que o loop for sobre o receptor termina quando todos os Senders (incluindo clones) forem dropados. Isso é útil para processar um fluxo finito de mensagens.

Padrões

Padrões comuns com canais incluem: roteamento de mensagens (usando enum para diferentes tipos), sincronização (enviar sinal de conclusão), e balanceamento de carga (vários workers consumindo de um canal).

enum Mensagem {
    Trabalho(String),
    Sair,
}

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

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

    let worker = thread::spawn(move || {
        loop {
            match rx.recv().unwrap() {
                Mensagem::Trabalho(t) => println!("Processando: {}", t),
                Mensagem::Sair => break,
            }
        }
    });

    tx.send(Mensagem::Trabalho("tarefa1".into())).unwrap();
    tx.send(Mensagem::Trabalho("tarefa2".into())).unwrap();
    tx.send(Mensagem::Sair).unwrap();

    worker.join().unwrap();
}

Outro padrão é usar canais para sincronizar o início de threads (barreira). Enviar um valor vazio como sinal.

Boas práticas

  • Sempre trate erros de send e recv; o canal pode fechar inesperadamente.
  • Use try_recv em loops não bloqueantes.
  • Clone o Sender apenas quando necessário; cada clone aumenta a contagem de produtores.
  • Para comunicação bidirecional, crie dois canais.

Referências

Exercícios

  1. Crie um programa que usa um canal para enviar 10 números inteiros de uma thread para a main. A thread produtora deve enviar os números de 1 a 10, e a main deve imprimi-los.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    
    fn main() {
        let (tx, rx) = mpsc::channel();
    
        thread::spawn(move || {
            for i in 1..=10 {
                tx.send(i).unwrap();
            }
        });
    
        for received in rx {
            println!("{} ", received);
        }
    }
  2. Modifique o exercício anterior para usar dois produtores: um envia números pares e outro ímpares. O consumidor (main) deve imprimir todos.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    
    fn main() {
        let (tx, rx) = mpsc::channel();
        let tx1 = tx.clone();
    
        thread::spawn(move || {
            for i in (2..=10).step_by(2) {
                tx.send(i).unwrap();
            }
        });
    
        thread::spawn(move || {
            for i in (1..=9).step_by(2) {
                tx1.send(i).unwrap();
            }
        });
    
        for received in rx {
            println!("{} ", received);
        }
    }
  3. Escreva um programa que usa um enum de mensagens para enviar comandos: Texto(String), Numero(i32) e Sair. Uma worker thread processa as mensagens até receber Sair.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    
    enum Mensagem {
        Texto(String),
        Numero(i32),
        Sair,
    }
    
    fn main() {
        let (tx, rx) = mpsc::channel::<Mensagem>();
    
        let worker = thread::spawn(move || {
            loop {
                match rx.recv().unwrap() {
                    Mensagem::Texto(t) => println!("Texto: {}", t),
                    Mensagem::Numero(n) => println!("Número: {}", n),
                    Mensagem::Sair => break,
                }
            }
        });
    
        tx.send(Mensagem::Texto("Olá".into())).unwrap();
        tx.send(Mensagem::Numero(42)).unwrap();
        tx.send(Mensagem::Sair).unwrap();
    
        worker.join().unwrap();
    }
  4. Use recv_timeout para ler mensagens com timeout de 1 segundo. Se não houver mensagem, imprima "timeout".

    ✓ 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("depois de 2s").unwrap();
        });
    
        loop {
            match rx.recv_timeout(Duration::from_secs(1)) {
                Ok(msg) => {
                    println!("Recebido: {}", msg);
                    break;
                }
                Err(mpsc::RecvTimeoutError::Timeout) => println!("timeout"),
                Err(mpsc::RecvTimeoutError::Disconnected) => {
                    println!("canal fechado");
                    break;
                }
            }
        }
    }
  5. Crie um programa que simula um pipeline: thread A gera números, envia para thread B que dobra e envia para thread C que imprime. Use dois canais.

    ✓ Resposta:
    use std::sync::mpsc;
    use std::thread;
    
    fn main() {
        let (tx_a, rx_a) = mpsc::channel();
        let (tx_b, rx_b) = mpsc::channel();
    
        // Thread A: gera números
        let produtor = thread::spawn(move || {
            for i in 1..=5 {
                tx_a.send(i).unwrap();
            }
        });
    
        // Thread B: dobra
        let processador = thread::spawn(move || {
            for num in rx_a {
                tx_b.send(num * 2).unwrap();
            }
        });
    
        // Thread C: imprime
        let consumidor = thread::spawn(move || {
            for num in rx_b {
                println!("Resultado: {}", num);
            }
        });
    
        produtor.join().unwrap();
        processador.join().unwrap();
        consumidor.join().unwrap();
    }