Design a real-time chat system: WebSockets and presence
Design a real-time chat system around long-lived connections: a gateway tier, a session registry, pub/sub fan-out, per-conversation sequences, and presence.
The chat feature ships on the existing API servers. A WebSocket endpoint, an in-memory dictionary of connections, and a loop that forwards messages: it works for the two hundred users in the beta. At twenty thousand users the first ordinary deploy disconnects all of them at once, they reconnect in the same second, and the reconnect storm takes down the login endpoint that had nothing to do with chat. A week later a customer reports that after switching from WiFi to cellular, a reply arrived above the message it answered.
Neither bug is about sending a message. Sending a message is a write and a push. The hard part of a chat system is that a connection is state, the state lives on one machine, and the machine will be replaced while the conversation is still going.
An HTTP request lives for milliseconds, carries everything the server needs, and can be answered by any server behind the load balancer. A WebSocket connection (RFC 6455) starts as an HTTP request, upgrades, and then lives for hours as one TCP connection, with a server process holding a socket, a buffer, and a user identity in memory the entire time.
That changes what the load balancer is allowed to be. It must either pass TCP through (layer 4) or
understand the Upgrade handshake and keep the connection pinned to the node that accepted it, and
its idle timeout must be longer than the client's heartbeat interval or it will quietly close
healthy connections. The load balancing algorithms that spread
requests across a fleet still apply to the handshake, but after it there is nothing left to
balance. The connection is where it is.
It also changes what a deploy is. Stopping a node drops every connection it holds. Draining means the node stops accepting new connections, sends each client a "reconnect" frame spread over a window with random jitter, and exits only when its socket count is near zero. Without the jitter, a rolling deploy is a series of self-inflicted reconnect storms.
Put the sockets on their own tier. Gateway nodes accept connections, authenticate the handshake, hold the socket, and do nothing else. API servers stay stateless: message history, search, profiles, and everything else a request can answer. A gateway node is sized by connections, an API server by requests, and the two deploy on different schedules.
concurrent connections at peak: 1,000,000
sockets per gateway node: 50,000 (at ~30 KB per socket, about 1.5 GB)
gateway nodes: 20, plus headroom for one draining and one failed
messages: 10M daily users x 40/day = 400M/day = 4,600/s average, plan for 5x peakThese are estimates to size the tiers, not measurements. The point they make is that connection count, not request rate, drives the gateway fleet.
Stateless above, stateful below
A request can go anywhere. A connection is where it is.
API tier
sized by requestshistory, search, signed upload URLs, the message log write: any request, any node
Session registry
shared store, TTL per entry- alice gw-1 c-0a91 ttl 84 s
- bob gw-2 c-77f3 ttl 61 s
- bob gw-3 c-b2c4 ttl 88 s
one entry per open socket; a crashed node's entries expire on their own
Gateway tier
sized by connections, holds sockets and nothing else
50,000 per node - gw-1 serving
48,900 sockets
alice's phone
- gw-2 serving
47,200 sockets
bob's laptop
- gw-3 serving
49,600 sockets
bob's phone
- gw-4 draining
3,100 sockets
deploy: reconnect frames with jitter
One message, two devices
- 1 alice sends on gw-1; the API writes seq 1,042 to the conversation log and acks her
- 2 registry lookup: bob has sockets on gw-2 and gw-3
- 3 publish one envelope to topic gw-2 and one to topic gw-3
- 4 each gateway drops the frame into bob's bounded outbound channel; a full channel closes that socket
gw-4 accepts nothing new and exits at zero sockets; without jitter its 50,000 clients would return in one second
If user B is connected to gateway 3 and user A sends from gateway 1, gateway 1 has to find gateway
3. The session registry answers that: a shared store mapping user_id to the set of (gateway node,
connection id) pairs currently open, written on connect, deleted on disconnect, and refreshed by
heartbeat with a TTL so a crashed node's entries expire on their own.
Fan-out between gateways is pub/sub. Each gateway subscribes to its own topic; a sender looks up the recipients in the registry and publishes one envelope per target node. Every device the recipient has open receives the message, because the registry holds one entry per connection. A recipient with no entry at all is offline, and the message goes to the notification system for a push instead. The bus is fire-and-forget by design; a queue and an event bus solve different problems, and durability comes from the log in the next section, not from the bus.
Ordering is the property people get wrong first. Two devices, a cellular handover, and a retry are enough to deliver frames in a different order than they were sent, and the network will never promise otherwise. So order is assigned by storage. Every conversation has a counter, every message takes the next value, and the message is committed to a log keyed by (conversation, sequence) before the sender receives an ack. Only then is it fanned out.
CREATE TABLE messages (
conversation_id uuid NOT NULL,
seq bigint NOT NULL,
message_id uuid NOT NULL UNIQUE, -- generated by the client
sender_id uuid NOT NULL,
body text NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (conversation_id, seq)
);
-- one transaction: the row lock on conversations serializes writers per conversation
UPDATE conversations SET last_seq = last_seq + 1
WHERE id = $1
RETURNING last_seq;
INSERT INTO messages (conversation_id, seq, message_id, sender_id, body)
VALUES ($1, $2, $3, $4, $5);The client generates message_id, so a retry after a lost ack carries the same id, hits the unique
index, rolls the transaction back (which also releases the sequence number), and the server answers
with the ack for the row already stored. That is at-least-once delivery made safe by an
idempotent receive: the sender retries until acked, the gateway may
push the same message twice after a reconnect, and every client keeps the ids it has already
rendered so a repeat is dropped on arrival.
Ordering is a property of the log, not the network. Give each conversation one counter, write to it before you ack, and let every client sort by it.
Presence is "is this user connected right now", and the naive implementation is correct and
unaffordable. Each connection sends a heartbeat every 30 seconds, the gateway refreshes a
presence:user-42 key with a 90 second TTL, and online means the key exists. That part is cheap.
The cost is telling other people. A user with 1,000 contacts who comes online produces 1,000
pushes; at 100,000 status changes a minute with 200 interested watchers each, that is 20 million
presence events a minute for a feature nobody reads to the millisecond.
Batch status changes per watcher into one frame every few seconds, and send presence only for contacts currently on screen: the client subscribes to the visible set and unsubscribes when it scrolls away. This is the same trade as fan-out on write versus fan-out on read, and presence is the case where reading on demand wins for large contact lists.
Inside a gateway, a fan-out loop that awaits SendAsync on every socket in turn is paced by the
slowest one, and a phone in a tunnel is slow for minutes. Give each connection a bounded outbound
channel with a single task draining it. The fan-out loop does a non-blocking TryWrite, and a full
channel means the client is not keeping up, so the socket is closed and the reconnect path below
backfills what it missed.
public sealed class ClientConnection
{
private readonly Channel<ServerFrame> _outbound =
Channel.CreateBounded<ServerFrame>(new BoundedChannelOptions(256)
{
SingleReader = true,
FullMode = BoundedChannelFullMode.Wait // TryWrite returns false when full
});
// Called by the fan-out loop for every recipient. Never awaits the network.
public bool TryEnqueue(ServerFrame frame)
{
if (_outbound.Writer.TryWrite(frame))
{
return true;
}
_outbound.Writer.TryComplete(); // slow consumer: drop the socket, not the loop
return false;
}
// One task per connection drains the channel into the socket.
public async Task PumpAsync(WebSocket socket, CancellationToken cancellationToken)
{
await foreach (var frame in _outbound.Reader.ReadAllAsync(cancellationToken))
{
await socket.SendAsync(frame.Payload, WebSocketMessageType.Text, true, cancellationToken);
}
await socket.CloseAsync(WebSocketCloseStatus.PolicyViolation, "client too slow", cancellationToken);
}
}The same shape applies to the heartbeat: a timer per connection that enqueues a ping frame and closes the socket when two replies in a row are missing. The ASP.NET Core WebSockets guide covers the accept loop that the pump attaches to.
An image is a megabyte; a chat frame is a hundred bytes. Pushing file bytes through the gateway means a socket that is busy for seconds, a gateway holding buffers, and a message that cannot be retried independently of its attachment. Instead the client asks the API for a short-lived signed upload URL scoped to one object key, uploads directly to object storage, and then sends an ordinary message whose body references the object. Recipients download through a signed URL of their own. The gateway never sees a byte of the file.
The per-conversation sequence pays for itself on reconnect. The client keeps the highest seq it
has seen per conversation. On reconnect (with jittered exponential backoff, so a node draining
50,000 sockets does not see them all return in the same second) it sends those high-water marks,
and the API returns everything above them from the log. A gap in the sequence on the client is a
detectable missing message rather than a silent one, and the same request serves a new device that
has seen nothing yet.
The topology decides whether any of this works: sockets on their own tier, a registry any node can read, storage that assigns order before the ack, and attachments that bypass the socket. Katabench's System Design Studio has a "Build Real-time Chat" series that builds that progression: "Establish a Chat Session", "Scale Real-time Connections", "Persist Chat Messages", "Deliver Chat Attachments", and "Operate the Real-time Chat Platform". Each challenge checks the canvas against authored rules (a registry the gateways share, no attachment bytes through the connection tier, no client talking straight to the message store), and the grading model explains what structural, deterministic feedback looks like next to the code tracks.
Practice what you just read
- Open the Pro kata: Establish a Chat Session
System Design Studio Easy Katabench Pro
Establish a Chat Session
A prototype lets mobile clients write messages straight to storage. Introduce a trusted real-time gateway and chat service before scaling connections.
More like this: System Design Studio →