Broadcasting Events with Tokio Broadcast Channels Permalink to this section
Part of Rust SSE with Axum and Tokio, under Backend Stream Generation & Connection Management.
tokio::sync::broadcast is the natural fan-out primitive for a Rust SSE server: one sender, any number of receivers, every receiver sees every message. It is also easy to use in a way that silently drops messages, copies every payload once per subscriber, or wakes ten thousand tasks for a message meant for one user. This guide builds a broadcast fan-out that sizes its buffer deliberately, reports lag honestly to clients, serialises once, and scopes delivery so only interested connections wake up.
Symptom & Developer Intent Permalink to this section
- Some clients occasionally miss events, with no error anywhere in the logs.
- CPU climbs linearly with subscribers even though the message rate is constant.
- Memory spikes during bursts in proportion to connections times payload size.
sendreturns an error at startup or during quiet periods and the producer panics onunwrap().- A message addressed to one user measurably delays delivery to everyone else.
The intent is a fan-out where every subscriber receives every message or is told explicitly that it missed some, where the cost of a message is paid once, and where per-user messages only touch that user’s connections.
Root Cause Analysis Permalink to this section
A broadcast channel is a fixed-size ring buffer shared by all receivers. send writes into the next slot and wakes every receiver; each receiver keeps its own read position. When a receiver is more than capacity messages behind, the slot it wanted has been overwritten, and its next recv returns RecvError::Lagged(n) before continuing from the oldest message still available.
The silent-miss symptom comes from code that ignores Lagged: BroadcastStream yields it as an Err, and a filter_map(|r| r.ok()) throws it away. The panic comes from send returning Err when there are zero receivers — which is normal on an idle server, not a failure. The CPU and memory symptoms come from putting large owned values in the channel: broadcast requires T: Clone and clones the value for each receiver on recv.
Step-by-Step Resolution Permalink to this section
Step 1 — Size the capacity from your burst profile Permalink to this section
Capacity is shared, so it is cheap to be generous. A useful estimate: peak messages per second × the longest stall you want a receiver to survive without lagging.
// 500 msg/s peaks, tolerate a 4-second stall (a phone switching networks): 2,000 → round up.
const CAPACITY: usize = 2048;
let (tx, _rx) = broadcast::channel::<Frame>(CAPACITY);
Memory held by the channel is at most CAPACITY messages, independent of the number of receivers.
Step 2 — Pre-serialise frames and share them Permalink to this section
use bytes::Bytes;
#[derive(Clone)]
struct Frame { seq: u64, event: &'static str, json: Arc<str> } // cheap to clone
fn publish(tx: &broadcast::Sender<Frame>, seq: u64, quote: &Quote) {
let json: Arc<str> = serde_json::to_string(quote).expect("serialisable").into();
// Err means "no receivers right now" — normal on an idle server, not a failure.
let _ = tx.send(Frame { seq, event: "price", json });
}
Each receiver clones an Arc (an atomic increment), not a Quote. Serialisation happens exactly once per message.
Step 3 — Turn Lagged into an explicit resync event Permalink to this section
let stream = BroadcastStream::new(tx.subscribe()).map(|r| -> Result<Event, Infallible> {
Ok(match r {
Ok(f) => Event::default().id(f.seq.to_string()).event(f.event).data(&*f.json),
Err(BroadcastStreamRecvError::Lagged(n)) => {
metrics::counter!("sse_lagged_total").increment(1);
Event::default().event("resync").data(format!("{{\"missed\":{n}}}"))
}
})
});
On the client, a resync event means “refetch the current state, then continue”:
es.addEventListener('resync', async () => {
const snapshot = await (await fetch('/api/prices/snapshot')).json();
replaceAll(snapshot); // the stream continues; no reconnect needed
});
Step 4 — Scope delivery with per-key channels Permalink to this section
For events addressed to a user or room, a global broadcast wakes every connection. Hold one small broadcast per key instead:
#[derive(Clone, Default)]
struct Topics(Arc<DashMap<String, broadcast::Sender<Frame>>>);
impl Topics {
fn subscribe(&self, key: &str) -> broadcast::Receiver<Frame> {
self.0.entry(key.to_owned()).or_insert_with(|| broadcast::channel(128).0).subscribe()
}
fn send(&self, key: &str, f: Frame) {
let gone = match self.0.get(key) {
Some(tx) => tx.send(f).is_err(), // Err: no receivers remain
None => false,
};
if gone { self.0.remove_if(key, |_, tx| tx.receiver_count() == 0); }
}
}
A connection for user 42 subscribes to "user:42" and, if it also shows global announcements, merges that receiver with one from the global channel using futures::stream::select.
Step 5 — Feed the channel from a broker in multi-instance deployments Permalink to this section
Each process runs one broker subscriber task that republishes into its local channels. Redis, NATS and Kafka clients for Tokio all expose async streams, so the relay is a loop:
tokio::spawn(async move {
let mut sub = redis_client.get_async_pubsub().await?;
sub.subscribe("prices").await?;
let mut msgs = sub.on_message();
while let Some(m) = msgs.next().await {
let payload: String = m.get_payload()?;
let _ = tx.send(Frame { seq: next_seq(), event: "price", json: payload.into() });
}
Ok::<_, anyhow::Error>(())
});
Validation & Monitoring Permalink to this section
# Force lag: hold one very slow reader while publishing a burst, then look for resync.
curl -sN --limit-rate 200 localhost:8080/api/prices/stream | grep -m1 -A1 'event: resync'
# Confirm idle-server publishes do not panic: restart with zero clients and watch logs.
A unit test can assert the lag behaviour deterministically with a tiny capacity:
#[tokio::test]
async fn lagging_receiver_is_told() {
let (tx, rx) = broadcast::channel::<u32>(2);
for i in 0..5 { tx.send(i).unwrap(); } // overwrite past capacity
let mut s = BroadcastStream::new(rx);
assert!(matches!(s.next().await, Some(Err(BroadcastStreamRecvError::Lagged(3)))));
assert_eq!(s.next().await.unwrap().unwrap(), 3); // continues from oldest kept
}
Export sse_lagged_total, open-stream count, and messages published per second. Lag events that correlate with bursts suggest raising capacity; lag concentrated on a few clients is simply slow networks, and the resync handles it.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Why does send fail when nobody is connected?
broadcast::Sender::send returns an error when there are no active receivers, because the message cannot be delivered to anyone. On an SSE server that is a normal state; discard the error.
Can a lagged receiver recover the missed messages?
Not from the channel — they were overwritten. Recover through your application: send a resync event so the client refetches state, or replay from a durable store using the last id it received.
Is broadcast fair to fast receivers when one is slow?
Yes. Each receiver has its own position, so a slow one never blocks the sender or other receivers. It only affects itself, by lagging.
When should I use watch instead of broadcast?
When only the latest value matters, such as a status or score. A watch receiver always sees the newest value and never lags, which is exactly conflation.