Live Comments / Score Updates — a worked solution
Pub/sub at scale. WebSocket vs SSE vs long polling. Approximate by design — mega-rooms drop comments on purpose.
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 Live Comments / Score Updates workspaceThe problem
MrBeast is live and 5 million phones are open at once. A goal scores in the World Cup final and 30 million pulse-vibrate. Somebody types "BRO" and a hundred thousand people see it in the same second — except, by design, the other four million nine hundred thousand don't. This is the dominant shape of "live comments + score updates": a write-explosive, read-explosive workload where O(n) publishers create O(n²) potential delivery work, where the right answer to "deliver every comment to every viewer" is no — sampling is the design, not a degradation.
This canonical is about the actual production fabric: where the long-lived socket terminates, where the sampler lives, how a per-room fanout actor survives a hot-room of 5M concurrent without a single shard saturating, how the publish leg never lets a moderation outage corrupt the user's send button, and how cross-region active-active maintains a single writer per room so the replay buffer doesn't fork on failover. It is also explicitly about the things that DON'T work at this scale — topic-per-room Kafka above 50K rooms, naïve EventSource 3-second reconnect, exact XTRIM MAXLEN, "dual-write for safety" during failover, the SSE-over-HTTP/2 HOL-blocking trap.
The two distinguishing design pressures, separate from any other real-time problem in the catalog, are (1) bidirectional: every viewer is both a publisher AND a subscriber on the same socket — not the asymmetric "broadcast-only" of sports scores — and (2) approximate by design: in a mega-room you SHIP a sampler that drops 99.99% of comments on a deliberate, ranked, fair-by-window policy, and the SLO for "did the user see this specific comment" is intentionally not 100%. Score events, conversely, are on a privileged lane that bypasses both the sampler and the moderation pipeline — they're low-volume, authoritative, and never sacrificed.
The reference architecture
What each component is for
- Mobile / Web clientiOS / Android / browser SDK (WS primary, SSE secondary, long-poll legacy)
Opens a single persistent WSS to the nearest edge PoP, subscribes to one room, publishes the user's own comments through the same socket, receives others' comments + privileged score events from that socket, renders last-revision-wins on corrections, persists last_event_id to disk for reconnect resume. Falls back to SSE (with X-Accel-Buffering: no + 15s keepalive) if WS Upgrade is blocked by a corporate proxy; falls back to long-poll (concurrency-capped) if SSE buffers; surfaces 'Connection problems' rather than infinite climb past rung 3. Renders the 'sampled view' UI badge when the server sends config: { sampled: true }.
Why it exists. Modeled as a first-class node because the *client's* state is what determines whether the system survives reconnect storms and failover windows: its decorrelated-jitter backoff, its sticky-cookie session, its last_event_id persistence, its snapshot-fetch fallback policy. A naive EventSource 3s-no-jitter reconnect is exactly the Slack/Discord storm shape; the client SDK is the design's first line of defense against it.
When it fails. Naive reconnect (no jitter, 3s default) turns a 60s ws-gateway rolling restart into a 5-minute self-DOS on the auth tier. App without persistent last_event_id silently drops every comment that happened during a 6s 4G handoff. Unbounded ladder climb (SSE failed → long-poll failed → retry WS infinitely) wedges the client in a connect-loop that drains battery and rotates IPs through carrier NAT. UI without sampler badge: viewers think the chat is broken when in fact it's working as designed at 0.002% keep-fraction.
- GSLB / Geo-DNSAWS Route53 traffic policy or NS1 (managed), 4-provider redundancy
Resolves live.example.com to the IP of the nearest healthy region's anycast prefix. Honors geographic policy (US users → us-east-1, EU → eu-west-1, APAC → ap-south-1) on steady-state; during operator-driven failover, control-plane flips the GSLB routing rule to redirect new resolutions away from the draining region within the DNS TTL window.
Why it exists. Anycast handles steady-state routing in microseconds, but anycast cannot do *graceful* regional drain — once a PoP advertises a prefix, every TCP SYN within its catchment lands there. The control-plane failover playbook depends on the GSLB to redirect *new* connections without affecting existing sockets. Without a GSLB layer, the only way to drain a region is BGP withdrawal — correct but blunt, and triggers the anycast-hairpin failure mode at the failover PoP.
When it fails. GSLB outage (whole DNS provider goes down) silently freezes regional steering: cached-resolution clients continue to land on their last-resolved region; if that's the failed one, new connects fail. Mitigated by 4-provider redundancy. DNS TTL too long (e.g. 5min) makes failover effect 5min late and breaks the 60s RTO claim. DNS TTL too short (5s) puts GSLB on every WS handshake and self-DOSes the DNS provider's edge.
- Anycast L4 LBBGP anycast + Cloudflare Magic Transit / AWS Global Accelerator
Terminates the TCP connection at the geographically nearest PoP via BGP anycast, runs an L4 admission-control token bucket on Upgrade: websocket requests (sized 1.2× steady-state — deliberately tight), then forwards the raw TCP stream to a co-located ws-gateway process. Issues Retry-After with decorrelated jitter when the local PoP is over its admission ceiling. Receives drain commands from control-plane to stop accepting new Upgrades without affecting existing sockets.
Why it exists. Pinning the TLS handshake at the PoP cuts ~80ms of cross-continent RTT off cold-connect latency, and the L4 admission cap is the *only* defense against a reconnect storm landing on a flapped-failover PoP at 5× normal rate. Putting the admission gate at L4 (before WS upgrade) means we shed at 200µs/decision instead of 20ms of TLS + Lua decision cost.
When it fails. BGP withdrawal in one PoP (Cloudflare Jul 17 2020 / 2022-06-21 shape) re-routes 1M sockets through one neighbor PoP that immediately saturates — anycast hairpin. Mitigation is pre-announce backup /24s + per-PoP admission caps that 503 cleanly. If admission cap is too high, we cascade to ws-gateway OOM. If too low, we self-DOS during legitimate peaks — cap MUST be sized to 90-percentile event concurrency, not steady-state. Cert expiry on this tier silently breaks every new connect; monitored by independent prober.
- WS / SSE / Long-poll GatewayCentrifugo (Go) on bare metal — primary; Cloudflare Durable Objects + WS Hibernation as roadmap alternative
Terminates TLS + WS Upgrade (or SSE event-stream / long-poll), validates the bearer JWT against a local JWKS cache + signed pin (24h offline tolerance), reads catalog for room entitlement + blackout + slow-mode, registers the subscription on the correct room-shard. On the publish side: receives {op:publish, client_msg_id, body} frames from the client, forwards via gRPC to comment-ingest, returns the publish_ack (carrying client_msg_id + a provisional spine (partition, offset); the authoritative room-global event_id arrives later on the fanout copy) on the same socket. On the fanout side: receives push from sampler/room-shard, flushes the JSON frame into every subscribed client's socket. Maintains a per-room in-pod ring-buffer (last 30s) so reconnects within the window resolve locally without hitting Redis.
Why it exists. The single biggest decoupling in the design: per-room fanout has to live wherever the *publishing* topic lives (one actor per room), but TLS termination has to live wherever the *subscriber* lives (200+ PoPs globally). Folding them together would force each room's fanout shard to terminate sockets from every continent — the topology that took Discord's pre-Manifold cluster to 900ms-2.1s hot-guild fanout. Slack's Flannel is the analog: an edge cache + WS terminator that absorbs reconnect storms and bootstrap fan-in without hammering the central message bus.
When it fails. Hits fd ceiling silently (default 1024) and accept_eagain_total climbs while new connects spin — must raise fs.nr_open to 4M+. TLS-resumption stampede after key rotation forces full handshakes; CPU melts. Slow-reader backpressure without the 32KB cap OOM-kills the pod and takes 100K sockets with it (then triggers reconnect storm — see anycast-lb). Coalescing window collapsing a revision is the SEV-1 bug class; defense is (roomId,eventId,revision) keying. WS idle disconnect under proxy fleet rollout (some MNO upgrades 30s→15s the proxy idle timeout overnight) is detected as bimodal disconnect histogram by geographic ASN; remediation is to drop the heartbeat to 10s for that ASN.
- Auth / JWT IssuerGo service + JWKS endpoint + KMS for signing keys + signed-JWKS-pin root
Issues 15-min RS256 JWTs when the mobile app authenticates (OAuth, Apple/Google sign-in, anonymous device id). Publishes the public JWKS endpoint that every ws-gateway caches with 15min TTL. Also serves a *signed JWKS pin* (the JWKS document signed by a long-lived root key) so ws-gateway can verify offline for up to 24h if auth is unreachable — closes the AWS Builders' Library 'avoid fallback' anti-pattern.
Why it exists. Existing as a *separate, off-the-WS-path* tier is load-bearing. If ws-gateway called back to auth on every WS open, mega-event kickoff (5M users simultaneously open the app) would self-DOS the auth tier, cascading to every other product sharing it. Short-lived JWT + local JWKS cache means auth only sees traffic when (a) a user first signs in or (b) JWKS keys rotate.
When it fails. JWKS rotation collides with edge cache invalidation (15-min worst-case staleness): some edges see signatures fail until they re-fetch. Auth outage during JWKS rotation extends that window — the signed JWKS pin keeps non-rotated tokens valid for 24h. If ws-gateway is misconfigured to fall back to remote verification on JWKS miss with exp-backoff and no circuit-breaker, an auth outage during a mega-event cascades — the fixed-retries=1, fast-fail policy on ws-gateway is the design's defense.
- Room catalog (Postgres HA)Postgres 16, Patroni-managed HA (1 primary + 2 sync followers)
Holds room metadata (room_id, kind, home_region, state, slow_mode_seconds, sampling_policy, blackout_regions, moderation_policy, owner_id), banned-user records (room_id, user_id, until), and pinned-message records. Queried by ws-gateway on every WS open (entitlement + slow-mode), by comment-ingest on every publish (slow-mode + ban check via 5s in-process LRU), and by snapshot-api on cold-load (room meta).
Why it exists. The metadata layer the rest of the design references but doesn't INLINE — a single source of truth for 'is this user banned in this room', 'is slow-mode on', and 'what's the moderation policy'. Folding it into ws-gateway or room-shard would either (a) make ban-changes propagate eventually (rather than within 5s), or (b) put a Postgres dependency on every WS frame — both wrong. Separate node, with aggressive caching everywhere it's called from, is the right shape.
When it fails. Patroni primary failover takes 30s — during the gap, writes (slow-mode toggle, ban add) BLOCK and ws-gateway returns 503 to clients that need a fresh catalog read (LRU miss). Replica-read for entitlement is acceptable (eventual consistency: a banned user might still get one more comment through if their ban arrived <1s ago — acceptable trade). PgBouncer transaction-mode prepared-statement footgun silently corrupts Go-driver queries unless protocol-handler mode is set — caught on deploy preflight. LRU stampede on slow-mode toggle for a viral room (every ws-gateway pod invalidates simultaneously) → Postgres sees a brief read-spike; mitigated by spine catalog.changes consumer-driven invalidation rather than TTL-driven.
- Comment IngestGo service; stateless; idempotent on (room_id, client_msg_id, user_id)
Receives publish requests from ws-gateway over gRPC. Enforces per-user, per-IP, per-room rate-limits (token bucket in Redis). Dedupes on (room_id, client_msg_id, user_id) — same request from a retry returns the same ack and never publishes a second event. Calls moderation synchronously for a classifier verdict (<80ms p99 budget; circuit-breaker opens on 50% errors / 30s). On block, returns moderation_failed to ws-gateway. On allow or soft-flag, publishes the event to spine with acks=all, enable.idempotence=true and partition = hash(room_id) % 1024, then writes the moderation decision to spine moderation-ledger topic (retained on the spine's tiered storage as the audit ledger). For score events on the privileged endpoint, validates against catalog (room is live) and publishes WITHOUT calling moderation (operator-typed is trusted; licensed feed is contractually accurate).
Why it exists. The publish-side write service. Separating it from ws-gateway is what allows ws-gateway to scale on socket count (150 pods at 100K sockets) and comment-ingest to scale on publish rate (80 pods at ~5K publishes/s/pod). Putting them in the same process couples two very different scaling shapes. Separating from spine producer logic is what allows the synchronous moderation call to live BEFORE the durable publish — no need for a saga / compensating retract path on the durable-side, just refuse the publish.
When it fails. Moderation classifier latency creep (30ms → 1.2s p99) couples publish-path latency 1:1; users retry; queue compounds. Circuit-breaker is the defense, but per-room policy decides what happens when it opens — wrong policy on a high-risk room means toxic content lands in user feeds. Dedup-cache eviction during a deploy can cause duplicate event_id assignment if Redis backup is also evicted — defense is producer-side enable.idempotence. Cross-region proxy adds ~120ms to writer's send latency — visible in the writer's UI as 'sending…' state; falls back to local-region 503 if home region unreachable (with retry, since RTO 60s).
- Moderation classifierPyTorch transformer (toxicity + spam) behind Triton inference server; ONNX export for CPU pods
Receives a comment payload over gRPC from comment-ingest; runs the toxicity model (multi-language transformer fine-tuned on platform-specific labeled data) + spam heuristics (URL detection, repetition, raid signature); returns { decision: allow|soft-flag|block, score, reasons[] } within 80ms p99. Bulkheaded: separate pod fleet from comment-ingest, separate consumer group on the spine moderation-ledger topic for the async retract path used by post-moderate-fast.
Why it exists. The single most-visible quality gate on the platform. Sits BETWEEN comment-ingest and spine so a block-decision never produces a durable record — there's nothing to retract. The bulkhead (separate pods, separate metrics, separate deploy cadence) prevents a moderation regression (slow model update, OOM during a vocabulary refresh) from coupling 1:1 with publish-path latency. The circuit-breaker on the COMMUNICATION between ingest and moderation (not inside moderation itself) is what lets ingest fall through per-room policy when moderation is degraded.
When it fails. Classifier latency creep — 30ms → 1.2s p99 — couples directly to publish-path latency if circuit-breaker is not tuned; this is the AWS Kinesis 2020-11-25 retry-amplification shape if clients retry. Model regression (a new training run misclassifies common words as toxic) blocks legitimate comments; defense is canary deployment per language with synthetic-comment probes + held A/B. OOM during model load (e.g. vocabulary expansion past pod RAM) takes the moderation tier down; defense is pre-deploy memory profiling + per-pod model warmup. Toxic content visible in a pre-moderate-strict room is the SEV-1 case — defense is fail-CLOSED on circuit-open for that policy.
- Score feed (Sportradar / Genius / live-blog editor)Third-party push API (Sportradar UOF) / Genius Sports WSS / internal live-blog editor PUT
Pushes score events into the system on the privileged lane that bypasses moderation and sampler. Multiple sources: Sportradar's licensed feed via persistent WSS (1-3s SLA from venue), Genius Sports WSS as backup (take-faster dedup at ingest), and an internal POST /v1/score-event endpoint for operator-typed updates (live-blog editor at ESPN.com, score-keeper at NBA arena). All sources authenticate via mTLS + HMAC-signed payload.
Why it exists. Marked kind: external because two of three sources are third-party — Sportradar/Genius are SLA-bound vendors, not internal services. Folding score-event ingestion into the user-comment flow would either (a) put score events behind the moderation classifier (adds 80ms; corrupts the live-feel SLO) or (b) route them through the sampler in mega-rooms (silently drops goals — SEV-1). Separating as a privileged lane with its own dedup window, its own retry policy (circuit-breaker on external, exp-backoff on internal), and its own auth (provider mTLS, not user JWT) is the design's defense.
When it fails. Provider 503 storm (Sportradar maintenance window collides with marquee event) → ingest routes to Genius backup; if BOTH out, score updates stop until provider recovers. Provider feed delayed (network blip at venue, scout's 4G tunnel) → provider_event_age > 3× sport-cadence fires (sport-aware threshold; a flat 8s over-fires on football's 30-60s cadence). Forged event from compromised provider credentials → mTLS rotation + HMAC validation defends. Operator typo — wrong score published — corrected via revision N+1 flow.
- Kafka spine — hash-routed shared topicsKafka (KRaft mode) with tiered storage to S3 (KIP-405); MirrorMaker 2 cross-region; KIP-429 cooperative-sticky
Receives all comment-events + score-events from comment-ingest; partitioned by hash(room_id) % 1024. Also carries moderation-ledger (audit), catalog.changes (cache invalidations), and audit.ops (control-plane operations). Consumed by room-shard (per-room fanout). MirrorMaker 2 replicates to peer regions for read-side fanout (NOT for redundant writes — single-writer-per-room is enforced upstream at ingest).
Why it exists. The durable, partitioned, ordered backbone. Critical that this is hash-routed shared topics, not topic-per-room. Past ~50K active rooms, Kafka's controller metadata snapshot grows past the controller-failover budget; KRaft pushes the ceiling but it's still finite at ~200K. Hash-routed shared topics scale to millions of rooms with constant controller cost — the per-partition consumer (room-shard) demultiplexes by room_id in-pod. YouTube Live Chat's pattern, Slack's pattern.
When it fails. Controller flap during a mega-event (broker process crash + ZK/KRaft re-election) → 10-30s of ActiveControllerCount flipping; producers see NotLeaderForPartition, retry, p99 produce climbs. Mitigated by KRaft + anti-affinity controller placement + 'do NOT bounce brokers during event' deploy rule. Hot partition for a mega-room saturates a single broker NIC; sub-sharding upstream (room-shard splits by (roomId, hash(userId) % subShardCount)) is the design's defense. ISR shrinkage during AZ blip → acks=all blocks producers; falls back to acks=1 via app-level fallback is WRONG (loses durability); correct response is brief 503 + retry. Topic-per-room misconfiguration past 50K rooms → controller metadata snapshot growth → controller failover takes 30s+. MirrorMaker 2 lag tail (60s+) under cross-region link stress → fanout in standby region sees stale events; defense is mm2_lag SLO + read-side replay-buffer regeneration after lag clears.
- Per-room fanout shard (Manifold-style)Elixir/OTP GenServer per room (Discord Manifold pattern); BEAM `min_bin_vheap_size` tuned per Maxjourney post
Owns the in-memory subscriber list for ONE roomId (or one sub-shard slice of a mega-room), consumes its assigned spine partition, demultiplexes events by room_id, checks the lease coordinator to confirm it still owns the single-writer fencing token BEFORE any side-effecting write, appends every event to the per-room replay stream (guarded by that fencing token), and pushes the event downstream — to sampler for mega-rooms, direct to ws-gateway for small/normal rooms. For rooms above 50K concurrent subs, the shard becomes the ROOT of a 2-tier Manifold relay tree: root publishes to ~10 regional relays, each relay to ws-gateway fanout of up to ~5K subs/relay.
Why it exists. Per-room consolidation is the only way the cross-PoP fanout-vs-receive ratio stays sane. Without it, every ws-gateway pod would consume the full per-partition stream and filter locally to its subscribers — wasting ~99% of inter-PoP bandwidth and pinning the spine consumer-group fan-out by N(PoPs). With per-room shards, only PoPs that have a subscriber for room-X get a connection to its shard. The Manifold sub-sharding for mega-rooms is what saved Discord from O(N²) actor cost during the breakout-streamer event — send/2 to N session pids is O(N) inside one scheduler, parallelizing via Manifold cuts it to O(N/M) per scheduler.
When it fails. Hot room (a breakout streamer or marquee event) saturates a single shard's CPU/NIC before sub-sharding pre-warm fires — reactive activation at 40K-subs threshold is the safety net; without it, per-event fanout climbs from 30ms to 1.5s (Discord pre-Manifold). Crashed shard takes its in-memory subscriber list AND in-flight (room_id, client_msg_id) dedup state with it; ws-gateways re-register subscriptions via spine re-consume + replay buffer covers ordering. Lease-check fail-OPEN instead of fail-CLOSED → dual-active publishing on a botched failover → replay-buffer fork (Slack 2018 message-ordering shape). Consumer-offset commit misordered (committed before publish-ack) → lease-loss silently skips events. BEAM min_bin_vheap_size not tuned → GC pathological on Manifold Offload (Maxjourney post-mortem). Mixed-version deploy on the same room (per-room sticky routing bug) → duplicate or missing fanout, detected by per-room delivery_rate divergence canary.
- Mega-room sampler (approximate-by-design)Go service collocated with room-shard for hot affinity; reservoir + engagement-rank + sticky-thread
Inserted between room-shard and ws-gateway egress for rooms in sampling-eligible mode. Runs a 200ms reservoir window for fairness (every incoming comment has equal probability of being kept regardless of within-window timestamp), an engagement-rank scorer (base + reply_count×5 + reaction_count×2 + badge_bonus + sticky_thread_bonus), and a per-viewer sticky-thread bonus for threads the viewer has replied to within 5min. Emits a per-subscriber fanout list at ≤5 msg/s/viewer. Score events are tagged in protobuf and BYPASS the sampler entirely — the privileged-lane invariant.
Why it exists. The product surface of mega-rooms. A 5M-concurrent room with 0.05 msg/s/viewer publish rate ingests 250K comments/s; phone radios cap at ~200 msg/s; eyes glaze above 5/s. Delivering every comment is infeasible AND pointless. The sampler is the deliberate, ranked, fair-by-window mechanism that makes the user experience tractable. Critically separated from room-shard so its failure mode (bad sample, silent room) is observable as a distinct tier — and so the bypass-for-score-events invariant is enforced at a type-tag layer (defense-in-depth: not a post-hoc filter that a bug could remove).
When it fails. Load signal lies (single-source CPU blip) → sampler engages on a small room → users see frozen chat → angry tweet. Defense: AND-gated dual signal + 50% floor + sampler_engaged{non-mega} alert. Loop-coupled feedback oscillation under sustained mild load: if queue_depth were measured at the sampler INPUT AND the sampler decision affects upstream room-shard's push rate (via the keep_fraction config feedback), engagement would oscillate at the sample period (200ms). Defense: queue-depth is measured at room-shard's outbound queue (the canonical pressure signal, NOT sampler input) AND engagement state has 5s hysteresis (engage on >threshold-for-5s, disengage on <threshold-for-30s). Score events sampled by a bug → SEV-1 (goal silently dropped to 5M viewers). Defense: bypass at type-tag ENTRY, not post-hoc filter. Engagement-rank bias (always-amplifies sub-badged users) — fairness regression discovered via per-window keep-rate-distribution audit. CPU saturation during a viral moment (5M subs × per-viewer sticky-thread eval) → fall back to global-reservoir-only (no per-viewer sticky) under load.
- Replay buffer (Redis Streams)Redis Cluster 7.x; Streams with MAXLEN ~ (approximate trim); per-room read-replicas
Holds replay:{roomId} Redis Stream per active room, written by room-shard on every event with XADD * MAXLEN ~ 1000 (approximate trim is O(1) per XADD; exact trim would serialize with writes and cap a hot room at ~10K writes/s). Read by ws-gateway on reconnect with XRANGE replay:{roomId} (Last-Event-ID + for the gap. Per-key TTL = replay_window_seconds from catalog (default 300s). Cluster-sharded by hash(room_id) % slot_count (Redis CRC16 slot mapping).
Why it exists. The 30-second-to-5-minute resume window. Every mobile tunnel, every WiFi-cellular handoff, every brief network blip is a reconnect that needs to know 'what did I miss in the last few seconds'. Without replay, every reconnect becomes a snapshot-fetch, which is 10× more expensive and doesn't deliver the per-comment ordering. The hot-key problem on a mega-room reconnect storm is what drives the read-replica + edge-ring-buffer defenses.
When it fails. Hot-key on a mega-room reconnect storm: 1M clients XRANGE the same replay:{roomId} simultaneously; single Redis shard CPU saturates; XREAD p99 climbs into seconds. Layered defenses (per-room read-replicas + edge ring-buffer + snapshot-fallback at gap > buffer) prevent the cascade, but a misconfiguration that bypasses the edge ring-buffer surfaces this within minutes of a regional ws-gateway flap. AOF rewrite during peak event amplifies write latency — defer to off-peak schedule. Cluster reshard mid-event drops keys briefly during slot migration — defense is 'no cluster ops during event'. Single-region failure during event → 'history unavailable, live only' for ~30s; the explicit accepted trade. Cross-region failover bug: clients reconnect to standby region whose replay is empty → falls through to snapshot-api correctly, but cold-load amplifies snapshot rate 20×. Exact-trim misconfiguration (forgot the ~) caps mega-room writes at ~10K/s → comment loss; canary catches this.
- Snapshot API (CDN-fronted)Go service behind Cloudflare/Fastly CDN; reads catalog + replay current-state
Serves GET /v1/rooms/{room_id}/snapshot — returns the room's current state: room metadata (slow-mode, sampled?), current score (if applicable), and last 50 post-moderated comments. Sized for two access patterns: (a) every cold app-open hits one snapshot per room (1-2 per app-open); (b) reconnects whose Last-Event-ID falls outside the replay buffer fetch a snapshot before re-subscribing. Reads room meta from catalog (5s LRU), current state from replay (XREVRANGE 50), and the deeper historical tail from the replay buffer (XREVRANGE bounded to the replay window) when requested. Assembles response in <30ms p99.
Why it exists. The system has TWO entry points for 'what's in this room right now': the live WS stream (dominant path) and the snapshot (cold connect + replay-gap fallback). Folding snapshot into ws-gateway would couple snapshot capacity to socket capacity — bad, because snapshot traffic spikes 20× when the replay cluster cold-restarts (the exact moment ws-gateway is also under pressure). Separating as its own service lets snapshot scale independently and lets the CDN absorb 90%+ of snapshot traffic with 1-5s TTL.
When it fails. Redis (replay) cold restart during a mega-event: snapshot rate spikes 20×, CDN absorbs most, origin still sees 5× normal. If snapshot-api is undersized, cold-connecting users see 503 → exponential retry → cascade. Mitigation: aggressive CDN TTL + pre-warmed pod fleet at gameday-forecast count. If catalog is also down (worst case), snapshot falls back to a stale-in-memory snapshot served as 200-OK with X-Snapshot-Stale: true header so UI can warn the user.
- Lease coordinator (etcd)etcd v3.6 (Raft), 5 nodes spanning 3 regions for cross-region quorum
Holds /leases/room/{room_id} keys with values { home_region, fencing_token, lease_ttl_30s }. Touched on every region failover by control-plane (Raft CAS for atomic transfer); refreshed by room-shard every 5s to maintain ownership. Linearizable reads/writes via Raft consensus across 5 nodes spanning 3 regions (2+1+2 or 2+2+1 typology).
Why it exists. The single load-bearing invariant of active-active multi-region. Without a linearizable lease, the system has no way to enforce single-writer-per-room — and without single-writer-per-room, the replay buffer forks on failover and the user sees their own comment twice/inverted/missing (Slack 2018 message-ordering bug / Jepsen MongoDB shape). etcd specifically — not Consul, not ZK — for its Raft-based linearizable semantics and the lightweight (~256B/key) lease primitive.
When it fails. Lose a full region → 3-node quorum survives, lease writes continue with higher latency (next-leader is in the surviving region). Lose TWO regions simultaneously → 1 node, no quorum → ALL lease ops block until a region returns. The design accepts this: dual-region-loss is a CONTINENT-level disaster, not a service problem. Leader-election storm under sustained CPU pressure → write latency p99 climbs into seconds; mitigation is dedicated etcd hosts (not co-tenant) + e5 timeout 100ms with leading-indicator alert (lease_renewal_lag > 10s P1). Operator pushes a bad lease via direct kubectl write (bypassing control-plane RBAC) → audit gap; defense is etcd ACLs that only the control-plane IAM principal can write /leases/*. Watch-connection storm during cluster failover (room-shards re-establish watches simultaneously) → brief Raft proposal queue spike; mitigated by watch-multiplexing in the etcd client.
- Control planeGo service + RBAC + audit log; hosts /admin endpoints
Hosts POST /admin/rooms/{id}/slow-mode, POST /admin/rooms/{id}/ban, POST /admin/region/failover, POST /admin/sampler/tune, and the deploy-blackout-gate endpoints. Every operator action is RBAC-gated (room-mod / region-ops / sampler-tuner / deploy-blackout-admin roles), audit-logged with operator identity, justification, and runbook reference, and broadcast via spine audit.ops topic to the audit-log consumer. Talks to catalog (state mutations), lease (etcd CAS for failover), anycast-lb (drain), and gslb (DNS flip).
Why it exists. Operator actions during a high-stakes mega-event (mid-stream cutover, raid response) require the same rigor as data-plane writes: authenticated, authorized, audited, ideally one-command. A kubectl edit against etcd directly has no audit trail and no RBAC for 'who is allowed to transfer a lease' as distinct from 'who can read etcd'. Separating this as a control-plane service is what makes monthly failover drills repeatable and what prevents Slack Jan 4 2021 — that incident's root cause was operator action with no enforcement layer.
When it fails. Control-plane unavailable during a real incident = manual etcd writes via kubectl — possible but slow, no audit, no preconditions. Mitigated by keeping control-plane on a separate deploy cadence + IAM + cluster from data-plane. Buggy precondition check that incorrectly says 'cannot drain' during a real cutover is the worst case — overrideable with a typed --force flag that is itself audited.
- Tracing (OTel + Tempo)OTel Collector + Grafana Tempo; tail-based sampling; always-sample-on-incident
Receives OTel spans from ws-gateway and room-shard (and ingest, moderation) over gRPC. Runs tail-based sampling at the collector — 100% for rooms ≤ 100 subs, ~0.1% for mega-rooms (inverse-room-size). Stores spans in Tempo (S3-backed). On incident=true tag, forces 100% sampling for the incident-tagged room_ids — necessary for per-comment forensics ('did THIS user see THIS comment') that 0.1% sampling can't answer.
Why it exists. Tempo runs on its own S3 bucket + its own Grafana org, deliberately isolated so an observability-backend outage never takes the data plane down with it. Tracing is also explicitly NOT on the synchronous response path — spans are emitted async with bounded local buffer; if Tempo is down, spans drop after buffer fill, but user latency is unaffected.
When it fails. Tempo/collector outage during a SEV-1 → if incident=true 100% sampling is enabled at the SAME moment the collector is down, spans drop after local buffer fills (1000-span ring × 150 ws-gateway pods = 150K spans). Defense is independent Prometheus federation for self-monitoring + multi-region Tempo cluster. Head-based sampling at the app would lose incident data forever; tail-based at the collector buys some resilience. If tracing were on the synchronous response path (bad: someone instruments span-flush as a blocking write), application latency couples to Tempo availability — defense is async buffer + bounded drop policy. Cardinality blowup (someone adds user_id as a span attribute) → Tempo storage cost explodes, alert via storage_growth_rate.
Stage by stage
The same 10 stages the workspace walks, answered.
01Clarifications
What would you ask before drawing a single box?
Clarifications worth surfacing on day 1 — every one of these flips a topology decision:
- Mega-room ceiling. Is the marquee event 1M concurrent (typical viral stream), 5M (Super Bowl / MrBeast / WC final), or 30M+ (cricket IPL final in India)? The Manifold relay-tree activation threshold and the sampler engagement curve flip at each tier.
- Comment provenance. Authenticated users only (Twitch login required), or anonymous + display name (a sportsbook live blog)? The auth tier sizing changes by 10×.
- Score-update provenance. Operator-typed (a live-blog editor at ESPN.com), licensed feed (Sportradar/Genius), or both? Determines whether the score-feed is
kind: external(third-party SLA) or an internalservice(we own the venue scouts). - Latency budgets, per audience. The same architecture serves a 100-viewer hobbyist stream AND a 5M Super Bowl room. The SLOs are DIFFERENT — 400ms p99 for small, 2s p99 for mega — and the design must make that explicit, not aspirational.
- WS vs SSE vs long-poll. Pure web (broad ladder needed), pure native mobile (almost always WS), or both? Determines whether you build the fallback infrastructure or skip it.
- Moderation policy. Pre-moderate (synchronous classifier in the publish path, ~80ms p99 budget), post-moderate (publish first, retract on classifier disagreement), or hybrid? High-risk rooms (children's content, EU regulated, sponsored brand-safety) typically require fail-closed pre-moderation; general chat fails open post-moderate.
- Replay window. Last 30s ("welcome to the conversation"), last 5 minutes ("catch up the last segment"), or full transcript? Drives Redis Streams
MAXLEN ~. - Approximate-by-design opt-in. Do creators / leagues get to opt OUT of sampling (a regulated betting feed must deliver every score event to every connected viewer)? If yes, that lane needs to be isolated cost-wise — a 100-viewer regulated feed shouldn't be sampled, but a 5M-viewer entertainment stream HAS to be.
02Functional reqs
What must this system actually do?
- Publish a comment over a persistent socket (WS primary, SSE fallback, long-poll legacy) and have it delivered to viewers in the same room.
- Subscribe to a room and receive comments + score events as they happen, with
Last-Event-IDresume after a network blip. - Receive score events for the subscribed room on a privileged lane that bypasses moderation latency and sampler drop (low-volume, authoritative).
- Approximate-by-design fanout in mega-rooms: per-viewer rate capped at a human-readable ceiling (~5 msg/s), sampled fairly by window + ranked by engagement, score-events never dropped.
- Honor slow-mode and bans: room-scoped slow-mode delays (3s, 30s, 5min) and per-user bans applied at ingest; client receives
rate_limitedcontrol frame. - Tombstone-and-retract a comment: post-moderate flow can retract a previously-delivered comment; client UI replaces it in place.
- Render
last-revision-winsfor score-event corrections (VAR rollback, scout typo). Revision N+1 always takes priority. - Operator slow-mode toggle without client app reload: control-plane → room-shard → ws-gateway propagates as a server-pushed config frame.
- Cold-load history: client opens a room, fetches the last N comments via
snapshot-api(CDN-cached) before opening the WS. - Regional failover in under 60s without dual-active publishing or replay-buffer divergence; clients see
feed_pausedthen auto-reconnect.
03Non-functional
What must it promise about speed, uptime and correctness?
- End-to-end latency (writer's keypress → other viewers' screens):
- Small room (≤ 1K conc): p50 < 150ms, p99 < 400ms, p99.9 < 800ms. Conversational — must feel synchronous.
- Normal room (10K conc): p50 < 200ms, p99 < 500ms, p99.9 < 1.2s.
- Mega room (1M+ conc, sampled): p50 < 800ms, p99 < 2s, p99.9 < 5s. The viewer cannot perceive 50K msg/s; perceptual freshness matters, not delivery freshness.
- Score event (any room): p50 < 150ms, p99 < 400ms, p99.9 < 800ms — always. Score events are authoritative, low-volume, never sampled.
- Availability on the subscribe path: 99.95% monthly (21.6 min budget). Asymmetric tail: stale-comments = warn; wrong score broadcast = SEV-1 with error budget = 0; toxic content delivered to a children's room = SEV-1.
- Durability:
- Moderation ledger (
moderation-ledgerspine topic): retained on the spine's tiered storage (14d hot / 90d cold) as the post-moderation audit trail. Tombstone (not hard-delete) on user GDPR delete. - Spine event log: 14 days hot tier (KIP-405),
acks=all,min.insync.replicas=2. - Replay buffer (Redis Streams): best-effort, 5-min window, lost on shard failure is acceptable (clients fall back to snapshot-api).
- Throughput at peak (reconciled):
- Mega-room (5M concurrent): 250K msg/s ingest from 5M viewers × 0.05 publish rate. Sampled fanout: 25M msg/s at the egress (5/s/viewer × 5M).
- Across 50K simultaneously live rooms (long tail, 100–10K subs each): ingest ~250K msg/s aggregate, fanout dominated by mid-size rooms, NOT the megas.
- Score-event QPS (all rooms): ~200 events/s (1.5 events/s/match × ~120 concurrent live games + admin-typed events).
- Egress (reconciled):
- Mega-room sampled egress: 25M msg/s × 250 B = 6.25 GB/s ≈ 50 Gbps for the marquee event alone.
- All-rooms steady-state global egress at peak: 80 Gbps across 4 regions.
- Olympics-gold-medal-moment burst: 130 Gbps peak; pre-warmed, not HPA-reactive (HPA is too slow for WS reconnects).
- Fanout factor flip: at room size ≥ 50K subscribers, raw
S × S × rfanout exceeds phone-radio bandwidth (LTE sustained ~5–10 Mbps); sampler engages. - Reconnect-storm resilience: survive a ws-gateway pod restart that dumps 1M sockets in 60s without auth-tier cascade.
- Multi-region: active-active with single-writer-per-room via etcd lease. Reads are active-active; writes for a given
roomIdroute deterministically to home region. RTO 5s on shard failure; RTO 60s on full regional failover. RPO is split into two distinct boundaries that the canonical previously conflated: (1) writer-durability RPO = 0 within the home-region — a comment is acknowledged to the writer ONLY afteracks=all, min.isr=2on the home-region spine; loss of an acknowledged comment is SEV-1. (2) peer-region reader-freshness RPO ≤ MM2 lag (typical 2-5s, stress tail 60s) — peer-region readers may lag MM2 replication; on regional failover we surfacefeed_pausedthen snapshot-required fallback rather than silently serving stale or claiming a 5s SLO we can't honor under stress. Score events use RPO 0 on the reader path via provider replay (Sportradar UOF 72h, Genius 24h) covering the MM2 gap. - Security: TLS 1.3 everywhere; short-lived RS256 JWT + signed JWKS pin (24h offline tolerance); mTLS internal; geographic blackout enforcement fail-CLOSED.
04Capacity estimation
How much load and data does this have to hold?
Inputs (defaults — tune per event):
- Peak concurrent in mega-room: 5M (Super Bowl / WC final / MrBeast). 30M+ exists (IPL) and triggers further sub-sharding.
- Avg publish rate per viewer (steady): 0.05 msg/s (= one comment every 20s, mega-room avg; spikes 10× during goal moments).
- Avg comment size: 250 B on the wire (JSON over WS + signature + user-id).
- Per-viewer sampled fanout ceiling: 5 msg/s (the "human-readable" rate — beyond which the eye glazes).
- Concurrent live rooms at peak: 50K (long tail of 100–10K subs each dominates broker partition decisions).
- Replay buffer per room: 1000 events (~5 min at typical rate).
Derivations — every step shown:
1. Mega-room ingest QPS: 5M × 0.05 = 250K msg/s. This is the LOAD on the comment-ingest tier + moderation.
2. Mega-room raw fanout (if NOT sampled): S × S × r = 5M × 5M × 0.05 = 1.25 × 10¹² msg/s. Infeasible on any topology — phone radios cap at ~50 KB/s (=200 msg/s) for chat. This is the math that justifies the sampler.
3. Mega-room sampled fanout: humanReadableRate × S = 5 × 5M = 25M msg/s. Per-viewer bandwidth: 5 × 250 B = 1.25 KB/s. Fits a 3G connection.
4. Mega-room sampler keep-fraction: 5 / 250000 = 2 × 10⁻⁵ = 0.002%. The sampler drops 99.998% of mega-room comments. Each viewer gets a different 5/s ranked by reservoir + engagement.
5. Mega-room egress: 25M msg/s × 250 B × 8 = 50 Gbps for the single marquee event. Across 4 regions: ~12.5 Gbps per region NIC — well within 25 Gbps NIC capacity. Without sampling, this would be 2.5 Pb/s — the sampler is a 50000× egress cost reduction.
6. Long-tail aggregate ingest: 50K rooms × 100 avg-subs × 0.05 = 250K msg/s aggregate ingest from the long tail (10M global concurrent, less the 5M in the mega-room, leaves ~5M long-tail viewers across ~50K rooms ≈ 100 avg-subs/room). Across 1024 spine partitions: ~250 msg/s/partition average. The mega-room (250K msg/s) is spread across ~100 sub-shard partitions via the composite key (room_id, hash(user_id) % 100) — per-room ordering is reconstructed at room-shard's in-pod sub-shard reducer using (ingest_ts, ingest_region) before fanout. Per-partition load = 250K / 100 = ~2.5K msg/s/partition for the mega-room slot, well below the 50 MB/s/partition (~200K msg/s at 250 B) Kafka p99-produce-latency floor. If we did NOT sub-shard, 250K msg/s on a single partition would exceed that floor and saturate one broker.
7. WS-gateway pod count: at 100K WS/pod operational (40% headroom on 250K theoretical ceiling), for global peak 10M concurrent: 10M / 100K = 100 pods minimum. Run 150 for N+1 + zonal redundancy across 3 AZs. Per-pod sysctl tuning: fs.nr_open=20M, fs.file-max=12M, somaxconn=65535, ip_local_port_range=1024-65535.
8. Mega-room sub-sharding: above 100K subscribers in a single room, a single room-shard actor's egress NIC and CPU saturate. Sub-shard by (roomId, hash(userId) % subShardCount) where subShardCount = ceil(subs / 50K). For 5M: 5M / 50K = 100 sub-shards. Each sub-shard fans out 5/s × 50K = 250K msg/s × 250 B = 62.5 MB/s — well within 25 Gbps NIC.
9. Spine partition count = 1024 with partition = hash(roomId) % 1024. Hash-routed shared topics — NOT topic-per-room. Past ~50K active rooms, Kafka controller metadata thrashes (~200K topic limit even on KRaft); shared topics scale to millions of rooms. In-pod demultiplex on consume.
10. Replay buffer footprint: per-room 1000 events × 250 B = 250 KB. 50K rooms × 250 KB = 12.5 GB hot; trivial in Redis Cluster. MAXLEN ~ 1000 (approximate trim — O(1) per XADD instead of O(log N) for exact).
11. etcd lease store: one key per active room × ~256 B/key × 50K = 12.8 MB. Comfortably below etcd's recommended 8 GB ceiling.
Architectural flip-points (with the rejected alternative + the citation):
- At 250K WS/pod Nginx+Lua tops out — flip to Centrifugo (Go) or Cloudflare Durable Objects + WS Hibernation. WhatsApp's 2M/box was Erlang on FreeBSD bare-metal; k8s + TLS terminates lower.
- At 50K subscribers per room the single-actor fanout (Erlang/Go) saturates — flip to Manifold sub-shards + relay tree (Discord 2017 Manifold blog + 2024 Maxjourney post).
- At 50K active rooms topic-per-room Kafka thrashes the controller — flip to hash-routed shared topics (
room.events.{0..1023}). - At 2M global concurrent single-region active-passive breaks under RTT — flip to multi-region active-active with single-writer-per-room fencing.
- At 50% sampler keep-fraction the "per-viewer 5/s" ceiling stops being the binding constraint and ingest QPS becomes the floor — sampler can be configured to relax above this point.
05API design
What does the outside world call, and what comes back?
# Subscribe + bidirectional comment exchange via WSS (the dominant entry point)
GET /v1/rooms/{roomId}/stream
Host: live.example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: ...
Authorization: Bearer <15-min JWT>
Last-Event-ID: 78145327 # optional, reconnect resume
101 Switching Protocols
Set-Cookie: edge_session=<signed cookie>; SameSite=Strict; Secure; HttpOnly
After the WS upgrade, the client SDK sends and receives JSON frames:
// Client → server (publish a comment; idempotent via clientMsgId)
{ "op": "publish", "client_msg_id": "abc123",
"room_id": "twitch:mrbeast:2026-05-17",
"body": "BRO IS LIVE", "reply_to": null }
// Server → client (provisional ack — reconcile the optimistic render on client_msg_id;
// the authoritative event_id arrives on the fanout copy, below)
{ "type": "publish_ack", "client_msg_id": "abc123", "partition": 512, "offset": 88991234 }
// Server → client (comment from another viewer)
{ "type": "comment", "event_id": 88991235, "room_id": "...",
"user": { "id": "u_q9z", "display": "...", "badges": ["sub"] },
"body": "no way", "ts": "2026-05-17T20:14:33.122Z" }
// Server → client (score event — privileged lane, never sampled)
{ "type": "score_event", "event_id": 88991240, "room_id": "...",
"revision": 1, "kind": "goal", "score": { "home": 1, "away": 0 },
"minute": "78:14" }
// Server → client (retraction — post-moderate disagreement)
{ "type": "retract", "event_id": 88991235, "reason": "policy.toxicity" }
// Server → client (slow-mode engaged — config push)
{ "type": "config", "slow_mode_seconds": 30, "until_event_id": 88991500 }
// Server → client (sampler engaged — UI badge "sampled view")
{ "type": "config", "sampled": true, "keep_fraction": 0.00002 }
// Server → client (feed paused — region failover or maintenance)
{ "type": "feed_paused", "reason": "regional_failover",
"snapshot_url": "/v1/rooms/.../snapshot", "retry_after_seconds": 5 }
// Server → client (rate-limited — caller exceeded slow-mode or per-IP cap)
{ "type": "rate_limited", "retry_after_ms": 30000, "reason": "slow_mode" }
# Cold-load snapshot (CDN edge-worker authenticates + geo-derives region BEFORE cache lookup)
GET /v1/rooms/{roomId}/snapshot
Authorization: Bearer <15-min JWT>
200 OK
Cache-Control: private, max-age=1, stale-while-revalidate=4 # edge worker caches internally, keyed on edge-derived (room, region, auth-tier); never public
{ "room_id": "...", "as_of_event_id": 88991500,
"current_score": { "home": 1, "away": 0 },
"recent_comments": [...last 50 post-moderated, with sampled=true marker if mega...] }
# SSE fallback (when WS Upgrade is blocked by corporate proxy)
GET /v1/rooms/{roomId}/sse
Authorization: Bearer <15-min JWT>
Accept: text/event-stream
Last-Event-ID: 78145327
200 OK
Content-Type: text/event-stream
X-Accel-Buffering: no # CRITICAL — disables nginx/proxy chunked buffering
event: comment
id: 88991235
data: { ... same payload as WS ... }
: keepalive # comment line every 15s, beats every middlebox idle timeout
# Long-poll legacy fallback (caps concurrency, last resort)
GET /v1/rooms/{roomId}/poll?since=78145327&timeout=25
Authorization: Bearer <15-min JWT>
200 OK
# next_since = max(event_id) in events (echo the request `since` on an empty batch);
# it is monotonically non-decreasing vs the request `since`. The client always feeds
# next_since back as `since` and NEVER sends a `since` lower than the previous next_since.
{ "events": [...last event_id 88991389...], "next_since": 88991389, "next_timeout_ms": 25000 }
# Operator slow-mode (control-plane, RBAC-gated, audit-logged)
POST /admin/rooms/{roomId}/slow-mode
Authorization: Bearer <operator JWT, role: room-mod>
{ "seconds": 30, "duration_ms": 600000, "justification": "raid in progress" }
# Score-event ingress (privileged write lane — Sportradar/Genius/operator)
# Bypasses moderation AND the sampler. Idempotent per provider on (matchId, eventId, provider);
# cross-provider take-faster dedups on the normalized (matchId, kind, minute) key — provider
# eventIds are namespaced (sr:goal:5512) and never match across providers.
POST /v1/score-event
X-Provider-Id: sportradar # provider identity (mTLS client cert MUST match)
X-Signature: HMAC-SHA256(provider_secret, body) # forged event from a non-provider source is rejected
Content-Type: application/json
{ "matchId": "epl:LIV-ARS:2026-05-17", "eventId": "sr:goal:5512",
"provider": "sportradar", "revision": 1,
"kind": "goal", "score": { "home": 1, "away": 0 }, "minute": "78:14" }
202 Accepted # queued to spine; fanout to every viewer, every revision
# 200 OK — idempotent replay of a known (matchId, eventId, provider); same result, not re-published
# 401 — HMAC signature or mTLS client-cert identity mismatch
# 409 — duplicate (matchId, eventId, provider) with a conflicting payload
# 409 / 422 — room state != live (no live room to attach the score to)
Transport ladder rationale:
| Rung | Use when | Don't trust because |
|---|---|---|
| WSS primary | Native mobile (always), modern web | none — preferred when available |
| SSE secondary | Corporate proxy blocks WS Upgrade | Buffer-by-proxy without X-Accel-Buffering: no; HTTP/2 HOL-blocking; browser per-host stream cap (~100) |
| Long-poll last | SSE also broken (mitmproxy-class DPI box, very old browser) | Quadratic server load; cap client concurrency; surface "Connection problems" rather than infinite climb |
06Data model
What gets stored, and what is it looked up by?
The data model is deliberately sparse — the system is a fanout fabric with a post-moderation audit trail on the spine, not a database. The load-bearing stores:
1. Room catalog (Postgres HA):
| field | type | notes |
|---|---|---|
| room_id | varchar(64) | PK, e.g. twitch:mrbeast:2026-05-17, epl:LIV-ARS:2026-05-17 |
| kind | enum | stream / match / live-blog / event |
| home_region | varchar(16) | consistent-hash on roomId; single-writer destination |
| state | enum | scheduled / live / ended / archived |
| slow_mode_seconds | int | 0 = off; pushed live to subscribers on change |
| sampling_policy | enum | none (regulated feeds), auto (engages > 50K subs), always (test) |
| blackout_regions | text[] | NFL/EPL geographic rules |
| created_at, ended_at | timestamp | |
| owner_id | varchar(64) | for the live-blog editor or stream creator |
| moderation_policy | enum | pre-moderate-strict (children's), pre-moderate (default), post-moderate-fast (mega-stream) |
Plus tables: banned_users (room_id, user_id, until), pinned_messages (room_id, event_id, until).
2. Comment event (spine wire format — Protobuf):
message Event {
uint64 event_id = 1; // room-global total order, assigned by the room-shard
// reducer (see below); ingest does NOT set it — publish_ack
// carries a provisional spine (partition, offset) instead
string room_id = 2; // hash-routed to spine partition
string client_msg_id = 3; // user-supplied idempotency key
oneof body {
Comment comment = 10;
ScoreEvent score = 11;
Retract retract = 12;
ConfigChange config = 13;
}
uint32 revision = 14; // bumped on correction; client renders last-revision-wins
google.protobuf.Timestamp ingest_ts = 15;
string ingest_region = 16; // for cross-region observability
ModerationDecision moderation = 17; // null on pre-moderate-fail (event suppressed)
}
Idempotency key for ingest dedupe: (room_id, client_msg_id, user_id). The same key submitted twice (browser retry, network blip) returns the same ack and never publishes a second event.
Who assigns event_id. Stateless, sub-sharded ingest (80 pods, mega-rooms spread across ~100 sub-shard partitions by (room_id, hash(user_id) % N)) cannot produce a single per-room counter — nor even a per-(room, sub_shard) one, since any of the 80 pods can produce for any user and there is no shared sequence allocator at that tier. The only monotonic per-partition sequence that exists at produce time is the Kafka offset. So ingest returns the publish_ack synchronously carrying client_msg_id plus the provisional spine coordinate (partition, offset) from the produce response — NOT a room-global id, which does not exist yet. The room-shard reducer (single-writer-per-room, holding the lease) then merges its sub-shard partitions into the room-global total order — ordering on (ingest_ts, ingest_region), the produce offset breaking ties within a sub-shard — and assigns the monotonic room-global event_id used in replay (XADD) and fanout. The client learns that authoritative id from the delivered fanout copy, reconciles its optimistic render on client_msg_id, and persists ONLY the room-shard-assigned id as Last-Event-ID — only that id is monotonic per room and safe to resume against.
3. Replay stream (Redis Streams): replay:{room_id} → XADD * <protobuf-encoded-event>, trimmed MAXLEN ~ 1000. Approximate trim (~) is O(1) per XADD; exact trim would serialize with writes and cap a hot room at ~10K writes/s. Per-key TTL = replay_window_seconds from catalog (default 300s).
4. Lease (etcd): /leases/room/{room_id} → { home_region, fencing_token, lease_ttl_30s }. Raft-linearizable. Touched on every region failover; otherwise idle.
07High-level design
Which components handle a request, and in what order?
Architecture summary — publish flow, fanout flow, score flow, cold-connect, reconnect, control plane:
Publish flow (the user posts a comment):
- client holds a single open WSS to ws-gateway (or SSE/long-poll fallback). Sends
{op:"publish", client_msg_id, room_id, body}. - ws-gateway authenticates the JWT against its local JWKS cache (signed-pin tolerates 24h auth outage), looks up the user's per-room ban + slow-mode state in catalog's rate-limit cache (in-process LRU, 5s TTL), and forwards the publish to comment-ingest over gRPC.
- comment-ingest enforces rate-limit (per-user, per-IP, per-room), dedupes on
(room_id, client_msg_id, user_id), calls moderation synchronously for a classifier decision (<80ms p99 budget; circuit-breaker opens on 50% errors / 30s). Onblock, returnsmoderation_failedto client. Onalloworsoft-flag, publishes the event to spine withacks=all, enable.idempotence=true, then writes the moderation decision to the spinemoderation-ledgertopic (the post-moderation audit trail, retained on the spine's tiered storage). - A
publish_ackflows back to the client via the same WS socket, carryingclient_msg_id+ a provisional spine(partition, offset)(the authoritative room-globalevent_idarrives with the fanout copy — step 9). The publisher SEES THEIR OWN COMMENT in the local UI immediately on ack (optimistic render), reconciled onclient_msg_idwhen the fanout copy arrives.
Fanout flow (others receive the comment):
- spine is partitioned by
hash(roomId) % 1024. Each partition is consumed by one room-shard actor that owns the room (or sub-shard, for rooms > 50K subs). - room-shard first checks lease to confirm it still owns the single-writer fencing token for this room (fail-CLOSED on lease loss → emits
feed_pausedto subscribers and makes NO further writes), and only THEN appends the event to replay (XADD replay:{roomId}withMAXLEN ~ 1000, carrying the fencing token so a stale writer's append is rejected at the store). - room-shard decides: if
room.size > 50Kandsampling_policy != "none", push the event to sampler; else push directly to ws-gateway fanout. - sampler (only on mega-rooms) runs a 200ms reservoir window + engagement-rank + sticky-thread scoring, emits a per-subscriber fanout list at ≤5 msg/s/viewer. Score events bypass the sampler entirely on the privileged lane.
- ws-gateway flushes the WS frame into the open socket of every subscribed client (per-socket coalescing on backpressure — coalesce only same-
(roomId, eventId, revision), never collapse a revision N+1). - client renders the comment, persisting
last_event_idto disk for reconnect.
Score flow (admin or licensed feed):
- score-feed (external — Sportradar UOF / Genius / live-blog editor's PUT) calls comment-ingest on the privileged HTTPS endpoint
POST /v1/score-event. Auth via mTLS + signed payload. - comment-ingest validates against catalog (matchId in
livestate), publishes to spine. Score events havebody.scoreand bypass the moderation classifier (operator-typed is trusted; licensed feed is contractually accurate). - Fanout flow steps 5–10 apply, with one difference: sampler is skipped on
body.scoreevents. Every connected viewer receives every score event, every revision.
Cold-connect / subscribe flow:
- client resolves
live.example.comvia gslb (DNS TTL 30s) to the nearest healthy region's anycast prefix. - client opens WSS to anycast-lb; admission control runs at L4 (1.2× steady-state token bucket, drainable by control-plane).
- anycast-lb forwards raw TCP to a co-located ws-gateway instance.
- ws-gateway terminates TLS + WS Upgrade, validates the JWT against local JWKS cache + signed pin, reads catalog for room entitlement + blackout + slow-mode (fail-CLOSED with coordinated
Retry-After). - ws-gateway registers the subscription on the correct room-shard (via gRPC
op: controlregister edge). - If client supplied
Last-Event-ID, ws-gateway doesXRANGEon replay for the gap; on gap > buffer head, returnssnapshot-requiredand the client fetches snapshot-api (CDN-fronted).
Reconnect + replay flow (the dominant resume path — every tunnel, every WiFi handoff):
- Client supplies
Last-Event-ID: Xon reconnect. ws-gateway runsXRANGE replay:{roomId} (X +against replay, replays each gap event into the socket, only THEN subscribes to live. Branch:replay.
Snapshot fallback (gap > replay buffer or first-load):
- client GETs
/v1/rooms/{roomId}/snapshotagainst snapshot-api (CDN-fronted with 1-5s TTL — absorbs 90%+ of cold loads). - snapshot-api reads catalog for room state, replay for the last 50 events (XREVRANGE), and the deeper replay tail for older comments if requested, assembles in <30ms p99. Branch:
snapshot-required.
Control plane (operator actions, failover, slow-mode):
- Operator hits control-plane
/admin/rooms/.../slow-modeor/admin/region/failover. RBAC-gated (room-modvsregion-opsroles), audit-logged. - Control-plane writes the new state to catalog AND publishes a control event to spine. Room-shard consumes, pushes a
configframe to subscribers via ws-gateway. For failover, control-plane wraps drain → lease-CAS → DNS flip in one transaction (single command, audited).
Observability:
- ws-gateway and room-shard emit OTel spans to tracing (tail-sampled, with always-sample-on-incident override).
Multi-region: active-active, single-writer-per-room by consistent-hash home_region(roomId). Reads from any region. MirrorMaker 2 replicates spine cross-region for read-side fanout (NOT for redundant writes). etcd lease coordinator spans 3 regions for quorum.
08Deep dives
Which part breaks first, and what do you do about it?
1. The MrBeast room — approximate-by-design is not a degradation, it's a feature.
A 5M-concurrent room with 0.05 msg/s/viewer publish rate ingests 250,000 comments per second. A 3G phone radio sustains ~50 KB/s on the downlink. At 250 B/comment that's 200 msg/s before the buffer fills — and the eye glazes over above ~5 msg/s anyway. The sampler is therefore not an emergency degradation; it is the product surface of a mega-room.
The policy lives in three layers:
- Fairness: 200ms reservoir window. Within each window, every incoming comment has equal probability of being kept, regardless of submission timestamp inside the window. This avoids the "early bird wins" bias of a naive head-cut.
- Value: engagement-rank scoring within the window — comments with replies, reactions, sponsored badges, moderator pins, or from accounts the viewer follows get a multiplicative bonus. The actual scoring is
score = base + reply_count × 5 + reaction_count × 2 + badge_bonus + sticky_thread_bonus. - Personalization: sticky-thread bonus — if a viewer has replied to a thread within the last 5 minutes, comments in that thread are 4× more likely to land in their sampled view. Costs CPU on the sampler per-viewer.
The load signal driving sampler engagement is two-source AND-gated: queue_depth > threshold AND cpu > threshold. Single-signal sampling has the classic Spectroscope failure mode — a CPU blip from a GC pause flips the sampler to "drop 90%" during a non-spiky moment and the room goes silent for ten minutes, viewers think it's broken. AND-logic means both signals must agree; alert fires if sampler engages outside expected windows. Floor at 50% drop — never sample more aggressively than that without operator override, so a sampler bug doesn't take the room completely dark.
Score events are tagged in the protobuf body and bypass the sampler entirely. The privileged-lane invariant is what makes the "goal at minute 78" never get sampled out — same as the live-sports-scores canonical, but here it coexists with a 99.998%-drop comment lane.
2. Topic-per-room breaks at 50K rooms — hash-routed shared topics are the answer.
The naive design: one Kafka topic per roomId. Easy mental model, lets consumers subscribe room-by-room. It breaks before you fill broker NICs because the Kafka controller's topic-metadata snapshot grows linearly with topic count, and controller failover time grows with snapshot size. Past ~200K topics (KRaft pushes higher than ZooKeeper, but it's still finite), controller failover takes 30s+ during which producers stall and partitions can't rebalance. We hit this at 50K active rooms, far below the broker-CPU or NIC ceiling.
The fix is hash-routed shared topics: 1024 partitions across 4 topics (room.events.{0..3}.{0..255}), with partition = hash(roomId) % 1024. Producers compute the partition deterministically; consumers (room-shards) own a partition and in-pod demultiplex the partition stream by roomId. This is YouTube Live Chat's pattern, and Slack's. The cost: a consumer can no longer "subscribe to just the rooms it cares about" — it gets the whole partition. The benefit: scale to millions of rooms without controller pain.
Per-partition load math: at 250K msg/s peak for a mega-room landing on ~100 partitions (one room sub-sharded, ~2.5K msg/s each), and a ~250 msg/s/partition baseline for the long tail, p99 produce latency stays well under 50ms (Kafka's well-known LinkedIn floor of ~50 MB/s/partition = ~200K msg/s at 250 B).
3. Single-writer-per-room is the load-bearing invariant that makes active-active honest.
The trap in any active-active fanout system is the "dual-write for safety" anti-pattern. If a comment for room=R is accepted in us-east-1 AND eu-west-1 simultaneously, the replay buffer in each region holds different orderings, and on cross-region read the user sees their own comment twice, or interleaved with itself, or in inverted causal order with a reply. This is the shape of Slack's 2018 message-ordering bug and the Jepsen-documented MongoDB / CockroachDB orderings.
Our defense: single-writer-per-room via etcd lease. Every roomId deterministically hashes to a home_region. Comment-ingest in the OTHER region (the reader-region) PROXIES the publish via mTLS+gRPC to the home-region's comment-ingest — it does NOT publish locally. Reads (consume from spine + push to ws-gateway) are active-active everywhere via MirrorMaker 2. Writes go to ONE region only.
The lease is what makes this safe under failover: when the operator runs score-ops region failover, the lease CAS in etcd atomically moves home_region from us-east-1 → eu-west-1. Old-region ingest sees lease-loss within 5s, fails CLOSED on local publishes (returns feed_paused-style 503 with Retry-After), forwards in-flight publishes to the new region. The drain order (drain L4 → CAS lease → flip DNS) is the same pattern that's spelled out in the live-sports-scores canonical and that prevents the Cloudflare Jun 21 2022 / Slack Jan 4 2021 dual-active footgun.
4. WS-SSE-long-poll fallback ladder — different infra per rung or you co-locate the failure.
The client SDK climbs the transport ladder: WS first, SSE second, long-poll last. The trap is putting all three rungs behind the same CDN, same anycast prefix, same edge — when the edge fails, every rung fails together. We split:
- WS terminates at our own anycast-lb + ws-gateway (Centrifugo / Durable Object).
- SSE terminates at Fastly Fanout (or our equivalent GRIP proxy), which converts the SSE subscription into a plain HTTP GET against ws-gateway — origin sends a
Grip-Holdinstruction and drops; Fastly holds the SSE client connection at the edge until it has something to push. Critical headers:X-Accel-Buffering: no+ a 15s keepalive comment line. Without these, corporate proxies and some nginx defaults buffer 8 KB of chunked output and SSE arrives in bursts. - Long-poll terminates at a third path — a separate CDN with a 25s server-side hold. Concurrency-capped on the client SDK to avoid quadratic load if the user opens 12 tabs.
The client SDK enforces a bounded climb: at most 3 rungs, with decorrelated-jitter backoff between rungs. On rung-3 failure, surface "Connection problems — last comment delivered at HH:MM" rather than recurse. Telemetry: transport=ws/sse/longpoll is a first-class label on every connect metric.
Mobile is the inversion: native apps almost always prefer WS (single multiplexed socket; opening a new TCP/TLS connection wakes the radio for ~10–20s of high-power state and burns battery).
HTTP/2 SSE has its own trap: multiple SSE streams share a TCP connection, a lost packet stalls every multiplexed stream (TCP HOL blocking that H2 multiplexing was supposed to fix but only fixed at the app layer). HTTP/3/QUIC isolates loss per QUIC stream and is the real fix, but at our scale we still see H2 in legacy regions — sharding SSE across sse-1.live.example.com, sse-2.live.example.com works around it.
5. Moderation in the publish path without putting it ON the publish-path SLO.
Moderation is synchronous in the default policy — comment-ingest waits for the classifier verdict before publishing to spine. That coupling is a 3 AM page waiting to happen: classifier p99 climbs from 30ms to 1.2s under load and the publish-path latency follows; users retry; the queue compounds.
The bulkhead policy is per-room, set in the catalog moderation_policy field:
pre-moderate-strict— children's rooms, regulated content, sponsored brand-safety. Fail-CLOSED on classifier timeout: comment is HELD, user seesmoderation_pendingwith retry. Classifier outage = comments stop. This is the correct behavior for these rooms.pre-moderate(default) — general chat. 80ms p99 budget. Circuit-breaker opens at 50% error rate / 30s; while open, comments fall through withflagged=pending, async classifier runs over them with retroactive retract on disagreement.post-moderate-fast— mega-stream chat. Comments publish immediately withmod=async; classifier runs over the spinemoderation-pendingtopic; onblockverdict aretractevent is published to the room and viewers see the comment replaced in place. The trade is that toxic content has a ~500ms window where it's visible.
The single critical invariant: the classifier never extends the per-comment publish-path latency more than the budget. If it can't decide in 80ms, fall through per policy. Don't queue the user's send button waiting for the classifier.
6. Replay buffer hot-key — the 90%-of-reconnect-users-hit-one-key problem.
When a region drops and 1M clients reconnect simultaneously to a mega-room, every one of them executes XRANGE replay:{roomId} (last_event_id +. That's one Redis key, on one shard, taking 1M reads in 60s. CPU on that Redis node saturates; XREAD p99 climbs into seconds; the reconnect spins.
Defenses, layered:
- Per-room read-replicas:
replay:{roomId}is replicated to 2 read-replicas (Redis Cluster with replica-read mode); ws-gateway picks one per reconnect via consistent-hash onuserId. Spreads the read load 3×. - Edge ring-buffer: ws-gateway maintains an in-pod ring buffer of the last 30s of events for every room it serves a subscriber for. On reconnect with
Last-Event-IDwithin that window, serve locally — never hits Redis. - Snapshot fallback at gap > buffer: if
Last-Event-ID < replay buffer head, returnsnapshot-requiredand route to snapshot-api (CDN-cached). Most reconnects after long tunnels (> 5 min) bypass replay entirely. - Reconnect-fanout token-bucket at anycast-lb: the L4 admission cap is
1.2× steady-state, deliberately undersized for storms. Storms spread across the decorrelated-jitter window naturally.
09Trade-offs
What did this design cost, and what breaks at 10×?
What we accepted, with the rejected alternative and the citation:
| Decision | Rejected alternative | Why |
|---|---|---|
| Approximate-by-design sampler (drop 99.998% in mega-rooms) | Deliver every comment to every viewer | Phone radios cap at ~200 msg/s; eye glazes at 5/s. Without sampling, mega-room egress is 2.5 Pb/s — infeasible. |
| Hash-routed shared topics (1024 partitions, in-pod demux) | Topic-per-room Kafka | Past ~50K rooms the controller metadata thrashes; failover takes 30s+ during which producers stall (YouTube Live Chat / Slack pattern). |
| Single-writer-per-room via etcd lease + cross-region proxy on publish | Multi-master accept-everywhere, reconcile on read | Multi-master causes replay-buffer divergence and out-of-order delivery (Slack 2018 / Jepsen MongoDB shape). Fencing accepts a small cross-region RTT cost. |
| Centrifugo (Go) primary + Cloudflare DO + WS Hibernation as roadmap | One vendor only | DO cost model differs 5×; we pick Centrifugo as primary, keep DO as strategic future. Both terminate WS at the edge; only one is operational today. |
| Synchronous pre-moderate (default) with per-room post-moderate-fast override | All-synchronous OR all-post-moderate | Children's rooms need fail-closed; mega-streams can't afford 80ms in the publish path. Per-room moderation_policy in catalog is the surface. |
| Sampler bypassed for score events on the privileged lane | Sample uniformly | Score events are low-volume, authoritative, never-droppable. Mixing them with comment sampling is the bug class. |
| WS → SSE → long-poll ladder over THREE different edge paths | Single edge for all three | Co-locating fallback infra means edge failure collapses every rung simultaneously. |
| SSE behind Fastly Fanout / GRIP (proxy-and-hold) | SSE direct to origin | GRIP terminates the streaming HTTP at the edge; origin returns a plain HTTP response with Grip-Hold and drops, letting backend stay stateless. (Fastly Fanout / Pushpin pattern.) |
Redis Streams MAXLEN ~ 1000 approximate trim | Exact MAXLEN 1000 | Exact trim is O(log N), serializes with writes, caps a hot room at ~10K writes/s. Approximate trim is O(1). |
| Decorrelated-jitter client reconnect, base 1s cap 30s (AWS Builders' Library) | Default EventSource 3s no-jitter | Default 3s no-jitter is the Slack/Discord reconnect-storm shape. |
| WS heartbeat at 15s | TCP keepalive | Most middleboxes ignore TCP keepalive; only application-layer payload counts. 15s beats every proxy idle timeout we've measured. |
| Signed JWKS pin (24h offline tolerance) | Fall back to remote auth on JWKS miss with exp-backoff | Exp-backoff cache-miss path is the AWS Builders' Library "avoid fallback" anti-pattern; auth becomes the cascade target. |
| Two-source AND-gated sampler engagement signal | CPU-only OR queue-only | Single-signal fires on a GC blip and the room goes silent for ten minutes. AND-logic + 50% floor is the SRE-Workbook load-shedding contract. |
| Postgres HA for room catalog | Cassandra / Dynamo | 50K active rooms × handful of mutations/room/day is well within Postgres HA; second datastore is an operational tax. |
| Edge ring-buffer in ws-gateway (last 30s per room) | Hit Redis every reconnect | Hot-key on replay:{roomId} during reconnect storm = Redis shard saturation. Local serving spreads the load. |
Deliberately scoped out (v1): mobile push (APNs/FCM), the approximate live-viewer-count badge, and the delivery-telemetry / long-retention comment archive are out of scope for this canonical, which focuses on the core publish → moderate → fanout spine. Each layers back on as a new consumer over the same Kafka spine — a push-dispatcher consumer group off a notif.push topic, a presence-count reducer fed by ws-gateway connect/disconnect events, and a ClickHouse/S3 archive consumer off moderation-ledger — without touching the core path.
Primary sources
- Slack — Flannel: an application-level edge cache
- Discord — Scaling Elixir to 5M concurrent + Manifold + Maxjourney
- Cloudflare — Durable Objects WebSocket Hibernation
- WhatsApp — 2M sockets/box on Erlang BEAM (FreeBSD tuning)
- LINE LIVE — sub-room sharding for celebrity streams
- Twitch — chat architecture (Room + Clue, Go rewrite)
- Fastly Fanout / Pushpin — GRIP proxy-and-hold
- Redis Streams — XADD MAXLEN ~ N approximate trim
- AWS Builders' Library — avoid retry storms / decorrelated jitter
- SRE Workbook Ch.22 — Addressing cascading failures
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 Live Comments / Score Updates 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.
- 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.
- Online IndicatorGreen dot for contacts. Mind the N² watch problem. Approximate by design — never quote presence more precisely than reality.
- 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.