Consistency modelsSession guarantees

100%

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.

Intermediate18 minUpdated 30 Sept 2026

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.

Notes
  • 9 Allowed by eventual consistency; breaks read your writes.
Timeline as a list

Save, then reload: 4 lanes, from 0 ms to 80 ms.

  1. 0–22 ms · Meera's app · save note (ok)
  2. 0 ms · Meera's app → API server, arriving 4 ms
  3. 4 ms · API server → Leader, arriving 10 ms
  4. 12 ms · Leader · commit 0/5A3F10
  5. 12–52 ms · Replica · window: replica 40 ms behind
  6. 12 ms · Leader → API server, arriving 18 ms
  7. 12 ms · Leader → Replica, arriving 52 ms: WAL stream
  8. 18 ms · API server → Meera's app, arriving 22 ms: 201
  9. 26–44 ms · Meera's app · reload, no note (stale) — Allowed by eventual consistency; breaks read your writes.
  10. 26 ms · Meera's app → API server, arriving 30 ms: GET
  11. 30 ms · API server → Replica, arriving 34 ms
  12. 36 ms · Replica → API server, arriving 40 ms: old recipe
  13. 40 ms · API server → Meera's app, arriving 44 ms
  14. 52 ms · Replica · applies 0/5A3F10 (ok)
Meera's save commits on the leader at 12 ms, but the replica applies it only at 52 ms. The reload at 26 ms reaches the replica inside that gap. The second tab, "With a version token", shows the same timing with a token: the API waits for the replica instead of answering from the past. Times are illustrative.

The four session guarantees (Terry et al., Bayou, 1994)

GuaranteePromiseAnomaly it prevents
Read your writesA read reflects every earlier write in the session.Meera's saved note vanishes on reload.
Monotonic readsA 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 writesThe 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 readsA 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.
Was this section helpful?

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?"

TechniqueHow it knowsGivesBreaks when
Leader for N secondsAfter 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 NLag exceeds N; the user switches device; every reread in the window hits the leader, needed or not
Sticky replicaThe 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 otherThe replica fails over to one that is further behind; hot users skew load
Version tokenThe 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 readsFailover loses the token's write; a cache in front ignores the token
Remote markerBefore 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 effortThe marker is evicted or lost before replication finishes

A read that carries a version token

A read that carries a version token, as an ordered list of steps:
A read that carries a version token13 steps between Tandoor app, API server, Leader database, Read replica. The steps are listed as text after the diagram.Read replicaLeader databaseAPI serverTandoor appwait, re-checking every 10 ms (budget 50 ms; if still behind, read the leader)PUT /recipes/812/notes1INSERT note2committed at LSN 0/5A3F103201 + X-Min-Version: 0/5A3F104GET /recipes/812 (X-Min-Version: 0/5A3F10)5applied ≥ 0/5A3F10?6no, at 0/5A3E807applied ≥ 0/5A3F10?8yes, at 0/5A3F109SELECT recipe 81210recipe with the note11200 + X-Min-Version: 0/5A3F10 (or newer)12
  1. Tandoor app → API server: PUT /recipes/812/notes
  2. API server → Leader database: INSERT note
  3. Leader database → API server (reply): committed at LSN 0/5A3F10
  4. API server → Tandoor app (reply): 201 + X-Min-Version: 0/5A3F10
  5. Tandoor app → API server: GET /recipes/812 (X-Min-Version: 0/5A3F10)
  6. API server → Read replica: applied ≥ 0/5A3F10?
  7. Read replica → API server (reply): no, at 0/5A3E80
  8. Note over API server: wait, re-checking every 10 ms (budget 50 ms; if still behind, read the leader)
  9. API server → Read replica: applied ≥ 0/5A3F10?
  10. Read replica → API server (reply): yes, at 0/5A3F10
  11. API server → Read replica: SELECT recipe 812
  12. Read replica → API server (reply): recipe with the note
  13. 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

What the API server does with one read. 6 states, 9 transitions. The table below lists them.
What the API server does with one readThe states of API server. 6 states, 9 transitions. The table below lists them.

no token / any replica

token present

replay ≥ token / read replica

replay ‹ token

re-check [under 50 ms]

replay ≥ token / read replica

50 ms passed

leader ≥ token / read leader

leader ‹ token / fail fast

Read arrives

Ask a replica its position

Replica behind, waiting

Read on the leader

Served, token raised

Token unsatisfiable

4 steps.

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.

Transitions of What the API server does with one read
From → ToEventGuardActionActor
Read arrives → Served, token raisedno tokenany replica
Read arrives → Ask a replica its positiontoken present
Ask a replica its position → Served, token raisedreplay ≥ tokenread replica
Ask a replica its position → Replica behind, waitingreplay < token
Replica behind, waiting → Replica behind, waitingre-checkunder 50 ms
Replica behind, waiting → Served, token raisedreplay ≥ tokenread replica
Replica behind, waiting → Read on the leader50 ms passed
Read on the leader → Served, token raisedleader ≥ tokenread leader
Read on the leader → Token unsatisfiableleader < tokenfail fastAPI
Read arrivesstart
Replica behind, waiting
Re-check every 10 ms
Served, token raisedend
Token unsatisfiableerror
Returned as an error the app can show
Was this section helpful?

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

Assumptions
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
Working
  1. Average writes2M × 10 ÷ 86,400 s~230/sfrom Daily active users and Edits per user per day
  2. Peak writes230 × 3~700/sfrom Average writes and Peak-to-average ratio
  3. 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.
  4. 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
  5. 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.
  6. Sticky replicano reads move0 extraBut it does not give read your writes here (writes go to the leader).
What it means
  • 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
Extra leader reads as replica lag growsThe token's cost climbs in steps as lag swallows each reread, and only matches the 5-second window's constant 2,100 reads/s once lag passes about 4 s.lag > 5 s05001k1.5k2k2.5k01 s2 s3 s4 s5 s6 sp99 lag 200 ms: ~700/sVersion tokenLeader for 5 sExtra leader reads (reads/s)Replica lag (ms)Extra leader reads as replica lag growsThe token's cost climbs in steps as lag swallows each reread, and only matches the 5-second window's constant 2,100 reads/s once lag passes about 4 s.lag >…05001k1.5k2k2.5k01 s2 s3 s4 s5 s6 sp99 lag 200 ms: ~700/sVersion tokenLeader for 5…Extra leader reads (reads/s)Replica lag (ms)
Illustrative: 700 writes/s at peak, rereads at about 100 ms, 2 s and 4 s after each write, and a 50 ms wait before falling back. Beyond 5 s of lag the time window stops protecting anyone (shaded).
Data
Replica lag (ms)Version token (reads/s)Leader for 5 s (reads/s)
002,100
150700no value
2,0501,400no value
4,0502,100no value
6,0002,1002,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

MongoDB causally consistent sessions
Each response carries operationTime and clusterTime; the driver sends afterClusterTime on the session's next read, and the node waits until it has caught up. The manual's table is the key detail: all four guarantees hold, even through failover, only with read concern "majority" and write concern "majority". With "local" reads and w: 1 writes, a failover can break every one of them.
Azure Cosmos DB, session level
After each write the SDK caches a session token bound to that physical partition. A replica that has not reached the token makes the client retry another replica, then another region. A client that never wrote to a partition holds no token for it, so its reads there behave as eventual. Microsoft calls session the most widely used level.
Memcache at Facebook
A web server in a replica region that updates key k first sets a remote marker for k in the regional cache, then writes to the master region, then deletes k locally. A later miss that finds the marker reads from the master region instead of the lagging local database: slower on that miss, far less likely stale. The paper calls the overall model best-effort eventual consistency.
Was this section helpful?

Trade-offs.

Tandoor's choices, with the options it passed over, then where each one breaks.

01
How Tandoor gives read your writes
Chosen:Version token (WAL position)
  • 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
Downside we accept:
  • 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
Ruled out:Leader for N seconds after a write

Over-routes: 2,100 leader reads/s instead of at most ~700; Still stale when lag exceeds N; Gives no monotonic reads

Ruled out:Sticky replica

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

02
Where the token lives
Chosen:In the user's server-side login session, raised with max()
  • Pro:Follows the user from phone to laptop
  • Pro:The API already loads the login session on every request
Downside we accept:
  • Con:A write to the session store on every response that raises the token
  • Con:Two devices racing need a max(), never an overwrite
Ruled out:Only in the app, as a response header it echoes back

Per device; the laptop does not know what the phone wrote; Lost when the app is reinstalled

Where session guarantees break

FailureImpactDetectionMitigationMeanwhile
Meera saves on her phone and opens the recipe on her laptop1Tandoor appThe laptop has no token, so read your writes does not hold across devicesSupport 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 databaseThe write is gone and no node can satisfy the token; LSNs on the new timeline can even pass it by coincidenceToken check fails on the leader; timeline change in the replication monitorCommit 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 hangThe app says the last change may not have saved
A recipe cache sits in front of the replicas5Recipe cacheA cache hit returns the old recipe no matter what token the read carriesStale reads with a satisfied token in tracesDelete 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 replicaEvery token read waits 50 ms and then lands on the leaderLag and fallback-rate alertsKeep the wait bounded; shed the lagging replica from the pool; cap fallback rate and serve a clear "still saving" state past itSlower reloads for recent writers only
A network partition cuts the API servers off from the leader2API serverReads that need the leader cannot be served; read your writes is only sticky availableLeader health checksClients that stay on a replica that already has their writes keep the guarantee; others get an error or wait

Before shipping read your writes

Every write response returns the new token.
Every read sends the token, and every read response returns the position it was served at.
The token is a map per shard or partition, and carries the leader's timeline or term.
The wait on a replica is bounded (tens of ms) and then falls back to the leader.
Caches in front of replicas are invalidated on write or checked against the token.
Replica lag and the fallback rate are alerted on, before the leader's CPU is.
An unsatisfiable token is an error the app can show, never a silent hang.
Was this section helpful?
Next in Core
Quorum reads and writes
Read next