Why Identity Platforms Become Distributed Systems Whether You Want To Or Not

Nobody sets out to build a distributed system. What people build is a service with a database, and then they solve six problems in a row, each of which has an obvious solution, and eighteen months later they are running a distributed system with none of the design properties of one.

Nobody sets out to build a distributed system. What people build is a service with a database, and then they solve six problems in a row, each of which has an obvious solution, and eighteen months later they are running a distributed system with none of the design properties of one.

The tell is a bug report that sounds impossible. "I disabled the user, but they made an API call four minutes later." "The admin UI shows the new redirect URI, but logins still fail with the old one." "Two sessions exist for the same login and only one of them can be logged out." Nobody can reproduce any of these locally, because locally there is one process and one database and none of these bugs can occur.

This article is the sequence — the six steps, each locally justified, and the specific invariant each one silently removes. The value in walking it isn't to argue against any step. Every step is correct. It's that each one has a known consequence, and if you can see it arriving you can design for it in an afternoon instead of debugging it for a quarter.

Step 0: one service, one database

Every request reads from and writes to Postgres. Configuration changes take effect on the next query. Sessions are rows. Read-your-writes is free. There is exactly one clock. Ordering is whatever the database says it is, and the database is authoritative about everything.

This system has properties that are extremely easy to take for granted, and the entire rest of this article is about losing them one at a time:

  • A write is visible to all subsequent reads. Immediately.
  • There is one source of truth for every fact.
  • "Now" means the same thing everywhere.
  • There is no partial failure. Either the request worked or it didn't.

Step 1: add a cache

Login latency is 180ms and 140ms of it is loading tenant configuration on every request. The config changes maybe twice a month. Caching it in-process for 60 seconds takes twenty minutes to implement and drops p50 to 45ms.

Obviously correct. Do it.

What you just lost: read-your-writes. An administrator changes a setting and their next request may be served by a pod holding the old value. If you have twelve pods with independent in-process caches, there are now up to twelve versions of the truth, and which one a request sees depends on load-balancer routing.

The bug this produces is the "I saved it but it didn't work" ticket, and it's maddening because it's intermittent by construction: refresh the page and you might land on a pod with fresh data. Support cannot reproduce it, and the engineer who investigates finds the correct value in the database and closes the ticket as user error.

What to do about it now rather than later: make the cached value's version visible. A config version number, returned in a debug header or logged with each request, converts "it's not working" into "this pod is serving v41, the database has v42." That single piece of instrumentation is the difference between a five-minute diagnosis and a week.

Step 2: share the cache

Twelve independent caches is wasteful and inconsistent, so move to Redis. Now there's one cache, one version of the cached truth, and a cache warm enough that misses are rare.

What you just lost: the ability to fail independently. Redis is now on the critical path of every login. Its availability multiplies into yours; its latency adds to yours; and a Redis failover of eight seconds is eight seconds of failed logins unless every caller handles it.

Which produces the first genuine distributed-systems decision, and it's one most teams make by default rather than deliberately: when the cache is unavailable, do you fail or fall through to the database? Falling through is obviously right — until you notice that at your traffic level the database cannot serve 100% of reads, so a cache outage becomes a database outage, and a database outage becomes a total outage. Falling through with a concurrency limit and a circuit breaker is the actual right answer, and it's meaningfully more code than a plain fallback.

You've also acquired the thundering herd: when a popular key expires, every in-flight request misses simultaneously and hits the database together. Single-flight or probabilistic early refresh, plus TTL jitter, are the standard fixes, and they're cheap when written on purpose and expensive when written during an incident.

Step 3: share the session store

Sticky sessions were fine until deploys started logging people out. So sessions move into the shared store — Redis, or a database table, or both.

What you just lost: session state is now a distributed object with a write path. Reads are frequent, writes are frequent (last-accessed timestamps, idle-timeout updates, sliding expiry), and two concurrent requests from the same user can now race. Two tabs refreshing simultaneously, with refresh-token rotation, is a genuine lost-update problem: one tab's rotation invalidates the other's token, the other retries, and the reuse-detection logic — which is supposed to be a theft alarm — kills the whole token family. The user gets logged out and your dashboard records a security event that wasn't one.

The fix is atomic operations for anything that mutates session state (Lua scripts, compare-and-swap, or a short grace window on rotation where the immediately-previous token is accepted once), and it needs to be designed in, because retrofitting atomicity into a code path that assumes it has exclusive access means auditing every mutation.

And you lost bounded blast radius. The session store is now the single component whose loss logs out every user simultaneously. Which is worth naming explicitly, because it changes what "graceful" means for that dependency: a session store failover that drops the dataset isn't a degradation, it's a company-wide logout.

Step 4: add a read replica

Read traffic outgrows the primary. Identity is overwhelmingly read-heavy, so replicas are the obvious lever.

What you just lost: a single point in time. Replication lag is normally 5ms and occasionally 30 seconds — during a bulk import, a schema migration, a vacuum, or a network hiccup. And now:

  • A user is created against the primary; the login attempt reads a replica and the account doesn't exist yet. The user's first-ever login fails.
  • A credential is revoked on the primary; validations against a replica keep succeeding.
  • A password is changed; the old one keeps working for the length of the lag.

That third one is the one that ends up in a security review, because "how long does a revoked credential keep working" now has an answer that includes replication lag, and nobody wrote it down.

The discipline that fixes this is unglamorous and mechanical: classify every read as "may be stale" or "must be current," and route accordingly. Session validation and profile reads tolerate staleness. Credential verification, revocation checks, and anything immediately following a write must go to the primary. Most codebases don't have this classification and instead have a single db handle that someone pointed at a replica, which means the classification exists implicitly and nobody knows what it says.

Step 5: add an event bus

Downstream systems need to know when users change. Provisioning needs to react to a new account, the audit pipeline needs every event, the cache needs to invalidate when config changes. Synchronous calls to all of them would couple your login path to their availability, so: publish events.

What you just lost: ordering and exactly-once delivery. Events arrive out of order, duplicated, or not at all. user.updated before user.created. user.deleted processed twice. A config.changed event dropped during a broker failover, leaving one region's cache stale indefinitely.

Every consumer now needs to be idempotent and order-insensitive, which is a real design constraint and one that's much easier to satisfy if events carry enough state to be self-describing: a version or sequence number per entity, so a consumer can discard an event older than what it already applied. Events that carry only "this entity changed, go look it up" are simpler to publish and push the ordering problem onto every consumer — which is fine until a consumer's lookup races the next change.

And you lost the ability to reason about "when." Cache invalidation by event means invalidation is now asynchronous and best-effort. Which is a good design — it's far better than TTL-only — but its correctness depends on the bus, so the TTL has to remain as a backstop. A system whose only invalidation mechanism is events has a permanent-staleness failure mode whenever an event is dropped.

Step 6: add a second region

Latency for European users, or a legal residency requirement, or a disaster-recovery commitment. Now you have two of everything and a wide-area link between them.

What you just lost: the option of being simple. Every question above returns, harder:

  • Cross-region replication lag is 50–200ms on a good day, and unbounded when a link degrades. Every staleness window widens.
  • Revocation must be global, and global means crossing the WAN, which means revocation is either slow or region-local.
  • Split brain is now possible: two regions accepting writes to the same entity during a partition, and someone has to decide what happens when they reconcile.
  • Data residency means some data legally cannot replicate, so "the same platform" now has different data in different places by design.

This is the step where teams finally admit they're running a distributed system, which is roughly four steps too late.

The pattern

flowchart TD
    A["1 service + 1 DB<br/><i>read-your-writes, one clock</i>"] -->|"add cache"| B["– immediate consistency"]
    B -->|"share cache"| C["– independent failure"]
    C -->|"share sessions"| D["– exclusive access to state"]
    D -->|"read replica"| E["– single point in time"]
    E -->|"event bus"| F["– ordering, exactly-once"]
    F -->|"second region"| G["– bounded staleness,<br/>single-writer"]

Each arrow is a good decision. Each arrow removes a guarantee that the code above it was written assuming. That's the mechanism: the invariant is removed from the infrastructure but the assumption stays in the code, and nobody audits for assumptions because they were never written down. They were just true.

Why identity gets hit harder than most services

Every service goes through some version of this. Identity has four properties that make each step more consequential.

Everything depends on it, so its staleness is everyone's staleness. A stale product catalogue shows a wrong price. A stale permission cache grants access that was revoked. Downstream systems can't compensate, because they have no independent way to know.

Its writes are security events. For most services, a delayed write is a UX problem. For identity, "revoke this credential" and "disable this account" are the operations whose latency a security team will ask about, and eventual consistency is a much harder thing to defend when the eventual outcome is "the attacker keeps working."

The read:write ratio hides the problem. At 10,000:1, almost nothing you do exercises the write path's consistency behaviour. Your tests pass, your load tests pass, and the 0.01% of operations that expose the inconsistency are administrative actions performed by a handful of people — who report bugs nobody can reproduce.

Correctness is bimodal. A recommendation engine that's slightly stale is slightly worse. An authorization decision that's stale is wrong, and it's wrong in one of the two directions that generate an incident: someone can't work, or someone can do something they shouldn't.

What to do with this

Not "avoid these steps" — you can't, and they're right. Four things, all cheap if done at the time:

Write down each invariant as you break it. A short document: "as of the config cache, a setting change may take up to 60 seconds to apply on all pods." Six lines, one per step. This document is the most valuable artifact in the system because it's the thing new engineers don't know and existing engineers have stopped noticing. It also becomes the honest answer to "what's your revocation window," which someone will ask.

Classify reads by staleness tolerance, explicitly, in code. Not by which handle happened to be injected. A mustBeCurrent flag or two distinct repository interfaces makes the decision visible at the call site and reviewable in a diff.

Make version visible everywhere. Config version, cache generation, replica lag, event lag — exported as metrics and attached to request logs. Nearly every impossible-sounding bug in this space is diagnosable in a minute if you can see which version served the request, and takes a week if you can't.

Design the degradation, per operation, before you need it. Cache down: fall through with a concurrency limit. Session store down: fail closed, because the alternative is unauthenticated access. Event bus down: TTLs are the backstop, and alert on consumer lag. Replica lagging: route critical reads to the primary automatically past a lag threshold. Each of these is a small amount of code and a decision that's much better made calmly.

The uncomfortable summary: you are going to build a distributed system. The only variable is whether you notice while you're doing it. Every one of those six steps takes an afternoon to implement and comes with a consequence that takes a quarter to debug — and the consequence is entirely predictable from the step, which means the whole cost is avoidable by writing one paragraph at the time you merge.