Payment / Wallet System — a worked solution
Idempotent, double-entry, sagas vs 2PC.
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 Payment / Wallet System workspaceThe problem
Build a payment / wallet system that accepts charges, moves money across accounts, and never loses or double-applies a transaction — even when the network drops, the database fails over, or a third-party acquirer returns 503 mid-charge. The architectural pressure here isn't QPS (Stripe-scale is ~10K charges/sec sustained, well within reach of a sharded Postgres) — it's correctness under retry, invariant preservation across services, and survivable failure modes.
Three load-bearing concepts:
- Idempotency — every mutating call is keyed and replayable, with state-machine semantics so a crashed-and-retried client gets the same answer (Stripe / Brandur / Airbnb Orpheus).
- Double-entry, immutable journal + projected balances — the ledger is append-only; balances are derived (Square Books, Monzo, Uber Gulfstream).
- Saga (not 2PC) for cross-service money movement — durable workflow engine drives multi-step charges with explicit compensations (Uber Cadence / Temporal, Airbnb Skipper).
The reference architecture
What each component is for
- Client
Two distinct callers share this entry point: (a) a merchant backend invoking POST /v1/charges with an Idempotency-Key header, and (b) the wallet user's mobile app issuing P2P transfers and reading their own balance. Receives webhook callbacks at a registered endpoint after charges complete asynchronously.
Why it exists. We considered modeling merchant and wallet as separate edge nodes, but the failure modes (idempotent retry, webhook dedup, 30 s charge SLO) are identical. Splitting them only matters when we model authn — restricted keys for merchants, OAuth + device binding for wallet users — which we draw at the auth node downstream rather than at the edge.
When it fails. Naïve clients regenerate idempotency keys on retry — that creates ghost charges. Detection: distinct-keys-per-charge p95 > 1 alerts SDK eng. Mitigation: SDK quickstarts ship the canonical retry helper; webhook signing keys are per-merchant so a leaked key doesn't compromise others.
- API GatewayEnvoy + Cloudflare WAF + TLS terminator
L7 ingress for every public payment API call. Terminates TLS with the public Stripe-style cert (m=8K req/sec/edge), runs the WAF ruleset (PCI 6.5 OWASP top-ten plus payment-specific rules: rejected card-number bodies, fingerprinted scrapers), enforces per-merchant token-bucket rate limits BEFORE forwarding so retry storms can't poison the idempotency store, and copies the Idempotency-Key header verbatim onto the upstream request.
Why it exists. Considered putting auth + rate limit inline in payment-api; rejected because (1) a 5xx in payment-api would skip the rate limiter and let the storm reach the ledger, (2) ratelimits are per-merchant which means we'd ship that data to every payment-api replica, vs centralizing it at the edge tier. Stripe's published 2019-07-10 retro shows what happens when ratelimits live in the wrong layer (cascade across products).
When it fails. Gateway pool exhaustion masks as 'merchant 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 merchant rather than 503'ing everyone.
- Auth / IdentityHashed restricted-key authn + OAuth2 for wallet
Validates the Authorization header on every public call: hashed key lookup against a sharded Postgres of API keys, scope check (charges:write vs balances:read vs refunds:write), per-key revocation list. Wallet-user calls go through OAuth2 + device binding instead. Returns a signed merchant_id + scopes envelope back to the gateway, valid for 60s.
Why it exists. We considered baking authn into payment-api. Rejected because (1) a CVE in the API service should not break authn for everyone, (2) restricted keys are a separate product surface (Stripe ships them as scoped credentials), and (3) the lookup hits a 100M-row keys table that doesn't belong on the charge hot path. Centralizing authn behind the gateway means payment-api can trust the signed envelope without a round trip.
When it fails. Auth down → all writes 503 (fail-closed); reads with cached scopes survive 60 s then 401. Detection: auth_5xx_rate, scope_cache_miss_rate. Square's 2023 cert-bundle outage (18 hours) is the cautionary tale — auth lookups must NOT depend on a freshly-fetched cert that the deploy process can break.
- Service · Payment API (Write Plane)Go + pgx + sqlc; Stripe / Brandur idempotency state machine
Hosts POST /v1/charges, /v1/refunds, /v1/transfers. For each request: hashes (method, path, body) into a request_fingerprint; runs the Brandur-style idempotency state machine — claim row → started → ledger_committed → psp_called → finished. Drives the synchronous slice of the charge (validate → idem.claim → ledger.write+outbox → balance-cache.invalidate); hands long-running PSP capture off to the saga orchestrator so the API call returns in <500 ms even when the acquirer is slow.
Why it exists. Considered: a single combined read+write service (Stripe's first version). Splitting Payment API (write plane) from Balance Reader (read plane) lets each scale on its own load shape, isolates blast radius (a balance-read storm can't crater charge writes), and makes per-endpoint rate limiting cleaner. The idempotency state machine MUST live here rather than as a sidecar — Brandur's writeup proves correctness only holds when the state transitions are interleaved with the business write in the same DB transaction.
When it fails. Idem store down → fail-CLOSED 503 (better to lose a request than double-charge). Ledger primary down → fail-CLOSED 503. PSP down → return 'requires_action' + persist a saga; client polls or waits for webhook. Replica AZ loss → traffic drains to surviving AZs in <30 s via gateway re-routing. Detection: charge_5xx_rate, idem_lock_wait_ms p99, ledger_write_p99, post_commit_psp_call_pending count. Stripe's 2019-07-10 lesson: a single shared shard's gray failure cascaded across services because health checks didn't notice — pin invariant SLOs, not just liveness.
- Service · Balance Reader (Read Plane)Go + pgx; reads from ledger followers + balance cache
Hosts GET /v1/balances/:account, GET /v1/charges/:id, list endpoints. Cache-aside: tries balance-cache first (p99 ~5 ms when warm); falls through to a ledger replica on miss (p99 ~80 ms). For read-your-writes (the user just made a transfer and immediately checks balance), routes to the leader using a session token issued by payment-api on commit.
Why it exists. Considered serving balance reads from payment-api; rejected because read QPS is ~10× write QPS and a stuck read connection pool would back-pressure writes. Considered a denormalized balance cache as the source of truth; rejected because Monzo's ledger retro (2023) shows projection-from-immutable-journal is the only design that survives auditors. Balance Reader holds NO state — just a cache lookup + a replica read.
When it fails. Cache cold → all reads hit the replica → replica saturates → consider degraded-mode that returns last-known balance + a 'may be stale' header for 30 s while the cache warms. Replica lag spike → RYW reads from leader survive but everyone else sees stale data. Detection: cache_hit_ratio, replica_lag_seconds, balance_read_p99.
- Cache · Balance CacheRedis 7 cluster mode, 32 slot groups
Caches balance:{account_id} and charge:{charge_id} for the read plane. Populated lazily by Balance Reader on miss; INVALIDATED (not written-through) by Payment API after each ledger commit because the projected balance can be derived only from the journal — writing into the cache directly would let a buggy writer poison the cache.
Why it exists. Considered making the cache the source of truth; that's the design that gave every payment team Monzo's drift retro. The cache is a strict performance optimization — cold cache is correct (just slow), and a stale entry self-corrects on the next invalidation. Considered cache-aside vs read-through; cache-aside (Balance Reader looks up the cache itself) is simpler and gives explicit miss accounting.
When it fails. Cache cluster cold restart → 100% miss → ledger replica saturates → balance reads slow but correct. Single shard down → 1/N keyspace cold; balance-reader degrades to leader for that slot group. Detection: hit_ratio, slot_health, cluster_state. Mitigation: warm-up job replays last 30 minutes of postings on startup.
- SQL DB · Idempotency StorePostgres 16 KV (Brandur pattern)
Per-row state for every Idempotency-Key the gateway has seen in the last 24 h. Schema: (key, request_fingerprint, locked_at, recovery_point ENUM, response_code, response_body_jsonb, expires_at). Three operations: CLAIM (INSERT … ON CONFLICT row-lock), ADVANCE (UPDATE recovery_point in same txn as the business write), and FETCH (SELECT for retry-after-finish).
Why it exists. Considered a separate Redis idempotency layer; rejected because the load-bearing correctness property (state transitions atomic with the business write) only holds if the idempotency row and the journal row commit in the SAME transaction. That forces same-engine, same-cluster — Postgres. Considered DynamoDB conditional writes; works for simple idempotency but not for the phased recovery_point machine that lets a crashed writer resume from the right step.
When it fails. Store down → payment-api fails CLOSED with 503 (cannot accept a charge without claim). Reaper accidentally deletes hot keys → the next retry re-executes as a fresh charge — a genuine double-charge, with no downstream guard to catch it. That is exactly why the reaper only sweeps rows 48 h past expires_at (the 72 h soft retain) and why its delete volume is alerted against expected expiry counts. Detection: claim_p99 < 50ms (1-min) → P2 page; lock_wait_ms_p99 < 100ms (1-min) → page if breached for 5min; reaper_lag > 1h → P3. Stripe's published lesson: the fail-closed posture is non-negotiable; fail-open trades a 503 for permanent ledger drift.
- SQL DB · Ledger PrimaryPostgres 16 + Citus, 64 hash shards by account_id
Source of truth for every money movement. Three tables co-resident on the same shard for transactional cohesion: journal_entries (append-only audit log — posting fields immutable, rows never deleted; only the state lifecycle column advances forward), balance_proj (mutable running balance per account, updated in the same txn as the journal write to maintain the invariant sum(debits)=sum(credits) per account_id), and outbox (event rows the CDC relay tails into Kafka). One SERIALIZABLE transaction writes all three; either everything commits or nothing does.
Why it exists. Considered: balance as the source of truth + journal derived (the design every team starts with). Rejected because Monzo's ledger drift retro and Square's pre-Books migration to Spanner both show derived-balance-from-immutable-journal is the only invariant that survives auditors and account-correction edits. Considered Spanner directly; rejected for in-region p99 (Spanner write p99 ~30-100ms vs Postgres 5-15ms) and operational familiarity. Considered Cassandra (Monzo); works at scale but loses cross-row transactionality which we need for the journal+balance+outbox triple-write.
When it fails. In-region primary failover (automatic, RTO=30s, RPO=0): 30 s of write outage per affected shard; payment-api 503s cleanly with Retry-After 30. Cross-region failover to ledger-db-dr (operator-driven, RTO ~10 min, data-loss budget 0 via reconciliation — NOT RPO=5s; the 5 s number is replication lag, not what gets lost): on regional outage, halt the API for the failover window, reconcile the WAL gap from idempotency-store + acquirer settlement, then resume. Replication lag spike: balance-reader detects lag > 1 s and falls back to the leader for that shard. Single shard down: 1/64 of accounts see write failures; the remaining 63/64 keep working. Detection: oldest_unreplicated_lsn, replica_lag_p99, prepared_xact_age (alerts on stuck txns), shard_health, dr_replication_lag_seconds. Cross-shard saga step failure → compensation; never partial commit.
- Service · PSP AdapterGo + per-acquirer circuit-breaker + HMAC signing
Translates our internal Charge envelope into the acquirer's native protocol (Stripe API, Adyen Online Payments, Worldpay XML, ACH NACHA file for bank transfers). Per-acquirer circuit-breaker, retry policy, idempotency key forwarding, signature verification on responses. Holds outbound TLS with PCI-compliant cipher suites and per-acquirer cert pinning.
Why it exists. Considered baking acquirer logic into payment-api; rejected because (1) per-acquirer outages need independent circuit-breakers, (2) cert/cred rotation is per-acquirer, (3) ten acquirers' SDKs in one binary is a dependency-hell guarantee. Uber's payments platform writeup names this exact split (collection vs disbursement adapters in front of 10+ processors). The adapter never holds state — it's a stateless protocol bridge so we can scale it independently.
When it fails. Single acquirer 503: circuit-breaker opens, payment-api returns 'requires_action', saga retries via fallback acquirer if configured. All acquirers down (Visa Europe 2018-style): full charge outage; we surface clear error to the merchant. Cert expiry: detection at 30 days; failover to a secondary cert (Square 2023 outage was exactly this — 18 hours of mTLS failures). Detection: per-acquirer error_rate, breaker_state, cert_days_remaining.
- External · PSP / Card NetworkStripe / Adyen / Worldpay / Visa / Mastercard / ACH
Out-of-system money movement: card auth via Visa/MC, bank ACH/SEPA, alt-payments (PayPal/iDEAL/UPI). Posts settlement files (T+1 or T+2). Sends webhooks for async events (chargeback, refund.succeeded, payout.paid). NOT under our SLO control; we treat their availability as an exogenous variable.
Why it exists. We considered acting as our own acquirer (Stripe's eventual end state) — out of scope for this design. The diagram pins this as a single external surface to make blast-radius reasoning explicit: 'when this goes down, X stops working, Y degrades.'
When it fails. Visa Europe 2018: 5.2M failed txns in 10 hours from a single switch defect. Mitigation: multi-acquirer routing detects sustained errors and re-routes traffic to the fallback within 60 s of breach. Per-acquirer chargeback latency varies; recon worker tolerates T+90 reversal of T+1-cleared txns (180-day chargeback reserve per merchant).
- Coordinator · Saga OrchestratorTemporal 1.x cluster (Cassandra-backed)
Owns long-running money-movement workflows: charge.authorize → wait → charge.capture → ledger.post-fees → notify-merchant. Each workflow step is an Activity that calls payment-api, psp-adapter, or recon-worker. Activities are idempotent (each carries a unique activity_id passed as the Idempotency-Key downstream). Compensations are explicit and registered up-front (Cadence's Saga.compensate pattern): a failed capture triggers a credit-reversal that posts an offsetting journal entry.
Why it exists. Considered: 2PC across payment-api / psp-adapter / ledger. Rejected because (a) 2PC across services needs a coordinator that's a SPOF without Paxos replication (PostgreSQL pg_prepared_xacts is the textbook liveness-failure store), (b) acquirer calls take seconds and 2PC would hold locks for that whole window, (c) cross-region 2PC needs Spanner-class storage. Considered application-level orchestration (a state machine in payment-api); rejected because workflow durability needs persistent retry + compensation that a stateless service can't provide. Temporal/Cadence give us 'durable execution' as a primitive — Uber Cadence and Airbnb Skipper both ship this exact shape.
When it fails. Worker dies mid-activity: Temporal re-schedules on heartbeat-timeout; activity_id ensures the downstream sees the same idempotency key and returns the cached response. Both forward AND compensation fail: workflow goes to DLQ + ops queue + page. Saga timeout: workflow ScheduleToClose fires → automatic compensation → 'requires_action' surfaced to merchant. Detection: workflow_failed_rate, compensation_failed_count, stuck_workflow_age_p99, DLQ depth.
- NoSQL DB · Saga State StoreCassandra 4 (Temporal default persistence)
Persists every saga: workflow execution rows, activity history events, mutable state, search-attribute indexes. Temporal worker reads/writes this on every workflow tick; the workflow is THE log — replay reconstructs state from history events. ~10 KB per workflow with ~100 events/charge lifecycle.
Why it exists. Considered Postgres for Temporal persistence (works fine up to ~1K wf/sec). Rejected at our scale: Cassandra was Temporal's default for a reason — leaderless writes scale linearly with cluster size, and the read pattern (replay history events for a single workflow) is well-suited to a partition-key-by-workflow_id design. Considered DynamoDB; works but locks us into AWS. Cassandra RF=3 LOCAL_QUORUM is the Monzo-published shape adapted for our saga store.
When it fails. Quorum loss (2/3 replicas unreachable): saga workers can't make progress; charge endpoints 503. Detection: hinted_handoff queue depth, read/write timeout rate. Mitigation: cross-AZ replica spread, sloppy quorum disabled (we want correctness over availability for money movement). Monzo's 2019-07-29 retro is the exact playbook: they brought up 6 nodes simultaneously with auto_bootstrap=false and quorum reads returned empty data — never again.
- Worker · Outbox CDC RelayDebezium + Kafka Connect
Tails the Postgres logical-replication slot on each ledger shard, decodes WAL records into Kafka events keyed by account_id, and commits the LSN to Kafka offsets. Reads ONLY the outbox table (not journal_entries directly) so the event payload is a clean domain message rather than a raw row diff. Outbox-row publication is tracked via Kafka offsets committed by the relay; outbox rows themselves are tombstoned by table partition retention after 7 days (NOT by an UPDATE in the WAL stream — that would be a feedback loop).
Why it exists. Considered an application-level outbox poller (SELECT … FROM outbox WHERE published=false LIMIT 1000); rejected at scale because polling burns DB IOPS proportional to throughput, not lag. Considered dual-write (write to DB + write to Kafka in the same code path); rejected because dual-write fails on a crash between writes — the textbook source of phantom or missing events. CDC + outbox is the Debezium-recommended pattern that gives at-least-once delivery without poisoning the DB hot path.
When it fails. Relay falls behind → WAL grows → ledger disk fills → write outage. Detection: replication_slot_lag_bytes (alert at 10 GB), oldest_unpublished_outbox_age (alert at 30 s). Worst case: relay DOWN for >24 h → slot GC → backfill from S3 + replay-with-dedup at consumers. Mitigation: scale relay horizontally per shard; alert on slot lag, not on outbox row count.
- Stream · Event BusKafka 3.7 (Confluent Platform)
Multi-topic event bus: ledger.posted (per-shard, 64 partitions), payment.captured, webhook.queued (256 partitions for fan-out), recon.completed, recon.drift_found, dead-letter (DLQ) topics for poison messages. RF=3, min.insync.replicas=2, acks=all on the producer side so a single broker loss never loses a published event.
Why it exists. Considered RabbitMQ; rejected at sustained 50K events/sec because broker memory blows up on backed-up queues and ordering across multiple producers is not guaranteed. Considered Kinesis; works but vendor-locks. Kafka's per-partition ordering is the load-bearing property for the webhook dispatcher — we need ordered delivery PER charge so 'created' lands before 'succeeded' at the merchant.
When it fails. Broker AZ loss: ISR shrinks but writes continue with min.insync=2. Two-broker loss in same AZ: producers block until quorum returns. Consumer lag spike: page when oldest-unconsumed-age > 60 s. Bad consumer poisoning a partition: DLQ after 5 failed deliveries; alert on DLQ depth growth.
- Worker · Webhook DispatcherGo + per-endpoint circuit-breaker + HMAC signing
Consumes from kafka.webhook.queued, fans out HTTPS POST to each subscribed merchant endpoint (signed with HMAC-SHA256 using the merchant's secret), retries with exponential backoff over 72 h on non-2xx responses, marks delivery success/failure in a per-endpoint state store, and gives up after 72 h to the merchant's dead-letter dashboard. Stripe's published webhook contract is the spec we follow.
Why it exists. Considered firing webhooks from payment-api directly; rejected because a 30 s merchant timeout would block the charge endpoint. Considered firing from the saga; rejected because webhook fan-out (5-20 endpoints/event × N events/charge = ~100 outbound calls per charge) doesn't fit the saga's per-step durability model. A dedicated dispatcher with its own consumer group, retries, and DLQ is the only design that survives a slow merchant.
When it fails. Slow merchant: backpressure handled by per-endpoint cap; other merchants unaffected. Cert expiry on merchant's endpoint: 4xx burst, retry exhausts, merchant's dead-letter dashboard shows the events. Dispatcher pool exhaustion: alerts at 90% saturation; auto-scale on consumer lag. Detection: webhook_delivery_p99, dlq_depth, merchant_endpoint_error_rate.
- Worker · ReconciliationSpark batch + Temporal cron-saga
Daily at 02:00 UTC: pulls each acquirer's settlement file (CSV/SFTP for legacy, JSON via API for modern), ingests into audit-store, runs a Spark job that joins our ledger postings (for that acquirer + that settlement date) against the file. Three outputs: (1) matched rows → mark as settled, (2) breaks (in our ledger but not in file, or vice versa) → exception queue, (3) FX-drift rows → P&L journal entry. Drift > $1k or > 0.01% triggers automatic settlement-freeze for that acquirer pending ops review.
Why it exists. Considered streaming reconciliation (compare each posting to a per-charge settlement event); rejected because acquirer settlement is fundamentally batched (T+1 file dropoff is the contract). Considered ad-hoc SQL queries; rejected because (a) reconciliation IS the audit story for SOX, so it has to be repeatable and append-only, (b) Adyen and Plaid both publish formal reconciliation reports that auditors expect to mirror.
When it fails. Acquirer file delayed: alert at T+1 + 4h; recon job retries every 30 min, falls back to API pull. Drift detected: ops queue + page + automatic settlement-freeze for that acquirer. Recon worker dies mid-batch: Temporal cron-saga resumes from last checkpoint. Detection: recon_completion_lag, drift_amount, exception_queue_depth.
- Object Store · Audit / SettlementS3 (versioned, SSE-KMS, Object Lock 7 yr)
Three buckets: (1) acquirer-settlement-files/ (raw uploads from PSPs, immutable), (2) ledger-cold/ (parquet snapshots of journal_entries older than 90 days, queryable via Athena), (3) recon-reports/ (per-day, per-acquirer JSON reports — the SOX audit artifact). All buckets versioned + Object Lock to prevent tampering for 7 years.
Why it exists. Considered keeping all 7 years of postings hot in Postgres; rejected because the storage cost dominates and 99% of queries are for the last 90 days. S3 + Athena gives us cold-tier query for compliance/forensic queries without inflating ledger DB. Object Lock is the auditor-mandated answer to 'how do we know your ledger hasn't been edited?'.
When it fails. S3 partial outage (1-2 hours, rare): recon-job pauses, alerts. Detection: s3_5xx_rate > 1% over 5m → P2 page; object_lock_compliance_drift > 0 → P1 (any object modified despite Object Lock — ledger tamper signal); settlement_file_pickup_lag > 4h past T+1 cron → P2 → fall back to API pull from psp-ext. Object Lock prevents accidental delete; tampering attempts trigger CloudTrail alarms with PagerDuty integration.
- MonitoringPrometheus + Grafana + PagerDuty + invariant-check job
Three planes: (1) RED metrics (rate/errors/duration) on every service, scraped at 15 s; (2) per-shard infra metrics (replication lag, disk, connection pool); (3) BUSINESS-INVARIANT alerts — a 60 s cron runs SELECT txn_group_id, currency, sum(amount) FROM journal_entries WHERE day=today AND state IN ('confirmed','reversed') GROUP BY txn_group_id, currency HAVING sum(amount) <> 0 and alerts on every group whose paired-posting balance is non-zero (Square Books / Monzo invariant). The invariant check is the load-bearing 'wake oncall up' signal because every other alert can be a gray failure; summing all amounts (the naïve query) would alert on any single-leg posting mid-write, so we group by txn_group_id and exclude pending groups (a cross-shard saga is legitimately unbalanced mid-flight).
Why it exists. Considered relying purely on RED metrics; rejected because Stripe's 2019-07-10 retro shows gray failures pass health checks but break correctness. The invariant-check is Monzo's and Square's published lesson: the ledger is the source of truth, so its invariants are the canonical SLOs.
When it fails. Monitoring down: services keep working but oncall is blind. Detection: prometheus self-monitoring + dead-man's-switch alert that fires if no metrics arrive for 2 min. Mitigation: cross-region prometheus federation; PagerDuty has its own monitoring (we don't depend on our own infra to alert).
- Cache · Rate Limit StoreRedis 7 cluster mode (token-bucket Lua scripts)
Holds per-merchant token-bucket counters (rl:{merchant_id}:tokens, rl:{merchant_id}:last_refill). Gateway issues an EVAL of a Lua script per request: refill tokens at the merchant's rate, decrement, return allow/deny in one round-trip. Sub-millisecond on hits.
Why it exists. Considered in-process token buckets at each gateway replica; rejected because a 6-replica gateway with per-replica state lets a merchant burst to 6× their nominal limit by load-balancer luck. A shared store is the only design that gives merchants a single, global rate limit. Stripe's published rate-limit infrastructure shows exactly this pattern.
When it fails. Store down → gateway fails OPEN to a permissive bucket (better to let traffic through than 5xx; the per-merchant rate limit is a soft guard, not a security control). Detection: rate_limit_store_p99, fail_open_count. Mitigation: gateway has a 100ms hold-timeout on the EVAL; circuit-breaker after 3 consecutive Redis errors.
- Service · Webhook ReceiverGo + HMAC verifier + idempotency on (psp, event_id)
Terminates POST /webhooks/{psp} from Stripe, Adyen, Visa DPS. Verifies the signature header (Stripe-Signature, X-Adyen-HMAC) using the per-PSP key from secrets, extracts the event_id, dedupes against the idempotency store keyed (psp, event_id) (idempotency claim with 30-day TTL because PSPs replay events for up to 30 days), and publishes psp.event onto kafka for the saga or recon-worker to consume.
Why it exists. Considered routing inbound webhooks through the same payment-api; rejected because (a) sig verification + dedup is a separate, sec-sensitive surface, (b) PSPs retry aggressively and a slow payment-api would back-pressure their delivery infrastructure, (c) the dedup key for inbound events is (psp, event_id) — DIFFERENT from the merchant-supplied Idempotency-Key on outbound, and conflating them in the same service is a bug factory.
When it fails. Receiver down → PSPs retry per their own schedules (Stripe = 3 days exp-backoff). Detection: receiver_5xx_rate, signature_mismatch_rate (security signal — a non-zero rate may mean key compromise). Mitigation: stateless replicas, fail-fast 5xx so PSPs back off.
- SQL DB · Ledger DR FollowerPostgres 16 + Citus (cross-region async follower)
Cross-region async follower of every ledger-db shard. Receives WAL via streaming replication; lag typically <5s. NEVER promoted automatically — operator-driven failover only, with a documented 10-minute RTO playbook. On promotion, payment-api re-points its connection strings via a config push and resumes; the WAL gap (up to ~5s of postings) is closed by reconciling against the idempotency store and acquirer settlement files.
Why it exists. Considered active-active cross-region with Spanner; out of scope (different storage engine, vendor lock). Considered no DR follower at all (re-bootstrap from base backup); RTO is hours, unacceptable for a payment system. The async-DR follower gives us a ~10-minute RTO with a recoverable RPO via reconciliation.
When it fails. DR follower itself dies: re-bootstrap from base backup + WAL archive (~2-4 hours, acceptable because primary still serves). DR region completely lost: revert to single-region operation while standing up a new DR region. Detection: dr_replication_lag_seconds (alert at 30s), dr_health. Note: an RPO=5s number on a DR follower is replication lag, NOT data-loss budget — the data-loss budget is 0 because reconciliation closes the gap on promotion.
- External · KMS / SecretsAWS KMS + HashiCorp Vault
Holds: (1) per-merchant HMAC signing keys for outbound webhooks, (2) per-acquirer mTLS client certs for the PSP adapter, (3) the seed/salt for hashing API keys at the auth layer, (4) DB connection passwords used by every service. Auto-rotation on a per-secret schedule. Each service caches secrets in-memory for 5 minutes to absorb KMS short-window unavailability.
Why it exists. Considered storing secrets in environment variables; rejected because (a) rotation requires deploys, (b) secrets in env are visible to anyone with kubectl exec, (c) audit story is non-existent. Square's 2023 cert-bundle outage (the doc cites it) is exactly the failure mode KMS + automated rotation prevents.
When it fails. KMS unreachable: services use cached secrets for up to 5 minutes; beyond that, fail-CLOSED on operations that need fresh secrets. Detection: secrets_fetch_5xx_rate, cert_days_remaining. Mitigation: cache + alert at <30 days expiry; pre-deploy hook validates all certs renewed before promote.
- Object Store · WAL ArchiveS3 + wal-g (continuous Postgres WAL backup)
Continuous Postgres WAL archive via wal-g (or pgbackrest). Captures every WAL segment from every ledger-db shard within ~30s of generation. Enables (a) point-in-time recovery for any past 30 days, (b) backfill replay for cdc-relay if its logical-replication slot is GC'd after a >24h outage.
Why it exists. Without continuous WAL archive, the cdc-relay's slot-GC failure mode is unrecoverable — we'd lose every unpublished outbox event. Considered relying on the 90-day cold ledger archive; rejected because that's a DAILY snapshot, not a continuous WAL stream, so it can't replay individual outbox events with their original ordering.
When it fails. Archive write fails: ledger-db's archive_command alarms; wal-g retries with backoff. Sustained failure (>1h) → page (we lose recovery capability). Detection: wal_archive_lag_seconds (alert at 60s), archive_command_failure_rate.
- SQL DB · Webhook SubscriptionsPostgres 16 (co-located on the auth cluster)
Holds (merchant_id, endpoint_url, event_types[], hmac_secret_ref, enabled, created_at). webhook-dispatcher SELECTs the active subscriptions for an event and fans out to each. Read-mostly: writes are merchant dashboard updates (~10 QPS). Reads are ~5K QPS at peak (one per outbound event).
Why it exists. Considered storing subscriptions as a Redis set; rejected because (a) merchant dashboards need ACID semantics on subscription edits, (b) audit log requirements (who added/removed an endpoint and when), (c) the data is small (~1M rows) and fits comfortably alongside auth on the same Postgres cluster.
When it fails. DB down: webhook-dispatcher serves from its 60s cache; beyond that, queues events into kafka (already happens) and replays when DB returns. Detection: subscriptions_read_p99, cache_staleness_seconds.
- Queue · Recon Exception QueueKafka topic + ops dashboard
Kafka topic recon.break (16 partitions, 90-day retention) plus an ops dashboard. Recon-worker publishes here when a settlement-file row doesn't match a ledger entry, when drift > $1k or > 0.01%, or when a chargeback at T+90 reverses a long-cleared txn. Each break carries the affected acquirer, charge_id, expected and actual amounts, and a runbook link.
Why it exists. Considered logging breaks to a flat file; rejected because reconciliation is the audit trail and it has to be queryable. Considered a Postgres table; rejected because at 100M txns/day even a 0.01% break rate produces 10K rows/day — Kafka's append-only model fits better, plus the ops dashboard wants a stream, not a snapshot.
When it fails. Topic unreachable: recon-worker holds breaks in local memory + spills to disk; alarms after 1h. Detection: exception_queue_depth, oldest_unacknowledged_break_age. Mitigation: human SLO is 24h to first triage on every break.
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:
- Scale tier? Card processor (Stripe-scale, ~10K charges/sec) vs neobank (Monzo-scale, ~1K charges/sec) vs marketplace add-on (Stripe Connect, mostly fan-out). Drives shard count and saga store choice.
- Card-present vs card-not-present? Card-present has 3 s consumer tolerance — sync auth only. CNP can be async (PaymentIntents pattern,
requires_actionUI). - Single currency or multi? Multi-currency means FX rate snapshot at txn time + drift recognition at settlement (separate P&L journal).
- In-region or cross-region? In-region with sync replication is the "easy" mode. Cross-region active-active for payments needs Spanner-class storage (Square Books picked this) or saga + outbox + reconciliation.
- Internal wallet only or external acquirer integration? Wallet-only collapses many failure modes; external PSP integration is where the production complexity lives.
- Reg jurisdiction? PCI DSS card data segregation, PSD2 Strong Customer Authentication in EU, RBI / Monetary Authority of Singapore caps. Affects routing and audit retention.
- Settlement cadence? T+0 (real-time, RTP/UPI), T+1 (typical card networks), T+2 (some bank rails). Drives reconciliation pipeline.
Assumptions to state:
- 100M charges/day → 1.16K charges/sec avg, 5–10K peak; ~$5B daily TPV (~$1.8T annual) at $50/charge.
- 7-year ledger retention (SOX 17 CFR 210.2-06).
- 99.95% charge-success SLO (5-min window); sync charge p99 < 3s (inline acquirer auth is ~800ms typical); ledger-commit slice p99 < 500ms.
- Single home region with cross-region async DR follower; multi-region active-active is out of scope.
02Functional reqs
What must this system actually do?
- Charge a card or wallet account; idempotent on Idempotency-Key.
- Authorize then capture (PaymentIntents pattern): 2-step flow for card-not-present.
- Refund (full or partial); compensation against an existing journal entry.
- Transfer between two accounts (P2P wallet, marketplace payout).
- Read balance (with read-your-writes for the user who just transferred).
- Receive webhooks from PSP for async events (chargeback, settlement complete).
- Send webhooks to merchants on charge state change.
- Reconcile ledger against acquirer settlement file daily (T+1).
03Non-functional
What must it promise about speed, uptime and correctness?
- Correctness: Idempotency + double-entry invariant non-negotiable. A retried charge MUST NOT double-apply.
sum(debits) = sum(credits)per account at all times. - Availability: 99.95% on the charge endpoint; 99.99% on the read endpoint. Charge endpoint may legitimately 503 under storage failover (fail-closed > drift-open).
- Latency: Charge p50 ~1s, p99 ~3s (sync — the inline acquirer auth alone is typically ~800ms, and 3s is the card-present tolerance); the ledger-commit slice budgets p99 < 500ms, which is also what the 202
requires_actionpath returns in when the PSP is slow. Balance read p99 < 80ms. - Durability: No accepted charge ever lost — RPO=0 for in-region ledger writes (sync replication mandatory). RPO~5s cross-region for DR-only.
- Auditability: Every state change recorded in the immutable journal; 7-year retention; tamper-evident (S3 Object Lock + cryptographic chain optional).
- Security: PCI DSS scope minimized — card data flows directly to PSP via tokenization, not stored in our ledger. mTLS on every internal hop. HMAC-signed webhooks.
04Capacity estimation
How much load and data does this have to hold?
Defaults: 100M charges/day, $50 avg ticket, 4 postings/charge, 5× peak multiplier.
- Charges QPS: 100M / 86,400 ≈ 1.16K avg, ~5.8K peak.
- Ledger postings QPS: 1.16K × 4 ≈ 4.6K avg, ~23K peak. Sharded Postgres+Citus 64 shards comfortably absorbs this (~360 writes/sec/shard peak).
- Annual TPV: 100M × 365 × $50 ≈ $1.83T. Stripe processed $1.4T in 2024 — same order of magnitude.
- Ledger storage (7 yr): 100M × 365 × 7 × 4 × 1KB ≈ 1.0 PB. Hot tier (90 days) ~36 TB on Postgres; rest goes to S3 cold tier.
- Idempotency hot KV: 1.16K × 24 × 3600 × 200B ≈ 20 GB for a 24h window. Comfortably in-memory if needed.
- Webhook fan-out (peak): 5.8K × 8 endpoints ≈ 46K outbound HTTPS/sec. Webhook dispatcher autoscales to absorb.
- Saga workflows in flight: A charge lifecycle averages 24h (auth → capture → settle → recon). 100M × 1 day ≈ 100M concurrent workflows worst case; Temporal sharded by workflow_id over Cassandra RF=3 absorbs.
What flips at 10× scale (~58K charges/sec peak, ~12K avg — for reference, Alipay Singles' Day 2020 hit 583K TPS):
- Ledger shards: 64 → 256+ (re-shard takes a quarter of work; plan ahead).
- Saga store: Cassandra at this scale is fine, but Temporal cluster topology matters — separate clusters per region.
- Idempotency: window stays at 24h but the hot KV moves to dedicated DynamoDB / Redis cluster.
- PSP: must integrate multiple acquirers per region with auto-failover.
05API design
What does the outside world call, and what comes back?
POST /v1/charges
Authorization: Bearer rk_live_<restricted>
Idempotency-Key: 5a8c1f2e-9b3d-4e6f-8a2b-1c3d4e5f6a7b
Content-Type: application/json
{
"amount": 5000,
"currency": "USD",
"source": "tok_visa_4242",
"merchant_account": "acct_1MZQz5",
"description": "Order #12345",
"capture": true
}
201 Created
{
"id": "ch_3MQz8jXY",
"status": "succeeded",
"amount": 5000,
"currency": "USD",
"captured": true,
"created": "2026-05-08T12:34:56Z",
"balance_transaction": "txn_1MQz8j"
}
# Retry with same key, same body — replays cached response
201 Created (same response, idempotent)
# Retry with same key, DIFFERENT body — error
400 Bad Request
{ "error": "idempotency_key_mismatch" }
# PSP slow / async confirm path
202 Accepted
{ "id": "ch_3MQz8j", "status": "requires_action", "next_action": { "type": "use_stripe_sdk" } }
# Storage failover in flight
503 Service Unavailable
Retry-After: 30
POST /v1/refunds
Idempotency-Key: <key>
{ "charge": "ch_3MQz8j", "amount": 5000 }
201 Created
POST /v1/transfers
Idempotency-Key: <key>
{
"source_account": "acct_A",
"destination_account": "acct_B",
"amount": 10000,
"currency": "USD"
}
# Cross-shard transfer runs as a saga (debit, then credit). The debit posts as
# 'pending' and the funds show as in-transit until the credit confirms.
202 Accepted
{
"id": "tr_3MQz8j",
"status": "pending", # funds in transit; debit posted, credit not yet confirmed
"amount": 10000,
"currency": "USD"
}
# Saga confirms (or compensates) asynchronously; merchant learns the outcome via webhook:
# POST {merchant_url} { "type": "transfer.succeeded" | "transfer.reversed", ... }
GET /v1/balances/{account}
200 { "available": 12345, "pending": 2500, "currency": "USD" }
# Inbound webhook from PSP
POST /webhooks/psp
Stripe-Signature: t=...,v1=...
{ "id": "evt_1MQz8j", "type": "charge.captured", ... }
# Outbound webhook to merchant
POST {merchant_url}
X-Webhook-Signature: t=...,v1=...
{ "id": "evt_1MQz8j", "type": "charge.succeeded", ... }06Data model
What gets stored, and what is it looked up by?
Idempotency rows (Postgres, sharded by hash(key)):
| field | type | notes |
|---|---|---|
| key | varchar(255) | PK (Idempotency-Key header) |
| merchant_id | bigint | FK; idempotency is per-merchant |
| request_fingerprint | bytea | hash(method, path, body) |
| recovery_point | enum | started, ledger_committed, psp_called, finished |
| response_code | int? | cached on completion |
| response_body | jsonb? | cached on completion |
| locked_at | timestamp? | NULL when unlocked; row lock held |
| created_at | timestamp | |
| expires_at | timestamp | now() + 24h; reaper sweeps |
Journal entries (Postgres + Citus, sharded by account_id, append-only):
| field | type | notes |
|---|---|---|
| journal_id | uuid | PK |
| account_id | bigint | shard key |
| amount | numeric(20,4) | signed (debit negative, credit positive) |
| currency | char(3) | ISO-4217 |
| txn_group_id | uuid | groups paired postings (debit + credit) |
| txn_type | enum | charge, refund, transfer, fee, fx_drift |
| state | enum | pending, confirmed, reversed — written pending on synchronous commit; the saga advances it forward to confirmed after PSP confirms, or to reversed on compensation. Forward-only; beyond it only the confirmed_at stamp is ever written |
| ref_charge_id | varchar? | null for fees / FX adjustments |
| posted_at | timestamp | indexed; partitioned by day |
| confirmed_at | timestamp? | null until saga confirms PSP success |
| metadata | jsonb |
State-machine note: The state column — plus the confirmed_at stamp that rides the confirm transition — is the only thing that ever changes on a journal row, and state only moves forward. A successful charge writes pending rows at e6; the saga UPDATEs the group to confirmed (stamping confirmed_at) after PSP confirms. Compensation flips the group to reversed AND appends offsetting postings in the same txn_group — money is undone by a new row, never by editing an amount. The invariant sum(amount) GROUP BY txn_group_id = 0 holds for confirmed and reversed groups; pending groups may be unbalanced for the seconds a cross-shard saga is mid-flight, which the invariant SQL accounts for by filtering WHERE state IN ('confirmed','reversed').
Balance projections (mutable, same shard as journal):
| field | type | notes |
|---|---|---|
| account_id | bigint | PK |
| currency | char(3) | composite PK with account_id |
| available | numeric(20,4) | |
| pending | numeric(20,4) | |
| version | bigint | optimistic-concurrency token |
| last_journal | uuid | last posted journal_id |
Outbox (same shard as journal, drained by Debezium):
| field | type | notes |
|---|---|---|
| event_id | uuid | PK |
| topic | varchar | ledger.posted, webhook.queued |
| key | varchar | partition key (charge_id, account_id) |
| payload | jsonb | |
| created_at | timestamp |
Indexes:
- Idempotency:
(merchant_id, key)PK (idempotency is per-merchant — a bare(key)PK would let two merchants that happen to pick the same UUID collide, so merchant B's CLAIM would land on merchant A's row); partial onWHERE expires_at < ...for reaper. - Journal: (account_id, posted_at DESC) covering for balance computation; (txn_group_id) for paired-posting lookup; (ref_charge_id) for charge-history.
- Outbox: none beyond the PK — delivery position is the relay's committed Kafka offset / replication-slot LSN, and rows age out via 7-day partition retention (no
publishedflag; flipping one would feed an UPDATE back into the WAL stream the relay is tailing).
Database choice — recommended:
Postgres 16 + Citus (sharded by account_id) for the ledger, idempotency, and outbox — co-located on the same engine so the journal+balance+outbox triple-write fits in one SERIALIZABLE transaction. Cross-shard transfers handled by saga + outbox, NOT 2PC across shards. Cassandra (RF=3 LOCAL_QUORUM) for the saga store (Temporal default; partition by workflow_id; Monzo-pattern for write linearity at scale). Redis 7 cluster for balance cache only — not source of truth. S3 + Object Lock for cold ledger archive and acquirer settlement files.
Why not Spanner everywhere? Cross-shard 2PC is real and tempting (Square Books took this trade). The cost: Spanner write p99 30-100ms vs Postgres 5-15ms in-region; vendor lock to GCP; loss of operational familiarity for ledger-shape queries. Take it only if cross-shard atomic txns dominate the workload — most payment systems live within single-account writes plus saga'd transfers.
Why not DynamoDB? Works for the idempotency store (single-key conditional writes are perfect for CLAIM). Doesn't work for the ledger because the journal+balance+outbox triple-write needs cross-row transactionality at scale (DynamoDB's TransactWriteItems caps at 25 items, and at peak we'd be hitting 4+ rows per charge × thousands/sec).
07High-level design
Which components handle a request, and in what order?
Architecture summary:
- API Gateway (Envoy + Cloudflare WAF): TLS terminate, per-merchant rate limit BEFORE forwarding (so retry storms can't poison the idempotency store), forward Idempotency-Key verbatim.
- Auth / Identity: Restricted-key validation behind the gateway; signed envelope to upstream services so payment-api doesn't re-validate per call.
- Payment API · Write Plane: Hosts every mutating endpoint. Owns the Brandur-style idempotency state machine. Splits sync (PaymentIntent + immediate capture) from async (saga-driven capture).
- Idempotency Store (Postgres KV): Per-row state per Idempotency-Key. SAME engine as ledger so the idem.row + journal commit in one txn.
- Ledger DB (Postgres + Citus, sharded): Source of truth. Three tables per shard: journal_entries (append-only audit log), balance_proj (running balance), outbox (event drain to Kafka). One SERIALIZABLE txn writes all three.
- Balance Reader · Read Plane: Cache-aside through Redis Cluster; falls back to ledger replicas. RYW token from payment-api routes recent-write reads to the leader.
- Saga Orchestrator (Temporal): Owns multi-step money movement (authorize → capture → settle → reconcile). Workflow durability + explicit compensations replace 2PC for cross-service writes.
- Saga State Store (Cassandra): Workflow history events. RF=3 LOCAL_QUORUM. Partitioned by workflow_id.
- PSP Adapter: Stateless protocol bridge to acquirers / card networks. Per-acquirer circuit-breaker; cert pinning; HMAC.
- Outbox CDC Relay (Debezium): Tails the ledger WAL and ships outbox rows to Kafka. The reason we don't dual-write.
- Event Bus (Kafka, RF=3, 256 partitions on the webhook topic): Per-charge ordering preserved. acks=all + min.insync.replicas=2 makes each publish durable; the idempotent producer removes producer-internal retry duplicates only — end-to-end delivery is at-least-once, and consumers dedupe by event_id.
- Webhook Dispatcher: HTTPS fan-out to merchants. Per-endpoint circuit-breaker + 72h retry schedule.
- Reconciliation Worker: T+1 cron-saga that pulls acquirer settlement files, joins against the ledger, and emits drift / settle events. The SOX audit trail.
- Audit / Settlement Store (S3 + Object Lock): 7-year tamper-proof archive of ledger snapshots, settlement files, recon reports.
- Monitoring (Prometheus + Grafana + invariant-check job): RED metrics + the 60-second
sum(debits)+sum(credits)=0invariant signal. The latter is what wakes oncall up at 3am.
Data flow on a synchronous charge: Client → Gateway (auth + rate limit) → Payment API → Idempotency Store (CLAIM) → Ledger DB (one SERIALIZABLE txn writes journal + balance + outbox) → Balance Cache (INVALIDATE) → PSP Adapter → External PSP. Idempotency row marked finished only after PSP returns; client gets the cached response on retry. Outbox row drains via CDC → Kafka → Webhook Dispatcher → Merchant.
Data flow on a saga-driven async charge: Client → Gateway → Payment API → Idem CLAIM → Ledger commit (PENDING posting) → Saga signal → return 202 requires_action to client. Saga then drives PSP capture activity; on success, payment-api activity posts the CAPTURED journal entry, balance cache invalidates, outbox publishes — webhook fires.
08Deep dives
Which part breaks first, and what do you do about it?
1. Idempotency-key state machine (the canonical bug.) The idempotency contract is "same key + same body = same response." The hard part is guaranteeing it under crashes. The Brandur recipe: every key carries a recovery_point enum that advances atomically with the business write. Started → ledger_committed → psp_called → finished. A retry sees the current state, resumes from the right step, and locks the row for the duration via a 25 s lease so concurrent retries serialize. Same-key + DIFFERENT body returns 400 idempotency_error (the body doesn't match the original request the key claimed; Airbnb Orpheus rejects this too). 409 is reserved for the distinct case where a second request arrives with the same key while the first is still in-flight (lock held) — the concurrent-retry conflict. The rare case where this fails: the request fingerprint is computed at the wrong granularity (e.g. excludes a metadata field that mattered). Test invariant: forall keys k: response(k, retry_n) = response(k, 0).
2. Double-entry invariant under failover. The core double-entry invariant is that every transaction balances: for each journal entry (txn_group), sum(postings) = 0 per currency — total debits equal total credits across the accounts it touches — and system-wide sum(all postings) = 0 per currency. (Per account the running sum is that account's balance, which is generally non-zero — an account that only ever receives credits never sums to zero.) Because both legs of a transfer are written in ONE local DB transaction they can't half-commit; the dangerous scenario is at the replication layer: primary commits both legs, replicates the debit to the follower, crashes BEFORE replicating the credit — now the follower's history holds an unbalanced transaction (a lone debit) = drift. Async replication = drift. Mitigation: in-region replication MUST be sync (RPO=0) so a transaction never half-replicates. The cross-region async follower is DR-only; we never serve writes from it without an explicit operator failover. Monzo's published 2019-07-29 retro is the textbook example: their Cassandra scale-out with auto_bootstrap=false returned empty quorum reads — same shape of bug, different storage layer.
3. Saga vs 2PC for cross-account transfers. Account A on shard 17 wants to send $100 to account B on shard 4. Two options:
- 2PC: coordinator opens distributed txn, both shards PREPARE, both COMMIT. Atomic, but coordinator is a SPOF (Postgres
pg_prepared_xactsis the canonical liveness-failure catalogue). Cross-region multiplies the lock-holding time. - Saga: workflow posts the debit on shard 17 first; on success, posts the credit on shard 4; on failure of credit, posts a COMPENSATING credit on shard 17 to undo. NOT atomic — for the seconds between the two postings, the balances disagree (debit shows up before credit). Mitigation: a "pending" intermediate state on the debit side; clients see the funds "in transit." Saga wins when liveness matters more than instantaneous atomicity, when cross-region is in scope, or when the destination is an external PSP (you can't 2PC with Stripe's API). Spanner-class storage gives you 2PC inside ONE storage engine; everything else is saga + reconciliation. Uber Cadence and Airbnb Skipper publish exactly this trade.
4. Reconciliation drift and TSB-2018 lessons. T+1 settlement file from an acquirer says they sent $1.030M; our ledger says we charged $1.034M. Drift could be: a chargeback we processed but acquirer hadn't, a fee we accrued differently, FX drift between charge-time and settle-time, or genuine missing funds. Run a per-charge join: matches → mark settled, breaks → exception queue, FX-only deltas → P&L journal. Drift > $1k OR > 0.01% triggers automatic settlement-freeze for that acquirer. TSB's 2018 multi-day drift was the catastrophic version of this: their migration produced inconsistent records that took weeks to reconcile, and the reconciliation pipeline didn't have the freeze-and-page mechanism in place.
5. Hot treasury / clearing account. Every charge debits the platform's clearing account and credits the merchant. That clearing account's row is touched by 100% of write QPS. Two mitigations:
- Sub-account sharding: clearing account_clear:0 .. account_clear:N, randomly chosen per charge; an hourly sweeper rolls them into the master account_clear. Trades exact-time-balance for write-distribution.
- Append-only running balance: never read the balance_proj; compute on demand from the journal stream. Costs read latency; eliminates the single-row contention. Square Books does some variant of this.
6. PSP 503 mid-charge (Visa Europe 2018, Stripe 2019-07-10). Acquirer returns 503 for 10 minutes. With circuit-breaker open on PSP Adapter, payment-api can't authorize sync; we return PaymentIntent requires_action and persist a saga. Saga retries via fallback acquirer if configured, or holds the workflow open until the primary recovers (up to 7 days). User experience: charge is in processing state on the dashboard; webhook fires when resolved. The load-bearing piece is that we NEVER write a successful journal entry for a charge whose PSP confirmation didn't arrive — the journal stays in pending until the saga confirms.
7. Outbox + CDC vs dual-write. Three patterns for "ledger write + downstream event":
- Dual-write: payment-api writes to ledger, then writes to Kafka. If the process crashes in between, event lost. Production never picks this.
- Application outbox poll: SELECT FROM outbox WHERE published=false. Burns IOPS proportional to throughput, not lag.
- CDC + outbox: Debezium tails the WAL, publishes outbox events. Ledger writes pay no event-publish cost; the relay is the only consumer of the WAL slot. Failure mode: relay falls behind, WAL grows, ledger disk fills. Detection: replication_slot_lag_bytes alert at 10 GB. Stripe-scale picks this; Debezium publishes the recipe.
09Trade-offs
What did this design cost, and what breaks at 10×?
What breaks at 10× scale (~58K charges/sec, ~232K postings/sec):
- Ledger shards: 64 → 512 shards (NOT 256 as a naïve scaling). At 10× postings × 12 row-writes per charge = 696K row-writes/sec. Per-shard SERIALIZABLE row-write rate at 256 shards would be ~2.7K/sec/shard — over Citus's published comfortable ceiling (~1.5K/sec/shard before serialization conflicts spike). Either re-shard to 512 (1.4K/sec/shard) OR relax journal_entries to SNAPSHOT isolation with the SERIALIZABLE invariant enforced on balance_proj only. Cross-shard saga rate climbs proportionally → saga store (Cassandra) needs cluster expansion + careful partition key review.
- Idempotency store: Postgres KV starts to feel hot at ~20K QPS sustained; consider DynamoDB conditional writes with an off-cluster Redis cache for fast hits.
- Webhook fan-out: 50K outbound HTTPS/sec stresses egress capacity and merchant endpoints; per-merchant rate limits become first-class.
- Reconciliation join: daily Spark job scales linearly with txn count; partition by acquirer + day so jobs stay parallelizable.
Per-component failure stories:
- Idempotency store down: payment-api fail-CLOSED, all writes 503. Better than fail-open (would risk drift). Mitigation: in-region sync replication on the idem store, just like the ledger.
- Ledger primary failover: 30 s of write outage per affected shard; clients see 503 + Retry-After. RPO=0 in-region. Cross-region failover is DR-only because RPO~5s is unacceptable in the steady state.
- PSP outage: circuit-breaker opens, sync charges return
requires_action, saga retries on fallback acquirer. 100% sync charge availability lost; 100% wallet-only / cached-PSP work continues. Visa Europe 2018 (5.2M txns over 10h) is the pattern. - Cassandra (saga store) quorum loss: sagas can't progress, charge endpoints 503. Mitigation: 3-AZ spread, never relax to ONE on writes.
- Kafka broker AZ loss: ISR shrinks; writes continue with min.insync=2. Two-broker loss in same AZ stops producers — webhook dispatcher backs up, retry budget consumed.
- Webhook dispatcher overwhelm: per-endpoint cap prevents one slow merchant from starving others; sustained merchant outage drains to DLQ + dashboard.
- Recon drift detected: automatic settlement-freeze for the affected acquirer; ops queue + P1 page. Better to halt settlement than to silently accumulate drift.
- Bad ledger deploy (Knight Capital pattern): Stripe-style online-migration shadow-write + Scientist read-compare gates promotion. Catches wrong-fee bugs before customers see them.
Trade-offs we accepted:
- Saga over 2PC: chose liveness over atomicity for cross-account / cross-region transfers. Brief windows of "funds in transit" are visible to clients; reconciled within seconds in the happy path, within saga ScheduleToClose (7 days) in the worst.
- Single home region: active-active multi-region for payments needs Spanner-class storage; we picked Postgres + saga. Cross-region DR failover is operator-driven (RTO ~10 min); the 5 s RPO number is replication lag, NOT data-loss budget — the data-loss budget is 0 because reconciliation closes the gap on promotion.
- Eventual webhook delivery: webhooks are at-least-once with up to 72 h retry; merchants must dedupe by event_id. Industry-standard contract (Stripe / Adyen / Plaid all behave this way).
- 24 h idempotency window: beyond 24 h, the same key replays as a new charge. Bug-fix headroom of 48 h soft retain (Brandur). At 7 days a misconfigured client could double-charge — clients are responsible for not retrying that long.
- Pending state visible: the journal carries a
pendingstate for the 25 s a synchronous PSP call may take. Reads filtered tostate IN ('confirmed','reversed')for invariant checks; reads to "what's the user's balance right now?" sumpendinginto a separate "in-flight" total displayed aspendingto the user (Stripe PaymentIntent UX).
Primary sources
- Brandur — Implementing Stripe-like Idempotency Keys in Postgres
- Stripe — Designing robust APIs with idempotency
- Square — Books: an immutable double-entry accounting database
- Monzo — Speeding up balance reads (ledger)
- Uber — Engineering Uber's Next-Gen Payments Platform (Gulfstream)
- Airbnb — Avoiding double payments (Orpheus)
- Garcia-Molina & Salem — Sagas (1987)
- Debezium — Outbox pattern
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 Payment / Wallet System yourselfMore in Transactions, Concurrency & Money
Correctness when two writers collide and money is involved: serializability, two-phase commit versus sagas, hold-then-confirm, single-writer matching, and the databases that give you external consistency.