Notification System — a worked solution
Push, email, SMS. Idempotent. Failover.
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 Notification System workspaceThe problem
Build a multi-channel transactional notification platform — push (APNs/FCM), email (SES/SendGrid), SMS (Twilio/Sinch). Producers (your own services or merchant integrations) POST a notification request with an Idempotency-Key; the platform fans out to the user's preferred channels under the user's preferences and quiet hours, retries on transient provider failures, dedupes end-to-end so no duplicate ever reaches the recipient, and survives a full-region failover without sending the same message twice.
This is the system that powers "your driver is arriving" SMS, "order shipped" email, "2FA code" push, and the merchant-webhook fan-out that other systems lean on. Three things define it: (1) end-to-end idempotency under at-least-once Kafka and retried HTTP, (2) per-channel bulkheading so one provider's outage doesn't cascade, and (3) active-region fencing so failover never duplicates.
The reference architecture
What each component is for
- Client
Two distinct callers share this entry point: (a) an internal service emitting a transactional event ('order shipped', 'ride arriving', 'password reset', '2FA code') via POST /v1/notifications with an Idempotency-Key; (b) a merchant integration POSTing webhook deliveries when this same platform doubles as a webhook fan-out tier. Both expect a 202 Accepted in <100 ms — the actual delivery happens async on the Kafka spine.
Why it exists. Considered modeling internal-callers and merchant-integrations as separate edge nodes. Rejected because their failure shapes are identical (idempotent retry, HMAC-signed body, exp backoff) and the only divergence is auth method — service-mesh JWT for internal, scoped API key for merchants — which we draw at the auth node downstream rather than at the edge.
When it fails. Naïve clients regenerate idempotency keys on retry → ghost notifications. Detection: distinct_keys_per_logical_event_p95 > 1 alerts SDK eng. Mitigation: official SDKs ship the canonical retry helper; broken clients are caught at the gateway via an HMAC-signed request fingerprint mismatch (Stripe pattern: same key + different body → 422).
- API Gateway · Public EdgeEnvoy + Cloudflare WAF + TLS terminator
L7 ingress for /v1/notifications and /webhooks/{provider}. Terminates TLS, runs the WAF ruleset (OWASP top-10 + PII rules — reject log payloads carrying card numbers), enforces token-bucket rate limits per tenant BEFORE forwarding so retry storms can't poison the idempotency store, copies the Idempotency-Key header verbatim onto the upstream request and signs the forwarded headers with HMAC so the acceptor trusts tenant_id without re-validating.
Why it exists. Considered putting auth + rate limit inline in the notification acceptor. Rejected because (1) a 5xx in the acceptor would skip the rate limiter and let the storm reach the idempotency store, blowing out conditional-write capacity; (2) per-tenant rate-limit data lives in a single Redis cluster which we'd otherwise ship to every acceptor replica. Stripe's 2019-07-10 retro showed what happens when rate limits live in the wrong layer — the cascade crosses product surfaces.
When it fails. Edge pool exhaustion masks as 'producer says we're slow' but is actually upstream load. Detection: gateway_active_streams, downstream_5xx_rate. Mitigation: drain mode + per-zone failover; circuit-break a misbehaving tenant rather than 503'ing everyone. Cloudflare's Jun 2022 retro warns: anycast misconfig can route traffic to the wrong region — we monitor region_active_set_size and page if a tenant has >1 active region simultaneously.
- Auth / IdentityHashed scoped-API-key + service-mesh JWT (SPIFFE)
Validates the Authorization header on every public call: hashed key lookup against a sharded Postgres keys table, scope check (notifications:write vs webhooks:read vs admin:read), per-key revocation list. Internal calls use service-mesh JWT (SPIFFE/SPIRE) and skip the keys lookup. Returns a signed (tenant_id, scopes, allowed_channels) envelope back to the gateway, valid for 60 s.
Why it exists. Considered baking authn into the acceptor. Rejected because (1) a CVE in the acceptor should not break authn for everyone (blast-radius isolation), (2) the 100M-row keys table is hot in cache and doesn't belong on the notification hot path, (3) scoped keys (e.g. notifications-only API keys for marketing tools) are a separate product surface that needs scope validation independent of any single service.
When it fails. Auth down → ALL writes fail-CLOSED 503 (better than fail-OPEN which would let unauth traffic through). Cached scopes survive 60 s then 401. Detection: auth_5xx_rate, scope_cache_miss_rate. Microsoft Teams 2020-02-03 cert-expiry outage (4 h global) is the cautionary tale — cert-manager + alert at T-30 days. Square 2023-09-07 cert-bundle outage is the same lesson at the service-mesh layer.
- Service · Notification AcceptorGo + AWS SDK v2 + librdkafka (idempotent producer)
Hosts POST /v1/notifications. For each request: (1) checks the region-lease — if our region is not the active leader for this tenant, return 503 + Retry-After BEFORE touching any durable state (claiming the key or publishing from a standby region would strand a committed message the client was told failed); (2) hashes (tenant_id, idempotency_key, body) into a request fingerprint; (3) attempts a DynamoDB conditional PutItem with attribute_not_exists on the idempotency table — wins → claim; loses → returns the cached response_id; (4) publishes a notif.requested event to Kafka with the idempotent producer (acks=all, enable.idempotence=true) keyed by recipient_id. Returns 202 in <100 ms.
Why it exists. Considered fronting Kafka directly with a thin gateway plus a separate dedup worker. Rejected because the dedupe-then-publish ordering MUST be tight — if dedup happens after publish, a Kafka producer retry creates a dup that the consumer can't dedupe without an end-to-end key. Considered Kafka transactions (exactly-once semantics within Kafka). Rejected because EOS doesn't span Kafka + DynamoDB writes; the cross-store dedupe must live OUTSIDE Kafka. The acceptor is where producer-side dedupe happens; the dispatcher does a SECOND dedupe at send-time as defense-in-depth (Stripe webhook pattern).
When it fails. Idempotency store down → fail-CLOSED 503 (better to lose a request than dup). Kafka producer cluster unreachable → fail-CLOSED 503. Region-lease unreachable (etcd partitioned) → fail-CLOSED 503 (better to refuse than risk both regions sending). Detection: claim_p99 (alert <50 ms 1-min P2), kafka_produce_ack_p99 (alert <200 ms P2), region_lease_check_p99 (alert <30 ms P2). Cloudflare 2022 / Slack 2021 lesson: never fail-OPEN on the fence check — that's the path to double-sends.
- Service · Preference / Opt-Out APIGo + Cassandra driver (LWT) + Redis client
Hosts the preference + opt-out write surface — PUT /v1/users/{id}/preferences and POST /v1/users/{id}/opt-outs — and is where the SMS dispatcher's inbound-'STOP' webhook lands. For an opt-out: (1) INSERT ... IF NOT EXISTS (Cassandra LWT / Paxos at LOCAL_SERIAL) into the user's home-region pref-store for linearizable consistency within that cluster — the TCPA primitive (cross-region is async-mirrored, so the pre-send check must read the home region synchronously); (2) invalidate pref:{user_id} in pref-cache so the router re-reads on the next event (invalidate-on-change, not write-through, to avoid the stale-but-fresh race). Distinct from the notification acceptor, which only writes the Kafka spine.
Why it exists. Considered folding preference writes into the notification acceptor. Rejected because the acceptor is tuned for the high-QPS fire-and-forget publish path (DynamoDB claim + Kafka produce), whereas opt-out writes are low-QPS but must be LINEARIZABLE (LWT/Paxos, ~30 ms) — mixing the two couples a rare-but-legally-critical write to the hot path. Separating the plane also keeps the TCPA-exposed code surface small and auditable.
When it fails. pref-store unreachable → opt-out write fails-CLOSED 503 (never ACK an opt-out we didn't durably record — TCPA exposure). Cache-invalidate drop → bounded by the 5 s opt-out TTL on pref-cache. Detection: optout_lwt_p99, optout_write_5xx_rate, pref_cache_invalidate_lag_seconds. Satterfield v Simon & Schuster is the precedent — a missed opt-out is a per-message statutory liability.
- KV Store · IdempotencyDynamoDB Global Tables (us-east-1 ↔ us-west-2)
Stores per-Idempotency-Key state for the producer-side dedupe AND per-(notif_id, channel, recipient) state for the dispatcher-side defense-in-depth dedupe. Schema: (PK = sha256(tenant_id, idempotency_key), request_fingerprint, response_id, response_code, response_body_compressed, ttl). Two operations: CLAIM = PutItem with ConditionExpression 'attribute_not_exists(PK)' (fresh) — wins → claim; FetchExisting = GetItem on conflict → returns cached response. Same store, different PK namespace for the dispatcher dedupe.
Why it exists. Considered Redis SETNX. Rejected because GitHub 2018-10-21 retro shows SETNX races during cluster reshards (we'd hit the same on a Redis cluster topology change), and Redis async failover has RPO=5 s — 5 seconds of dedupe state vanishing on AZ outage is enough to leak dupes. DynamoDB conditional PutItem is the AWS Builders' Library 'Idempotency at Scale' standard answer (re:Invent ARC403 2021). Considered Postgres unique constraint — works but cross-region replication via logical-decoding is more operationally complex than DynamoDB Global Tables, which AWS manages.
When it fails. Cross-region replication lag > 5 s (AWS-published worst case 30 s) opens a dedupe-bypass window — a producer retry that lands in the standby region before its idem-row replicates from the home region wouldn't dedupe. Mitigation: regional ownership — each tenant has a HOME region (consistent-hash on tenant_id); the OTHER region 503s for that tenant unless an explicit failover ramp has fired. Defense-in-depth dispatcher dedupe (in-region) catches same-region replays even if cross-region leaks. Detection: dynamodb_replication_lag_seconds (alert >5 s P2, >15 s P1), ConditionalCheckFailed_rate (high = busy), throttled_writes_total (high = capacity issue). Single-row blast radius — never larger than one Idempotency-Key.
- Coordinator · Region Leaseetcd 3.5 (5-node Raft cluster across 3 AZ)
Holds a per-tenant 'active region' lease at /lease/{tenant_id} → region_id, with a 30 s TTL renewed by a 10 s heartbeat from the active region. Notif-acceptor and dispatchers query before accepting / sending: if our region != lease holder, fail-CLOSED (acceptor returns 503; dispatcher drops the in-flight message — the active region's dispatcher will pick it up). Operator-flippable for planned failover via a typed runbook.
Why it exists. Considered DNS-based region failover. Rejected because DNS TTL plus client-side resolver caching means both regions can be active for several minutes during a cutover — a guaranteed double-send window for transactional channels. Considered DynamoDB Global Tables for the lease (single-row conditional writes). Rejected because Global Tables replication lag is documented up to 30 s, and the lease IS the safety primitive — the lease store must have STRONGER consistency than the dedupe store. etcd / Consul give us linearizable leases via Raft. Cloudflare Jun 21 2022 and Slack Jan 4 2021 retros are the precedents — fencing tokens for active-region are the *Database Reliability Engineering* (Campbell & Majors, ch. 9) answer.
When it fails. etcd quorum lost (3 of 5 nodes failed) → all writes fail-CLOSED → notification system halts (no fence = no send). Severe SPOF. Mitigation: 5-node cluster across 3 AZs + cross-region observers (operator-driven failover RTO ~10 min); per-tenant cached lease state survives 30 s on etcd unavailability so brief blips don't halt the system. Detection: etcd_has_leader, etcd_proposals_committed_total, lease_renewal_lag_seconds. Trade: we accept full-system halt over double-send risk (lower-cost incident than TCPA / merchant-trust damage).
- Stream · Notification SpineApache Kafka 3.7 (KRaft mode, 24 brokers, 3 AZ)
Per-region Kafka cluster carries every notification through the system — the tenant's home-region cluster is the spine of record, async-mirrored to the standby region via MirrorMaker 2 so a fenced cutover resumes from a local copy. Topics: notif.requested (post-acceptor), notif.{push,email,sms}.queued (post-router, per-channel for bulkheading), notif.{push,email,sms}.retry (delayed-retry topics with timestamp-keyed exp backoff), provider.feedback (webhook receiver), notif.delivered (audit). Partitioned by hash(recipient_id) so per-recipient ordering is preserved (a 'started' notif precedes the 'completed' for the same trip).
Why it exists. Considered SQS / RabbitMQ. Rejected because (a) SQS FIFO groups cap at 300 msg/sec/group which can't carry our peak 92 K QPS without painful sharding gymnastics, (b) replay-from-offset for bug-fix backfills is impossible with SQS's 14-day point-of-no-return retention, (c) multi-consumer fan-out (push + email + sms + audit + analytics all reading the same notif.requested) wants Kafka's pub-sub model, not point-to-point. Per Uber's Push Platform, DoorDash's notification platform, LinkedIn ATC, and Slack's Job Queue — Kafka is the production answer.
When it fails. Single broker loss → controller re-elects leaders for affected partitions, ~5 s gap on the affected partition's writes. min.ISR=2 → single-broker loss never blocks; two-broker loss blocks writes (correct — better than risking data loss). Hot partition (10× QPS on one recipient — viral content / mass corporate user) lights up monitoring; mitigation is partition-level rate limiting + recipient-group coalescing at the dispatcher. Consumer-group rebalance during deploys is the LinkedIn KIP-429 lesson — cooperative rebalancing turned a 30-60 s STW into seconds; static membership (group.instance.id) makes restart-in-place rebalance-free.
- Service · Router & RendererJava + Kafka Streams + handlebars/Mustache templates
Consumes notif.requested. For each event: (1) fetch user preferences + device tokens from preference cache (miss → Cassandra); (2) apply frequency cap (LinkedIn ATC rule engine — e.g. max 5 notif/user/hour for non-critical); (3) apply quiet-hours filter in user-local timezone (NOT UTC); (4) render the template per channel (handlebars — Twitter '{firstName}' bug prevention via golden-dataset render-diff in CI); (5) for each channel that survives filtering, publish a notif.{push,email,sms}.queued event back to Kafka with the rendered payload. Stateless — all state is in pref-store/cache.
Why it exists. Considered doing routing inside the channel dispatchers. Rejected because that path makes every dispatcher re-render and re-fetch preferences — 3× the read load on pref-store and N× the duplicated render code paths. Centralizing routing means ONE preference fetch per business event regardless of how many channels fire. Considered making this a Kafka Streams app (stateful join). Rejected for operational simplicity — stateless service + cache hits the same 90%-of-events case at lower complexity. LinkedIn Air Traffic Controller (Espresso + Couchbase + Samza dispatcher) is the canonical reference; we adapted the shape to Cassandra + Redis + plain Kafka consumers.
When it fails. Pref-cache cold restart → all reads hit Cassandra → LOCAL_QUORUM saturates → render p99 climbs from 30 ms to 300 ms. Detection: pref_cache_hit_ratio < 0.85 P2 page; render_p99 > 100 ms P2. Mitigation: warmup job replays last 30 min of events on restart; circuit breaker fast-fails when Cassandra > 200 ms p99 (degrade to 'best effort default preferences' for non-critical channels). Stateless — scaling is horizontal, no local state to lose. Malformed-template deploy (Twitter '{firstName}' literal) prevented by golden-dataset render-diff in CI + canary rollout.
- Cache · PreferencesRedis 7 cluster mode, 16 shards
Caches pref:{user_id} (preferences blob, TTL 60 s) and freq:{user_id} (rolling-window frequency counter, TTL 1 h). Populated lazily by the router on miss; INVALIDATED (not written-through) by the user-API when the user changes preferences — invalidate-on-change avoids the double-write race that gives stale-but-cache-fresh data.
Why it exists. Considered putting the cache inside the router process (in-memory). Rejected because (1) cache hit ratio depends on cross-replica sharing — in-memory caches partition by replica and give 1/N hit ratio at our 16-replica fleet, dropping our effective hit rate to ~6 %; (2) opt-out invalidation needs to touch ALL caches atomically, which is easier with a centralized Redis. Considered making Cassandra the only store. Rejected because at 100 K events/sec peak that's 100 K Cassandra LOCAL_QUORUM reads/sec — Cassandra prefers <50 K reads/sec/cluster at this footprint without expensive horizontal scaling.
When it fails. Cluster cold restart → 100 % miss → router falls through to Cassandra → see pref-store failure mode (saturates LOCAL_QUORUM reads). Single shard down → 1/16 of users go cold; router degrades to Cassandra for those slot groups. Detection: hit_ratio, cluster_state, slot_health. Mitigation: warmup job replays preferences for the top-1M most-active users on startup. NEVER trust the cache for opt-outs — those go straight to pref-store via Cassandra LWT (TCPA exposure if the cache leaks).
- NoSQL DB · Preferences & TokensCassandra 4.1, RF=3, LOCAL_QUORUM
Source-of-truth for: (a) per-user preferences (channels enabled, quiet hours, opt-outs, frequency caps); (b) device-token registry (one row per (user_id, app_id, device_id) with token + last_seen_at + status); (c) suppression lists (bounced emails, complained addresses, blocked numbers). Read-mostly (~95 % reads from the router hot path, ~5 % writes from webhook GC + user opt-out flows).
Why it exists. Considered Postgres + read replicas. Rejected because at 200M users × 2.5 devices × 3 apps the row count breaks single-shard Postgres; cross-shard reads on the router hot path would need a fan-out which Cassandra avoids via partition key = user_id. Considered DynamoDB. Rejected because per-user query patterns include 'all my devices for app X' which is a sort-key range scan in DynamoDB (provisioned-RCU expensive at this scale) or a Cassandra clustering-key scan (cheap). Considered Aerospike — works but our team owns Cassandra ops. The Monzo 2019-07-29 retro on auto_bootstrap is the cautionary tale we learn from for cluster scale-out.
When it fails. AZ loss → quorum still met (2 of 3 replicas survive). Region loss → operator-driven cross-region promote (we deliberately don't auto-promote because cross-region cassandra-stream-mirroring lag varies under load). Hinted handoff handles brief replica unavailability. Monzo 2019-07-29 lesson: never relax to consistency=ONE on writes that must be authoritative — for opt-outs we use LWT (Paxos) and ACCEPT the latency cost. Detection: hinted_handoff_queue_depth, read/write timeout rate, cluster status, optout_replication_lag_seconds (P1 alert). Single-partition blast radius — 1/N users.
- Worker · Push DispatcherGo + APNs HTTP/2 client + FCM Admin SDK
Consumes notif.push.queued from Kafka. For each message: (1) defense-in-depth dedupe — conditional PutItem on idem with PK=(notif_id, channel='push', recipient_token); (2) check the region-lease — if not the active region for this user, drop the message (active region's dispatcher will pick it up); (3) send via APNs HTTP/2 (apns-priority based on notification class) or FCM HTTP v1 (multicast up to 500 tokens for fan-out coalescing); (4) on success, publish notif.delivered to Kafka for audit. Per-provider circuit breaker with 50 % error threshold over 30 s window.
Why it exists. Considered one polyglot dispatcher for all channels. Rejected because the Twilio Jul 18 2023 retro shows EXACTLY what happens when channels share a thread pool — Twilio's degradation took down push and email in teams that didn't bulkhead. Per-channel dispatchers give us the bulkhead. Considered making this a Lambda. Rejected because per-call cold-start tail (200-800 ms) violates push p99 budget, and the long-lived APNs HTTP/2 connections (worth keeping warm — TLS handshake is 50 ms) want a long-lived process.
When it fails. APNs/FCM 5xx storm → circuit opens → messages divert to notif.push.retry (delayed Kafka topic with timestamp-keyed scheduling). Discord 2018 push-retry-storm + AWS Builders' Library precedents. Detection: apns_5xx_rate, circuit_breaker_state, push_p99, push_dispatcher_pool_saturation. Token churn: APNs returns 410 Unregistered (and FCM returns UNREGISTERED) SYNCHRONOUSLY on the dispatcher's own HTTP/2 send — the dispatcher detects it inline and emits a token-GC event that deletes the token from the pref-store (no inbound webhook; APNs/FCM don't POST callbacks). Daily Spark job over pref-store reaps tokens with last_seen > 90 d as backstop. Per-channel bulkhead means a push outage doesn't drag SMS or email down.
- Worker · Email DispatcherGo + AWS SES SDK + SendGrid v3 client
Consumes notif.email.queued. For each message: (1) defense-in-depth dedupe; (2) region-fence check; (3) suppression-list lookup against pref-store (skip if email is suppressed due to bounce/complaint); (4) route to the right IP pool — transactional pool (high reputation, dedicated IPs) vs marketing pool (separate reputation domain); (5) send via SES; on Throttling, vendor-failover to SendGrid (manual approval — too easy to make worse); (6) publish notif.delivered to Kafka.
Why it exists. Considered making the email path a thin Lambda over SES SendEmail. Rejected because (1) suppression-list lookups want a long-lived connection to pref-store (Lambda spins one up per cold start), (2) IP-pool routing is non-trivial state, (3) SES Throttling fast-fails benefit from a circuit breaker held by a long-lived process. Considered SES alone. Rejected because we need vendor failover — when SES has a regional outage we route to SendGrid; relying on one provider for the most-trafficked channel is the SPOF that retro pages get written about (e.g. Mailgun's 2018 outage took down Slack invites).
When it fails. SES Throttling → circuit opens, messages divert to notif.email.retry. Bounce-rate spike (>5 %) → SES auto-suspends sending → vendor failover to SendGrid (manual approval). SendGrid IP-reputation crash from a bad campaign — quarantine campaigns at >2 % bounce in 1 h. Detection: bounce_rate_15m, complaint_rate, ses_5xx_rate. Per-channel bulkhead means an email outage doesn't drag push or SMS down.
- Worker · SMS DispatcherGo + Twilio Programmable SMS + Sinch client
Consumes notif.sms.queued. For each message: (1) defense-in-depth dedupe; (2) region-fence check; (3) carrier lookup (E.164 normalization + carrier prefix) to pick the right Twilio messaging-service-SID for that destination's regulatory regime (US 10DLC, UK shortcode, IN DLT-registered, EU sender-ID); (4) per-MSISDN rate-shape (Twilio long-code = 1 MPS, short-code = 100 MPS, A2P 10DLC tiers up to 75 MPS); (5) SYNCHRONOUS opt-out check against pref-store (TCPA compliance — never trust the cache); (6) send via Twilio Messages API; on 503, vendor-failover to Sinch.
Why it exists. Considered putting SMS into the same dispatcher as push. Rejected because SMS rate-shaping is fundamentally different — push has a pool of HTTP/2 connections, SMS has a pool of SENDER NUMBERS each with its own carrier-imposed cap. The number-allocation logic is non-trivial (carrier-prefix → messaging-service-SID lookup) and doesn't belong in a polyglot dispatcher. Lyft's published 'SMS Reliability' post documents exactly this split (2019).
When it fails. Twilio 503 → vendor failover to Sinch (manual + automatic depending on incident length). Twilio's Jul 18 2023 incident is the precedent (global SMS queue). Detection: twilio_5xx_rate, sinch_5xx_rate, sms_p99_carrier_ack. Per-channel bulkhead (separate consumer group + topic + dispatcher pool) means a Twilio outage doesn't take down push or email. Carrier-filtering escalation: if 30007 rate spikes, auto-quarantine the template and page content team. Opt-out replication lag → fail-CLOSED with retry-after; legal cost of one missed opt-out >> latency cost.
- External · Provider GatewaysAPNs HTTP/2 / FCM HTTP v1 / SES / SendGrid / Twilio / Sinch
Receives push, email, and SMS messages from our dispatchers and delivers to the recipient's device, mailbox, or carrier. Sends asynchronous webhook callbacks to our webhook receiver for delivery state changes (delivered, bounced, complained, unregistered).
Why it exists. We don't run our own carrier interconnect or MTA fleet — at this volume, multi-vendor providers (Twilio + Sinch for SMS, SES + SendGrid for email, APNs + FCM for push) give us higher delivery success than home-rolled would, plus they handle the regulatory minefield (10DLC TCR registration, GDPR, DKIM key management, APNs cert provisioning). Considered self-hosted Postfix + carrier SS7 connections. Rejected because that path needs an 18-month build, ~$5M/yr ops cost, and the regulatory risk (a misrouted SMS to a country we lack a sender-ID in is an immediate carrier ban) is the kind of thing a payment-system retro is written about.
When it fails. Single provider down → vendor failover (preferred). Multi-provider correlated failure (rare — a regulatory or DNS-level event) → channel goes dark; circuit breakers fast-fail; messages divert to retry topic. Detection: per-provider 5xx rate, delivery success rate (NOT what the provider reports — the GAP between provider-reported-success and device-acked-success is the gray-failure tell flagged by Uber and LinkedIn). Twilio 2023, FCM Aug 2022 outages both hit this pattern.
- Service · Provider Webhook ReceiverGo + AWS SDK + per-provider signature verification
Hosts /webhooks/{ses,sendgrid,twilio,sinch} — the providers that actually POST callbacks. APNs and FCM do NOT send delivery/unregister webhooks: their 410 Unregistered / UNREGISTERED comes back SYNCHRONOUSLY on the push-dispatcher's own HTTP/2 send, so push-token GC is initiated at the dispatcher, not here. For each inbound callback: (1) verify the provider's signature (SNS RSA-certificate signature for SES events, ECDSA P-256 for SendGrid's Signed Event Webhook, HMAC-SHA1 X-Twilio-Signature for Twilio); (2) dedupe against a separate idempotency PK = (provider, provider_event_id) — 30-day TTL covers SES/Twilio replay windows which can stretch up to 30 d; (3) classify: delivered / bounced / complained; (4) for hard bounces / unsubscribes, update the suppression list (and delete the dead email/SMS address row); (5) publish provider.feedback to Kafka for downstream consumers (audit, suppression-list updater, analytics).
Why it exists. Considered routing provider webhooks straight into the API gateway. Rejected because providers' source IPs vary, signature verification is per-provider (different algorithms), and the rate-limit profile is different (providers can dump 10 K webhooks/sec on a delivery batch — completely different from our per-tenant 200/sec ingest cap). Centralizing into a dedicated receiver lets us tune rate-limits and signature checks per provider. Also, the receiver MUST be CRITICAL-AVAILABILITY — if it's down, providers retry for hours and our token-GC stalls (apns_5xx_rate climbs from a flood of dead tokens).
When it fails. Webhook receiver down → providers retry for hours; we backfill once it's back. Token-GC stalls → apns_5xx_rate climbs; mitigation is the daily Spark sweeper over pref-store that reaps last_seen > 90 d tokens regardless of webhook state. Detection: provider_webhook_5xx_rate, signature_verify_fail_rate, token_gc_lag_seconds. SES bounce surge (>5 % → SES auto-suspends sending) is a P1 — paged immediately.
- SQL DB · Notification HistoryPostgres 16 + Citus, 64 hash shards by recipient_id
Append-only history of every notification: (notif_id, tenant_id, recipient_id, channel, template_version, sent_at, provider, provider_message_id, status, attempt_count, opt_in_proof_id). Powers /v1/notifications/{id} status (a merchant polling 'did my notification deliver?'), the per-user history UI ('your notifications'), and compliance audits — TCPA / GDPR — 'show me every SMS we sent to this number with timestamp + opt-in proof'.
Why it exists. Considered ClickHouse for analytics + nothing for OLTP queries. Rejected because the /v1/notifications/{id} status endpoint is sub-100 ms p99 (a merchant polling) — ClickHouse doesn't do single-row OLTP at p99. Considered keeping audit in Cassandra alongside preferences. Rejected because audit queries cross-reference opt-in proofs with delivery logs (legal compliance: prove we had consent at send time) and that's a SQL JOIN — Cassandra discourages JOINs and forces app-level fan-out which is slow and error-prone. Sharded Postgres + Citus is the production choice (Stripe, Square Books, Monzo all use this pattern for audit logs).
When it fails. Single shard down → 1/64 of recipients see status-endpoint 503; sends still happen (writes are async via Kafka consumer — the audit write isn't on the critical path). Detection: shard_health, replica_lag_p99, audit_consumer_lag. Replication lag spike: status endpoint reads from leader for fresh data (RYW); analytics reads from follower (eventual). Cold-tier S3 query latency is fine for compliance subpoenas (hours acceptable).
Stage by stage
The same 10 stages the workspace walks, answered.
01Clarifications
What would you ask before drawing a single box?
Typical clarifications to surface:
- Channels? Push + email + SMS — confirm. Webhook fan-out for merchants is the same machinery; same idempotency story.
- Read traffic? A merchant polling /v1/notifications/{id} to check delivery status — yes; that's the audit-store read path.
- Multi-tenant? Yes. Per-tenant rate limits at the gateway. Per-tenant home-region pinning for the active-region lease.
- Multi-region? Active-active for the public ingress (any region can accept), but per-tenant ACTIVE-STANDBY for SEND — only the tenant's home region dispatches to providers, with a fenced failover. This is the single most important architectural decision and the only thing that prevents the Cloudflare-2022 / Slack-2021 double-send incident.
- Latency targets? Push p99 < 5 s (4 s of which is provider+radio; ours is < 1 s). SMS p99 < 30 s (carrier-bounded). Email p99 < 5 min. Webhook outbound is retry-budget-bounded, not latency-bounded.
- Idempotency contract? Producer-supplied V4 UUID Idempotency-Key, 7-day TTL. Same key + different body = 422.
- Compliance? TCPA for SMS (opt-out is non-negotiable), CAN-SPAM for email, GDPR for EU users.
Assumptions to state:
- 1 B notifications/day total, mix push:email:SMS = 90:9:1.
- 11.5 K avg QPS aggregate; 92 K peak (8× burst on Black Friday / breaking news).
- 200 M MAU × 2.5 devices × 3 apps = 1.5 B device-token rows.
- 5-year retention default for compliance; 7 years for finance-adjacent traffic.
- Per-recipient ordering must be preserved (a "started" precedes the "completed" for the same trip).
02Functional reqs
What must this system actually do?
- Accept POST /v1/notifications with Idempotency-Key and route to the user's preferred channels.
- Apply user preferences (channels enabled, quiet hours, frequency caps, suppression lists) before send.
- Render templates per channel with locale + personalization.
- Dispatch to APNs/FCM/SES/SendGrid/Twilio/Sinch with vendor failover.
- Receive provider feedback (delivered, bounced, complained, unregistered) and update suppression lists + token registry.
- Expose GET /v1/notifications/{id} for delivery status.
- Audit log every send with opt-in proof for TCPA / GDPR compliance.
03Non-functional
What must it promise about speed, uptime and correctness?
- Idempotency: End-to-end. No duplicate ever reaches the recipient under producer retry, Kafka rebalance, DLQ replay, or region failover.
- Availability: 99.95 % accept (the API). 99.9 % delivery (push/email/SMS — provider-dependent for the last mile).
- Latency: Accept p99 < 100 ms. Push delivered p99 < 5 s. SMS accepted-by-carrier p99 < 30 s. Email accepted-by-MTA p99 < 5 min.
- Durability: Once we 202 the producer, the notification is durably committed to Kafka. RPO = 0 within a region; cross-region RPO ≤ 5 s on async mirror.
- Failover: Single-region failure → operator-driven cross-region cutover with fenced lease handoff. RTO 10 min; ZERO duplicate sends during the cutover.
- Compliance: TCPA opt-out check on every SMS (synchronous, against authoritative store). 4-year audit retention. DKIM/SPF/DMARC alignment for email.
04Capacity estimation
How much load and data does this have to hold?
Assumptions: 1 B notifs/day total, 90 % push / 9 % email / 1 % SMS, 8× peak multiplier.
- Aggregate QPS: 1 B / 86 400 ≈ 11.5 K avg, 92 K peak.
- Push QPS peak: 92 K × 0.90 ≈ 84 K. APNs serves ~9 K/sec/HTTP-2 connection → 10 connections ≈ 1 dispatcher replica with headroom; we run 16 for fault-tolerance and burst. 84 MB/s outbound to APNs/FCM aggregated, ~5 MB/s/replica.
- Email QPS peak: 92 K × 0.09 ≈ 8.3 K. SES production tier 10 K/sec on dedicated IP pool; we run 8 dispatcher replicas across 2 IP pools for transactional + marketing isolation.
- SMS QPS peak: 92 K × 0.01 ≈ 920. A2P 10DLC T1 = 75 MPS per messaging-service-SID; we hold 14 SIDs across 6 dispatcher replicas (headroom 4×).
- Kafka ingress at peak: 92 K msg/s × 1 KB = 92 MB/s ingress; with RF=3 = 276 MB/s broker write bandwidth across 24 brokers ≈ 12 MB/s/broker. With per-channel fan-out (one inbound becomes push + email + sms), peak total topic write = ~230 MB/s across the cluster.
- Idempotency hot storage: steady-state day volume sets this, not peak QPS — 1 B/day × 7 d × 100-150 B/row ≈ 0.7-1.05 TB. At ~1 TB the table spans many DynamoDB partitions; key is hash-distributed (sha256(tenant_id, key)) so no single partition goes hot.
- Device-token storage: 200 M users × 2.5 devices × 3 apps × 200 B = 300 GB in Cassandra.
- Audit history hot: 1 B/day × 90 d × 256 B = 23 TB across 64 Postgres shards (~360 GB/shard).
- Audit history cold: 1 B × 365 × 7 × 256 B = 654 TB on S3 Parquet (~$15 K/yr at S3 IA pricing).
- Per-channel provider cost (peak day): Push free; email 90M × $0.0001 = $9 K/day; SMS 10M × $0.0075 = $75 K/day. SMS dominates the bill.
Bytes-economics check (write-path): Notification bodies are 1 KB typical, 64 KB max. Service-tier mediated is correct — the acceptor's per-replica bandwidth (~7 MB/s) is far below the 100 MB body-economics threshold for direct-to-storage. No need for pre-signed URLs.
05API design
What does the outside world call, and what comes back?
POST /v1/notifications
Authorization: Bearer <scoped-key>
Idempotency-Key: 8a7b3e2c-...
{
"tenantId": "merchant_abc",
"recipientId": "user_42",
"channels": ["push", "email", "sms"],
"templateId": "order_shipped_v3",
"data": { "orderId": "ord_91", "trackingUrl": "https://..." },
"priority": "transactional",
"ttlSeconds": 3600
}
202 Accepted
{ "notificationId": "notif_xyz", "status": "queued" }
# Retry with same key + same body:
202 Accepted (same notificationId)
# Retry with same key + different body:
422 { "error": "idempotency_key_mismatch" }
GET /v1/notifications/notif_xyz
200 {
"notificationId": "notif_xyz",
"status": "delivered",
"channels": [
{ "channel": "push", "status": "delivered", "providerMessageId": "...", "deliveredAt": "..." },
{ "channel": "email", "status": "delivered", ... },
{ "channel": "sms", "status": "carrier_filtered", "errorCode": "30007" }
]
}
# Polled before any notif.delivered event has landed (i.e. right after the 202):
425 { "error": "too_early", "notificationId": "notif_xyz" }
This endpoint reads the audit store, which today is written only by the dispatchers' post-send notif.delivered event (edges e18/e19/e20). No edge inserts a row at accept/route time, so a merchant polling immediately after the 202 finds nothing — hence GET is delivery-state-only: pre-delivery it returns 425 Too Early, not the queued row the 202 body implies. To make queued observable you need a Kafka consumer edge from notif.requested (or a notif.{channel}.queued fan-in) into audit that inserts the row in queued status, which the notif.delivered event then UPSERTs to delivered/bounced (see # hld graph gap).
POST /webhooks/{ses,sendgrid,twilio,sinch} # Provider delivery/bounce/complaint callbacks
# signed (SNS for SES, ECDSA P-256 for SendGrid, HMAC for Twilio); body schema is per-provider
# NOTE: APNs and FCM have NO delivery webhook — 410 Unregistered / UNREGISTERED is returned
# synchronously on the dispatcher's own push send, which triggers token GC inline.
200 OK
Preferences & opt-out (the TCPA-critical write surface):
PUT /v1/users/{id}/preferences # channels enabled, quiet hours, frequency caps
Authorization: Bearer <scoped-key>
{ "channels": { "sms": false }, "quietHours": { "tz": "America/New_York", "from": "22:00", "to": "08:00" } }
200 OK
POST /v1/users/{id}/opt-outs # record a statutory opt-out (TCPA / CAN-SPAM)
{ "channel": "sms", "identifier": "+1415...", "source": "stop_keyword" }
202 Accepted # write goes to pref-store via Cassandra LWT (Paxos), then pref-cache is invalidated
# Inbound carrier STOP (TCPA): the SMS dispatcher's inbound-message webhook lands here
POST /v1/users/{id}/opt-outs { "channel": "sms", "source": "stop_keyword" }
These writes are the headline compliance primitive: the opt-out path performs the INSERT ... IF NOT EXISTS LWT into pref-store and invalidates pref:{user_id} in the cache. They are served by the user-API / preference-service plane (see # hld) — distinct from the notification acceptor, which only writes the Kafka spine.
06Data model
What gets stored, and what is it looked up by?
Idempotency table (DynamoDB Global Tables):
| field | type | notes |
|---|---|---|
| pk | string | sha256(tenant_id, idempotency_key) — partition key |
| request_fingerprint | binary | sha256(method, path, body) — collision check |
| response_id | string | the assigned notif_id |
| response_code | int | the cached HTTP status |
| response_body_gz | binary | gzip'd response body for replay |
| ttl | number | epoch seconds; DynamoDB auto-expires (7 day window) |
A SECOND PK namespace (notif_id, channel, recipient_token) covers the dispatcher-side defense-in-depth dedupe.
Preferences & device tokens (Cassandra):
CREATE TABLE preferences (
user_id uuid,
...preferences blob...,
PRIMARY KEY (user_id)
);
CREATE TABLE device_tokens (
user_id uuid,
app_id text,
device_id text,
token text,
status text, // 'active', 'unregistered'
last_seen timestamp,
PRIMARY KEY (user_id, (app_id, device_id))
);
CREATE TABLE suppressions (
channel text, // 'email', 'sms'
identifier text, // email or E.164
reason text, // 'bounce', 'complaint', 'opt-out'
created_at timestamp,
PRIMARY KEY ((channel, identifier))
);
Opt-out writes use LWT (Paxos) — INSERT ... IF NOT EXISTS at LOCAL_SERIAL — to guarantee linearizable consistency WITHIN the user's home-region cluster. Cross-region is async-mirrored (eventually consistent, not linearizable — the regions are separate clusters); the compliance guarantee comes from routing the user to their home region and running the pre-send opt-out check as a synchronous home-region read, so it never misses an opt-out. The latency cost (~30 ms) is dwarfed by TCPA statutory damages ($500-$1500 per missed opt-out).
Audit history (Postgres + Citus, sharded by recipient_id):
| field | type | notes |
|---|---|---|
| notif_id | uuid | PK part |
| tenant_id | uuid | |
| recipient_id | uuid | shard key (hash 64) — leading PK column |
| channel | text | 'push' / 'email' / 'sms' — PK part |
| template_id | text | |
| template_version | int | |
| sent_at | timestamptz | |
| provider | text | 'apns' / 'fcm' / 'ses' / 'sendgrid' / ... |
| provider_message_id | text | |
| status | text | 'queued' / 'delivered' / 'bounced' / ... |
| attempt_count | int | |
| opt_in_proof_id | uuid | refs the consent record (TCPA / CAN-SPAM) |
One row per (notif_id, channel) — that's what lets GET /v1/notifications/{id} return the per-channel status array. PRIMARY KEY (recipient_id, notif_id, channel): Citus requires the distribution column in every PK / unique constraint on a distributed table. And because the shard key is recipient_id but the status poll looks up by notif_id, the notif_id embeds the recipient hash (notif_{recipHash}_{uuid}) — the API layer extracts it and routes to one shard instead of scatter-gathering all 64.
Database choice — recommended:
- DynamoDB Global Tables for idempotency (cross-region cond PutItem is the AWS Builders' Library standard answer for at-scale dedupe; Redis SETNX is too lossy under cluster reshards and Redis async failover RPO=5s).
- Cassandra (RF=3, LOCAL_QUORUM) for preferences + tokens + suppressions (per-user partition, opt-out via LWT for linearizable correctness, async cross-region mirror for survivability).
- Postgres + Citus for audit history (SQL JOIN with opt-in-proof table for TCPA queries; ClickHouse can't do single-row OLTP at the < 100 ms p99 the merchant status-poll requires).
- Kafka 3.7 for the durable spine (replay-from-offset, partition-keyed ordering, multi-consumer fan-out — none of which SQS/RabbitMQ give you).
- Redis 7 cluster for the preference cache (cache, not authority).
- etcd 3.5 for the active-region lease (linearizable Raft is the safety primitive; can't be DynamoDB which has 30s replication lag).
07High-level design
Which components handle a request, and in what order?
Architecture summary:
- API Gateway (Envoy + WAF + per-tenant rate limit). Anycast across 2 regions.
- Auth / Identity (scoped API keys + service-mesh JWT). Validates before the acceptor sees the request.
- Notification Acceptor (the write plane). Claims idempotency in DynamoDB Global Tables, publishes notif.requested to Kafka, returns 202.
- Idempotency Store (DynamoDB Global Tables). Cross-region dedupe. Two PK namespaces: producer-side (tenant_id, idempotency_key) and dispatcher-side (notif_id, channel, recipient).
- Region Lease (etcd). Per-tenant active-region fencing token. Both the acceptor and the dispatchers check this before doing user-visible work.
- Kafka spine (per-region, 24 brokers, 3 AZ; async-mirrored cross-region via MirrorMaker 2 — RPO ≤ 5 s). Topics: notif.requested, notif.{push,email,sms}.queued, notif.{push,email,sms}.retry, provider.feedback, notif.delivered.
- Router & Renderer consumes notif.requested, fetches preferences from Redis (miss → Cassandra LOCAL_QUORUM), applies frequency cap + quiet hours, renders templates per channel, publishes notif.{push,email,sms}.queued.
- Per-channel Dispatchers (Push, Email, SMS — separate pools, separate consumer groups, separate provider clients) consume the per-channel topics, defense-in-depth dedupe, fence-check, send to providers with circuit breakers + exp backoff + jitter.
- External Providers (APNs, FCM, SES, SendGrid, Twilio, Sinch) — multi-vendor per channel for failover.
- Webhook Receiver ingests provider callbacks (delivery, bounce, complaint, unregistered), dedupes against (provider, provider_event_id), synchronously deletes dead tokens from pref-store.
- Notification History (Postgres + Citus) is the audit log. Powers /v1/notifications/{id} status, the user history UI, and TCPA / GDPR compliance subpoenas.
- Preference / Opt-Out API (the compliance write plane, separate from the acceptor). Serves PUT /v1/users/{id}/preferences and POST /v1/users/{id}/opt-outs (and the inbound-STOP handler): writes pref-store via Cassandra LWT (the TCPA primitive) and invalidates pref:{user_id} in pref-cache.
> Preference / opt-out write plane. Beyond the send path, the preference + opt-out write surface (PUT /v1/users/{id}/preferences, POST /v1/users/{id}/opt-outs, and the SMS dispatcher's inbound-STOP handler) is served by the user-API / preference-service node fronting two edges: user-API → pref-store (e26 — Cassandra LWT write, the TCPA primitive) and user-API → pref-cache (e27 — invalidate pref:{user_id}). This backs the "opt-outs are written by the user-API" claims in deep-dive 5 and the pref-cache node.
> Graph gap (audit row at accept time). The audit store is written only by the dispatchers' notif.delivered event (edges e18/e19/e20), so no row exists until a send succeeds. For GET /v1/notifications/{id} to return a queued status (as the 202 body implies), the graph needs a Kafka → audit consumer edge off notif.requested (or a notif.{channel}.queued fan-in) that inserts the row in queued status, which notif.delivered then UPSERTs to delivered/bounced. Until that edge exists, GET is delivery-state-only and returns 425 pre-delivery.
Data flow on accept (sync, < 100 ms):
Client → Gateway → Acceptor → [region-lease check, idem CLAIM, kafka publish] → 202
Data flow on push delivery (async, p99 ~ 2 s):
Kafka(notif.requested) → Router → [pref-cache hit / miss → Cassandra] → Kafka(notif.push.queued) → Push Dispatcher → [idem dedupe, region fence, APNs/FCM]
Data flow on token GC (sync, ~ 100 ms):
Provider → Webhook Receiver → [verify sig, dedupe, classify Unregistered] → Cassandra DELETE token
trace catalogue
The simulator runs these specific journeys against your diagram. Each is a deterministic walk; clicking between them lets you see how the canonical answers change shape.
| ID | Journey | Budget |
|---|---|---|
notification-system:accept-fast-path | Multi-segment: (1) CLAIM idem (DynamoDB) + (2) publish notif.requested (Kafka) → 202 | 200 ms |
notification-system:idempotent-retry-replay | Same Idempotency-Key → cached response, no re-publish | 80 ms |
notification-system:push-render-and-send | Kafka → router (cache hit) → Kafka → push dispatcher → APNs/FCM | 2.5 s |
notification-system:preference-cache-miss | Router falls through to Cassandra LOCAL_QUORUM | 300 ms |
notification-system:webhook-feedback-token-gc | Multi-segment: (1) DELETE stale token (Cassandra) + (2) publish provider.feedback (Kafka) | 400 ms |
notification-system:retry-storm-circuit-break | APNs 5xx → circuit opens → notif.push.retry topic | 200 ms |
notification-system:region-failover-fence-block | Standby region's dispatcher checks fence → drops in-flight | 100 ms |
failure scenarios we model
The simulator exercises each of these chaos cards against the canonical and reports per-trace impact. Each is grounded in a published real-world precedent.
| ID | Scenario | Precedent |
|---|---|---|
notification-system:twilio-2023-sms-pile-up | Twilio 503 for 30 min; SMS backs up, push/email survive bulkhead | Twilio Jul 18 2023 status incident |
notification-system:apns-retry-storm | APNs 5xx → naïve retry pile-up self-DOSes the dispatcher | AWS Builders' Library timeouts/retries/backoff |
notification-system:redis-pref-cache-stampede | Pref cache cold restart → router falls through; Cassandra saturates | Twitter Redis-cluster failover retros |
notification-system:dynamodb-global-table-replication-lag | Cross-region link blip; idem replication lag opens dedupe window | AWS re:Invent DAT328 Global Tables trade-off |
notification-system:kafka-rebalance-during-deploy | Rolling deploy triggers consumer-group rebalance; 30-60 s STW | LinkedIn KIP-429 cooperative-rebalance retro |
notification-system:both-regions-active-double-send | Operator-driven failover with old region undrained → double-send | Cloudflare 2022-06-21, Slack 2021-01-04 retros |
notification-system:opt-out-replication-lag-tcpa | Opt-out write hasn't replicated; standby region sends → TCPA exposure | Satterfield v Simon & Schuster |
08Deep dives
Which part breaks first, and what do you do about it?
1. End-to-end idempotency under at-least-once Kafka. Producer-side dedupe at the acceptor (DynamoDB cond PutItem on (tenant_id, idempotency_key)). Dispatcher-side defense-in-depth dedupe at send time (cond PutItem on (notif_id, channel, recipient_token)). Both stored in DynamoDB Global Tables with 7-day TTL. DLQ retention (up to 14 days) can EXCEED the dedupe window, so replay safety comes from the operator-drain rule: the drain script refuses messages older than dedupe-TTL-minus-margin. If we ever wanted unconditional replay, we'd extend the TTL past max DLQ retention (≥16 d) instead. The two namespaces are intentionally different: producer-side catches naïve client retries; dispatcher-side catches Kafka producer retries (idempotent producer is per-broker, not end-to-end across consumer groups). GitHub 2020's webhook-DLQ retro is the precedent for getting the TTL ordering wrong.
2. Active-region fencing tokens. DNS failover lets both regions stay live for minutes (resolver caching). The fence check has to be at the SEND boundary in a CP store (etcd or Consul, NOT DynamoDB which has up-to-30-s replication lag). Each tenant has a HOME region; the lease at /lease/{tenant_id} names it. Acceptor + dispatchers query before doing user-visible work; non-active regions fail-CLOSED. The Cloudflare 2022 + Slack 2021 retros are the precedents — both shipped the "old region kept consuming after failover" incident. The fix is Database Reliability Engineering ch. 9's fencing-token pattern.
3. Per-channel bulkheading. Separate Kafka topics (notif.{push,email,sms}.queued), separate consumer groups, separate dispatcher pools, separate provider clients. The Twilio 2023-07-18 retro shows what happens without this: a 30-minute SMS provider outage took down push + email at teams that shared a thread pool. The bulkhead is the load-bearing property; per-channel topics also lets us tune partition count to per-channel QPS (push 128 partitions, sms 64).
4. Token GC lifecycle. APNs returns 410 Unregistered on a stale token — synchronously, on the dispatcher's own HTTP/2 push (FCM likewise returns UNREGISTERED on the send). There is no APNs/FCM delivery webhook to receive. The dispatcher detects the 410 inline, dedupes against (provider, provider_event_id), and SYNCHRONOUSLY deletes the token from Cassandra (pref-store). Daily Spark sweeper backstops with last_seen > 90 d. Without GC, apns_5xx_rate climbs and Apple eventually IP-throttles us. (Email/SMS providers — SES, SendGrid, Twilio — DO post real callbacks to /webhooks/*; that's what the webhook receiver handles.) Slack and Discord both run this exact pipeline.
5. TCPA opt-out semantics. Opt-out write uses Cassandra LWT (Paxos) at LOCAL_SERIAL for linearizable consistency within the user's home-region cluster — accepts the ~30 ms write latency because TCPA statutory damages are $500-$1500 per missed opt-out (Satterfield v Simon & Schuster). Cross-region is async-mirrored (eventually consistent), so the opt-out check at the SMS dispatcher is SYNCHRONOUS against the user's authoritative HOME-region pref-store, NEVER the cache and NEVER a remote async replica. Quiet hours evaluated in user-local IANA timezone (NOT UTC — that's the 3 AM bug industry folklore).
6. Vendor failover without double-send. Each channel has a primary + standby provider (APNs+FCM, SES+SendGrid, Twilio+Sinch). On primary 5xx storm, dispatcher's circuit opens and routes to standby. The dispatcher-side dedupe ensures that if a message was half-sent on the primary before the cutover (the primary received but didn't ack), the standby attempt is short-circuited. Idempotency is namespaced per (notif_id, channel, recipient_token), not per-provider, so vendor swaps are transparent to dedupe.
7. Marketing blast vs transactional. A 10M-user marketing email blast on the same Kafka topic + SES sub-account as transactional email would back up transactional delivery. We separate at every layer: separate topic (notif.email.bulk vs notif.email.transactional), separate SES sub-account, separate IP pool with separate reputation. Producer-side admission control at the gateway throttles bulk independently of transactional. Mailchimp's published IP-warmup curve guides the bulk-pool ramp.
09Trade-offs
What did this design cost, and what breaks at 10×?
What breaks at 10× scale (~10 B notifs/day, ~1 M peak QPS):
- DynamoDB conditional writes start hitting per-partition limits → shard the idempotency key over more partition prefixes.
- Single Kafka cluster can't hold the full topic count → split per channel into separate clusters; cross-cluster mirror via MirrorMaker 2 for the audit pipeline.
- Cassandra LWT at 1M opt-out QPS becomes the bottleneck → batch opt-outs per user where multiple settings change at once.
- APNs HTTP/2 connection pool at 1M push QPS would need 100+ concurrent connections per Apple cert → registers a second cert and load-balances.
Per-component failure stories (mapped to chaos cards):
- Twilio 503 (chaos:
twilio-2023-sms-pile-up) — bulkhead + Sinch failover. Push and email keep flowing. - APNs 5xx storm (chaos:
apns-retry-storm) — circuit breaker + retry budget capped at 10 % normal traffic. - Pref-cache cold restart (chaos:
redis-pref-cache-stampede) — warmup job + degraded-mode default preferences for non-critical channels. - DynamoDB cross-region lag (chaos:
dynamodb-global-table-replication-lag) — home-region pinning + standby fail-CLOSED with retry-after. - Kafka deploy rebalance (chaos:
kafka-rebalance-during-deploy) — KIP-429 cooperative rebalance + static membership. - Both regions active (chaos:
both-regions-active-double-send) — etcd fencing token + dispatcher dedupe defense-in-depth. - Opt-out replication lag (chaos:
opt-out-replication-lag-tcpa) — synchronous LWT writes; fail-CLOSED on lag.
trade-offs we accepted
- Etcd cluster is a SPOF for active-region designation. Quorum loss halts the system. We accepted this over the alternative (DynamoDB Global Tables) because the 30-s replication lag on Global Tables is too weak a consistency primitive for the safety-critical fence. Mitigated by 5-node etcd across 3 AZs + cross-region observers.
- 30-s worst-case double-active-region window during a hard region partition. The lease TTL has to balance fast renewal against quorum reads. The in-region defense-in-depth dispatcher dedupe catches same-region replays during the window; cross-region dupes during a hard partition are accepted as a residual risk (~30 s × peak QPS = 2.7M maximum exposure, but real-world has been <1000 dupes per recorded incident).
- Email delivery to last-mile is provider-bounded. SES greylisting, MX retries, recipient mailbox-full — none of these are under our control. Our 5-min p99 is "accepted-by-MTA", not "in-inbox". Compliance-grade audit log shows what we did; what the recipient mailbox did is the recipient mailbox's story.
- No exactly-once across regions. We get exactly-once-EXPECTED via the idempotency contract; physically a 30-s cross-region partition during a failover is the documented residual risk. If that risk is unacceptable for a future tenant (e.g. financial regulator messaging), we'd run that tenant single-region with no failover.
Primary sources
- Uber — Real-time Push Platform (uber.com/blog/real-time-push-platform)
- LinkedIn — Air Traffic Controller: Member-First Notifications (2016)
- LinkedIn — Incremental cooperative rebalancing in Kafka (KIP-429)
- Stripe — Idempotency Keys + Webhooks docs
- Brandur — Implementing Stripe-like Idempotency Keys in Postgres
- Slack — Scaling Slack's Job Queue (slack.engineering)
- Slack — Outage on January 4th 2021 retro
- Cloudflare — Outage on June 21 2022 retro
- AWS Builders' Library — Timeouts, retries, backoff with jitter
- AWS re:Invent ARC403 — Idempotency at Scale (2021)
- AWS re:Invent DAT328 — DynamoDB Global Tables (2022)
- Netflix — Active-Active for Multi-Regional Resiliency (2013)
- Twilio — A2P 10DLC Registration docs
- Apple — APNs HTTP/2 provider API docs
- Firebase — FCM scaling guide
- Campbell & Majors — Database Reliability Engineering (ch. 9)
- Beyer et al. — SRE Workbook (Managing Load + 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 Notification System yourselfMore in Queues, Pub/Sub & Event Streaming
Moving events between services without losing them: logs, work queues, fanout, change data capture, stream processing, and the delivery systems built on top.
- Build Build KafkaA partitioned, replicated, append-only log. The log is the database — internalize that, and a dozen product designs get easier.
- Build Build a Message Queue (RabbitMQ / SQS)A point-to-point work queue — the messaging primitive Kafka is NOT. Each message goes to one consumer, ack deletes, retries push to a dead-letter queue, and a poisoned message is everyone's problem. Internalize ack vs visibility timeout vs DLQ vs prefetch vs FIFO groups — and learn to tell when Kafka is the wrong tool and when a queue is.
- Build Build a CDC pipeline (Debezium + outbox)Your service writes to its DB and publishes to Kafka — and any crash between those two writes is permanent inconsistency. Build a Change Data Capture pipeline (modeled on Debezium + the outbox pattern) that closes the gap by making the database itself the event source.
- Ad Click AggregatorStream processing with watermarks and exactly-once.
- Distributed Cron — Mass Scheduled EmailSingle trigger, 50M recipient idempotency, catch-up.