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.
<!-- 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:
- **
movehands each worker its own world.** The closure takesspec,inbox_rxand a clonedevents_txby value. Nothing is shared, so nothing needs a lock. - **
sendtakes ownership, so fan-out costs a clone.** OneMessagecannot 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. - 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:
| Channel | Senders | Where each dies |
|---|---|---|
| each inbox | one inbox_tx, held in inboxes | drop(inboxes) after the publish loop |
| events | one clone per worker, plus the original in main | worker 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
- Holding the original sender in main while iterating the receiver in main.
- Cloning a sender "just in case" and storing it in a struct that outlives the loop.
- Sharing one inbox between two workers and expecting fan-out. That is Kleppmann's load balancing: each message goes to one of them.
- 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
- Duha #18: threads and message passing (this lesson applies it)
- Ops: Pub/Sub subscription expiry census (same trio)
- Cert: PCA Pub/Sub delivery, ordering, seek (same trio)
- Prior Dev: k8s-openapi drain preflight
- — closure rule and Listing 16-11