Reject stale workers and recover uncertain external effectsLESSON 16.02 · 2 OF 9 IN CHAPTER
PART E / Live migrations
Step 216 of 252
LESSON 16.02 · 2 OF 9 IN CHAPTERHands-on

Reject stale workers and recover uncertain external effects

Application and assignment

A background worker fetches a bookmark’s display title. Worker A pauses long enough for its claim on the job to expire. Worker B takes over and saves a newer result. A resumes and tries to overwrite it. The expired claim did not stop A’s code from running.

Use the local reference to identify the authority checked at the final write. Then extend the reasoning to a provider that may complete a billable operation before its response is lost. A local result constraint and a provider’s idempotency record protect different effects.

Contract and starting evidence

Constructed candidate brief: “Worker A fetches an old page title and pauses. Its five-second lease expires; B takes the job and stores the new title. A resumes. Keep B's result. Next, replace the fetch with a billable provider call whose successful response can disappear. What can the service honestly promise?”

Prerequisites: the atomic-result queue lab and transaction concepts. This is a separate extension: the existing worker's entire effect remains one conditional result item. Do not insert the provider into that worker and claim its guarantee automatically grew. Attempt first; reference.py (download file, source below) and assessor.md contain the proposed boundaries and review criteria.

Read the supplied code · reference.py
reference.py · reference.py
"""Local protocol boundaries: real SQLite atomic intent; simulated remote failures."""
from dataclasses import dataclass
import hashlib
import json
import sqlite3


class Conflict(Exception):
    pass


class HistoryExpired(Exception):
    pass


def digest(payload):
    return hashlib.sha256(json.dumps(payload, sort_keys=True, separators=(",", ":")).encode()).hexdigest()


class FencedResult:
    """One atomic storage boundary; clock represents authoritative lease time."""
    def __init__(self):
        self.epoch = 0
        self.until = 0
        self.result = None

    def claim(self, now, ttl=5):
        if now < self.until:
            raise Conflict("lease held")
        self.epoch += 1
        self.until = now + ttl
        return self.epoch

    def complete(self, epoch, now, result):
        if epoch != self.epoch or now >= self.until:
            raise Conflict("stale or expired owner")
        self.result = result


class Provider:
    def __init__(self):
        self.receipts = {}
        self.effects = 0
        self.available = True
        self.lose_next_response = False

    def send(self, operation_id, payload):
        if not self.available:
            raise TimeoutError("provider unavailable")
        identity = digest(payload)
        previous = self.receipts.get(operation_id)
        if previous and previous[0] != identity:
            raise Conflict("same ID, changed intent")
        if previous is None:
            self.effects += 1
            self.receipts[operation_id] = (identity, f"receipt-{self.effects}")
        if self.lose_next_response:
            self.lose_next_response = False
            raise TimeoutError("effect happened, response lost")
        return self.receipts[operation_id][1]


class IntentStore:
    """Business row + outbox are committed together in a real local transaction."""
    def __init__(self):
        self.db = sqlite3.connect(":memory:")
        self.db.executescript("""
            CREATE TABLE intent(id TEXT PRIMARY KEY, digest TEXT NOT NULL, payload TEXT NOT NULL);
            CREATE TABLE outbox(id TEXT PRIMARY KEY REFERENCES intent(id), state TEXT NOT NULL, receipt TEXT);
        """)

    def close(self):
        self.db.close()

    def accept(self, operation_id, payload, crash_before_commit=False):
        with self.db:
            old = self.db.execute("SELECT digest FROM intent WHERE id=?", (operation_id,)).fetchone()
            if old:
                if old[0] != digest(payload):
                    raise Conflict("quarantine changed intent")
                return "duplicate"
            self.db.execute("INSERT INTO intent VALUES (?,?,?)", (operation_id, digest(payload), json.dumps(payload)))
            if crash_before_commit:
                raise RuntimeError("crash: transaction rolls back")
            self.db.execute("INSERT INTO outbox VALUES (?, 'pending', NULL)", (operation_id,))
        return "accepted"

    def dispatch(self, operation_id, provider, crash_after_send=False):
        payload, state = self.db.execute("SELECT payload,state FROM intent JOIN outbox USING(id) WHERE id=?", (operation_id,)).fetchone()
        if state == "done":
            return "done"
        try:
            receipt = provider.send(operation_id, json.loads(payload))
        except TimeoutError:
            with self.db:
                self.db.execute("UPDATE outbox SET state='uncertain' WHERE id=?", (operation_id,))
            return "uncertain"
        if crash_after_send:
            raise RuntimeError("crash before local receipt")
        with self.db:
            self.db.execute("UPDATE outbox SET state='done', receipt=? WHERE id=?", (receipt, operation_id))
        return "done"


@dataclass(frozen=True)
class Row:
    version: int
    value: str | None
    deleted: bool = False


def apply(rows, key, row):
    previous = rows.get(key)
    if previous is not None and previous.version == row.version and previous != row:
        raise Conflict("same version has different contents")
    if previous is None or row.version > previous.version:
        rows[key] = row


class ChangeSource:
    """Single-thread atomic source mutation + log model, not an actual CDC client."""
    def __init__(self):
        self.rows, self.log = {}, []
        self.sequence = self.retained_after = 0

    def write(self, key, value=None, deleted=False):
        version = self.rows[key].version + 1 if key in self.rows else 1
        row = Row(version, value, deleted)
        self.sequence += 1
        self.rows[key] = row
        self.log.append((self.sequence, key, row))
        return row

    def snapshot(self):
        # In a real source, snapshot and position must describe the SAME database state.
        return dict(self.rows), self.sequence

    def events_after(self, checkpoint):
        if checkpoint < self.retained_after:
            raise HistoryExpired("resnapshot required before continuing")
        return [event for event in self.log if event[0] > checkpoint]

    def expire_through(self, sequence):
        self.log = [event for event in self.log if event[0] > sequence]
        self.retained_after = sequence


class Replica:
    def __init__(self):
        self.rows, self.checkpoint = {}, 0

    def resume(self, source, crash_after_apply=False):
        for sequence, key, row in source.events_after(self.checkpoint):
            apply(self.rows, key, row)
            if crash_after_apply:
                raise RuntimeError("row durable, checkpoint not yet advanced")
            self.checkpoint = sequence

    def resnapshot(self, source):
        self.rows, self.checkpoint = source.snapshot()

    def differences(self, source):
        return {key for key in self.rows.keys() | source.rows.keys() if self.rows.get(key) != source.rows.get(key)}


class RoutingRegistry:
    def __init__(self):
        self.epoch, self.owner = 1, "old"

    def cutover(self, target, source):
        if target.checkpoint != source.sequence or target.differences(source):
            raise Conflict("cutover requires full reconciliation at write barrier")
        # A real implementation atomically fences old writes and changes authority.
        self.epoch += 1
        self.owner = "new"

    def check_write(self, owner, epoch):
        if owner != self.owner or epoch != self.epoch:
            raise Conflict("refresh routing; stale authority")


def rollback_ready(old_rows, new_rows):
    return old_rows == new_rows
Input/schedule Expected value Scope
A epoch 1 at t=0; B epoch 2 at t=6; B finishes t=7; A finishes t=8 Result new; A rejected Epoch and lease checked at protected write
Store business intent then crash before transaction commit 0 intent rows and 0 outbox rows Real local SQLite transaction
Provider commits; response lost; replay same ID/body 1 effect, eventually done Provider retains and enforces idempotency key
Replay same ID with changed body Conflict; original effect unchanged Changed intent is quarantined
Provider forgets key before a replay 2 effects possible Dedup retention is part of the guarantee

Baseline and counterexample

Diagram: Baseline and counterexample

A lease grants temporary authority; it cannot stop a suspended process from running later. “Writing twice is safe” assumes both writes contain the same intent and the same result. A fetched title can change between attempts. A crash after a fetch but before recording completion may force another fetch even when the database ultimately contains one result.

Add an authority check at the write

  1. Define safety: a stale owner never changes the current result. Define liveness: after lease expiry, a healthy worker may acquire a new epoch.
  2. Have the authoritative store issue a monotonically increasing fencing epoch when it grants ownership. A client-generated timestamp is not this authority.
  3. Require the result store to compare the submitted epoch with the current epoch atomically at the write. The sample also rejects an expired current lease. A separate lock service and result store need a real protocol ensuring the protected destination knows which epoch is valid.
  4. Preserve logical operation ID and payload digest. A matching replay is safe; a conflicting payload needs investigation, not overwriting the earlier intent.
Diagram: Add an authority check at the write

Follow-up: the effect lives at a provider

Draw the crash on each side of the remote call before opening the diagram. If you write “done” first, a crash can omit the effect. If you call first, a crash can repeat it. No local marker removes both gaps for an uncooperative provider.

Diagram: Follow-up: the effect lives at a provider

The fake provider commits one effect plus receipt and enforces matching identity. The real SQLite transaction verifies rollback of business intent without a matching outbox. Relay delivery can repeat. If the real provider offers neither idempotency nor a reliable receipt query, keep uncertain, stop automatic replay at a bounded attempt/deadline, and route reconciliation to an owner. The business must choose duplicate risk, omission risk, or manual resolution. Never turn “unknown” into “failed” merely because a socket timed out.

Run from the repository root (Python 3.11+, standard library only):

python -m unittest discover -s curriculum/04-scale-and-evolution/04-migrations/labs/recovery-migration -p 'test_*.py' -v

Tests inject stale-owner completion, pre-commit rollback, success with lost response, crash after receipt, changed intent, and expired provider dedup history. This model has no real distributed clock, lease service, or AWS account. Failures are deterministic call boundaries, not operating-system kill or network tests.

Senior follow-ups: what if provider key retention is one day but the DLQ is replayed a week later? What if all workers pause and lease time is measured on different machines? Expected reasoning names replay retention and authoritative time, respectively. Lead: specify the reconciliation owner, deadline, volume budget, and customer-facing uncertain state. Continue to live migration and rollback and rebalancing and region loss.