chat-system-designWebSocketinbox-patternfan-outpresenceAPNs-FCMhorizontal-scalingmessage-routingsystem-design-interviewreal-time-messaging
TL;DR Here's the thing most tutorials miss: Server-Sent Events (SSE) is a third option between long polling and WebSockets. SSE is HTTP-based (works through standard HTTP/2 multiplexing), server-push only (client-to-server still uses regular HTTP requests), and has automatic reconnection built in.
WhatsApp processes over 100 billion messages a day. Every message that arrives instantly, every green dot that appears correctly, every notification that wakes up a sleeping phone — each one is the product of deliberate architectural decisions made by engineers who understood the tradeoffs. Here's how to design the same system from scratch.
Read the Deep Dive ↓ Open the Lab 💬 WebSocket + HTTP hybrid Inbox pattern for reliability Fan-out for group messages Heartbeat presence Horizontal scaling strategy Table of ContentsYour first instinct when designing a chat system might be to use HTTP for everything — it's familiar, it scales well, every load balancer knows how to handle it. And for sending a message, HTTP works perfectly: client makes a POST request, server receives the message and returns a 200 OK. Clean, simple, stateless. But receiving messages instantly is where HTTP fundamentally breaks down.
HTTP is client-initiated. The server cannot push a message to a client without being asked first. The two naive workarounds are short polling (client asks "any new messages?" every few seconds — wasteful, generates enormous load for empty responses) and long polling (client opens a request and the server holds it open until a message arrives or a timeout occurs — better, but each pending request still consumes a server connection). Neither is acceptable for a real-time chat system at scale.
WebSockets solve this properly. After an initial HTTP handshake (which upgrades the connection), a WebSocket is a persistent, full-duplex channel — the server can push messages to the client any time, without the client asking. Latency drops to milliseconds. But WebSocket connections are stateful: each connection ties up a server process or thread, scaling requires careful architecture, and load balancers need WebSocket awareness (sticky sessions or a shared state layer).
The industry-standard solution is a hybrid approach: HTTP for sending (stateless, scales trivially with horizontal load balancers), WebSocket for receiving (persistent connection managed by dedicated chat servers). This is exactly what WhatsApp, Facebook Messenger, and Discord use. You get HTTP's scaling simplicity for the high-volume write path and WebSocket's push capability for the latency-critical receive path. Don't try to over-engineer a single-protocol solution — the hybrid is the answer.
💡 Server-Sent Events: The Middle GroundHere's the thing most tutorials miss: Server-Sent Events (SSE) is a third option between long polling and WebSockets. SSE is HTTP-based (works through standard HTTP/2 multiplexing), server-push only (client-to-server still uses regular HTTP requests), and has automatic reconnection built in. For read-heavy, server-push scenarios like notifications, activity feeds, or dashboard updates, SSE is simpler to scale than WebSockets because it works over standard HTTP/2 connections without the complexity of WebSocket upgrade management. The tradeoff: SSE is one-directional (server → client only), while WebSockets are full-duplex. For chat, use WebSockets. For notifications, SSE is worth considering.
websocket_chat.js — client connectionHybrid approach: HTTP for sending, WebSocket for receiving
1. Connect to chat server via WebSocket (persistent receive channel)
const ws = new WebSocket('wss://chat.example.com/ws?token=' + authToken);
ws.onopen = () => {
console.log('Connected to chat server');
};
ws.onmessage = (event) => {
const msg = JSON.parse(event.data);
displayMessage(msg);
Send acknowledgement back to server
ws.send(JSON.stringify({ type: 'ack', messageId: msg.id }));
};
Heartbeat: send ping every 30 seconds to keep connection alive
setInterval(() => {
if (ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify({ type: 'ping' }));
}
}, 30000);
2. Send message via HTTP POST (stateless, load-balanced separately)
async function sendMessage(recipientId, text) {
const response = await fetch('/api/messages', {
method: 'POST',
headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${authToken}` },
body: JSON.stringify({ to: recipientId, text, clientMsgId: generateId() })
});
return response.json();
}
A chat system at scale isn't one server — it's a collection of specialized services organized into layers. The clean architectural separation is what makes it possible to scale each component independently as different bottlenecks emerge. Imagine the entire system divided into three functional layers, each with a distinct responsibility.
The stateless services layer handles everything that doesn't require persistent connections: user authentication, profile management, friend lists, group membership, and REST API endpoints for sending messages. These services are classic horizontally-scalable web servers behind a load balancer. You can add or remove instances freely without coordination — no state, no problem. This is where most of your write traffic lands (sending messages), and it's where you get the simplest scaling story.
The stateful services layer — the chat servers — is the interesting part. Each client maintains a persistent WebSocket connection to exactly one chat server. Modern servers can hold 100,000 to 1,000,000 concurrent WebSocket connections depending on hardware and configuration. As your user base grows, you add more chat servers. The critical constraint: since users are pinned to specific chat servers, routing messages between servers requires coordination — which we'll address in section 4.
The third-party integration layer handles what happens when users are offline: push notifications via Apple's APNs (Apple Push Notification service) and Google's FCM (Firebase Cloud Messaging). These services maintain persistent connections to literally hundreds of millions of devices and know how to deliver notifications even when apps are not running. Reinventing this infrastructure would be a multi-year project. You leverage it instead.
⚡ Service Discovery Is the Hidden Critical ComponentThe stateful chat server layer requires a service discovery mechanism that the architecture overview often glosses over. Every time a user connects or disconnects, the system needs to update a registry mapping "user X is on chat server Y." When Alice sends a message to Bob, Alice's chat server must look up which server Bob is currently connected to. This registry (often backed by Redis or etcd) is queried millions of times per minute. Its availability and latency directly determines message delivery latency. Design it for high read throughput, geographic distribution near your chat server regions, and automatic cleanup when connections drop.
Here's the fundamental problem: if you only deliver messages over live WebSocket connections, you lose messages whenever a user is offline, enters a tunnel, or their app crashes mid-delivery. You need a reliability mechanism that survives disconnections. The inbox pattern is the industry solution.
Every user has a personal message queue called an inbox. When a message is sent to you, the system first attempts immediate delivery via WebSocket. If you're online and connected, your device receives the message through the WebSocket and sends back an acknowledgement (ACK) — a confirmation that delivery succeeded. If the ACK is received, the message is never persisted to your inbox; the happy path is zero-storage. Only when you're offline or when delivery fails does the message get written to your inbox, where it waits until you reconnect.
The ACK mechanism is what provides the "at-least-once" delivery guarantee. The server only considers a message delivered when it receives the ACK. No ACK means the server will retry. When you come back online and reconnect to a chat server, your server immediately checks your inbox, downloads any waiting messages to your device in order, collects ACKs for each, and only removes messages from the inbox after successful ACK. This ensures that even a crash mid-download is recoverable — reconnect and the download resumes from the last un-ACKed message.
The design choice to store messages only temporarily (for offline delivery) rather than permanently is intentional and different from email. This keeps the inbox small, cheap to query, and avoids the indefinite storage costs of maintaining complete message history for billions of users. If persistent message history is a product requirement, it's a separate, explicitly designed storage system — not the inbox.
✅ Design for Idempotency: ACKs Can Be Lost TooThe ACK itself can fail. What if the server sends a message, the client receives it and sends an ACK, but the ACK gets lost? The server retries, the client receives a duplicate. Your client must deduplicate messages by message ID. Generate message IDs that are globally unique and deterministic (include sender ID, timestamp, and sequence number). The client-side message store should be keyed by message ID, making duplicate delivery a no-op rather than a duplicate display. This idempotency requirement applies to every message operation in the system — including the client's own sent messages, which should use client-generated IDs to prevent duplicates from network retries.
inbox_pattern.py — server-side delivery logicfrom dataclasses import dataclass
from enum import Enum
import asyncio
class DeliveryStatus(Enum):
DELIVERED = "delivered" # ACK received — don't store in inbox
PENDING = "pending" # in inbox, waiting for ACK
FAILED = "failed" # retries exhausted
async def deliver_message(server, msg, recipient_id: str):
# Try immediate delivery via WebSocket first
ws_connection = server.get_connection(recipient_id)
if ws_connection and ws_connection.is_alive:
try:
await ws_connection.send_json(msg.to_dict())
# Wait for ACK with timeout
ack = await asyncio.wait_for(
server.wait_for_ack(msg.id), timeout=5.0
)
return DeliveryStatus.DELIVERED # ACK received — done!
except asyncio.TimeoutError:
pass # fall through to inbox
# Delivery failed or user offline — put in inbox
await inbox.store(recipient_id, msg, ttl_hours=168) # 7 day TTL
# Trigger push notification for offline users
await push_service.notify(recipient_id, {
"title": ff"New message from {msg.sender_name}",
"body": msg.preview_text,
"data": {"message_id": msg.id}
})
return DeliveryStatus.PENDING
async def sync_inbox_on_connect(server, user_id: str):
"""Called when user reconnects — deliver all queued messages"""
messages = await inbox.get_pending(user_id)
for msg in messages: # in order
await ws.send_json(msg.to_dict())
ack = await wait_for_ack(msg.id, timeout=10.0)
if ack:
await inbox.remove(user_id, msg.id) # only remove after ACK!
With multiple chat servers each holding thousands of WebSocket connections, a new problem emerges: when Alice (connected to Chat Server 1) sends a message to Bob (connected to Chat Server 3), how does the message get to Bob's server? This is the inter-server routing problem, and it's where naive architectures fail at scale.
The solution uses a user presence service — a fast key-value store (Redis works well here) that maps each online user to their current chat server. When a user connects, their chat server registers them: presence.set("user:bob", "chat-server-3"). When a user disconnects, the registration is removed. When Alice's server needs to deliver a message to Bob, it first queries the presence service to find Bob's server, then makes a direct RPC (Remote Procedure Call) to that server to push the message through Bob's WebSocket connection. This direct server-to-server approach minimizes latency — the message takes one extra hop instead of going through a central broker.
The counterintuitive insight: a message queue (like Kafka) between chat servers seems like a natural fit but actually adds latency for the online delivery path. For real-time chat, you want the lowest possible latency — one direct RPC hop between servers. Kafka is excellent for the offline path (storing messages durably before inbox processing) but adds unnecessary latency to the critical hot path. Discord's architecture is a well-documented example of this direct server-to-server RPC approach.
⚠️ The Split-Brain Problem with PresenceWhat happens when the presence service reports Bob is on Chat Server 3, but Bob's connection actually dropped 2 seconds ago before the presence entry expired? Alice's server makes an RPC to Chat Server 3, which has no connection for Bob. This is the split-brain window — the gap between when a connection drops and when the presence service reflects it. The solution: Chat Server 3 should respond to the "no connection found" case by queuing the message to Bob's inbox immediately and returning success to Alice's server. Never fail silently when presence is stale — always fall through to the inbox as the reliability backstop.
One-on-one messaging is straightforward: one sender, one recipient, one delivery path. Group chat with 100 members is fundamentally different: one message must be delivered to potentially 100 different users, each possibly connected to different chat servers, each possibly online or offline. The naive solution — a for-loop that runs the one-on-one delivery logic 100 times — actually works at this scale, but requires careful design to handle the fan-out efficiently.
The fan-out pattern works as follows: when a message is sent to a group, the server retrieves the list of group members. For each online member, it makes an RPC to their chat server to push the message through their WebSocket. For each offline member, the message is written to their inbox. The fan-out to 100 members can be parallelized — all RPC calls go out concurrently, not sequentially, keeping latency low. At 100 members, this is manageable. At 1,000 members (like Slack workspaces), the fan-out requires more careful batching and queuing. At 10,000+ members (broadcast channels), a different architecture — publish-subscribe with per-channel queues — is necessary.
Here's the thing most system design explanations miss: group chat requires sender-side deduplication. If Alice sends to a group, Alice's own device should receive the message back (so she sees it in the conversation thread) — but only once. The server must include the sender in the fan-out with a flag indicating "this is your own message" so the client handles it correctly, and use the client-generated message ID to prevent the client from displaying a duplicate alongside the locally-rendered optimistic message.
🔬 Fan-Out at WhatsApp Scale Is a Distributed Systems ProblemFor groups with millions of members (broadcast channels, community groups), the fan-out pattern must shift from write-time fan-out to read-time fan-out. Instead of writing the message to each member's inbox at send time, you write the message once to a channel queue and let each member's client "pull" it when they reconnect. This is the distinction between "fan-out on write" (hot path, high write amplification) and "fan-out on read" (lazy, lower write cost but more complex read logic). Most chat systems use write-time fan-out for small groups (< 100 members) and read-time fan-out for large channels (> 1000 members).
The green "online" indicator seems trivial — just check if the user has an active WebSocket connection, right? The problem is that WebSocket connection state is an unreliable indicator of user presence. A phone screen turns off, the app goes to the background, the phone enters a tunnel — the WebSocket connection can appear alive at the TCP level for minutes after the user has effectively gone away. You don't want to show your friend as "online" when they're actually asleep and their phone is in airplane mode.
The solution is application-level heartbeats. Clients send a WebSocket ping frame (or a small JSON message) to their chat server every 30 seconds. The chat server tracks the timestamp of the last received ping for each connected user. If no ping has been received within 60 seconds, the server marks the user as offline — even if the underlying TCP connection is still technically open. This 60-second grace period smooths out brief disconnections: a 10-second tunnel or elevator doesn't show your friend as offline and back online repeatedly, which would be annoying and consume bandwidth updating all their contacts.
Presence updates need to be broadcast to relevant contacts efficiently. When Alice goes online, her contacts need their presence indicators updated. The fan-out here is the same problem as message fan-out: efficiently propagating state changes to a potentially large number of subscribers. For most chat apps, presence updates use eventual consistency — there's no guarantee that all of Alice's contacts see the green dot at exactly the same millisecond, and that's fine. The 30-second heartbeat interval already means presence has 30 second granularity by design.
💡 Presence Updates Are High-Volume: Batch ThemFor a user with 500 contacts, going online triggers 500 presence updates. Multiplied by millions of users doing this every time they open the app, presence becomes one of the highest-volume write operations in the system. The optimization: batch presence updates and use pub-sub channels rather than individual writes. Each user subscribes to a "contacts" presence channel. When Alice goes online, a single event is published to that channel, and a fan-out service distributes it. Rate limit and debounce presence broadcasts — don't broadcast "online" status more than once per 30 seconds per user, and batch updates to avoid a thundering herd when millions of users wake up simultaneously after an outage.
The architecture above works at 10,000 users. At 10 million users, several components will break under load, and you need to anticipate where. The first bottleneck is almost always chat server connection capacity. Each chat server handles a finite number of concurrent WebSocket connections. Scale this by adding more chat servers, updating service discovery to route new connections to servers with capacity. The bottleneck is linear and predictable — the easiest scaling problem in the system.
The inbox database is the second bottleneck. As millions of users go offline and come back online, their inboxes are written to and emptied constantly. The access pattern is heavily user-partitioned: user A's inbox is completely independent of user B's inbox. This makes horizontal partitioning (sharding) by user ID natural and effective. Shard the inbox database by user ID hash. Each shard handles a fraction of users. Add shards as load grows. This scales linearly and avoids cross-shard coordination.
For global reach, deploy chat servers in multiple geographic regions — US East, Europe, Asia Pacific, etc. Users connect to their nearest region for minimum latency. But cross-region messaging — Alice in London sending to Bob in Singapore — adds complexity: the message must travel from the London chat server cluster to the Singapore cluster. This is solved with a global message routing layer (often a Kafka-based backbone between regions) that handles the slower, higher-latency cross-region delivery path while keeping intra-region delivery fast and direct.
⚡ Connection Limits: Linux Kernel TuningA default Linux server configuration limits far fewer than 1 million concurrent connections. The default file descriptor limit (65,536 on most systems) caps the number of open sockets. To run high-concurrency WebSocket servers, tune: ulimit -n 1000000 (increase file descriptor limit), adjust net.core.somaxconn and net.ipv4.tcp_max_syn_backlog in kernel settings, and consider a socket-sharing architecture where multiple worker processes share the same listening port. With proper tuning, a single modern server (64+ cores, 256GB RAM) can sustain 500,000 to 1,000,000 concurrent WebSocket connections.
Follow a single message from Alice to Bob: Alice types "Hello" and hits send. Her client makes an HTTP POST to the stateless API servers (load-balanced, stateless). The API server receives the message, looks up Bob's presence in the presence service. Bob is online on Chat Server 3. The API server makes an RPC to Chat Server 3. Chat Server 3 pushes the message through Bob's WebSocket connection. Bob's client receives it, displays it, and sends an ACK. Chat Server 3 receives the ACK and confirms delivery. Total time: under 100 milliseconds for users in the same region.
If Bob is offline: the API server writes to Bob's inbox in the sharded inbox database and triggers a push notification via FCM/APNs. When Bob comes back online, he connects to a chat server (probably not Chat Server 3 — that's fine), the server checks his inbox, syncs all pending messages, collects ACKs, and cleans up the inbox. Bob sees the message immediately upon opening the app.
Four experiments: protocol comparison, inbox delivery, fan-out calculator, and server capacity planner.
Message delivery latency comparison across protocols
Protocol Latency Simulator Base network latency (ms) 50 Poll interval — short polling (s) 5s — Short polling (avg) — Long polling — WebSocket — WinnerMessage lifecycle — click buttons to simulate scenarios
Inbox Pattern Simulator 0 Delivered (WS) 0 In inbox 0 Push notifications 0 RetriesFan-out cost at different group sizes
Fan-Out Calculator Group members 50 % Online 60% RPC latency (ms) 5ms — Online (direct RPC) — Offline (inbox) — Fan-out latency — StrategyWebSocket servers needed as DAU grows
Chat Server Capacity Planner Daily Active Users (DAU) 1M Concurrent % online 10% WS connections per server 100K — Concurrent users — Chat servers needed — Est. server cost/mo — Headroom