06 · Real-Time Systems: WebSockets & Pub/Sub¶
Chat messages, live scores, collaborative cursors, ride tracking, trading dashboards — all need the server to push updates to clients as they happen. HTTP's request-response model is client-initiated, so real-time systems need either a way to simulate push or a persistent connection. The hard parts are not the protocol; they are holding millions of connections, routing a message to the right one, and recovering gracefully when connections drop (which, on mobile networks, is constantly).
Four ways to push¶
| Technique | How | Latency | Cost / complexity |
|---|---|---|---|
| Short polling | Client asks every N seconds | Up to N seconds | Wasteful: most responses are empty |
| Long polling | Server holds the request until data arrives or timeout | Near real time | One request per message; reconnect overhead |
| Server-Sent Events (SSE) | One long-lived HTTP response streaming events | Real time, server→client only | Simple, works with HTTP infrastructure, auto-reconnect with last event ID |
| WebSockets | Upgraded persistent, bidirectional connection | Real time, both directions | Stateful connections; needs L7 support for upgrades and sticky routing |
Choose the simplest one that fits. A notifications badge updating every 30 seconds is fine with polling. A live-score page (server → client only) is a natural fit for SSE. Chat, games, and collaborative editing — frequent messages in both directions — suit WebSockets.
The architecture: gateways plus pub/sub¶
Persistent connections make the connection tier stateful: user 42's socket lives on one specific gateway server. Something must route "deliver to user 42" to that server.
flowchart LR
C1[Client A] <--> G1[Gateway 1]
C2[Client B] <--> G2[Gateway 2]
G1 <--> PS[[Pub/Sub layer]]
G2 <--> PS
S[App services] --> PS
G1 -. register user→gateway .-> REG[(Connection registry)]
G2 -.-> REG
- Gateways do nothing but hold connections, authenticate, and relay frames. Keeping them thin lets each hold a very large number of mostly idle connections.
- Business logic lives in ordinary stateless services.
- Routing, two common options:
- Registry: a store mapping
user_id → gateway_id. To deliver, look up the gateway and send to it directly. - Topic subscription: each gateway subscribes to pub/sub channels for the users (or
rooms) connected to it; publishers publish to
user:42orroom:7and the pub/sub layer delivers to subscribed gateways.
Worked example: a WebSocket-free pub/sub fan-out model¶
The core routing logic is independent of the wire protocol. This model shows rooms, multiple gateways, and per-connection delivery — and why a message sent while someone is disconnected needs a separate path.
# pubsub_model.py — gateways subscribe to room topics; publishers never see sockets
from collections import defaultdict
class PubSub:
def __init__(self):
self.subs = defaultdict(set) # topic -> gateways
def subscribe(self, topic, gw): self.subs[topic].add(gw)
def unsubscribe(self, topic, gw): self.subs[topic].discard(gw)
def publish(self, topic, msg):
for gw in list(self.subs[topic]):
gw.deliver(topic, msg)
class Gateway:
def __init__(self, name, bus):
self.name, self.bus = name, bus
self.conns = defaultdict(set) # topic -> local user ids
self.inbox = defaultdict(list) # user -> frames "sent" on the socket
def join(self, user, room):
if not self.conns[room]:
self.bus.subscribe(room, self) # first local member: subscribe once
self.conns[room].add(user)
def leave(self, user, room):
self.conns[room].discard(user)
if not self.conns[room]:
self.bus.unsubscribe(room, self)
def deliver(self, topic, msg):
for user in self.conns[topic]:
self.inbox[user].append(msg)
bus = PubSub()
g1, g2 = Gateway("g1", bus), Gateway("g2", bus)
g1.join("ana", "room:7"); g1.join("ben", "room:7"); g2.join("cy", "room:7")
bus.publish("room:7", {"from": "ana", "text": "hi all", "seq": 1})
g2.leave("cy", "room:7") # cy's phone loses signal
bus.publish("room:7", {"from": "ben", "text": "cy?", "seq": 2})
print({u: [m["seq"] for m in g1.inbox[u]] for u in ["ana", "ben"]}) # both got 1 and 2
print({"cy": [m["seq"] for m in g2.inbox["cy"]]}) # only 1
Note two design points. First, a gateway subscribes to a room once, no matter how many
local members it has — so a message to a 10,000-member room crosses the pub/sub layer once
per gateway, not once per member. Second, pub/sub is fire-and-forget: cy missed
message 2. Real-time delivery is an optimization; durable delivery needs a store plus a
catch-up protocol.
Reconnection and catch-up¶
Clients disconnect all the time. A robust protocol:
- Every message has a per-conversation (or per-stream) sequence number, assigned when it is durably stored.
- The client remembers the last sequence it has seen.
- On reconnect, it asks for everything after that sequence (SSE's
Last-Event-IDis a built-in version of this), fetched from the durable store; then live delivery resumes. - Clients deduplicate by sequence/message ID, since catch-up and live delivery may overlap.
Presence¶
"Online" indicators look trivial and are expensive at scale: every status change could fan out to every friend. Typical techniques:
- Heartbeats with TTLs: the gateway refreshes
presence:user42with a short TTL; no refresh means offline. No explicit "went offline" message is needed for crashes. - Lazy fetch: fetch presence for the users currently visible on screen rather than pushing every change to everyone.
- Debounce: ignore flaps shorter than a few seconds.
Operational concerns¶
- Load balancing: WebSocket upgrades need an L7 balancer that supports them; once established, a connection stays on its gateway. Balance on new connections, and watch for imbalance after deploys.
- Deploys: restarting a gateway drops all its connections, and they all reconnect at once. Drain gradually and make clients reconnect with randomized backoff to avoid a thundering herd.
- Idle timeouts in proxies close quiet connections; send periodic pings.
- Back-pressure: a slow client cannot be allowed to make the gateway buffer unbounded data — cap per-connection buffers and drop or disconnect.
How It Actually Works¶
A WebSocket starts life as an HTTP/1.1 request with Upgrade: websocket and a random key
header. The server replies 101 Switching Protocols with a hash derived from that key,
proving it understood the handshake. From then on, the same TCP connection carries
WebSocket frames — small headers indicating type (text, binary, ping, pong, close) and
length — in both directions, with no per-message HTTP overhead. Client-to-server frames are
masked with a random key, a measure designed to stop the protocol from being abused to
poison intermediary caches.
A gateway can hold many connections because an idle connection costs only a socket and some buffers, and an event loop (epoll/kqueue-style readiness notification) lets one thread watch a huge number of sockets and act only on the few with data ready. The practical limits are memory per connection, file-descriptor limits, and the CPU spent on TLS and heartbeats — which is why gateways are kept free of business logic.
Common mistakes¶
- Treating pub/sub delivery as durable. Store first, then publish; recover with sequence-based catch-up.
- Business logic in the gateways, making them heavy and hard to deploy.
- Synchronized reconnect storms after a deploy or network blip.
- Pushing every presence change to every contact.
- WebSockets where SSE or polling would do, adding stateful infrastructure for one-way, low-frequency updates.
Exercise¶
- Extend
pubsub_model.pywith a durable per-room message log and areconnect(user, last_seq)that replays missed messages, then resumes live delivery without duplicates. - Estimate gateway count for 5M concurrent connections if one gateway comfortably holds 100,000 (state this as an assumption), with capacity to lose one availability zone of three.
- Design live tracking of a delivery courier's position (updates every 3 seconds, viewed by one customer). Which push technique, and why? What happens when the customer's app is in the background?