Slack / Discord
Worked solution

Slack / Discord — a worked solution

Channels and history. Push or pull — and how a hot-channel fanout doesn't melt the gateway.

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 Slack / Discord workspace

The problem

Build a production reference architecture for channel-based real-time messaging at hyperscale (think Slack workspaces or Discord guilds). The shape: tens of millions of long-lived TLS sockets, sub-second channel-broadcast delivery, server-side history that survives restarts, and a fanout path that does not melt when 50M users come online at once. Every component below is specified well enough that an SRE team would defend the choices in a 5-year incident retro — Slack and Discord have both already had those retros, public, and the canonical bakes their lessons in.

The reference architecture

Reference architecture for Slack / Discord: 15 components — Client (Web + Mobile + Desktop), Recipient Client (peer device), Anycast L4 LB + WAF, Edge Cache (Flannel-style), WebSocket Gateway (ws-gw), Control API Gateway, Auth / Identity Service, Channel Server (channel-svc), Mega-Channel Relay Service, Presence Cache, Idempotency / Dedup KV Store, History Read Service, Message Store, Metadata DB (workspaces, channels, members), Coordinator (etcd / Consul) — connected by 27 flows.Client (Web + Mobile …React (web), Electron (…Recipient Client (pee…Same client codebaseAnycast L4 LB + WAFAWS NLB (Slack) / Cloud…Edge Cache (Flannel-s…Slack Flannel (custom E…WebSocket Gateway (ws…Slack Gatewayserver (Go…Control API GatewayEnvoy + WAF + per-token…Auth / Identity Servi…Custom Go service + JWT…Channel Server (chann…Slack Channelserver (Go…Mega-Channel Relay Se…Discord-style relay ser…Presence CacheRedis 7 cluster (ETS in…Idempotency / Dedup K…Redis 7 Cluster (SETNX …History Read ServiceDiscord Rust 'data serv…Message StoreScyllaDB 6.x (Discord, …Metadata DB (workspac…Slack: Vitess + MySQL (…Coordinator (etcd / C…etcd 3.x (Discord) / Co…
15 components, 27 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Client (Web + Mobile + Desktop)React (web), Electron (desktop), Swift / Kotlin (mobile)

Holds the user's long-lived WebSocket to ws-gw (one socket per device, multiplexes every workspace and channel the user belongs to). Maintains a per-channel sequence cursor for resume, tracks the visible-channel viewport for selective subscription, and runs a small in-memory outbox that retries unacked sends with the same client-generated message_id.

Why it exists. The product claim 'channel I'm watching feels live' depends on a single stateful WS — opening a new HTTP connection per channel would cost a TLS handshake per chat switch and burn mobile battery. The viewport tracking is what makes Discord-style selective subscription possible; without it the client would pull live fanout for every channel of every workspace it belongs to and the gateway would melt.

When it fails. Stale resume token after a long sleep → client falls back to IDENTIFY which costs more on the server. Client crash mid-send → retried message_id collapses on idem-check at channel-svc. Aggressive carrier NAT-rebind kills the socket → reconnect with same session_id resumes from last_seq.

Recipient Client (peer device)Same client codebase

A different end-user's device subscribed to the same channel. Holds its own long-lived WebSocket (likely on a different ws-gw box, possibly a different region). Receives MESSAGE_CREATE frames pushed from its assigned ws-gw and surfaces them in the UI. Delivery is one-way server-push — the recipient renders the frame; there is no delivered/read ACK on the hot path.

Why it exists. Drawn explicitly so the end-to-end fanout hop ('the message arrives on the other person's screen') is traceable in the simulator. Earlier drafts hid recipient delivery behind a WS edge from ws-gw to itself — a reader couldn't tell whether fanout was modelled or hand-waved.

When it fails. Recipient socket dies mid-deliver → on reconnect, ws-gw replays from the per-(user, channel) cursor stored in the subscriber-set; client dedupes by message_id. Longer offline → the message surfaces on the next history backfill when the client reopens the channel.

Anycast L4 LB + WAFAWS NLB (Slack) / Cloudflare-fronted GLB (Discord) + Envoy WSS

Terminates BGP-anycast IPs and forwards WebSocket Upgrade requests at L4 to a fleet of Envoy WSS nodes that handle TLS termination, per-IP token-bucket rate-limit, and WAF rules. Pure L7 passthrough on /v1/gateway and /v1/api/*; the long-lived socket pinning is done by Envoy's hash policy on team_id (Slack) or by Discord's anycast → POP routing that lands at a region-local Envoy.

Why it exists. Pinning a WS to one Envoy at handshake decouples connection density from the LB layer (NLB is L4-only, so it has no per-connection state to lose); WAF + per-IP rate limit at edge prevents the connection-storm attack surface (a botnet hammering /v1/gateway can't flood the gateway). Slack's 2021 outage explicitly cited TGW (transit-gateway) saturation upstream of this layer — keeping the L4 path simple is what lets cellular DR work.

When it fails. TGW or upstream peering saturates (Slack 2021-01-04: TGW saturation drove cascade) → handshake p99 cliff; reconnects pile up. Mitigation: pre-warm capacity at known traffic-floor times (Monday-after-holiday is the documented worst case); decouple monitoring VPC from production VPC so on-call can still see the dashboards when production traffic is hosed.

Edge Cache (Flannel-style)Slack Flannel (custom Erlang/Go, Redis-L2) — Discord embeds equivalent in gateway

Application-level edge cache deployed at every POP. On WS connect, ws-gw consults flannel for the user's workspace metadata (channel list, member roster) and pre-pushes the bootstrap blob over the same socket so the client renders the workspace shell without round-tripping origin. Holds ~600K qps and ~4M concurrent at Slack's published peak; absorbs 80–95% of cold-start metadata reads.

Why it exists. A cold reconnect to a workspace with 50K members and 5K channels would issue tens of separate origin reads against meta-db — the IDENTIFY storm in Slack's 2021 outage was exactly this. Flannel collapses workspace bootstrap into one POP-local read with a pre-pushed payload; without it, every IDENTIFY hits the metadata tier and a reconnect storm immediately becomes a meta-db storm.

When it fails. Cold flannel after deploy or post-failover → IDENTIFY storm cascades to meta-db (Slack 2022-02-22 was the inverse: mcrouter cache churn during Consul rollout caused a Vitess feedback loop). Mitigation: synthetic warm-up on rolling deploy, single-flight on cache fills (per (team_id, key) lock), and the per-team rate limit at api-gw so one tenant can't blow the cache.

WebSocket Gateway (ws-gw)Slack Gatewayserver (Go) / Discord Cowboy (Elixir/BEAM)

Holds 200–500K stable WSs per box (Discord Elixir reaches 1M/box on the very heaviest machines; we size for the conservative 200K Go figure). Each session is a per-user GenServer (Discord) or goroutine pool (Slack) that owns: identity, the channel subscription set, the sequence cursor for RESUME, the inbound mailbox, and the subscription-anchor for selective subscription. On send: dispatches to the channel-svc owning the target channel. On RECEIVE from channel-svc / relay: writes a MESSAGE_CREATE frame back through edge-lb to the WS.

Why it exists. Long-lived stateful WS 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 registration spike or bad auth deploy would take chat down too. Slack and Discord both maintain this split for exactly this reason.

When it fails. Reconnect storm during deploy if RESUME protocol regressed (every client falls back to IDENTIFY → 5–10× cost); Slack's Disasterpiece Theater rehearses this. Mitigation: dual-version protocol support across deploys, jittered drain, the connection-count token bucket at edge-lb. BEAM scheduler / GC pause stalls a box → per-box scheduler_utilization > 95% alert; sessions migrate on next reconnect.

Control API GatewayEnvoy + WAF + per-token bucket

Stateless HTTPS gateway for the control plane: login (rtm.connect), workspace + channel CRUD, member management, slash-commands. Authenticates the bearer/session token, applies per-token + per-IP rate limits, and forwards to auth or meta-db over mTLS with circuit-breakers in front of both. Search and file-uploads are explicitly out of scope for this canonical (separate sub-canonicals).

Why it exists. Control-plane traffic is bursty, stateless, and CPU-bound — a different operational shape from the long-lived stateful ws-gw. Separating the two means a registration storm, a bad auth-deploy, or a slash-command misconfig can't take chat down. Slack's split between Webapp/Adminserver and Gatewayserver/Channelserver is the same principle.

When it fails. Bad config push (Slack 2022-02-22 — Consul-driven config regression hit prod, mcrouter cache churn cascaded to Vitess) → 50% errors. Mitigation: staged rollout 1%→10%→100% with a 10-min soak; kill-switch on out-of-band etcd (separate path, can't be disabled by its own bug). Credential-rotation storm during a security incident → per-token bucket at this layer is the only thing keeping auth alive.

Auth / Identity ServiceCustom Go service + JWT + workspace-scoped sessions

Mints short-lived WS session URLs (Slack rtm.connect returns a one-time URL that contains the bearer; Discord IDENTIFY validates a long-lived user token), evaluates per-workspace ACL on every privileged op, and writes the sessions table in meta-db.

Why it exists. Auth must be a chokepoint so per-user rate limits actually rate-limit (you can't enforce a token quota if every service self-mints sessions); a single tier means token revocation propagates in seconds, not minutes. Auth-tier latency is on the WS connect critical path, so it's separated from api-gw and replicated independently.

When it fails. Auth tier degraded → ws-gw cannot mint new sessions, but existing sockets keep working (the JWT verification is local on ws-gw). Mitigation: long JWT TTL (24h) with refresh on activity — even 30 minutes of full auth outage doesn't drop existing sockets; api-gw circuit-breaker on auth so login attempts fail fast instead of stacking.

Channel Server (channel-svc)Slack Channelserver (Go) / Discord guild GenServer (Elixir)

Owns the per-channel actor: holds the recent-history ring buffer (last ~200 messages in memory), the subscriber set (which ws-gw boxes have sessions in this channel), and the ordered sequence allocator. On send from ws-gw: dedup-checks idem-cache via SETNX on message_id (returns the original ack on duplicate), validates ACL via meta-db, allocates the next channel-local seq, writes the canonical row to msgstore, multiplexes the MESSAGE_CREATE event to every ws-gw with subscribers (small/medium channels). On large channels (>100K subs) routes the fanout through relay instead of dispatching directly.

Why it exists. Per-channel ordering must be enforced at one place (FIFO is a load-bearing UX claim — 'reply-to-message' would break under reordered delivery). The fanout work is too heavy for ws-gw (sender's box would do O(N-subscribers) work on every send, blocking that user's other sends) so we centralize it on the channel actor and shard *channels*, not users. Slack chose this with Channelserver in 2014; Discord arrived at the same place with the guild GenServer.

When it fails. Hot channel (Maxjourney; #general on a 1M-member workspace) saturates the single owning actor → BEAM mailbox grows unbounded, fanout p99 spirals (Discord 2017 #general meltdown, Maxjourney Nov 2023). Mitigation: above 100K subscribers route through relay — passive sessions filtered there, fanout parallelized; per-channel send rate-limit (1k msgs/min) as a backstop. Owning-CS crash → coordinator detects via lease expiry, reassigns the channel; clients see a 1–3s delivery gap during reassignment which is within budget.

Mega-Channel Relay ServiceDiscord-style relay service (Elixir + Rust NIF for SortedSet)

Inserted between channel-svc and ws-gw for any channel with >100K subscribers. Each relay owns a slice of subscribers (~15K sessions per relay, the published Discord Maxjourney number) and applies the passive-session filter: subscribers whose client viewport doesn't currently include this channel get dropped at the relay rather than at the gateway. On a hot-channel post: channel-svc sends one message to each owning relay; relays parallel-fan-out to their slice of ws-gw boxes; ws-gw multiplexes to the actual sockets.

Why it exists. A naive single-actor fanout to 1M subscribers is O(N) work blocking the channel actor — Discord measured this in Maxjourney as the cliff at ~100K subs (relay tier was their fix). Without relays, the 90%+ passive-session ratio (subscribers who 'follow' a channel but don't have it focused) all pay full WS-deliver cost. Relays drop them upstream of ws-gw, cutting fanout cost ~10×.

When it fails. All relays for a hot channel die simultaneously (deploy bug; coordinator pause) → channel-svc detects relay-loss heartbeat and falls back to direct dispatch (degraded — fanout latency 2× normal, passive-session filter inactive). The fallback is the load-bearing safety net; without it, a relay deploy that hits all instances at once would be a P0. Backpressure: relay subscriber index growing unbounded → relay sheds new subscribers via 'try again later' frames at 80% capacity.

Presence CacheRedis 7 cluster (ETS in-process for Discord guild subscriber-set)

Holds the per-channel subscriber set on a Redis cluster (and ETS in Discord's per-guild process heap for hot data): subs:{channel_id} → set of {user_id, ws_gw_id, viewport_focus} — the subscriber-set that ws-gw and channel-svc both consult to compute fanout targets. Updated by ws-gw on connect / subscribe / unsubscribe / viewport-change; read by channel-svc and relay when computing which sockets a message must reach.

Why it exists. The fanout decision has to know, cheaply, which subscribers currently have a channel in their visible viewport — that viewport_focus bit is what lets relay/channel-svc skip the 90%+ of subscribers who follow a channel but aren't looking at it (Discord's published passive ratio). Without a fast in-memory cache holding the (channel_id → {user_id, ws_gw_id, viewport_focus}) set, the passive-session filter can't be cheap to read on every fanout decision.

When it fails. Subscriber-set shard primary fails over (RTO ~15s) → during gap, fanout decisions can't read viewport bits, so relay falls back to 'send to all subscribers' for that channel partition (degraded fanout cost, not data loss). Mitigation: passive-fallback heuristic at relay (deliver-then-let-client-drop); alert on subscriber_shard_failover_count correlated with relay_passive_filter_skipped.

Idempotency / Dedup KV StoreRedis 7 Cluster (SETNX + AOF every-second, 5-min TTL)

Hot-path message_id → {seq, server_ts, ack_blob} written by channel-svc on every accepted send before the durable msgstore write. SETNX returns the original ack on a duplicate message_id (client retried because no ack arrived in time, e.g. region failover mid-send) so the durable INSERT and fanout dispatch happen exactly once. 5-minute TTL window — clients retry within seconds; longer windows would over-retain hot keys.

Why it exists. Without this single hot-path SETNX, edge e7 (ws-gw → channel-svc) using retryPolicy: fixed-retries is unsafe — a transient blip would double-write to msgstore, double-fanout to recipients, and surface a duplicate message UI. The product claim 'message arrives once' depends on this dedup store. An in-process LRU per channel-svc box would lose dedup state on the natural CS-restart that follows a region-failover-rebalance — the failure mode the dedup must specifically handle.

When it fails. Primary loses last AOF segment + fails over → small dedup gap during which a retry can double-apply. Mitigation: clients also dedup by message_id at the recipient timeline (defense-in-depth); alert on idem_failover_during_aof_write_count. Loss of all 16 shards (regional outage) → channel-svc falls open with retryPolicy: none on e7 (cooperative graceful degradation — duplicate-message risk traded for availability) and an explicit page.

History Read ServiceDiscord Rust 'data services' tier (Tonic/gRPC + tokio) / Slack history-api

gRPC tier in front of msgstore that serves channel history paginated by (channel_id, before_seq, limit). The load-bearing trick is request coalescing: concurrent reads for the same (channel_id, day_bucket, page) collapse to a single ScyllaDB read whose result is multicast back to all waiters. A 200K-user thundering herd on a popular channel becomes one DB read.

Why it exists. Without coalescing, a channel-open after deploy (everybody's clients reconnect → everybody re-fetches history) would issue tens of thousands of identical ScyllaDB reads in a 5-second window — Discord's pre-Rust-tier failure mode. The fix is famous: their blog explicitly cites the data-services tier doing dedup of in-flight reads as the difference between 'works at scale' and 'falls over on every deploy.'

When it fails. All replicas serving a hot channel die together → coalescing collapses, downstream ScyllaDB sees the un-deduped storm. Mitigation: process-level cache TTL = 30s (a brief restart still warms within seconds); the coalescing window itself is single-process so a partial outage degrades the channel partition, not all reads.

Message StoreScyllaDB 6.x (Discord, post-2022 migration) — alternative: Vitess+MySQL (Slack)

Wide-column store of every accepted message, partitioned by ((channel_id, day_bucket), message_id) and clustered by monotonic per-channel seq. RF=3 with quorum W=R=2 in-region + 1 async observer per partition cross-region. Discord's published peak post-Scylla-migration: 72-node cluster handles billions of messages/day at p99 read 15ms / write 5ms. Holds the canonical record for both DM and channel messages; CDC log feeds search-indexer + analytics.

Why it exists. Hyperscale chat history is fundamentally an LSM-tree problem — sequential write + tail-read access patterns on 1KB rows at petabyte scale. A single-leader RDBMS crushes under write amplification at >10K msgs/sec/node; serializable cross-shard transactions are not what messaging needs. Discord migrated Cassandra → ScyllaDB explicitly because JVM GC pauses + LSM compaction starvation made p99 unpredictable; Scylla's shard-per-core C++ runtime gives deterministic tail.

When it fails. Hot mega-channel on an un-bucketed partition (Discord 2017 → tombstone wall → coordinator stall on every read). Mitigation: time-bucketed partition key already in place; per-channel send rate-limit at channel-svc; alert on tombstones_scanned_per_query_p99 > 1000. Shard-loss (a single ScyllaDB node fails) → quorum W=R=2 still met by the surviving 2 replicas (degraded p99); two-node loss → quorum unmet on that token range, blast radius ≈ 1/96.

Metadata DB (workspaces, channels, members)Slack: Vitess + MySQL (sharded by team_id) — Discord: Postgres/Citus on user_id

Sharded relational store for accounts, workspaces (Slack teams / Discord guilds), channels, channel_members (with role + last_seen_seq), abuse_disabled flag. Slack's Vitess setup runs >50B queries/day at peak 2.3M QPS with p99 11ms; sharded by team_id so a workspace's data co-locates. Serves channel-svc's ACL check on every send, history-svc's permission check, and api-gw's CRUD operations.

Why it exists. Membership has strict invariants ('alice is a member of #eng' must be either true or false at any instant — never flickering during a write) and is small enough (~30TB at Slack scale) that a sharded relational store with serializable transactions is the right shape. The ACL check is on the send hot path; this is why it ships as Vitess+MySQL with sync replication, not eventual NoSQL — the cost of a stale ACL is a privacy violation, not just stale UX.

When it fails. Both followers in an AZ trip → sync-quorum unmet → writes block on that shard. Mitigation: ANY-1-of-2 fallback documented as the deliberate availability/RPO trade-off (RPO ≤1s during recovery). Cross-region failover is paged because the cross-region replica is async (RPO ≤30s). Hot shard during a viral workspace signup → Vitess split-shard primitive (the 2020 capability that saved Slack during pandemic onboarding).

Coordinator (etcd / Consul)etcd 3.x (Discord) / Consul (Slack) — Raft-based

Holds the channel→channel-svc owner consistent-hash ring as a watched keyspace, the (user_id → ws-gw box) affinity table for reconnect routing, ws-gw / channel-svc / relay membership, and per-service config (rate limits, kill-switches). Every routing lookup at the gateway is local against an in-memory mirror of this state; coordinator only sees writes (membership change, owner reassignment, config push) and watch fan-out.

Why it exists. Slack's 2020-05-12 outage (HAProxy/Consul synchronizer bug) and 2022-02-22 outage (Consul-driven config rollout) both came from this layer — getting the membership-source-of-truth wrong cascades to every routing decision in the system. A dedicated, Raft-backed coordinator with explicit watches is the only defensible answer; ad-hoc service registries have all failed in production.

When it fails. Coordinator quorum-loss (e.g. 3 of 5 nodes in one AZ during AZ outage) → membership writes block, but reads continue from any survivor. Routing decisions on existing topology keep working; only new CS box assignments / channel reassignments pause. Mitigation: 5-node quorum across 3 AZs (tolerates a single-AZ loss); the steady-state read path doesn't depend on the coordinator being up — only changes do.

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? No. Slack and Discord both store plaintext messages on the server (server-side search depends on it; abuse pipelines depend on it; admin compliance depends on it). This is the load-bearing differentiator from WhatsApp / Signal. A separate canonical handles E2EE chat.
  • Channels vs DMs? Channels are the dominant traffic shape — typically 10–10K subscribers each, with a long fat tail (Slack #general on a 50K-employee tenant, Discord #general on a 1M-member guild). DMs are 1:1 channels with subscriber-count=2 and identical infrastructure.
  • Selective subscription? Yes — clients declare which channels are in their visible viewport, and the fanout path only delivers to subscribers whose viewport includes the channel. Naive eager delivery to every follower is the most-cited scaling failure in chat history (Discord's Elixir 5M post).
  • History on chat-open? Yes — server-side, paginated. Up to several years' worth on a free tier. This rules out 'history lives on device' (WhatsApp's model) and forces the request-coalescing tier in front of msgstore.
  • Push vs pull? Push for online recipients (server-initiated WS frame from channel-svc → ws-gw → edge-lb → device); pull-only for history (client requests page on chat-open). No long-poll fallback in scope.
  • Mega-channels? Yes — channels can grow to 1M+ subscribers (Maxjourney moment). The relay tier is load-bearing for these, with the cliff at ~100K subs.
  • Voice/video? Out of scope. Signaling shares the WS, but media plane is a separate SFU + TURN tier with its own DR posture. Reference but don't unify.
  • Search? Yes, but as an async-populated Elasticsearch tier off msgstore CDC — explicitly modelled outside this diagram (would deserve its own subgraph).
  • Multi-region? Active-active across 18 regions. Workspace home-region selected on workspace creation; a single workspace's authoritative writes serialize through its home region.

Assumptions baked in (defaults — see capacity for the math):

  • 150M DAU, 400M MAU.
  • 25 messages / DAU / day → 3.75B msgs/day.
  • Avg fanout per send: 60 subscribers × 35% online = 21 deliveries/send average, with a fat tail.
  • Peak / avg multiplier: 3× (lunch + evening + Mondays).
  • Concurrent WS sockets: 80M at peak (55% of DAU online).
  • Hot retention: 30 days for the data services LRU; cold storage is indefinite (ScyllaDB's compaction handles long tail).

02Functional reqs

What must this system actually do?

  • Send a message to a channel (or DM treated as 2-subscriber channel) of arbitrary size; ≤4KB body, optional file/media reference.
  • Deliver in real time to every online subscriber via server-initiated WS push from channel-svc → ws-gw → edge-lb → device. Sub-200ms p99 to recipient.
  • Persist server-side (the differentiator): every accepted message lands in msgstore before fanout, indexed by (channel_id, day_bucket, message_id).
  • Fetch history on chat-open (paginated before_seq, 50 messages/page) via the request-coalescing data-services tier.
  • Channel CRUD: create, archive, member add/remove, role change, invite expansion. All via api-gw → meta-db.
  • Replies / threads: replies treated as channel-scoped events (same fanout); thread is a sub-channel with its own subscriber set.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Availability: 99.99% on the WS data path (52 min/yr error budget, per-region not global); 99.95% on control APIs.
  • Latency (online recipient, same-region):
  • WS connect (warm, TLS resumption): p50 80ms / p99 350ms / p99.9 1.5s.
  • Send → server ack: p50 60ms / p99 200ms / p99.9 500ms.
  • Send → recipient delivered (perceived): p50 120ms / p99 350ms / p99.9 1s. (Slack's published end-to-end target is ~500ms global.)
  • Cross-region recipient: add ~80ms RTT.
  • History page (50 msgs): p50 60ms / p99 250ms / p99.9 700ms (Discord's published post-Scylla read p99 is ~15ms at the storage tier; the rest is request-coalescer + network).
  • 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). Cross-region RPO ≤30s on msgstore. RPO ≤1s on the subscriber-set cache (semi-sync — losing 1s of subscriber-set updates is invisible).
  • Consistency:
  • Per-channel FIFO, enforced by the channel actor's monotonic seq allocator.
  • Read-your-writes for the sender on ACL/membership reads (route via meta-db leader for the sending user's home shard).
  • Eventual cross-region for everything else; cap at 30s replication lag with a paging burn-rate alert.
  • Scalability: every tier horizontal except the per-channel actor (single-writer per channel by design — per-channel ordering invariant) and the per-user WS session (single-writer per user — each user has exactly one WS per device). Both scale by sharding the keyspace, not the writer.
  • Security: mTLS between every internal service; per-token + per-IP rate limits at api-gw; WAF at edge-lb on control APIs.
  • Compliance: GDPR erasure ≤30d cascading across meta-db (account row delete), msgstore (insert tombstone with TTL), search index (rebuild). Eventually-atomic, audited monthly.

04Capacity estimation

How much load and data does this have to hold?

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

Send QPS at peak

peak_send_qps = DAU × msgs_per_DAU_per_day × peak_factor / 86400
              = 150M × 25 × 3 / 86400
              ≈ 130K sends/sec

(Slack's Vitess peak of 2.3M QPS is the combined read+write tier; sends are a small fraction.)

Fanout deliveries at peak

peak_deliveries = peak_send_qps × avg_subs_per_msg × online_fraction
                = 130K × 60 × 0.35
                ≈ 2.7M deliveries/sec

The fat tail dominates: a 1M-subscriber Maxjourney channel posting 1 msg/sec at 30% online = 300K deliveries/sec from one writer. The relay tier exists to absorb this single-channel volume without back-pressuring the rest of the system.

Concurrent WS sockets

sockets = DAU × concurrent_online_fraction = 150M × 0.55 = 80M
sockets_per_box (Linux + Envoy WSS terminator at edge-lb, ws-gw is consumer):
  CPU + WS-state ceiling ≈ 250K (Go gw — Slack figure)
  Discord's Elixir/BEAM: 1M+ (blog "Real-time communication at scale", 2020)

→ 1200 ws-gw boxes globally at the conservative Go figure (250K each), spread across 18 regions (≈70/region) with 4× headroom for region-failover and the Slack-2021-Monday-after-holiday case. The replicas number reflects this.

Channel-svc shard count

Each box owns ~8M channels (Slack quotes 16M/host; we size for 8M with headroom). At 50M channels active in the month, 600 channel-svc boxes are sufficient. Re-shard via etcd ring rebalance is a deploy-time op, not steady-state.

Relay tier (mega-channels only)

Discord Maxjourney published number: 15K sessions per relay. For a Maxjourney-class channel (1M concurrent subs), a single channel needs 1M/15K ≈ 70 relays simultaneously; we provision 120 relays across the fleet to handle 1–2 mega-channels concurrently. Most channels never need a relay (cliff is at ~100K subs).

Message store sharding

ScyllaDB partition ((channel_id, day_bucket), message_id) where bucket = 10-day window (Discord's number). At 3.75B msgs/day × 280B (RF=3 incl. indexes, post-Zstd compression) ≈ 0.4 PB hot-tier text per year (raw, before lifecycle to warm tier; cumulative storage at year-5 is ~2 PB once warm-tier rows pile up). With 96 ScyllaDB nodes (Discord ran 72 post-migration; we add headroom for 1.5× their MAU):

physical_write_qps = peak_send_qps × RF = 130K × 3 = 390K phys writes/s
sustainable per i4i.4xlarge node ≈ 30K writes/s
nodes = 390K / 30K = 13 minimum for ingest

Plus reads + compaction headroom (×7) ≈ 96 nodes. Discord's published post-migration p99: 15ms read / 5ms write.

History request coalescing

peak_history_qps = DAU × channel_opens_per_DAU_per_day × peak_factor / 86400
                ≈ 150M × 30 × 3 / 86400 ≈ 156K history reads/sec
coalescing_ratio = empirical 8–20× during normal traffic, 100×+ during deploy stampede
post-coalesce_qps_to_msgstore ≈ 8K–20K reads/s in steady state

The coalescing factor is what makes the 96-node ScyllaDB cluster sufficient. Without the data-services tier, the same workload would crush 200+ nodes during deploy.

Presence cache (subscriber-set)

subs_writes_per_sec = sockets × viewport_change_freq ≈ 80M / 30s = 2.7M writes/s
working_set = subscriber_rows × 200B/entry ≈ 64 GB across 64 Redis primaries (~1 GB/primary)

The subs:{channel_id} set carries the viewport_focus bit read on every fanout decision — the load-bearing data for the passive-session filter.

Idem-cache (dedup KV)

idem_writes_per_sec = peak_send_qps = 130K writes/s
key = message_id (UUIDv7), value = {seq, server_ts, ack_blob} ≈ 200B
working set = 130K × 300s × 200B ≈ 8 GB across 16 Redis primaries (~500 MB/primary)

The 5-minute TTL is short enough that working set is bounded; a 7-day window like WhatsApp would push us into 2-3 TB and cost more than it earns for a non-E2EE chat where defense-in-depth dedup also lives at the recipient client.

05API design

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

Chat protocol is a custom binary frame over WebSocket with a JSON-equivalent shown for clarity. Slack uses RTM; Discord uses gateway opcodes. We sketch as JSON.

client ⇄ ws-gw  (wss://gateway.arc.ly/v1/gateway)

# IDENTIFY (cold connect)
C→S { "op": "IDENTIFY", "token": "...", "intents": 1234, "shard": [0, 16] }
S→C { "op": "READY", "session_id": "...", "user": {...}, "guilds": [...] }

# RESUME (warm reconnect)
C→S { "op": "RESUME", "token": "...", "session_id": "...", "seq": 1842 }
S→C { "op": "RESUMED" }   ← 0 backfill cost vs IDENTIFY's full state push

06Data model

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

meta-db (Vitess+MySQL, sharded by team_id / guild_id) — leader+2 sync followers per shard, 16 shards.

tablekeynotes
accountsaccount_id PKidentity, registration_ts, abuse_disabled
workspacesworkspace_id PKname, plan, region, owner_id
channels(workspace_id, channel_id) PKname, type, retention, archived
channel_members(channel_id, account_id) PKrole, last_seen_seq, notification_pref
sessionssession_id PKaccount_id, expires_at

Why Vitess+MySQL: small data (~30 TB at this scale), strict membership invariants, need for live resharding (Slack's 2020 capability that survived the COVID surge). Failover RPO=0.

msgstore (ScyllaDB) — partition ((channel_id, day_bucket), message_id), RF=3 quorum.

columntypenotes
channel_iduuidpartition
day_bucketintpartition (rolls every 10d)
message_iduuidclustering, time-sortable (snowflake)
seqbigintper-channel monotonic
author_iduuid
contenttextup to 4KB
attachmentsjsonfile refs (ids only; bytes live elsewhere)
replied_touuid?thread parent
edited_tstimestamp?
deleted_tstimestamp?tombstone for delete-for-everyone

Why ScyllaDB: shard-per-core, predictable p99, no JVM compaction pause (the Discord migration thesis). Partition key bucketed by 10-day window so a 1M-subscriber channel's partition rolls every 10 days, capping width at 5–50MB. TWCS compacts old buckets and drops them as units at retention TTL. LWT only on the rare admin paths (channel-rename, pin-message); never in the per-message hot path.

presence (Redis cluster) — per-channel subscriber set subs:{channel_id} → {user_id, ws_gw_id, viewport_focus}. 64 shards, leader-follower with semi-sync replication.

07High-level design

Which components handle a request, and in what order?

The architecture in one paragraph. Clients hold a long-lived WebSocket to ws-gw via Anycast L4 LBs (NLB→Envoy WSS). On connect, ws-gw consults flannel for workspace bootstrap (cache hit absorbs >90% of the IDENTIFY load) and registers its channel subscriptions. On send, ws-gw routes the message to the channel-svc owning that channel (consistent-hash ring in coordinator); channel-svc validates ACL via meta-db, allocates the per-channel monotonic seq, writes the canonical row to msgstore (the durable commit), and fans out based on subscriber count: small/medium channels (<100K subs) get direct dispatch from channel-svc → each ws-gw with a subscriber → WS frame to the recipient device; mega-channels (>100K) route through the relay tier, which applies the passive-session filter (drop subscribers whose viewport doesn't include this channel) before dispatching. History fetch on chat-open routes through history-svc, which performs request coalescing on (channel_id, before_seq, page) so a 200K-user thundering herd becomes one ScyllaDB read.

Why two gateways. ws-gw is stateful (per-session GenServers / goroutine pools, sequence cursors, subscription anchors) and traffic is long-lived sticky-session WebSocket — a connection storm during a deploy can't be papered over by Envoy alone. api-gw is stateless HTTPS — it autoscales on CPU and absorbs auth/registration/CRUD bursts without affecting the chat data path. Slack and Discord both maintain this split for exactly this reason; conflating them means a search-storm or a bad auth-deploy takes down chat.

Why a relay tier (and where the cliff is). Below ~100K subscribers, channel-svc fans out directly: it iterates the subscriber set, batches by destination ws-gw, and dispatches over gRPC. Above ~100K, the channel actor's mailbox grows unbounded during the dispatch loop and BEAM scheduler share collapses (Discord 2017 #general meltdown, Maxjourney Nov 2023). The relay tier inserts itself between channel-svc and ws-gw, owns ~15K sessions/relay, and applies the passive-session filter (90%+ of subscribers are inactive viewers and should not receive a WS frame). This is the single most important pedagogical lesson in production chat scaling.

Why the viewport filter. Naive fanout delivers every message to every subscriber of a channel, whether or not they're looking at it. On a 1M-subscriber channel where ~90% follow but aren't focused, that's a 10× waste of WS-deliver work. Discord's fix: clients declare which channels are in their visible viewport (~20 entries on a phone, ~50 on desktop), and every send's fanout decision reads viewport_focus on the subscriber row and skips subscribers whose viewport doesn't include this channel. Without this single bit, a hot-channel fanout melts the gateway on every post.

Why request coalescing on history. A deploy that bumps 200K users' clients → 200K fresh IDENTIFYs → 200K history fetches in a 5-second window. Without coalescing, that's 200K identical ScyllaDB reads against the same partition — exactly the hot-partition pattern that drove Discord's pre-Scylla failure mode. The data-services tier collapses concurrent reads for the same key into one DB read whose result is multicast back. This is the load-bearing tier between deploy storms and database stability.

Why server-side history (not device-side). Slack and Discord both expose months/years of channel history that survives every device reinstall and is searchable workspace-wide. This rules out the WhatsApp / Signal model where history lives on devices and the server only caches undelivered. Server-side history forces:

  • A persistent message store (msgstore) that scales to PB-tier hot data.
  • A read tier in front of it (history-svc) to absorb cold-channel-open thundering herds.
  • A separate search index off CDC (modelled outside the diagram, referenced).
  • Per-channel retention configurable at admin level.

Why the recipient is drawn separately. The sender's ws-gw shard ≠ the recipient's ws-gw shard (consistent-hashed onto different boxes by user_id), and the WS carries traffic in both directions per device. Modeling the recipient as a separate client node makes the end-to-end fanout hop ('the message arrives on the other person's screen') explicit and traceable. Without it, fanout is hidden behind a self-loop on ws-gw.

Failure modes & mitigations

These come 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
1Cold-cache reconnect storm Monday-after-holiday (Slack 2021-01-04 — TGW saturation)TGW packet drops; WS handshake p99 cliff; IDENTIFY/RESUME ratio > 5:1Pre-warm capacity at known traffic-floor times; flannel POP-local cache absorbs IDENTIFY metadata reads; observability decoupled from production VPC
2Service-discovery sync bug → routing to dead boxes (Slack 2020-05-12)coordinator_membership_drift_count > 0 for 2mCoordinator (etcd/Consul) Raft is source of truth; ws-gw reconciliation loop bounded; alarms on stale-backend ratio
3Cache rollout cascades to DB (Slack 2022-02-22 — mcrouter / Vitess feedback)flannel hit-rate drop + meta-db QPS spike both risingThrottled config rollout (1%→10%→100%, 10-min soak); single-flight on flannel cache fills; per-team rate limit at api-gw
4Hot mega-channel fanout (Discord Maxjourney Nov 2023; #general 2017)Channel actor BEAM mailbox length; per-channel fanout p99; relay capacity utilizationRelay tier above 100K subs; passive-session filter at relay; per-channel send rate-limit (1k msgs/min)
5Tombstone-wall on un-bucketed partition (Discord pre-Scylla migration)tombstones_scanned_per_query_p99 > 1000; partition-size histogram p99Time-bucketed partition key already in place; TWCS compaction on hot tables; soft-delete via TTL bucket, not DELETE
6WS deploy breaks RESUME → IDENTIFY stormIDENTIFY/RESUME ratio > baseline×3; gateway CPU spikeBackwards-compatible RESUME across versions; jittered drain (60–600s); Slack's Disasterpiece Theater rehearses this exact path
7Hot history partition during deploy stampedehistory-svc cache-miss rate; per-(channel,day_bucket) read QPSRequest coalescing in history-svc; process-level cache TTL; ScyllaDB shard-loss tolerant via quorum
8meta-db sync-quorum self-DOS (both followers in AZ trip)vitess_sync_quorum_unmet_seconds > 10 per shardPatroni synchronous_standby_names = ANY 1 of 2 so a single follower healthy is sufficient; documented availability/RPO trade
9Edge-provider BGP outage (Cloudflare 2020-07-17 → Discord, Shopify down)Synthetic external probe failure across multiple ASNsMulti-CDN failover, anycast diversity; out-of-band management plane (Slack's 'cellular' migration was this lesson)
10Vitess hot shard from viral workspace signupPer-shard QPS skew; replica-lag spikeLive shard-split capability (Slack 2020 capability); pre-emptive hot-shard alarms via team_id QPS top-N
Observability

SLIs (RED + USE):

  • ws-gw: connection establish-rate, p99 connect latency, sockets-per-box, drain-duration, scheduler_utilization (BEAM), send→ack p99, IDENTIFY/RESUME ratio.
  • channel-svc: ACL-check p99, fanout p99 by channel-size cohort (small, medium, large, whale), per-channel mailbox length.
  • relay: subscriber-slice utilization, passive-filter skip ratio (target ≥85%), fanout dispatch latency.
  • history-svc: coalescing factor (post-coalesce QPS to msgstore vs incoming), cache hit-rate, per-(channel,day) heat map.
  • msgstore: write quorum-form rate, partition size histogram, compaction queue depth, tombstone scan ratio, p99 read/write per node.
  • meta-db: replication lag (must be ≤1s on hot path), failover-test cadence (weekly), shard QPS top-N skew.
  • presence (subscriber-set): shard primary failover count, viewport-update-write rate, subscriber-set eviction-rate.
  • coordinator: watch-fanout latency p99, membership-drift count.

SLOs (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 WS connect (warm) p99 < 350ms.
  • 99.99% on send → server ack p99 < 500ms (52 min/yr error budget).
  • 99.95% on send → recipient delivered p99 < 1s.
  • 99.99% on history-page p99 < 700ms.
  • Per-channel-cohort tail-latency SLI: p99 send-to-deliver bucketed by channel into (small <100, medium <10K, large <100K, whale ≥100K) cohorts. Anomaly alert when any cohort's 7-day p99 grows ≥2× baseline — catches a hot-channel creep before failure mode #4 actually fires. Propagation requirement (load-bearing): every send-fanout span carries a channel_size_cohort tag computed at channel-svc by snapshotting subs:{channel_id} cardinality at dispatch time; ws-gw and relay propagate the tag end-to-end so the SLI is computable across the whole fanout. Without explicit propagation, the cohort signal is only visible at channel-svc and the SLI is a paper claim.

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 (~14× error rate).
  • 5% in 6h → P2.
  • 10% in 3 days → ticket.
  • Cross-cohort gray failure: per-channel-cohort tail-latency p99 ≥2× 7-day baseline for 30m → P2 page even when no SLO is violated. Catches msgstore p99 doubling silently or a relay capacity problem before the whale cohort SLO actually breaches.

Dashboards on-call opens during a page:

  1. Edge & connection: connect-rate, p99 connect latency, sockets-per-box, IDENTIFY/RESUME ratio, TLS error-rate by SNI, TGW saturation.
  2. Send → ack pipeline: ws-gw → channel-svc → msgstore latency by hop; channel-svc per-cohort fanout p99.
  3. Fanout / Relay: relay subscriber-slice utilization, passive-filter skip-rate, channel actor mailbox length top-N.
  4. Storage: msgstore partition-size histogram, compaction queue, repl-lag, meta-db Vitess shard QPS skew.
Deploy & runtime
  • Rollout: ws-gw deploys via blue/green per-AZ — never region-wide simultaneous (the reconnect-storm risk is exactly the Slack-2021 lesson). Each AZ drained over 30 min via lame-duck reconnect-after: jittered(60–600s). Channel-svc and relay deploy with cooperative-sticky channel-ownership transfer (channels migrate to a sibling box over 10s without dropping fanout).
  • Schema migrations on Vitess: dual-write + shadow-read + backfill + flag-flip; online ALTER is forbidden on hot tables.
  • ScyllaDB schema: additive only; new columns default; column drop = ignore-on-read for one major version, then physical drop in maintenance window.
  • Disasterpiece Theater (Slack's chaos rehearsal practice): every quarter, on-call walks through a documented failure-mode scenario in a controlled production-mirror, before the actual incident hits. The chaos card list in this canonical maps directly to the rehearsal scripts.
  • Config plane (intentionally not a request-path node): Statsig-style flags with 1%→10%→100% staged push, 10-min soak between stages, kill-switch on out-of-band etcd so a bad config push can't disable its own off-button — this is the explicit Slack-2022-02-22 lesson.
Multi-region / DR

> Diagram convention: cross-region edges and Slack-2024 cell-router paths are intentionally not drawn — the canonical foregrounds the channels/history/fanout main flow. The RPO / RTO contracts below reference observer/cell-router topology that exists behind the visible diagram nodes; chaos card slack-discord:cloudflare-bgp-outage tests whether a deployer has multi-CDN diversity even though only one edge-lb node is drawn.

  • Active-active across 18 regions. Workspace home_region set on workspace creation; ws-gw in any region can serve any user, but writes funnel home-region for sync-replicated metadata.
  • ws-gw: stateless w.r.t. persistent data, replicated per region. Anycast routes to nearest healthy POP.
  • channel-svc: in-region; channels owned by their workspace's home-region CS. Cross-region routing via coordinator's hash ring (workspace_id → home_region → CS).
  • 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 ≤ 5 min.
  • meta-db: in-region 1 leader + 2 sync followers (RPO=0 / RTO ~30s automated). Cross-region async replication (Vitess vreplication) for follow-the-sun fallback. Cross-region failover is the only operation that human-confirms.
  • flannel: regional (POP-local). Loss of a region drops cache; rebuilt from meta-db on next IDENTIFY (the cold-cache risk).
  • Cellular architecture (Slack 2024): each region partitioned into independently-failing cells of ~5K teams each. A cell can be drained in 5 minutes. The InfoQ 2024 article documents this explicitly; we don't model individual cells in the diagram but they exist as a sub-deployment of every regional tier.
Security posture
  • NOT E2EE on content: server reads message bodies (search, abuse pipeline, admin compliance). Documented as the load-bearing differentiator from WhatsApp / Signal.
  • mTLS between every internal service.
  • JWT on the WS for session validation (verified locally on ws-gw with cached public key — auth tier is not on every-frame critical path).
  • WAF + per-IP rate limit at edge-lb on control APIs; per-token rate limit at api-gw.
  • GDPR erasure cascades through meta-db → msgstore (tombstone with TTL) → search index rebuild. Eventually-atomic, audited monthly.

08Deep dives

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

1. Hot-channel fanout — the cliff at ~100K subs and what replaces it. Below ~100K subscribers, channel-svc iterates the subscriber set, batches by destination ws-gw, and dispatches over gRPC — O(N/batch_size) hops on the channel actor's main loop. At ~100K subs the actor's BEAM mailbox starts growing during the loop because incoming messages on the same channel queue while the dispatch is running. Past that, fanout p99 spirals into seconds (Discord 2017 #general meltdown — exactly this). The fix: relay tier inserts a fan-out-of-fan-outs. channel-svc dispatches once per relay (~70 dispatches for a 1M-sub channel @ 15K sessions/relay); each relay parallel-dispatches to its slice of ws-gw boxes. The passive-session filter at relay is what really lets it scale — Discord's published 90% passive ratio means a 1M-sub channel becomes a 100K effective fanout, not 1M.

2. Request coalescing on history reads. After a deploy, 200K users reconnect → 200K IDENTIFYs → 200K cold history fetches in a 5-second window. Without coalescing, that's 200K identical ScyllaDB reads against the same (channel_id, day_bucket) partition — exactly the hot-partition pattern that drove Discord's pre-Scylla failure mode. The data-services tier collapses concurrent reads for the same key into one DB read whose result is multicast back. The coalescer is single-process per host, keyed on (channel_id, before_seq, limit) with a 1-second waiter window (kept under ws-gw's e8 1500ms timeout minus msgstore p99 read + serialization). Empirical coalescing factor: 8–20× normal, 100×+ during deploy stampede. Without this tier, the canonical needs 200+ ScyllaDB nodes; with it, 96.

3. Idempotent retries through region failover. Client sent a message; got no ack because edge-lb's anycast IP shifted regions mid-send. Client retries with the same message_id. The new region's ws-gw → channel-svc dedup-checks: channel-svc keeps a 5-minute LRU of message_id → seq for accepted sends. Redis sees the SETNX failed → returns the original ack. No double-write to msgstore, no double-fanout. The 5-minute window is short by chat-system standards because the recipient-side dedup also runs (clients dedup by message_id in the visible-channel timeline). 7-day TTL like WhatsApp would be over-engineering for non-E2EE chat.

4. ScyllaDB partition-key bucketing — the single most-important storage decision. Discord 2017: #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 = ((channel_id, day_bucket), message_id) where day_bucket = floor(server_ts / 10d). A 100K-member channel sending 10/sec produces ~864K messages/day = ~8.6M/10-day-window → ~2.4GB partition (still large, but bounded; new partition every 10 days). TWCS compacts old buckets and drops them as units at TTL. Without this, no other mitigation matters.

5. The channel-svc-mid-fanout-crash trade and how the seq-gap reconciler closes it. Unlike WhatsApp (which has a cdc-relay worker tailing the storage CDC log to close an outbox transactionally under chatd-crash), this canonical accepts that a channel-svc crash between e10 (msgstore commit) and e13/e14 (in-process WS fanout) loses the in-flight fanout for that one message. The message is durable in msgstore — the loss is liveness, not durability. Three recovery cases:

  • Reconnecting recipient (lost their socket entirely) → on RESUME, ws-gw replays from last_received_seq via history-svc. Covers naturally.
  • Offline recipient → user opens app on next session, history-svc backfills from the last visible seq. Covers naturally because clients pull history on chat-open.
  • Online focused recipient (socket alive, was actively viewing the channel) — this is the case the no-CDC-closer trade misses. The reconciler is client-side: every client tracks last_received_seq per channel, and on detection of a seq gap (next received > expected+1) issues a history-svc backfill for the missing range. The gap is detected on the next successful fanout to the same channel (which arrives whenever the next user posts), bounded by typical channel cadence at <1 minute.

This is the explicit trade we accepted in §tradeoffs: ~30s–1min visible delay on the rare channel-svc-crash case in exchange for not building a CDC-based closer (operational cost we avoid since chat history is non-E2EE and cheap to re-fetch). Chaos card slack-discord:channel-svc-mid-fanout-crash rehearses exactly this path.

09Trade-offs

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

The big architecture trade-offs we accepted:

  • Channel-svc as single-writer per channel. Per-channel ordering is a load-bearing UX claim (reply-to-message would break under reordered delivery). We accept that a hot channel is bottlenecked on its owning actor — and pay for it with the relay tier above 100K subs. The alternative (multi-writer with cross-actor ordering protocol) was rejected: too much complexity, and ordering is the product.
  • Server-side history (not E2EE). The single biggest difference from WhatsApp/Signal. Slack and Discord both made this product call: server-side search, admin compliance, and shared-history-on-new-device require it. The cost: messages are visible to ops and to legal subpoena. We accept this; the user-facing privacy claim is "encrypted in transit + at rest, accessible to admins of your workspace."
  • Selective subscription by viewport (not deliver-to-all). We accept that a subscriber whose viewport doesn't include a channel gets no live WS frame and only sees the message when they focus the channel (history pull). The alternative (deliver every message to every follower) melts the gateway on hot channels.
  • Relay tier is mostly idle. 120 relays exist for the 0.001% of channels that need them. We could share the boxes with another workload in lower-traffic times, but the cold-start cost at relay-cliff time is too high — keep them warm.
  • No CDC-closure on the in-process fanout. Unlike WhatsApp (which has cdc-relay to close an outbox transactionally under chatd failure), we accept that a channel-svc crash between the msgstore commit (e10) and the in-process WS fanout (e13/e14) may lose the fanout for that one message. The recipients converge via the client-side seq-gap reconciler (deep-dive #5): every client tracks last_received_seq per channel and on a detected gap issues a history-svc backfill. Reconnecting and offline recipients are covered automatically; online focused recipients see ≤30s–1min delay until the next channel fanout exposes the gap. The cost is acceptable for non-E2EE chat where history-fetch is cheap; chaos card slack-discord:channel-svc-mid-fanout-crash rehearses this exact path.
  • Offline push / audit deliberately scoped out. This canonical stops at online WS fanout + server-side history; the offline-notification pipeline (channel-svc → Kafka push.requests → a push worker → APNs/FCM) and the compliance audit.events topic are out of scope. Both layer back on the same send spine without touching the core: channel-svc gains one more fire-and-forget publish and a new consumer tier drains it — the send path's durability and per-channel ordering are unaffected.

What we rejected and why:

  • Cassandra (pre-2023 Discord): JVM GC pauses + LSM compaction starvation made p99 unpredictable. Discord's published migration thesis stands.
  • DynamoDB for messages: rigid query patterns (no WHERE seq > X), $$ at 2.7M deliveries/sec, harder hot-partition mitigation.
  • MongoDB: was Discord's original choice (pre-2017); they outgrew it. Documented in their blog.
  • Spanner / CockroachDB for messages: serializable cross-shard txns are not what messaging needs. Pay 5× write latency for nothing.
  • Long-poll fallback: explicitly rejected — modern clients are WS-capable; supporting long-poll would couple to legacy infrastructure that's already painful at Slack scale.
  • Direct WS from edge-lb to client (no Envoy WSS layer): NLB has no TLS termination; we'd need to terminate in the gateway box, doubling its CPU load. Envoy WSS is the explicit Slack lesson (their blog "Migrating Millions of Concurrent WebSockets to Envoy").
Trade-offs we accepted (open questions)
  • Search latency on hot channels: the async-populated Elasticsearch tier (off msgstore CDC) lags by 30s–5min during heavy ingest. Users searching for a message they sent 30s ago may miss it. Acceptable for chat search; not acceptable for "find channels with a recent name change" which uses meta-db directly.
  • Cell-failure isolation: Slack's cellular architecture (2024) limits blast radius to ~5K teams per cell. We don't model individual cells in the diagram; one regional tier represents many cells. The actual blast-radius math depends on cell size + drain time.
  • Mega-channel relay cold-start: when a previously-medium channel crosses 100K subs, channel-svc must dynamically promote it to the relay tier. The promotion takes ~3s and during that window, fanout is degraded (small-channel path overwhelmed). Acceptable trade-off; rare in practice.
  • Search index rebuild on GDPR erasure: deleting one user's messages from Elasticsearch requires a full re-index of affected channels. This is slow (hours). Acceptable since GDPR's 30-day window comfortably absorbs it.
Trace catalogue

The trace cards exposed by the simulator for this canonical:

  • slack-discord:ws-connect-resume — Warm WS reconnect via RESUME (cheap; 0 backfill). Goes client → edge-lb → ws-gw → flannel → ws-gw and ends at the gateway. Budget 350ms.
  • slack-discord:send-channel-fanout — Send to a small/medium channel with online subscribers; durable write to msgstore, in-process fanout dispatch (no relay). Path ends at msgstore. Budget 500ms.
  • slack-discord:deliver-megachannel-via-relay — Same send shape but channel has >100K subs; routes through the relay tier, passive-session filter applied. Path ends at recipient via the relay path. Budget 1s.
  • slack-discord:fetch-history-coalesced — User opens a channel; history-svc collapses concurrent waiters into one ScyllaDB read. Ends at msgstore via history-svc. Budget 700ms.
Failure scenarios we model

The chaos cards exposed by the simulator for this canonical, each grounded in a published incident:

  • slack-discord:hot-channel-fanout-meltdown — Channel-svc actor mailbox saturation on a Maxjourney-scale channel (Discord Maxjourney Nov 2023). Mitigated by relay tier; tests whether the user's diagram has one.
  • slack-discord:cassandra-tombstone-wall — Tombstone explosion on un-bucketed history partition (Discord pre-2023 migration). Tests whether the user's msgstore has time-bucketed partition keys.
  • slack-discord:cold-cache-reconnect-storm — Monday-after-holiday TGW saturation + IDENTIFY storm (Slack 2021-01-04). Tests whether flannel-style edge cache absorbs the storm.
  • slack-discord:vitess-cache-feedback-loop — Mcrouter rollout cascades to Vitess (Slack 2022-02-22). Tests whether the user has staged config rollout + single-flight cache fills.
  • slack-discord:cloudflare-bgp-outage — Edge-provider BGP failure takes down anycast (Cloudflare 2020-07-17 → Discord). Tests whether the user has multi-CDN failover. Note: only one edge-lb is drawn in the canonical; multi-CDN diversity is an implicit topology assumption read out of edge-lb.failureMode.
  • slack-discord:ws-deploy-resume-storm — Bad gateway deploy invalidates RESUME tokens; every client falls back to IDENTIFY (5–10× cost). Tests whether the user's deploy story has dual-version RESUME compatibility.
  • slack-discord:channel-svc-mid-fanout-crash — channel-svc commits to msgstore (e10) then crashes before the in-process fanout (e13/e14). Tests whether the user has the client-side seq-gap reconciler (see deep-dive #5) or a CDC-based outbox closer — the explicit hand-wave from §tradeoffs.

Primary sources

  • Discord — How Discord Stores Trillions of Messages (Cassandra → ScyllaDB)
  • Discord — How Discord Scaled Elixir to 5,000,000 Concurrent Users
  • Discord — Maxjourney: 1M+ Online in a Single Server (relay tier)
  • Discord — Using Rust to Scale Elixir for 11M Concurrent Users (SortedSet NIF)
  • Slack — Flannel: an Application-Level Edge Cache to Make Slack Scale
  • Slack — Real-time messaging (Channel Server + Gatewayserver)
  • Slack — Scaling Datastores at Slack with Vitess
  • Slack — Slack's Outage on January 4th, 2021 (TGW saturation)
  • Slack — Slack's Incident on 2-22-22 (Consul / Vitess feedback)
  • Slack — A Terrible, Horrible, No-Good, Very Bad Day (May 2020 HAProxy)
  • Slack — Migration to a Cellular Architecture (InfoQ 2024)
  • Slack — Tracing Notifications & How Slack Rebuilt Notifications
  • Slack — Migrating Millions of Concurrent WebSockets to Envoy
  • elixir-lang.org — Real-Time Communication at Scale with Elixir at Discord (2020)
  • Cloudflare — July 17 2020 BGP 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 Slack / Discord yourself

Build the primitives this design leans on

Each one is an animated curriculum that constructs the system from scratch.

More in Real-Time: Chat, Presence & Live Updates

Long-lived connections and the things you push down them — messages, cursors, green dots, scores, prices and bids.