Part 8 · 6 chapters · ~45 min

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.

50

The brief and the questions

the brief

"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."

the questions, and the answers
  1. 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.
  2. How many clients, how many topics each? 20k concurrent; 5 to 15 topics per page.
  3. What happens on bad networks? Phones on trains; office wifi that drops for 30 s; laptops that sleep for an hour.
  4. What must not be lost? Chat messages (sent and received); editor operations; job completion. Prices and dashboard samples may be sampled.
  5. Ordering? Within a chat room and a document, strictly. Across topics, no.
  6. Presence? Who is in a document, who is in a chat, with "away" states.
  7. 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.
  8. What is "one thing"? One client library every team uses; one server protocol; one place to debug.
the requirements, with numbers
  1. 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.
  2. 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.
the thesis
Real-time is a subscription model with sequence numbers, backpressure and reconnection, or it is a demo. The transport is the easy part; the protocol on top of it, and the policies per topic, are the design.
51

v1: the transports, and one socket per widget

code
// 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
the transports, by what they give
  1. Polling: simple, late by half the interval, cost per client per interval. Right for a badge count that may be a minute stale.
  2. Long-polling: held requests; near-zero latency; superseded by SSE except where SSE cannot pass.
  3. SSE: one HTTP stream, server to client, with reconnection and Last-Event-ID replay 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.
  4. 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.
  5. 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.
the sentence for the transport choice
  1. 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.
THE TRANSPORTS
polling, long-polling, SSE and WebSockets against latency, direction and cost
swipe the figure sideways, or tap expand for full screen
1/6
polling
Polling: GET /updates?since=… every 10 s. Latency: 5 s average, 10 s worst. Cost: a request per client per interval whether or not anything changed: 10k clients at 10 s is 1k requests a second of mostly nothing. Works through every proxy and CDN. Right for low rates and loose freshness (a badge count); wrong for a chat.
52

Round two: the subscription model

code
// 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
v2
  1. 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.
  2. 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).
  3. 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.
  4. 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.
  5. 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.
  6. Into React: useRealtime(topic, handler) is an effect; the handler updates a store read by useSyncExternalStore or patches the query cache (the React course part 8's "server state is a cache", fed by events instead of refetches).
the sentence
  1. 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).
V2: THE SUBSCRIPTION MODEL
one connection, many topics, a sequence per topic, and resubscribe from where you were
swipe the figure sideways, or tap expand for full screen
1/6
one connection
One connection per tab (not per component: a page with twelve live widgets opening twelve sockets is twelve heartbeats, twelve reconnects, and a server that sees twelve clients). A client-side connection manager owns the socket; components subscribe to topics through it and unsubscribe on unmount; the manager reference-counts topics and sends subscribe and unsubscribe frames.
53

Round three: backpressure

code
// 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 break
  1. 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.
v3: a policy per topic, by data kind
  1. 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.
  2. 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.
  3. Throttle at the source: the subscription carries maxRate or a conflation flag; the server conflates per client; the cheapest events never leave the server.
  4. Hidden tabs unsubscribe high-rate topics on hide and resubscribe with since on show (replay or snapshot); slow devices measure their apply time and request a lower rate, the way M3's player adapts bitrate.
  5. 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.
the sentence
  1. 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).
V3: BACKPRESSURE
when events arrive faster than the client can apply them, something must give, on purpose
swipe the figure sideways, or tap expand for full screen
1/6
without
Without backpressure: 200 events a second on a price topic; each event triggers a store update and a render; the main thread spends 100% on applying; the socket's receive buffer backs up; events are applied seconds late; the displayed price is old while looking live. bufferedAmount on the send side and a growing application queue on the receive side are the signals.
54

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.

v4
  1. 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.
  2. 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.
  3. 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).
  4. 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.
  5. 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.
the sentence, and the stop
  1. 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).
  2. 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.
V4: PRESENCE AND DELIVERY GUARANTEES
who is online, what was delivered, and what the client promises the user
swipe the figure sideways, or tap expand for full screen
1/6
presence
Presence: on subscribe to a room the server adds the client to the room's member set with a TTL refreshed by heartbeats; on disconnect or TTL expiry it removes them and broadcasts a leave. The client renders the member list from join and leave events plus an initial snapshot. Flapping connections (a phone on the edge of coverage) produce join/leave storms: debounce leaves by a grace period (a user is "away" for 30 s before "offline").
55

The whole board, and the exercise

RoundThe numberThe breakThe designPaid in
v1Five sockets per page; a 60 s balancer; a proxyNo heartbeat, reconnect or resume; nothing through the proxyOne WebSocket per tab; SSE fallback for one-way topics; polling last; a connect-time probeA transport negotiation
v215 topics per page; 30 s drops; hour-long sleepsLost events; stale state after reconnect; no orderingA connection manager with refcounts; a versioned protocol with seq per topic; gap detection; backoff with jitter; resubscribe with since; replay buffers; snapshot fallback; heartbeatsA manager; a contract; server buffers; a fetchable snapshot for every topic
v350 to 200 events/s per topic; hidden tabsApply lag; queue growth; a 30 s freeze on foregroundPolicy per topic: coalesce state, batch events, sample for display, throttle at the source, downgrade when hidden; lag metricsA registry; server conflation; an agreed sampling
v4Flapping phones; sends across drops; 5,000-person roomsPhantom presence; lost and duplicated messagesPresence with TTL, grace and a shared store; client ids, acks and a persisted outbox; receipts bounded by room sizeA presence protocol; an outbox; quadratic receipts avoided
what the sequence teaches
  1. The transport is a choice of direction and rate; the protocol above it is the design.
  2. Sequence numbers per topic are the whole correctness story: they make loss visible and resume possible.
  3. Backpressure is a policy per kind of data, chosen on purpose, measured as lag, with the server as the cheapest place to drop.
  4. 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.
the exercise
Open a live page you own, count its connections in the Network panel, sleep the laptop for ten minutes, wake it, and watch. Then turn off wifi for thirty seconds while sending a message. Every surprise is a round above; the first one you hit is the next thing to build.