Online Indicator — a worked solution
Green dot for contacts. Mind the N² watch problem. Approximate by design — never quote presence more precisely than reality.
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 Online Indicator workspaceThe problem
Build the Online Indicator — the green dot next to every contact's name in your chat app, plus the "last seen 12 min ago" line that follows it. Used by Slack, Discord, WhatsApp, Messenger, Teams, Instagram, LinkedIn — every product with a contact list. Sounds simple. It is not.
Three things make this problem hard at production scale:
- N² watch problem. Every user is both a presence publisher and a presence subscriber. At 1B DAU each with ~500 contacts, the naive eager-broadcast design produces 500B subscription edges and ~250M frames/s of fanout. Physically impossible.
- Ephemeral state at billion-user RAM scale. ~500M concurrent online users × ~100 B/record = ~50 GB per region — hot, in RAM, all the time, fully replicated.
- Approximate by design. The product line "last seen 12 min ago" is a privacy + capacity invariant, not a UX choice. Quoting
last_seenprecisely (±1 s) leaks PII and forces hot-path persistence; quantizing to ±60 s solves both.
The architecture is the canonical Discord/Slack/WhatsApp shape: a long-lived WebSocket to a stateful gateway, a sharded actor tier holding hot presence in RAM, a separate fanout relay tier as the N² solver, an ephemeral cache for cross-shard reads, a durable spine for the audit trail, and a quantized last-seen store for the offline case.
The reference architecture
What each component is for
- Client (SDK)iOS / Android / Web SDK
Maintains a single persistent WebSocket to the edge for the entire app session. Multiplexes three logical streams over that one socket: (1) outbound app-level PING every 60 s; (2) outbound viewport subscribe / unsubscribe messages as the user scrolls a contact list, opens a DM, or focuses a channel; (3) inbound presence_update push frames from fanout-relay. Renders the green dot, the 'last seen X min ago' label, and the 'typing…' indicator from the same socket. Persists the last-rendered presence map to local SQLite so a foreground/background bounce doesn't blank every dot for 200 ms before the resubscribe completes.
Why it exists. Without a long-lived WebSocket, every viewport scroll would be a fresh HTTPS handshake and a fresh subscription round-trip — at 1B DAU × ~10 scroll-into-view events per minute that is 167 M handshakes/s, an order of magnitude beyond any sane gateway tier. The WS keeps the subscription set warm; only deltas (subscribe-N, unsubscribe-M) cross the wire. The SDK is also where the N² watch problem first gets bled off: the viewport observer reports at most ~30 contacts to subscribe to, never the full 500-person contact list.
When it fails. If decorrelated jitter is missing or the cap is set too low, the client itself becomes the post-region-heal weapon: synchronized retries hammer the gateway in 1-second bursts and the system never recovers. If the SDK doesn't persist last-rendered state, every backgrounded mobile app shows a flicker of 'everyone offline' on foreground — users perceive this as 'the app is broken' even though the backend is healthy. If viewport-tracking misses an unsubscribe, the subscription-store grows unbounded and the user effectively eager-broadcasts.
- Edge LB (anycast)Cloudflare Spectrum / AWS Global Accelerator + Envoy
Anycast-routed L4/L7 layer that terminates TLS for the WebSocket Upgrade, applies per-IP and per-ASN rate limits on new connection attempts, runs the WAF rule set on the HTTP Upgrade headers, and forwards the established WS to the nearest healthy ws-gw shard. Carries the connection for its entire lifetime — typically 30 min median, 12 h p99 — using long-lived L4 tunnels back to the regional gateway tier so TCP state survives Envoy hot-restarts.
Why it exists. Two reasons production teams put a dedicated anycast layer ahead of the WS gateway: (a) absorb volumetric DDoS before it touches stateful gateway processes — a SYN flood at the gateway tier triggers connection-table exhaustion and takes down legitimate sessions, while at the anycast scrubber it's silently bucketed into per-prefix rate limits; (b) route the user to their geographically-nearest gateway POP for sub-100 ms WS handshake latency. The 80-replica count is sized to absorb 2× the 500 M peak concurrent at ~12.5 M L4 tunnels per box — pass-through flow state is a conntrack entry, not an app-layer socket, so the per-box ceiling is bandwidth-bound rather than the 5 M Discord / 2 M WhatsApp gateway-process ceilings (those apply at ws-gw).
When it fails. Provider BGP withdrawal (Cloudflare 2020-07-17 type) drops the whole edge layer; in-flight WS sessions survive on the LB's TCP backbone for ~minutes but new connections fail. Mitigation: multi-CDN / multi-anycast provider with DNS-level failover (cap 60 s TTL). If a single POP loses anycast announcement, traffic shifts to neighbor POP with cold cache at downstream — gateway sees a spike in IDENTIFY (not RESUME) ratio because the new POP's gateway shard doesn't know these users. Connection-table exhaustion at the tunnel-table ceiling manifests as ECONNREFUSED, not slow handshakes — discrete failure mode, easy to detect.
- WS GatewayElixir Phoenix Channels (or Go + gorilla/websocket)
Authenticates the WS Upgrade by validating a short-lived JWT (issued by identity-svc on app login), establishes a session record {conn_id, user_id, device_id, region, capabilities} in local process memory, and routes every inbound message frame to the correct backend by message type: heartbeats and presence-write go to presence-svc on a sharded gRPC stream (shard key = hash(user_id) mod presence_svc_shards), subscribe/unsubscribe go to presence-svc, and inbound presence_update frames coming back from fanout-relay are matched to the connection by conn_id and pushed back to the client. Drains gracefully on rolling deploy with an application-level RECONNECT frame carrying a jittered 60–600 s delay, followed by Close code 1001.
Why it exists. Hands stateful connection ownership to a dedicated tier so the backend services (presence-svc, fanout-relay) can be addressed by user_id, not by network address. Without a gateway, every backend service would need to maintain its own WS connection map — duplicated state and a synchronization problem on every connect/disconnect. The gateway also gives us a single chokepoint to enforce per-user rate limits (heartbeat spam protection, subscribe-flood protection) before the backend ever sees the traffic. 200 replicas at ~5 M sockets per box (BEAM/Elixir, per Discord's 2017 published number) gives ~1B concurrent ceiling — sized for 500 M concurrent online at half-saturation per region (50% headroom for 1-AZ-out + Monday-morning reconnect storm).
When it fails. A bad deploy that rejects RESUME tokens minted by the previous version (forgetting the N±1 compat rule) causes every client to fall back to IDENTIFY, which is 5–10× more expensive (full session push instead of delta resume). Gateway CPU saturates within 60 s of the rollout; if the rollout was not throttled, all shards cliff together. Detection: identify_rate / resume_rate ratio > 3× baseline. Mitigation: staged rollout 1% → 10% → 100% with 10-min soak per stage and an automated rollback on the ratio. Per-shard mailbox growth from broadcast amplification storms (covered by the fanout-relay's existence) is the failure mode this tier explicitly delegates downstream.
- Coordinator (etcd / Raft)etcd 3.x (5-node Raft), multi-AZ
Holds three load-bearing keyspaces, all replicated via Raft for linearizable reads: (1) presence_svc_ring/{shard_id} → {owner_instance, epoch, lease_expires} — the per-shard ownership lease that ws-gw and fanout-relay watch to route traffic; (2) broadcaster_salts/{user_id} → {salt_count} for celebrity users requiring sub-sharded watcher fanout; (3) feature_flags/* — kill-switches and gradual-rollout config. Every other tier WATCHES the relevant keyspace prefix; values are versioned with mod_revision so a watch can resume exactly from the last seen revision without missing updates. Read traffic is ~5K reads/s (mostly cache fills from ws-gw discovery), write traffic is ~50 writes/s (lease renewals + rare rebalances) — comfortably within etcd's published 10K writes/s ceiling.
Why it exists. Without an explicit coordinator, shard ownership is implicit in DNS / config-management / process-startup — and that is precisely the failure mode Slack's 2022-02-22 incident retro documents (mcrouter rollout cascading because Consul ownership data wasn't where the doc claimed). At 64 presence-svc shards, 128 fanout-relay shards, 32 subscription-store shards, and 200 ws-gw replicas, NO HUMAN can keep this in sync; only a Raft-backed registry can. The Slack 2022-02-22 retro specifically says the cure for cache-rollout cascade was a kill-switch on an OUT-OF-BAND etcd path — we model that explicitly here.
When it fails. Quorum loss (2 of 5 nodes lost simultaneously) → etcd reverts to read-only; existing shard ownership leases stay valid for their TTL but no new leases can be granted. Shard rebalances pause; presence-svc shard crashes during the gap cause that shard's users to be unservable until quorum returns. RTO: 5–15 min p99 (etcd recovery procedure). The cure is the out-of-band emergency-config etcd which can be promoted to primary by a human in a documented runbook. Slow leader (GC pause, disk IOPS saturation) is the more common failure — watch latency climbs from 1ms to 1s, all tiers' routing tables go stale, fanout drops up to 10s of frames per tier. Detection: etcd_disk_wal_fsync_duration_seconds_p99 > 100ms.
- Presence ServiceElixir GenServer per-shard, sharded by hash(user_id)
Receives heartbeats (op=write, refresh TTL), online/offline transitions (op=write, plus publish to events-bus), subscribe/unsubscribe (op=write to subscription-store, op=read of presence-cache for the initial snapshot), and presence reads (op=read). Each shard owns 1/64 of the user keyspace by hash(user_id). Per-shard GenServer process keeps a hot in-process LRU of the most recently touched 100 K presence records, falling through to presence-cache (Redis Cluster) for the cold tail. Detects offline transitions in two ways: (a) explicit client message on graceful disconnect; (b) an in-memory heartbeat-timeout sweeper in the per-shard actor that fires every 5 s and emits an offline event for any user whose last in-memory heartbeat is older than the 120 s liveness threshold (heartbeat × 2). The latter is the load-bearing path — 70% of offline events come from heartbeat timeout, not graceful disconnect. Detecting offline in the actor's memory (not via a Redis TTL) is what lets heartbeats skip Redis entirely — Redis is written only on transitions.
Why it exists. The actor tier is the answer to 'where does presence state live?' The naive answer is 'in Redis' — but every heartbeat would then be a Redis round-trip and every subscribe would be a Redis multi-get. At 16 M heartbeats/s and 500 K transitions/s, that's beyond what a Redis Cluster can sustain even with aggressive pipelining (per the capacity analysis: ~100 K ops/s/shard × shards). Holding hot state in a sharded per-user actor lets us serve the 99% case from process-local ETS / in-memory map and only touch Redis for cold reads and durable trail writes. This is the Discord/WhatsApp pattern. The 64-shard count is sized so that one shard owning ~8 M concurrent users at peak fits comfortably in one BEAM node's RAM (~10 GB hot state).
When it fails. A bad shard rebalance during deploy causes presence ownership to bounce between two shards; users observed see flapping 'online/offline' for ~30 s as the two shards disagree about TTL state. Mitigation: explicit drain protocol — the coordinator (etcd) holds the shard lock; shard A finishes serving in-flight before shard B accepts. Per-shard mailbox unbounded growth on a celebrity-going-online event (100M watchers' subscriptions all pointing at the same shard) is mitigated by debouncing transitions at the source and by fanout-relay being the actual broadcaster, NOT presence-svc — presence-svc just publishes ONE event regardless of watcher count. TTL-expiry sweeper falling behind under load (the heartbeat-backpressure cascade) causes mass false-offline; backstop is the explicit-disconnect signal from ws-gw which short-circuits the sweeper.
- Fanout RelayElixir GenStage, sharded by hash(watched_uid)
Subscribes to presence_transitions from presence-svc and, for each transition, looks up the watcher set in subscription-store (reverse index watchers:{watched_uid} → set of {conn_id, ws_gw_shard}), applies the passive-session filter (90%+ of watchers' viewport doesn't currently include this user; those frames are dropped at the relay, NOT at the client), coalesces multiple transitions of the same user within a 1 s window into a single delta frame, and dispatches the surviving frames to ws-gw shards by conn_id. Each relay shard owns 1/128 of the watched-user keyspace.
Why it exists. This is the architecture's load-bearing answer to the N² watch problem. If presence-svc dispatched fanout directly, every transition for a celebrity user (50 M watchers) would land in one presence-svc actor's mailbox, growing it unbounded and starving every co-located user on the same shard — the Discord Maxjourney moment recreated. Splitting fanout into its own sharded tier — sharded by watched_uid, NOT by watcher_uid — gives every popular user their own dispatcher that can be horizontally scaled and that has a bounded mailbox per-actor. The passive-session filter is the second N² lever: 95%+ of subscriptions are inactive at any moment (the user's viewport is on a different screen), so dropping those frames at the relay rather than pushing them to ws-gw and the client cuts wire traffic by 20×.
When it fails. Subscription-store unavailable → relay can't look up watcher set → frames are dropped at the relay (fail closed, never fail open with a wildcard fanout). Subscription-store partial outage (one shard down) means 1/N of watchers don't see the transition — bounded blast radius, lazy reconciliation: when the user's next viewport scroll triggers a fresh subscription, the initial snapshot read covers the gap. Relay shard crash means its in-flight transitions are lost (RPO ~1 s, the coalesce window); the durable trail in events-bus is the recovery vehicle for any consumer that needs replay (analytics, abuse). The relay does NOT replay events-bus on restart — presence is approximate by design and a 5 s gap during failover is within budget.
- Presence CacheRedis Cluster 7.x with io_threads, AZ-pinned
Stores one Redis hash per online user: presence:{uid} → {status: 'online'|'away'|'invisible', last_seen_ts, region, device_hint, privacy_mask} with a long backstop TTL (~15 min) — NOT a per-heartbeat 120 s TTL. Liveness lives in the presence-svc actor's memory (the 120 s heartbeat×2 timeout is the actor's in-memory sweeper threshold, not this key's TTL); the key is set on the online-transition and flipped/removed on the offline-transition the actor raises. The backstop TTL only reaps keys leaked by an actor crash that lost its recovery. This is the hot copy of 'is user X online RIGHT NOW' and of the last_seen cached value (the persisted-to-disk last_seen lives in lastseen-db). Presence-svc writes are SETEX on transitions (single round-trip atomic write+TTL); heartbeats do not write here. Presence-svc reads on cache-cold are GET → fallthrough to lastseen-db. 24 shards × ~4.2 GB primary working set on 16 GB instances gives ~100 GB primary + two followers ≈ ~300 GB cluster-wide RAM per region (the B1-corrected capacity math), with the 16 GB instance size providing 2× growth + warm-restart headroom.
Why it exists. Couldn't this live in presence-svc's in-process memory? Yes, for the 99% hot path, and that's exactly where it does live (the actor-tier ETS LRU). The cache exists to handle three cases the actor-tier alone can't: (1) cold-start a presence-svc shard after deploy or crash without losing every user's online state; (2) cross-shard reads — presence-svc shard A asking 'is user X (on shard B) online?' goes via cache, NOT via cross-actor RPC, because that would couple shards and break the no-cross-shard invariant; (3) lookup-by-many-keys on subscribe — the initial snapshot for a viewport of 30 contacts is one MGET, not 30 cross-shard RPCs. Redis Cluster (not Redis Sentinel) because the keyspace grows past one box's RAM at 1B DAU scale.
When it fails. Whole-cluster cache restart (cold-cache scenario) drops the hot working set; presence-svc actors get every read as a fallthrough to lastseen-db, which is sized for 1% of cache load not 100%, so lastseen-db cliffs within seconds; explicit mitigation = read-through circuit breaker at presence-svc that, on detected cache-cluster failure, serves last-rendered state from the actor's in-process LRU and serves 'unknown' (NOT 'online') for misses. Per-shard primary crash → 15 s RTO during which writes to that 1/24 of users are blocked; clients see their heartbeat fail and retry with jitter. Hot-key (Cristiano Ronaldo) is NOT a problem on this tier because we never store the fanout list in the hot-key — we store presence(user) once, watchers are in subscription-store, fanout-relay reads subscription-store NOT presence-cache for the fanout list.
- Subscription StoreRedis Cluster 7.x with AOF, multi-region
The inverse index that powers the fanout. Stores watchers:{watched_uid} as a Redis hash mapping conn_id → {ws_gw_shard, viewport_active_bit, subscribed_at}. On subscribe (presence-svc adds a viewer to the watchers set), this is an HSET. On unsubscribe, HDEL. On every transition, fanout-relay HGETALLs the set, filters by viewport_active=1, and dispatches. Cleanup is field-level, not key-level: when ws-gw reports a session close, presence-svc HDELs that conn_id's field from every key the session subscribed to, and a background sweeper HDELs fields whose subscribed_at watermark has gone stale (ungraceful disconnects the close signal missed). The key-level TTL (1800 s, re-extended on writes) is only a whole-set backstop — Redis hashes have no per-field expiry, so a key TTL alone could never reap one dead watcher while other watchers keep the key alive.
Why it exists. Distinct from presence-cache because the access pattern is fundamentally different: presence-cache is keyed by presence:{uid} for 'is X online?' (one-row reads at heartbeat QPS); subscription-store is keyed by watchers:{uid} for 'who watches X?' (one-set reads at transition QPS). Co-locating them on the same Redis Cluster would mean a hot transition would compete with hot heartbeats for the same shard's CPU — separated, we can tune each independently. Also distinct from graph-db because graph-db holds the DURABLE contact-adjacency (mutual-contact + privacy preferences), while subscription-store holds the SESSION subscription (who's actively watching right now). Subscribing is gated by graph-db at subscribe-time; once subscribed, fanout reads subscription-store at transition-time.
When it fails. Subscription-store shard down → fanout-relay can't read watchers for that 1/32 of the keyspace → those users' transitions silently don't fan out → observers see the green dot for up to 5 s longer (within SLO budget). Lazy reconciliation: when the observer's next viewport scroll triggers a fresh subscribe, the initial snapshot read covers the gap. AOF rewrite under load causes write latency to spike — mitigation is rewrite during low-traffic windows + per-shard rotation. Subscription-store full (memory pressure) means subscribe writes get rejected with OOM — fail closed at presence-svc with a 503 to ws-gw, client retries with jitter. Hot-shard CPU on a celebrity-going-online: the salted-sub-shard split mitigates the read fanout but the WRITE fanout (50 M watchers all subscribing in a viewport-scroll storm) requires per-broadcaster write throttling at presence-svc.
- Identity GraphSharded Dgraph (36 shards) with read-through graph-cache at presence-svc
Stores the durable social graph: nodes are users, edges are contact/friend/follow relationships with edge properties {mutual: bool, blocked: bool}. Each user node also carries privacy preferences: presence_visibility ∈ {everyone, contacts, mutual_contacts, nobody}, last_seen_visibility ∈ {everyone, contacts, nobody, except_blocked}, and a per-user 'invisible mode' override. Queried by presence-svc on every SUBSCRIBE gesture to determine whether the viewer is allowed to subscribe to the target's presence — viewers in violation get a synthetic 'unknown' status, NOT a 403 (privacy by obscurity, not by error code, so the abuser can't probe).
Why it exists. Privacy is first-class for any presence service — see WhatsApp's 2012 privacy retrofit (the 'two ticks' / 'blue ticks' lineage). A flat KV store would force privacy checks at every fanout, on the hot path; a graph store keeps the adjacency natural and lets us cache the resolved 'visible_to_user' set per (target_uid, viewer_uid) pair with long TTL. Neo4j (or sharded Dgraph) for the contact graph specifically because the dominant query is 1-hop: 'is viewer mutual-contact of target?' — that's a single edge traversal, exactly what graph-dbs are built for. NOT used for the fanout watchers-list query (that's subscription-store) — the graph-db is the DURABLE contact relationship; subscription-store is the EPHEMERAL session subscription.
When it fails. Graph-db unavailable → presence-svc fail closed on subscribe (reject with 503, retry-after 30 s, NOT 'allow by default') — the 99% case where graph-db is up but slow degrades to longer subscribe latency, not to permission bypass. Causal cluster leader election (5–15 s RTO) during failover means new subscribes hang briefly; existing subscriptions are unaffected (subscription-store is independent of graph-db at runtime). Privacy-cache staleness window: a 'block this user' takes up to 5 min to propagate to fanout-relay (cache TTL); explicit push-invalidate on block/unblock cuts that to <1 s for the immediate-impact path. Bulk-import abuse — a malicious user attempting to read every contact in the graph — is rate-limited at presence-svc per-user (the 'no probe API' rule).
- Transition SpineKafka with idempotent producers + tiered storage
Single durable topic presence.transitions carries every online/offline/away/invisible transition emitted by presence-svc. Partition key = user_id so all transitions for one user are ordered. Three consumer groups subscribe: persist-worker (writes durable last_seen to lastseen-db), analytics-worker (writes minutely rollups to analytics-db — implicit in this canonical, folded into persist-worker for diagram clarity), and abuse-worker (off-canonical here; consumes anomalous flap patterns for trust & safety). Retention 7 days hot (Kafka brokers) + 30 days cold (tiered storage on S3) — the abuse and analytics consumers can rewind for backfill within 7 days; beyond that they pay the S3 read cost.
Why it exists. The durable trail is the architecture's event-sourced spine: presence-svc publishes a transition to events-bus FIRST (ack=all), then projects it to presence-cache and fanout-relay — the log is the source of truth and every other store is a derived projection. The persist-worker is the deferred consumer that pulls from events-bus and writes lastseen-db, decoupling the durable write from the hot path. Without the spine, presence-svc would have to dual-write to presence-cache AND lastseen-db synchronously, doubling write latency and creating a cross-DB consistency problem. With it, presence-svc publishes to Kafka (fast, ack=all at 10 ms), projects to presence-cache (fast, 5 ms) and returns; lastseen-db lags by ~2 s but that's well within the 'last seen ±60 s' SLO.
When it fails. Kafka cluster down → transitions queue in presence-svc's bounded producer buffer (the hot read path stays healthy — liveness is actor-tier state); a prolonged outage overflows the buffer and drops the oldest transitions from the durable trail, and the e14b replay reconverges fanout-relay projections once the cluster recovers. Consumer lag spike on persist-worker → durable last_seen lags by minutes-to-hours; users see slightly stale 'last seen' values but presence (online/offline) is unaffected. Partition leader migration during broker bounce: producers see ~5 s of higher latency; events are buffered in producer-side queue, no events lost. Poison message in the stream — the persist-worker's DLQ pattern (separate topic presence.transitions.dlq) captures the bad event so the consumer doesn't get stuck.
- Persist WorkerGo consumer pool, one per Kafka partition
Single-purpose consumer pool that subscribes to presence.transitions (op=subscribe, one consumer per partition for predictable parallelism), QUANTIZES the timestamp to the nearest 60 s bucket (the 'approximate by design' lever — never quote last_seen more precisely than ±60 s), applies the privacy_mask filter (drop the persist write if the user is in 'no last_seen visible' mode but emit the transition for analytics with the user_id stripped), writes to lastseen-db (op=write), and simultaneously writes a minutely aggregate row to analytics-db (op=write, batched at 10 s). Maintains per-partition checkpoint offsets in Kafka's __consumer_offsets.
Why it exists. Decouples the durable last_seen write from the hot transition path. Without this worker, presence-svc would have to write last_seen synchronously, paying a 20 ms ScyllaDB write on every transition — at 500 K transitions/s peak that's 500 K write IOPS with ~10 K writes in flight (Little's law at 20 ms), every one blocking a presence-svc actor, so the hot path inherits every ScyllaDB latency spike. The worker consumes the same stream off Kafka instead — batched, retryable, and off the latency-critical path. The worker is also where the quantization happens: a single place to enforce the ±60 s rounding, so every downstream (lastseen-db, analytics-db, abuse signals) sees the same bucketed values.
When it fails. Worker pool falls behind (consumer lag > 60 s) → lastseen-db is stale; users see 'last seen 12 min ago' even when the user was online 30 s ago. Mitigation: per-partition lag SLO at 30 s, auto-scaling on consumer lag, dedicated DLQ for slow partitions. Worker crash mid-batch → on restart, replays from the last committed offset; idempotent writes ensure no double-count. Schema-version skew on deploy — old worker version sees a new event field it doesn't understand — handled by forward-compatible Protobuf schemas (always-add-fields-with-defaults, never remove) and a strict schema-registry contract.
- Last-Seen StoreScyllaDB (or Cassandra), regional, eventually replicated
Stores one row per user: (user_id PK, last_seen_bucket_ts, last_status, privacy_mask, region) where last_seen_bucket_ts is the timestamp rounded to the nearest 60 s. Updated by persist-worker; read by presence-svc on cache-cold (presence-cache miss) and by the read-fallback degraded-mode path. The 'last seen' the user sees in the UI ('Last seen 12 min ago') comes from this row when presence-cache shows offline. NEVER read on the hot heartbeat path — only the cold-fallback and the degraded-mode read.
Why it exists. Presence-cache is ephemeral (no AOF/RDB) — when a user goes offline and their TTL expires, presence-cache forgets them entirely. But the product spec says 'last seen 12 min ago' must survive the user's disconnect. So we need a durable store for the most recent timestamp per user. Why ScyllaDB and not Postgres? At 1B users × ~10 transitions/day = 10B writes/day = ~115 K writes/s avg / 500 K peak — Postgres at this write rate would need 10+ shards and aggressive partitioning; ScyllaDB at this rate is 4 nodes per cluster. Why ScyllaDB and not DynamoDB? Cost: at 500 K writes/s on DynamoDB the bill is ~$20K/month per region; on ScyllaDB on self-hosted infra it's ~$3K/month per region (per Discord's public migration cost analysis for trillions of messages, scaled to this load).
When it fails. ScyllaDB shard down → 1/36 of users can't have their last_seen read OR written; the read path falls through to a default last_seen='a long time ago' (fail closed); the write path queues in persist-worker (Kafka offsets don't advance) and catches up when the shard recovers. Replication-lag spike (cross-AZ link flap) means a recently-written last_seen value isn't visible to readers for ~30 s; bounded staleness is acceptable here because the product spec only quotes ±60 s anyway. Tombstone-wall problem (the Discord-published pre-Scylla failure mode) is avoided here by NOT clustering on (user_id, ts) — one row per user, OVERWRITE semantics, no tombstones at all in the steady state.
- Analytics StoreClickHouse with ReplacingMergeTree, regional
Stores per-minute rolled-up presence aggregates: (minute_ts, region, status_bucket, count) and a presence_minutely materialized view that feeds the on-call dashboards (concurrent online users, online by region, daily/monthly active users, presence freshness SLO percentiles). The persist-worker writes one batch per partition per minute; ClickHouse merges with ReplacingMergeTree so late-arriving events for a minute (a 30 s Kafka lag, say) over-write the prior aggregate idempotently. Reads are exclusively from product dashboards and from the SLO burn-rate alert pipeline — NEVER on the hot user-facing path.
Why it exists. Presence is a product metric (DAU, concurrent users, regional online rates) and an SLO metric (presence-freshness p99, broadcast-amplification ratio). Querying lastseen-db for these is wrong: ScyllaDB is optimized for point-lookups not for GROUP BY minute over 1B rows. ClickHouse is built for exactly this access pattern — columnar storage, vectorized aggregations, sub-second SELECT count() WHERE minute_ts BETWEEN ... GROUP BY region. Without analytics-db, the product team would write ad-hoc ScyllaDB scans that take down the operational store; with it, analytics is isolated on its own infra with its own SLO.
When it fails. ClickHouse merge-storm under high ingest QPS can cause query latency to climb 10× — mitigation is to batch ingest at 10 s windows (already done in persist-worker) and to size the merge thread pool generously. Replication lag between ClickHouse replicas can cause dashboards to flicker as different replicas serve slightly different aggregates — solved by sticky read-from-leader for dashboards in incident mode. Cardinality storm — a bad label like user_id accidentally added as a dimension — would push series count past memory; mitigation is strict schema review and per-dashboard cardinality budget alarms.
- Metrics (TSDB)VictoriaMetrics (cluster mode) with Prometheus scrape
Scrapes every component's /metrics endpoint at 15 s intervals: gateway connection count, presence-svc shard-mailbox length, presence-cache hit rate, fanout-relay drop-ratio, persist-worker consumer-lag, lastseen-db p99 latency. Stores at full 15 s resolution for 14 days, then downsamples to 5 min resolution for 90 days, then to 1 h for 1 year. Powers the four on-call dashboards (presence freshness, gateway health, fanout health, persist health) and the burn-rate SLO alerts (2%-budget over 1 h fast burn; 5%-budget over 6 h slow burn — Google SRE Workbook chapter 5).
Why it exists. Tracing answers 'why is THIS request slow?'; metrics answer 'is the SYSTEM healthy?' Both are required. The TSDB is the meta-monitoring layer that detects the failure modes the canonical's failure-modes table enumerates — broadcast amplification storm (fanout_frames_per_s × 50 baseline), heartbeat backpressure (presence_offline_transitions/s × 10 baseline), holiday-floor reconnect storm (identify/resume ratio > 3), TTL drift (negative TTL observed > 0). Without this tier, on-call is blind during exactly the incidents where they need visibility.
When it fails. Cardinality explosion (someone adds user_id as a metric label) → head-block memory pressure → ingest cliffs → alerting silent → 'silent outage' where the data plane is degraded but no page fires. Mitigation: per-metric cardinality budget enforced at scrape time + an independent watchdog. Scrape-storm thundering herd (every 15 s every component is scraped simultaneously) — VictoriaMetrics handles this natively with consistent-hash request routing. Long-term storage cost: 1 year of 15 s data is too much; the downsampling tiers (15 s → 5 min → 1 h) is non-negotiable for cost containment at this scale.
- TracingOpenTelemetry Collector → Tempo (or Honeycomb)
Every component emits OTel spans on its hot paths: ws-gw spans the WS frame parse + routing, presence-svc spans the heartbeat process + presence-cache RPC, fanout-relay spans the subscription lookup + dispatch, persist-worker spans the Kafka consume + DB write. The OTel Collector tail-samples at 1% in steady-state, escalating to 100% sampling for traces flagged as 'slow' (p99-band exceeded) or 'error' — the Honeycomb/Cribl pattern. Traces stored in Tempo (or Honeycomb) for 7 days, queryable by trace_id from incident reviews.
Why it exists. When a user reports 'my green dot is stuck on, fix it', the only way to find the root cause is to follow one specific user's heartbeat through every component. Without distributed tracing, on-call has to grep 12 components' logs by user_id and reconstruct the timeline manually — an hour of MTTR per incident. With tracing, it's a 30-second trace_id query. The 'tail-sampled at 1% + 100% on slow/error' policy is the cost trade: full sampling at this scale is petabytes/day; tail-sampling captures the interesting traces and discards the boring ones.
When it fails. Tracing backend down → collector queues fill, then drops; on-call is debugging blind during the very incident that needs the most visibility. Mitigation: 10 s collector buffer + independent watchdog + on-host fallback. Trace cardinality explosion (every span tagged with conn_id as a metric label by mistake) → same as TSDB's cardinality storm. Span-format incompatibility between OTel versions during a rolling deploy → resolved by schema-registry-enforced version negotiation in the OTel collector.
Stage by stage
The same 10 stages the workspace walks, answered.
01Clarifications
What would you ask before drawing a single box?
Typical clarifications to surface:
- Privacy levels? Three-tier minimum:
everyone/contacts/nobodyfor presence, separate setting forlast_seen, plus per-userinvisiblemode. This shapes the entire fanout design. - Cross-region semantics? Is "online in EU, viewed from US" supposed to be precise, or is "active recently" bucket acceptable? (The honest answer at scale is the latter.)
- What counts as "online"? App in foreground? Backgrounded with notifications enabled? Just connected to WS? Each definition has a different capacity model.
- last_seen granularity? ±1 s is hostile (PII leak + hot writes), ±60 s is industry standard, ±5 min for "active recently" buckets.
- Heartbeat interval? 60 s is the production sweet spot — 30 s doubles wire cost without meaningfully improving freshness.
- Latency targets? Presence change → visible to watchers p99 < 5 s. Subscribe initial snapshot p99 < 200 ms.
Assumptions to state:
- 1B DAU, 500M concurrent at peak (Slack/Teams/Discord-tier).
- ~10 online/offline transitions per user per day.
- Viewport-bound subscriptions cap at ~30 active contacts per viewer.
- ±60 s
last_seenprecision is acceptable and is the privacy lever.
02Functional reqs
What must this system actually do?
- Maintain a long-lived WebSocket per active user session.
- Receive heartbeats every 60 s; refresh TTL on the presence record.
- Detect online → offline transition via explicit message OR TTL expiry.
- Detect privacy changes (appear offline) and propagate within 5 s.
- For each transition: persist the durable trail; fan out to viewport-active watchers; update the quantized
last_seenstore. - Subscribe gesture (viewport scroll): viewer asks for current state of ≤30 contacts; system verifies viewer-can-see-target privacy, returns snapshot, registers the watcher in the reverse index.
- Read gesture (initial app load): viewer asks for the state of their visible contacts; system serves from presence-cache with degraded fallback to last-seen-db.
03Non-functional
What must it promise about speed, uptime and correctness?
- Availability: 99.95% on the read path (degraded mode = serve stale or "unknown" status, NEVER fake "online"). 99.9% on the heartbeat/write path is acceptable — a missed heartbeat is recovered by the next one within 60 s.
- Latency:
- Presence change → visible to watchers (steady-state): p99 < 5 s, hard cap 15 s.
- Privacy-flip propagation: p99 < 1 s (the harassment-escape SLO — tighter than steady-state and routed via the high-priority lane at fanout-relay with a 100 ms coalesce, NOT the 1 s steady-state coalesce). Hard cap 5 s.
last_seenaccuracy: ±60 s.- Subscribe initial snapshot: p99 < 200 ms.
- Heartbeat round-trip: p99 < 100 ms.
- Durability:
last_seensurvives WS disconnect via the events-bus → persist-worker → lastseen-db pipeline (RPO 60 s, RTO 0 s for reads thanks to leaderless replication). Presence-cache is explicitly NOT durable. - Consistency: Eventually consistent across regions for presence reads; read-your-write within a single region for the heartbeat-then-read path in steady state (presence-cache RYW guarantees this). During presence-cache failover (≤15 s RTO window), the read path falls through to lastseen-db's eventually-consistent ±60 s bucket and the response is flagged
WS_DEGRADED— the SDK renders 'unknown' rather than stale-but-claimed-fresh state. Privacy-flip is linearizable globally (the harassment-escape SLO). - Scalability: Horizontal at every tier. Sharded by
user_ideverywhere except fanout-relay (sharded bywatched_uid, the N² solver split). - Security: Three-tier privacy preferences, per-user invisible mode, no PII leak (last_seen quantized), no probe API (privacy violations return "unknown" not "403"), rate-limited subscribe to prevent contact-list enumeration.
04Capacity estimation
How much load and data does this have to hold?
Walk-through from product spec to architectural choice. The numbers come from the Stage-1 capacity research (Discord 5M concurrent + WhatsApp 2M-per-box published figures).
Heartbeat economics. 1B DAU × 50% concurrent at peak = 500M concurrent online sessions. At 60 s heartbeat interval, that's 8.3M PING/s steady-state, and at the 4× diurnal/reconnect-burst peak multiplier, ~33M PING/s peak going INTO the gateway tier (the number we actually size for). Wire cost at 40 B/frame = 1.3 GB/s = 10.5 Gbps at peak — absorbed by the BEAM-tuned ws-gw tier.
Why 60 s and not 30 s. Liveness threshold = heartbeat × 2 = 120 s — one missed beat is tolerated, two means gone; no jitter margin needed because the actor's sweeper compares server-side receive times. A 30 s heartbeat would double wire cost (~21 Gbps at peak) without meaningfully reducing detection latency — the dominant detection mode is graceful disconnect, not the timeout sweeper.
Transitions QPS. 1B DAU × 10 transitions/day = 10B/day = 115K/s avg, ~500K/s peak (4× peak:avg multiplier, standard diurnal).
Fanout math — eager broadcast (the design we REJECT). 500K transitions/s × 500 contacts each = 250M frames/s. At 30 B/frame = 7.5 GB/s = 60 Gbps. Physically infeasible on commodity infra.
Fanout math — viewport-bounded (the design we ACCEPT). 500K transitions/s × 30 active watchers = 15M frames/s. At 30 B/frame = 450 MB/s = 3.6 Gbps. Comfortably absorbable on the fanout-relay tier. This is the load-bearing capacity invariant — without viewport bounding, the entire architecture collapses.
Presence-cache RAM (B1 corrected). 500 M concurrent × 100 B/record = 50 GB working set per region. With Redis hash overhead + value padding (~1.5×) + jemalloc fragmentation (~1.3×) = ~100 GB primary RAM per region. RF=3 (one primary + two followers, defended below) → ~300 GB total cluster RAM per region. Across 24 Redis Cluster primaries = ~4.2 GB/shard primary, sized at 16 GB instances for 2× headroom + warm-restart slack. (NOT 19 GB/shard — that number conflated per-box capacity with per-shard working set.) Two followers per shard (not one) because the AZ-failover RTO=15 s SLO requires a hot follower in a third AZ to survive one-AZ-down-during-deploy. With only one follower, a deploy that bounces the primary's AZ during a separate AZ outage breaks failover entirely.
Heartbeat → Redis write load (B5 corrected). A short per-user coalesce window can not cut Redis load here: heartbeats are 60 s apart (see "Why 60 s"), so any window shorter than the interval holds ≤1 heartbeat — there is nothing to collapse, and a "10 s window → 10× fewer SETEX" claim is arithmetically impossible. The real lever is the actor tier (see presence-svc): the per-beat heartbeat firehose is absorbed in the actor's in-memory (ETS) liveness state; Redis is written on state transitions (online/away/offline) plus a slow backstop-TTL keepalive far below the beat rate. Offline-by-timeout is raised by the actor's in-memory sweeper, not by a short Redis TTL, so a still-online user does not need a SETEX per beat. Redis write load is therefore dominated by the ~500 K/s transition rate (+ a lazy keepalive) — ~600 K–1.3 M writes/s, not the 8.3 M/s a per-heartbeat SETEX would cost. Across 24 presence-cache shards = well under ~55 K writes/s/shard, comfortably below Redis Cluster's ~100 K ops/s/shard ceiling, with headroom for a 4× post-partition transition storm. (Keeping 500 M keys warm under a short 120 s TTL would instead cost ~8 M refresh-SETEX/s — impossible on 24 shards — which is exactly why liveness lives in the actor's memory and the Redis key carries a long backstop TTL refreshed lazily, not a 120 s per-heartbeat TTL.)
Identity-svc JWT refresh QPS. 500 M concurrent sessions × (1 / 300 s refresh interval) = 1.66 M refresh-token-exchange req/s sustained against identity-svc. ws-gw caches the JWKS for 1 h (so JWT signature validation on every WS Upgrade is local, no identity-svc round-trip), but the refresh exchange (long-lived refresh token → new short-lived JWT) is a real call. Identity-svc must be sized for 1.66 M req/s — out-of-scope for this canonical (separate identity-svc design owns it), but stated here so it does not get hand-waved.
Subscription-store RAM. 500M concurrent × 30 active watchers each, watched_uid keyspace = effectively 500M rows × ~1 KB each (the watcher set hash) = ~500 GB / region. Sharded across 32 Redis Cluster shards = ~16 GB/shard, on 64 GB instances. Celebrity p99 caveat: a top-1000 broadcaster (50 M watchers) holds a 1.5 GB watcher hash; salted across 4 sub-shards is 375 MB on the heaviest sub-shard. Fits in 64 GB easily — but the HGETALL latency on a 12.5 M-entry hash is seconds, not ms. This is why fanout-relay paginates with HSCAN for broadcasters >10 K watchers (see fanout-relay.keyChoices).
Kafka throughput. Transitions + audit + privacy-changes = ~600K events/s peak. At 200 B/event compressed = 120 MB/s = 960 Mbps. Single 18-broker cluster per region absorbs this comfortably. Tiered storage (S3 cold) caps disk cost.
ScyllaDB last-seen writes. Quantization dedupes only same-user-same-minute flaps, so persist-worker's write rate tracks the transition rate: 1B DAU × 10 transitions/day / 86400 ≈ ~115K writes/s sustained, ~500K peak (4× multiplier). ScyllaDB cluster at 36 nodes × ~50K writes/s capacity = 1.8M writes/s — ~3.5× headroom at peak; the cluster also serves reads on cache-miss (~5K reads/s) and the leaderless RF=3 layout is the floor for production durability.
WS gateway box count. We commit to BEAM/Elixir at the gateway tier — Discord's published 5M sockets/box is the production-defended ceiling, and the per-process memory footprint AND preemptive scheduler are what make this work. 500M concurrent / 5M per box = 100 boxes per region at steady state, and we size at 200 boxes (2× headroom) so a single-AZ loss leaves the remaining 2-of-3 AZs at <75% saturation. At the Go/Rust 1M-per-box alternative the canonical would need 1000 boxes — 5× the operational footprint for no production benefit. The 5M number is non-negotiable for this canonical's economics.
The scale tiers where each architectural choice flips:
| From → To | Threshold | Citation |
|---|---|---|
| Single Redis → Redis Cluster | ~100K presence writes/s | Redis Cluster docs |
| Eager broadcast → viewport-bound subscription | ~100K concurrent users / channel | Discord 2017 Elixir blog |
| Region-local → multi-region replication | ~10M global users OR first non-US POP | WhatsApp Erlang Factory |
| In-memory only → persisted last_seen | Day 1 if spec says "last seen 12 min ago" survives offline | Product spec |
| Naive heartbeat → coalesced/debounced | ~10M concurrent (above this, presence-svc shard mailbox grows) | This canonical's actor-tier choice |
05API design
What does the outside world call, and what comes back?
GET wss://presence.example.com/v1/connect
Sec-WebSocket-Protocol: presence.v2
Authorization: Bearer <jwt>
# Server confirms upgrade and sends an initial frame:
{ "type": "WELCOME", "session_id": "...", "resume_token": "...", "heartbeat_interval_seconds": 60 }
// Client → Server (over WS)
// First frame after connect: authenticate the session.
{ "type": "IDENTIFY", "auth": "<jwt>", "client": { "app_version": "...", "device": "mobile" } }
// Warm reconnect after a drop: resume the prior session instead of a cold IDENTIFY.
{ "type": "RESUME", "resume_token": "<from WELCOME>", "last_seq": 12344 }
{ "type": "HEARTBEAT", "seq": 12345 }
{ "type": "SUBSCRIBE", "uids": ["u_alice", "u_bob", "u_carol"], "viewport_active": true }
{ "type": "UNSUBSCRIBE", "uids": ["u_alice"] }
{ "type": "STATUS_SET", "status": "away" } // user-initiated
{ "type": "STATUS_SET", "status": "invisible" } // privacy override
// Server → Client (push)
{ "type": "RESUMED", "session_id": "...", "missed_from_seq": 12344 } // RESUME accepted; deltas replayed
{ "type": "INVALID_SESSION" } // resume_token stale → client must re-IDENTIFY
{ "type": "PRESENCE_UPDATE", "uid": "u_bob", "status": "online", "last_seen_bucket": 1748120460 }
{ "type": "PRESENCE_SNAPSHOT", "states": [{ "uid": "u_alice", "status": "offline", "last_seen_bucket": 1748120400 }, ...] }
# Fallback REST API for clients that can't keep WS open (push-notification cold-start, etc.)
GET /v1/presence/batch?uids=u_alice,u_bob,u_carol
Authorization: Bearer <jwt>
200 OK
{
"results": {
"u_alice": { "status": "online", "last_seen_bucket": 1748120460 },
"u_bob": { "status": "offline", "last_seen_bucket": 1748120400 },
"u_carol": { "status": "unknown" } // privacy-mask hides last_seen, OR cache miss
}
}
Privacy-enforcement contract:
"status": "unknown"is returned for (a) viewer not allowed (privacy mask), (b) presence-cache miss + lastseen-db miss, (c) target in invisible mode. Same response — no probe-by-error-code.last_seen_bucketis always quantized to 60 s; clients MUST NOT display sub-minute precision.
06Data model
What gets stored, and what is it looked up by?
Presence-cache (Redis hash, ephemeral):
KEY: presence:{uid} TTL: 900s (leak backstop — liveness is the actor's 120s in-memory threshold)
HASH:
status : "online" | "away" | "invisible"
last_seen_ts : int64 (raw, will be bucketed by persist-worker)
region : "us-east-1" | "eu-west-1" | ...
device_hint : "mobile" | "web" | "desktop"
privacy_mask : "everyone" | "contacts" | "mutual" | "nobody"
Subscription-store (Redis hash, AOF-persisted):
KEY: watchers:{watched_uid} TTL: 1800s (whole-set backstop — cleanup is field-level HDEL on session close + stale-field sweeper)
HASH (one field per watching connection):
{conn_id} → {ws_gw_shard, viewport_active, subscribed_at}
Last-seen-db (ScyllaDB, durable, one row per user):
CREATE TABLE last_seen (
user_id bigint PRIMARY KEY,
last_seen_bucket bigint, -- unix_ts // 60 * 60 (quantized to 60s)
last_status text,
privacy_mask text,
region text,
updated_at timestamp
) WITH compaction = {'class': 'LeveledCompactionStrategy'};
Identity graph (Neo4j):
(:User {user_id, presence_visibility, last_seen_visibility, invisible_mode})
-[:CONTACT {mutual: bool, blocked: bool, since_ts}]->
(:User)
Events-bus topic (Kafka, presence.transitions):
message Transition {
fixed64 user_id = 1;
enum Status { ONLINE = 1; OFFLINE = 2; AWAY = 3; INVISIBLE = 4; }
Status status = 2;
fixed64 ts_unix_ms = 3;
string region = 4;
fixed32 conn_id = 5;
enum Reason { GRACEFUL = 1; TTL_EXPIRY = 2; CLIENT_SET = 3; PRIVACY_FLIP = 4; }
Reason reason = 6;
}
Database choice — recommended:
For last_seen: ScyllaDB for production new-builds at this scale. Same wire-protocol as Cassandra, ~4× lower per-node cost (per the Discord trillions-of-messages migration), and the OVERWRITE-only access pattern means no tombstone-wall risk (the failure mode that drove Discord off Cassandra). Postgres single-leader works to ~10K writes/s; we're at ~500K peak, so we'd need sharded Postgres (Citus / Vitess) — workable but operationally heavier than ScyllaDB for this access pattern.
For presence-cache: Redis Cluster with allkeys-lru, AOF off. The ephemeral nature is the load-bearing choice — losing the cache means clients re-establish state within 60 s.
For subscription-store: Redis Cluster with AOF everysec, sharded by watched_uid. AOF (not RDB) because losing the subscription set forces every client to re-subscribe their viewport — a 100M-subscriber reconnect storm. The 1 s RPO is acceptable.
For graph-db: Neo4j Causal Cluster OR sharded Dgraph. The 1-hop "is viewer mutual contact of target?" is exactly the access pattern graph-dbs are built for; storing this in Postgres requires a join on every check.
Why not just one Postgres for everything? At 1B users × ~10 transitions/day, you'd need 50+ shards, cross-shard subscription queries would be a mess, and ephemeral presence state in Postgres pays unnecessary durability overhead on every heartbeat.
07High-level design
Which components handle a request, and in what order?
Read/subscribe path. Client → edge-lb (anycast WS) → ws-gw → presence-svc → presence-cache (read) → response back through the same chain. On cache miss, presence-svc falls through to lastseen-db (degraded mode, surfaced as last_seen_bucket). On subscribe, presence-svc additionally writes the watcher entry to subscription-store and consults graph-db for privacy.
Heartbeat path. Client → edge-lb → ws-gw → (batched 10 ms windows) presence-svc, which refreshes the actor's in-memory liveness timer. Returns to client through ws-gw with sub-100ms RTT. Heartbeats DON'T write presence-cache — Redis is written on transitions (SETEX with the ~15 min backstop TTL).
Transition fanout path. Triggered by either (a) explicit STATUS_SET, (b) TTL expiry sweeper detecting an offline transition, (c) privacy-flip. presence-svc publishes to events-bus FIRST (durable trail, ack=all — the event-sourced source of truth), then projects to presence-cache AND fanout-relay (in-process publish). fanout-relay reads subscription-store, filters by viewport_active=1, dispatches per-session push frames back through ws-gw → recipient clients.
Async persistence path. events-bus → persist-worker (Kafka consumer) → lastseen-db (UPSERT quantized) + analytics-db (batch insert minutely).
Why this shape and not [insert simpler shape]:
- Why a separate fanout-relay tier and not have presence-svc do it directly? Because at the celebrity tier (50M watchers), the dispatch loop would saturate a single presence-svc actor's mailbox. Splitting fanout onto a separately-sharded tier (sharded by watched_uid) gives each popular user their own dispatcher with a bounded mailbox.
- Why a separate subscription-store and not derive watchers from graph-db every time? Because graph-db's 1.5M traversals/s on subscribe is fine, but graph-db's
would-be 500K-traversals-on-every-transition × 30-watchers = 15M traversals/sis well past Neo4j's published throughput. The subscription-store is the materialized denormalization for the hot transition path; graph-db is the slow read on the rare subscribe path.
- Why persist via the Kafka spine and not direct-write to lastseen-db? Because synchronous writes to lastseen-db on every transition would double presence-svc's write latency and quadruple ScyllaDB's required size. The event-sourced spine is the standard decoupling — the hot path publishes to the durable log and projects to the ephemeral cache; the deferred consumer persists durably.
08Deep dives
Which part breaks first, and what do you do about it?
1. The N² watch problem.
The architecture must guarantee that fanout_frames_per_second does NOT scale as O(users × avg_friend_count). Three load-bearing levers:
- Viewport at the client. SDK reports at most ~30 contacts to subscribe to, never the full contact list. Citation: Discord 2017 — "How Discord Scaled Elixir".
- Subscription-store as the inverse index. Per-watched-uid hash of active watchers. Fanout reads this, not the full friend-list.
- Passive-session filter at fanout-relay. Even within the subscribed set, ~90% of watchers' viewport doesn't currently include this user at this moment — those frames are dropped at the relay, not at the client.
Drop any one of these and the system melts. The math: 500K transitions/s × 30 active watchers × 0.1 viewport-active filter = 1.5M actually-delivered frames/s — versus the 250M frames/s a naive eager-broadcast design would generate.
2. Approximate by design — the ±60s last_seen lever.
Quoting last_seen precisely (±1 s) is a design defect, not a feature. Two reasons:
- Privacy. A precise last_seen ("last seen 14 s ago") is a stalking signal. WhatsApp's 2012 privacy retrofit was driven by this exact harassment vector.
- Capacity. Persisting
last_seenat sub-minute precision forces a durable write on every transition — at 500K transitions/s, that's the difference between a 4-node ScyllaDB cluster and a 36-node one.
Bucketing at 60 s (last_seen_bucket = floor(unix_ts / 60) * 60) gives natural write deduplication (multiple transitions in the same minute become one row update) AND privacy-friendly UX strings ("last seen 12 min ago" reads more humanely than "last seen 14:32:17 UTC"). The persist-worker is where the quantization happens — one place, enforced.
3. The TTL-vs-heartbeat invariant.
liveness_threshold = heartbeat_interval × 2 = 60 × 2 = 120s, enforced by the presence-svc actor's in-memory sweeper.
Why × 2: one missed heartbeat is normal (cellular handoff, GC pause); two missed heartbeats means the client is genuinely gone. No added safety margin is needed because the sweeper compares server-side receive times (the last_byte_recv watermark) — client clock skew never enters the comparison. The presence-cache key carries a separate long backstop TTL (~15 min) that only reaps records leaked by an actor crash; it is NOT the liveness mechanism.
The other half of this invariant: presence_cache.TTL < subscription_store.TTL. Otherwise the push-vs-pull race lets the relay fan out a transition AFTER presence-cache has forgotten the user — observers receive a green-dot frame for a user whose authoritative state has expired. The canonical sets presence-cache backstop TTL = 900 s, subscription-store TTL = 1800 s (30 min) — wide safety margin. (The "friend list" — the durable social graph — lives in graph-db, NOT in a hot cache; it's read at subscribe time, cached in-process at presence-svc for 5 min, and the 5-min cache is bounded by the 5-min privacy-flip propagation hard cap.)
4. The broadcast amplification storm on region heal.
50M users come back online in the same 60 s after a region partition heals. Naive eager-broadcast: 50M × 500 = 25B presence events in 60 s = 417M events/s. Beyond physical capacity by ~80×.
The canonical's three mitigations:
- Decorrelated jitter at the client. Reconnect window is randomized 30–300 s (AWS Builders' Library algorithm). 50M reconnects spread over 5 min = ~167K/s at peak.
- Viewport-bounded subscription. Each reconnect generates ~30 subscribes, not 500.
- Coalesce window at fanout-relay. 1 s window collapses N transitions of the same user into one frame. During the storm, many users will go online → offline → online as their flaky-network reconnect attempts succeed/fail — the coalesce drops 60–80% of those.
Result: 167K/s × 30 × 0.3 (post-coalesce) = 1.5M fanout frames/s during the storm peak — within 10× of steady-state. Survivable.
5. The privacy-flip race.
When a user toggles 'appear offline' while harassment is in progress, the system must propagate within 5 s globally. This is its own trace:
- SDK → ws-gw → presence-svc (op=write, STATUS_SET invisible)
- presence-svc → graph-db (op=write, durable user property)
- presence-svc → presence-cache (op=invalidate, mark privacy_mask)
- presence-svc → events-bus (op=publish, transition_reason=PRIVACY_FLIP)
- presence-svc → fanout-relay (op=publish, high-priority queue)
- fanout-relay → ws-gw → recipients (synthetic PRESENCE_UPDATE with status=unknown)
The privacy-flip uses a separate high-priority lane at fanout-relay (100 ms coalesce, not the 1 s steady-state coalesce) — the harassment-escape SLO is tighter than the steady-state SLO.
6. The celebrity hot-shard problem.
Cristiano Ronaldo (or your platform's equivalent) has 100M followers. When he opens his app:
- His presence record is on ONE presence-cache shard. Read traffic for
presence:{ronaldo}would spike to ~100M reads/s — instantly hot-shard meltdown. - His watcher set is on ONE subscription-store shard. The HGETALL would return 100M entries — single-shard memory + bandwidth meltdown.
The canonical's mitigations:
- No fanout list in the hot-key.
presence:{ronaldo}is just{status, last_seen}— 100 bytes. The 100M follower fanout is NOT stored against this key. - Subscription-store salted sub-shards. Top-1000 "broadcaster" users get their watcher set split across 4 sub-shards (
watchers:{ronaldo}:salt0..3). HGETALL becomes 4 parallel reads against different shards. - Pre-aggregated active-watcher count at presence-svc. Instead of presence-cache reads on every viewer (100M of them), the presence-svc actor for
ronaldoanswers from in-process state for 99% of reads. - Per-broadcaster write throttling. When ronaldo's followers all open his profile at once (viewport-scroll storm on a celebrity announcement), the subscribe write rate is capped per-broadcaster at presence-svc; excess subscribes get a fast 429 + jittered retry.
7. Recovery-after-region-failure semantics.
Two regions: us-east-1 (primary) and eu-west-1 (secondary). Presence is region-local — there is no global presence read. A user in EU goes online; only EU subscribers see them within 5 s. A US viewer querying the EU user sees status: 'active recently' with the last cross-region replicated last_seen_bucket from lastseen-db (replicated cross-region with ~30 s lag).
This is the honest answer at scale. Pretending to do globally-consistent presence costs an order of magnitude more infrastructure for a benefit users don't perceive — they can't tell that their friend went online in EU 2 seconds ago vs 30 seconds ago.
09Trade-offs
What did this design cost, and what breaks at 10×?
What we ACCEPTED:
| Trade-off | What we accepted | What we rejected | Why |
|---|---|---|---|
| Approximate last_seen | ±60 s quantization | Sub-second precision | Privacy + capacity. Quantization makes both better. |
| Eventually-consistent cross-region | "active recently" bucket cross-region | Globally-consistent presence | 10× infra cost for a UX benefit users can't perceive. |
| Lossy fanout | Drop oldest frames on relay overload | Backpressure all the way to client | Losing "X went online 3s ago" is better than gateway brownout. |
| Event-sourced (not outbox, deliberate) | Presence-svc writes events-bus FIRST (durable source of truth), then projects to presence-cache and fanout-relay. On presence-svc crash mid-write, the events-bus → fanout-relay catch-up consumer (edge e14b, branchKey=outbox-recovery) re-projects from the durable log. No CDC-relay sidecar, no per-shard outbox table. | True outbox (local outbox table + CDC relay) OR 2PC across cache + Kafka | Event-sourced is operationally simpler than outbox-with-CDC: ONE source of truth (events-bus), every other store is a derived projection. Recovery is "replay from the last consumed offset" — the SAME mechanism we already use for persist-worker. No new node, no new sidecar. The doc honestly labels e10.txn.mode as event-sourced and the recovery edge is drawn (e14b). |
| Region-local subscription-store | No cross-region replication of subscriptions | Global subscription mesh | Subscriptions are session-state; users in one region don't watch users in another via the WS session. |
| Server-side coalesce window | 1 s coalesce drops transitions | Pass through every flip | Client can't perceive sub-second flicker; coalesce saves wire. |
What breaks at 10× scale (10B DAU):
- The 64-shard presence-svc tier needs to shard further — at 8M concurrent/shard we're at ~30 K updates/s/shard, fine; at 80M/shard the BEAM scheduler can't keep up. Re-shard to 256.
- The 18-broker-per-region Kafka cluster gets pushed to 100+ brokers, or migrate to Pulsar's tiered storage model natively.
- Subscription-store at 10× the rows might overflow Redis Cluster — migrate the inverse index to an embedded LSM (RocksDB-per-shard) or to a leaderless wide-column.
- The graph-db tier needs Dgraph or a custom sharded graph store; Neo4j Causal Cluster doesn't scale past ~50 shards.
Primary sources
- How Discord Scaled Elixir to 5,000,000 Concurrent Users (discord.com/blog, 2017)
- How Discord Handles Push Request Bursts with Elixir's GenStage (discord.com/blog)
- How Discord Stores Trillions of Messages (discord.com/blog, 2023)
- Slack's Disasterpiece Theater — Approachable Chaos Engineering (slack.engineering)
- Slack's Outage on January 4th 2021 (slack.engineering)
- Slack's Incident on 2-22-22 (slack.engineering)
- Rick Reed — WhatsApp Scaling, Erlang Factory 2014 (erlang-factory.com)
- WhatsApp — Giving You More Control Over Your Privacy (blog.whatsapp.com, 2012/2018)
- Facebook TAO: A Distributed Data Store for the Social Graph (USENIX ATC 2013)
- Google SRE Workbook — Alerting on SLOs (burn-rate alerts)
- AWS Builders' Library — Timeouts, Retries and Backoff with Jitter
- Cloudflare 2020-07-17 BGP-withdrawal post-mortem (blog.cloudflare.com)
- Datadog 2023-03-08 multi-region connectivity outage retro (datadoghq.com)
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 Online Indicator yourselfMore 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.
- WhatsApp / MessengerHundreds of millions of long-lived sockets, sub-second 1:1 + group delivery, E2E-encrypted, multi-device, multi-region active-active.
- Slack / DiscordChannels and history. Push or pull — and how a hot-channel fanout doesn't melt the gateway.
- Live Comments / Score UpdatesPub/sub at scale. WebSocket vs SSE vs long polling. Approximate by design — mega-rooms drop comments on purpose.
- Collaborative Editor (Google Docs)OT vs CRDT. Causal ordering. Real conflict-freedom. Server-authoritative single-writer doc-actor with WAL-before-ack — per-doc serialization, persisted before broadcast, pinned to a home region.
- Concurrent Hotel Viewers"X users viewing this right now." Hot keys, HLL.
- Live Viewer Count (YouTube/Twitch)Millions of viewers on one entity. Approximate by design.