Distributed Cron — Mass Scheduled Email
Worked solution

Distributed Cron — Mass Scheduled Email — a worked solution

Single trigger, 50M recipient idempotency, catch-up.

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 Distributed Cron — Mass Scheduled Email workspace

The problem

A campaign-owner uploads a 50-million-recipient list and schedules a single send for 9:00 AM tomorrow. At 9:00 AM, one scheduler tick fires — and that single event must become 50 million emails leaving SES within 60 minutes, with the strict invariant that no recipient receives the same campaign twice, ever, regardless of scheduler crashes, region failovers, Kafka rebalances, fan-out worker deaths, SES throttles, or operator-initiated catch-up replays.

This is the simplest-sounding distributed-systems problem that turns out to be one of the hardest. Three things define it: (1) at-least-once everything except the very tail — the scheduler, the fan-out controller, the Kafka pipeline, every retry, every catch-up replay is at-least-once because making them exactly-once is intractable; (2) a durable per-recipient idempotency tail that collapses at-least-once into at-most-once-per-recipient at the SES send boundary; (3) a rate-paced catch-up controller that knows the difference between "fire the missed run now" and "the email is too stale to send" — a 9 AM "good morning" email arriving at 1 PM is reputation damage, not delivery.

The architecture comes straight from Google's Borgcron paper (Sundheim, ACM Queue 2015) for the scheduler tier, LinkedIn ATC and Pinterest NEP for the fan-out tier, Stripe / Brandur's idempotency-key contract for the dedupe tier, and DoorDash's Cadence-as-fallback pattern for the catch-up reconciler. Every load-bearing decision below has a published precedent — this is not speculation.

The reference architecture

Reference architecture for Distributed Cron — Mass Scheduled Email: 18 components — Trigger Source, API Gateway · WAF + Auth + Rate-limit, Service · Campaign Acceptor, SQL · Campaign Store + Trigger Log, Scheduler · Borgcron-style, Coordinator · etcd Lease, NoSQL · Recipient Store, Worker · Fan-out Controller, Stream · Kafka Spine, Worker · Render & Dispatch, Cache · Hot Dedupe (Redis Cluster), NoSQL · Idempotency Store (LWT), NoSQL · Suppression (TCPA-grade), External · Amazon SES, Service · Webhook Receiver (SES), Worker · Catch-up Reconciler (Temporal), Analytics DB · ClickHouse, Tracing · OTel Collector — connected by 31 flows.Trigger SourceclientAPI Gateway · WAF + A…Envoy 1.30 + Cloudflare…Service · Campaign Ac…Go + sqlx + librdkafka …SQL · Campaign Store …Postgres 16 + Citus 12 …Scheduler · Borgcron-…Custom Go service + Raf…Coordinator · etcd Le…etcd 3.5 (5-node Raft, …NoSQL · Recipient Sto…ScyllaDB 5.4, RF=3, LOC…Worker · Fan-out Cont…Go + Kafka consumer + S…Stream · Kafka SpineApache Kafka 3.7 (KRaft…Worker · Render & Dis…Go + handlebars rendere…Cache · Hot Dedupe (R…Redis 7 cluster mode, 1…NoSQL · Idempotency S…Cassandra 4.1, RF=3, LO…NoSQL · Suppression (…Cassandra 4.1, RF=3, LW…External · Amazon SESAmazon SES (per-region …Service · Webhook Rec…Go + AWS SDK v2 (SNS-si…Worker · Catch-up Rec…Temporal 1.22 workflows…Analytics DB · ClickH…ClickHouse 24.x (replic…Tracing · OTel Collec…OpenTelemetry Collector…
18 components, 31 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Trigger Source

Two distinct callers share this entry point: (a) the campaign-owner UI / internal service POSTing a campaign definition (template_id, recipient_list_ref, scheduled_at, priority) — happens dozens to hundreds of times per day; (b) the scheduler tier ITSELF self-firing at scheduled_at — happens internally and does not flow through the public gateway. We model both as 'clients' of the system because both originate writes that mutate state; the scheduler's self-fires are drawn as a separate edge from the scheduler node directly into the Kafka spine.

Why it exists. Considered modeling the cron tick as an implicit internal-only event (not a node). Rejected because the load-bearing semantic of this whole problem is 'a single event becomes 50M', and hiding the originator obscures where the at-least-once-then-dedupe contract begins. Surfacing the trigger as a node makes the Borgcron pattern visible: the trigger LANDS in a durable log before fan-out, and that log is what survives scheduler failover.

When it fails. Naïve callers regenerate campaign_id on retry → ghost campaigns + double-fires. Detection: distinct_campaign_ids_per_logical_submit_p95 > 1 alerts SDK eng. Mitigation: official SDKs ship the canonical retry helper; broken callers caught at gateway via HMAC fingerprint mismatch (Stripe pattern). The cron-tick caller (internal scheduler self-fire) cannot mis-key — its fire_id = (campaign_id, scheduled_at) is deterministic.

API Gateway · WAF + Auth + Rate-limitEnvoy 1.30 + Cloudflare WAF + per-tenant token-bucket

L7 ingress for the campaign control plane and the SES webhook intake. Terminates TLS, runs the WAF (OWASP top-10 + bulk-email-spam detector), enforces token-bucket rate limits per tenant BEFORE forwarding so a misbehaving caller cannot poison the idempotency store, validates the bearer token against the auth service (cached 60 s), copies tenant_id + Idempotency-Key headers verbatim onto the upstream request, and HMAC-signs forwarded headers so the acceptor trusts tenant_id without re-validating. SES bounce webhooks land on a separate path (/webhooks/ses) that verifies SES SNS signatures before forwarding to webhook-rx.

Why it exists. Considered putting auth + rate limit inline in the campaign-api. Rejected because (1) a 5xx in the acceptor would skip the rate limiter and let a retry storm reach Cassandra LWT — at $0.000005/LWT write, even a small storm wastes budget AND triggers Paxos contention; (2) per-tenant rate-limit state lives in Redis which we don't want to ship to every campaign-api replica. Stripe's 2019-07-10 retro is the precedent: rate limits in the wrong layer cascade across product surfaces.

When it fails. Edge pool exhaustion masks as 'campaign owner says we're slow' but is actually upstream load. Detection: gateway_active_streams, downstream_5xx_rate. Mitigation: drain mode + per-AZ 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.

Service · Campaign AcceptorGo + sqlx + librdkafka (idempotent producer)

Hosts POST /v1/campaigns (create + schedule) and PUT /v1/lists/{id}/append (chunked recipient-list staging). For create: (1) validates the template renders cleanly against a 100-recipient golden-dataset canary; (2) writes the campaign row to Postgres (UNIQUE constraint on campaign_id is the producer-side dedupe — duplicate submit returns the cached response_id); (3) registers the campaign with the scheduler tier (writes a pending_fire row in the trigger log). For list-append: streams chunks of recipients into ScyllaDB partitioned by (tenant_id, list_id, shard). Returns 201 in <200 ms for create; bulk-append is server-streaming with backpressure.

Why it exists. Considered fronting Postgres directly with a thin gateway plus a separate dedup worker. Rejected because the dedupe-then-publish ordering MUST be tight — a duplicate campaign_id submit must collapse to one row BEFORE the scheduler can pick it up, or we double-fire. Considered making the scheduler register triggers directly. Rejected because the scheduler is a sensitive Paxos cluster; exposing it to public traffic risks a CVE in the scheduler crashing the cron tier. The acceptor is the blast-radius firewall between public ingress and the scheduler.

When it fails. Postgres campaign-store down → fail-CLOSED 503 on campaign-create (better to refuse than half-create). ScyllaDB recipient-store down → list-append fails partway; the saga marks list_state='failed' via the janitor (NOT inline — saves the next caller's chunk upload from blocking on Postgres). Client retry with the same list_id lands chunks via UPSERT, so re-uploaded rows replace stale ones idempotently. The 'half-uploaded list with abandoned campaign' scenario (Critic Round 1 issue 5a): janitor reaps orphan upload state every 30 min; campaign_create rejects 'failed' lists at the API layer; if ops re-classifies a list mid-life the orphan-check worker catches the dangling campaign. Detection: campaign_create_p99 (alert >500 ms P2), list_append_chunk_throughput (alert <50 MB/s/replica P2), list_orphan_uploading_age_minutes_p99 (alert >60 min P2 — the janitor is itself the safety net here, monitor that it's running). Canary-render-failed deploy (Twitter '{firstName}' literal) prevented by golden-dataset render-diff in CI + a 100-recipient pre-flight render.

SQL · Campaign Store + Trigger LogPostgres 16 + Citus 12 (sharded by tenant_id, 16 shards)

Two logical tables: (1) campaigns(tenant_id, campaign_id, template_id, list_id, scheduled_at, priority, status) — the campaign definition; (2) trigger_log(tenant_id, campaign_id, scheduled_at, fired_at, fired_by_leader, fire_status, paxos_term) — the Borgcron-style durable log of every fire decision. Sharded by tenant_id (Citus distributed table) so a single tenant's hot campaign list doesn't bottleneck another tenant. The trigger log is the load-bearing artifact for catch-up: the scheduler reads it on failover to decide what to re-fire.

Why it exists. Considered putting the trigger log in the scheduler's local Raft state. Rejected because (a) the log must survive scheduler-tier full-cluster loss (DR scenario), (b) the catch-up controller needs to query it (SQL JOIN with delivered counts), (c) on-call needs ad-hoc 'what did we fire last Tuesday at 9 AM' queries which Raft state can't answer. Considered DynamoDB. Rejected because Citus gives us per-tenant SQL semantics + transactional pending_fire → fired state transitions with a single UPDATE — much harder in DynamoDB without conditional writes plus a status-index GSI. Postgres + Citus is what Notion (eng-blog 2023), Heap, and Cloudflare DO Resilience use for this shape.

When it fails. Per-shard primary failure → 30 s write gap for that 1/16 of tenants; reads continue from sync follower. Detection: pg_repl_lag_seconds (alert >100 ms P2), patroni_master_role per shard. Cross-shard query for catch-up (SELECT * FROM trigger_log WHERE fire_status='pending') is fan-out across all 16 shards — bounded latency ~200 ms, acceptable for the reconciler. Schema migrations use pg_repack-style online rewrites + Citus's two-phase metadata commit to avoid downtime.

Scheduler · Borgcron-styleCustom Go service + Raft (3 replicas across 3 AZ)

Leader-elected (Raft, 3 replicas) scheduler that wakes every second, reads pending_fire rows whose scheduled_at <= now() from the campaign-store, and for each one: (1) takes the lease from the lease coordinator (etcd) to assert it's still the active leader; (2) writes a firing row to the trigger log (durable write, gates everything that follows); (3) publishes a campaign.fired event to Kafka with idempotent producer enabled; (4) on producer ack, updates the trigger row to fired. On leader failover the new leader scans firing-status rows older than 30 s and DELIBERATELY re-fires them — Borgcron-style at-least-once, leaning on the per-recipient idempotency tail to make duplicates safe. Followers do NOT fire; they just shadow Raft state for failover.

Why it exists. Considered Kubernetes CronJob. Rejected per Vallery Lancey's 2020 retro — K8s CronJob stops scheduling after 100 missed runs and silently skips on concurrency conflicts. Considered Quartz clustered. Rejected because Quartz's row-lock-on-QRTZ_LOCKS approach degrades under load and its 60 s misfire threshold is too coarse for a 1-minute trigger cadence. Considered Temporal as the scheduler. Kept Temporal for the catch-up reconciler (where its workflow semantics shine) but used a dedicated Raft scheduler for the firing tier because the firing decision needs sub-second latency and direct integration with the trigger log — Temporal's queue-then-poll model adds 100–500 ms per fire. Google's Borgcron paper (Sundheim, ACM Queue 2015) is the canonical reference for this exact pattern: a Paxos-replicated leader that durably logs before firing.

When it fails. Leader dies mid-fire → up-to-5 s leadership gap; new leader replays trigger log and may re-fire rows in firing status > 30 s old (deliberate at-least-once). Trigger-log shard (Citus) primary failover (~30 s) is longer than the 1 s scheduler tick — mitigation: scheduler treats a single tick's trigger_log write failure as 'skip this tick, retry the WHOLE fire next second' (NOT block on the shard). If the shard outage exceeds the lease TTL (5 s), the scheduler voluntarily releases the lease and another replica takes over; the surviving replica's tick will also hit the same shard outage, but the catch-up controller's 5-min reconcile will absorb the missed fires once the shard recovers. Detection: scheduler_leader_count (alert >1 P1 — split-brain), trigger_fire_lag_seconds_p99 (alert >5 s P2), pending_fires_backlog (alert >100 P2, but ONLY when no shard primary is mid-failover — joint condition). Mitigation: etcd lease + fencing token; the in-flight firing row stamps fired_by_leader so the next leader can detect 'someone else was here' and not blindly re-fire. Catch-up controller is the safety net for whole-region scheduler loss. The single failure mode this whole problem exists to prevent (double-fire after partition heal) is mitigated three-deep: (a) etcd lease fences the writer, (b) trigger_log UNIQUE on (campaign_id, scheduled_at) rejects duplicate fires at the storage layer, (c) per-recipient idem-store rejects duplicate sends if (a) and (b) both fail.

Coordinator · etcd Leaseetcd 3.5 (5-node Raft, 3 AZ)

Two leases co-located on the same etcd cluster: (1) /lease/scheduler/leader — the scheduler's leader-election lease, 5 s TTL, 2 s heartbeat; the scheduler holding it is authorized to fire triggers, others stand by. (2) /lease/region/{tenant_id} — per-tenant active-region fence, 30 s TTL, 10 s heartbeat; dispatchers and catch-up workers in the active region renew it, dispatchers in the standby region check it BEFORE sending and fail-CLOSED if they don't hold it. Operator-flippable for planned failover via a typed runbook.

Why it exists. Considered DNS-based region failover. Rejected because DNS TTL + client-side resolver caching means both regions can be active for several minutes during a cutover — a guaranteed double-send window. Considered putting the lease in Postgres with SELECT FOR UPDATE. Rejected because Postgres failover (30 s) is longer than the lease TTL (5 s) and would cascade scheduler unavailability into a wider outage. Considered using Postgres advisory locks. Same problem. etcd / Consul / ZooKeeper are the right tier for sub-second leader election with linearizable reads — that's why the entire industry uses them for this exact purpose (Kafka KRaft, Kubernetes, HashiCorp Vault, every modern coordination tier).

When it fails. etcd quorum lost (3 of 5 nodes failed) → all writes fail-CLOSED → scheduler halts (no lease = no fire) and dispatchers halt (no fence = no send). Severe SPOF — we accept this over the alternative (a weaker fence) because losing 5 minutes of campaigns is recoverable via catch-up, but a 5-minute double-send window is not (TCPA + reputation damage). Detection: etcd_has_leader, etcd_proposals_committed_total, lease_renewal_lag_seconds (alert >2 s P1). Mitigation: 5-node cluster across 3 AZs absorbs single-AZ loss + 1 unrelated node; cross-region observers buy an operator-driven 10-min DR path.

NoSQL · Recipient StoreScyllaDB 5.4, RF=3, LOCAL_QUORUM, 24 nodes

Stores per-tenant recipient lists: (tenant_id, list_id, segment, recipient_id, email, locale, opt_in_proof_id, personalization_blob). Partition key = (tenant_id, list_id, segment) where segment = hash(recipient_id) % 1024 — so a single 50M list is naturally chunked into 1024 partitions of ~50K rows each, each fitting comfortably in a Scylla partition's recommended size. Clustering key = recipient_id. The fan-out controller paginates by segment for predictable parallelism: 1024 segments × ~50K recipients = 50M, fanned out by 64 concurrent workers reading 16 segments each.

Why it exists. Considered Postgres + Citus. Rejected because at 200M MAU × 5 lists avg = 1B rows, the cross-shard scan on WHERE list_id=? is expensive even with sharding; ScyllaDB's per-partition contiguous read is ~50× faster for this access pattern (Discord's 2023 Cassandra→Scylla migration retro is the precedent — they hit the exact same bottleneck on message history reads). Considered DynamoDB. Rejected because BatchGetItem is capped at 100 items, and a list with 50K items per segment would require 500 RoundTrips per segment to read fully — Scylla streams it as one range read in 200 ms. Considered S3 (recipient list as JSONL on object storage). Rejected because we need to UPDATE individual recipients (opt-out propagation, email change) and S3 doesn't support partial updates. The right pick is a wide-column store with range-scan + per-row update — that's Scylla/Cassandra.

When it fails. Single node loss → RF=3 LOCAL_QUORUM survives (W=R=2 still met); fan-out latency unchanged. Two-node loss in same AZ → some partitions go to ONE which we refuse on writes (Monzo 2019-07-29's published lesson — never relax CL=ONE on writes). Hot partition (a tenant's '100M-recipient' list creates 1024 segments of 100K each, each at the upper edge of Scylla's partition size) — mitigation: list-append API rejects appends that would push a segment over 200K rows, forces the caller to spread across more list_ids. Detection: scylla_partition_size_p99, scylla_compaction_pending, read_repair_rate. Cross-region replication lag (stream-mirroring) does NOT affect fan-out because each region runs its own recipient-store; lists are uploaded redundantly per region by the saga at campaign-create time.

Worker · Fan-out ControllerGo + Kafka consumer + Scylla driver

Consumes campaign.fired events from Kafka. For each fired campaign: (1) reads the campaign metadata + template_id + list_id from campaign-store; (2) for each of the 1024 list segments, BFS-reads ScyllaDB by segment (range scan ~5 MB / 50K rows / 200 ms); (3) for each recipient batch, produces N send.requested Kafka events keyed by (tenant_id, recipient_id) with idempotent producer enabled; (4) writes a checkpoint (campaign_id, segment, last_recipient_id, batch_acked_at) to campaign-store every 1000 recipients so resume-after-crash skips already-enqueued recipients. The controller is sticky-partitioned: one Kafka partition of campaign.fired belongs to one fanout replica, so a single campaign is owned by one worker (in-order checkpoint advancement) but 64 concurrent campaigns spread across the 32-replica fleet.

Why it exists. Considered doing the fan-out inline in the scheduler. Rejected because the scheduler's hot path is sub-200 ms; a 60-min fan-out for 50M recipients would block the scheduler tier completely (no other campaigns fire). Considered making the dispatcher do the recipient lookup per-event (no fan-out controller). Rejected because dispatchers would then read ScyllaDB once per recipient (50M random-key reads/campaign = catastrophic) vs the controller's 1024 segment range-scans. Considered Temporal for the fan-out itself. Plausible — Temporal handles long-running workflows well. We use it for catch-up (where workflow durability is load-bearing) but use a plain Kafka consumer for the fan-out tier because (a) fan-out is high-throughput stateless work that benefits from Kafka's partition-level parallelism, and (b) the checkpoint in campaign-store is durable enough — we don't need Temporal's workflow state. LinkedIn's ATC and Pinterest's NEP both run the same pattern: stateless fan-out workers + external checkpoint store.

When it fails. Single replica crash mid-fan-out → consumer-group rebalance hands the partition to another replica which resumes from the last checkpoint (worst case 5 s / 1 K recipients of re-enqueue, harmless under dedupe). ScyllaDB segment read slow (>1 s) → backs up the per-partition pipeline; mitigation: per-replica concurrency limit (16 concurrent segment reads) prevents one slow segment from blocking the whole replica. Detection: fanout_progress_recipients_per_sec (alert <80% of campaign target P2), fanout_checkpoint_lag (alert >10 K recipients P2). Kafka producer back-pressure → the controller paces its read rate to match Kafka's accept rate, never overruns the broker. The classic 'crash after enqueueing 23M of 50M' scenario (Failure Mode #3 from the spec) is exactly the case checkpointing solves — resume is bounded by the checkpoint window.

Stream · Kafka SpineApache Kafka 3.7 (KRaft mode, 24 brokers, 3 AZ)

Single Kafka cluster carrying every event the cron tier produces. Topics: campaign.fired (64 partitions, post-scheduler) → consumed by fanout; send.requested.bulk and send.requested.transactional (separate topics for bulkheading — a bulk blast must not back up a transactional 'password reset', 128 partitions each) → consumed by dispatchers; send.retry (32 partitions, post-circuit-breaker delayed retries with timestamp-keyed exponential backoff); send.delivered and send.failed (64 partitions, post-SES result, consumed by analytics-db); provider.feedback (32 partitions, post-webhook, consumed by suppression-store + analytics). Partitioned by (tenant_id, recipient_id) so per-recipient ordering is preserved (a 'pending' event precedes the 'delivered' for the same recipient).

Why it exists. Considered SQS / RabbitMQ. Rejected because (a) SQS FIFO groups cap at 300 msg/sec/group which can't carry our 14 K msg/s peak without painful sharding gymnastics, (b) replay-from-offset for catch-up backfills is impossible with SQS's 14-day point-of-no-return retention, (c) multi-consumer fan-out (dispatcher + analytics + audit all reading the same send.delivered) wants Kafka's pub-sub model, not point-to-point. Per LinkedIn ATC, Pinterest NEP, Shopify (which peaks at 66M msg/s on Kafka), DoorDash, Discord — Kafka is the production answer for this shape.

When it fails. Single broker loss → controller re-elects leaders for affected partitions, ~5 s gap on those partition writes. min.ISR=2 → single-broker loss never blocks; two-broker loss in same AZ blocks writes (correct — better than risking durability). Hot partition (one tenant blasts a single mega-list — all 50M recipients hash to nearby partitions) → mitigation: bulk topic uses (tenant_id, recipient_id) key so a single tenant spreads across all 128 partitions; per-tenant rate limit at the gateway prevents the noisy-neighbor scenario. 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.

Worker · Render & DispatchGo + handlebars renderer + SES SDK v2 + circuit-breaker

The send tail. Consumes send.requested.{bulk,transactional} from Kafka. For each event: (1) reads template from campaign-store (cached 5-min in-process, hit ratio >99%); (2) checks the active-region fence in etcd — if not ours, drops the event silently (active region will pick it up); (3) checks suppression store for opt-out / hard-bounce on this recipient — if suppressed, skips and publishes a send.failed{reason=suppressed} event; (4) checks idem-cache for a HIT on (tenant_id, campaign_id, recipient_id) — if HIT, this is a duplicate enqueue from a Borgcron re-fire or a fanout-resume; skip and publish dedupe-skipped metric; (5) on MISS, CLAIM the idem-store via Cassandra LWT (INSERT ... IF NOT EXISTS); (6) on CLAIM success, render the template with the recipient's personalization data; (7) call SES with MessageDeduplicationId = sha256(campaign_id, recipient_id); (8) on SES success, UPDATE idem-store status to sent + publish send.delivered. The render+SES call IS the critical path for the user's 'when does my email arrive' SLO.

Why it exists. Considered making this a separate render service + dispatch service. Rejected because the render → CLAIM → SES sequence is single-threaded per recipient and splitting it adds a Kafka hop (50 ms+) per send — at 14 K sends/s that's 14 K × 50 ms = 700 s of accumulated added latency per peak campaign. Considered putting rendering on the fanout controller. Rejected because rendering needs the per-recipient personalization data which lives in the recipient-store — the fanout controller already reads it for the recipient_id, but rendering also needs the template fetch and the suppression check, which happen here. Single service, sequential pipeline, is correct.

When it fails. SES 5xx storm → circuit-breaker opens, in-flight requests divert to send.retry topic with exp-backoff (1 s, 5 s, 30 s, 5 min, 30 min, 6 h, then DLQ). Workers self-throttle so the storm doesn't poison IP reputation. Detection: ses_5xx_rate (alert >2% P2, >10% P1), dispatcher_circuit_open (alert any P2), send_p99_seconds. Idem-store down → fail-CLOSED (better to delay than dup); the failed event sits in Kafka and the worker retries on next poll. Suppression-store down → fail-CLOSED on send (TCPA compliance is non-negotiable; a missed send is recoverable, a TCPA fine is not). Render error on a single template → CI golden-dataset render-diff prevents most; runtime errors land in send.failed{reason=render-error} for the campaign-owner to inspect.

Cache · Hot Dedupe (Redis Cluster)Redis 7 cluster mode, 16 shards, 1 follower/shard

Caches idem:{tenant_id}:{campaign_id}:{recipient_id} → claim_status for the SET of currently-active or recently-active campaigns (typically the last 24 h). Populated WRITE-THROUGH by the dispatcher when it WINS a CLAIM at idem-store. Read by the dispatcher on every send as the fast-path dedupe — a HIT means 'already CLAIMED in this region, skip', a MISS means 'fall through to Cassandra LWT'. Eviction is allkeys-LFU + 14-day TTL. NOT the source of truth (Cassandra is); the cache exists ONLY to take the load off Cassandra LWT for the 95% of dispatches that are first-time-sends.

Why it exists. Considered making Cassandra the only dedupe store. Rejected because Cassandra LWT (Paxos) costs ~30 ms per write and ~5 ms per read — at 14 K sends/s that's 14 K × 30 ms = 420 s of accumulated LWT latency for one campaign, plus heavy CPU on the Cassandra nodes. Two-tier dedupe is the standard pattern: hot cache for the 95% of straight-through dispatches + durable LWT for the 5% of conflicts and the eventual source-of-truth. Considered making this a per-worker in-process cache. Rejected because (a) cache hit ratio depends on cross-replica sharing — in-memory partitions by replica and gives ~2% hit rate at our 50-replica fleet, (b) a re-fire from Borgcron failover MUST hit a shared cache to dedupe, not a per-replica one.

When it fails. Single shard down → 1/16 of dispatches fall through to Cassandra LWT (5–10× slower, ~30 ms/op vs ~1 ms/op) — degraded but correct; Cassandra absorbs the spike at ~10× the cost. Full cluster restart → ALL dispatches fall through; Cassandra LWT saturates within seconds — page on idem_cache_hit_ratio < 0.85. Mitigation: cluster-mode auto-failover handles single-shard within 15 s; full-cluster restart is rare (last incident: never in this design's history). Hot key: a single campaign's recipients all carry (tenant_id,campaign_id) hash tag → cluster pins them to ONE slot pair → single shard absorbs 100% of the campaign's lookups. Mitigation: at >10K sends/s on a single campaign, the dispatcher salt-shards by (tenant_id, campaign_id, recipient_id % 16) instead of by (tenant_id,campaign_id) (per Failure Mode #5 in the spec). NEVER trust this cache for opt-out semantics — those go straight to suppression-store via LWT.

NoSQL · Idempotency Store (LWT)Cassandra 4.1, RF=3, LOCAL_QUORUM + LWT (Paxos) on CLAIM

Source of truth for per-recipient dedupe. Schema: PK (tenant_id, campaign_id, recipient_shard) where recipient_shard = hash(recipient_id) mod 8, CK recipient_id, columns (claim_status: claimed|sent|failed|suppressed, claimed_at, claimed_by_pod, ses_message_id, attempt_count, ttl). The dispatcher CLAIMs a recipient via Cassandra LWT: INSERT INTO idem (...) VALUES (...) IF NOT EXISTS — Paxos-coordinated linearizable write. Wins → row inserted, dispatcher proceeds with SES call. Loses (applied=false) → row already exists, dispatcher dedupes the send. On SES success, dispatcher UPDATEs claim_status='sent', ses_message_id=? (no LWT on the update, only the initial CLAIM). 14-day TTL on every row. The recipient_shard mod 8 in the PK is the load-bearing decision: without it, a single 50M-campaign first-fire lands 14 K LWT/s on ONE Cassandra partition coordinator and exceeds the ~5 K/s/partition Paxos ceiling by 2.8×; with mod 8, the load spreads across 8 partition groups (~1.7 K/s each, comfortable). This is not a 10× scale fix — it is needed on day one and the Critic Round 1 review flagged the original single-partition design as broken under first-fire.

Why it exists. Considered Redis SETNX as the source of truth. Rejected because Redis async failover has RPO=5 s — 5 seconds of dedupe state vanishing on AZ outage is enough to leak millions of dupes for an in-flight 50M campaign (Failure Mode #4 in the spec). Considered DynamoDB conditional PutItem. Plausible — AWS Builders' Library ARC403 endorses this exact pattern. We picked Cassandra over DynamoDB because (a) catch-up controller needs RANGE SCANS over (tenant_id, campaign_id, recipient_shard) to compare intended vs delivered — Cassandra's clustering key gives O(log N) range read across the 8 shards (8 small scans, parallelized); DynamoDB needs a GSI or N×GetItem; (b) the team already runs Cassandra for suppression — operational compounding. Both Cassandra-LWT and DynamoDB-conditional-PutItem are defensible; the engineering trade is operational topology vs managed-service convenience.

When it fails. Single node loss → RF=3 LWT survives (Paxos quorum is 2 of 3); CLAIM latency unchanged. Two-node loss in same AZ → some partitions have only ONE replica, CLAIM blocks (correct — better than risking dupe). 8-way recipient-shard prevents the single-partition hot-coordinator problem on first-fire (1.7 K LWT/s/partition vs ~5 K/s/partition ceiling). Orphan claimed rows from a dispatcher crash between CLAIM-win and SES-success — janitor reclaims every 5 min by re-publishing to send.retry; the second dispatcher gets applied=false on its CLAIM (the row still exists) and proceeds to SES with the existing CLAIM, no double-send. Detection: lwt_p99_seconds (alert >100 ms P2), lwt_paxos_contention_rate (alert any P2 — applied=false is NOT a failure, only LWT timeouts/errors count), cassandra_pending_compactions, idem_claimed_orphan_age_seconds_p99 (alert >10 min P2). Cross-region replication via async stream-mirror; on planned failover, catch-up does cross-region range-scan before re-dispatch so already-sent recipients are skipped (~80 ms p99 cross-region read penalty, paid once per failover not per send).

NoSQL · Suppression (TCPA-grade)Cassandra 4.1, RF=3, LWT on opt-out writes

Source of truth for 'do NOT send to this recipient'. Three categories: (1) opt-outs — user clicked unsubscribe (CAN-SPAM requires processing within 10 business days, we do it inline within seconds); (2) hard-bounces — recipient address doesn't exist (SES feedback); (3) complaints — recipient flagged as spam (SES feedback). Schema: PK (channel, identifier) where channel='email' and identifier is the lowercased email address. Opt-out writes use Cassandra LWT (INSERT ... IF NOT EXISTS) for linearizable cross-region consistency. Reads on the send path use LOCAL_QUORUM (W=R=2) — strong within region, sufficient because regional ownership ensures the recipient was last seen by the home region.

Why it exists. Considered keeping suppression in Postgres + Citus alongside campaigns. Rejected because suppression reads are on the send hot path (14 K/s peak per campaign, 1 read per send → 14 K reads/s sustained per active blast), and Citus's cross-shard fan-out for WHERE email=? is wasteful for a point-lookup. Considered making it a Redis-only cache. Rejected because TCPA statutory damages are $500–$1500 per opt-out violation — the cost of a cache miss on opt-out is too high; the source of truth must be durable with linearizable opt-out writes. Considered DynamoDB. Plausible — but we already run Cassandra for idem-store, and LWT semantics are identical.

When it fails. Suppression-store down → dispatcher fail-CLOSED on send (a missed send is recoverable, a TCPA violation is $46 K). Detection: suppression_check_p99 (alert >50 ms P2), suppression_lwt_timeout_or_error_rate (alert any P1 because that means opt-outs are not landing — direct legal exposure). NOTE: applied=false from LWT (the recipient already opted out) is NOT a failure and is excluded from this metric — alerting on raw applied=false would alert-storm during every mass-unsubscribe event. Hot key (mass-unsubscribe campaign) — mitigation: opt-out clicks are rate-limited at the unsubscribe-handler service to <1 K/s per tenant, so suppression LWT contention is bounded. Cross-region LWT failures during a region partition → opt-out writes block (correct — better to delay confirmation than to silently lose the opt-out).

External · Amazon SESAmazon SES (per-region sub-accounts + dedicated IP pools)

Third-party email delivery. We hold multiple SES configurations: (a) transactional sub-account in us-east-1 with a dedicated IP pool warmed to the 50 K/s tier — for 'password reset' / 'order shipped' / '2FA' traffic; (b) marketing-bulk sub-account in us-east-1 with a separate dedicated IP pool warmed to the 14 K/s tier — for the 50M-recipient blasts; (c) marketing-bulk-west sub-account in us-west-2 with its own dedicated IP pool — for regional failover and load-spreading on truly massive campaigns. Each call: SES validates SPF/DKIM/DMARC alignment, scans for spam-triggering content via its content classifier, accepts the message, returns a MessageId, then handles MTA-level retries and bounce processing on its side. Bounces and complaints land at our webhook-rx via SES → SNS → HTTPS.

Why it exists. Considered self-hosted Postfix + IP block. Rejected because IP-reputation management is a full-time team — SES + SendGrid + Mailgun do it for us, and recovery from a reputation hit takes weeks. Considered SendGrid as primary. Plausible — SendGrid's published 1.5M/hr ≈ 416/s sustained is workable; we keep SendGrid as the secondary provider for failover. SES wins as primary because it's the AWS-native cheapest option at scale ($0.0001/email).

When it fails. SES 5xx storm (regional API degradation) → dispatcher circuit-breaker diverts to send.retry topic with exp-backoff; campaign throughput drops; if storm lasts >30 min, the catch-up controller will detect (delivered count vs intended count diverges) and pace replay. Bounce rate spike (>5%) → SES auto-pauses the sub-account; we get a Reputation.BounceRate CloudWatch alarm. Mitigation: pre-flight list hygiene (suppress hard-bounces from previous campaigns), automated pause when bounce rate >3%, sub-account isolation prevents whole-tenant outage. Detection: ses_throttle_rate (alert >2% P2), ses_reputation_bounce (alert >3% P1), ses_reputation_complaint (alert >0.05% P1). 429 backoff cascade (Failure Mode #6) — AIMD token bucket per-pool plus circuit-breaker open is the documented mitigation.

Service · Webhook Receiver (SES)Go + AWS SDK v2 (SNS-signed payload verify)

Receives bounce, complaint, and delivery webhooks from SES via SNS HTTPS POST. For each: (1) verifies the SNS signature (X.509 cert validated against AWS's published cert URL, cached); (2) parses the payload (bounce_type, recipient, original_message_id); (3) dedupes against (provider, provider_event_id) in idem-cache; (4) publishes a provider.feedback event to Kafka; returns 200 OK in <50 ms (SES retries on non-200 with backoff up to 24 h). Does NOT do the suppression-store write inline — that happens async in the suppression-update worker that consumes provider.feedback. Keeping this tier minimal absorbs bounce storms without backing up.

Why it exists. Considered having the dispatcher poll SES for delivery results. Rejected because SES doesn't support polling at scale — webhook is the only practical feedback channel. Considered making the webhook receiver write suppression-store synchronously. Rejected because a bounce storm (12M bounces in 10 minutes after a deliverability hit — Failure Mode #13) would back up Cassandra LWT and 503 the webhook, which would cause SES to retry, which would compound the storm. Decoupling via Kafka means the webhook is a thin durable hand-off; suppression updates drain at their own pace.

When it fails. SES webhook payload validation fails (cert rotation issue, malformed JSON) → log + return 400 (do NOT 500, or SES retries forever). Kafka publish fails → 500 to SES → SES retries with backoff. Webhook storm (bounce campaign hits 12M bounces in 10 min) → consumer-side backpressure absorbs via Kafka buffering — webhook-rx itself stays p99 <50 ms because all it does is verify + publish. Detection: webhook_p99 (alert >100 ms P2), webhook_5xx_rate (alert >1% P1), webhook_dedupe_hit_rate. Mitigation: SES has its own 24-h retry budget; a brief webhook-rx outage is recoverable via SES's retry.

Worker · Catch-up Reconciler (Temporal)Temporal 1.22 workflows + custom diff worker

Long-running Temporal workflow per campaign. Triggered (a) on a schedule every 5 min for in-flight campaigns; (b) on-demand by ops after a known outage (scheduler down, region failover, SES outage). Algorithm: (1) acquire the catch-up lease in etcd at /lease/catchup/{campaign_id} (prevents two workflows reconciling the same campaign); (2) read intended recipient count from campaign-store (SELECT segment, recipient_count FROM list_segments WHERE list_id=?); (3) read delivered/claimed state DIRECTLY from idem-store via range-scan across all 8 recipient_shard partitions for the campaign (SELECT recipient_id, claim_status FROM idem WHERE (tenant_id, campaign_id, recipient_shard) IN (?, 0..7)) — idem-store is the source of truth for 'did we send this'; analytics-db is informational only and is NOT read on the reconciler path (Critic Round 1 issue 6c: analytics-db lag would false-positive-trigger a send.retry storm — direct idem-store read avoids the circular dependency); (4) for each recipient with claim_status NOT IN ('sent', 'suppressed', 'failed'), publish a send.retry event with attempt_origin='catch-up' — at a paced rate of ≤ 50% of baseline send rate (so SES isn't overwhelmed); (5) when intended ≈ (sent + suppressed + failed), mark campaign complete. Workflow can resume from any step on worker crash (Temporal's durable workflow state).

Why it exists. Considered making the scheduler do catch-up. Rejected because scheduler hot path is sub-200 ms; catch-up scans take minutes to hours and don't belong in that latency budget. Considered a cron-triggered batch job. Rejected because catch-up needs to resume mid-scan after a worker crash (a 50M-recipient diff takes hours to complete), and a plain cron job has no durable state — it would restart from scratch on every crash. Temporal's workflow durability is the right abstraction: write step-by-step state to Temporal's history, resume on crash, retry per-step. DoorDash's published 'Cadence as fallback' post is the precedent for using a workflow engine for exactly this 'safety-net reconciler' role.

When it fails. Catch-up workflow stuck waiting on a slow analytics-db query → Temporal's per-activity timeout fires, retries on a different worker, eventually escalates to ops via a P2 alert. Two catch-up workflows for the same campaign (lease race) → etcd lease prevents (whichever workflow gets the lease wins; the other waits or fails-fast). Catch-up storm (Failure Mode #9 in the spec) — pacing prevents; if pacing is mis-configured, the SES throttle + circuit-breaker absorbs the secondary storm. Detection: catchup_workflow_age_seconds (alert >2 h P2), catchup_lag_recipients (alert >100 K behind expected P2), catchup_relevance_drops (informational metric for the campaign-owner UI). Temporal cluster down → catch-up halts but happy-path is unaffected; ops can replay missed catch-ups manually from trigger_log.

Analytics DB · ClickHouseClickHouse 24.x (replicated, 12 nodes, ZooKeeper-Keeper)

Consumes send.delivered, send.failed, and provider.feedback from Kafka via the ClickHouse Kafka engine. Stores per-event rows: (event_time, tenant_id, campaign_id, recipient_id, segment, channel, provider, message_id, status, error_code, attempt_count). Materialized views aggregate to: (a) per-campaign delivered/failed/bounced/complained counters (powering the campaign-owner UI); (b) per-tenant rolling bounce rate (driving the SES reputation alarm); (c) per-segment delivery progress (informational input to the catch-up controller's freshness check — the reconcile diff itself reads idem-store). 30-day hot retention, 4-year cold retention via S3 + Parquet for CAN-SPAM audit.

Why it exists. Considered Postgres for analytics. Rejected because aggregate queries over 50M-recipient × 4K-campaign tables are O(N) scans that crush an OLTP engine; ClickHouse's columnar engine answers them in ~100 ms vs minutes. Considered putting analytics in Snowflake / BigQuery. Plausible for the cold-archive tier; we use S3 + Parquet because we already export there for compliance, and ClickHouse can query S3 Parquet via its s3() function for the rare cold-archive query.

When it fails. ClickHouse node loss → reads fail over to replica within 60 s; ingest pauses on the affected shard for 30 s. Kafka engine resumes from the last committed offset on restart, so no data loss for ingest. Cardinality storm if a campaign accidentally includes recipient_id as a dimension on a metric (Failure Mode #20) — mitigation: per-tenant series limits at the ClickHouse ingest layer, plus we explicitly forbid recipient_id as a dimension at the metric library level. Catch-up controller depends on analytics-db reads being fresh; a 60-s lag is acceptable, anything beyond → catch-up workflow retries until fresh.

Tracing · OTel CollectorOpenTelemetry Collector + Tempo backend

Receives OTLP spans from scheduler, fanout, dispatcher, webhook-rx, catch-up workflows. Per-trace context propagated via traceparent header on every Kafka message (in the message headers, not the body — keeps deserialization independent of tracing). Stored 7 days hot in Tempo, sampled to 4 weeks for high-cardinality investigations. Powers Slack's tracing-notifications pattern (eng-blog 2017): a customer support agent pastes recipient_id + campaign_id into a forensics UI and gets back the full timeline — fire → fanout → dispatch → SES → delivery (or where it went wrong).

Why it exists. Considered relying on logs alone. Rejected per Slack's 2017 post — for fan-out systems, logs without trace correlation are unsearchable at scale; on-call spends 30+ min reconstructing 'what happened to recipient X' from 6 different services' logs. Tracing makes this a 30-second lookup. Considered head-based sampling at 1% to control cost. Rejected because we need 100% sampling on send.failed and the catch-up workflows — tail-based sampling at the OTel collector keeps the per-trace decision local while still keeping the rare-and-interesting traces.

When it fails. Tracing backend unreachable (Tempo down, network issue) → OTel collector buffers 2 min locally, then drops with a tracing_spans_dropped_total metric (NOT a page — observability is degradable). NEVER on the sync send path (OTel emit is async UDP). Detection: tracing_export_failure_rate (informational P3), tracing_collector_buffer_pct. Critically — we do NOT page on tracing-down during an incident because the same root cause may be taking down monitoring (the Failure Mode #20 scenario: cardinality death takes down both tracing and Prometheus simultaneously). Mitigation: self-monitor tracing via an INDEPENDENT lightweight metrics pipeline.

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:

  • Who's the user? Two distinct users: (a) the campaign owner (marketing ops, transactional-product engineer) who submits campaigns and watches delivery metrics — they need a clean API and a status dashboard; (b) the email recipient who must receive at-most-one of any given campaign — they never see our system, just the email. The whole architecture exists to serve recipient #b at scale.
  • Peak burst shape? Black-Friday-style: one tenant fires a 50M-recipient marketing blast while 100s of small transactional sends ('password reset', '2FA') flow through concurrently. Bulk MUST NOT starve transactional.
  • Catch-up policy? Define per-campaign: a 9 AM 'good morning' email has a 4-hour relevance window (drop past 1 PM); a 'shipment delivered' transactional email has a 24-hour relevance window. The campaign-owner picks at submit time.
  • Compliance? CAN-SPAM ($46 K/email statutory damages on opt-out violation) for marketing; GDPR for EU recipients; opt-in proof required per email for audit.
  • Multi-tenant? Yes — strict per-tenant Kafka bulkheading. A noisy tenant must not back up other tenants' sends.
  • Idempotency contract? Per-recipient: (tenant_id, campaign_id, recipient_id) is the dedupe key. TTL 14 d (≥ 2× max catch-up window of 6 h, with 56× headroom).
  • Multi-region? Single-region scheduler (one Borgcron cluster per environment — simplicity wins over the marginal availability gain); active-active fan-out + dispatch with per-tenant active-region pinning via etcd lease.
  • Latency targets? Trigger-to-first-send p99 ≤ 5 s; trigger-to-95%-delivered ≤ 60 min for 50M-recipient peak; transactional sends p99 ≤ 30 s end-to-end.

Assumptions to state:

  • 4,000 active campaigns/day across the tenant base.
  • One peak 50M-recipient blast per tenant per day, scheduled around 9 AM in the tenant's home timezone.
  • Peak aggregate send rate 14 K emails/s sustained, 25 K/s burst.
  • 200M MAU × 5 lists avg = 1B recipient-list rows in ScyllaDB.
  • 7-day Kafka retention on hot topics; 14-day on send.retry (catch-up may reach back this far).
  • Idempotency hot storage: 1 peak-50M campaign/day × 14 d × 50M × 120 B + 3,999 long-tail × 14 d × 10K avg × 120 B = ~151 GB cluster-wide hot working set on the 14-d TTL window (the original "33 TB" figure assumed every campaign is 50M-recipient — see capacity section for the corrected derivation).

02Functional reqs

What must this system actually do?

  • Accept POST /v1/campaigns from a campaign-owner: validate template, register schedule, return campaign_id.
  • Accept PUT /v1/lists/{id}/append for bulk recipient-list staging (server-streamed chunks into ScyllaDB).
  • Fire a registered campaign at its scheduled_at time, exactly once per logical schedule (at-least-once at the scheduler + at-most-once at the per-recipient tail).
  • Expand the recipient list and dispatch one email per recipient with the rendered template.
  • Apply per-recipient suppression (opt-out, hard-bounce, complaint) BEFORE sending.
  • Retry transient failures with exponential backoff; cap retries at 6 attempts then DLQ.
  • Catch up missed campaigns after scheduler outage / region failover, within a per-campaign relevance window.
  • Ingest SES delivery / bounce / complaint webhooks and update suppression-store.
  • Expose GET /v1/campaigns/{id}/status: returns intended count, delivered count, failed count, in-flight count.
  • Audit-log every send attempt with opt_in_proof_id for CAN-SPAM / GDPR compliance subpoenas.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Idempotency: Per-recipient. No recipient ever receives the same campaign twice under scheduler re-fire, fan-out resume, dispatcher retry, send.retry replay, or catch-up reconcile.
  • Availability: 99.99% on the scheduler tier (the load-bearing 'will it fire' guarantee). 99.9% on the send tier (retries mask blips). 99.95% on the campaign control plane.
  • Latency: Trigger-to-first-send p99 ≤ 5 s. Trigger-to-95%-delivered ≤ 60 min for 50M-recipient peak. Transactional p99 ≤ 30 s end-to-end. Campaign-create p99 ≤ 500 ms (recipient-list upload separate, 5–30 min for 50M).
  • Durability: Once we 201 the campaign-create, the campaign is durably persisted (Postgres sync follower). Once the scheduler 'fires' a campaign, the trigger log is durably committed BEFORE the Kafka publish (outbox pattern). RPO=0 within region; cross-region RPO ≤ 30 s on async mirror (catch-up replays from trigger log).
  • Failover: Per-tenant active-region cutover via etcd lease + dispatcher fencing. RTO 10 min for full-region failover; ZERO duplicate sends during cutover (defense-in-depth: lease + idem-cache + idem-store).
  • Compliance: CAN-SPAM opt-out processed inline within seconds. 4-year audit retention (CAN-SPAM safe-harbor requirement). DKIM/SPF/DMARC alignment enforced per sub-account.
  • Scheduling correctness: No missed fires (catch-up safety net). At-most-one-per-recipient at SES (per-recipient idem). Sub-second leadership failover on scheduler crash (etcd lease).

04Capacity estimation

How much load and data does this have to hold?

Defaults: 4,000 campaigns/day, peak 50M-recipient blast, 60-min delivery target, 15 KB rendered body, 14-day idem TTL.

  • Peak send QPS: 50M / (60 × 60) = 13.9 K emails/s sustained for 60 min. Sits in the SES "warm 14 K/s tier" — achievable per AWS docs on a pre-warmed dedicated IP pool. Going to 15 min instead of 60 min would require 55 K/s, which exceeds SES's max 50 K/s tier without multi-account fan-out — so we picked 60 min as the SLO.
  • Worker count: 14 K sends/s, SES p50 350 ms / p99 800 ms (TLS handshake + DATA), async pods at 200 concurrent SES connections/pod. p50-sized: 14 K / (200 / 0.35) ≈ 25 pods. p99-sized: 14 K / (200 / 0.80) ≈ 56 pods. We provision 50 pods which sits between — p99-feasible with 12% spare under sustained p99 latency, comfortable under steady-state p50, and gives reputation isolation (2 pods/sending-IP across 25 IPs). The remaining gap during a p99-latency excursion is absorbed by Kafka consumer backpressure (events queue, not drop) without violating the 60-min SLO.
  • Recipient-state write QPS: 50M × (1 INSERT enqueue + 1 UPDATE attempt + 1 UPDATE terminal + ~0.05 retries) = ~152M writes / 60 min = ~42 K writes/s sustained, ~80 K peak. Distributed across 36 Cassandra nodes ≈ 1.2 K writes/s/node — well within Cassandra's ~10 K writes/s/node ceiling.
  • Idempotency hot storage (corrected — Critic Round 1 issue 2.1). The original formula campaignsPerDay × idemDays × peakRecipients is wrong by a factor of 4 K because it assumes every campaign is 50M-recipient. Realistic mix: 1 peak-50M campaign/day × 14 d × 50M × 120 B = 84 GB from peak campaigns; 3,999 long-tail campaigns/day × 14 d × 10 K avg recipients × 120 B = 67 GB from long-tail = ~151 GB cluster-wide hot working set across the 14-day TTL window. On a 36-node Cassandra cluster that's ~5 GB/node — within recommended <2 TB/node with 100× headroom before a node-add is needed. The original 33 TB figure assumed 200 peak-50M campaigns/day which is implausible at our tenant base.
  • Idempotency hot CACHE size: Only IN-FLIGHT campaigns are cached (95% hit ratio for active sends, fall-through for resumes). Worst case: 2 concurrent peak campaigns × 50M × 30 B = ~3 GB hot working set. Fits in one Redis cluster shard with headroom.
  • Kafka ingress at peak: 14 K send.requested/s × 500 B (control-plane bytes only — rendered body is NOT in Kafka) = 7 MB/s. With fan-out factor (one campaign.fired becomes 50M send.requested) and RF=3: ~21 MB/s broker write per 50M campaign, ~1 MB/s/broker on 24 brokers. Healthy. Kafka is nowhere near its bottleneck.
  • Bounce webhook QPS: 14 K sends/s × 2% bounce rate = ~280 bounces/s sustained, ~600/s peak. Trivial for the 4-replica webhook-rx tier.
  • Audit annual storage: the naive 4 K × 365 × 50M × 256 B ≈ 18.7 PB/yr assumes every campaign is a 50M blast — the same factor-of-4K error the idem-storage math corrects. Realistic mix (1 × 50M + 3,999 × 10 K ≈ 90M send events/day): 90M × 365 × 256 B = ~8.4 TB/yr raw. ClickHouse + zstd compresses ~5×; the 4-year CAN-SPAM retention window accumulates ~34 TB raw → ~7 TB as S3 Parquet cold archive — a rounding error on the storage bill.
  • The scale tier where each architectural choice flips:
  • Single-Postgres campaign-store works to ~500 campaigns/day before the trigger-log write rate (4 inserts/fire) saturates one shard — we use Citus from day one for blast-radius isolation.
  • Cassandra LWT on a SINGLE partition coordinator works to ~5 K sends/s before Paxos contention dominates → we shard idem-store by (tenant_id, campaign_id, recipient_shard mod 8) from day one (load spreads across 8 coordinators, 1.7 K LWT/s each). The Redis hot cache helps on re-fires + fan-out resumes but does NOT help first-fire (cache MISS on every recipient until CLAIM populates) — the recipient-shard split is the structural fix; the cache is the steady-state optimization.
  • Single-region SES sub-account works to ~14 K sends/s warm-pool tier → we provision a second region's sub-account from day one for failover and to spread truly massive (100M+) campaigns.

05API design

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

POST /v1/campaigns
Authorization: Bearer <tenant-scoped-key>
Idempotency-Key: 8a7b3e2c-...

{
  "tenantId": "merchant_abc",
  "campaignId": "campaign_2026q1_promo",   // V4 UUID, client-supplied for dedupe
  "templateId": "promo_template_v3",
  "listId": "list_us_subscribers",
  "scheduledAt": "2026-05-26T13:00:00Z",
  "priority": "bulk",                       // bulk | transactional
  "relevanceWindowHours": 4                 // drop catch-up past this
}

201 Created
{
  "campaignId": "campaign_2026q1_promo",
  "status": "scheduled",
  "fireAt": "2026-05-26T13:00:00Z",
  "recipientCount": 50000000
}

# Retry with same campaign_id + same body:
201 Created (same campaign_id, idempotent — UNIQUE constraint at campaign-store)

# Retry with same campaign_id (Idempotency-Key) + DIFFERENT body — key reuse:
409 Conflict { "error": "idempotency_key_reuse" }
# (422 stays reserved for a genuinely unprocessable body — bad scheduledAt, unknown priority.)
PUT /v1/lists/list_us_subscribers/append
Content-Type: application/jsonl
Transfer-Encoding: chunked

{"recipientId":"u_001","email":"alice@example.com","locale":"en-US",...}
{"recipientId":"u_002","email":"bob@example.com","locale":"en-US",...}
...

200 OK
{
  "listId": "list_us_subscribers",
  "appended": 50000000,
  "totalRows": 50000000,
  "status": "ready"
}
GET /v1/campaigns/campaign_2026q1_promo/status
200 {
  "campaignId": "campaign_2026q1_promo",
  "status": "in-flight",
  "firedAt": "2026-05-26T13:00:00Z",
  "recipientCount": 50000000,
  "delivered": 47812341,
  "bounced": 1023456,
  "complained": 412,
  "suppressed": 893218,
  "failed": 271573,
  "inFlight": 32500,
  "etaDelivery95Pct": "2026-05-26T13:43:00Z"
}
POST /webhooks/ses              # SES → SNS → here (SNS X.509-signed; verify against AWS cert URL)
# Body schema is SES standard JSON
200 OK   (must be <50 ms or SES will retry)

06Data model

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

Campaign metadata + trigger log (Postgres + Citus, sharded by tenant_id, 16 shards):

CREATE TABLE campaigns (
  tenant_id          uuid NOT NULL,
  campaign_id        uuid NOT NULL,                     -- client-supplied for dedupe
  template_id        text NOT NULL,
  list_id            uuid NOT NULL,
  scheduled_at       timestamptz NOT NULL,
  priority           text NOT NULL,                     -- 'bulk' | 'transactional'
  relevance_hours    int NOT NULL DEFAULT 4,
  status             text NOT NULL,                     -- 'scheduled'|'firing'|'fired'|'complete'|'failed'
  created_at         timestamptz NOT NULL,
  created_by         uuid NOT NULL,
  PRIMARY KEY (tenant_id, campaign_id)                 -- distributed by tenant_id (Citus)
);

CREATE INDEX campaigns_scheduled ON campaigns(scheduled_at) WHERE status='scheduled';

CREATE TABLE trigger_log (
  tenant_id          uuid NOT NULL,
  campaign_id        uuid NOT NULL,
  scheduled_at       timestamptz NOT NULL,
  fired_at           timestamptz,
  fired_by_leader    text,                              -- scheduler pod identity
  paxos_term         bigint,                            -- etcd lease term — fences old leaders
  fire_status        text NOT NULL,                    -- 'pending'|'firing'|'fired'|'failed'
  kafka_offset       bigint,                            -- the offset of the campaign.fired publish (audit)
  PRIMARY KEY (tenant_id, campaign_id, scheduled_at)   -- distributed by tenant_id (Citus)
);

CREATE TABLE fanout_checkpoint (
  tenant_id          uuid NOT NULL,
  campaign_id        uuid NOT NULL,
  segment            int NOT NULL,                      -- 0..1023
  last_recipient_id  text,
  batch_acked_at     timestamptz,
  PRIMARY KEY (tenant_id, campaign_id, segment)
);

CREATE TABLE list_segments (                            -- per-list recipient counts (source for recipientCount)
  tenant_id          uuid NOT NULL,
  list_id            uuid NOT NULL,
  segment            int NOT NULL,                      -- 0..1023
  recipient_count    bigint NOT NULL,                   -- rows in this segment, frozen at list close
  PRIMARY KEY (tenant_id, list_id, segment)            -- distributed by tenant_id (Citus)
);
-- Written when the recipient-list upload closes (uploading → ready): the close step counts rows
-- per segment in recipient_list and records them here, plus a list-level recipient_count_checksum.
-- recipientCount on the create/status API = SELECT SUM(recipient_count) FROM list_segments WHERE list_id=?;
-- this is the authoritative 'intended' total the catch-up reconciler reads.

Recipient store (ScyllaDB, RF=3 LOCAL_QUORUM, 1024 segments per list):

CREATE TABLE recipient_list (
  tenant_id          uuid,
  list_id            uuid,
  segment            int,         -- hash(recipient_id) % 1024
  recipient_id       text,
  email              text,
  locale             text,
  opt_in_proof_id    uuid,
  personalization    blob,        -- per-recipient template variables
  PRIMARY KEY ((tenant_id, list_id, segment), recipient_id)
) WITH compaction = {'class':'LeveledCompactionStrategy'};

Idempotency store (Cassandra LWT, RF=3, TTL 14 d, 8-way recipient shard):

CREATE TABLE idem (
  tenant_id          uuid,
  campaign_id        uuid,
  recipient_shard    int,          -- hash(recipient_id) mod 8 — spreads LWT load across 8 partition coordinators
  recipient_id       text,
  claim_status       text,         -- 'claimed'|'sent'|'failed'|'suppressed'
  claimed_at         timestamp,
  claimed_by_pod     text,
  ses_message_id     text,
  attempt_count      int,
  PRIMARY KEY ((tenant_id, campaign_id, recipient_shard), recipient_id)
) WITH default_time_to_live = 1209600;   -- 14 days

The recipient_shard element in the partition key is mandatory — without it, a single 50M-recipient first-fire would land 14 K LWT/s on ONE Cassandra coordinator and exceed the ~5 K/s/partition Paxos ceiling. With 8 shards, load spreads to 1.7 K LWT/s/partition. The catch-up range-scan fans out across all 8 shards (IN (0..7)) — bounded latency.

CLAIM uses LWT (Paxos) for linearizable cross-replica consistency:

INSERT INTO idem (...) VALUES (...) IF NOT EXISTS;   -- wins → first sender, loses → duplicate

Suppression store (Cassandra, RF=3, LWT on opt-out writes):

CREATE TABLE suppression (
  channel            text,         -- 'email'
  identifier         text,         -- lowercased email
  reason             text,         -- 'opt-out' | 'hard-bounce' | 'complaint'
  created_at         timestamp,
  proof_id           uuid,
  PRIMARY KEY ((channel, identifier))
);

Database choice — recommended:

  • Postgres + Citus for campaign + trigger log + fan-out checkpoint. SQL JOIN with opt_in_proof for compliance queries; transactional pending → firing → fired state transitions; per-tenant shard isolation.
  • ScyllaDB (RF=3 LOCAL_QUORUM) for recipient store. Per-segment range scans at 50× the throughput of an equivalent Postgres table at this scale (Discord's published migration retro).
  • Cassandra + LWT for idem store and suppression. LWT (Paxos) for linearizable CLAIM + opt-out writes. Separate clusters for blast-radius isolation.
  • Redis Cluster 7 for the hot dedupe cache. Cache, not authority. Cold restart is correct (Cassandra absorbs the spike).
  • Kafka 3.7 KRaft for the spine. Separate topics for bulk + transactional + retry — the bulkhead is load-bearing.
  • etcd 3.5 for scheduler leader + active-region lease. Linearizable Raft is the safety primitive — never weaken to DynamoDB Global Tables for the fence (30-s replication lag is too soft).
  • ClickHouse for analytics + reputation tracking. Columnar engine answers aggregate-by-campaign queries in ~100 ms vs minutes for Postgres.
  • Temporal for the catch-up reconciler. Durable workflow state is exactly the abstraction this safety-net role needs.

07High-level design

Which components handle a request, and in what order?

Architecture summary:

  1. API Gateway (Envoy + WAF + per-tenant rate limit + auth). Single-region accept; cross-region failover via DNS.
  2. Campaign Acceptor is the write plane. Validates the template (golden-dataset canary render), inserts the campaign row (UNIQUE on campaign_id is the producer-side dedupe), registers the scheduled fire.
  3. Campaign Store + Trigger Log (Postgres + Citus, 16 shards). The campaign metadata + the Borgcron-style durable log of every fire decision. Sharded by tenant_id.
  4. Recipient Store (ScyllaDB, 1024 segments per list). The 50M+ recipient rows.
  5. Scheduler (Borgcron-style, 3-replica Raft with etcd lease). Wakes every second, reads pending_fire rows, takes the lease, writes the trigger log, publishes campaign.fired to Kafka.
  6. etcd Lease holds the scheduler leader lease AND per-tenant active-region fencing lease.
  7. Kafka Spine (KRaft, 24 brokers, 3 AZ). Multi-topic: campaign.fired, send.requested.{bulk,transactional}, send.retry, send.delivered, provider.feedback.
  8. Fan-out Controller (32 replicas, sticky-partitioned). Consumes campaign.fired, paginates ScyllaDB by segment, publishes send.requested batches, checkpoints every 1 K recipients to fanout_checkpoint.
  9. Render & Dispatch Workers (50 replicas, separate pools for bulk + transactional). Consume send.requested, fence-check, suppression-check, idem-cache check, idem-store CLAIM via LWT, render template, call SES.
  10. Idempotency Hot Cache (Redis Cluster 16 shards). Write-through on CLAIM, 95% hit ratio target.
  11. Idempotency Store (Cassandra LWT, 36 nodes). Durable per-recipient dedupe; source of truth.
  12. Suppression Store (Cassandra LWT). Opt-outs / bounces / complaints; TCPA-grade linearizable opt-out writes.
  13. External SES (per-region sub-accounts + dedicated IP pools). Multiple sub-accounts for bulk + transactional + regional failover bulkheading.
  14. Webhook Receiver (4 replicas). Ingests SES bounce + complaint webhooks via SNS, dedupes, publishes to provider.feedback.
  15. Catch-up Reconciler (Temporal workflows). Diffs intended vs delivered, paces gap-replay to send.retry.
  16. Analytics DB (ClickHouse, 12 nodes). Powers the campaign-owner status UI and the catch-up reconciler's delivered-count lookup.
  17. Tracing (OTel Collector + Tempo). 100% sampling on send.failed and catch-up traces; 5% on happy-path.

Data flow at trigger fire (the heart of the system):

scheduler (every 1s tick)
  → lease (renew leader lease)
  → campaign-store (trigger_log INSERT 'firing', single-row LWT — outbox prepare)
  → kafka (publish campaign.fired, acks=all)
  → campaign-store (trigger_log UPDATE 'fired', records kafka_offset)

This is the outbox pattern: the trigger log row is written FIRST, then the Kafka publish. If the publish fails, the row stays firing and the next scheduler tick will retry. If the scheduler crashes between the log write and the publish, the new leader sees the firing row > 30s old and DELIBERATELY re-publishes. Duplicate publishes are dedupedat the per-recipient tail.

Data flow at fan-out (the 1 → 50M expansion):

kafka (consume campaign.fired, sticky-partitioned 1 partition → 1 fanout pod)
  → campaign-store (read campaign metadata)
  → recipient-store (range-scan segment N, 5 MB / 50K rows, ~200ms)
  → kafka (publish 100 send.requested batches × 500 events each)
  → campaign-store (checkpoint cursor: 'segment N up to recipient X')
  ↻ next segment (parallel up to 16 concurrent segments/pod)

Single 50M campaign: 1024 segments × 50K recipients / 64 concurrent workers = ~16 segments/worker × ~5 MB / ~200 ms each = ~16s × 16 = 256s of Scylla read work / worker. Total wall-clock: ~5 min to enqueue all 50M into Kafka. Send completes within the 60-min SLO.

Data flow at dispatch (per-recipient, the critical path):

kafka (consume one send.requested event, ~5 ms)
  → lease (fence-check: are we the active region? ~20 ms)
  → suppression (point-read: is this recipient opted-out? ~5 ms cached / ~30 ms LWT)
  → idem-cache (point-read: already CLAIMED? ~1 ms)
  ⤷ HIT → drop (duplicate, branchKey='cache-hit-dup'). END.
  → idem-store (LWT INSERT IF NOT EXISTS, ~30 ms) — branchKey='cache-miss'
  ⤷ applied=false → drop (concurrent CLAIM won by another pod). END.
  → idem-cache (write-through SET, ~1 ms)
  → ses (SendEmail, ~350 ms p50)
  → idem-store (UPDATE claim_status='sent', ses_message_id=?, ~5 ms)
  → kafka (publish send.delivered, ~10 ms — follows-from, not on critical path)

p50 critical path: ~430 ms per send. With 50 pods × 200 concurrent each = 10 K concurrent in-flight → 14 K sends/s sustained, matching the SES warm-pool tier.

Data flow at catch-up (the safety net):

catch-up-temporal-workflow (every 5 min for in-flight, or on-demand after incident)
  → lease (acquire /lease/catchup/{campaign_id}, 5-min TTL)
  → campaign-store (read intended segment-count + recipient-count)
  → idem-store (range-scan across 8 recipient_shard partitions for claim_status NOT IN ('sent','suppressed','failed'))
            ▲
            └── source of truth — NOT analytics-db. Reading analytics-db here would
                false-positive on analytics lag and trigger a send.retry storm
                (Critic Round 1 issue 6c).
  → kafka (publish send.retry events for unsent recipients, paced ≤ 50% baseline send rate)
  → idem-store (eventually: all rows have claim_status='sent'|'suppressed'|'failed')
  → analytics-db (read campaign analytics for the campaign-owner UI freshness check — informational only)
  → campaign-store (mark campaign 'complete')

trace catalogue

The simulator runs these specific journeys against your diagram. Each is a deterministic walk; clicking between them shows how the canonical answers shift shape.

IDJourneyBudget
mass-scheduled-email:campaign-createMulti-segment: (1) INSERT campaign metadata into Postgres + (2) bulk-append recipient list into ScyllaDB. Saga: half-uploaded list times out and the caller retries.25 s
mass-scheduled-email:scheduler-fires-triggerMulti-segment: (1) scheduler writes firing to trigger_log (outbox prepare) + (2) publishes campaign.fired to Kafka (outbox dispatch). The Borgcron at-least-once invariant lives here.300 ms
mass-scheduled-email:fanout-controller-chunksSingle segment with revisit. Fan-out worker consumes campaign.fired, range-scans recipient-store one segment at a time, publishes send.requested batches, checkpoints to campaign-store. Revisit allowed because the controller cycles segments.600 s
mass-scheduled-email:send-happy-pathMulti-segment: (1) pre-flight checks (fence + suppression + idem-cache MISS), (2) durable CLAIM (idem-store LWT, then write-through cache), (3) SES call + delivered publish. The critical path the user's email rides.1500 ms
mass-scheduled-email:send-duplicate-suppressedSingle segment: dispatcher consumes send.requested → idem-cache HIT → drop. The load-bearing dedupe story — when a re-fire or fan-out resume creates a duplicate enqueue, this is where it gets caught harmlessly.80 ms
mass-scheduled-email:webhook-bounce-ingestMulti-segment: (1) SES → webhook-rx (SNS signature verify, fast 200), (2) webhook-rx publishes provider.feedback to Kafka, (3) consumer downstream upserts suppression-store.5 s
mass-scheduled-email:catch-up-reconcilesMulti-segment with revisit: (1) catch-up acquires lease, (2) reads intended (campaign-store), (3) range-scans idem-store for delivered/claimed state (source of truth — NOT analytics-db), (4) publishes gap to send.retry at paced rate. The Temporal-backed safety net.30 s

failure scenarios we model

Each grounded in a real public precedent. The simulator replays these against the canonical to surface the on-call story.

IDScenarioPrecedent
scheduler-double-fireTwo scheduler leaders both fire after partition heal; per-recipient idem catches it.Borgcron retro
scheduler-missed-fireScheduler down 4 h; catch-up controller replays missed fires within relevance window.K8s CronJob 24-day retro (Lancey 2020)
fanout-controller-crashFan-out crashes after enqueueing 23M of 50M; resume from checkpoint, per-recipient idem catches re-enqueue overlap.Shopify high-availability background jobs
idem-cache-cluster-lossRedis cluster cold restart mid-campaign; dispatcher falls through to Cassandra LWT; ~10× latency, no correctness loss.Twitter Redis-cluster failover retros
ses-throttle-cascadeSES 429 rate climbs; circuit-breaker opens; divert to send.retry with exp-backoff; reputation preserved.AWS Builders' Library timeouts/retries/backoff
catch-up-stormRecovery worker fires 8 h of missed campaigns at once; pacing prevents SES throttle.Slack Jan 2021 thundering-herd retro
tcpa-opt-out-lagOpt-out written in us-east-1; us-west-2 dispatcher tries to send before replication; Cassandra LWT linearizable opt-out check blocks.Satterfield v Simon & Schuster
idem-ttl-too-shortCatch-up runs at +25 h with TTL=24 h; recipient re-receives. Our 14-d TTL prevents this.GitHub Oct 2020 webhook-DLQ retro
noisy-neighborTenant A blasts 100M while tenant B's transactional starves. Bulk-vs-transactional topic split prevents.Twilio Jul 2023 SMS pile-up
webhook-bounce-storm12M bounces in 10 min after deliverability hit; webhook-rx absorbs via Kafka, suppression updates drain async.Datadog Mar 2023 multi-region

08Deep dives

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

1. Why Borgcron (Paxos-replicated scheduler) instead of Kubernetes CronJob.

K8s CronJob stops scheduling after 100 missed runs and silently skips on concurrencyPolicy: Forbid overlaps (Lancey 2020 retro). Quartz clustered runs into row-lock contention on QRTZ_LOCKS at our QPS. A single-leader scheduler is a SPOF — when that node dies, every campaign owed an "at 9 AM" send misses its window. Borgcron's answer: a 3-replica Paxos/Raft cluster where the leader fires triggers and durably logs every fire decision to an external store BEFORE publishing. On leader failover, the new leader reads the log and DELIBERATELY re-fires rows that are firing but > 30 s old. The system favors duplicates over misses — duplicates are loud (idempotency store rejects them) and recoverable; misses are silent and unrecoverable. Google's published paper (Sundheim, ACM Queue 2015) is the canonical reference.

2. The per-recipient idempotency tail — two tiers, why.

Cassandra LWT (Paxos) is the source of truth: linearizable cross-replica INSERT ... IF NOT EXISTS per (tenant_id, campaign_id, recipient_id). But LWT costs ~30 ms per write and is CPU-heavy on Cassandra; at 14 K sends/s, sustained LWT load would saturate the cluster. We pair it with a Redis hot cache (95% hit ratio): the dispatcher checks the cache first (~1 ms point read), only falls through to Cassandra on MISS. The cache is WRITE-THROUGH on CLAIM win, so cache and store stay consistent. Cache cold restart is correct — Cassandra absorbs the spike at ~10× latency degradation. Critically: the cache is an OPTIMIZATION not a correctness primitive. Losing the cache leaks NO duplicates; it just makes the system slower.

3. The outbox at trigger fire.

The scheduler can't atomically write the trigger log AND publish to Kafka — they're separate systems. The outbox pattern solves it: (a) write trigger_log row with fire_status='firing' (durable Postgres single-row LWT); (b) publish campaign.fired to Kafka with acks=all; (c) on producer ack, UPDATE fire_status='fired'. If the scheduler crashes between (a) and (b), the next leader sees a firing row older than 30 s and re-publishes — the per-recipient idem tail catches the duplicate at the dispatcher. If the publish fails after retries exhaust, the row stays firing and ops can manually retrigger. NEVER batch (a) and (b) into a 2PC — Kafka isn't a 2PC participant and forcing it to be one is the path to 30-second commits.

4. Catch-up controller — Temporal's load-bearing role.

Catch-up is the safety net for everything the scheduler tier missed. A 50M-recipient diff between (intended from campaign-store) vs (delivered/claimed from idem-store) takes minutes to compute and the resulting send.retry publishes take hours to drain. A naive cron job restarts from scratch on every crash — useless. Temporal's workflow durability is exactly the abstraction: each step's state writes to Temporal's history; a worker crash resumes from the last step. DoorDash's published "Cadence as fallback" post is the precedent — workflow engines shine when you need durable state across hours of execution. We use Temporal HERE and not for the scheduler tier because the scheduler needs sub-200-ms firing latency, which Temporal's queue-then-poll model can't meet.

5. Active-region fencing — the etcd lease.

The single failure mode this whole problem exists to prevent is "two regions both think they're active and double-send 50M recipients each." Cloudflare 2022 and Slack 2021 are the precedents. The fence is an etcd lease at /lease/region/{tenant_id} with a 30 s TTL; the active region renews every 10 s, the standby region fails-CLOSED on send if it doesn't hold the lease. Dispatcher checks BEFORE every send — yes, that's a ~20 ms tax on every send, and yes, we pay it because the alternative is TCPA fines. etcd's Raft gives us linearizable lease reads; DynamoDB Global Tables would be too soft (30 s replication lag = 30 s of double-active risk).

6. SES reputation as a resource.

You can't buy IP reputation back once you've lost it. A bounce spike >5% triggers SES auto-pause; a complaint spike >0.1% does the same. The dispatcher's circuit-breaker is the primary defense: SES 429-rate >5% opens the circuit and diverts to send.retry with exp-backoff. Pre-flight list hygiene (suppress hard-bounces from previous campaigns BEFORE the fan-out) is the second defense — sending to known-bad addresses is the fastest way to crater reputation. The catch-up controller's pacing is the third defense: never replay faster than baseline send rate. Mailchimp's published IP-warmup curve is the playbook for new pools.

7. Per-tenant Kafka bulkheading — why separate topics.

A 100M-recipient marketing blast and a 1-recipient "password reset" cannot share a Kafka topic. The bulk blast's volume would back up the partition the transactional message hashes to, blowing the 30 s transactional SLO. We separate at every layer: separate Kafka topics (send.requested.bulk vs send.requested.transactional), separate consumer groups, separate dispatcher pools, separate SES sub-accounts, separate IP pools, separate reputation. Producer-side admission control at the gateway throttles bulk independently. Twilio's Jul 2023 SMS pile-up retro is the textbook story of what happens without this.

8. Why the recipient-list lives in ScyllaDB, not Postgres.

A single 50M-recipient list is too big for a single Postgres partition; sharding by list_id puts all 50M rows on one shard which doesn't help. We segment the list into 1024 partitions of ~50K rows each via segment = hash(recipient_id) % 1024, and ScyllaDB's per-partition contiguous read is ~50× faster than Postgres's equivalent at this access pattern. Discord's 2023 Cassandra→Scylla migration retro is the precedent — they hit the exact same bottleneck on message history reads (177→72 nodes during migration, sustained 3.2M writes/s).

09Trade-offs

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

What breaks at 10× scale (~500M-recipient peak campaigns, ~140 K sends/s):

  • Single Cassandra LWT cluster for idem-store hits per-partition Paxos contention → shard idem-store by (tenant_id, campaign_id, recipient_id_mod_8) instead of just (tenant_id, campaign_id) so each campaign's 500M rows spread across 8 partition groups.
  • Single SES dedicated IP pool caps at 50 K/s warm-pool tier → split across 3+ AWS regions with cross-region Kafka mirror; dispatcher reads in-region only, eliminating cross-region SES API latency.
  • Single Kafka cluster can't hold the topic count → split per priority into separate clusters; cross-cluster mirror via MirrorMaker 2 for the audit pipeline only.
  • Single scheduler tier at 200 fires/s ceiling → shard the scheduler tier by tenant_id_mod_N so each scheduler shard handles 1/N of fires; etcd lease per shard.
  • Cassandra cluster size at 500M rows × 14 d × 120 B = 840 GB / shard → grow node count from 36 to 100+ to keep <2 TB/node; node-add operationally heavy.

Per-component failure stories (mapped to chaos cards):

  • Scheduler double-fire (scheduler-double-fire) — Borgcron at-least-once + trigger_log UNIQUE constraint + per-recipient idem-store all defend.
  • Scheduler missed fire (scheduler-missed-fire) — catch-up controller replays within relevance window.
  • Fan-out controller crash (fanout-controller-crash) — checkpoint cursor + per-recipient idem at the tail.
  • Idem cache loss (idem-cache-cluster-loss) — fall-through to Cassandra LWT, ~10× latency but no correctness loss.
  • SES 429 cascade (ses-throttle-cascade) — circuit-breaker + AIMD + send.retry exp-backoff.
  • Catch-up storm (catch-up-storm) — rate-pacing at ≤ baseline + relevance-window drop.
  • TCPA opt-out lag (tcpa-opt-out-lag) — synchronous Cassandra LWT on opt-out write + dispatcher synchronous read.
  • Idempotency TTL too short (idem-ttl-too-short) — TTL=14 d vs 6 h catch-up window = 56× headroom.
  • Noisy neighbor (noisy-neighbor) — per-tenant + per-priority Kafka topics + SES sub-account isolation.
  • Webhook bounce storm (webhook-bounce-storm) — webhook-rx is a thin durable hand-off; suppression updates drain at their own pace via Kafka.

observability — what on-call actually opens during a page

SLIs (RED + USE per component):

  • Scheduler tier: trigger_fire_lag_seconds_p99 (target <5 s), scheduler_leader_count (target =1), pending_fires_backlog (target <100), paxos_term_changes_per_hour (target <5).
  • Fan-out tier: fanout_progress_recipients_per_sec (target ≥ campaign_target × 0.95), fanout_checkpoint_lag_recipients (target <10 K).
  • Dispatcher tier: send_p99_seconds (target <2 s), send_critical_path_p99 (target <1 s), ses_5xx_rate (target <2%), dispatcher_circuit_open (binary, page on any).
  • Idempotency: idem_cache_hit_ratio (target ≥0.95, scoped to T+5min after campaign start — first-fire is inherently cache-MISS, alerting from T=0 would alert-storm), idem_cache_lwt_fallthrough_rate (target <5% — pages BEFORE the LWT cluster saturates, gives root-cause-pointing signal vs the LWT cluster being the canary), lwt_p99_seconds (target <100 ms), lwt_paxos_contention_rate (target <0.5% — counts only LWT errors/timeouts, NOT applied=false), idem_claimed_orphan_age_seconds_p99 (target <10 min — the janitor's own health), idem_double_send_sample_rate (target <1-in-1M — hashes (campaign_id, lowercase(email), date) from send.delivered, counts dupes in ClickHouse, catches gray-failure bugs like email-case-normalization mismatches that no other SLI sees).
  • Suppression: suppression_check_p99 (target <50 ms), suppression_lwt_timeout_or_error_rate (binary P1 — excludes applied=false which is the normal already-opted-out path).
  • SES: ses_throttle_rate (target <2%), ses_reputation_bounce (target <3%), ses_reputation_complaint (target <0.05%).
  • Catch-up: catchup_workflow_age_seconds (target <2 h), catchup_lag_recipients (target <100 K).
  • Kafka: consumer_lag_seconds per topic (target <30 s on hot topics).

SLO burn-rate alerts (Google SRE Workbook ch. 5):

  • 99.99% scheduler-fire-on-time → 1-h burn 1× = P2, 5-min burn 14.4× = P1.
  • 99.9% per-recipient at-most-once → any duplicate above 1-in-1M is a P1 incident review.
  • 99.95% campaign-control-plane availability → 1-h burn 2× = P2.

Dashboards on-call opens during a page (in this order):

  1. Campaign health — per-campaign progress: intended vs delivered vs failed, ETA-95%, last 1 h of fire decisions.
  2. Scheduler + lease health — leader pod, paxos term, lease renewals, etcd quorum, last 1 h of fires.
  3. Send pipeline — Kafka consumer lag per topic, dispatcher pool utilization, SES throttle + 5xx, circuit-breaker state.
  4. Dedupe + suppression — idem-cache hit ratio, Cassandra LWT p99, suppression check p99 + LWT failures.
  5. Tracing forensics — paste recipient_id + campaign_id → full timeline (fire → fanout → dispatch → SES → delivery).

Tracing carries traceparent on every Kafka message header (not body) — so deserialization stays independent of tracing. Tail-based sampling at the OTel collector: 100% of send.failed, 100% of catch-up workflows, 5% of happy-path.

deploy & runtime

  • Rollout strategy: Blue/green for the scheduler tier (a bad scheduler deploy can miss every fire — we accept the 2× capacity cost for safety). Canary for everything else: 1% → 10% → 50% → 100% with 5-min hold at each step, automated rollback on send_p99_seconds regression or ses_5xx_rate spike.
  • Schema migrations: Citus's two-phase metadata commit + pg_repack for online column adds; no downtime. ScyllaDB additive-only schema changes deployed via cqlsh rolling.
  • Secrets rotation: SES API keys rotated quarterly via AWS Secrets Manager with 24-h overlap window (both keys valid). mTLS certs auto-rotated by cert-manager weekly; Square 2023-09-07 cert-bundle outage is the cautionary tale (alert at T-30 days on cert expiry).
  • Consumer-group rebalance: Kafka KIP-429 cooperative rebalance + group.instance.id static membership → near-zero STW on rolling deploys.
  • Recipient list updates: Mid-campaign list changes are NOT supported — the recipient set is frozen at campaign-create time. Opt-outs propagate via suppression-store (checked synchronously per send).

multi-region / DR

Active-active for the campaign control plane. Both regions accept campaign submissions, route to the home region via consistent hash on tenant_id. The scheduler tier is single-region active (one Borgcron cluster per environment — operational simplicity wins; the fan-out + dispatch tiers are where active-active matters). On full-region scheduler loss, an operator-driven failover promotes the standby region's scheduler via a typed runbook. Honest RTO: 30 min p99, including page-time + operator-context-build + etcd observer promotion to voting member + scheduler cluster activation. The earlier "10 min" claim was aspirational (a clean-runbook best case); 30 min is what we measure in GameDay drills.

Per-tenant active-region pinning for fan-out + dispatch via etcd lease. Each tenant has a HOME region; sends go through that region's dispatcher pool + SES sub-account. Standby region's dispatchers fence-check on every send and fail-CLOSED if they don't hold the lease.

RTO/RPO posture:

  • Scheduler tier RTO 30 min p99 honest (operator-driven cross-region promotion + etcd observer-to-voter switch), RPO 0 (trigger log replicated synchronously within region; catch-up replays from log).
  • Dispatch tier RTO 5 min, RPO 0 (per-tenant lease cutover; per-recipient idem prevents double-send during cutover).
  • Recipient store + idem store RTO 0 (LOCAL_QUORUM survives single-node loss instantly).
  • Cross-region link loss → tenants pinned to the still-reachable region; catch-up reconciles missed sends on heal.

Defensible because: every RTO is tied to a specific recovery mechanism (etcd lease TTL, Paxos election, Kafka leader re-election), not aspirational. RPO=0 on the safety-critical stores (campaign-store trigger log, idem-store) because losing fire records or dedupe records = double-send or missed send. RPO=30 s on analytics-db is acceptable because catch-up reconciler retries until the count is fresh.

security posture

  • Auth: Tenant-scoped bearer tokens (rotated quarterly); service-mesh mTLS (SPIFFE/SPIRE) between internal services.
  • Secrets: SES API keys in AWS Secrets Manager with 24-h overlap on rotation. mTLS certs auto-rotated by cert-manager weekly.
  • PII: Email addresses are PII under GDPR — encrypted at rest in ScyllaDB (SSE-KMS) and Cassandra (TDE); access logged. Right-to-erasure via a separate worker that scrubs recipient rows by opt_in_proof_id across recipient-store + idem-store + suppression + analytics-db.
  • Threat model: (1) Compromised tenant key → blast-radius bounded by per-tenant rate limit at gateway + per-tenant Kafka topic isolation. (2) Insider data exfil → all ScyllaDB / Cassandra access mediated through audited service-level APIs; no direct DB credentials in dispatcher pods. (3) SES sub-account compromise → IP-pool isolation contains to that pool's reputation; other tenants unaffected. (4) Catch-up worker abuse (firing replays of OLD campaigns) → catch-up requires acquiring the lease + the relevance-window check drops stale campaigns automatically.

open questions

  • Per-tenant scheduler sharding for 10× scale? Current design has a single global scheduler tier. At 10× scale (400 fires/s sustained) it bottlenecks. Plan B is sharding by tenant_id_mod_N — needs operational design for the N etcd leases and rolling deploys across N scheduler clusters. Worth doing pre-emptively or only when the bottleneck is visible?
  • Should the catch-up reconciler also handle SES-side bounce reprocessing? Currently webhook-rx handles real-time SES feedback. The catch-up controller could also reconcile against provider.feedback for completeness (e.g. a webhook delivery failure leaves the suppression-store stale). Plausible — adds workflow complexity.
  • Cross-region idempotency for cross-region tenants? If a tenant goes multi-region for redundancy, do we run cross-region LWT for the idem store (~80 ms p99 cross-region) or rely on regional ownership? Currently the latter, but a future tenant might demand stronger guarantees.
  • DLQ replay window vs idempotency TTL. We picked TTL=14 d for a 6 h catch-up window. If we ever extend the catch-up window past 7 d (e.g. for compliance retros), we need to bump TTL to ≥ 2× the new window.
  • Template versioning under in-flight campaigns. If a campaign-owner updates a template mid-fan-out, the in-flight sends use the OLD version (template_id is captured at fire time). That's the right semantic for marketing but might surprise transactional users. Worth surfacing as a campaign-create-time option?
  • Components deliberately not modeled at this canonical level. A template/content registry (S3 + Postgres metadata with content-addressable hash) is referenced indirectly through campaign-store.templates + dispatcher's 5-min in-process template cache, but isn't drawn as a separate node — it's a load-bearing real component and would be added on the next canonical pass. Similarly: a runtime config / feature-flag plane (LaunchDarkly or internal) for per-tenant rate limits, AIMD parameters, IP-pool routing — exists in production, currently absorbed into the gateway's static config. The DLQ for send-retry exhaustion is modeled as the send.retry topic with max-retries enforced at the dispatcher rather than a separate dlq-store node; the topic IS the DLQ at our scale (per Stripe's webhook pattern). All three are conscious omissions for diagram readability, not blast-radius oversights.
  • Critic-acknowledged residual risks accepted in this version. Round 1 surfaced: (a) the 30-s scheduler-leader re-fire window is deliberate at-least-once (mitigated by per-recipient idem), (b) the 5-min orphan-claim janitor cadence means up-to-5-min stalled claimed rows are temporarily un-retriable (acceptable: the per-recipient TTL is 14 d and the catch-up workflow surfaces these eventually), (c) the cross-region failover RTO is honestly 30 min not 10 min (operator-driven typed runbook). These are accepted trade-offs documented above; if any prove unacceptable for a specific tenant we'd revisit the design for that tenant.

Primary sources

  • Sundheim — Reliable Cron across the Planet (ACM Queue 2015, Borgcron paper)
  • LinkedIn — Air Traffic Controller: Member-First Notifications (2016)
  • LinkedIn — Hermes mass email retros (eng blog)
  • Pinterest — NEP Notification System and Relevance (Medium 2017)
  • Uber — Cherami: Uber's durable distributed task queue
  • Uber — Announcing Cadence (Temporal predecessor)
  • DoorDash — Cadence as a Fallback for Event-Driven Processing
  • Discord — How we store trillions of messages (Cassandra → ScyllaDB)
  • Stripe — Designing robust idempotency keys (Brandur)
  • Brandur — Implementing Stripe-like Idempotency Keys in Postgres
  • AWS Builders' Library — Idempotency at Scale (re:Invent ARC403 2021)
  • AWS — Handling SES throttling (Maximum sending rate exceeded)
  • AWS — SES sending quotas & dedicated IP warmup
  • Mailgun — Bulk email sending with queue management
  • Slack — Tracing notifications (Go → Kafka → Elasticsearch)
  • Shopify — High availability background jobs
  • Vallery Lancey — Kubernetes CronJob Failed For 24 Days (case study against naive cron)
  • Quartz Scheduler — JDBC JobStore clustering docs (the misfire-threshold story)
  • Cloudflare — Cron Triggers internals (Nomad-distributed schedulers)
  • Apache Kafka — KIP-429 cooperative incremental rebalancing
  • Beyer et al. — SRE Workbook ch. 21–23 (overload, cascading failures, critical state)
  • Campbell & Majors — Database Reliability Engineering ch. 9 (fencing tokens)

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 Distributed Cron — Mass Scheduled Email yourself