Bound cache misses across instances and preserve fresh reads
Application and assignment
A shared reading-list page is cached to avoid repeating its database query. When the entry expires, many readers can request the same missing value. Ana can also save a new bookmark while a read replica still contains the previous version of the list.
Use the supplied local models to distinguish shared loading within a process, fleet-wide admission, and read freshness. Predict origin calls and returned versions before reading the implementation. No Redis instance or replica cluster is started by this lab.
Contract and starting evidence
Constructed candidate brief: “Ana opens a group page just as its cache entry expires. Two hundred requests arrive together on ten API instances. Preserve the page's tenant boundary and keep refresh work inside a database budget. What changes if the cache is unavailable or Ana has just saved a newer version?”
Prerequisites: cache and replica concepts, Python 3.11+,
and understanding that an asyncio task may be cancelled while another task waits
on its result. Attempt the design and write the expected counters before opening
the reference (download file, source below) or assessor notes.
Read the supplied code · reference.py
"""Deterministic in-process boundary model, not a Redis or fleet-lock client."""
import asyncio
from collections import deque
class Unavailable(Exception):
pass
class Clock:
def __init__(self):
self.now = 0.0
def advance(self, seconds):
self.now += seconds
class OriginBudget:
"""Shared model of atomic admission: rolling one-second cap plus in-flight cap."""
def __init__(self, clock, per_second=100, concurrent=10):
self.clock, self.per_second, self.concurrent = clock, per_second, concurrent
self.starts = deque()
self.active = self.peak = self.total = 0
def enter(self):
while self.starts and self.starts[0] <= self.clock.now - 1:
self.starts.popleft()
if len(self.starts) >= self.per_second or self.active >= self.concurrent:
raise Unavailable("origin admission exhausted")
self.starts.append(self.clock.now)
self.active += 1
self.total += 1
self.peak = max(self.peak, self.active)
def leave(self):
self.active -= 1
class CacheAside:
def __init__(self, clock, budget, ttl=5):
self.clock, self.budget, self.ttl = clock, budget, ttl
self.values, self.flights = {}, {}
self.cache_available = True
async def get(self, key, loader):
cached = self.values.get(key) if self.cache_available else None
if cached and self.clock.now < cached[0]:
return cached[1]
task = self.flights.get(key)
if task is None:
task = asyncio.create_task(self._load(key, loader))
self.flights[key] = task
# Retrieve even an unobserved exception if every waiter leaves.
task.add_done_callback(lambda done: None if done.cancelled() else done.exception())
# A disconnected waiter must not cancel work shared with other waiters.
return await asyncio.shield(task)
async def _load(self, key, loader):
entered = False
try:
self.budget.enter()
entered = True
value = await loader(key)
if self.cache_available:
self.values[key] = (self.clock.now + self.ttl, value)
return value
finally:
if entered:
self.budget.leave()
self.flights.pop(key, None)
def sticky_read(primary, replica, now, pin_until):
return primary if now < pin_until else replica
def strict_read(primary, replica, minimum_version):
if replica is not None and replica >= minimum_version:
return replica
if primary is not None and primary >= minimum_version:
return primary
raise Unavailable("no source has the session watermark")
class LinkDelivery:
"""Explicitly separate authorization-decision age from object cache age."""
def __init__(self, clock, policy="immediate", auth_ttl=5):
if policy not in {"immediate", "bounded"}:
raise ValueError(policy)
self.clock, self.policy, self.auth_ttl = clock, policy, auth_ttl
self.tokens = {"A": ("tenant-a", "schedule-1"), "B": ("tenant-b", "schedule-1")}
self.objects = {("tenant-a", "schedule-1"): "Ana 09:00", ("tenant-b", "schedule-1"): "Ben 10:00"}
self.decisions, self.cache = {}, {}
self.auth_available = self.origin_available = True
self.origin_calls = 0
def get(self, token):
decision = self.decisions.get(token)
fresh = self.policy == "bounded" and decision and self.clock.now < decision[0]
if fresh:
identity = decision[1]
else:
if not self.auth_available:
return 503, None # Never extend an expired decision during outage.
identity = self.tokens.get(token)
if identity is None:
return 403, None
self.decisions[token] = (self.clock.now + self.auth_ttl, identity)
# Authorization precedes EVERY object-cache lookup; identity scopes bytes.
if identity in self.cache:
return 200, self.cache[identity]
if not self.origin_available:
return 503, None
self.origin_calls += 1
value = self.objects[identity]
self.cache[identity] = value
return 200, value
| Teaching input | Expected value | Excluded guarantee |
|---|---|---|
| 200 same-key overlapping misses, naive baseline | 200 origin calls | Jitter supplies no exclusion |
| Same misses, one local flight map | 1 successful load | Persistent/fleet coordination |
| Ten independent flight maps, one shared cache miss | Up to 10 loads; fixture observes 10 | One loader across processes |
| Ten waves of 100 unique misses at one fake second | 100 admitted, 900 unavailable, peak 10 | Production rate-limiter distribution |
| Primary v8, replica v7, pin ends 5s, read at 6s | Timer-only read returns v7 | Unbounded lag cannot fit a finite timer |
These are synthetic values, not AWS limits. The supplied local model does no
network access. Its shared OriginBudget represents an atomic fleet admission
boundary; copying that object into each process would multiply the budget.
Start with the failure
Walk a single key through the miss path. “Cache absent” is an observation, not a claim on the right to refresh. All 200 requests can observe the same absence.
Predict: if each reader adds a different TTL after loading, how many loads have already happened? All 200. Jitter helps different keys expire at different times. Early refresh is probabilistic; neither is mutual exclusion.
Change the boundary in order
- Name the invariant: one overlapping load per key within one process; at most 100 origin admissions in a rolling second and ten in flight in this fixture.
- Put a task in the flight map before waiting on it. Followers await that task. Keep keys tenant-scoped. Distinct keys cannot share values.
- Shield shared work from one waiter's cancellation. Remove the flight and free its budget on success, exception, and loader cancellation. If every waiter leaves, the sample deliberately lets the bounded load finish; production loaders also need deadlines and lifecycle shutdown.
- Apply admission before calling the origin. When exhausted, return the declared unavailable result. An HTTP adapter can map capacity overload to 429 or 503; safe stale data is an option only with a separate freshness and auth policy.
- Add instances and count again. If you require fleet exclusion, put arbitration in a shared service and reason about owner pause, expiry, and fencing. A shared rate gate bounds load even when exclusion is deliberately weaker.
Run the local boundary tests
From the repository root:
python -m unittest discover -s curriculum/04-scale-and-evolution/01-data-at-scale/labs/cache-consistency -p 'test_*.py' -v
Tests observe actual loader calls, all waiter results, cleanup, ten independent
instances, distinct keys, outage caps, and the lag counterexample. They use event
barriers and a fake clock; asyncio.sleep(0) yields scheduling without wall-clock
delays. Green tests verify this model, not a live Redis lock or global rate gate.
Follow-up: the cache disappears during a traffic burst
Before revealing this diagram, choose a user-visible result for 900 excess requests. “Same page, slower” is feasible only inside the origin's capacity.
Senior: distinguish a rate cap from a concurrency cap; neither implies the other when latency changes. Lead: allocate tenant shares and spare database capacity across services; describe what happens when the shared limiter itself fails. This model fails closed on exhausted admission and does not implement a distributed limiter. Validate real limiter availability/consistency separately.
Follow-up: the user's write is newer than the replica
Define a session watermark as the minimum version this session has acknowledged. For one row this model uses an integer. A real multi-row system needs a database replication position and a defined failover epoch; do not compare unrelated row versions as though they were a global log position.
Senior: show why a five-second sticky route returns v7 at t=6 when lag lasts ten seconds. Lead: explain the write acknowledged only in a failed region and choose between synchronous replication availability costs and an explicit RPO.
Continue to cached-link revocation and the lease/recovery lab.