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.

From a published event to the socket in axum Flow from a producer sending into a tokio broadcast channel, through a per-connection BroadcastStream, a map to axum Event, the Sse response with KeepAlive, and Hyper writing to the socket. From a published event to the socket in axum Producer tx.send(msg) send broadcast ring buffer recv BroadcastStream per connection map Sse + KeepAlive Event frames poll Hyper chunk / DATA
Each connection owns its receiver and its stream adapter. Dropping the response on disconnect drops them too, which is the whole cleanup story in Rust.

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(&quote)?;                // 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.

A broadcast receiver under load State diagram of a tokio broadcast receiver moving between keeping up, falling behind within capacity, and lagged, where it skips to the oldest buffered message and the handler emits a resync event. A broadcast receiver under load CURRENT reads promptly BEHIND < capacity queued LAGGED skipped n RESYNC SENT client refetches slow socket > capacity Err(Lagged) catches up
Falling behind within capacity costs nothing extra. Exceeding it drops messages for that receiver only, and the Lagged error tells you exactly how many.

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.

Per-room senders instead of one global broadcast A publisher looks up a room in a concurrent map and sends into that room's broadcast sender, which reaches only the connections subscribed to that room. Per-room senders instead of one global broadcast publish(room 42) pre-serialised frame DashMap lookup room → Sender get(42) Connection A room 42 Connection B room 42 Connection C room 42 send
Only the listeners of room 42 wake up. With one global broadcast, every connection on the process would be polled and would discard the message.

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’s CompressionLayer compresses responses by content type. Exclude text/event-stream with a predicate, or events sit in the compressor.
  • Timeout layers. A global TimeoutLayer ends 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: no with a SetResponseHeaderLayer on 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 shutdown watch channel 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.

Resident memory for 100,000 idle SSE connections Bar chart comparing approximate resident memory for one hundred thousand idle streams across axum on Tokio, Go net/http, and Node.js. Resident memory for 100,000 idle SSE connections axum + Tokio ~0.7 GB Go net/http ~1.5 GB Node.js http ~2.7 GB approximate resident set size, 100k idle streams, heartbeats only
Idle connections are cheap in all three. Rust's advantage shows up most under broadcast load, where there is no GC and no per-subscriber payload copy.

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.

Deep Dives