Copy live rows while preserving new writes and deletions
Application and assignment
You are moving a tenant’s bookmarks to another store while members keep editing. A snapshot carries yesterday’s title, but a live update already contains today’s title. Applying the snapshot last must not erase the newer edit. A deletion needs the same protection against an old copy arriving later.
Follow the supplied protocol model through partial writes, replay, and reader cutover. Record which store currently owns writes and what durable evidence lets the other catch up. This fixture models the migration protocol. Connecting a production change-data-capture system is additional implementation work.
Contract and starting evidence
Constructed candidate brief: “Ana edits bookmark A while you move her tenant to a new store. A slow backfill still carries yesterday's title. Ben deletes a bookmark while the new-store write is failing. Preserve acknowledged updates and deletes, then show how reader cutover can safely roll back.”
Prerequisites: intent/outbox recovery, row versions, and the PostgreSQL lab. This is a runnable protocol model, not a production CDC connector. Attempt the schedule before reading the reference (download file, source below) and assessor guide.
Read the supplied code · 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
| Tiny schedule | Expected after repair | Counterexample |
|---|---|---|
| Snapshot A v1; live A v2; v2 copied before v1 | A remains v2 | Unconditional backfill overwrites new data |
| Delete D v2; then old snapshot D v1 arrives | D retains v2 tombstone | Physical absence alone permits resurrection |
| Source write succeeds; target unavailable | Source authoritative; durable change pending | Comparing rows cannot itself repair |
| Crash after target apply before checkpoint | Replay same version; unchanged result | Checkpoint-before-apply can skip data |
| New-store-only A v3; route readers back | Block until old store has v3 | Routing cannot undo missing data |
Baseline: two independent writes
Start with authority. During this phase, accepted writes commit to the old store and a durable change record together (transactional outbox or supported CDC). Readers may move independently; writers must not accidentally create two uncoordinated authorities.
Order the copy and the changes
- Expand the schema and document reader/writer compatibility before rollout.
- Capture a consistent baseline and its log position. Begin retaining live changes early enough that snapshot work cannot create a gap.
- Apply each row conditionally by version. For this fixture, versions increase per key under one writer; equal version/different contents is corruption. Preserve tombstones until all replay/backfill windows are closed.
- Advance the replay checkpoint only after the destination apply is durable. A crash in between repeats an idempotent apply. Checkpoint-before-apply risks silent omission. The local fixture explicitly injects this gap.
- Reconcile versions, values, ownership, and deletion state across the full cohort. Samples can monitor ongoing health; they do not prove zero mismatches.
The in-memory source/log mutation models one atomic boundary; the companion outbox test uses SQLite to execute an actual transaction. The migration model does not establish durability of a production CDC service. If a paused consumer's checkpoint predates retained history, refuse to continue from an incomplete log. Take a fresh consistent snapshot with a new position, then resume live replay.
Follow-up: a rollback after new-only writes
| Phase | Old reader | New reader | Accepted writer | Rollback requirement |
|---|---|---|---|---|
| Expand/copy | Supported | Shadow only | Old authority | None beyond compatible schema |
| Read cutover | Supported | Supported | Old authority + capture | Route reads back, drain caches/connections |
| Writer cutover | Needs current data | Supported | New authority after fencing old writers | Reverse capture and compatible transformation |
| Contract/retire | Unsupported | Supported | New only | Restore/forward repair; simple route rollback ended |
Run the command in the recovery lab. Supplied tests cover live updates/deletes, simulated partial failure, stale backfill rejection, replay after an apply/checkpoint crash, expired-history resnapshot, and reverse repair before rollback. Equality checks include tombstones and versions, not just counts.
DNS constraint: a load-balancer rule changes new admissions at that balancer; it still needs propagation and connection draining. DNS answers can remain in resolver caches for their TTL, and established connections may survive longer. Route 53 weighted routing is not a promise that all requests can reverse in seconds. See Route 53 TTL, undated technical documentation accessed 2026-09-22. The fixture independently asserts the cached old answer and an existing old connection after registry change.
Lead decision exercise: three teams, two dates
This is a build-and-review brief beyond the model. Security needs tenant isolation in four weeks. Mobile cannot remove an old writer for twelve weeks. Data platform can fund double storage for six weeks. Backfill uses 20% of spare I/O; normal peak already uses 85% of capacity. Produce a one-page decision memo with owners, critical path, a bounded cohort, data authority, stop criteria, and alternatives.
Senior acceptance: no stale overwrite or deleted-row resurrection; zero full cohort mismatches after repair; rollback exercised after a new-only update. Lead acceptance: resolve the 85% + 20% capacity conflict, retain or adapt old writers explicitly, set a budget extension/scope reduction decision, and name who can block cutover. Migrated cohorts may already gain isolation or capacity; retiring the old stack captures additional cost/complexity benefits. A credible memo distinguishes both and discusses blind spots: unused old clients, offline jobs, unobserved deletion paths, and compatibility outside the sampled traffic.