Lesson 6 of 8 · 55 min

Worked design — realtime chat & notifications

WebSocket gateways, presence, per-channel ordering, multi-device catch-up, push reliability, and message persistence with Discord’s reported production anchors.

Lesson 6 · Worked design

Realtime chat & notifications

Connections are state

Chat is not a CRUD API with lipstick. You must design a WebSocket (or long-poll) gateway, presence, per-channel ordering, multi-device delivery, push for offline, and a message store that survives trillions of messages. Discord’s public engineering numbers are reported anchors (2023 blog): Cassandra era grew from ~12 nodes (billions of msgs, 2017) to ~177 nodes (trillions, early 2022); migration to ScyllaDB brought the same workload to ~72 nodes with improved p99 (blog cites Cassandra historical read p99 ~40–125 ms and write ~5–70 ms; Scylla steady p99 read ~15 ms / write ~5 ms class). Slack’s WS→Envoy migration (reported) speaks to millions of concurrent connections. Label all as reported, not our metrics. Message edit/delete are events in the same per-channel order: clients apply them by seq. Search indexes (L7) consume the same event stream asynchronously. Never invent a second ordering domain for edits.
MVP scope: 1:1 and group channels, message send/history, delivery/read receipts (eventual), multi-device, push notifications. Stretch: E2E encryption, threads, search (L7). NFRs (example heuristics): p99 send-to-receive for online users < 200ms same-region; durable messages; at-least-once delivery to devices; ordering per channel; 10s of millions concurrent connections class for large products.

Capacity sketch (decision-linked)

python Delivery receipts vs read receipts: delivery means the device acknowledged the frame or catch-up; read means the user opened the channel past that seq — store them separately. Group chats should not rewrite the message row on every reader. For media messages, upload to object storage first, then send a message referencing the blob id so the chat path stays small; CDN the blobs. These choices keep the message store append-optimized and the realtime path thin under load.
1# Heuristic capacity — state assumptions2concurrent_ws = 50_000_0003per_gateway = 50_000                 # heuristic conns/box4gateways = concurrent_ws / per_gateway  # ~1000 gateway boxes class56msgs_day = 100_000_000_000           # 100B/day class product (illustrative)7msg_qps_avg = msgs_day / 86_400      # ~1.2M msg/s avg8msg_qps_peak = msg_qps_avg * 8       # large peak factor9bytes_msg = 50010# storage 5y × 3 replicas → multi-hundred-PB class cold → MUST tier cold to object store1112# DECISIONS FORCED:13# - gateway fleet + sticky/resume sessions14# - partition message log by channel_id (+ time bucket)15# - pub/sub for live fan-out separate from durable store16# - cold tier for history older than weeks/months17# - push path separate with coalesce/throttle

High-level topology

code
1Clients (mobile/web)2  │  WebSocket3  ▼4Gateway fleet  ── session: user/device → conn  · heartbeats · reconnect resume5  │6  ├─► Channel/router service  (who is in channel? which gateways?)7  ├─► Pub/sub / Kafka (channel partition key)8  ├─► Message service → durable store (Cassandra/Scylla-class or equiv)9  ├─► Presence service (often Redis / ephemeral)10  └─► Push worker → FCM/APNs (offline devices)1112History API (HTTP) for catch-up / infinite scroll
WebSocket gateway responsibilities. Terminate TLS/WS, authenticate, maintain conn map, heartbeats, apply backpressure (slow client), route outbound frames. Horizontal scale: any gateway may hold a user’s socket; fan-out uses pub/sub so a message published on channel X reaches all gateways hosting members of X. Sticky LB optional for efficiency, not correctness. After outages, reconnect storms are a first-class failure — jitter resume and rate-limit history pull.
Ordering guarantees. Usual promise: per-channel total order (or per-DM), not global order across all chats. Implementation: partition message log by channel_id; single-threaded virtual consumer per partition; assign monotonic channel_seq on assign. Clients may still see retries — use server seq + client nonce for idempotent send. Discord-style buckets: partition key (channel_id, time_bucket) so a single channel’s history does not create one unbounded partition forever.
python
1def send_message(channel_id, sender_id, body, client_nonce):2	if dedupe.exists(channel_id, sender_id, client_nonce):3		return dedupe.get(...)  # retry-safe4	seq = next_seq(channel_id)  # from partition owner / atomic counter5	msg = Message(id=new_id(), channel_id, seq, sender_id, body, ts=now())6	store.append(msg)           # durable first or after assign — state choice7	pubsub.publish(channel_id, msg)8	dedupe.put(...)9	return msg

Presence, typing, receipts

Presence is lossy by design: heartbeats update last_seen; TTL expires online status. Do not put presence on the critical path of message durability. Typing indicators: ephemeral (Redis TTL seconds), debounced — not persisted. Read receipts for groups: per-user last_read_seq stored separately, aggregated carefully; do not rewrite the message row per reader. Large channels: broadcast storms need aggregation (“N people typing”), not per-keystroke global floods.

Multi-device & offline

Each device has a connection and a push token. Online devices get WS frames; offline devices get push (FCM/APNs) with collapse keys for noisy channels. On reconnect: client sends last_seq per channel → server streams missed messages (catch-up). Multi-device read receipts are eventual; max-seq wins. Push is best-effort via third parties — durable truth is the message store + catch-up.

Persistence — Discord-reported lessons

Design lessons from Discord’s public migration story: messages are append-heavy; partition by channel; bound partitions (time windows / buckets) so hot channels do not create unbounded rows; compaction and tombstones matter; historical cold data can live on cheaper tiers; GC/tail latency can force a storage engine change even when the data model is “right.” Interviewers care that you know why a wide-column/log-ish store fits chat better than a naive relational row-per-message without partitioning.
code
1# Message storage sketch (channel-bucketed)2# PK: (channel_id, bucket, msg_seq)  bucket = e.g. yyyyww or snowflake range3# Access: latest page = max bucket for channel; scroll up = previous bucket45STORAGE CHOICE TRADEOFF6Store                 Pros                          Cons                     Pick when7--------------------  -----------------------------  -----------------------  -------------------8Cassandra/Scylla      write scale, PK lookups        ops/compaction care      high-write chat9Postgres + replicas   familiar, joins, tx            write ceiling            small team chat10Kafka→S3 + search     replay + audit                 cold read latency        regulated history1112# Hot channel mitigation: separate outbox/fanout tier for large member lists

Notifications fan-out & 10× failures

Not every message should push: user prefs, mute, mention-only. Pipeline: message event → notify workers → dedupe → preference check → provider. Debounce noisy group chats. Discord’s GenStage push post (classic, still cited) is a backpressure story under million-per-minute bursts. At 10× QPS: autoscale gateways; shed non-critical presence; more pub/sub partitions; protect store with backpressure; coalesce push — never drop durability first.
Senior close: “Per-channel order, durable append store, WS gateway + pub/sub fan-out, catch-up by seq, push as best-effort, presence as ephemeral.”

Interview answers — chat

  1. 01Exactly-once delivery? → you can’t; at-least-once + client/server dedupe by msg id/nonce.
  2. 02Join channel with 1M history? → don’t backfill all; recent + lazy scroll.
  3. 03Shard by channel not user? → per-channel order needs single writer/partition.
  4. 04Typing without flood? → debounce + short TTL pub/sub.
  5. 05Group read receipts? → per-user counters, not message row rewrites.
  6. 06Offline delivery? → durable store + push wake + catch-up by last_seq.
  7. 07Search messages? → async inverted index (L7), not OLTP LIKE.
  8. 08E2E encryption? → client keys; server metadata limited; different design space.
  9. 09Scale WS gateways? → shard conns, resume tokens, Envoy/L7 routing class.
  10. 10Admin broadcast to 10M? → async fanout workers + badge batching.
  11. 11Why cite Discord migration? → reported anchor for append-heavy storage + p99/ops.
  12. 12Hot channel? → time buckets + fanout tier; throttle degraded mode.
articleHow Discord Stores Trillions of MessagesDiscord EngineeringarticleHow Discord Indexes Trillions of MessagesDiscord EngineeringarticleDiscord — push bursts with Elixir GenStageDiscord EngineeringarticleSlack Engineering — Real-time MessagingSlack EngineeringarticleSlack — Migrating millions of concurrent WebSockets to EnvoySlack Engineering

Message API and history pagination

code
1WS /ws  — auth, subscribe channels, send, ack, resume(last_seq)2POST /channels/{id}/messages  — REST fallback3GET  /channels/{id}/messages?cursor&limit=504POST /channels/{id}/typing56messages PK idea: (channel_id, bucket, seq)7memberships: (user_id, channel_id, last_read_seq)8cursor = (bucket, seq) — scroll upward loads older bucket

Large channels and fan-out amplification

A channel with 10k online members turns each message into 10k gateway deliveries. Mitigations: separate realtime fanout tier, capacity-aware subscription, collapse presence, and for extreme sizes treat as broadcast with server-side fanout workers (similar to celebrity feed). Storage remains one append stream; delivery amplification is the separate problem. Senior answers split history storage tier from live delivery tier.
code
1CHAT SENIOR SIGNALS2MID FAILURE                              SENIOR SIGNAL3--------------------------------------  ------------------------------------------4messages are just rows                   per-channel order is first-class5one fanout approach for DM and 10k room  different delivery strategies6mix history store with live pubsub       separate tiers, shared message id/seq7push = delivery                          push = wake; store = truth8ignore E2EE when asked                   E2EE changes metadata visibility9just add a column on 1T messages         migration/backfill strategy
E2E encryption (WhatsApp/Signal class) when asked: client holds keys; server stores ciphertext; metadata (timestamps, group membership) still leaks unless sealed sender style designs; key rotation and multi-device are the hard parts. Do not derail MVP if not required — mention as a fork that changes server visibility.

Reconnect, resume, and catch-up storms

Session resume token maps to user/device and last acknowledged seqs. After a blip, millions of clients reconnect — exponential backoff with jitter is mandatory. Catch-up should be rate-limited per user and prefer push of missed seq ranges over full history scans. Gateways shed non-critical events (typing) under load while preserving message frames. Observability: gateway connection count, send p99, store write p99, pubsub lag, push provider error rate, catch-up QPS after incidents. Cardinality: label by channel size tier, not channel_id, on hot metrics.

Checkpoint

What ordering guarantee should you claim by default for group chat?

ATotal order across all channels product-wideBPer-channel (or per-DM) ordering via partition/seq; no global orderCNo ordering — clients sort by local clock only
Sign up free to answer and see why

Checkpoint

User has phone offline and laptop online. Message arrives. Correct path?

AOnly save on the phone later; laptop can miss itBLaptop gets WS delivery; message durable in store; phone gets push + catch-up on reconnect by last_seqCPush both devices only; no store
Sign up free to answer and see why

Checkpoint

Why mention Discord’s Cassandra→Scylla migration in an interview?

ATo recite node counts as if they were universal lawsBAs a reported production anchor that message storage is append-heavy, partition-sensitive, and latency/ops-driven — then map lessons to your designCBecause every chat must use Scylla
Sign up free to answer and see why

Checkpoint

Send API is retried after a gateway blip. How do you avoid duplicate messages in-channel?

ATell users not to retryBClient nonce + server dedupe per (channel, sender, nonce) before assigning seqCUse random message text as the key
Sign up free to answer and see why

Checkpoint

QPS jumps 10× after a viral event. First healthy degradation?

ADrop message durability to memory-onlyBAutoscale gateways; coalesce push; rate-limit catch-up/presence; protect store write path with backpressureCDelete old messages live
Sign up free to answer and see why

End-to-end narration script (15 minutes)

Frame: 1:1 + group chat, multi-device, push offline, per-channel order. Estimate: concurrent WS drives gateway count; message QPS drives store; history retention forces cold tier — say the PB-class warning with replication so you do not pretend one Postgres holds forever. Design: gateway, pub/sub, message service, bucketed store, presence ephemeral, push workers. Deep dive: send idempotency + seq assignment + Discord-reported storage lessons (node counts and p99 as reported). Close: reconnect storms, large-channel fanout, 10× degradation order (preserve durability).
code
1CHAT DECISION LOG2Signal                      Forces3--------------------------  ------------------------------------4Per-channel order           partition by channel_id (+ bucket)5At-least-once networks      client nonce dedupe6Multi-device                 durable store + per-device push/WS7Offline                     push best-effort + catch-up by seq8Long retention              hot/cold tiers; object store cold9Large rooms                 delivery fanout tier ≠ storage tier10Outage recovery             jittered resume; rate-limit catch-up1112Reported anchor: Discord Cassandra→Scylla ~177→72 nodes; p99 improved (2023 blog).
Slack’s public realtime and incident posts are useful for operability color: long-lived connections stress L4 load balancers; migrating to Envoy-class L7 (as reported) is a sticky-connection scale story. Use it to justify gateway as a first-class subsystem, not “just our API servers with WS enabled.”

Chat math + decision card

code
1HEURISTIC WORKED EXAMPLE250M concurrent WS / 50k per gateway → ~1000 gateways3100B msgs/day → ~1.2M msg/s avg (illustrative product scale)4500B bytes/msg → multi-PB retention with RF → cold tier required5THEREFORE: gateway fleet, channel partitions, hot/cold storage, push separate67REPORTED ANCHORS (Discord 2023 blog — not our metrics)8Cassandra ~12 nodes (2017) → ~177 nodes (trillions msgs early 2022)9Scylla migration → ~72 nodes; p99 read ~15ms / write ~5ms class steady10partition key (channel_id, time_bucket) bounds hot channels1112GUARANTEES13order: per-channel not global14delivery: at-least-once + nonce dedupe15push: best-effort wake; store is truth16presence/typing: ephemeral TTLs1718SEND PATH19dedupe nonce → assign seq → durable append → pubsub → gateways20REST history for catch-up with cursor (bucket, seq)212210x DEGRADE ORDER23scale gateways → coalesce push → rate-limit catch-up/presence24→ backpressure sends → NEVER drop durability first2526LARGE ROOM27storage still one stream; delivery fanout is separate amplification problem

Can you design chat with WS gateway, per-channel order, multi-device catch-up, and a durable store story with cited production anchors?

New to itGetting thereConfident

Takeaways

  • Gateway + pub/sub + durable message service is the spine.
  • Order per channel via partitions/seq; idempotent sends with nonces.
  • Presence is ephemeral; push is best-effort; store is truth.
  • Discord blog numbers are reported anchors for storage/latency lessons.
  • 10×: scale gateways, backpressure, coalesce notify — don’t drop durability.

Next: two high-frequency component designs — distributed rate limiting and search/inverted index.

Sources

Free to read · better with Enzo

Learn it with Enzo

Save your progress, answer the checkpoints, and let Enzo quiz you on what you just read.