Coordination APIs: Locks, Leases, Watches, and Compare-And-Swap
LESSON
Coordination APIs: Locks, Leases, Watches, and Compare-And-Swap
By the end of this lesson, you will be able to...
Design a controller-election API that exposes the preconditions and evidence its clients need.
Trace a safe recovery from a missed watch event through snapshot, revision, and conditional update.
Explain why a lease grants temporary coordination authority but fencing protects the resource that can be harmed.
Idea in one sentence: A coordination API is safe when it turns hidden consensus evidence into explicit revisions, conditions, expiry rules, recovery paths, and tokens that clients and protected resources can enforce.
Core Insight
Atlas offers a raw key-value store backed by consensus. A deployment team needs one active controller for payments. The quickest interface looks harmless:
write /leaders/payments = controller-a
read /leaders/payments
It leaves the crucial questions to every client. Did controller-a win a race or overwrite another owner? What happens when it crashes? How does a watcher know it missed an event? What tells the workload database that a resumed former controller is stale?
Consensus can protect the store's history, but a raw get and put API does not automatically give application code a safe authority contract. A coordination API must expose the evidence that makes the safe path usable.
For Atlas, that means four connected pieces:
- compare-and-swap (CAS) turns “write this” into “write this only if the state I inspected is still current”;
- a lease or session gives ownership an expiry and loss rule;
- a watch with revisions lets a client recover a current view after a disconnect;
- a monotonically increasing fencing token lets the workload database reject stale controller actions.
The goal is not to make clients incapable of mistakes. The goal is to make stale assumptions and recovery boundaries visible instead of letting every team rebuild them from a convenient put.
The Promise We Need to Keep
The deployment platform promises two things:
- At most one controller has current coordination authority for
paymentsat a time. - A previous controller that pauses, loses its session, or misses watch events cannot successfully overwrite work accepted from a newer controller.
Those promises involve different actors.
| Actor | Needs to know | API evidence it needs |
|---|---|---|
| Controller | Whether it won ownership and when that claim becomes stale | Conditional result, revision, lease state, token |
| Standby controller | When ownership changed or disappeared | Snapshot revision, ordered watch events, resync rule |
| Workload database | Which controller request may change assignments | Monotonic fencing token on each protected write |
| Operator | Whether recovery is working or a controller is guessing | Watch cancellation, compaction, session-loss, and token-rejection signals |
The API's job is to make these evidence paths explicit. It should not make a controller infer authority from an old cache or a successful TCP connection.
The Naive Design
The first design gives every controller these operations:
GET(key)
PUT(key, value)
WATCH(key)
Two controllers can both read “no leader,” then both write themselves as leader. A process that crashes may leave its key forever. A watch can disconnect after a change, leaving a controller with an old local cache. Even if the store serializes writes, the application cannot tell whether its own decision was based on the same state that the store later changed.
Adding a lock name alone does not repair this design. A lock that says only “acquired” still leaves open the same questions: which revision records it, when it expires, how a client detects loss, and whether an external database refuses an old owner.
The pressure becomes visible after a pause:
controller A writes itself as owner
A pauses
controller B writes itself as owner
A resumes and writes workload assignments
The store may correctly order both writes. The workload database still receives a dangerous action from A unless the authority evidence reaches that database.
A Better Boundary: Guarded Ownership Records
Atlas defines an ownership record with fields that clients can use:
/leaders/payments = {
owner: controller-a,
lease_id: 17,
token: 41
}
revision: 900
The controller acquires it through a transaction, not an unconditional write:
if key does not exist at current committed revision:
create ownership record attached to lease 17
allocate token 41
return revision 900 and token 41
else:
fail with the current revision or conflict result
This is compare-and-swap in its useful form. The precise comparison may be “key absent,” “modification revision equals R,” or another documented condition. The meaning is stable: the store decides whether the precondition and the write occur together in the committed order.
A failed CAS is not a mysterious transient error. It says the controller's local story is stale. The controller must read or resync before deciding again. Retrying the same unconditional write would turn visible evidence into a hidden race.
The token is separate from the revision. A revision records where this particular metadata change sits in store history. A fencing token represents a monotonically increasing authority generation that the protected database can compare. An API may derive or allocate these differently, but it must state their ordering and retention rules.
Lease Loss Is a State Change, Not a Hint
The ownership record is attached to a lease. The controller renews the lease while it is healthy. If renewal stops long enough under the service's documented session semantics, the coordination service expires the lease and removes or invalidates the attached ownership record.
The client rule should be blunt:
if lease renewal fails, session closes, or connectivity is uncertain:
stop protected work
treat ownership as lost until a new acquisition succeeds
This rule can feel conservative. It is necessary because a client cannot distinguish every pause from every partition using local time alone. A current lease record transfers coordination authority, but it does not stop a former owner from continuing to execute in memory. The fencing token from that record is what lets the downstream system reject the former owner later.
The trade-off is liveness versus safety at the client boundary. Fast expiry releases work sooner after a failure but is more sensitive to pauses and renewal delays. Slow expiry reduces false loss but delays takeover. The API should expose the expiry and loss semantics rather than implying that “lock held” is a timeless fact.
Watches Need a Recovery Contract
A watch is not a guarantee that a client has seen every change forever. It is a stream beginning at a known point in committed history.
Atlas gives a controller this recovery rule:
1. Read a snapshot and remember its revision R.
2. Build or replace the local cache from that snapshot.
3. Watch changes after R.
4. If the watch disconnects, resumes unsuccessfully, or reports compaction,
discard the assumed gap and repeat from a new snapshot.
The R boundary prevents the client from mixing an old snapshot with an arbitrary stream. A watch event later than R updates a known base. If the service has compacted the missing event history, it cannot honestly replay the gap. The correct response is a new snapshot, not a guessed cursor or an eternal retry.
This is a clear example of a useful API failure: a compaction response tells the controller exactly when its cache is no longer a defensible reconstruction. The safe API gives that failure a name and a recovery path.
Worked Trace: One Election, One Lost Watch
The following revisions, lease IDs, and tokens are illustrative.
| Step | Store or client state | API result | Why it matters |
|---|---|---|---|
| 1 | A creates lease 17 |
Lease is live | Ownership will have expiry semantics |
| 2 | A CASes an absent /leaders/payments key |
Store creates owner A, token 41, revision 900 |
Only one conditional acquisition can win |
| 3 | B reads a snapshot at revision 900 |
B watches from after 900 |
B has a recoverable baseline |
| 4 | A pauses and stops renewal |
Lease 17 expires; the owner record changes at revision 904 |
Coordination authority moves or becomes available |
| 5 | B's watch was disconnected and history through 904 is compacted |
API reports that it cannot resume | B must not trust its old cache |
| 6 | B reads a fresh snapshot at revision 910 |
It sees no owner, then acquires lease 22 with token 42 |
Takeover uses current state and a new token |
| 7 | B writes an assignment with token 42 |
Database accepts it and records high-water token 42 |
The protected resource knows newer authority |
| 8 | A resumes and sends a delayed write with token 41 |
Database rejects 41 because 41 < 42 |
A stale process cannot win by waking up |
Notice the distinct safeguards. CAS prevents two acquisitions from both succeeding against the same precondition. The lease makes a crashed owner's record expire. Snapshot plus watch lets B recover its view. Fencing protects the action outside the coordination store. Removing any one leaves a different hole.
Operational Consequences
An API that exposes these primitives has costs. Conditional operations add retries when contention is real. Lease expiry can delay takeover. Watch streams require retention, compaction, and client resync. Tokens require every protected write path to perform a comparison. The alternative is not free simplicity; it is pushing the same hard decisions into every client library.
Make the operational signals part of the contract:
- CAS conflict rate shows competing or stale decisions;
- lease-renewal failures and session closes show controllers that must stop work;
- watch cancellations, compacted revisions, and resync duration show cache-recovery pressure;
- rejected fencing tokens show delayed or stale actors reaching protected resources;
- unbounded lease or deduplication records show retention rules that are not being enforced.
These signals help an operator distinguish normal contention from an API client that is replaying unsafe local assumptions.
Design Review
Before exposing a coordination primitive, ask:
| Question | A safe answer includes |
|---|---|
| What operation creates authority? | A linearizable conditional write and its success evidence |
| How does authority end? | Lease/session loss, expiry behavior, and client stop-work rule |
| How does a client learn it is stale? | CAS conflict, revision change, watch failure, or session loss |
| How does a client recover a missed stream? | Snapshot at R, watch after R, and resync after compaction |
| How is a stale actor stopped at the resource? | Monotonic token checked on every protected action |
This is a situated API preference. A simple feature flag may need only CAS and revision. A scheduler that mutates an outside database needs the full chain because stale external work is harmful.
Practice: Review an Incomplete Lock API
An API offers lock(name) and returns true or false. It provides no lease ID, expiry event, revision, resume point, or fencing token. A billing worker uses true to charge a customer.
Review the API.
A good answer should mention:
- a conditional acquisition result tied to a revision or explicit precondition;
- session or lease loss and the rule that the worker stops on uncertainty;
- a snapshot-and-watch recovery path for standby workers;
- a monotonic token that the billing system checks before accepting a charge-related side effect;
- an idempotent operation identity, because fencing alone does not turn a lost reply into exactly-once billing.
Connections
The previous lesson established why a current authority read and fencing belong to different boundaries. This lesson turns that evidence into client-facing operations and recovery rules.
The next lesson asks whether the consensus cluster can serve these contracts within a usable latency, disk, network, and failure-domain envelope.
Resources
- [DOC] etcd API — Focus: Transactions, leases, revisions, watches, compaction, and explicit consistency options.
- [DOC] ZooKeeper Programmer's Guide — Focus: Sessions, ephemeral nodes, watches, and client recovery behavior.
- [PAPER] The Chubby Lock Service for Loosely-Coupled Distributed Systems — Focus: Why coordination APIs need explicit client-visible semantics beyond a raw consensus log.
Key Takeaways
- CAS turns a local assumption into a condition the coordination service evaluates atomically in its committed order.
- A lease gives ownership an expiry rule, but the client must stop work when it loses confidence in that lease.
- Watches are safe only with a snapshot revision and an honest resync path after a missed or compacted stream.
- Fencing tokens must reach the protected resource, which rejects stale actions after newer authority succeeds.
- A coordination API earns its extra surface area when it makes these safety and recovery rules easier to use than to bypass.