Why Identity Systems Fail Gradually, Then All at Once

The database failover took nineteen seconds.

I want to be precise about that, because it's the whole point of the story. At 14:12 a managed Postgres instance behind an identity platform promoted a replica. Nineteen seconds later the new primary was accepting connections, the query p99 was back to 3ms, and the dashboard for the database — every panel on it — was green and stayed green for the rest of the day.

At 14:31, login was still failing. At 15:05, login was still failing. At 15:40 the incident channel had thirty-one people in it, four of whom were now arguing about whether the failover had really completed, because the alternative explanation was that a healthy database, a healthy network, a healthy Kubernetes cluster, and a fleet of application servers pinned at 100% CPU were somehow producing a total authentication outage with no cause anywhere on any graph.

The failover was not the cause. The failover was the trigger, and it had been over for eighty-eight minutes. What was keeping the platform down was the platform's own clients, doing exactly what they had been carefully engineered to do: retrying.

This failure mode has a name. It's called metastable failure, and if you have run a service of any size you have almost certainly lived through one. What most engineers don't have is the model — and without the model, the incident is genuinely unsolvable, because every instinct you have during it is wrong. You look for a cause that no longer exists. You add capacity that gets eaten instantly. You wait for it to drain, and it doesn't drain. Eventually someone does something drastic and slightly illegal-feeling, and it comes back, and the postmortem says "we restarted the fleet" under Resolution, which explains nothing.

This article is about the mechanism, the arithmetic, and the defenses — with the trade-offs each defense actually costs you.

Two stable states

Here is the definition, stated carefully, because the precision is what makes it useful.

A metastable failure occurs in a system that has two stable states: a healthy state, in which work arrives and is served, and a congested state, in which the system is fully occupied and producing little or no useful output. Both states are self-sustaining. A trigger — a burst, a failover, a slow dependency, a bad deploy — pushes the system from the first into the second. A sustaining effect — almost always retries — holds it in the second state after the trigger is gone.

The defining property, the one that makes these incidents feel supernatural: removing the original cause does not restore service. The system does not return on its own. It has found a different equilibrium and it is perfectly happy there, burning every core you own, forever.

This is hysteresis. The load at which the system falls into congestion is much higher than the load at which it will climb out, and the gap between those two numbers is where your outage lives.

The word "metastable" is borrowed from physics, and the analogy is worth ten paragraphs of explanation. Supercooled water is liquid below 0°C. It is stable — you can hold it there indefinitely. Tap the glass, and it flashes to ice in a second. Now warm it back to the temperature it was at before you tapped it. It does not turn back into water. The trigger is gone; the state persists. To get liquid water back you have to go somewhere the system has never been: well above the freezing point, which costs energy you were not spending before.

For a service, "well above freezing" means dropping work you have already accepted. That is the entire recovery strategy, and we'll get to why nothing else works.

Why identity is unusually prone to this

Every request-serving system can go metastable. Identity systems are structurally worse, for four reasons that compound.

The work is expensive and looks non-sheddable. A password verification is 50–100ms of deliberately-hard CPU. You cannot make it cheaper without making it weaker, so your capacity per core is small and fixed, and the gap between "fine" and "saturated" is narrow. The cost of password hashing covers why that number is what it is. What matters here is that it makes each wasted request expensive: a discarded HTTP 404 costs you microseconds, a discarded login costs you 90 milliseconds of a core you will never get back.

Every client retries authentication automatically. Not as a policy decision — as a default. Mobile SDKs retry. HTTP clients retry. Service meshes retry, sometimes invisibly and sometimes twice at two different layers. Nobody wrote a design document approving this; it is simply what the libraries do out of the box, and it is normally correct.

Failures are synchronous and user-visible, so humans retry too. This is the multiplier nobody models. A user staring at a spinner on a login screen will hit refresh, then hit it again, then close the tab and open a new one. Your beautifully-tuned client retry budget does not govern a person's index finger. If your SDK does three attempts and your user does three page loads, one intended login is nine requests, and the human layer has no backoff, no jitter, and no cap.

Fan-in converts one slowdown into retries from everything at once. Identity sits under every application in the estate. When it slows down, it does not receive a retry storm from one client population — it receives simultaneous, uncorrelated retry storms from forty applications, each with its own retry policy, none of which coordinate with each other. This is the same fan-in that makes tail latency so brutal in authentication (percentiles do not compose the way people assume), and it has the same root: identity is the shared dependency everyone forgets they have.

The arithmetic

Numbers make this concrete in a way prose can't. Take a plausible system:

  • 100 cores dedicated to credential verification, Argon2id at 90ms per hash.
  • Capacity therefore ≈ 1,000 logins/second. Call it C = 1,000.
  • Steady-state offered load A = 800/s, so utilization ρ = 0.8. Comfortable-looking. Most capacity plans target exactly this.
  • Client SDK: 2-second timeout, up to 3 total attempts.

Now the trigger. Nineteen seconds during which effective capacity drops to 400/s. Simple accumulation: arrivals continue at 800/s, service runs at 400/s, so the backlog grows by 400/s for 19 seconds — about 7,600 requests sitting in a queue in front of a pool that can drain 1,000/s. Queue wait is now 7.6 seconds.

The client timeout is 2 seconds.

That single inequality is the entire failure. Every request in that queue will be abandoned by its client before the server gets to it — and the server will then hash it anyway, spending 90ms of a core producing a correct answer for a socket that closed five seconds ago. Capacity spent on abandoned work is capacity permanently removed from the system. It is worse than idle capacity, because idle capacity is available.

Meanwhile each of those abandoned requests generates a retry. Let p be the probability a request fails, and let offered load be:

L = A × (1 + p + p²) — the original attempt plus up to two retries.

Now find the equilibria:

p (failure rate) Offered load L vs. capacity C = 1,000 Consistent?
0 800 under capacity → requests succeed → p stays 0 Yes — healthy stable state
0.5 1,400 over capacity → failures continue → p rises No, unstable
1.0 2,400 2.4× capacity → everything fails → p stays 1 Yes — congested stable state

Two fixed points. Both self-consistent. The system is a bistable switch, and a nineteen-second failover is enough to flip it.

Look at what the congested state requires to escape. To get back under capacity you need L < 1,000, which means 800(1 + p + p²) < 1,000, which means p < 0.22. But p cannot fall below 1 until queue wait drops under 2 seconds, and queue wait cannot drop until L < C. Each condition is waiting on the other. That circularity is not a metaphor — it is why the system sits there for eighty-eight minutes with a healthy database.

And notice the number that makes it inescapable: the retry multiplier at full failure is 3, and 3 × 800 = 2,400, which is 2.4× capacity. If the client had done one retry instead of two, the congested load would be 1,600 — still above capacity, still stuck. If the base load had been 300/s instead of 800/s, congested load would be 900/s, below capacity, and the system would have drained on its own in a couple of minutes and nobody would have opened an incident. The same code, the same trigger, a different steady-state utilization, and the outage simply doesn't happen. That is why "we've had this trigger before and it was fine" is not evidence of anything.

The curve that turns over

Here is the property that breaks people's intuition. Plot goodput — responses that arrive before the client gave up — against offered load:

Offered load Queue wait Throughput (completions/s) Goodput (useful/s)
500/s ~ms 500 500
800/s ~360ms 800 800
950/s ~1.7s 950 ~950
1,000/s unbounded growth 1,000 falling
1,200/s > 2s 1,000 0
2,400/s > 2s 1,000 0

Throughput is monotonic and then flat — the server completes 1,000 requests per second no matter how hard you push it, and every instrument you own will report that it is working at full capacity. Goodput rises, peaks somewhere just below saturation, and then falls off a cliff to zero.

That non-monotonicity is the whole problem, and it has two consequences worth stating separately.

First: the peak of the goodput curve is not at maximum load. There is a load at which you serve the most users, and pushing past it serves fewer. A system operating past that peak is strictly worse than the same system with a bouncer at the door turning people away. This is the formal version of the argument in tail latency about running at 30–50% utilization: the headroom isn't waste, it's the distance between you and the cliff edge.

Second: every dashboard measuring throughput shows a system at 100% health during a total outage. CPU: 100%. Requests completed: 1,000/s, same as always. Error rate at the server: possibly zero, because the server successfully computed every answer — the failures happened in clients, which are somebody else's telemetry. This is the shape that produces the nobody's-layer-broke incident, where four teams each demonstrate their component is healthy and the user still can't log in.

flowchart LR
    T["Trigger<br/>failover / burst / slow dep"] -->|"one-time"| Q["Queue wait exceeds<br/>client timeout"]
    Q --> W["Work completed for<br/>departed clients<br/>(goodput → 0)"]
    W --> F["Clients see failure"]
    F --> R["Retries: L = A x 3"]
    R -->|"sustaining loop"| Q
    T -.->|"trigger removed —<br/>loop continues"| Q

The dotted edge is the article. Cut the solid loop and the system recovers; cut the trigger and nothing happens at all.

The standard reliability toolkit is the mechanism

Timeouts and retries are the first two things anyone learns about building resilient distributed systems, and in a metastable failure they are not the defense — they are the engine.

The reason is that retries only make sense against independent failures. If a request fails because it hit one bad node out of fifty, or a transient network blip, retrying is close to free and very likely to work: the second attempt lands somewhere else and the failure probability multiplies down. This is the case the textbooks describe, and it is real.

Saturation failures are the opposite. They are perfectly correlated: the reason your request failed is that the system is full, and the system being full is exactly the condition your retry makes worse. Retrying against a saturated dependency has a success probability that decreases with the number of clients doing it. You are not sampling a second time from a distribution; you are voting to make the distribution worse.

And a timeout followed by a retry has a specific, nasty property in this regime: the original work usually completed. The server hashed the password, built the token, wrote the audit record, and put the response on a closed socket. From the server's perspective it did the job. From the client's perspective nothing happened, so it asks again, and now the server does the same job a second time. Under sustained overload with a 2-second timeout and a 7-second queue, essentially 100% of your CPU is spent on work that has already been done or will never be read.

The uncomfortable summary: a timeout is a decision to stop waiting, not a decision to stop paying. Unless something on the server side acts on that timeout, the cost is still yours.

The defenses, and what each one costs

Five mechanisms actually address this. None is free.

Retry budgets

A retry budget is a global cap on the fraction of outgoing requests that may be retries — typically 10–20% — enforced at the client library or sidecar, across all callers of a dependency, as a rolling ratio. When the budget is exhausted, failures are returned to the caller immediately without a retry.

Why this beats per-client caps is arithmetic, and it's the most useful single idea in this article. A per-client cap of "3 attempts" bounds offered load at L ≤ 3A — a multiplicative bound that is only reached when everything is failing, which is to say it is loosest exactly when you need it tightest. A 10% budget bounds offered load at L ≤ 1.1A — an additive bound that holds regardless of the failure rate.

Run our numbers through it. Per-client cap: congested load 2,400/s against capacity 1,000/s, no escape. 10% budget: congested load 880/s, which is below capacity, so the queue drains, latency falls under the timeout, requests succeed, and the retry rate falls back to near zero on its own. The congested fixed point ceases to exist. That is a much stronger property than "recovers faster" — it means the system is no longer bistable, and no trigger can flip it.

There's a second, subtler virtue: a budget discriminates automatically between the two failure regimes described above. When failures are rare and independent, the budget is never binding and you get full retry benefit. When failures are correlated and widespread, the budget binds immediately and retries stop. You do not have to detect which situation you're in.

The cost: during a genuine partial failure — one bad node, a brief network fault — some requests that would have succeeded on a retry now fail. You have deliberately traded a small amount of single-request reliability for system-wide stability, and someone will file a bug about it. Budgets are also only as good as their scope: a budget per process, in a fleet of 500 processes, is 500 independent budgets, which is better than nothing but weaker than it looks.

Jittered backoff — and jitter matters more than backoff

Exponential backoff is the famous half. Jitter is the half that does the work.

Take 10,000 clients that all failed at t=0 because the service was briefly down. With clean exponential backoff — 1s, 2s, 4s, 8s — all 10,000 retry at t=1.0, together, inside the same few tens of milliseconds. If they land in a 100ms window, that's an instantaneous arrival rate of 100,000/s against a service that can do 1,000/s. Backoff has reduced the average retry rate and simultaneously made the peak sharper, by keeping the herd synchronized. It is a scheduling mechanism for thundering herds.

Now apply full jitter — sleep for random(0, cap) rather than cap. The same 10,000 clients spread across a 4-second window arrive at 2,500/s. Same average rate, 40× lower peak, and the peak is what determines whether you re-enter the congested state.

The rule is: backoff reduces average load; jitter reduces correlation, and correlation is what kills you. If you can only fix one, fix jitter. And note that this applies well beyond retries — synchronized token lifetimes and lockstep cache TTLs create the same herd on a schedule, which the art of load testing covers as the self-inflicted spike.

The cost: jitter increases worst-case latency for individual requests, and it makes behavior non-reproducible, which complicates testing. Both are cheap at the price.

Circuit breakers — which have their own metastable failure

A circuit breaker stops sending requests to a dependency that is failing, cutting the retry loop at the source. It works, and it is the mechanism that makes degradation fast rather than merely correct — a fallback that takes eight seconds per request to decide is not a fallback (what identity should do when the database is down makes this case in detail).

But breakers introduce a bistability of their own, and it is under-appreciated. Every breaker in your fleet trips at roughly the same moment, because they are all watching the same failing dependency. They all have the same reset timeout — say 30 seconds, from the same config file. So at t+30 they all half-open simultaneously, and the dependency, which has had 30 seconds of zero load and is now cold, receives 100% of offered traffic in one step. It fails. Every breaker trips again. At t+60, repeat.

You have built an oscillator: a system that is down 100% of the time in a sawtooth pattern, with a perfectly regular period that looks like a cron job and gets misdiagnosed as one. I have watched a team spend an hour hunting for the scheduled task causing 30-second outage cycles.

Two fixes, and you want both. Jitter the reset interval — the same argument as above, applied to breakers. And make half-open gradual rather than binary: admit 1% of traffic, then 5%, then 20%, ramping only while success rates hold. A breaker that admits a controlled fraction is a load shedder with a feedback loop, which is what you wanted all along.

Bounded queues, deadline propagation, and LIFO

The mechanics of bounded queues, deadline checks at dequeue, and LIFO ordering under overload are worked through in rate limiting is a capacity problem, and I won't repeat them. Two things belong here specifically because they are about the loop rather than about admission.

Deadlines must propagate across services, not just within one. Stamp the request with an absolute deadline at the edge, pass it on every hop, and have each service check it before starting expensive work and abandon immediately if it has passed. Without propagation, service D happily spends 90ms hashing for a request that service A gave up on two seconds ago — the retry-waste problem, reproduced internally at every layer. Propagate the remaining duration rather than a wall-clock timestamp unless your clocks are genuinely disciplined, because a 200ms clock skew silently turns into a 200ms budget error at every hop.

The shed path must be genuinely cheap. This is where implementations quietly fail. At 2,400/s of shed traffic, a rejection costing 50µs is 0.12 cores — free. A rejection that takes a Redis round trip is 2.4 cores — acceptable. A rejection that constructs an exception with a stack trace, logs it at ERROR, and emits a metric tagged with tenant and username is easily 2ms and unbounded cardinality, which is 5 cores and a monitoring outage on top of your identity outage. Load shedding only helps if the shed path cannot itself become the bottleneck. Measure it.

Recovery: why you have to drop work you already accepted

Now the part that determines how long the incident lasts.

"Just add capacity" usually fails during the event. It's the first thing everyone reaches for, and in a metastable state it is close to useless for three reasons. The offered load is 2,400/s against 1,000/s, so returning to health requires more than doubling the fleet, not adding 20%. New instances arrive with cold caches, so they serve at a fraction of nominal capacity while warming — but the load balancer gives them a full share of traffic immediately, meaning you have added instances that consume their share of requests and convert fewer of them into answers (identity platforms are mostly caching problems is the same observation from the other direction). And a cold instance's first act is often to hammer the database that just failed over. Ten nodes at 100/s plus five cold nodes at 40/s is 1,200/s against 2,400/s offered: still stuck, now with a bigger bill.

Capacity is how you avoid the cliff. It is not how you climb back up it.

What actually works is dropping accepted work. The queue holds 7,600 requests, every one of them past its deadline. Serving them costs 7,600 × 90ms = 684 core-seconds — nearly seven seconds of your entire fleet's total capacity — to produce exactly zero goodput. The correct action is to discard the queue, in full, immediately.

This feels indefensible in the moment. You are deliberately destroying thousands of real user requests. Say it out loud in an incident channel and someone will object. The answer is that those requests are already lost — their clients left minutes ago — and the only remaining question is whether you also spend seven core-seconds proving it.

The recovery sequence that works, in order:

  1. Flush the queues. Everything past deadline is dropped without processing. Free the capacity.
  2. Clamp admission below capacity — admit at 50–60% of C and reject the rest fast with 429 and Retry-After. You need admitted requests to complete inside the client timeout, because a completion that arrives late generates a retry and re-enters the loop. Serving 500/s successfully beats serving 1,000/s uselessly.
  3. Ramp the clamp gradually, watching queue wait rather than error rate. Queue wait is the leading signal; error rate lags it by however long your timeout is.
  4. Only then consider capacity, if the steady state genuinely needs it.

Step 2 is worth dwelling on. Returning a fast 429 is not just a rejection, it is backpressure that reaches into the clients — a good SDK honors Retry-After and stops hammering you. Many don't, and treat 429 as just another retryable error, which is why step 2 depends on the shed path being cheap. Design for the pessimistic case and be pleasantly surprised.

The deliberate-rejection argument

All of this rests on a claim worth making explicitly, because it is where engineering intuition and product intuition diverge hardest:

A system that cleanly rejects 20% of requests is worth far more than one that serves 100% of them at 40 seconds.

The second is not "degraded." Every client timeout in your estate is under 40 seconds, so a 40-second response is a failure that also consumed a full unit of capacity. The user experience is identical to a total outage; the difference is that the outage version at least leaves the CPU free. And the 40-second system is actively self-harming, because every one of those timeouts becomes a retry.

The rejecting system, meanwhile, is telling the truth. 80% of users log in normally. 20% get a fast, honest error with a retry hint, and most of them succeed on the next attempt seconds later. The system stays on the good side of the goodput curve, the retry loop never closes, and the incident ends when the trigger ends.

There is a trust dimension too, and it isn't decoration. A platform that fails fast and clearly is a platform people can build a fallback against. A platform that hangs for forty seconds poisons every caller's thread pool and turns one team's incident into everyone's incident — the blast-radius asymmetry that makes every identity bug a trust bug. Predictable failure is a feature, in the same family as fast, boring, and predictable success.

The counterpoint

An honest accounting of where this argument is weak.

Every mechanism here is a new failure mode. Circuit breakers cause outages. Retry budgets cause spurious failures. Load shedders reject traffic they shouldn't when a metric goes stale. I have seen more incidents caused by a misconfigured breaker than prevented by a correct one. These mechanisms earn their keep in systems where saturation is a real risk; in a comfortably over-provisioned service they are complexity with a negative return.

Sometimes you are just under-provisioned, and shedding is a way to avoid noticing. If your median tenant is being shed during normal business hours, you do not have a metastability problem, you have a capacity problem wearing a control loop as a disguise. The tell is whether sheds correlate with events or with growth. Rate limiting is a capacity problem makes this case at length and it applies with equal force here.

Not all shedding is equal, and identity makes that sharp. Shedding a login is shedding a person's entire ability to work — there is no partial degradation, no "browse in read-only mode." So prioritize within the shed: first attempts before retries, token refreshes (which keep large populations of already-working sessions alive, and are cheap) before new logins, and interactive traffic before batch. A shedder that treats all requests as fungible discards the cheapest, highest-value work alongside the rest.

LIFO is genuinely unfair. Under sustained overload it will starve some users indefinitely while others sail through, and "some users can never log in while the system reports 95% success" is a hard sentence to say to a customer. It's the right trade under sustained overload with deadlines, and the wrong default the rest of the time.

And retries are usually right. The entire argument above is about the correlated-failure regime. The overwhelming majority of the time, your failures are independent, retries convert them into successes, and a system that never retries is needlessly fragile. The goal is not fewer retries. It is retries that switch themselves off when they stop helping.

What to watch for, before it happens

Your existing alerts will not catch this. Throughput is flat and CPU is high, which is what a healthy system at peak looks like.

Retry ratio, as a first-class metric. The fraction of incoming requests that are retries. This is the single best leading indicator of metastability, and almost nobody has it, because it requires clients to tell you — an attempt-number header or a retry flag, propagated by your SDKs and honored by your gateway. That's a small protocol change with an enormous diagnostic payoff: it turns the sustaining effect from an inference into a number. A retry ratio drifting from 2% to 15% is the loop closing, and it moves before anything else does.

Queue wait, separated from service time, at p99. Total duration mixes the two and they have opposite fixes. Queue wait crossing your client timeout is the exact moment goodput becomes zero — alert on that threshold directly, by name.

Goodput, not throughput. Count responses that arrived before the client's deadline. This is the only metric that would have shown 14:12 for what it was, and it requires you to know the deadline, which is another argument for propagating it.

Hysteresis, tested deliberately. Drive the system into overload in a staging environment, remove the overload, and see whether it comes back without intervention. If it doesn't, you have a latent outage regardless of what your peak throughput number says. This is the highest-value experiment in chaos engineering for identity systems, and it is one of the few tests whose result is a straight yes or no.

Where this lands

The lesson of a metastable failure is not that retries are bad or that timeouts are wrong. It's that the properties that make a system resilient to small, independent failures are the same properties that make it fragile to large, correlated ones — and that the transition between those two regimes is a cliff, not a slope.

So build for the cliff. Bound the retry multiplier globally rather than per client, so the congested equilibrium does not exist. Jitter everything that could otherwise synchronize. Propagate deadlines and stop working the instant they pass. Make rejection cheap enough to survive a storm, and build the switch before you need it, because nobody builds a load shedder at 3am.

And know the shape well enough to recognize it in the first five minutes rather than the eighty-eighth. When the dashboards are green, the trigger is long gone, and the system is still down — stop looking for the cause. The cause is the thing your clients are doing right now, and it is the only thing you can still change.