Hedronite · Dev Lesson · Polyglot-Dev / Rust · Mon 2026-09-28

Rust mpsc fan-out — the sender that keeps the loop open

A for loop over a receiver ends when the last sender is gone. Count your senders.

Lesson Class: Dev (Rust concurrency, std only)
Topic: T1 Systems core
Crates: none (std::thread, std::sync::mpsc)
Lag rule: Duha #18 ch16-01/02; no Mutex/Arc, no async
Verification: cargo test 3 passed · clippy clean · hang reproduced (exit 124)
Paired Ops: Pub/Sub subscription expiry census
Sender Count
List every live Sender and where it dies.
Clone per inbox
send takes ownership; fan-out costs one clone each.
Empty vs Disconnected
try_recv tells you which bug you have.
Drop the sender you are not using.

<!-- hal:authoritative:yaml -->

*A for loop over a receiver ends when the last sender is gone. Count your senders, because the compiler will not.*

§I. Frame

The Ops lesson in this trio reads Pub/Sub subscriptions from the API. This lesson builds the shape of Pub/Sub delivery offline, with nothing but std::thread and std::sync::mpsc from this morning's Duha session: one topic, several subscriptions, each subscription a thread with its own inbox.

Two Pub/Sub rules set the behavior. Every subscription gets every message (Kleppmann's fan-out, DDIA ch11, printed p.430). And a subscription with a filter delivers only matching messages; the Pub/Sub filter page says the service "automatically acknowledges the messages that don't match the filter."

The crate pubsub-fanout sits in this bundle under pkg/. It has no dependencies.

§II. The shape

pub fn fan_out(subs: Vec<SubSpec>, messages: Vec<Message>) -> Vec<Event> {
    let (events_tx, events_rx) = mpsc::channel::<Event>();
    let mut inboxes = Vec::new();
    let mut handles = Vec::new();

    for spec in subs {
        let (inbox_tx, inbox_rx) = mpsc::channel::<Message>();
        let events_tx = events_tx.clone();
        handles.push(thread::spawn(move || {
            for m in inbox_rx {
                let sub = spec.name.clone();
                let event = match &spec.filter {
                    Some(f) if !f.matches(&m) => Event::AutoAcked { sub, id: m.id },
                    _ => Event::Delivered { sub, id: m.id },
                };
                events_tx.send(event).unwrap();
            }
        }));
        inboxes.push(inbox_tx);
    }

    // Main still owns the original events sender. Drop it, or the collect below never ends.
    drop(events_tx);

    for m in messages {
        for inbox in &inboxes {
            inbox.send(m.clone()).unwrap();
        }
    }
    // Close every inbox so each worker's `for m in inbox_rx` loop can finish.
    drop(inboxes);

    let events: Vec<Event> = events_rx.iter().collect();
    for h in handles {
        h.join().unwrap();
    }
    events
}

Three ch16 moves carry the whole function:

  1. **move hands each worker its own world.** The closure takes spec, inbox_rx and a cloned events_tx by value. Nothing is shared, so nothing needs a lock.
  2. **send takes ownership, so fan-out costs a clone.** One Message cannot go into three channels. m.clone() per inbox is the honest price of "every subscription gets its own copy," which mirrors Pub/Sub, where each subscription keeps its own backlog.
  3. One sender per inbox keeps order. Each inbox has exactly one producer (main), and a channel delivers in send order, so each subscription sees messages in publish order. The first test asserts it.

§III. Sender Count

Sender Count (named technique). Before any for x in rx, list every live Sender for that channel and say where each one dies.

TRPL ch16-02 gives the rule: a channel "is said to be closed if either the transmitter or receiver half is dropped," and when rx is used as an iterator, "when the channel is closed, iteration will end." With cloned senders, "the transmitter half" means all of them. Listing 16-11 clones tx for the first thread and moves the original into the second, so both copies die when their threads finish.

fan_out has two channels to count:

ChannelSendersWhere each dies
each inboxone inbox_tx, held in inboxesdrop(inboxes) after the publish loop
eventsone clone per worker, plus the original in mainworker clones at thread exit; the original at drop(events_tx)

Remove the drop(events_tx) line and the count is off by one. Every worker finishes and drops its clone, but main still holds the original, so events_rx.iter() waits for a message that can never come. I ran exactly that on the box: the edited copy compiled with no warning and timeout 5 killed it with exit 124.

The third test proves the rule without hanging, using try_recv, which returns at once:

let (tx, rx) = mpsc::channel::<u32>();
let worker_tx = tx.clone();
thread::spawn(move || {
    worker_tx.send(1).unwrap();
    worker_tx.send(2).unwrap();
})
.join()
.unwrap();

assert_eq!(rx.try_recv(), Ok(1));
assert_eq!(rx.try_recv(), Ok(2));
// The worker is gone, but `tx` is alive: a `for x in rx` loop would block here forever.
assert_eq!(rx.try_recv(), Err(TryRecvError::Empty));
drop(tx);
assert_eq!(rx.try_recv(), Err(TryRecvError::Disconnected));

Empty means "no message now, a sender still exists." Disconnected means "no sender left." The difference is the whole bug.

§IV. The filter, kept honest

The Pub/Sub filter language has :, =, !=, hasPrefix and NOT, among others. This crate supports one form, attributes.KEY = "VALUE", and rejects everything else with an error instead of guessing. It also enforces the documented 256-byte limit:

if src.len() > MAX_FILTER_BYTES {
    return Err(format!("filter is {} bytes; the limit is {MAX_FILTER_BYTES}", src.len()));
}

A parser that silently treated != as = would invert a subscription, and in real Pub/Sub the filter cannot be changed after creation.

§V. Run it

On the lab Mac (cargo 1.96.0) with the bundled fixture:

$ cargo test
test tests::filter_parser_accepts_one_form_and_rejects_the_rest ... ok
test tests::a_held_sender_keeps_the_channel_open ... ok
test tests::every_subscription_sees_every_message_in_publish_order ... ok
test result: ok. 3 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out

$ cargo run -q -- fixture/messages.txt fixture/subscriptions.txt
published=6 subscriptions=3
audit-sink delivered=[101, 102, 103, 104, 105, 106] auto_acked=[]
eu-region delivered=[101, 103, 106] auto_acked=[102, 104, 105]
refunds delivered=[103, 105] auto_acked=[101, 102, 104, 106]

Every message appears once per subscription, split between delivered and auto-acked. A != filter in bad-subscriptions.txt exits 65 with "unsupported operator"; a wrong argument count exits 64. cargo clippy --all-targets -- -D warnings is clean.

§VI. Traps

  1. Holding the original sender in main while iterating the receiver in main.
  2. Cloning a sender "just in case" and storing it in a struct that outlives the loop.
  3. Sharing one inbox between two workers and expecting fan-out. That is Kleppmann's load balancing: each message goes to one of them.
  4. Treating a clean compile as proof the channels close. Closure is a runtime fact.

§VII. Close instruction

Add a fourth subscription whose worker forwards every delivered message to a second channel read by a new thread. Write its row in the Sender Count table first, then add the drop that makes the new thread's loop end, and prove it with a try_recv test.

Related