WhatsApp / Messenger
Worked solution

WhatsApp / Messenger — a worked solution

Hundreds of millions of long-lived sockets, sub-second 1:1 + group delivery, E2E-encrypted, multi-device, multi-region active-active.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the WhatsApp / Messenger workspace

The problem

Build the production reference architecture for a 1:1 + group messaging service at hyperscale (think WhatsApp / Messenger / Discord / Signal). The shape: hundreds of millions of long-lived TLS sockets, sub-second 1:1 + group delivery, end-to-end encrypted with multi-device fanout, multi-region active-active with regional-failover RPO=0 on metadata. Every component below is specified well enough that a staff SRE could file the implementation tickets — and would defend the choices in a 5-year incident retro.

The reference architecture

Reference architecture for WhatsApp / Messenger: 16 components — Mobile / Web Client, Anycast L4 LB + WAF, Media CDN, Chat Gateway (chatd, WS), Control API Gateway, Kafka (outbox + abuse), Identity + Key Transparency, Routing & Fanout Service, Inbox Cache, Idempotency / Dedup KV, Message Store, Metadata DB (accounts, groups, ACLs), Media Object Store, Abuse / Safety Worker, Outbox CDC Relay, Recipient Device (WS) — connected by 28 flows.Mobile / Web ClientiOS, Android, Web (libs…Anycast L4 LB + WAFKatran / GLB (BPF) + En…Media CDNFastly / CloudFrontChat Gateway (chatd, …Erlang/OTP + Cowboy WS,…Control API GatewayEnvoy + WAF + per-user …Kafka (outbox + abuse)Kafka 3.x (KRaft mode)Identity + Key Transp…Custom + libsignal prek…Routing & Fanout Serv…Rust / ElixirInbox CacheRedis 7 Cluster (Stream…Idempotency / Dedup KVRedis 7 Cluster (SETNX …Message StoreScyllaDB 6.x (shard-per…Metadata DB (accounts…Postgres 15 + Patroni, …Media Object StoreS3 (multi-bucket, rando…Abuse / Safety WorkerFlink job, fingerprint …Outbox CDC RelayScyllaDB CDC consumer +…Recipient Device (WS)iOS, Android, Web (libs…
16 components, 28 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Mobile / Web ClientiOS, Android, Web (libsignal)

Holds the user's long-lived WebSocket to chatd. Generates a UUIDv7 message_id per send, performs all libsignal X3DH + Double Ratchet encryption locally for 1:1 and Sender-Keys for groups, and decrypts inbound envelopes on-device. Reconnects with jittered backoff after a drop and replays unacked sends with the same message_id.

Why it exists. E2EE puts the cryptographic boundary on the device — the server cannot decrypt content. The client also amortizes connection cost across every chat the user is in: one WS, not one per conversation. Without a stateful client, every send would re-handshake TLS and re-establish per-conversation key state.

When it fails. Phone offline → messages queue in the recipient's inbox (≤30d) until reconnect; on reconnect the client replays unacked sends. Client crash mid-encrypt → retry with same message_id is caught by idem (SETNX miss returns the original ack). Server drops the WS with reconnect-after: jittered(60–600s) so reconnects don't thunder.

Anycast L4 LB + WAFKatran / GLB (BPF) + Envoy WAF

Terminates the user's BGP-anycast IP and forwards L4 packets at line rate via XDP/eBPF (Katran-style) to a chatd shard chosen by Maglev hashing on the connection's 4-tuple. Per-IP token-bucket rate-limit and WAF rules on /v1/* control APIs; the WS upgrade is pure passthrough so the long-lived socket never re-traverses the LB after handshake.

Why it exists. Terminating TLS on the LB would double the handshake cost at 900M sockets — every TLS handshake's keypair and session-ticket cache would live both on the LB and on chatd. L4-only keeps per-connection state on chatd where it belongs and lets the LB do nothing more than ECMP+Maglev forwarding.

When it fails. BGP withdrawal storm (WhatsApp Oct 4 2021 — backbone routes pulled) → full global outage. Mitigation: out-of-band management plane that doesn't itself depend on the data-plane DNS/BGP, multi-provider DNS, peering diversity, and a tested 'restore from console' runbook.

Media CDNFastly / CloudFront

Edge-cached media GET gated by a signed URL (HMAC over blob_id + expiry + recipient_account_id). Each blob carries a surrogate key = blob_id, so 'delete for everyone' purges only that key globally instead of bumping a TTL and waiting. Stale-while-revalidate on miss so a single slow origin fetch doesn't cliff p99.

Why it exists. 525 PB/yr media egress is unsustainable from origin — even at S3's class pricing the egress line item dwarfs storage. The CDN absorbs ~95% of media reads at a global PoP within a few hops of the user; origin only sees cold blobs and recently-purged keys.

When it fails. Mass purge storm on a viral 'delete for everyone' wave → origin fetch spike → S3 throttle. Mitigation: staged purge (batch surrogate keys), S3 multi-bucket sharding (caps blast radius to one bucket), origin shielding so only one PoP fetches per blob.

Chat Gateway (chatd, WS)Erlang/OTP + Cowboy WS, BoringSSL, Noise framing

Holds 500K stable WS sockets per box, terminates TLS, and runs the per-session sequence cursor. On send: dedup-checks idem, writes the durable row to msgstore, publishes msg.outbox to Kafka, and publishes abuse.signals envelope-meta. On inbound XADD into a recipient's inbox: reads the envelope and writes a deliver frame back through edge-lb to the recipient's WS.

Why it exists. Long-lived stateful socket termination is fundamentally different from stateless HTTP — needs a dedicated tier with consistent-hash routing, lame-duck draining, cooperative load shedding, and per-session memory. Conflating with api-gw couples deploy risk: a control-plane spike or bad auth deploy would take chat down too.

When it fails. Reconnect storm on hard drain (Slack Jan 2021 mTLS thundering herd) — the only mitigation is the protocol: lame-duck reconnect-after: jittered(60–600s) header + connection-count token bucket at edge-lb. Crash mid-publish is closed by cdc-relay (failure mode #15).

Control API GatewayEnvoy + WAF + per-user token bucket

Stateless HTTPS gateway for the control plane: group create/membership ops, profile changes, presigned media-upload URLs, prekey / KT proof requests. Authenticates the bearer token, applies per-token + per-IP rate limits, and forwards to auth or meta-db over mTLS with circuit-breakers in front of both.

Why it exists. Control-plane traffic is bursty, stateless, and CPU-bound — a different operational shape from the long-lived stateful chatd. Separating the two means a control-plane storm or a bad auth-deploy can't take chat down. Conflating them was rejected explicitly in §hld.

When it fails. Bad config push or bad deploy (Slack Feb 2022 — config regression hit prod with no canary) → 50% errors. Mitigation: staged rollout 1%→10%→100% with a 10-min soak between stages, kill-switch on out-of-band etcd that doesn't depend on the broken plane.

Kafka (outbox + abuse)Kafka 3.x (KRaft mode)

Single Kafka cluster shared across two logical topics — msg.outbox (compacted, the durable-fanout source of truth) and abuse.signals (90d, envelope-meta only). 1024 partitions hash-keyed by conversation_id so per-conversation FIFO lands on a single partition; RF=3 with sync replication.

Why it exists. Decouples the durable write from the fanout work. Without an outbox, every send would block on the slowest of N recipient inbox writes — coupling availability of every storage shard to a single user's latency. With it, fanout autoscales independently and a slow inbox shard never bleeds into send→ack.

When it fails. Consumer-group rebalance during deploy → paused consumption window. Mitigation: cooperative-sticky + static membership (consumers keep their partitions across restart). Cluster-wide failure (ZK-era was the 2018 outages, KRaft mitigates) → outbox events queue at producers; cdc-relay catches anything that wasn't published.

Identity + Key TransparencyCustom + libsignal prekey bundles + KT Merkle log

Owns the E2EE trust root: libsignal prekey bundle storage and rotation, per-device identity-key serving, and the Key Transparency Merkle log — including signed root publication and per-account proof generation. Reads/writes the accounts, devices, and kt_log tables in meta-db.

Why it exists. Identity is the cryptographic trust root for the whole E2EE system. KT must be auditable so a compromised server can't silently swap a user's identity key; prekeys must be served per-device (not per-account) so multi-device fanout works.

When it fails. KT log inconsistency → user-visible identity-key change without proof. Mitigation: signed root + client-side audit on every contact-key fetch, with the gossip layer (cross-user proof exchange) listed as an open question in §tradeoffs.

Routing & Fanout ServiceRust / Elixir

Consumes msg.outbox from Kafka. For each event: looks up the recipient device list in meta-db and XADDs an envelope into every recipient device's inbox stream in Redis. Idempotent on (message_id, recipient_id) so replays from the chatd-direct path or from cdc-relay are safe.

Why it exists. Fanout is the heaviest work in the system — 9.4M deliveries/sec at peak, with N=avg-fanout 3.5× per send. It absolutely cannot run on the sender's hot path: the send→ack budget is 400ms, fanout latency is 1.5s, and they're on different SLO contracts. Decoupling lets fanout autoscale on Kafka lag while send→ack stays bounded.

When it fails. Kafka consumer lag spike on a mega-group's partition → head-of-line blocking on every other conversation hashed there (Discord 2017 #general meltdown). Mitigation: sub-shard whales onto a dedicated route consumer pool with reserved Kafka partitions; per-conversation send rate-limit at chatd as a backstop.

Inbox CacheRedis 7 Cluster (Streams)

Per-(user, device) XSTREAM inbox:{user, device} of pending envelopes (TTL 30d, XDEL on ACK so the working set stays small). Recipient chatd shards XREADGROUP BLOCK on their assigned inbox shards; route XADDs an envelope; Redis wakes the subscription; chatd writes the WS deliver frame.

Why it exists. Recipients must receive a delivery within sub-second when online — pulling envelopes from msgstore per delivery would burn ScyllaDB and bottleneck on its quorum reads. Inbox is a per-user mailbox optimized for tail-write + tail-read on hot data, with the sub/pub semantics WS server-push needs.

When it fails. Primary fails over → semi-sync loses ≤1s of XADDs → route at-least-once replays from Kafka → recipient receives duplicate deliver frames; client dedupes by message_id. Alert on redis_failover_count and kafka_consumer_replay_rate correlated.

Idempotency / Dedup KVRedis 7 Cluster (SETNX + WAIT 1, AOF, 7d TTL)

Hot-path message_id → {status, original_ack} set by the sender's chatd before the durable write. SETNX returns the original ack on a duplicate retry; the durable msgstore write happens only on first arrival. 7-day TTL window so a client retry through airplane-mode + region failover still resolves to its original ack.

Why it exists. Without idempotency, a client retry through region failover (different chatd box, different Postgres-resolved auth) would double-write to msgstore, double-fanout to recipients, and surface a duplicate message UI. The product claim 'message arrives once' depends on this single hot-path SETNX.

When it fails. Primary loses last AOF segment + fails over → small dedup gap during which a retry can double-apply. Mitigation: AOF every-write (not every-second); the gap is ≤1s and clients typically retry slower than that. Alerting on redis_failover_during_aof_write_ratio.

Message StoreScyllaDB 6.x (shard-per-core)

Wide-column store of every accepted message envelope, partitioned by (conversation_id, day_bucket) and clustered by monotonic per-conversation seq. RF=3 with quorum W=R=2 in-region + one async cross-region observer per partition for DR. Holds the durable record for both 1:1 and group sends; ScyllaDB CDC log feeds the cdc-relay outbox closer.

Why it exists. 100B msgs/day × 1KB/row × RF=3 = 36 PB/year of small immutable rows with sequential write + tail-read access patterns — the canonical wide-column shape. A single-leader RDBMS would need painful manual sharding at this volume; serializable cross-shard transactions are not what messaging needs.

When it fails. Hot mega-group head-of-line blocking on its partition (Discord 2017) → reads timeout, coordinator stalls. Mitigation: time-bucketed partition key (caps width); per-conversation send rate-limit at chatd; sub-shard whales via dedicated route consumer pool with reserved Kafka partitions.

Metadata DB (accounts, groups, ACLs)Postgres 15 + Patroni, Citus sharding

Sharded relational store for accounts, devices, groups, group_members (with role + last_seq_seen), the abuse-disabled flag, and the KT log roots. 32 hash shards on account_id × (1 leader + 2 sync followers) per shard. Serves chatd's ACL check on send, route's recipient device-list expansion, and api-gw's group ops.

Why it exists. Membership has strict invariants (alice can't be both a member and not a member at the same logical instant) and is small enough (~8TB) that a sharded relational store with serializable transactions is the right shape. Wide-column or KV would force the same invariants into the application — a famous source of bugs at messaging scale.

When it fails. Both followers in an AZ trip → sync-quorum unmet → writes block. Mitigation: ANY-1-of-2 fallback documented as the deliberate availability/RPO trade-off (RPO ≤1s during recovery). Cross-region failover is the only paged operation — admins confirm because the cross-region replica is async (RPO ≤30s).

Media Object StoreS3 (multi-bucket, random key prefix)

Stores opaque AES-256-GCM ciphertext blobs uploaded via presigned PUT through api-gw. Multi-bucket layout (media-{0..31}) with random 4-char key prefix per blob. Versioning on so 'delete for everyone' produces a delete-marker; lifecycle rule expires prior versions ≤30d to satisfy GDPR.

Why it exists. Media at 525 PB/year is fundamentally an object-store problem, not a database problem. E2EE means the server stores opaque bytes — there's nothing to query, nothing to index, nothing to inspect. Putting media in a DB would burn money for no gain.

When it fails. S3 SlowDown (503) on a hot prefix — typically a viral upload re-fetched from origin during CDN purge. Mitigation: random prefix already in place; client retries with exponential backoff; multi-bucket caps the blast radius.

Abuse / Safety WorkerFlink job, fingerprint + reputation scoring

Flink job that consumes abuse.signals from Kafka — envelope metadata only (sender, recipient list cardinality, fanout count, timestamp, message_id). Reputation-scores per sender against a sliding-window model; flagged accounts get abuse_disabled=true written back to meta-db, which chatd's ACL check on the send path picks up.

Why it exists. E2EE forbids server-side content scoring, but envelope-meta + behavioral patterns catch >95% of spam waves and account-takeover-driven fanout abuse. It must run as a decoupled consumer group so abuse pipeline lag never blocks user-facing fanout (a poison message here cannot pause route).

When it fails. Poison message crashes worker in tight loop. Mitigation: DLQ + max-deliveries=5 + schema validation at producer; alerting on abuse_consumer_lag_seconds > 30m (business-hours SLO, not user-facing page).

Outbox CDC RelayScyllaDB CDC consumer + Kafka producer

Tails the ScyllaDB CDC log per msgstore partition, computes the message_id for each committed row, and republishes any msg.outbox event to Kafka that hasn't been seen there yet. Consumes both the durable commit log and Kafka's outbox topic to compute the diff; idempotent so duplicates from the chatd-direct path collapse on the route consumer's (message_id, recipient_id) dedup.

Why it exists. A bare txn.group: outbox label on edges e7 (durable commit) and e9 (kafka publish) is aspirational — chatd crashing between commit and publish would silently lose fanout for that message. CDC is the only mechanism that closes the outbox transactionally under chatd failure: the commit log is the source of truth and the relay is the eventually-consistent publisher.

When it fails. Relay falls behind → fanout latency spike for any chatd-crashed-mid-send events. Mitigation: alert on cdc_relay_lag_seconds > 30; auto-scale consumer count for catch-up; the user-facing impact is bounded because the typical case (chatd healthy) doesn't depend on cdc-relay at all.

Recipient Device (WS)iOS, Android, Web (libsignal)

A different end-user's device holding its own long-lived WebSocket to its assigned chatd shard (consistent-hashed onto a different box from the sender's, often in a different region). Receives deliver frames pushed by chatd, decrypts them locally with its per-device Sender Key, surfaces the message in the UI, and sends a delivered ACK back over the same socket — chatd then XACKs the inbox stream entry, which XDELs the envelope.

Why it exists. Drawn explicitly as a separate node because the diagram's most user-visible hop — 'message arrives on the other person's phone' — was hidden inside a bidirectional WS edge in earlier drafts. Modeling sender and recipient as separate mobile nodes makes the end-to-end delivery flow traceable in the simulator and unambiguous to a reader.

When it fails. Recipient socket dies between deliver and ACK → on reconnect, chatd reads pending from inbox and re-delivers; client dedupes by message_id (the same idempotency window that protects the sender). Phone offline beyond the inbox TTL window (30d) → message lost on the server side; sender shows 'undelivered' state.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Surface these before drawing a single box:

  • E2EE in scope? Yes — Signal Protocol (X3DH + Double Ratchet) for 1:1 and Sender Keys for groups. The server only ever sees opaque ciphertext envelopes. This rules out server-side message search and constrains abuse pipelines to envelope metadata.
  • Multi-device? Yes — primary + N linked devices. Each device has its own identity key; fanout is per-device, not per-account.
  • Calls (voice/video)? Out of scope for this canonical (signaling shares the WebSocket, but media plane is a separate SFU + TURN tier — point at it but don't unify).
  • Server-side history? Yes for undelivered (≤30d), no for delivered E2EE messages — those live on devices. Backups are device-driven (encrypted to user-held key).
  • Group size cap? Hard cap at 1024 members for fanout-on-write; communities/channels above that flip to fanout-on-read.
  • Disappearing messages? Yes — server TTL of 24h / 7d / 90d, enforced by ScyllaDB TWCS.
  • Multi-region? Active-active across at least 18 regions. User home-region pinned on account creation, follow-the-sun fallback for outage.

Assumptions baked into the canonical (defaults — see capacity walk-through for the math):

  • 1.5B DAU, 2B MAU (hyperscale tier, WhatsApp-class).
  • 67 messages / DAU / day → 100B msgs/day.
  • Avg fanout ~3.5× (most chats are 1:1 or small).
  • Peak / avg multiplier 2.3× (regional lunch + evening waves).
  • 60% of DAU online at peak → 900M concurrent sockets globally.
  • Hot-tier retention: 30 days for undelivered + delivered metadata; cold tier for legal-hold.

02Functional reqs

What must this system actually do?

  • Send a message to a 1:1 conversation or group of ≤1024 members.
  • Deliver in real time to recipients with an active socket — server-initiated WS push from the recipient's chatd shard down through edge-lb to the recipient device, with delivered ACK back over the same socket. Queue for ≤30 days for offline recipients.
  • Multi-device delivery: every linked device of every recipient gets its own ciphertext envelope.
  • Idempotent retry: client-generated message_id (UUIDv7); duplicate sends collapse to a single record.
  • Fetch history (last N pages of a conversation) on chat open.
  • Upload / download media (image, voice, video) via opaque AES-256-GCM blobs in S3-compatible storage, distributed via signed URLs.
  • Group operations: create, add/remove member, change admin, leave; key rotation on membership change.
  • Delete for everyone: tombstone replicated to all recipients' inboxes, including offline ones (replay-then-delete on reconnect).
  • Key transparency: proof on identity-key change, so a compromised server can't silently swap a user's identity key.
  • Abuse pipeline: envelope-metadata-only signal stream feeds spam/safety scoring; flagged accounts get abuse_disabled=true.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Availability: 99.99% on the WS data path (52 min/yr error budget); 99.95% on control APIs.
  • Latency (same-region recipient online):
  • Socket connect (warm, TLS resumption): p50 30ms / p99 120ms / p99.9 400ms.
  • Send → server ack: p50 60ms / p99 180ms / p99.9 400ms.
  • Send → recipient delivered (perceived): p50 120ms / p99 350ms / p99.9 800ms (sub-second is the load-bearing product claim).
  • Cross-region recipient: add ~80ms RTT to all of the above.
  • Fetch-history page (50 msgs): p50 80ms / p99 250ms / p99.9 600ms.
  • Durability: zero loss on accepted (ack'd) messages. In-region RPO=0 on metadata (sync replication, with sync-quorum-of-1-of-2 fallback under degraded follower set — see failure mode #16); cross-region RPO ≤30s on metadata (async). RPO ≤1s on dedup (Redis AOF every-write fsync + WAIT 1, best-effort); RPO ≤1s on inbox cache (semi-sync — losing 1s of ACKs replays at most a few duplicates per user, which the client dedupes by message_id). RPO ≤30s on msgstore cross-region (1 async observer per partition).
  • Consistency:
  • Per-conversation FIFO, enforced by a sequencer in the partition's owning shard.
  • Read-your-writes for the sender on ACL/membership reads (route via meta-db leader).
  • Eventual cross-region for everything else; cap at 30s replication lag with a paging burn-rate alert.
  • Scalability: every tier horizontal except the per-user inbox stream and the per-conversation sequencer (each is naturally single-writer per key — so they scale by sharding the keyspace, not the writer).
  • Security: E2EE on message content; mTLS between every internal service; per-token + per-IP rate limits at api-gw; key transparency on identity keys; envelope-metadata abuse pipeline (never decrypts).
  • Compliance: GDPR erasure ≤30d, cascading across every store that holds the user's identifiers: meta-db (account row delete), inbox (drop pending streams), msgstore (insert tombstone with TTL), media-store (object delete + version expiry via lifecycle rule ≤30d), msg.outbox Kafka topic (compaction with cleanup.policy=compact,delete and per-account-id null tombstone), abuse.signals (90d retention then expiry). Erasure is eventually atomic — the SLA is "all stores reach consistent erased state within 30 days," and a verification job audits this monthly. India IT Rules 2021 traceability via metadata hash chain (no plaintext); EU DMA interop logs ≥6mo.

04Capacity estimation

How much load and data does this have to hold?

Capacity walk-through. Hyperscale tier defaults — every shard count below is derived, not chosen.

Send QPS at peak

peak_send_qps = DAU × msgs_per_DAU_per_day × peak_factor / 86400
              = 1.5B × 67 × 2.3 / 86400
              ≈ 2.7M sends/sec

Fanout QPS

peak_fanout_qps = peak_send_qps × avg_fanout
                = 2.7M × 3.5
                ≈ 9.4M deliveries/sec

Concurrent WS sockets

sockets = DAU × online_ratio = 1.5B × 0.6 = 900M
sockets_per_box (i4i.4xlarge, Linux + BoringSSL TLS termination):
  RAM ceiling     = 256GB / 50KB ≈ 5M (raw memory only)
  CPU + TLS-state = ~500K (epoll wait, OpenSSL session cache, BPF egress filters)
  → 500K is the load-bearing practical ceiling; the 1M number assumed
    L4-only TLS offload (which we deliberately rejected).

→ 1800 chatd boxes globally, spread across 18 regions (≈100/region) with 2× headroom to absorb a single-region failover. WhatsApp's 2014 "2M/box" number was FreeBSD + custom TCP stack + zero TLS — our shape (Linux + TLS at the box) lands at 500K. The replicas in the diagram reflects this.

chatd shard count

Each chatd box holds 500K sockets; 1800 boxes → 1800 shards in the connection-balancer's consistent-hash ring. Lame-duck rebalance is jittered over 60–600s — without that, the reconnect storm in §failure-modes 1.1 is a P1 within 90s.

Message store sharding

ScyllaDB partition key (conversation_id, day_bucket) where bucket = 7-day window (Discord uses 10d; we tighten because mega-groups grow faster). 256 vnodes/node × 940 nodes (RF=3). Why 940? With RF=3 every logical write becomes 3 physical writes:

physical_write_qps = peak_send_qps × RF = 2.7M × 3 = 8.1M phys writes/s (sender-only)
+ fanout_phys_qps  = 9.4M × 3 / fanout = ... handled at fanout (consumer side)
sustainable per i4i.4xlarge node ≈ 30K writes/s
nodes = 8.1M / 30K = 270 minimum for sender-side ingest

Plus reads + compaction headroom (×3.5) ≈ 940 nodes. The earlier headline "9.4M deliveries/sec" is the fanout number — those land in inbox (Redis), not msgstore. ScyllaDB ingest is bounded by unique sends (2.7M/s), not fanout.

Storage growth

row_size = 200B logical + 80B metadata + 40B index = 320B × RF3 ≈ 1KB effective
text_per_year = 1.5B × 67 × 365 × 1KB = 36 PB

Hot tier (NVMe, 0–30d) ≈ 3 PB; warm (HDD/EC8+3, 30–365d) ≈ 33 PB; cold (Glacier, archive) for compliance.

Media bytes

media_per_year = sends × media_ratio × avg_blob_size × 365
               = 100B × 0.08 × 180KB × 365
               ≈ 525 PB / year

Stored at RF=2 with erasure coding in S3. CDN egress dominates cost (3–5× the storage line item).

Inbox cache (Redis Streams)

The inbox stream holds only undelivered envelopes — once a recipient device XACKs, the entry is XDEL'd. The 30-day TTL applies to envelopes that were never delivered (the recipient was offline beyond their phone's typical reach). Working set ≈ 5 unread/online user × 900M sockets × 200B ≈ 900GB across 256 Redis primaries (3.5GB/primary). The full retained set (offline-recipient pending, capped by TTL) is bounded by offline_users × avg_unread_at_30d × 200B — empirically <2× the working-set figure, so 8–10GB/primary peak. Without XDEL on ACK, the math would be 6PB and this design wouldn't fit Redis at all — that XDEL is load-bearing.

Idempotency cache

Working set = message_id → status for last 7d sends. 100B × 7 × 32 bytes/entry ≈ 22TB hot (7-day window). Sharded over 64 Redis primaries (350GB each, AOF persistence).

Kafka

write_throughput = 9.4M fanout msgs/s × 1KB envelope = 9.4 GB/s
read_throughput  = 9.4M × 2 consumer groups (route + abuse) = 19 GB/s
partitions       = 19 GB/s / (50 MB/s per partition for headroom) = 380 minimum

We use 1024 partitions total, sized for the read multiplier on consumer groups plus per-conversation FIFO headroom, not just write throughput. RF=3, sync replication, KRaft mode (ZooKeeper retired in Kafka 3.x). Two logical topics share the cluster:

topicretention policyreason
msg.outboxcleanup.policy=compact,delete; retention.ms=24hcompaction collapses message_id-keyed events on retry; deletion sweeps stale tombstones
abuse.signalsretention.ms=90dabuse pipeline replay window; envelope-meta only

The route consumer's per-message budget is 100ms (consume + meta-db lookup + N×XADD inbox). Lag SLO: p99(consume_to_xadd_latency) < 500ms.

Postgres metadata

Account+membership row size ~500B, 2B accounts × 500B = 1TB; with group memberships and device lists ≈ 8TB. Sharded into 32 shards × (1 leader + 2 sync followers); each shard ≈ 250GB, comfortably hot-cached.

05API design

What does the outside world call, and what comes back?

Chat protocol is a custom binary frame over WebSocket with a Noise-derived handshake; here we sketch the wire shape as JSON for clarity.

client ⇄ chatd  (ws://wss-edge/v1/chat)

C→S { "op": "send", "conversation_id": "...", "message_id": "uuidv7", "envelope": "<ciphertext>", "ts": 1730000000 }
S→C { "op": "ack",  "message_id": "uuidv7", "seq": 18234, "server_ts": 1730000001 }
S→C { "op": "deliver", "from": "...", "conversation_id": "...", "message_id": "...", "envelope": "<ciphertext>", "seq": 18235 }
C→S { "op": "delivered", "message_id": "..." }
C→S { "op": "fetch", "conversation_id": "...", "before_seq": 18234, "limit": 50 }
S→C { "op": "history", "conversation_id": "...", "messages": [ { "message_id": "...", "seq": 18234, "envelope": "<ciphertext>" } ] }

History pages over the open socket (chatd checks the ACL in meta-db, then paginates the conversation partition out of msgstore) — the control plane has no read path to the message store.

Control APIs (HTTPS to api-gw):

POST /v1/groups                { name, members[] }
POST /v1/groups/:id/members    { add[], remove[] }
POST /v1/media/upload-url      { content_length } -> { presigned_put_url, blob_id }
GET  /v1/keys/:account_id      -> { device_list_signed, bundles: [ { device_id, identity_key, signed_prekey, one_time_prekey } ] }
# X3DH prekey fetch: the sender pulls the recipient's per-device bundles before the first
# 1:1 message. The server CONSUMES one one-time prekey per device per fetch (falls back to
# the signed prekey when the one-time pool is drained). Reads accounts+devices; rate-limited.
POST /v1/keys/transparency/proof { account_id } -> { merkle_proof }

Idempotency: every state-changing request carries Idempotency-Key (separate from message_id for retries at the API layer). Server stores idem_key → response for 24h.

06Data model

What gets stored, and what is it looked up by?

Postgres meta-db (sharded by account_id) — leader+2 sync followers per shard, 32 shards.

tablekeynotes
accountsaccount_id PK, phone uniqueidentity_key, created_ts, abuse_disabled
devices(account_id, device_id) PKplatform, identity_key, linked_ts
groupsgroup_id PKcreated_by, name, member_count
group_members(group_id, account_id) PKrole, added_ts, last_seq_seen
kt_logkt_seq PKmerkle hash, signed root, append-only

Why Postgres + sync replication: small data (8TB), strict membership invariants ("alice can't be removed AND post in the same logical order"), need serializable transactions for group ops. Failover RPO=0 because losing a recently-added member produces "ghost messages" on the next send.

ScyllaDB msgstore — partition (conversation_id, day_bucket), RF=3, leaderless quorum.

columntypenotes
conversation_iduuidpartition
day_bucketintpartition (rolls every 7d)
seqbigintclustering, monotonic per-partition
message_iduuidunique in partition
sender_iduuid
envelopeblobE2EE ciphertext, ~200B–4KB
ts_servertimestamp

Why ScyllaDB: shard-per-core, predictable p99, no JVM compaction pause. Partition key bucketed by time so a mega-group's partition rolls every 7 days, capping width and tombstone load. Quorum reads/writes (W=R=2, N=3) for in-region; one observer in a second region for DR. LWT only on the rare conversation-rename / pin-message paths — never in the per-message hot path.

Redis inbox — per-user XSTREAM inbox:{account_id, device_id}; entries are (seq, envelope). TTL = 30 days. Cluster mode with 256 shards, leader-follower with semi-sync replication.

Redis idem — SETNX message_id:{id} status with 7-day TTL. 64 shards, leader-follower with AOF every-write fsync + WAIT 1 (best-effort); RPO ≤1s on primary loss — Redis has no true synchronous replication.

S3 media-store — buckets named media-{0..31}, blob keys are {rand_4_chars}/{blob_id}. Random prefix prevents S3 partition hotspots; multi-bucket sharding caps blast radius. Versioning on for "delete for everyone".

07High-level design

Which components handle a request, and in what order?

The architecture in one paragraph. Clients hold a long-lived WebSocket to chatd via Anycast L4 LBs. On send, chatd dedup-checks idem, writes the durable row to ScyllaDB, and publishes a msg.outbox event to Kafka (these two writes form an outbox txn group — local commit + event publish, no 2PC). The Routing & Fanout service consumes the outbox, looks up recipient devices in Postgres, and XADDs an envelope into each recipient's Redis inbox stream. Online recipients receive via server-initiated WS push: each recipient's chatd shard holds a XREADGROUP BLOCK subscription on its assigned inbox shards (e29); when route XADDs (e17), Redis wakes the subscription, chatd reads the envelope and writes a deliver frame back through edge-lb (e30) → recipient device (e31); the recipient's delivered ACK (e32) reverses through the same socket and chatd XACKs the stream entry, which XDEL's it from the inbox. Offline recipients keep their envelopes queued in the Redis inbox stream (TTL ≤30d) and drain them on reconnect. Media is encrypted client-side, uploaded via api-gw → S3 (presigned PUT), and downloaded directly from a CDN (signed URL). The abuse-worker consumes a separate envelope-metadata stream and writes flags back to Postgres — it never sees plaintext.

Why the recipient is drawn as its own node. The sender's chatd shard ≠ the recipient's chatd shard (they're consistent-hashed onto different boxes), and the WebSocket carries traffic in both directions for each device. Modeling the recipient as a separate mobile node with its own reverse-direction edges (e30/e31/e32) makes the end-to-end delivery hop explicit and traceable; without it, "the message arrives at the other user" is hidden behind a single bidirectional edge and easily missed at review.

Why outbox, not 2PC. A 2PC across ScyllaDB + Kafka would block every send on coordinator availability — unacceptable when fanout is 9M/sec and the coordinator is a SPOF. Outbox commits the durable row first; the publish is best-effort fast-path on e9 with at-least-once semantics absorbed by the route consumer's idempotency on (message_id, recipient_id). The transactional outbox is closed by the cdc-relay worker, which tails the ScyllaDB CDC log and re-publishes any msg.outbox event the chatd missed (e.g. crashed between e7 durable commit and e9 publish). Without this relay, "chatd dies after commit, before publish" silently loses fanout — the failure mode that makes interview-style outbox stories fail in production. The route consumer dedupes the chatd-direct + cdc-relay double-publish on message_id.

Why fanout-on-write, mostly. At avg group size 3.5, write-amp is trivial and read-amp on a "fanout-on-read" design would crush mobile cellular plans (every recipient polling on chat open). We flip to fanout-on-read above ~512 members — communities/channels write the message once and members poll on visibility (Discord guild model). The threshold is configurable and observability-driven.

Why two gateways. chatd is stateful (per-session cursors, sequence numbers) and traffic is sticky-session WebSocket — a connection storm during a deploy can't be papered over by Envoy. api-gw is stateless HTTPS — it autoscales on CPU and absorbs auth/group/profile bursts without affecting the chat data path. Conflating them means a control-plane-spike DDOS takes down chat.

Where the abuse pipeline gets its data. chatd publishes envelope-metadata-only events (sender, recipient list, fanout count, timestamp, message_id — never plaintext, since the server can't decrypt) on the abuse.signals Kafka topic via e28. The abuse-worker consumes this stream and does Flink-based reputation scoring; a flagged account writes abuse_disabled=true back to meta-db (e24). The signal stream is separate from msg.outbox (different topic, different consumer group) so a poison message in abuse pipeline can never block fanout. Abuse signals are best-effort by design — cdc-relay covers msg.outbox only (CDC tails the durable store), so a Kafka broker hiccup that drops an abuse.signals event is not recovered. We accept this: spam scoring is statistical, not transactional; missing 0.1% of envelope-meta events does not change a reputation score that aggregates over millions.

The E2EE identity root. The auth service serves libsignal prekey bundles and per-device identity keys and maintains the Key Transparency Merkle log (signed root + per-account proofs on e11). Steady-state chat traffic never touches auth; a degraded auth blocks new prekey fetches and KT proofs, never the send/deliver data path.

Failure modes & mitigations

These came out of real production post-mortems (citations in §reading). Each mitigation is a concrete piece of the canonical, not "monitor it harder."

#FailureDetectionMitigation in this design
1Reconnect storm after edge replica drain (Slack 2021)socket_reconnects/sec > 5× baseline for 60s burn-rateLame-duck reconnect-after: jittered(60–600s) header on chatd; client honors it. ulimit + connection-count token bucket at edge-lb
2Mega-group HoL blocking (Discord 2017 #general)Per-partition Kafka consume_lag_p99 > 30s; tail-latency by conversation_idSub-shard whales: groups >5K members map to a dedicated route consumer pool with reserved Kafka partitions; ScyllaDB partition key bucketed by 7d window
3Cache stampede post-deploy (Discord, GitHub 2018)cache_miss_rate > 10× baseline + downstream DB QPS spikeSingle-flight on chatd inbox reads (per-user lock); pre-warm inbox shards on rolling deploy via a synthetic touch job
4Retry amplification (AWS Dec 2021 us-east-1)fanout_msgs_out / msgs_in > 2σ; rising 5xx + retry counts bothRetry budget ≤10% of inbound at every RPC boundary; idempotency key on every state-changing edge; circuit-breakers around meta-db + msgstore so a slow downstream sheds fast
5TLS / config push misses an edge box (Slack Jan 2021)TLS handshake error-rate >0.1%; config-drift alarmAutomated cert rotation (cert-manager) with 30-day burn alarm; staged config push 1%→10%→100% with 10-min soak; kill-switch on a separate path
6Hot conversation single-shard meltdownScyllaDB partition size histogram p99 > 100MBTime-bucketed partition key already caps width; per-conversation send rate-limit at chatd (1k msgs/min/group)
7Cross-DC replication lag → "sent but missing"repl_lag_seconds > 30 per keyspaceMulti-device device-list reads via quorum across DCs; user-bound to home region with read-your-writes header on cross-region failover
8Tombstone accumulation (ScyllaDB)tombstones_scanned_per_query_p99 > 1000TWCS (Time-Window Compaction Strategy) on msgstore; soft-delete via TTL bucket, not DELETE
9Multi-device key-sync racekt_proof_failure_rate > 0Sender-key sync barrier: reject send until all recipient devices' prekeys are fetched; KT Merkle proof on identity-key change
10Fanout abuse (spam wave)messages_per_sender_p99 > 10σPer-sender fanout-rate limit at route service; reputation-weighted admission; envelope-metadata abuse pipeline
12Region failover botched (LINE/KakaoTalk SK C&C fire 2022)Synthetic probe failure from 3+ external ASNsActive-active routing; user home region in account; failover ramp with warm-cache hit rate gate, not pure traffic shift
13Bad config push (Slack Feb 2022)Error-rate SLO breach + config_revision deltaStaged config rollout with kill-switch on separate path; canary on config templates
15chatd dies between durable commit and outbox publishcdc_relay_lag_seconds > 30; mismatch between msgstore CDC offset and Kafka outbox offsetcdc-relay worker tails ScyllaDB CDC log and re-publishes any missed msg.outbox event. Route consumer dedupes the chatd-direct + cdc-relay double-publish on message_id. Without this, the outbox label on e7+e9 is aspirational, not load-bearing.
16meta-db sync-quorum self-DOS (both followers in an AZ trip; sync writes block)pg_sync_quorum_unmet_seconds > 10 on any shardSync replication degrades to async-with-paging via Patroni synchronous_standby_names = ANY 1 (replica_a, replica_b) so a single follower healthy is sufficient; writes resume on degraded RPO ≤ 1s during recovery. Documented as the deliberate availability/RPO trade.
Observability

SLIs (RED + USE):

  • chatd: connection establish-rate, p99 connect latency, sockets-per-box, drain-duration, send→ack p99, sequence-gap detection.
  • route: kafka-lag-by-partition (p99 < 5s), per-conversation fanout latency, idempotent-replay rate.
  • msgstore: write quorum-form rate, partition size histogram, compaction queue depth, tombstone scan ratio.
  • inbox: keyspace size per shard, eviction rate (must be 0 under TTL-only), AOF fsync latency.
  • meta-db: replication lag (must be ≤1s on hot path), failover-test cadence (weekly).
  • abuse-worker: lag (business-hours SLO, not user-facing page).

SLOs (each split per-region — global is not the alerting target; a single regional outage is invisible against a global budget until 1/18th of the budget burns):

  • 99.99% on socket connect (warm) p99 < 120ms — connect-storm gate so deploy-induced regressions page before users churn.
  • 99.99% on send → server ack p99 < 400ms (52 min/yr error budget).
  • 99.95% on send → recipient delivered p99 < 1.5s (4.4 hr/yr budget).
  • 99.99% on history-page p99 < 600ms.
  • Per-conversation tail-latency SLI: p99 send-to-deliver bucketed by conversation_id into (small, medium, large, whale) cohorts. Anomaly alert when any cohort's 7-day p99 grows ≥2× rolling baseline — catches a hot-conversation creep before failure-mode #2 actually fires.

Burn-rate alerts (Google SRE Workbook Ch. 5, per-region — multi-window-multi-burn-rate):

  • 2% of 30-day per-region budget consumed in 1h → P1 page (sustained ~14× error rate).
  • 5% in 6h → P2 / Slack high-prio.
  • 10% in 3 days → ticket for capacity / dependency review.
  • Cross-cohort gray failure: per-conversation tail-latency cohort p99 ≥2× 7-day baseline for 30m → P2 page even when no SLO is violated. Catches msgstore p99 doubling silently.

Dashboards on-call opens during a page:

  1. Edge & connection: connect-rate, p99 connect latency, sockets-per-box, reconnect rate, TLS error-rate by SNI.
  2. Send → ack pipeline: chatd → idem → msgstore → kafka latency by hop, Kafka producer error-rate.
  3. Fanout: route consumer lag by partition, inbox XADD latency, inbox keyspace size.
  4. Storage: msgstore partition-size histogram, compaction queue, repl-lag, meta-db failover state.
Deploy & runtime
  • Rollout: chatd deploys via blue/green per-AZ — never region-wide simultaneous (the reconnect-storm risk). Each AZ drained over 30 min via lame-duck. kafka and msgstore are config-only changes; binary rolls go region-by-region with a 1h soak.
  • Schema migrations on Postgres: dual-write + shadow-read + backfill + flag-flip. Online ALTER is forbidden on hot tables.
  • ScyllaDB schema: additive only; new columns get a default; column drop = ignore-on-read for one major version, then physical drop in a maintenance window.
  • Kafka consumer rebalances: cooperative-sticky assignor + static membership so deploys don't pause consumption.
  • Secrets: rotated via cert-manager / Vault; clients re-auth on cert change with a 30-day overlap window.
  • Config plane (intentionally not a request-path node): Statsig-style flags with 1%→10%→100% staged push, 10-min soak between stages, and a kill-switch on a separate transport (out-of-band etcd) so a bad config push can't disable its own off-button.
Multi-region / DR
  • Active-active across 18 regions. User's home_region set on account creation; chatd in any region can serve any user, but writes funnel home-region for sync-replicated metadata.
  • chatd: stateless, replicated per region. Anycast routes to nearest healthy POP.
  • msgstore: in-region RF=3 (quorum W=R=2) + 1 cross-region async observer per partition. Region-loss => observer promoted; reads continue, writes block until quorum reforms (typically 30-60s human-driven). RPO ≤ 30s, RTO ≤ 5min.
  • meta-db: in-region 1 leader + 2 sync followers (RPO=0 / RTO ~30s automated). Cross-region async replication for follow-the-sun fallback. Cross-region failover is the only operation that human-confirms (paged).
  • inbox: in-region only — losing a region drops at most pending envelopes for that region's users. Sender re-ACKs after socket reconnect; client dedupes by message_id. RPO ~ 1s effective from user's perspective; RTO < 1min (Redis sentinel).
  • kafka: MirrorMaker-2 to a passive cross-region cluster. The route consumer can resume from passive on regional failure.
Security posture
  • E2EE everywhere on content: libsignal X3DH + Double Ratchet on 1:1, Sender Keys on groups. Prekey bundles served via auth + Postgres. Server never holds plaintext after the client encrypts.
  • Key Transparency: identity-key changes append to a Merkle log; clients verify a proof on every contact key seen. Detects MITM via a compromised server.
  • mTLS between every internal service (auth: mtls on all internal edges; auth: bearer only on client→edge). Cert-manager handles rotation.
  • WAF + per-IP rate limit at edge-lb on control APIs; per-token rate limit at api-gw.
  • Abuse pipeline never decrypts: scoring is on envelope metadata (sender, recipient list, frequency, link count from URL preview which is client-side metadata). E2EE ⇒ search / scan must be client-side; the server side cannot.
  • Compliance: GDPR erasure pipeline triggered by user request → cascades through meta-db (account row), inbox (drop pending), msgstore (tombstone insert with TTL=30d), media-store (object delete with version). Legal-hold flag overrides erasure for the limited rows under hold.

08Deep dives

Which part breaks first, and what do you do about it?

1. Outbox pattern, end-to-end. Sender's chatd issues two writes to two systems (ScyllaDB row + Kafka event). They cannot be 2PC'd at scale (coordinator availability would gate every send, and Kafka has no 2PC). Instead: ScyllaDB write is the durable commit; Kafka publish is at-least-once (retry until ack, accept duplicates). The route consumer is idempotent on (message_id, recipient_id) — it XADDs to the recipient's inbox stream only if message_id isn't already in the inbox's recent-window dedup set. Worst case: route replays an event after a Kafka rebalance; recipients see no duplicate because the XADD is no-op'd. The txn:outbox group is declared on edges e7 (msgstore) and e9 (kafka) so the simulator narrates the contract.

2. Hot mega-group: why bucketing the partition key is the single most important decision. Discord 2017 lost a server because #general on a giant guild had a Cassandra partition that grew to 100GB+. Every read incurred a tombstone scan that timed out the coordinator. The fix: partition key = (conversation_id, day_bucket) where day_bucket = floor(server_ts / 7d). A 100K-member group sending 10/sec produces 6M messages/day = 42M/week — bounded. New partition every Monday. TWCS compacts old buckets and drops them as a unit at TTL. Without this, no other mitigation matters.

3. Multi-device key sync. On send, the sender's chatd asks auth for the recipient's device list (signed by their primary). For each device, the sender encrypts a Sender Key and packs it into the envelope. The route service stores per-device envelopes; recipients fetch only their own. New device linked? Primary signs an updated device list; KT log appends; future sends fanout to the new device. The race: a sender encrypts to a device list of {A, B} just as B unlinks. Recipient B's server-side opaque envelope is delivered but B's device is gone. Resolution: sender retries with new device list on receiving a "device-not-found" delivery error.

4. Idempotent retries through region failover. Client sent a message; got no ack because the edge-lb's anycast IP shifted regions mid-send. Client retries with the same message_id. The new region's chatd checks idem first — Redis sees the SETNX failed → returns the original ack. No double-write to ScyllaDB, no double-fanout. The 7-day TTL on idem is generous: covers a phone left in airplane mode through a region outage.

5. Capacity arithmetic for fanout vs read. Threshold derivation: write-amp W = group_size, read-amp R = unread_pages × per-page_cost / write_cost. At avg unread depth 5 pages × cost-ratio ~1, the crossover is empirically ~512–1000 (Iris threshold; Discord channel flip). Communities >5K must be read-fanout — else one celebrity broadcast at 1M members would mean 1M XADDs in a single fanout job, ~5 minutes of latency. We hardcode fanout_strategy = read above 512 members and document it as configurable.

09Trade-offs

What did this design cost, and what breaks at 10×?

The big architecture trade-offs we accepted:

  • Outbox over 2PC. We accept eventual consistency between durable record and fanout. Worst case: a 1–5s fanout-lag visible as "checkmark drifts in slowly." We rejected 2PC because the coordinator is a SPOF and every send pays its availability tax. Saga is overkill — there's nothing to compensate.
  • Sync metadata replication, async message-store replication. Postgres meta-db is sync (RPO=0, ~30ms write tail) because membership invariants must not flicker. ScyllaDB msgstore is quorum (sequential, RPO=0 in practice when 2/3 replicas survive) because it's the volume tier and locking on cross-AZ would crater throughput.
  • Two gateways, not one. The control-plane traffic shape (HTTPS bursts, stateless) is fundamentally different from the data-plane (long sockets, stateful). Combining them couples deploy risk; we'd take chat down on a control-plane-handler bug.
  • No server-side message search. E2EE forces client-side index. We accept that "search messages" is a per-device experience; can't be a server feature without breaking the security claim.
  • Single Kafka cluster with logical topics. We could split outbox / abuse into separate clusters for blast-radius isolation. We chose one cluster (cheaper, simpler ops) because the consumers already isolate per consumer group; a poison message on outbox doesn't pause abuse.
  • Peripheral surfaces scoped out to keep the messaging spine teachable. Offline push, presence / last-seen, read-receipts / typing, and SMS-OTP account registration are deliberately cut from this canonical. Each layers back on the same spine without reshaping it: an offline-push consumer subscribes to a third Kafka topic (push.requests) driven off the same fanout and drives APNs/FCM; presence and read-receipts ride the existing WebSocket + inbox as ephemeral frames; SMS-OTP registration hangs off auth via an external provider. The core we keep — durable write → outbox fanout → WS deliver, media, and the E2EE identity root — is the load-bearing part; the rest is a consumer or a frame away.

What we rejected and why:

  • DynamoDB for messages: rigid query patterns (no WHERE seq > X), $$ at 9M deliveries/sec, harder hot-partition mitigation than ScyllaDB. (Still a valid pick at 1/10th this scale.)
  • MyRocks (Iris): defensible if you already run MySQL at Meta scale. Greenfield, ScyllaDB wins on operational simplicity (shard-per-core, no JVM compaction).
  • Spanner / CockroachDB for messages: serializable cross-shard txns are not what messaging needs. Pay 5× write latency for nothing.
  • gRPC bidi-streaming instead of WebSocket: gRPC is fine for service-to-service but mobile NAT traversal + corporate proxies still favor WebSocket. We use gRPC internally; WS at the edge.
  • 2PC msgstore + kafka: see above. A no.
  • Single-leader Erlang Mnesia for state (WhatsApp original): great at 2012 scale, doesn't scale to 1.5B DAU. Replaced by sharded Postgres + Redis + ScyllaDB for the modern shape.
Trade-offs we accepted (open questions)
  • Disappearing-message TTL precision: TWCS gives us approximate TTL (within bucket boundaries). Users who set "1 day" may see the message vanish at 1d±3.5d depending on partition-bucket alignment. Acceptable because the threat model is "trace stays out of long-term backups," not "vanishes within ±1ms."
  • Calls / RTC tier: out of scope here. The signaling path reuses chatd's WS, but the media plane is a separate SFU + TURN tier with its own DR posture. A complete canonical needs to add it; we explicitly defer.
  • Server-side analytics on E2EE traffic: deliberately blind. Product analytics ("how many groups have >100 members") rely entirely on metadata queries against meta-db, which is not E2EE. We accept the limitation.
  • Hot-shard rebalance on msgstore: ScyllaDB's vnode rebalance under sustained ingest is a known operational sharp edge. We commit to a maintenance-window rebalance + 2× headroom rather than online rebalance under live load.
  • Key Transparency adoption: KT proofs are verified by clients on identity-key change, not on every contact-key fetch. A passive MITM by a compromised server is detectable by users who manually check; not auto-detected. Active-MITM detection requires a gossip layer we haven't built.

Primary sources

  • Rick Reed — That's 'Billion' with a B (Erlang Factory 2014)
  • Migrating Messenger storage to optimize performance — engineering.fb.com (Iris/MyRocks)
  • How Discord stores trillions of messages — discord.com/blog (ScyllaDB migration)
  • Slack's Outage on January 4th 2021 — slack.engineering
  • Deploying Key Transparency at WhatsApp — engineering.fb.com (2023)
  • Signal Protocol whitepaper
  • Cloudflare — October 2021 Facebook outage post-mortem
  • Google SRE Workbook, Ch. 5 (alerting on SLOs)

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work WhatsApp / Messenger yourself