← Disaster Recovery: prepare, recover, and validate
11 / 12 · 60 MIN

Rebuild the cache and bridge to the watch

Rebuild a teaching consumer after etcd restore, validate origin and completeness, and recover changes occurring before subscription.

Define the state that needs rebuilding

The exercise starts with rule=snapshot and obsolete=present-at-snapshot. After saving the artifact, it changes rule, deletes obsolete, and creates ghost. The cache retains that later state, while the recovered cluster returns to the earlier point. Before running, write both expected sets on paper and identify differences in values, additions, and removals. Rebuilding should make the cache consistent with the authorized source at the selected point. That does not recover business instructions outside the snapshot. In a fictional funds service, the external acknowledgement record still needs its own reconciliation even when the technical cache is correct.

Check identity and completeness before replacement

The consumer receives a response containing cluster_id and revision. It first tries applying it while retaining authorization for the old source: the guard rejects the change and preserves data, cursor, and identity. The exercise also presents a response marked more=true to check rejection of an incomplete set. Only then does it explicitly authorize the cluster created by the script and apply the complete read. In production, accepting any identity returned by an endpoint is not an authorization process. Establish expected origin through controlled configuration and decisions. Pagination requires obtaining the whole scope at a consistent boundary before publishing replacement.

Replace the set and check absences

A simple merge updates rule and adds obsolete, but may retain ghost. The count might even match expectations if another key is missing. The lab replaces the entire dictionary for the authorized prefix, removing entries absent from the read. It separately checks that ghost disappears, obsolete returns, and rule takes the snapshot value. This operation belongs to a rebuildable cache: it must not be applied indiscriminately to business data that is the only source of truth. The design needs to identify derived state, authoritative state, and external effects that cannot be erased through simple local replacement.

Cover the gap between reading and subscribing

After obtaining a read at revision R, the script executes a transaction before opening the watch. The transaction updates one key, deletes another, and creates a third. The consumer subscribes at R+1 while history remains available and receives all three events. The experiment exposes the gap that a watch opened only for future changes would leave uncovered. Later, the script compacts history needed by another attempt and observes rejection. It then repeats the complete read and rebuilds its continuity point. Increasing the timeout for a compacted cursor does not make removed history available again.

Read the format actually returned by the tool

The first adapter version looked for lowercase events and a textual DELETE type. The actual etcdctl 3.6.15 capture contained Events, Header, and numeric 1 for deletion; the program counted zero events despite their delivery. The lab corrects the adapter and normalizes the result before passing it to the consumer. The compact operation also returns textual confirmation that should not be parsed as JSON. These details belong to the executed tool and version. Keep a synthetic capture for parser testing and avoid concluding server-side loss when evidence shows a local interpretation failure.

Run the code and produce an explainable comparison

Save the complete code below as run.py and run python3 run.py --bin-dir /path/to/binaries --output evidence.json with etcd, etcdctl, and etcdutl 3.6.15. The exercise uses sequential clusters, up to three simultaneous processes on one host, and synthetic data only. It needs local ports and terminates processes it started. It accepts no existing endpoints or data directories. Compare observations from the ten groups with predictions from the first section and explain every difference. The bump=1000 value covers only this bounded fixture. The consumer is in memory; durable checkpoints, Kubernetes informers, and external business effects were not exercised.

"""Original cache recovery lab: only its own disposable loopback etcd processes."""
import argparse
import base64
from datetime import datetime, timezone
import hashlib
import json
import os
from pathlib import Path
import socket
import subprocess
import sys
import tempfile
import time
import uuid


def until(action, timeout=15):
 end = time.monotonic + timeout
 last = None
 while time.monotonic < end:
 try:
 result = action
 if result:
 return result
 except (AssertionError, RuntimeError, subprocess.TimeoutExpired) as exc:
 last = str(exc)
 time.sleep(.15)
 raise AssertionError('Condition not observed before deadline: ' + str(last))


class Cluster:
 def __init__(self, folder, binary_dir):
 self.root, self.bin = Path(folder), Path(binary_dir)
 self.token = 'dr-local-' + uuid.uuid4.hex
 self.procs, self.logs, self.clients, self.peers = {}, [], [], []
 reservations = []
 for _ in range(6):
 s = socket.socket; s.bind(('127.0.0.1', 0)); reservations.append(s)
 self.clients = ['http://127.0.0.1:' + str(s.getsockname[1]) for s in reservations[:3]]
 self.peers = ['http://127.0.0.1:' + str(s.getsockname[1]) for s in reservations[3:]]
 for s in reservations:
 s.close
 self.members = ','.join(f'n{i}={self.peers[i]}' for i in range(3))
 # Ignore inherited etcd configuration; never address external endpoints.
 self.env = {k: v for k, v in os.environ.items if not k.startswith(('ETCD_', 'ETCDCTL_'))}

 def start(self, i):
 assert i not in self.procs or self.procs[i].poll is not None
 log = open(self.root/f'n{i}.log', 'ab'); self.logs.append(log)
 args = [str(self.bin/'etcd'), '--name', f'n{i}', '--data-dir', str(self.root/f'n{i}'),
 '--listen-client-urls', self.clients[i], '--advertise-client-urls', self.clients[i],
 '--listen-peer-urls', self.peers[i], '--initial-advertise-peer-urls', self.peers[i],
 '--initial-cluster', self.members, '--initial-cluster-token', self.token,
 '--initial-cluster-state', 'new', '--heartbeat-interval', '100', '--election-timeout', '1000',
 '--log-level', 'error']
 self.procs[i] = subprocess.Popen(args, env=self.env, stdout=log, stderr=log)

 def stop(self, i, abrupt=False):
 p = self.procs[i]
 if p.poll is None:
 p.kill if abrupt else p.terminate
 p.wait(timeout=5)

 def ctl(self, i, *args, raw=False, data=None, expect=True):
 r = subprocess.run([str(self.bin/'etcdctl'), '--endpoints='+self.clients[i], '--dial-timeout=1s',
 '--command-timeout=2s', '--write-out=json', *args], input=data,
 text=True, capture_output=True, timeout=5, env=self.env)
 if expect and r.returncode:
 raise RuntimeError(r.stderr[-600:])
 return r if raw else json.loads(r.stdout)

 def get(self, i, key, serial=False):
 args = ('get', key, '--consistency=s') if serial else ('get', key)
 data = self.ctl(i, *args)
 kv = data.get('kvs', [])
 return None if not kv else base64.b64decode(kv[0]['value']).decode

 def status(self, i):
 return self.ctl(i, 'endpoint', 'status')[0]['Status']

 def leader(self, members):
 def find:
 states = [(i, self.status(i)) for i in members]
 ids = {s.get('leader', 0) for _, s in states}
 if len(ids)!= 1 or 0 in ids:
 return False
 return next(([i] for i, s in states if s['header']['member_id'] == s['leader']), False)
 return until(find)[0]

 def close(self):
 for i in self.procs:
 self.stop(i)
 for log in self.logs:
 log.close



PREFIX = '/dr/cache/'
def decode(x):
 return base64.b64decode(x).decode

def values(response):
 return {decode(k['key']): decode(k.get('value', '')) for k in response.get('kvs', [])}

class Consumer:
 """In-memory teaching consumer; no persistent checkpoints or external effects."""
 def __init__(self, cluster_id, cache, cursor):
 self.cluster_id, self.cache, self.cursor = cluster_id, dict(cache), cursor

 def replace(self, response, authorized_cluster):
 assert response['header']['cluster_id'] == authorized_cluster, 'unexpected cluster identity'
 assert not response.get('more', False), 'incomplete range cannot replace the cache'
 replacement = values(response)
 assert all(k.startswith(PREFIX) for k in replacement), 'unexpected key scope'
 self.cluster_id, self.cache, self.cursor = authorized_cluster, replacement, response['header']['revision']

 def apply(self, response, fail_after=None):
 assert response['header']['cluster_id'] == self.cluster_id, 'unexpected event cluster'
 if response.get('canceled') or response.get('compact_revision'):
 raise ValueError('watch reset required')
 staged, cursor, applied = dict(self.cache), self.cursor, 0
 events = response.get('events', [])
 for event in events:
 kv = event['kv']; revision = kv['mod_revision']; key = decode(kv['key'])
 # Compare with the committed cursor, not each earlier event in this response.
 if revision <= self.cursor:
 continue
 assert key.startswith(PREFIX), 'unexpected event scope'
 kind = event.get('type', 'PUT')
 if kind == 'DELETE': staged.pop(key, None)
 elif kind == 'PUT': staged[key] = decode(kv.get('value', ''))
 else: raise ValueError('unsupported event type')
 cursor = max(cursor, revision); applied += 1
 if fail_after is not None and applied == fail_after:
 raise ValueError('injected handler failure before in-memory commit')
 self.cache, self.cursor = staged, cursor
 return applied


def run(binary_dir):
 binary_dir = Path(binary_dir)
 version = subprocess.run([str(binary_dir/'etcd'), '--version'],text=True,capture_output=True,check=True).stdout
 utility = subprocess.run([str(binary_dir/'etcdutl'), 'version'],text=True,capture_output=True,check=True).stdout
 assert 'etcd Version: 3.6.15' in version and 'etcdutl version: 3.6.15' in utility
 checks = []
 def record(name, **obs): checks.append(dict(name=name, passed=True, observations=obs))
 with tempfile.TemporaryDirectory(prefix='dr-recovery-consumer-') as directory:
 root = Path(directory); clusters=[]
 def cluster(name):
 folder=root/name;folder.mkdir;c=Cluster(folder,binary_dir);clusters.append(c);return c
 def watch(c, revision):
 # History is deliberately written before subscribing, exercising the list/watch gap.
 p=subprocess.Popen([str(binary_dir/'etcdctl'),'--endpoints='+c.clients[0],
 '--write-out=json','watch',PREFIX,'--prefix','--rev='+str(revision)],
 text=True,stdout=subprocess.PIPE,stderr=subprocess.PIPE,env=c.env)
 try:
 try:out,err=p.communicate(timeout=2)
 except subprocess.TimeoutExpired:
 p.terminate;out,err=p.communicate(timeout=3)
 messages=[]
 for line in out.splitlines:
 if not line.strip:continue
 raw=json.loads(line)
 # etcdctl 3.6.15 watch JSON uses exported Go field names and numeric event enums.
 events=[]
 for event in raw.get('Events',[]):
 kind=event.get('type',0);assert kind in (0,1)
 events.append({**event,'type':'DELETE' if kind==1 else 'PUT'})
 messages.append(dict(header=raw['Header'],events=events,canceled=raw.get('Canceled',False),
 compact_revision=raw.get('CompactRevision',0),created=raw.get('Created',False)))
 return messages,err
 finally:
 if p.poll is None:p.kill;p.wait(timeout=3)
 def listing(c):return c.ctl(0,'get',PREFIX,'--prefix')
 def rejected(action):
 try:action
 except (AssertionError,ValueError):return True
 raise AssertionError('Expected rejection was not observed')
 try:
 source=cluster('source')
 for i in range(3):source.start(i)
 until(lambda:all(source.ctl(i,'endpoint','health')[0]['health']for i in range(3)))
 source.ctl(0,'put',PREFIX+'rule','snapshot')
 source.ctl(0,'put',PREFIX+'obsolete','present-at-snapshot')
 snapshot=root/'snapshot.db'source.ctl(0,'snapshot','save',str(snapshot),raw=True)
 digest=hashlib.sha256(snapshot.read_bytes).hexdigest
 source.ctl(0,'put',PREFIX+'rule','later-source')
 source.ctl(0,'del',PREFIX+'obsolete')
 source.ctl(0,'put',PREFIX+'ghost','only-after-snapshot')
 old=listing(source);consumer=Consumer(old['header']['cluster_id'],values(old),old['header']['revision'])
 assert consumer.cache=={PREFIX+'rule':'later-source',PREFIX+'ghost':'only-after-snapshot'}
 record('stale-consumer-retains-post-snapshot-state',cachedKeys=2,postSnapshotOnlyKeyPresent=True)
 for i in range(3):source.stop(i)
 restored=cluster('restored')
 for i in range(3):
 subprocess.run([str(binary_dir/'etcdutl'),'snapshot','restore',str(snapshot),
 '--name',f'n{i}','--data-dir',str(restored.root/f'n{i}'),'--initial-cluster',restored.members,
 '--initial-cluster-token',restored.token,'--initial-advertise-peer-urls',restored.peers[i],
 '--bump-revision','1000','--mark-compacted'],check=True,text=True,capture_output=True,timeout=20)
 for i in range(3):restored.start(i)
 until(lambda:all(restored.ctl(i,'endpoint','health')[0]['health']for i in range(3)))
 msgs,err=watch(restored,consumer.cursor+1)
 assert 'compacted' in (json.dumps(msgs)+err).lower
 record('old-consumer-cursor-observes-compaction',compactionObserved=True,cacheStillStale=True)
 fresh=listing(restored);new_id=fresh['header']['cluster_id'];assert new_id!=consumer.cluster_id
 before=(dict(consumer.cache),consumer.cursor,consumer.cluster_id)
 assert rejected(lambda:consumer.replace(fresh,consumer.cluster_id))
 assert (consumer.cache,consumer.cursor,consumer.cluster_id)==before
 incomplete={**fresh,'more':True}
 assert rejected(lambda:consumer.replace(incomplete,new_id))
 assert (consumer.cache,consumer.cursor,consumer.cluster_id)==before
 # Only the cluster created by this script is authorized here; this is not a production approval mechanism.
 record('identity-and-completeness-guards-preserve-old-cache',wrongIdentityRejected=True,incompleteRangeRejected=True,priorStateUnchanged=True)
 consumer.replace(fresh,new_id);listed_revision=consumer.cursor;listed_cache=dict(consumer.cache)
 assert consumer.cache=={PREFIX+'rule':'snapshot',PREFIX+'obsolete':'present-at-snapshot'}
 record('complete-relist-replaces-instead-of-merging',ghostRemoved=True,snapshotDeletionReversed=True,valueReconciled=True)
 tx='\nput "'+PREFIX+'rule" "gap-update"\ndel "'+PREFIX+'obsolete"\nput "'+PREFIX+'fresh" "gap-create"\n\n\n'
 txn=restored.ctl(0,'txn',data=tx);assert txn['succeeded']
 revision=txn['header']['revision'];assert revision>listed_revision
 msgs,err=watch(restored,listed_revision+1)
 event_messages=[m for m in msgs if m.get('events')];assert event_messages,(txn,msgs,err,listing(restored))
 batch=next(m for m in event_messages if len(m['events'])==3)
 assert {e['kv']['mod_revision']for e in batch['events']}=={revision}
 record('watch-replays-three-events-across-list-gap',events=3,sameTransactionRevision=True,historyWrittenBeforeWatch=True)
 before=(dict(consumer.cache),consumer.cursor)
 assert rejected(lambda:consumer.apply(batch,fail_after=1))
 assert (consumer.cache,consumer.cursor)==before
 record('handler-failure-does-not-publish-partial-cache',failureAfterStagedEvents=1,cacheUnchanged=True,cursorUnchanged=True,persistentCrashRecoveryTested=False)
 assert consumer.apply(batch)==3 and consumer.cursor==revision
 assert consumer.cache==values(listing(restored))
 replay_before=(dict(consumer.cache),consumer.cursor)
 assert consumer.apply(batch)==0 and (consumer.cache,consumer.cursor)==replay_before
 record('whole-batch-applies-once-to-cache',eventsApplied=3,replayedEventsApplied=0,cacheMatchesRange=True,externalEffects=0)
 restored.ctl(0,'put',PREFIX+'rule','after-disconnect')
 restored.ctl(0,'del',PREFIX+'fresh')
 msgs,err=watch(restored,consumer.cursor+1);applied=sum(consumer.apply(m)for m in msgs if m.get('events'))
 assert applied==2 and consumer.cache==values(listing(restored))
 record('reconnect-with-retained-history-catches-up',missedEventsApplied=2,deleteApplied=True,cacheMatchesRange=True)
 old_read=listing(restored);old_rev=old_read['header']['revision']
 restored.ctl(0,'put',PREFIX+'rule','after-old-read')
 newer=restored.ctl(0,'put',PREFIX+'after-compaction','current')['header']['revision']
 compacted=restored.ctl(0,'compact',str(newer),raw=True)
 assert 'compacted revision' in compacted.stdout.lower,compacted.stdout
 msgs,err=watch(restored,old_rev+1)
 assert 'compacted' in (json.dumps(msgs)+err).lower
 consumer.replace(listing(restored),new_id)
 assert consumer.cache==values(listing(restored))
 record('compacted-list-gap-restarts-from-new-complete-range',oldWatchRejected=True,completeRelistApplied=True,cacheMatchesRange=True)
 restored.ctl(0,'put','/dr/outside/ignored','not-in-cache')
 restored.ctl(0,'put',PREFIX+'acceptance','observed')
 msgs,err=watch(restored,consumer.cursor+1)
 applied=sum(consumer.apply(m)for m in msgs if m.get('events'))
 assert applied==1 and consumer.cache==values(listing(restored))
 assert '/dr/outside/ignored' not in consumer.cache and hashlib.sha256(snapshot.read_bytes).hexdigest==digest
 record('scoped-continuity-and-artifact-preservation',scopedEventsApplied=1,outsideKeyAbsent=True,cacheMatchesRange=True,snapshotUnchanged=True)
 finally:
 for c in clusters:c.close
 return dict(executedAt=datetime.now(timezone.utc).isoformat,version=version.strip,utilityVersion=utility.strip,passed=len(checks),failed=0,checks=checks,scriptSha256=hashlib.sha256(Path(__file__).read_bytes).hexdigest,scope='Original in-memory Python consumer over actual disposable etcd 3.6.15 restore and watches. Sequential three-member clusters on one Darwin ARM64 host, synthetic keys, loopback HTTP. Full-range replacement, retained-history replay, three-event transaction, injected pre-commit handler failure, reconnect and compaction reset observed. No existing cluster, Kubernetes informer, persistent checkpoint crash recovery, external business effects, production workload, credentials, physical fencing or RPO/RTO benchmark.')

if __name__=='__main__':
 p=argparse.ArgumentParser;p.add_argument('--bin-dir',required=True,type=Path);p.add_argument('--output',type=Path);a=p.parse_args;result=run(a.bin_dir);payload=json.dumps(result,indent=2)+'\n'
 if a.output:a.output.write_text(payload)
 print(payload)
IN PRACTICE

The recovered source contains rule and obsolete; the old cache contains rule and ghost. The exercise demonstrates that replacement removes ghost, recovers obsolete, and corrects rule, while a merge could retain inappropriate state.

Common pitfalls

Merging without removing absences; accepting one page as a complete set; watching only future changes; trusting any returned cluster_id; assuming JSON format without observing it.

Related topics: Restore: recovered point and integrity · Snapshot, identity, and recovered point · Revisions, watches, and controlled resumption

Take this idea with you

Rebuilding needs authorized origin, a complete set, and a verifiable bridge to subsequent changes. Comparison should include values, additions, and removals within the same scope.

Create account

Reference: Disaster recovery · BigSavant recovery 2026-09; PostgreSQL 18, etcd 3.6 and selected AWS/Azure behavior