Concurrent Hotel Viewers
Worked solution

Concurrent Hotel Viewers — a worked solution

"X users viewing this right now." Hot keys, HLL.

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 Concurrent Hotel Viewers workspace

The problem

When a user opens a hotel detail page, show "127 people are viewing this right now." The display refreshes as users arrive and leave. Sounds trivial — until you realize that on a viral hotel during a holiday weekend, you have one entity (one hotel ID) that 100K users are pinging every few seconds. Every detail of this design — TTL choice, push vs poll, exactness, where you keep state — flips when you account for the hot-key reality.

This canonical is sized for 5M concurrent active hotel pages globally (production-tier 3) with viral hotels reaching 100K concurrent viewers on a single ID. The design tilts heavily on three load-bearing tricks: edge fan-in (per-PoP coalescing collapses each PoP's writes into one batched origin call per second), hashtag-pinned sharded HyperLogLog (16–64 parallel HLLs per viral hotel co-located in one Redis slot for atomic PFMERGE), and a 2-bucket sliding window on top of the HLL to kill the 30-second roll-over cliff. Push fan-out follows Pusher's published coalesce-and-broadcast pattern (≤100 subscribers broadcast every change; >100 broadcast at most once per 5 s) layered over Discord-style sticky-hash gateway routing.

The reference architecture

Reference architecture for Concurrent Hotel Viewers: 17 components — Client (page), CDN, Function (FaaS) · Edge Aggregator, Load Balancer, API Gateway, Service · Heartbeat, Service · Count, Service · WS/SSE Fanout, Cache · Redis HLL, Stream · Count Pub/Sub, Stream · Heartbeat Events, Worker · Abuse Scanner, SQL DB · Metadata, Cache · Rate-Limit Redis, Service · Hot-Key Controller, Monitoring · Observability, Service · Peer Region Count — connected by 25 flows.Client (page)clientCDNFastly (VCL) + origin s…Function (FaaS) · Edg…Cloudflare Workers + Du…Load BalancerEnvoy ring-hash tier (G…API GatewayEnvoy + custom WASM fil…Service · HeartbeatRust + Tokio (async run…Service · CountRust + Tokio + tower-si…Service · WS/SSE Fano…Go + nhooyr/websocket, …Cache · Redis HLLRedis 7.4 Cluster, m7g.…Stream · Count Pub/SubRedis 7.0 Sharded Pub/S…Stream · Heartbeat Ev…Apache Kafka 3.7, MSK m…Worker · Abuse ScannerGo + sliding-window cou…SQL DB · MetadataPostgres 16 HA (Patroni…Cache · Rate-Limit Re…Redis 7.4 single primar…Service · Hot-Key Con…Go + leader-elected via…Monitoring · Observab…Prometheus 2.x + Thanos…Service · Peer Region…Stub representing a pee…
17 components, 25 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Client (page)

Renders the hotel detail page; on load fires POST /presence/heartbeat every 15s; subscribes via WebSocket to count.{hotel:H} for live updates and falls back to a 10s GET /presence/count poll if the WS handshake fails. Tags every request with a fingerprinted session cookie set on first paint.

Why it exists. It is the source of every write and read in this design — the heartbeat *is* the presence signal. We chose a client-pushed heartbeat over server-tracked TCP-state (the BookingChat-style 'count subscriber-set on the gateway') because TCP-state breaks the moment a CDN sits in the middle: the gateway sees PoP IPs, not user identities, and you've inflated your count to 200.

When it fails. If the heartbeat fails, the user's own count drifts down by one within 30s; if the count widget itself fails, render —. The page MUST NOT be on the critical-rendering path of the count widget — synthetic checks watch for blank-page time on hotel detail pages and block any deploy that puts the count XHR on a await before first paint.

CDNFastly (VCL) + origin shield

Serves GET /presence/count/:hotel from edge cache with 5s TTL and stale-while-revalidate up to 30s. Forwards POST /presence/heartbeat uncached but with WAF + per-IP/ASN rate-limit. Terminates TLS, ships logs to the abuse pipeline. Origin shield centralises misses to one PoP per region so the origin sees one miss per region per TTL window, not one per edge POP.

Why it exists. The cardinal trick of this whole design: with T_edge=5s, origin reads per hotel are bounded at 1/T_edge = 0.2 reads/s regardless of viewer count. Without the CDN, 5M concurrent viewers polling every 10s = 500K reads/s slamming the count service. We considered putting count behind an internal cache only (skip the CDN), but that means every viewer pays a full origin RTT per poll — UX research shows the count widget perceptibly slows the page above 60ms.

When it fails. PoP outage (Fastly 2021-06-08, 49 min global): origin sees 20–50× normal QPS within seconds. Detection: cdn_5xx_rate{pop=$x}>5% for 2m. Mitigation tiered — (1) origin must hold ≥3× steady-state for 60 min, (2) the count widget JS catches and shows —, (3) regional DNS failover to alternate edge provider. The widget MUST NOT block the page if the CDN is down.

Function (FaaS) · Edge AggregatorCloudflare Workers + Durable Objects (per-PoP)

Two responsibilities, one runtime. WRITE PATH: receives POST /presence/heartbeat at the PoP, verifies the session-fingerprint cookie signature (rotating Ed25519 verify key pushed to every PoP — verification needs no signing secret) and applies the PoP-local 10 heartbeats/sec per-fingerprint token bucket (fingerprints are PoP-sticky, so local enforcement is global enforcement), computes x-fingerprint-h64, holds the (hotel_id, fingerprint) tuple in an in-memory PoP-local Bloom filter, and once per flush window per (hotel_id, PoP) flushes a single PFADD {hotel:H}:shard:R hash1,hash2,... batch upstream, tagged with per-ASN entry counts for the gateway's global buckets. Collapses each hot (hotel, PoP) pair's writes to ≤1 batched call/s. READ PATH: on cache miss, applies XFetch probabilistic-early-expiry and singleflight coalescing — all simultaneous misses for the same hotel within one PoP collapse to one origin fetch.

Why it exists. Without edge fan-in, viral-hotel write QPS bottoms out the hot Redis slot at V_viral / H = 6.6K writes/s for V=100K — and scales linearly to catastrophe at V=1M. The edge layer absorbs that fan-in: 280 PoPs × 1 batched write/s = 280 batched writes/s on the hot slot regardless of V_viral. We considered (a) sharding the hot key into N parallel HLLs at origin (still does N writes per heartbeat) and (b) client-side sampling (loses precision below 1K viewers). Edge fan-in is the only option that avoids per-heartbeat origin writes AND keeps a flat origin write rate — at the cost of the PoP Bloom's bounded ~1% undercount, still well inside the ±5% SLO (unlike client sampling, which degrades badly below 1K viewers).

When it fails. Cold PoP (newly-spun region, post-deploy bounce): no Bloom state, no DO warmth — every heartbeat goes through unbatched. Origin sees a transient ~24× write spike on hot slots (280 batched writes/s → up to V_viral / H ≈ 6.7K/s) for ~30s. Mitigation: the LB capacity-plans for 3× steady-state and the heartbeat path circuit-breaks at 50K writes/s per upstream LB pool. Health signal: edge_fn_dispatch_p50{pop=$x} plus do_storage_hit_ratio{pop=$x}. We accept that for a region's first 30s, count freshness lags by min(H, batch_window) ≈ 1s — invisible at the 5s edge cache TTL.

Load BalancerEnvoy ring-hash tier (Global Accelerator anycast front, NLB internal)

Anycast-fronted L7 ingress; routes POST /presence/heartbeat/:hotel by consistent-hash(hotel_id) to the heartbeat-service pool, and GET /presence/count/:hotel by least-conn to the count-service pool. Health-checks with /healthz?verify=cluster once per second; ejects unhealthy targets within 10s. TLS 1.3, mTLS to all internal pools. Caps origin connections per upstream pool at 50K writes/s (slot-ceiling head-room) and sheds with 503 + Retry-After jitter when exceeded.

Why it exists. Read traffic is symmetric — least-conn is fine. Write traffic isn't: routing heartbeats by hotel_id to the same heartbeat-service replica means each replica's local hot-key Bloom and per-(hotel, shard) write batcher actually function. Without affinity, the per-replica batchers see fragments of the traffic and you lose the per-replica coalescing on top of edge coalescing. We considered putting the affinity at the gateway (slower, more L7 cost), but the LB is already terminating connections — cheapest place.

When it fails. Single-AZ LB outage: anycast re-converges in ~10s; in-flight requests 5xx. Detection: lb_target_5xx_rate>1% for 30s. Heartbeats are retry-tolerant (next one in 15s); count reads ride SWR=30s. RTO 10s automatic, no data loss (LB is stateless). Cold-PoP affinity collapse (the load-bearing failure mode): when edge-fn is cold in a region, the ~24× spike documented on edge-fn.failureMode (280 batched writes/s → up to V_viral / H ≈ 6.7K/s) collides with consistent-hash(hotel_id) — ALL writes for a single hotel land on ONE hb-svc replica that must briefly absorb that ~6.7K writes/s (vs its ~1.3K/s nominal share). Mitigation is dual: (1) per-target circuit-breaker on e5 lb→gw and on the hb-svc upstream pool ejects the saturated replica within 5s — heartbeats for that hotel decay for ~10s during the gap (visible count drift, NOT a page break); (2) the LB falls back from consistent-hash to least-conn for the affected hotel via the x-fallback: 1 header set by edge-fn when its local Bloom is empty, distributing load at the cost of breaking per-replica batching for that hotel. Pages on-call when lb_per_target_breaker_open_rate>0.5% for 2m.

API GatewayEnvoy + custom WASM filter (fingerprint validation)

On the coalesced write path, per-fingerprint HMAC validation and its 10-heartbeats/sec token bucket already ran at edge-fn (PoP-local, before coalescing), so the gateway consumes the per-ASN entry counts edge-fn tags on each batch and enforces the GLOBAL rate-limits that need cross-PoP coordination — 1000 heartbeats/sec per ASN and per-IP-/24 (a botnet spread across 280 PoPs is invisible to any single PoP's local bucket). Strips PII headers and attaches x-bucket-30s at the gateway clock so client clock skew can't shift the HLL bucket. On the cold-PoP fallback path (edge-fn bypassed, unbatched heartbeats) it does the full per-request check itself: validates the session-fingerprint cookie (HMAC of cookie || IP-prefix-/24 || UA-fingerprint, signed with rotating Ed25519) and computes x-fingerprint-h64 (BLAKE3 of cookie, 64-bit) — the hash, not the cookie, is what gets PFADD'd.

Why it exists. Heartbeat is unauthenticated by design (anonymous browsing of hotel pages must work), so we can't hide behind bearer tokens. The gateway is the only layer that can cheaply do per-(IP, ASN) rate-limit with global coordination — putting it at the edge means each PoP rate-limits independently and a botnet across 280 PoPs gets 280× headroom. We considered moving fingerprint validation into the heartbeat service (saves a hop) but then every replica needs the rotating signing key and the blast radius of a key leak grows from 12 → 64 nodes.

When it fails. Rate-limit Redis (small, separate cluster) failure: the gateway fail-opens — accepts the heartbeat without rate-limit. Detection: gw_rl_redis_timeout_rate>0.5% for 1m. Blast radius is bounded — even an unlimited heartbeat rate is capped per-PoP by the LB connection ceiling (50K writes/s), so a botnet can spike origin briefly but can't melt it. If the *signing key* leaks, we rotate immediately and accept that all sessions need to re-handshake — count drops by ~30% for ~10s as old fingerprints get rejected.

Service · HeartbeatRust + Tokio (async runtime, hot path)

Receives POST /presence/heartbeat already pre-coalesced by the edge layer; looks up hotel shard config in meta-db (cached locally with 60s TTL); computes {hotel:H}:shard:rand(0..N_H-1):bucket:{floor(now/30)} where N_H ∈ {1, 16, 64} per hotel's hot-key tier; PFADD's the fingerprint hash. Emits one event per heartbeat to presence.heartbeats (Kafka) for abuse + analytics. Total path is ≤ 30ms p99; lock-free in the hot path; no synchronous fan-out.

Why it exists. We split heartbeat from count because their failure modes are different: heartbeat loss is invisible (next one in 15s); count read failure is visible (page shows —). Splitting lets us scale them independently — the raw heartbeat aggregate is 333K writes/s, but edge fan-in reduces what hb-svc actually receives to ~83K req/s (5M req/min), which 64 replicas handle at 1.3K req/s/replica; count at 50K reads/s post-edge needs 8. Sharing a service would force the same replica count and saturate the count path during write spikes.

When it fails. Per-replica deadlock or memory leak: with consistent-hash routing, a single bad replica drops 1/64 of heartbeats — a 1.6% precision hit, invisible at HLL's 0.81% baseline error. The replica autoscaler (HPA on hb_writes_inflight) provisions a replacement in ~45s. Worst case is *all* 64 replicas crash simultaneously (deploy bug): heartbeat path dark for ~3min until rollback, counts decay to zero in 30s, page shows —. The page does not break.

Service · CountRust + Tokio + tower-singleflight

On GET /presence/count/:hotel cache-miss from edge, computes the live count via PFCOUNT {hotel:H}:shard:0..N_H-1:bucket:{cur} UNION {hotel:H}:shard:0..N_H-1:bucket:{prev} — two-bucket sliding window. Applies XFetch (Vattani 2015) probabilistic early expiry to refresh hot keys before the edge TTL hits zero. Singleflights all in-flight requests for the same hotel via tower-singleflight so origin sees one PFCOUNT per hotel per second. Publishes count deltas to count-pubsub for the WS push fleet.

Why it exists. PFCOUNT is O(N_shards) — for N=64 it's 64 register reads + 1 union, all within one Redis slot thanks to the hashtag pinning. Doing this on every request from the edge would hammer the slot. We considered (a) caching PFCOUNT result inside count-svc (race-with-write, stale by up to 5s), (b) precomputing on a write-time keyspace-notification listener (extra Redis dependency, complex). XFetch + singleflight inside count-svc is the cleanest: one Redis read per hotel per second, served to all simultaneous requesters from the in-flight future.

When it fails. Redis hot-shard contention spike (viral hotel): PFCOUNT p99 climbs from 8ms → 80ms; the breaker on e8 opens fast and we serve last-known-good from the 1s in-process cache. Detection: count_svc_redis_p99>30ms for 2m (auto-mitigates) → breaker_open_rate{target=redis-hll}>5% for 5m (pages on-call). The known hand-wave: tower-singleflight is per-REPLICA, not cluster-wide — across 8 replicas during rollout or least-conn fan-out, a hot hotel can land 8 concurrent PFCOUNTs on the slot. We accept that 8× burst as part of the redis-hll mitigation tree (the slot's headroom is sized for it: 50K writes/s ceiling vs 280 batched origin writes/s steady-state) rather than push singleflight into Redis via Lua. Worst case all PFCOUNT calls timeout: widget renders stale count with ?approximate=true&stale=true; page never blanks.

Service · WS/SSE FanoutGo + nhooyr/websocket, sticky-hash routing

Holds long-lived WebSocket connections from clients (SSE fallback for restrictive networks). Subscribes to count-pubsub channels for hotels with active subscribers. Coalesces broadcasts per Pusher's published rules: hotels with ≤100 subscribers broadcast every count change; hotels >100 broadcast at most once per 5s with bundled deltas (matches Pusher's published 5s lock + 30s bundle window). Routes incoming subscriptions by consistent-hash(hotel_id) so all subscribers for the same hotel land on the same replica — Discord's 'guild → gateway shard' pattern translated to hotels.

Why it exists. Pure polling at 5M concurrent pages × 1 read / 10s = 500K reads/s; even after 90% edge absorption, that's 50K reads/s flowing through count-svc constantly, including for hotels nobody is watching. WS fanout flips the model: a hotel with active viewers gets one push per ~5s regardless of subscriber count. We considered SSE-only (one-way, no reconnect handshake); we kept WS because we need the client→server channel for keep-alive frames that prove connection liveness — these carry no presence; presence is written solely by the separate POST /presence/heartbeat on its own path, keeping the read (push) and write (heartbeat) planes independent.

When it fails. Reconnect storm after rolling deploy (Slack 2018-06-27 pattern): 5M clients reconnect within 30s after fleet bounce. Detection: ws_fleet_handshake_rate > 5× baseline. Mitigation: server sends Reconnect-After: $jittered_seconds in the close frame (decorrelated jitter, AWS Builders' Library); rolling deploy paces ≤ 5% of fleet at a time; LB sheds new connections at 1.2× steady-state. Cert expiry on the WS edge: hard outage, all clients fall back to 10s polling. We have cert-expiry alerts at T-30/14/7/1 days (post-Spotify 2022).

Cache · Redis HLLRedis 7.4 Cluster, m7g.4xlarge × 16

Stores per-hotel HyperLogLog state at precision 14 (16384 registers, 12 KB/key, 0.81% standard error per Heule et al. 2013). Two buckets per hotel for the sliding window (bucket:{floor(t/30)}, bucket:{floor(t/30)-1}). For viral hotels, splits each bucket across N ∈ {16, 64} parallel HLLs co-located in one slot via the {hotel:H} hashtag — this lets PFCOUNT/PFMERGE remain atomic and intra-slot, while spreading PFADD writes across N register-arrays' CPU. TTL = 60s on every key.

Why it exists. HLL gives us the cardinality of a viral hotel for 12 KB regardless of viewer count — 100K viewers and 100M viewers cost the same memory. The 0.81% standard error is well below the product's ±5% SLO. We rejected ZSET (exact, but per-element 80 B → 4 MB at V=50K with O(log N) ZADD on the hot slot) and Memcached CAS counters (no probabilistic primitive, would need app-level HLL outside the store). Redis HLL is in-engine: PFADD is one O(1) register update.

When it fails. Hot-shard CPU pegged on viral hotel (Booking.com-style spike): one slot at 100% while cluster sits at 5%. Detection: redis_cpu_seconds{slot=$hot}>0.9 while cluster_cpu<0.10 for 30s — drives hotkey-controller auto-promotion. Two-stage mitigation: (1) controller bumps N_shards via meta-db (count-svc picks up first, hb-svc 90s later — see keyChoices); (2) edge fan-in cadence drops 1s → 250ms. CLUSTER REBALANCE on the hot slot: PFADD/PFCOUNT return ASK redirects during CLUSTER SETSLOT MIGRATING/IMPORTING; clients follow and counts continue, but the 'atomic intra-slot PFMERGE' guarantee briefly breaks. We gate slot-migration on hot-tier hotels via a slot-pin policy in hotkey-controller and a runbook that disables auto-rebalance for slots holding any V≥10K hotel — migration becomes operator-driven planned-maintenance. Worst case slot loss: leader-failover RTO=15-30s during which PFADDs 503; widget shows last-cached count for the gap, then resumes. RPO=5s of heartbeats lost — invisible because next heartbeat is in 15s and HLL is idempotent on retry. Pages on-call when failover RTO > 60s OR when hot-shard CPU>0.9 sustained 3m without controller mitigation; otherwise auto-mitigates.

Stream · Count Pub/SubRedis 7.0 Sharded Pub/Sub (SSUBSCRIBE)

Fast-path broadcast channel for count deltas. Count-svc SPUBLISH count.{hotel:H} <count> on every recompute; the WS fanout fleet SSUBSCRIBEs only for hotels with active local subscribers. Sharded Pub/Sub (Redis 7.0+) routes by hashtag so subscribers connect to the slot owner directly — no proxy hop. Zero retention; subscribers that disconnect don't replay (this is fine: the next count read or next push will catch them up).

Why it exists. We need a fan-out channel for count → ws-fanout that does not couple their lifecycles. Without pub/sub, count-svc would need a registry of active subscribers per hotel (load-bearing state), or ws-fanout would poll count-svc per hotel per cycle (rebuilds the read fanout problem we just solved). Redis Sharded Pub/Sub piggy-backs on the cluster we already operate, with one hop. We rejected Kafka here — Kafka topics with retention=0 are wasteful, and Kafka's per-topic overhead at 10M hotel-channels would be 10M consumer-group offsets, untenable.

When it fails. Subscriber-side lag (ws-fanout replica GC pause): SUBSCRIBE buffers fill, Redis evicts the slow subscriber per client-output-buffer-limit pubsub 32mb 8mb 60. Detection: pubsub_dropped_clients>0. Mitigation: ws-fanout replicas auto-resubscribe and the next count tick pushes again — at most 5s of staleness. Pub/Sub server outage: ws-fanout falls back to polling count-svc directly at 5s/hotel. Capacity headroom on count-svc is 5×, so this is survivable for the duration of the failover.

Stream · Heartbeat EventsApache Kafka 3.7, MSK m7g.2xlarge × 9

Persistent event log of every heartbeat ({hotel_id, fingerprint_h64, asn, geo, ts}) — *not* the HLL writes, just the metadata for off-line scoring. Kafka topic presence.heartbeats, 64 partitions keyed by hotel_id, RF=3, acks=all on producer, 72h retention. Feeds the abuse worker (real-time scoring), the analytics warehouse (Snowflake nightly batch), and a synthetic-traffic detector. Producer is fire-and-forget — hb-svc returns 204 to the client before the Kafka ack.

Why it exists. We need durable, replayable, ordered metadata for abuse detection that the volatile HLL Redis intentionally throws away. Putting abuse scoring in the hot path (hb-svc → meta-db sync write) would double heartbeat p99 and couple the count's availability to the metadata DB's. The stream gives us 24h of replay if a scoring rule changes, without re-processing live traffic. We rejected (a) writing scoring features to meta-db on each heartbeat (sync coupling) and (b) sampling — sampling misses the rare attacker patterns we need to catch.

When it fails. Broker outage (single-AZ): producer retries to in-sync replicas, no client-visible impact (acks=all + min.isr=2 survives 1 broker loss). Consumer-side lag (abuse worker falls behind): degrades scoring freshness, doesn't break heartbeats. Detection: kafka_consumer_lag{topic=presence.heartbeats}>1M for 5m. Worst case complete topic loss: scoring degrades but heartbeats and counts work normally. The cardinal contract: this stream's failure must NOT be visible on the page.

Worker · Abuse ScannerGo + sliding-window counters in BadgerDB

Consumes presence.heartbeats partitioned by hotel_id; maintains per-(fingerprint, ASN, hotel) sliding-window counters in embedded BadgerDB; scores patterns (heartbeat-rate spikes from unseen ASNs, fingerprint reuse across geo-impossible IPs, UA fingerprints with rotating cookies) using a hand-tuned rule set + a downstream gradient-boosted model that gets updated weekly. On a positive verdict, writes (fingerprint, expires_at, reason) to meta-db's abuse_denylist table — hb-svc reads this with its 60s cache.

Why it exists. Bot armies inflate counts (or deflate competitor hotels) — we need a feedback loop that makes the heartbeat path hostile to attackers without making it expensive for honest users. Putting scoring in the hot path costs every legitimate user a model inference (~1ms) for the rare attack signal. The async stream lets us spend ~250ms per heartbeat on scoring without anybody noticing. Hot-path denylist lookup (the cheap part) stays in the gateway via meta-db cache.

When it fails. Worker crash loop on poison message (e.g. malformed UA from a buggy SDK): one partition stalls, consumer-group lag grows. Detection: kafka_consumer_lag{partition=$x}>1M for 10m. Mitigation: DLQ pattern — after 3 retries, route to presence.heartbeats.dlq and continue. Worst case all 12 replicas stall: scoring goes silent for hours, attackers get a free window. The page does not break — denylist lookups still hit the last known state in meta-db. We page on-call when consumer lag exceeds 30 min.

SQL DB · MetadataPostgres 16 HA (Patroni + 2 sync followers)

Authoritative source for low-write, high-read metadata: hotel registry (hotel_id, name, region), per-hotel hot-key config (n_shards_read / n_shards_write ∈ {1, 16, 64}, last-promoted ts), abuse denylist (fingerprint_h64, expires_at, reason). Read by hb-svc on every heartbeat (cached 60s), by gateway for fingerprint-revocation checks (cached 5s), by abuse worker for scoring rules. Write paths: a control-plane operator service auto-promotes hotels between hot-key tiers based on observed write QPS (metric → control loop → 2-phase UPDATE hotel_shards).

Why it exists. Some metadata genuinely needs durability — abuse denylist must survive Redis bounces, hot-key tier decisions need a record so hb-svc and count-svc agree on N_shards, hotel registry is shared with the broader booking platform. Postgres for the same reason every team picks it: SQL escape hatch for ad-hoc operator queries (SELECT hotel_id, count(*) FROM abuse_denylist WHERE created_at > now()-'1 hour'::interval), and we already operate it. We rejected DynamoDB (rigid query patterns) and putting this in Redis (no durability semantic for abuse denylist).

When it fails. Primary failover (Patroni-orchestrated): writes block for ~30s; hb-svc and gateway continue reading from local 60s caches — no user-visible impact. Detection: pg_replication_lag>5s or patroni_state!=running for 1m. Worst case primary + both sync followers down: control-plane writes block until manual promotion; hot-key tier promotions and abuse denylist updates pause; hb-svc and gateway run on stale config for hours, which is fine — counts stay accurate, attackers caught yesterday remain caught.

Cache · Rate-Limit RedisRedis 7.4 single primary + 1 follower (NOT clustered)

Backs Envoy's ratelimit-service. Holds the short-TTL GLOBAL token-bucket counters the gateway coordinates — keyed by (IPv4-/24) and (ASN) — two bucket queries per gateway check (once per batch on the coalesced path; the per-fingerprint bucket lives PoP-local at edge-fn, not here). Returns ALLOW/DENY in <2ms p99 over a connection-pooled TCP path; gateway fail-opens (ALLOW) on timeout >5ms.

Why it exists. Putting rate-limit state in the same Redis cluster as the HLL data means a hot-shard meltdown on a viral hotel's HLL slot also melts rate-limit reads for every other hotel — one viral incident knocks out global abuse defense. Splitting blast radius is the load-bearing reason this is a separate cluster, even though it costs operating two Redises. Single-instance + 1 follower (not clustered) because the working set (per-/24 and per-ASN buckets with 60s TTL, ~100MB) fits in one node and clustering would just add slot-routing cost.

When it fails. RL Redis failure: gateway fail-opens, accepts every heartbeat without rate-limit. Detection: gw_rl_redis_timeout_rate>0.5% for 1m. Blast radius bounded — even unlimited rate, the LB connection ceiling (50K writes/s per upstream pool) caps per-pool damage; HLL itself is idempotent on the same fingerprint. Pages on-call when gw_rl_redis_timeout_rate>5% for 5m (sustained fail-open is bot-army-friendly and we want a human to look). Single-AZ failure: Sentinel promotes the follower in 10s; brief <1s write-stall at the moment of cutover during which a few heartbeats fail-open.

Service · Hot-Key ControllerGo + leader-elected via Patroni advisory lock on meta-db

Continuously samples per-slot write QPS and CPU from obs (PromQL: sum by (slot) (rate(redis_pfadd_total[30s])) + redis_cpu_seconds); when a hotel's slot exceeds the (N_shards, V_viewers) threshold table — 1→16 at V≥1K OR pfadd_qps>1K, 16→64 at V≥10K OR pfadd_qps>10K — issues the 2-phase tier promotion: STEP 1 UPDATE hotel_shards SET n_shards_read=$new (count-svc snapshot picks up at next 60s rotation), wait 90s, STEP 2 UPDATE hotel_shards SET n_shards_write=$new. Demotes (64→16→1) symmetrically after V drops below tier-floor for >10 min. Also enforces the slot-pinning runbook for V≥10K hotels (refuses a manual CLUSTER REBALANCE request on those slots).

Why it exists. The auto-promotion N=1→16→64 is the load-bearing safety story for hot keys, and we need a service that *runs* it — not implicit logic on hb-svc (each hb-svc replica would race) and not a manual operator runbook (60s reaction time is the difference between healthy auto-promotion and a melted slot). Leader-elected so only one replica drives the control loop at a time; the 2nd replica is hot-standby with a 5s lock-acquisition timeout. The promotion threshold V=10K (not V=50K) gives the metric-loop budget headroom: slot saturates in 5s at V_viral, our scrape interval is 15s — promoting at V=10K leaves 4× margin while the slow loop catches up.

When it fails. Both replicas crash: promotion stops; hot keys saturate without intervention. Detection: hotkey_controller_leader_age_seconds > 60. Pages on-call within 60s — this is *the* control-plane that prevents melt-downs and a human MUST take over with manual UPDATE hotel_shards. Replication-lag > 30s on the meta-db read replica controller uses for its sample: controller refuses to issue promotion (stale data is worse than no promotion); pages immediately. Detection: hotkey_controller_data_age_seconds > 30.

Monitoring · ObservabilityPrometheus 2.x + Thanos + Grafana + Alertmanager + PagerDuty

Scrapes /metrics from every service/cache/stream/db node every 15s; long-term storage in Thanos with 30-day retention. Hosts the SLO dashboards (count-read p99, heartbeat-write success-rate, hot-shard CPU heatmap, ws-fanout connection count, abuse consumer lag). Drives Alertmanager → PagerDuty for the page-on-call thresholds named in each node's failureMode. Exposes a query API the hotkey-controller consumes for its 15s control loop.

Why it exists. Every failureMode in this canonical names a metric (cdn_5xx_rate, redis_cpu_seconds, count_svc_redis_p99, kafka_consumer_lag, gw_rl_redis_timeout_rate, pg_replication_lag, ws_fleet_handshake_rate, lb_per_target_breaker_open_rate). These metrics drive auto-mitigation (controller-driven shard promotion, breaker decisions, autoscale signals) — they are not optional decoration. Putting obs on the diagram makes the metric pipeline drawable; without it, the chaos scenarios that say 'Detection: $metric' have no graph target. Prometheus over Datadog because we already operate it; Thanos for cross-region federation when we add multi-region.

When it fails. Prometheus replica failure: 6-replica federation; Thanos sidecar fronts queries — single replica loss is invisible. Quorum loss (4/6 replicas down): metrics gap until quorum recovers; Alertmanager flushes pending alerts on recovery. Detection: prometheus_up < 4 for 1m. Pages immediately because we lose detection on every other failure. Worst case complete obs outage: the controller can't promote (refuses on stale data, see hotkey-controller); manual operator promotion via runbook; slot melts only if a viral hotel happens to spike during the obs gap, which is the genuine 3am page.

Service · Peer Region CountStub representing a peer-region count-svc (full mirror of this canonical)

Represents one peer region's count-svc for the cross-region count-merge story. Local count-svc fetches GET /presence/count-merge-bundle/:region from each peer region every 1s for hotels with cross-region viewers, PFMERGE's the snapshots, returns the union count. Stub-modeled here as one 'peer' node so the multi-region chaos scenario (region-partition-split-brain) has a graph target — in reality this is a full mirror of the entire canonical in another region.

Why it exists. The non-functional spec claims active-active per-region with ≤1s cross-region staleness. A single-region canonical can't defend that. Stubbing a peer-region node is the minimum graph representation needed to: (a) draw the cross-region merge edge that count-svc consumes, (b) let the chaos region-partition-split-brain exercise a real path. We considered mirroring the entire canonical for each region (infeasible — 17 nodes × N regions = unreadable diagram) vs downgrading to single-region (loses the active-active claim). The stub is the honest middle.

When it fails. Cross-region link partition: peer-region request times out at 1000ms; count-svc returns local-region-only count (visible drift in 'global' display proportional to peer-region's share of viewers). Detection: cross_region_merge_timeout_rate>10% for 1m. Pages on-call when sustained 5m — most cross-region partitions are <1m blips, and we accept regional-only display for the duration. Split-brain on heal: the redis-hll PFMERGE is associative and commutative — counts converge upward to UNION on next merge cycle, no manual reconciliation needed.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Worth surfacing:

  • What does "viewing" mean? Page loaded? Tab focused? Scroll activity? Defines the heartbeat semantics. We track page loaded + visible (Page Visibility API gates the heartbeat) — minimized tabs stop heartbeating after 30s, count expires naturally.
  • How fresh? Update once per ~5–10 s. Per-second freshness is a different design — the edge cache TTL alone forces ≥ T_edge staleness.
  • Exact or approximate? "127" vs "100+" — approximate unlocks much cheaper designs (HLL, sampled counters). We commit to ±5 % tolerance; HLL p=14 gives 0.81 % standard error, well inside.
  • Per-property hot keys. Booking has historically reported (QCon 2018 talks) that ~5 % of properties carry ~80 % of pageviews; viral hotels (concert venues during the Eras Tour, Olympics host city listings) hit 100 K concurrent viewers on a single ID.
  • Bots / abuse? A scraper hitting refresh shouldn't pump the count. HLL deduplicates by id; bots need rotating fingerprints to inflate. We score off the heartbeat metadata stream, not in the hot path.
  • Push vs poll? Push for hotels with active subscribers (WebSocket / SSE); poll-with-edge-cache for cold pages.

Assumptions:

  • 5 M concurrent active hotel pages globally during peak.
  • Avg hotel has ~5 viewers; popular hotel has ~10 K; viral has 100 K.
  • 30 s TTL on "viewing" status (= 2 × heartbeat). Heartbeat every 15 s while page is visible.
  • Approximate count is acceptable (±5 %).
  • Single-region active-active (per-region Redis state, lazy cross-region merge for the global view — 1 s staleness acceptable).

02Functional reqs

What must this system actually do?

  • Mark a user as "currently viewing" hotel H when they load the page.
  • Heartbeat to keep the marker alive while they remain.
  • Remove the marker when they leave / disconnect / time out.
  • Return the current viewer count for hotel H.
  • Stream the count live to the open page (push channel for hot hotels, poll for cold).
  • Rate-limit and fingerprint-dedupe heartbeats to defend against count inflation.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Latency: Heartbeat write p50 < 30 ms, p99 < 100 ms, p99.9 < 250 ms. Count read p50 < 20 ms (CDN hit), p99 < 60 ms (origin), p99.9 < 200 ms (hot-shard contention). Push fan-out e2e ≤ 5 s.
  • Freshness: Count reflects activity within the last 30 s. Edge cache adds up to 5 s staleness (with SWR up to 30 s on origin error).
  • Consistency: Strongly inexact-but-eventually-correct. ±5 % is fine; 0.81 % standard error from HLL p=14 is well inside.
  • Availability: Count failure must NOT break the hotel page. Render — and move on. The widget is decorative; the booking flow is the product.
  • Cost: Viewing-count must be near-free per request. Edge absorbs ≥ 90 % of reads; per-key origin reads bounded at 1/T_edge regardless of viewer count.
  • Durability: Presence is ephemeral — TTL is the source of truth. Never persist live presence to a durable store. Persist heartbeat metadata (fingerprint, ASN, ts) to the Kafka stream for abuse only.
  • Security: Fingerprint cookies bound by HMAC of cookie || /24 IP || UA hash. Per-IP and per-ASN rate limits at the gateway. Async fraud scoring writes to denylist.

04Capacity estimation

How much load and data does this have to hold?

Inputs: C = 5 M concurrent pages, H = 15 s heartbeat, P = 10 s poll, T_edge = 5 s, V_viral = 100 K, N_hotels_active = 250 K (5 % of 5 M).

Aggregate workload:

  • Heartbeat writes: C / H = 5M / 15 = 333K writes/s. Tier-4 (50 M concurrent) → 3.33 M writes/s — at that scale we'd add a tertiary edge-aggregation layer.
  • Count reads pre-edge: C / P = 500K reads/s.
  • Origin count reads post-edge: N_hotels_active / T_edge = 250K / 5 = 50K reads/s. Edge absorbs ≈ 90 %.

Per-key load (the bottleneck):

  • Viral hotel write QPS: V_viral / H = 100K / 15 = 6.7K writes/s on ONE Redis slot. Redis primary tops out around 100–150 K ops/s aggregate, slot-CPU hits the wall ~50 K writes/s — V=100 K is comfortable, V=750 K is the cliff. Edge fan-in at 280 PoPs cuts this to ~280 batched writes/s on the slot (≈24×).
  • Viral hotel reads at origin: 1 / T_edge = 0.2 reads/s per hotel. Bounded by TTL, NOT by viewer count. Reads are a solved problem; writes are the thing that breaks.

Storage:

  • HLL p=14 = 12 KB / register array. 2 buckets (sliding window) × N_shards (avg 4 across all hotels, weighted) ≈ 96 KB / hotel.
  • 10 M hotels-with-recent-activity × 96 KB ≈ 960 GB Redis footprint. Ridiculous? No — sized for HLL keys with 60s TTL, the working set is more like the 250 K hotels currently active = 24 GB. Fits in 16 × m7g.4xlarge nodes (16 GB usable / node post AOF buffers).

Egress:

  • 500 K reads/s × 100 B response = 50 MB/s aggregate; CDN absorbs 90 % → origin egress ≈ 5 MB/s.

Push channel connections:

  • 5 M concurrent → 50–100 K WS / replica realistic ceiling under fanout load → 50–100 ws-fanout replicas. Plan 96 nodes baseline + autoscale headroom 2×.

05API design

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

POST /api/v1/presence/heartbeat
Cookie: sid=<fingerprinted-session>
Content-Type: application/json

{
  "hotelId": "h-abc123",
  "bucket": 1739283150          # client-computed floor(now/30); server validates
}

204 No Content
Cache-Control: no-store
GET /api/v1/presence/count/:hotelId
200 { "count": 127, "approximate": true, "ts": "2026-05-08T12:34:50Z" }
Cache-Control: public, max-age=5, stale-while-revalidate=30, stale-if-error=300

WebSocket / SSE for active pages:

WS /ws/presence/hotel/:hotelId
→ server pushes { "count": 127, "ts": "..." } every ~5 s when count changes
→ client also sends a WS keep-alive ping every 15 s (connection liveness only — NOT the presence heartbeat; presence is written solely by the separate `POST /presence/heartbeat` on its own path)

Why short edge cache + SWR? With T_edge=5 s, every PoP fetches origin once per hotel per 5 s — that's the load-bearing arithmetic that makes per-key origin reads bounded at 1/T regardless of V_viral. SWR=30 s and stale-if-error=300 s cover origin blips without ever blanking the widget.

Why not 301-style permanent cache? Counts must update — Cache-Control: public, max-age=5 is the sweet spot.

06Data model

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

Two designs depending on exactness — we pick HLL.

A. Set-based (exact, expensive on hot keys): Redis ZSET per hotel, score = expiry_ts, member = fingerprint_h64. Heartbeat = ZADD key now+30 fp; read = ZRANGEBYSCORE key now +inf count; sweep lazily on read.

  • Pros: exact.
  • Cons: 80 B/member; at V=50 K → 4 MB key with O(log N) ZADD on the hot slot. Slot saturates above ~50 K writes/s.

B. HLL-based (approximate, hot-key-friendly): (chosen) Per-hotel HLL at precision 14 (16 384 registers, 12 KB, 0.81 % stderr). Two buckets per hotel for sliding window: {hotel:H}:bucket:{floor(t/30)} and {hotel:H}:bucket:{floor(t/30)-1}. Heartbeat = PFADD {hotel:H}:bucket:{cur}:shard:{rand} fp; read = PFCOUNT bucket:{cur}:shard:0..N-1 bucket:{prev}:shard:0..N-1 (intra-slot via hashtag).

  • Pros: O(12 KB) regardless of viewer count; PFADD O(1).
  • Cons: approximate (well inside 5 % SLO); can't list viewers.

Hot-key tiering (config in meta-db, snapshotted to hb-svc memory):

  • Cold (V < 1 K): N_shards = 1.
  • Hot (1 K ≤ V < 10 K): N_shards = 16.
  • Viral (V ≥ 10 K): N_shards = 64. Auto-promoted pre-emptively by control-plane based on observed viewers / write QPS on the slot — promoting at V = 10 K keeps 5× headroom under the ~50 K writes/s slot ceiling; reacting at V = 50 K would race saturation.

meta-db schema (Postgres):

tablekey columnsnotes
hotelhotel_id (PK)name, region, created_at
hotel_shardshotel_id (PK)n_shards_read int, n_shards_write int, last_promoted_at, reason text
abuse_denylistfingerprint_h64 (PK)expires_at, asn, reason

Composite index: abuse_denylist (fingerprint_h64, expires_at); the hot lookup keeps WHERE expires_at > now() in the query, not the index predicate (a partial-index predicate can't call now() — it is STABLE, not IMMUTABLE), with a periodic job purging expired rows. Snapshot of hotel_shards rotates into hb-svc memory every 60 s.

07High-level design

Which components handle a request, and in what order?

Architecture summary (the load-bearing claims, in dependency order):

  1. Client (page). Heartbeat every 15 s while visible; subscribe via WS for live count, fall back to 10 s polling.
  2. CDN (Fastly). 5 s TTL + 30 s SWR on count GET. Absorbs ≥ 90 % of read traffic. WAF and per-IP/ASN rate limit on heartbeat POST. Origin shield centralises misses to one PoP / region.
  3. Edge Aggregator (Cloudflare Workers + Durable Objects). Per-PoP heartbeat coalescer: PoP-local Bloom + DO-per-(hotel, PoP), flushes one batched PFADD per second upstream. This is the load-bearing layer that makes the hot key tractable — without it, viral hotels exceed the slot ceiling at V_viral = 750 K. With it, origin sees 280 PoPs × 1 batch/s = 280 batches/s on the hot slot regardless of V_viral (far under the ~50 K/s slot ceiling).
  4. Load Balancer. Anycast L7; consistent-hash on hotel_id for write-affinity to heartbeat-svc replicas (so per-replica batchers see coherent traffic per hotel); least-conn for read traffic.
  5. API Gateway (Envoy + WASM). Global per-IP/ASN rate-limit (per-fingerprint validation + x-fingerprint-h64 run at edge-fn before coalescing; the gateway re-does them only on the cold-PoP fallback path), attaches x-bucket-30s. Rate-limit Redis is a separate cluster from the HLL Redis (different blast radius).
  6. Heartbeat Service (Rust). Stateless. Per-replica per-(hotel, shard) batcher with 100 ms flush; PFADD batched into one Redis command of up to 256 args. Reads hot-key tier from local snapshot of meta-db.
  7. Count Service (Rust + tower-singleflight). PFCOUNT 2 buckets × N shards (within slot via hashtag); XFetch β=1.0 probabilistic early expiry; singleflight coalesces concurrent misses; in-process cache 1 s; publishes count delta to pub/sub on every recompute.
  8. WS/SSE Fanout (Go). Long-lived client connections; subscribes to count pub/sub by hotel_id (sticky-hashed); coalesces broadcasts (≤100 subs: every change; >100: ≤1/5 s with bundled deltas, per Pusher); SSE fallback for restrictive networks.
  9. Redis HLL Cluster (Redis 7.4). 16 nodes m7g.4xlarge, p=14 HLL, 2-bucket sliding window, hashtag-pinned sharded HLL per hotel-tier, ttl-only eviction, 60 s TTL, async replication (1 follower per primary), cluster auto-failover RTO=15-30 s.
  10. Count Pub/Sub (Redis 7.0 Sharded Pub/Sub). SPUBLISH/SSUBSCRIBE on count.{hotel:H} channels — sharded routes to slot owners (no N×N cluster broadcast). Zero retention.
  11. Heartbeat Stream (Kafka). Durable async log of heartbeat metadata for abuse + analytics. Producer fire-and-forget; never blocks client. RF=3, acks=all, 64 partitions keyed by hotel_id, 72 h retention. Poison messages route to presence.heartbeats.dlq (separate topic, RF=3, retention=7 d) after 3 retries.
  12. Abuse Worker (Go). Stream consumer; per-(fingerprint, ASN, hotel) sliding-window scoring; writes denylist to meta-db with idempotent UPSERT semantic; offset commit AFTER meta-db write returns 200 → effectively-exactly-once.
  13. Metadata DB (Postgres 16 HA). Hotel registry, hot-key tier config, abuse denylist. Patroni HA, sync-replicated, RTO=30 s, RPO=0. Off the heartbeat hot path via 60 s in-memory snapshot in hb-svc.
  14. Rate-Limit Redis (separate cluster). Token-bucket counters keyed by (fingerprint, /24, ASN) for the gateway. Single-instance + 1 follower (NOT clustered, ~200 MB working set). Separate from HLL Redis specifically so a hot-shard meltdown on a viral hotel cannot also melt the abuse-defense layer — different blast radii.
  15. Hot-Key Controller (leader-elected service). The control loop that drives auto-promotion N=1 → 16 → 64 based on observed slot QPS / CPU from obs. 2-phase promotion: count-svc first, hb-svc 90 s later — without this ordering, mid-promotion the count widget visibly drops 75 % on a 16→64 step. Leader-elected via Postgres advisory lock (no extra infra).
  16. Observability (Prometheus + Thanos + Grafana + Alertmanager). Scrapes every node every 15 s. Drives the hotkey-controller control loop (PromQL on slot_cpu_p99_window_30s + pfadd_qps) and pages on-call per the thresholds named in each node's failureMode. Without this node, the metric pipeline is invisible and the auto-mitigation claims are aspirational.
  17. Peer Region Count (stub). Represents one peer region's count-svc for the cross-region merge story. count-svc fetches each peer region's PFMERGE-bundle async every 1 s and unions; on partition, returns local-only count with ?regional=true.

Data flow on count read (warm hotel, edge hit): Client → CDN (cached). Done. ≈ 20 ms p50 e2e.

Data flow on count read (cold / TTL-expired): Client → CDN (miss) → Edge-fn (singleflight) → LB → Gateway → Count-svc → Redis (PFCOUNT 2 buckets × N shards). p99 ≈ 60 ms.

Data flow on heartbeat (warm PoP): Client → CDN → Edge-fn (fingerprint validate + PoP-local Bloom + DO batcher; coalesces to 1 upstream call per flush window per (hotel, PoP)) → LB (consistent-hash hotel_id) → Gateway (global per-ASN rate-limit) → Heartbeat-svc (batch + PFADD). p99 client-observed ≈ 100 ms.

Data flow on push: Heartbeat lands in HLL → Count-svc XFetch eventually recomputes → SPUBLISH count.{hotel:H} → WS-fanout SSUBSCRIBEs → coalesced broadcast frame to all subscribers within 5 s.

Data flow on abuse: Heartbeat → hb-svc → Kafka (presence.heartbeats) async → abuse-worker (BadgerDB-backed scoring) → meta-db UPDATE abuse_denylist → gateway re-reads on its 5 s cache cycle. End-to-end: 1–5 minutes from attack signal to denied heartbeats.

08Deep dives

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

1. Hot-key on a viral hotel. Even with HLL, 6.7 K writes/s on one key can saturate a Redis node. Three layers, engaged tier by tier:

  • Edge fan-in (always). Per-PoP DO batches all heartbeats for (hotel, PoP) into one upstream call per 1 s flush window. 280 PoPs × 1 batch/s = 280 origin writes/s on the hot slot regardless of V_viral.
  • Sharded HLL (N=16 at V ≥ 1 K, N=64 at V ≥ 10 K). N parallel HLLs per hotel co-located in one slot via hashtag {hotel:H}. Atomic PFMERGE on read. Spreads register-array CPU within slot. Limit: still one slot.
  • Cross-slot key sharding (V ≥ 750 K). Drop the hashtag, distribute shards across slots, accept multi-slot PFMERGE (cluster-aware client gathers and unions in count-svc). Loses atomicity; gain horizontal scale.

2. Push vs poll for the count update channel. Polling alone is wasteful (5 M reads/s for an adornment). WebSocket is leaner but adds sticky-connection state. We use both: WS for hotels with active subscribers (the long-tail hot pages), polling-with-5 s-edge-cache for cold pages. The WS channel carries only count pushes down and connection keep-alive pings up — presence is written exclusively by the separate POST /presence/heartbeat, which keeps the read (push/poll) and write (heartbeat) planes independent: the WS path terminates at ws-fanout and never touches the heartbeat write path. SSE fallback for restrictive networks.

3. The 30-second bucket cliff. Single tumbling 30 s bucket: at wall-clock second 30 the count drops 127 → 0 visibly on every connected client at once. Two-bucket sliding window costs 2× memory (24 KB / hotel base + 24 KB × N_shards on hot) and one extra union op per read. Worth every byte.

4. What if the user closes the tab? Heartbeat stops; TTL expires entry in 30 s. Don't try to detect tab close — beforeunload is unreliable on mobile Safari especially. The TTL is the source of truth.

5. Cache stampede on TTL expiry. With T_edge=5 s, every PoP cache for the hot hotel expires within the same wall-clock second. Origin sees 50 K simultaneous misses on one key. Three layers (Vattani 2015):

  • TTL jitter ±20 % (decorrelates expiry).
  • Singleflight at edge-fn (coalesces in-PoP misses).
  • XFetch (Vattani 2015) at count-svc (probabilistic early expiry — recomputes with rising probability as TTL nears zero).

6. Bot resistance. Fingerprint = HMAC(cookie || IP-/24 || UA hash, rotating Ed25519). HLL deduplicates same-fingerprint heartbeats naturally — bots need rotating fingerprints to inflate. Per-ASN rate limit at gateway catches the cheap attacks (single ASN's botnets); async fraud scoring catches the expensive ones (rotating IPs, residential proxies).

7. Multi-region. Each region runs its own Redis HLL state. The "global" count for a hotel is per-region merge: count-svc fetches each region's PFCOUNT, sums (or PFMERGE'd snapshots) async every 1 s, returns. Stale by ≤ 1 s — acceptable. Cross-region failover for a region outage falls back to remaining regions' counts; numbers drop 1/N visibly during the outage, page does not break.

09Trade-offs

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

Trace catalogue (each authored against this canonical's hldNodes/hldEdges; verified by findPath):

trace idjourneybudget
concurrent-hotel-viewers:count-edge-hitRead · CDN cache hit (90 %+ of reads)30 ms
concurrent-hotel-viewers:count-edge-miss-originRead · CDN miss → Count-svc → Redis PFCOUNT (cold path)80 ms
concurrent-hotel-viewers:count-viral-hot-shardRead · viral hotel · sharded HLL PFMERGE on hot slot150 ms
concurrent-hotel-viewers:heartbeat-write-coalescedWrite · normal path · edge fan-in PoP coalescer100 ms
concurrent-hotel-viewers:heartbeat-write-uncoalescedWrite · cold-PoP fallback · no edge coalescing150 ms
concurrent-hotel-viewers:push-count-fanoutAsync · count delta → pub/sub → WS push5000 ms
concurrent-hotel-viewers:abuse-fingerprint-scanAsync · heartbeat metadata → Kafka → abuse worker → DB30000 ms
concurrent-hotel-viewers:ws-subscribeControl · client opens WebSocket subscription1500 ms

Failure scenarios we model (each grounded in a real-world precedent — ids in <slug>:<scenario> form):

chaos idcategoryprecedent
concurrent-hotel-viewers:viral-hot-shard-meltdowntrafficBooking.com holiday-weekend single-property spike (10–50×)
concurrent-hotel-viewers:cdn-pop-blackoutdepsFastly 2021-06-08 49 min global outage
concurrent-hotel-viewers:redis-leader-failoverdataAWS ElastiCache Multi-AZ failover, RTO 30–120 s
concurrent-hotel-viewers:count-stampede-on-ttltrafficMailchimp dogpile / Vattani 2015 XFetch paper
concurrent-hotel-viewers:ws-fleet-reconnect-stormprocessSlack 2018-06-27 reconnect storm + Discord gateway bounces
concurrent-hotel-viewers:abuse-consumer-backlogprocessGeneric Kafka consumer-lag pattern, SRE Workbook Ch. 22
concurrent-hotel-viewers:cert-expiry-ws-edgedepsSpotify 2022-05-15 root-cert expiry / Microsoft Teams 2020
concurrent-hotel-viewers:region-partition-split-braininfraMulti-region active-active divergence on cross-region link

What breaks at 10× scale (50M concurrent active pages):

  • Redis Cluster footprint. 24 GB working set → 240 GB; need 64-node cluster or move to Redis Enterprise / KeyDB Pro.
  • Per-region throughput. 333 K writes/s → 3.3 M; even with edge fan-in, the tertiary aggregation tier (regional rollup before origin) becomes mandatory.
  • Hot-key controller scrape latency. At 10× scale, the 15 s scrape interval becomes the bottleneck — V_viral can spike past the slot ceiling between two scrapes. Move the controller to a push-based signal from hb-svc (pfadd_qps_per_hotel reported every 5 s).
  • WebSocket connection count. 50 M live connections × 100 K/replica = 500 ws-fanout nodes. Operationally heavy — consider falling back to SSE-only for the long-tail and reserving WS for top-1000 hottest hotels.
  • Edge cache invalidation cadence. 5 s TTL × millions of hot pages = many origin requests despite singleflight. Lengthen to 10 s on the long tail; keep 5 s for hot.

Trade-offs we accepted:

  • 0.81 % HLL standard error (well below 5 % SLO) for O(12 KB) per hotel regardless of V.
  • ≤ 5 s count staleness from edge cache (the load-bearing trick that makes per-key reads tractable).
  • ≤ 5 s RPO of heartbeat data on Redis async replication (next heartbeat in 15 s anyway).
  • WS as the push transport even though SSE is leaner — paid for in connection state to gain the client→server channel for keep-alive liveness and clean reconnect handshakes (presence still flows on the separate heartbeat POST, not this channel).
  • Sticky-hash routing on hotel_id at the LB and ws-fanout — reduces operational flexibility (can't scale a hot hotel's connections independently of others) but is the only way Pusher-style coalescing actually works.

Per-component failure stories (one-line each, the longer story is in each node's failureMode):

  • CDN PoP outage — origin sees 20–50× spike; widget falls open to —.
  • Redis hot-shard meltdown — auto-promote N_shards 16 → 64; edge fan-in cadence 1 s → 250 ms.
  • Redis leader failover — 15-30 s gap; PFADDs 503 during gap; counts decay; widget shows last-known-good.
  • Cache stampede — XFetch + singleflight + jitter; layered, not picked one of.
  • WS reconnect storm — Reconnect-After: $jittered; rolling-deploy ≤5 % at a time.
  • Abuse worker backlog — scoring goes silent; heartbeats and counts unaffected; page on 30 min lag.
  • Cert expiry on WS — clients fall back to 10 s polling; cert-expiry alerts at T-30/14/7/1 days.
  • Region split-brain — counts diverge per region; on heal, PFMERGE union — counts reconcile up.

Primary sources

  • Heule/Nunkesser/Hall — HyperLogLog in Practice (HLL++, EDBT 2013)
  • Vattani et al. — Optimal Probabilistic Cache Stampede Prevention (VLDB 2015)
  • Pusher — How we built subscription counting at scale
  • Discord — How we scaled Elixir to 5M concurrent users
  • Cloudflare — Durable Objects: easy / fast / correct, choose three
  • Google SRE Workbook Ch. 22 — 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 Concurrent Hotel Viewers yourself