Session guarantees.
Most complaints about eventual consistency come from one person contradicting themselves: I saved it, and it's gone. Session guarantees fix exactly that case. They don't make replicas agree faster; they make each client avoid the replicas that haven't caught up with it, usually by carrying a small version token.
Builds on Strong and eventual consistency.
The idea.
Eventual consistency lets any reader see an old value for a while. Session guarantees carve out one exception that users notice most, their own actions.
The note that vanished
Tandoor writes every recipe to one leader database in Mumbai and serves most reads from read replicas that copy the leader's log a little later: a few in Mumbai to take read load off the leader, and one in each other region. Meera, in Mumbai, adds a note to her dal recipe, "use kasuri methi", and taps Save. The leader commits it. The app then reloads the recipe, the load balancer sends that read to a replica running 40 ms behind, and the note is not there. Meera assumes the save failed and adds the note again. A few seconds later she has two identical notes.
Nothing was lost, and the database did exactly what "eventually consistent" promises: every replica got the note soon enough. The defect is that the system contradicted the one person who knew better. A stranger in another city seeing the old recipe for 40 ms would never notice; the author always does.
A session, in this topic, is the sequence of reads and writes one client makes, not a login cookie. The guarantees below are about that sequence.
Save, then reload
Scenario 1 of 2: As described.
- 9 Allowed by eventual consistency; breaks read your writes.
Timeline as a list
Save, then reload: 4 lanes, from 0 ms to 80 ms.
- 0–22 ms · Meera's app · save note (ok)
- 0 ms · Meera's app → API server, arriving 4 ms
- 4 ms · API server → Leader, arriving 10 ms
- 12 ms · Leader · commit 0/5A3F10
- 12–52 ms · Replica · window: replica 40 ms behind
- 12 ms · Leader → API server, arriving 18 ms
- 12 ms · Leader → Replica, arriving 52 ms: WAL stream
- 18 ms · API server → Meera's app, arriving 22 ms: 201
- 26–44 ms · Meera's app · reload, no note (stale) — Allowed by eventual consistency; breaks read your writes.
- 26 ms · Meera's app → API server, arriving 30 ms: GET
- 30 ms · API server → Replica, arriving 34 ms
- 36 ms · Replica → API server, arriving 40 ms: old recipe
- 40 ms · API server → Meera's app, arriving 44 ms
- 52 ms · Replica · applies 0/5A3F10 (ok)
The four session guarantees (Terry et al., Bayou, 1994)
| Guarantee | Promise | Anomaly it prevents |
|---|---|---|
| Read your writes | A read reflects every earlier write in the session. | Meera's saved note vanishes on reload. |
| Monotonic reads | A later read never shows an older state than an earlier read. | A recipe shows 41 likes, then 38 on refresh because the second read hit a replica further behind. |
| Monotonic writes | The session's writes apply everywhere in the order they were issued. | If Tandoor added write regions: renames to "Dal v1", then "Dal v2", land in different regions, and one region applies them in reverse and ends on v1. |
| Writes follow reads | A write is ordered after every write the session had read before making it. | If Tandoor added write regions: a reply "yes, 2 tsp" shows up in one region before the question it answers. |
How it works.
Every technique answers one question per read: has this replica seen everything this session depends on? They differ in how they know.
Four ways to answer "is this replica new enough?"
| Technique | How it knows | Gives | Breaks when |
|---|---|---|---|
| Leader for N seconds | After a write, a cookie or last-write timestamp sends the session's reads to the leader for, say, 5 s. | Read your writes, while lag stays under N | Lag exceeds N; the user switches device; every reread in the window hits the leader, needed or not |
| Sticky replica | The session is pinned to one replica (hash of user ID). | Monotonic reads. Not read your writes: writes still go to the leader, and the pinned replica lags like any other | The replica fails over to one that is further behind; hot users skew load |
| Version token | The client carries the highest log position it wrote or read; a replica serves only if it has applied that far. | Read your writes, plus monotonic reads if the token also rises on reads | Failover loses the token's write; a cache in front ignores the token |
| Remote marker | Before a write, set a per-key marker in the cache; a later miss that finds it reads from the primary region (Facebook memcache). | Read your writes for that key, best effort | The marker is evicted or lost before replication finishes |
A read that carries a version token
- Tandoor app → API server: PUT /recipes/812/notes
- API server → Leader database: INSERT note
- Leader database → API server (reply): committed at LSN 0/5A3F10
- API server → Tandoor app (reply): 201 + X-Min-Version: 0/5A3F10
- Tandoor app → API server: GET /recipes/812 (X-Min-Version: 0/5A3F10)
- API server → Read replica: applied ≥ 0/5A3F10?
- Read replica → API server (reply): no, at 0/5A3E80
- Note over API server: wait, re-checking every 10 ms (budget 50 ms; if still behind, read the leader)
- API server → Read replica: applied ≥ 0/5A3F10?
- Read replica → API server (reply): yes, at 0/5A3F10
- API server → Read replica: SELECT recipe 812
- Read replica → API server (reply): recipe with the note
- API server → Tandoor app (reply): 200 + X-Min-Version: 0/5A3F10 (or newer)
The token in PostgreSQL
-- Run on the leader, after the INSERT has committed.
-- The WAL insert position is at or past our commit record even
-- with synchronous_commit = off (pg_current_wal_lsn() is the write
-- position, which can trail it then), so this is a safe token.
SELECT pg_current_wal_insert_lsn()::text AS token; -- e.g. '0/5A3F10'// minLsn comes from the X-Min-Version header (may be absent).
async function readRecipe(id: number, minLsn?: string) {
// One connection, so the check and the read hit the same replica.
const conn = await replicaPool.connect();
try {
if (!minLsn) return await replicaRead(conn, id); // nothing to wait for
const deadline = Date.now() + 50; // wait budget
do {
const { rows } = await conn.query(
'SELECT pg_last_wal_replay_lsn() >= $1::pg_lsn AS ok', [minLsn]);
if (rows[0].ok) return await replicaRead(conn, id); // caught up
await sleep(10);
} while (Date.now() < deadline);
} finally {
conn.release();
}
return leaderRead(id); // lag outlasted the wait
}
// replicaRead / leaderRead also return the position they served at
// (pg_last_wal_replay_lsn() or pg_current_wal_insert_lsn()), which becomes
// the new token: that is what makes reads monotonic too.One token, more guarantees
MySQL has the same idea built in: WAIT_FOR_EXECUTED_GTID_SET(gtid_set, timeout) blocks on a replica until it has applied the given transactions or the timeout passes, so the token is a GTID set instead of an LSN.
- Monotonic reads: the client raises its token to the position of every response it gets, not only its writes. A later read can then never land on a replica older than one it already saw.
- Monotonic writes: with one leader and a client that waits for each write's acknowledgement, the leader's log already orders them. With several leaders (multi-region writes), a leader must first apply the session's earlier writes before accepting the next one, so the token travels with writes too.
- Writes follow reads: the same rule with the read set. With one leader it is free, because every replica holds a prefix of the leader's log. Across leaders or shards, the write carries the token and the receiving node orders it after everything the token covers.
Bayou tracked exactly this: a read set and a write set of write IDs per session, compacted into version vectors, and a server may serve the session only if its own vector covers them. A single-leader LSN is that vector with one entry. With shards, the token becomes a small map, one position per shard the session touched.
What the API server does with one read
States of2API server
Every read ends served, except when even the leader is behind the token, which only happens after a failover lost writes. That case fails fast instead of hanging.
| From → To | Event | Guard | Action | Actor |
|---|---|---|---|---|
| Read arrives → Served, token raised | no token | any replica | ||
| Read arrives → Ask a replica its position | token present | |||
| Ask a replica its position → Served, token raised | replay ≥ token | read replica | ||
| Ask a replica its position → Replica behind, waiting | replay < token | |||
| Replica behind, waiting → Replica behind, waiting | re-check | under 50 ms | ||
| Replica behind, waiting → Served, token raised | replay ≥ token | read replica | ||
| Replica behind, waiting → Read on the leader | 50 ms passed | |||
| Read on the leader → Served, token raised | leader ≥ token | read leader | ||
| Read on the leader → Token unsatisfiable | leader < token | fail fast | API |
- Read arrivesstart
- Replica behind, waiting
- Re-check every 10 ms
- Served, token raisedend
- Token unsatisfiableerror
- Returned as an error the app can show
In practice.
What each technique costs the leader at Tandoor's scale, and how three real systems package the same idea.
What each technique costs the leader
- Daily active users
- 2Millustrative
- Edits per user per day
- 10notes, renames, likes
- Peak-to-average ratio
- 3×illustrative
- Reads by the writer within 5 s of a write
- 3reload at ~100 ms, then ~2 s and ~4 s (illustrative)
- Replica lag at p99
- 200 msillustrative; the same figure as the strong-vs-eventual estimate
- Average writes2M × 10 ÷ 86,400 s~230/sfrom Daily active users and Edits per user per day
- Peak writes230 × 3~700/sfrom Average writes and Peak-to-average ratio
- Leader for 5 s after a write700 × 3 rereads2,100 extra leader reads/sfrom Peak writes and Reads by the writer within 5 s of a write · Every reread goes to the leader, whether the replica was behind or not.
- Rereads that find the replica behindonly the ~100 ms reload is inside a 200 ms lag (+50 ms wait still short)1 of 3from Reads by the writer within 5 s of a write and Replica lag at p99
- Version token, worst case at p99 lag700 × 1~700 extra leader reads/sfrom Peak writes and Rereads that find the replica behind · At median lag (tens of ms) the 50 ms wait absorbs almost all of these, so the real figure is far lower.
- Sticky replicano reads move0 extraBut it does not give read your writes here (writes go to the leader).
- The time window over-routes by about 3×; the token sends only the reads that truly need the leader.
- Leader load under tokens tracks actual lag, so it rises exactly when replicas are in trouble. Alert on lag before the leader feels it.
Extra leader reads as replica lag grows
- Version token
- Leader for 5 s
Data
| Replica lag (ms) | Version token (reads/s) | Leader for 5 s (reads/s) |
|---|---|---|
| 0 | 0 | 2,100 |
| 150 | 700 | no value |
| 2,050 | 1,400 | no value |
| 4,050 | 2,100 | no value |
| 6,000 | 2,100 | 2,100 |
- lag > 5 s: Replica lag (ms) from 5,000 to 6,000
- At 200: p99 lag 200 ms: ~700/s
The same idea in three real systems
Trade-offs.
Tandoor's choices, with the options it passed over, then where each one breaks.
- Pro:Precise; only reads that land on a lagging replica move to the leader
- Pro:Leader load tracks real lag, about 3× less than a time window here
- Pro:Works across any number of stateless API servers
- Pro:Also gives monotonic reads if raised on every response
- Con:Every response must return the token and every request must send it
- Con:One position per shard, so the token grows into a small map
- Con:Meaningless across a failover unless it carries the leader's epoch or timeline
Over-routes: 2,100 leader reads/s instead of at most ~700; Still stale when lag exceeds N; Gives no monotonic reads
Does not give read your writes when writes go to a leader; Failover moves users to a replica that may be further behind; Uneven load from heavy users
- Pro:Follows the user from phone to laptop
- Pro:The API already loads the login session on every request
- Con:A write to the session store on every response that raises the token
- Con:Two devices racing need a max(), never an overwrite
Per device; the laptop does not know what the phone wrote; Lost when the app is reinstalled
Where session guarantees break
| Failure | Impact | Detection | Mitigation | Meanwhile |
|---|---|---|---|---|
| Meera saves on her phone and opens the recipe on her laptop1Tandoor app | The laptop has no token, so read your writes does not hold across devices | Support tickets about edits that "didn't save" | Keep the token per user on the server, not only per device | |
| Failover promotes an asynchronous replica that never received the token's write3Leader database | The write is gone and no node can satisfy the token; LSNs on the new timeline can even pass it by coincidence | Token check fails on the leader; timeline change in the replication monitor | Commit to a majority (or a synchronous standby); put the timeline or term in the token and treat an unsatisfiable token as a hard error, not a hang | The app says the last change may not have saved |
| A recipe cache sits in front of the replicas5Recipe cache | A cache hit returns the old recipe no matter what token the read carries | Stale reads with a satisfied token in traces | Delete the entry on write, and store the position each entry was filled at; serve a hit only if it is at or past the token | |
| Replica lag jumps to seconds (a long vacuum, a bulk import)4Read replica | Every token read waits 50 ms and then lands on the leader | Lag and fallback-rate alerts | Keep the wait bounded; shed the lagging replica from the pool; cap fallback rate and serve a clear "still saving" state past it | Slower reloads for recent writers only |
| A network partition cuts the API servers off from the leader2API server | Reads that need the leader cannot be served; read your writes is only sticky available | Leader health checks | Clients that stay on a replica that already has their writes keep the guarantee; others get an error or wait |