M8: Real-Time Design
One real-time layer for a product that had five, in four rounds: the transports and what the infrastructure lets through, a subscription model with sequence numbers, replay and reconnection, backpressure as a policy per kind of data, and presence and delivery guarantees the UI makes visible.
The brief and the questions
"Our product has live bits everywhere: the collaborative editor, the notification badge, a chat per project, progress on long jobs, and a dashboard that should update by itself. Right now each team did its own thing and it is a mess: five sockets per page, nothing reconnects properly, and chat loses messages on bad wifi. Make real-time one thing."
- What flows, in which direction, at what rate? Notifications: server to client, a few a minute. Progress: server to client, 10 a second per job. Dashboard: server to client, 50 a second. Editor: both ways, keystroke rate. Chat: both ways, bursts.
- How many clients, how many topics each? 20k concurrent; 5 to 15 topics per page.
- What happens on bad networks? Phones on trains; office wifi that drops for 30 s; laptops that sleep for an hour.
- What must not be lost? Chat messages (sent and received); editor operations; job completion. Prices and dashboard samples may be sampled.
- Ordering? Within a chat room and a document, strictly. Across topics, no.
- Presence? Who is in a document, who is in a chat, with "away" states.
- Infrastructure? A load balancer with a 60 s idle timeout; a corporate customer behind a buffering proxy; several socket servers behind the balancer; Redis available.
- What is "one thing"? One client library every team uses; one server protocol; one place to debug.
- FR: subscribe to topics; receive ordered events per topic; send messages with delivery confirmation; presence per room; reconnection with resume; works through the customer's proxy.
- NFR: one connection per tab; 20k concurrent clients; 15 topics per page; events visible under 200 ms on a good link; reconnect within 2 s of the network returning, resuming without a page reload; no lost chat messages across a 30 s drop or a reload; a 50-per-second topic without dropped frames; a hidden tab that does not accumulate; heartbeats under the balancer's 60 s timeout; a fallback transport for the buffering proxy.
v1: the transports, and one socket per widget
// v1, as found: each widget opens its own connection and hopes
useEffect(() => {
const ws = new WebSocket(`${WS}/notifications`) // one of five on this page; no reconnect; no heartbeat; no resume
ws.onmessage = e => setCount(JSON.parse(e.data).count)
return () => ws.close()
}, [])
// what happens: the balancer closes idle sockets at 60 s (no heartbeat); a sleep kills all five; the proxy customer gets nothing;
// chat on a drop loses whatever was in flight; a tab left open all day has reconnected zero times and shows yesterday- Polling: simple, late by half the interval, cost per client per interval. Right for a badge count that may be a minute stale.
- Long-polling: held requests; near-zero latency; superseded by SSE except where SSE cannot pass.
- SSE: one HTTP stream, server to client, with reconnection and
Last-Event-IDreplay built into the browser. Needs HTTP/2 (HTTP/1.1 caps at ~6 connections per origin) and a server that flushes and disables proxy buffering. Right for notification, progress and price streams. - WebSockets: bidirectional, binary-capable, a held connection per client the server pays for in memory; no reconnection or replay for free. Right for chat and the editor.
- The infrastructure decides what passes: heartbeats under the balancer's idle timeout; a route around the CDN for sockets; a polling or SSE fallback for the buffering proxy, selected by a quick probe at connect time.
- One WebSocket per tab for the interactive session, with SSE as the fallback for server-to-client topics and polling as the last resort; pays: a transport negotiation at connect (complexity), and everything above the transport, which is the next round.
Round two: the subscription model
// v2: one connection per tab; topics with refcounts and lastSeq; backoff; resubscribe; heartbeat; a snapshot fallback
class Realtime {
private ws?: WebSocket; private topics = new Map<string, { refs: number; seq: number; subs: Set<(e: Event) => void> }>()
private attempt = 0; private pingTimer?: number; private pongDeadline?: number
subscribe(topic: string, onEvent: (e: Event) => void, since = 0) {
const t = this.topics.get(topic) ?? { refs: 0, seq: since, subs: new Set() }
t.refs++; t.subs.add(onEvent); this.topics.set(topic, t)
if (t.refs === 1 && this.ws?.readyState === WebSocket.OPEN) this.send({ type: 'sub', topic, since: t.seq })
this.ensureOpen()
return () => { t.refs--; t.subs.delete(onEvent); if (t.refs === 0) { this.topics.delete(topic); this.send({ type: 'unsub', topic }) } }
}
private ensureOpen() {
if (this.ws && this.ws.readyState <= WebSocket.OPEN) return
const ws = this.ws = new WebSocket(URL)
ws.onopen = () => { this.attempt = 0; for (const [topic, t] of this.topics) this.send({ type: 'sub', topic, since: t.seq }); this.heartbeat() }
ws.onmessage = ({ data }) => {
const m = JSON.parse(data)
if (m.type === 'pong') { this.pongDeadline = undefined; return }
if (m.type === 'snapshot_required') { this.topics.get(m.topic)?.subs.forEach(fn => fn({ kind: 'resnapshot', topic: m.topic })); return } // caller refetches via the query layer, then resumes
if (m.type === 'event') {
const t = this.topics.get(m.topic); if (!t) return
if (t.seq && m.seq > t.seq + 1) { this.send({ type: 'sub', topic: m.topic, since: t.seq }); return } // gap: ask for replay; drop this one (it will come back in order)
t.seq = m.seq; t.subs.forEach(fn => fn(m))
}
}
ws.onclose = () => { clearInterval(this.pingTimer); const delay = Math.min(30_000, 1000 * 2 ** this.attempt++) * (0.7 + Math.random() * 0.6); setTimeout(() => this.ensureOpen(), delay) }
}
private heartbeat() { this.pingTimer = window.setInterval(() => { if (this.pongDeadline && Date.now() > this.pongDeadline) return this.ws?.close(); this.pongDeadline = Date.now() + 10_000; this.send({ type: 'ping' }) }, 25_000) }
private send(m: unknown) { if (this.ws?.readyState === WebSocket.OPEN) this.ws.send(JSON.stringify(m)) }
}
document.addEventListener('visibilitychange', () => { if (!document.hidden) rt.ensureOpenNow() }) // foreground: assume dead, reconnect now
// React: useRealtime(topic, handler) = useEffect(() => rt.subscribe(topic, handler), [topic]); the handler updates a store or patches the query cache- One connection manager per tab; components subscribe to topics and get an unsubscribe for their effect cleanup; the manager reference-counts topics and sends subscribe and unsubscribe frames; twelve widgets are one socket, one heartbeat, one reconnect.
- A small, versioned protocol: sub, unsub, event, ack, ping, pong, snapshot_required. Every event carries topic and seq; the client keeps lastSeq per topic. JSON until the profile says otherwise (M9 moves to binary).
- Ordering within a topic, none across: seq in order on one connection; a jump is a gap, and the client asks for replay from its lastSeq (or refetches a snapshot through the query layer and resumes). Across topics the client assumes nothing.
- Reconnect with exponential backoff and jitter (1 s doubling to 30 s, ± 30%), resubscribe every held topic with its lastSeq; the server replays from a per-topic buffer (the last N events or minutes) or answers snapshot_required.
- Liveness: ping every 25 s under the balancer's 60 s; a missed pong closes and reconnects; a tab returning to the foreground assumes the socket is dead and reconnects now rather than waiting out the timer.
- Into React:
useRealtime(topic, handler)is an effect; the handler updates a store read byuseSyncExternalStoreor patches the query cache (the React course part 8's "server state is a cache", fed by events instead of refetches).
- v2 buys topics over one connection with order, gap detection and resume; pays: a manager with reference counting (complexity), a protocol both sides implement and version (a contract), per-topic replay buffers on the server (money), and a snapshot fallback that requires every topic's state to be fetchable on demand (complexity, and the reason the query layer exists).
Round three: backpressure
// v3: per-topic policy. state-like topics coalesce by key; event-like topics batch per frame; hidden tabs downgrade
const latest = new Map<string, unknown>(); const batches = new Map<string, unknown[]>(); let raf = 0
function onEvent(topic: string, e: Event) {
const kind = policy(topic) // 'state' | 'events' from the topic registry
if (kind === 'state') latest.set(e.key, e.data) ; else (batches.get(topic) ?? batches.set(topic, []).get(topic)!).push(e.data)
if (!raf) raf = requestAnimationFrame(flush)
}
function flush() {
raf = 0
if (latest.size) { store.applyLatest(latest); latest.clear() } // one apply per frame, newest only
for (const [topic, items] of batches) { store.applyBatch(topic, items); batches.delete(topic) } // every item, one commit
}
// throttle at the source: rt.subscribe('px:EURUSD', h, { maxRate: 10 }) → the server conflates per client
// hidden tab: on hide, unsubscribe high-rate topics (keep their lastSeq); on show, resubscribe with since (replay or snapshot)
// instrumentation: apply lag = now − event.ts (server clock, offset-corrected); queue size; alert at 1 s of lag- The dashboard topic at 50 a second and a price topic at 200 a second each trigger a store update and a render per event; the main thread saturates; the application queue grows; the displayed values are seconds old while looking live; a tab hidden over lunch has 400k queued events and freezes for 30 s on return. The number is events per second against frames per second, and the failure is silent.
- State-like topics (a price, a presence entry, a progress percentage) coalesce by key: a Map of key → latest, applied once per frame (M4's ring and rAF); intermediates are history, not state; the display is always the newest.
- Event-like topics (chat, trades, log lines) batch: the server sends arrays per 50 ms; the client applies a batch per frame. Past what the UI can show, sample for display ("+1,200 messages") and route the full stream to a worker (M2) if something needs it all.
- Throttle at the source: the subscription carries
maxRateor a conflation flag; the server conflates per client; the cheapest events never leave the server. - Hidden tabs unsubscribe high-rate topics on hide and resubscribe with
sinceon show (replay or snapshot); slow devices measure their apply time and request a lower rate, the way M3's player adapts bitrate. - Instrument it: apply lag (now minus the event's server timestamp, clock-offset corrected) and queue size as metrics; alert at one second of lag. The failure is silent without this.
- v3 buys bounded memory, a responsive page and the newest state under any rate; pays: a topic registry with a policy per kind (complexity), server-side conflation and batching (complexity, a contract), and the acceptance that the display is a sample of the stream, which product agrees to everywhere except the tape (M9).
Round four: presence and delivery guarantees
The chat shows people as online after they closed the laptop; messages sent on a train vanish; a retried send appears twice. The numbers that broke are flapping connections per user and messages in flight across a drop. v4 adds the two things a subscription model does not give: honest presence, and guarantees for what the client sends.
- Presence is server truth with a TTL: subscribe to a room adds the client to its member set with a TTL refreshed by heartbeats; close or expiry removes them and broadcasts a leave; the client renders from an initial snapshot plus join and leave events. A grace period (30 s "away" before "offline") keeps flapping phones from flickering the list. Across socket servers the member set lives in Redis with TTL keys and keyspace events; the client sees one room.
- Client ids and acks: every outgoing message gets a UUID and a client sequence, sits in a pending outbox as "sending", and is resent with the same id if no ack arrives in 5 s; the server deduplicates by client id and acks with its own id and seq: at-least-once delivery, idempotent application, no duplicates.
- The outbox persists (IndexedDB) so a reload or a crash does not lose a message; on reconnect, pendings are resent in order before resubscribing; the UI shows them greyed with a clock until acked (M5's local-first idea, for a queue).
- Delivered and seen are recipient acks and batched read receipts (the last seen seq per room per view session), each a message per recipient: linear if batched, quadratic in a 5,000-person room if not. Direct messages and small groups get receipts; large rooms do not; the numbers decide.
- The UI states are the protocol: clock (sending), one tick (server ack), two ticks (delivered), blue (seen). A user can read the guarantees off the screen.
- v4 buys presence that is honest and messages that are neither lost nor duplicated; pays: a presence protocol with TTLs, grace periods and a shared store (complexity, money), an outbox with ids, acks and persistence (complexity), and receipts bounded by room size (capability).
- Stop: one library, one protocol, one place to debug; every number in the brief is met. End-to-end encryption, voice, and a trading-grade tape (M9) are the next systems on the same base.
The whole board, and the exercise
| Round | The number | The break | The design | Paid in |
|---|---|---|---|---|
| v1 | Five sockets per page; a 60 s balancer; a proxy | No heartbeat, reconnect or resume; nothing through the proxy | One WebSocket per tab; SSE fallback for one-way topics; polling last; a connect-time probe | A transport negotiation |
| v2 | 15 topics per page; 30 s drops; hour-long sleeps | Lost events; stale state after reconnect; no ordering | A connection manager with refcounts; a versioned protocol with seq per topic; gap detection; backoff with jitter; resubscribe with since; replay buffers; snapshot fallback; heartbeats | A manager; a contract; server buffers; a fetchable snapshot for every topic |
| v3 | 50 to 200 events/s per topic; hidden tabs | Apply lag; queue growth; a 30 s freeze on foreground | Policy per topic: coalesce state, batch events, sample for display, throttle at the source, downgrade when hidden; lag metrics | A registry; server conflation; an agreed sampling |
| v4 | Flapping phones; sends across drops; 5,000-person rooms | Phantom presence; lost and duplicated messages | Presence with TTL, grace and a shared store; client ids, acks and a persisted outbox; receipts bounded by room size | A presence protocol; an outbox; quadratic receipts avoided |
- The transport is a choice of direction and rate; the protocol above it is the design.
- Sequence numbers per topic are the whole correctness story: they make loss visible and resume possible.
- Backpressure is a policy per kind of data, chosen on purpose, measured as lag, with the server as the cheapest place to drop.
- Presence and delivery are protocols the UI makes visible: a TTL and a grace period, an outbox and an ack, a clock and two ticks.