Sharpened prompt. Design a messaging platform for 2B users and 500M concurrent connections delivering 100B messages/day, where the server can never read a message, a phone offline for a week receives everything in the right order on reconnect, a 500-member group does not require 500 encryptions per message, and four of a user's devices stay in sync without any of them being always online.
Two constraints fight each other on every page of this design: the server must route messages it cannot read, and it must do so for hundreds of millions of simultaneously connected devices. Almost every interesting decision comes from that tension.
1. Problem framing
Functional requirements
- One-to-one and group messaging with end-to-end encryption.
- Delivery and read receipts; typing indicators; presence.
- Offline delivery: messages queue until the recipient reconnects.
- Media (photos, video, voice notes), also end-to-end encrypted.
- Multiple devices per account, all in sync.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Concurrent connections | 500M | Cheap per-connection cost — Erlang/Go, not threads |
| Delivery latency | P99 < 500ms when both online | Persistent connections, regional edge |
| Message durability | Never lose an undelivered message | Durable per-user inbox until acknowledged |
| Server knowledge | Zero plaintext, ever | E2EE by default; metadata minimised |
| Ordering | Per-conversation, causally consistent | Client-assigned sequence, server preserves it |
Back-of-the-envelope
Users: 2B, 500M concurrent connections
Messages: 100B/day = 1.16M/sec average, ~3M/sec peak
Per-conn: ~10 KB of kernel + app state (Erlang processes are ~2 KB each)
500M × 10 KB = 5 TB of RAM across the connection tier
At 2M connections/node (achievable with tuned FreeBSD/Linux + BEAM)
-> ~250 connection nodes. WhatsApp famously ran ~1M/server in 2012.
Inbox: undelivered messages only. ~2% of messages queue = 2B rows/day,
deleted on ack. Steady state is small — this is a QUEUE, not an archive.
Groups: 500-member group, 1 message -> pairwise encryption = 500 ciphertexts.
With sender keys: 1 ciphertext + occasional key distribution. See 4.4.
The sender-key line is the number that matters most: naive group E2EE is O(members) per message, and sender keys make it O(1).
2. High-level architecture
The router handles opaque ciphertext. It knows sender, recipient, size, and timing — that metadata is unavoidable — but never content. Saying which metadata you cannot hide is more honest and more impressive than claiming perfect privacy.
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Connection tier | Erlang/OTP (BEAM) or Go | Millions of lightweight processes/goroutines; preemptive scheduling; per-connection isolation |
| Transport | TLS 1.3 + Noise Protocol framing over TCP | Noise gives a compact, well-analysed handshake; WebSocket fallback where TCP is blocked |
| Session registry | Redis or a distributed hash table | "Which node holds device D?" must be a sub-millisecond lookup |
| Inbox | Cassandra/ScyllaDB, partition (user_id, device_id) | Append-heavy, ordered range read, delete on ack. Perfect LSM workload |
| Key server | Simple KV, high availability | Stores only public keys; a breach reveals nothing decryptable |
| Crypto | Signal protocol (X3DH + Double Ratchet) | The standard; do not invent your own |
| Media | Object storage + CDN, client-side encrypted | Server stores random bytes; keys travel in the message |
| Wake-up | APNs/FCM with content-free pushes | The OS will not let you hold a socket in the background |
4. The toughest parts
4.1 Five hundred million open sockets
Why it's hard. A thread-per-connection model dies immediately: 500M threads at 1 MB of stack each is 500 TB. Even an event-loop model in a typical language struggles — per-connection heap objects, TLS session state, and buffers add up, and a single node handling a million connections must survive garbage collection without stalling every socket.
Solution — an actor-per-connection runtime with per-actor isolation, plus aggressive kernel tuning.
%% Erlang: one lightweight process per connection. ~2 KB each, isolated heaps,
%% so GC is per-process and never stops the world for other connections.
%% This isolation property is the reason BEAM fits this problem specifically.
handle_connection(Socket, User, Device) ->
ok = registry:bind(User, Device, self()),
loop(Socket, #state{user = User, device = Device}).
loop(Socket, State) ->
receive
{tcp, Socket, Data} ->
handle_frame(decode(Data), State),
loop(Socket, State);
{deliver, Envelope} -> % routed from another node
gen_tcp:send(Socket, encode(Envelope)),
loop(Socket, State);
{tcp_closed, Socket} ->
registry:unbind(State#state.user, State#state.device)
after 60000 ->
gen_tcp:send(Socket, ping()), % keepalive; also NAT hole-punching
loop(Socket, State)
end.
The properties that make this work, and are worth naming:
- Per-process heaps. A GC pause affects one connection, not a million. In a JVM or Go runtime, a large heap of connection state creates stop-the-world or scan pressure that shows up as latency on every connection.
- Preemptive scheduling. BEAM preempts processes by reduction count, so one misbehaving connection cannot starve the scheduler.
- Kernel tuning is part of the design: raise
fs.file-maxand per-process FD limits, tunenet.ipv4.tcp_memand socket buffers down (a million × 256 KB of buffers is 256 GB), enableSO_REUSEPORTso multiple acceptor threads share a listening socket, and useTCP_NODELAYsince messages are small and latency-sensitive. - Keepalive interval is a battery decision. Too frequent drains the phone; too infrequent and carrier NAT drops the mapping. WhatsApp tuned this per-network; the correct answer is "adaptive, learned per carrier," not a constant.
The session registry is the coordination point: (user, device) → node. Keep it in Redis or a DHT with a short TTL and heartbeats, so a node loss expires its sessions automatically and clients reconnect elsewhere.
4.2 End-to-end encryption when the recipient is offline
Why it's hard. A naive key exchange (Diffie-Hellman) requires both parties online simultaneously. But most messages are sent to someone whose phone is asleep. You need asynchronous key agreement: Alice must be able to encrypt to Bob without Bob participating. And you need forward secrecy — compromising today's key must not decrypt last month's messages — plus post-compromise recovery.
Solution — the Signal protocol: X3DH for asynchronous session setup, Double Ratchet for ongoing messages.
# --- X3DH: Alice establishes a session with an OFFLINE Bob ---
# Bob published to the key server, once, at registration:
# IK_B identity key (long-term)
# SPK_B signed prekey (rotated weekly, signed by IK_B)
# OPK_B one-time prekeys (a batch of ~100, each used once and deleted)
def x3dh_initiate(IK_A, EK_A, bundle) -> bytes:
verify(bundle.SPK_B, bundle.SPK_sig, bundle.IK_B) # prevents substitution
dh1 = DH(IK_A, bundle.SPK_B) # binds Alice's identity to Bob's prekey
dh2 = DH(EK_A, bundle.IK_B) # binds Bob's identity to Alice's ephemeral
dh3 = DH(EK_A, bundle.SPK_B) # forward secrecy from the ephemeral
dh4 = DH(EK_A, bundle.OPK_B) # one-time: protects against replay
return HKDF(dh1 + dh2 + dh3 + dh4) # the shared secret
# --- Double Ratchet: every message gets a fresh key ---
class DoubleRatchet:
def encrypt(self, plaintext):
# Symmetric ratchet: advance the sending chain for THIS message.
self.CKs, mk = kdf_chain(self.CKs)
header = Header(self.DHs.public, self.PN, self.Ns); self.Ns += 1
return header, aead_encrypt(mk, plaintext, associated_data=header)
def decrypt(self, header, ciphertext):
# DH ratchet: a new public key from the peer means a new root key —
# this is what gives POST-COMPROMISE recovery.
if header.dh != self.DHr:
self.skip_message_keys(header.PN) # store keys for out-of-order msgs
self.dh_ratchet(header)
self.CKr, mk = kdf_chain(self.CKr)
return aead_decrypt(mk, ciphertext, associated_data=header)
Three guarantees worth naming precisely, because "we use E2EE" is not an answer:
- Forward secrecy. Each message key is derived and then deleted. Compromising the device today does not decrypt yesterday's messages (assuming they were also deleted).
- Post-compromise security. The DH ratchet mixes fresh entropy on each round trip, so an attacker who stole the state loses access once the parties exchange messages again.
- Asynchrony. One-time prekeys let Alice complete the handshake alone. The server hands out one prekey per request and deletes it — running out means falling back to the signed prekey only, which weakens replay protection, so clients must replenish their prekey batch.
The server never sees a private key, so a full server breach yields ciphertext and metadata. Be explicit about what it does yield: who talks to whom, when, and how often. That metadata is the real privacy surface, and mitigations (sealed sender, where the sender identity is itself encrypted to the recipient) are worth mentioning.
The unavoidable trade-off to volunteer: E2EE means the server cannot do server-side search, server-side spam filtering on content, or server-side backup without a user-held key. Every one of those becomes a client-side problem, and pretending otherwise is the most common mistake on this question.
4.3 Offline delivery and ordering
Why it's hard. A phone is off for a week. Thousands of messages arrive. On reconnect, all must be delivered, exactly once, in an order that makes conversational sense — and a message must never be dropped because the server assumed delivery succeeded when the socket closed mid-write.
Solution — a durable per-device inbox queue, drained by explicit acknowledgement, with client-assigned ordering.
# The inbox is a QUEUE, not an archive: rows exist only until acknowledged.
# CREATE TABLE inbox (
# user_id, device_id, seq bigint, envelope blob, created_at,
# PRIMARY KEY ((user_id, device_id), seq)
# ) WITH CLUSTERING ORDER BY (seq ASC);
async def route(envelope):
for device in await devices_of(envelope.recipient):
seq = await seqgen.next(envelope.recipient, device) # monotonic per device
await inbox.insert(envelope.recipient, device, seq, envelope) # durable FIRST
if node := await registry.lookup(envelope.recipient, device):
await node.deliver(device, seq, envelope) # best-effort push
else:
await push.wake(device) # content-free wake-up
async def on_ack(user, device, up_to_seq):
# Delete only what the client CONFIRMS. A closed socket is not an ack.
await inbox.delete_range(user, device, up_to=up_to_seq)
Two ordering concerns, and they need different answers:
Server ordering is a per-device monotonic sequence number, which makes reconnect trivially resumable: the client says "I have up to 4,171," and the server sends everything after that. Gaps are detectable; duplicates are ignorable.
Conversational ordering is not the server's to decide. Two messages sent from two devices in the same second have no meaningful global order. Clients attach a Lamport-style counter plus their own timestamp and resolve display order locally, breaking ties deterministically by device ID so every participant renders the same sequence. Explaining why the server should not impose conversational order — it does not know the causal relationships and cannot read the content — is a sharp point.
Bound the queue: messages undelivered after 30 days are dropped (with the sender informed if they are still waiting on a receipt). An inbox is a queue with a TTL, not storage, and this is also a privacy property — the server holds as little as possible for as short as possible.
4.4 Groups: one message, five hundred recipients
Why it's hard. Pairwise E2EE means encrypting the message separately for each member: a 500-member group turns one send into 500 encryptions and 500 envelopes, on a phone, on mobile data. Sending a photo caption to a large group would take seconds and meaningful battery. But the server cannot help, because it cannot decrypt.
Solution — sender keys: encrypt the message once with a symmetric group key, and distribute that key pairwise only when membership changes.
# Each member generates their OWN sender key for the group and distributes it
# once, over the pairwise Signal sessions, to every other member.
class SenderKeySession:
def __init__(self):
self.chain_key = os.urandom(32) # ratchets forward per message
self.signing_key = Ed25519.generate() # so recipients can verify the sender
def encrypt(self, plaintext):
self.chain_key, mk = kdf_chain(self.chain_key) # forward secrecy preserved
ct = aead_encrypt(mk, plaintext)
return ct, self.signing_key.sign(ct) # ONE ciphertext for everyone
# Cost model:
# Send a message: 1 encryption, 1 ciphertext -> server fans out to 500 inboxes
# Join / leave / rotate: N pairwise key distributions (the expensive, rare operation)
The trade is explicit and worth stating: you move the O(N) cost from the frequent operation (sending) to the rare one (membership change). Sending becomes O(1) for the client; the server still writes 500 inbox rows, but that is cheap server-side fan-out of an opaque blob, which is exactly the work a server should be doing.
The security consequence to name: when a member leaves, everyone must rotate their sender key, or the departed member can still decrypt future messages. That rotation is N pairwise distributions and it must happen before the next message. Forgetting this is a real vulnerability, and mentioning it demonstrates you understand the protocol rather than the buzzword.
Group membership is server-managed metadata (the server must know who to fan out to), which means the server knows group composition even though it cannot read content. Again — say what leaks.
4.5 Four devices, one account, none of them always on
Why it's hard. A user has a phone, a laptop, a tablet, and a web session. Every device must receive every message and be able to send. But E2EE binds sessions to keys, not accounts, and the phone — historically the key holder — may be off. The old design (desktop as a dumb mirror of the phone) fails the moment the phone's battery dies.
Solution — treat each device as an independent Signal identity, and fan out at the sender.
Each device has its own identity key and its own prekey bundle.
An account is a SET of devices, published to the key server.
Alice (device A1) sends to Bob:
fetch Bob's device list: [B1, B2, B3]
fetch Alice's OWN other devices: [A2] <- so her laptop sees what she sent
encrypt separately to each of B1, B2, B3, A2 <- 4 ciphertexts, 4 sessions
server routes each to its device inbox
Device added: the new device is provisioned by scanning a QR from an existing one,
which transfers the account identity and message history over a
temporary encrypted channel. New Signal sessions are then established
with every contact lazily, on first message.
Device removed: revoke it from the device list; peers stop encrypting to it.
Sender keys for every group must rotate.
Two consequences worth discussing. Sending cost scales with the recipient's device count, not just recipient count — a 500-member group where everyone has three devices is 1,500 sender-key distributions on membership change. That is why device counts are capped in practice.
And history sync is genuinely hard: a newly added device has no past messages, because the server cannot supply them (it never had plaintext). The options are transferring history device-to-device over an encrypted channel at provisioning time, or an encrypted cloud backup whose key the user holds (a passphrase or a hardware-backed key). Both have real usability costs — losing the key means losing the history — and being honest that there is no good answer here, only trade-offs, is a better response than pretending it is solved.
4.6 Receipts, presence, and typing indicators at 3M messages/second
Why it's hard. Every message generates a delivery receipt and often a read receipt, so receipts roughly double or triple message volume. Typing indicators fire on every keystroke. Presence ("last seen") changes constantly for 500M connections. Treating these like messages — durable, acknowledged, queued — would multiply the system's cost several times over for data that is worthless a second later.
Solution — classify by durability requirement and treat each class differently.
| Signal | Durable? | Delivery | Rate control |
|---|---|---|---|
| Message | Yes, until acked | Guaranteed, ordered | None — every one matters |
| Delivery receipt | Yes, but batched | Guaranteed, coalesced | Batch per conversation per few seconds |
| Read receipt | Yes, batched | Guaranteed, coalesced | Batch; a single "read up to seq N" replaces N receipts |
| Typing indicator | No | Best effort, dropped if offline | Throttled to ~1 event / 3s, expires in 5s |
| Presence | No | Best effort, subscription-based | Only pushed to devices with an open chat with that user |
# Receipts collapse: "read up to sequence N" replaces N individual receipts.
async def on_read(user, conversation, up_to_seq):
await receipt_batcher.record(user, conversation, up_to_seq) # overwrite, not append
# Flushed every ~2s: one receipt per conversation, regardless of how many messages.
Presence subscription scoping is the highest-leverage optimisation and the one candidates miss: do not broadcast presence to everyone in a user's contact list. Push it only to devices that currently have that conversation open — typically zero or one. That turns an O(contacts) broadcast on every state change into approximately nothing.
4.7 Media the server cannot read
Why it's hard. A 30 MB video must be end-to-end encrypted, delivered to 500 group members, cached at the CDN for efficiency, and re-downloadable later. But if it is encrypted per recipient, you store and transfer 500 copies. If it is encrypted once, everyone shares a key — which is fine, but the key must travel securely and the CDN must be able to cache the ciphertext.
Solution — encrypt the blob once with a random key, upload the ciphertext, and send the key inside the E2EE message.
# Sender
key = os.urandom(32)
iv = os.urandom(16)
ct = aes_cbc_encrypt(key, iv, media_bytes)
mac = hmac_sha256(mac_key_from(key), iv + ct) # integrity, checked before decrypt
url = await blobstore.put(ct) # server sees random bytes
envelope = {"type": "image", "url": url, "key": key, "iv": iv,
"mac": mac, "sha256": sha256(media_bytes), "size": len(media_bytes),
"thumb": aes_encrypt(key, generate_thumbnail(media_bytes))}
await send_e2ee(recipients, envelope) # the KEY is inside the E2EE message
Three properties make this work. One ciphertext, cached by the CDN — all 500 recipients download the same object, so bandwidth and storage are O(1) rather than O(members). The key rides inside the encrypted message, so only recipients can decrypt, and the server's blob store holds indistinguishable-from-random bytes. And the MAC is verified before decryption, preventing chosen-ciphertext attacks — a small detail that shows crypto literacy.
Include an encrypted thumbnail inline in the message so a preview renders instantly without downloading 30 MB, and expire blobs after ~30 days (recipients who downloaded keep their local copy; the server is not an archive, which is both a cost and a privacy decision).
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| Connection node loss | 2M clients reconnect at once | Full-jitter backoff; session registry TTL expires cleanly; capacity headroom of one node per zone |
| New Year's Eve | 5–10× message rate | Connection tier is the bottleneck, not routing; pre-scale; degrade presence and typing first |
| Large group message storm | 500 inbox writes per message | Batch inbox writes; groups are capped; fan-out is async |
| Prekey exhaustion | Weaker session setup for that user | Clients replenish proactively; server alarms below a threshold |
| Push provider outage | Offline users not woken | Messages stay durably queued; delivered on next foreground; see the notification system |
| Inbox partition hot | One user with enormous backlog | Partition by (user, device); cap and TTL the queue |
| Device added to a big group | N pairwise key distributions | Rate-limit membership churn; batch distributions |
6. Cheat sheet
- Connections: actor-per-connection on BEAM (or goroutines), ~2 KB each, per-process GC, ~2M sockets/node, adaptive keepalive tuned per carrier.
- Crypto: X3DH for asynchronous setup (identity + signed prekey + one-time prekey), Double Ratchet for per-message keys. Forward secrecy and post-compromise security.
- Offline: durable per-device inbox queue, delete only on explicit ack, monotonic seq for resumable reconnect, 30-day TTL.
- Ordering: server orders per device; conversational order is decided by clients with Lamport counters and a deterministic tie-break.
- Groups: sender keys make sending O(1); membership change costs O(N) key distributions, and leaving requires rotation.
- Multi-device: every device is its own Signal identity; sender fans out per device; history sync has no clean answer, only trade-offs.
- Ephemeral signals: typing and presence are best-effort, throttled, and scoped to open conversations only.
- Media: one AES key per blob, ciphertext on the CDN, key inside the E2EE envelope, MAC verified before decrypt.
- The one-liner: "The server is a dumb, extremely well-tuned router for opaque bytes — all the intelligence lives in the clients, which is exactly what end-to-end encryption forces, and every scaling decision follows from keeping 500 million sockets cheap."