Rust SSE with Axum and Tokio Permalink to this section
Part of Backend Stream Generation & Connection Management.
Rust is an excellent fit for a streaming tier. A Tokio task parked on a channel costs a few hundred bytes, there is no garbage collector to pause every stream at once, and the type system forces you to decide what happens when a subscriber falls behind — a decision other runtimes let you postpone until production makes it for you. Axum, built on Tokio and Hyper, has first-class Server-Sent Events support in axum::response::sse: an Sse response wraps any Stream of events and a KeepAlive configuration adds heartbeat comments. This guide covers how that machinery maps onto the protocol, a complete fan-out server with tokio::sync::broadcast, the behaviour of lagging receivers, disconnect cleanup through Drop, Rust clients, proxy and deployment details, and how far one process scales. It targets axum 0.8 and Tokio 1.x.
How It Works Permalink to this section
An axum handler returns Sse<S> where S is a Stream<Item = Result<Event, E>>. Axum sets Content-Type: text/event-stream and Cache-Control: no-cache, polls the stream, serialises each Event into its wire form and hands the bytes to Hyper, which sends them immediately as chunks (HTTP/1.1) or DATA frames (HTTP/2). When the client disconnects, Hyper drops the response body, which drops your stream — and with it anything the stream owns.
The Event builder mirrors the wire fields:
use axum::response::sse::Event;
let ev = Event::default()
.id("42") // id: 42
.event("price") // event: price
.retry(std::time::Duration::from_secs(3)) // retry: 3000
.json_data("e)?; // data: {"sym":"ACME",...} (serde_json, one line)
json_data serialises with serde and fails rather than producing invalid framing; data accepts a string and splits embedded newlines into multiple data: lines for you, which keeps multi-line text correct per the event stream format. Event::default().comment("hb") produces a comment line.
HTTP/1.1 200 OK
content-type: text/event-stream
cache-control: no-cache
retry: 3000
id: 42
event: price
data: {"sym":"ACME","bid":101.2}
:
The final : is axum’s default keep-alive comment: an empty comment line, sent every 15 seconds when the stream has produced nothing.
Server-Side Implementation Permalink to this section
The canonical fan-out uses a broadcast channel in shared state. Every connection subscribes, getting its own receiver that reads from a shared ring buffer.
// main.rs — axum 0.8, tokio 1, tokio-stream 0.1 (feature "sync"), serde
use axum::{extract::State, http::HeaderMap, response::sse::{Event, KeepAlive, Sse}, routing::get, Router};
use futures::stream::{Stream, StreamExt};
use serde::Serialize;
use std::{convert::Infallible, sync::Arc, time::Duration};
use tokio::sync::broadcast;
use tokio_stream::wrappers::{errors::BroadcastStreamRecvError, BroadcastStream};
#[derive(Clone, Serialize)]
struct Quote { seq: u64, sym: String, bid: f64 }
#[derive(Clone)]
struct AppState { tx: broadcast::Sender<Arc<Quote>> }
async fn prices(State(st): State<AppState>, headers: HeaderMap)
-> Sse<impl Stream<Item = Result<Event, Infallible>>>
{
let _last_id = headers.get("last-event-id").and_then(|v| v.to_str().ok()).map(str::to_owned);
let rx = st.tx.subscribe(); // this connection's receiver
let stream = BroadcastStream::new(rx).filter_map(|msg| async move {
match msg {
Ok(q) => Some(Ok(Event::default()
.id(q.seq.to_string())
.event("price")
.json_data(&*q)
.unwrap_or_else(|_| Event::default().comment("encode error")))),
// The receiver fell more than `capacity` messages behind and skipped ahead.
Err(BroadcastStreamRecvError::Lagged(n)) => Some(Ok(Event::default()
.event("resync")
.data(n.to_string()))),
}
});
Sse::new(stream).keep_alive(KeepAlive::new().interval(Duration::from_secs(15)).text("hb"))
}
#[tokio::main]
async fn main() {
let (tx, _) = broadcast::channel::<Arc<Quote>>(1024); // capacity = lag tolerance
let state = AppState { tx: tx.clone() };
tokio::spawn(feed(tx)); // producer task
let app = Router::new().route("/api/prices/stream", get(prices)).with_state(state);
let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await.unwrap();
axum::serve(listener, app).await.unwrap();
}
Three design choices are embedded here. Messages are wrapped in Arc so the broadcast clones a pointer per receiver, not the payload. The channel capacity bounds memory: the broadcast holds at most 1,024 messages regardless of how many subscribers exist or how slow they are. And a receiver that falls more than 1,024 messages behind gets Lagged(n) — it has missed n messages and continues from the oldest still buffered — which the handler turns into a resync event so the client can refetch state instead of silently showing a gap. Broadcasting events with tokio broadcast channels goes deeper into capacity, lag and per-user filtering.
Choosing the right Tokio channel Permalink to this section
Tokio offers several channel types, and each corresponds to a different stream semantics. Picking by data shape avoids most of the subtle bugs:
| Channel | Semantics | Slow receiver | Use it for |
|---|---|---|---|
broadcast |
every receiver gets every message | lags, then skips with Lagged(n) |
shared feeds: prices, dashboards, public events |
watch |
receivers see only the latest value | never lags; intermediate values vanish | single current state: a config, a status, a score |
mpsc per connection |
one producer path to one connection | sender waits (send) or fails (try_send) |
per-user notifications that must not drop |
broadcast per key |
broadcast sharded by user, room or topic | as broadcast, scoped to the key |
many small audiences: documents, chat rooms |
watch deserves more use than it gets. A stream of a single state value — a deployment’s status, a match score, a job’s progress — is exactly a watch channel: a receiver that was busy simply sees the newest value when it next polls, which is conflation for free. tokio_stream::wrappers::WatchStream turns it into a stream for Sse::new in one line.
Routing to one user or room Permalink to this section
A single broadcast wakes every connection for every message. When each message is for one user or one room, keep a map from key to a small per-key sender, created on first subscribe and removed when its last receiver goes away:
use dashmap::DashMap;
#[derive(Clone, Default)]
struct Rooms(Arc<DashMap<String, broadcast::Sender<Arc<str>>>>);
impl Rooms {
fn subscribe(&self, room: &str) -> broadcast::Receiver<Arc<str>> {
self.0.entry(room.to_owned())
.or_insert_with(|| broadcast::channel(256).0)
.subscribe()
}
fn publish(&self, room: &str, frame: Arc<str>) {
if let Some(tx) = self.0.get(room) {
if tx.send(frame).is_err() { // no receivers left
drop(tx);
self.0.remove_if(room, |_, tx| tx.receiver_count() == 0);
}
}
}
}
Publishing to a room with no listeners costs one map lookup, and rooms with no receivers are removed lazily, so the map’s size tracks active audiences rather than every room ever seen.
Replay from Last-Event-ID Permalink to this section
The broadcast channel is a live feed, not a history. For resumable streams, keep recent events in a bounded store and chain a replay stream before the live one — subscribing to the broadcast first, so nothing published during the replay query is lost:
let rx = st.tx.subscribe(); // 1. subscribe first
let after: u64 = last_id.and_then(|s| s.parse().ok()).unwrap_or(0);
let history = st.recent.after(after).await; // 2. then query history
let high = history.last().map(|q| q.seq).unwrap_or(after);
let replay = futures::stream::iter(history.into_iter().map(|q| Ok(to_event(&q))));
let live = BroadcastStream::new(rx).filter_map(move |m| async move {
match m {
Ok(q) if q.seq > high => Some(Ok(to_event(&q))), // 3. skip overlap
Ok(_) => None,
Err(BroadcastStreamRecvError::Lagged(n)) => Some(Ok(resync(n))),
}
});
Sse::new(replay.chain(live)).keep_alive(KeepAlive::default())
Cleanup is Drop Permalink to this section
There is no close handler to remember. When the client disconnects, the stream — including the BroadcastStream, which owns the receiver — is dropped, and the receiver unsubscribes as part of its Drop. Per-connection resources you create yourself (a registry entry, a metrics gauge, a presence lease) should live in a guard value moved into the stream so they are released the same way; keep-alive and disconnects in axum SSE shows the pattern and its one caveat: the drop happens only when Hyper notices the disconnect, which for a silent peer requires a write — hence keep-alive.
Client-Side Consumption Permalink to this section
Browsers use EventSource as usual. For Rust clients, reqwest with the stream feature exposes the body as a byte stream; the eventsource-stream crate turns that into parsed events, and reqwest-eventsource adds reconnection with Last-Event-ID:
use eventsource_stream::Eventsource;
use futures::StreamExt;
let mut last_id: Option<String> = None;
loop {
let mut req = reqwest::Client::new()
.get("https://prices.example.com/api/prices/stream")
.header("accept", "text/event-stream");
if let Some(id) = &last_id { req = req.header("last-event-id", id); }
match req.send().await {
Ok(res) if res.status().is_success() => {
let mut events = res.bytes_stream().eventsource();
while let Some(Ok(ev)) = events.next().await {
if !ev.id.is_empty() { last_id = Some(ev.id.clone()); }
handle(&ev.event, &ev.data);
}
}
_ => {}
}
tokio::time::sleep(backoff.next()).await; // reconnect with backoff
}
Set no overall request timeout on the client, or the stream is cut at the timeout. More client runtimes are covered in non-browser SSE clients.
Edge Cases & Network Interference Permalink to this section
- Compression layers.
tower-http’sCompressionLayercompresses responses by content type. Excludetext/event-streamwith a predicate, or events sit in the compressor. - Timeout layers. A global
TimeoutLayerends every request after its duration, including streams. Apply it per route, not around the stream router. - Reverse proxies. nginx buffers by default. Add
X-Accel-Buffering: nowith aSetResponseHeaderLayeron the stream route, or configure the proxy as in proxy and CDN configuration for SSE. - HTTP/2 and connection limits. Behind a TLS terminator speaking HTTP/2 to browsers, many streams share one connection and the six-per-origin limit disappears; directly served HTTP/1.1 still has it.
- Graceful shutdown.
axum::serve(...).with_graceful_shutdown(signal)waits for in-flight requests — which, for infinite streams, is forever. Combine it with a shutdownwatchchannel that the stream selects on, so streams end promptly and clients reconnect elsewhere.
// End every stream when shutdown is signalled, so graceful shutdown can complete.
let mut shutdown = st.shutdown.clone(); // tokio::sync::watch::Receiver<bool>
let live = live.take_until(async move { let _ = shutdown.changed().await; });
Mitigation checklist:
Performance & Scale Considerations Permalink to this section
A Tokio-based SSE server’s cost per idle connection is small and predictable: a task’s state machine, the receiver (a cursor into the shared ring), Hyper’s connection buffers, and the socket.
Broadcast cost is dominated by serialisation and syscalls. Serialise once per message, not once per subscriber: build the event’s JSON before sending into the channel (send an Arc<str> or pre-built Bytes), and have each connection wrap the shared string in Event::default().data(...). With a multi-threaded runtime, Tokio spreads connection tasks across cores automatically.
A useful rule of thumb for the broadcast path: the producer should do work proportional to messages, and each connection should do work proportional to the bytes it writes — nothing proportional to messages times connections except the unavoidable socket writes. Filtering per connection breaks that rule when most messages are irrelevant to most connections, which is exactly when per-room senders pay off. Measure with a load test that holds realistic connection counts and publishes at your peak rate; the k6 load-testing guide shows how to hold thousands of streams and measure delivery latency per event rather than requests per second.
Operationally, raise the file descriptor limit (LimitNOFILE in systemd) above your target connection count, and size the broadcast capacity for the burstiest second you expect times a comfortable lag margin. For multiple instances, feed each process’s broadcast channel from a broker subscriber task — Redis, NATS or Kafka — as in Redis pub/sub fan-out.
Validation & Debugging Permalink to this section
# Frames and keep-alives arrive on schedule.
curl -sN localhost:8080/api/prices/stream | while IFS= read -r l; do echo "$(date +%T) $l"; done
# Many connections: hold 10k streams and watch RSS and task count.
for i in $(seq 1 10000); do curl -sN localhost:8080/api/prices/stream > /dev/null & done
ps -o rss= -p $(pgrep my-sse-server)
tokio-console shows live tasks, how long each has been idle, and wakeup counts; with the tracing feature, each connection’s task is visible, which makes leaks (tasks that outlive their connections) obvious. Instrument with the metrics crate: a gauge of open streams incremented in the connection guard’s constructor and decremented in its Drop, and a counter of Lagged events.
#[tokio::test]
async fn emits_price_events() {
let (tx, _) = broadcast::channel(16);
let app = Router::new().route("/s", get(prices)).with_state(AppState { tx: tx.clone() });
let server = axum_test::TestServer::new(app).unwrap();
let res = server.get("/s").await; // or drive with reqwest against a bound port
assert_eq!(res.header("content-type"), "text/event-stream");
}
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
How does axum know the client disconnected?
Hyper detects the closed connection when it reads the FIN or when a write fails, then drops the response body and your stream with it. Keep-alive comments ensure there is a write to fail on idle streams.
What broadcast capacity should I choose?
Enough to absorb your largest burst plus the delay of a moderately slow client — often 256 to 4,096 messages. Capacity bounds memory for the whole channel, not per subscriber, so it can be generous.
Is tokio::sync::broadcast suitable for per-user streams?
For a few thousand users, one broadcast with a filter per connection is fine. Beyond that, keep a map from user to a small broadcast or mpsc sender, so publishing to one user does not wake every connection.
Can I serve SSE from actix-web or warp instead?
Yes. Both support streaming bodies with text/event-stream, and warp has an sse module. The design — bounded fan-out, lag handling, keep-alive and cleanup on drop — is identical across frameworks.
How do I authenticate an axum SSE endpoint?
Use an extractor that validates a session cookie or a short-lived query token, since browsers' EventSource cannot send an Authorization header. Reject with 401 before returning the Sse response; a non-200 status makes EventSource stop instead of retrying.
Should events be serialised before or after the channel?
Before. Serialise once in the producer and send an Arc<str> or Bytes through the channel, so each connection only wraps the shared string in an Event. Serialising per connection repeats identical work for every subscriber.
Do I need a separate thread per connection?
No. Each connection is a lightweight async task scheduled on Tokio's worker threads, typically one per core. Hundreds of thousands of idle streams fit on a handful of threads.