Definir o estado que precisa de ser reconstruído
O exercício começa com rule=snapshot e obsolete=present-at-snapshot. Depois de guardar o artefacto, altera rule, elimina obsolete e cria ghost. A cache guarda esse estado posterior, mas o cluster recuperado volta ao ponto anterior. Antes de executar, escreve os dois conjuntos esperados numa folha e identifica as diferenças de valor, inclusão e exclusão. A reconstrução deve tornar a cache coerente com a origem autorizada no ponto selecionado. Isso não recupera instruções de negócio que ficaram fora do snapshot. Num serviço fictício de fundos, o registo externo das confirmações continua a exigir reconciliação própria, mesmo que a cache técnica já esteja correta.
Verificar identidade e completude antes de substituir
O consumidor recebe uma resposta com cluster_id e revisão. Primeiro tenta aplicá-la mantendo a autorização da origem antiga: a guarda rejeita a mudança e preserva dados, cursor e identidade. O exercício também apresenta uma resposta marcada more=true para verificar a recusa de um conjunto incompleto. Só depois autoriza explicitamente o cluster criado pelo próprio script e aplica a leitura completa. Em produção, aceitar qualquer identidade devolvida por um endpoint não é um processo de autorização. Define a origem esperada através de configuração e decisão controladas. Se houver paginação, é necessário obter todo o âmbito numa fronteira coerente antes de publicar a substituição.
Substituir o conjunto e conferir as ausências
Uma fusão simples atualiza rule e acrescenta obsolete, mas pode conservar ghost. A contagem pode até coincidir com a esperada se outra chave estiver ausente. O laboratório substitui o dicionário completo do prefixo autorizado, removendo entradas que não existem na leitura. Verifica separadamente que ghost desaparece, obsolete regressa e rule assume o valor do snapshot. Esta operação é própria de uma cache reconstruível: não deve ser aplicada indiscriminadamente a dados de negócio que sejam a única fonte de verdade. O desenho precisa de identificar o que é derivado, o que é autoritativo e que efeitos externos não podem ser apagados por uma simples substituição local.
Cobrir o intervalo entre leitura e subscrição
Depois de obter uma leitura na revisão R, o script executa uma transação antes de abrir o watch. A transação atualiza uma chave, elimina outra e cria uma terceira. O consumidor abre a subscrição em R+1, com historial ainda disponível, e recebe os três eventos. A experiência torna visível o intervalo que um watch aberto apenas para alterações futuras deixaria por cobrir. Mais adiante, o script compacta o historial necessário a outra tentativa e observa a rejeição. Nesse caso, repete a leitura completa e reconstrói o ponto de continuidade. Aumentar o timeout de um cursor compactado não volta a disponibilizar o historial removido.
Ler o formato realmente devolvido pela ferramenta
A primeira versão do adaptador procurava events em minúsculas e um tipo DELETE textual. A captura real de etcdctl 3.6.15 apresentou Events, Header e o número 1 para uma eliminação; o programa contava zero eventos apesar de estes terem sido entregues. O laboratório corrige o adaptador e normaliza o resultado antes de o passar ao consumidor. A operação compact também devolve uma confirmação textual que não deve ser interpretada como JSON. Estes detalhes pertencem à ferramenta e versão executadas. Conserva uma captura sintética para testar o parser e evita concluir perda no servidor quando a evidência mostra uma falha local de interpretação.
Executar o código e produzir uma comparação explicável
Guarda o código integral abaixo como run.py e executa python3 run.py --bin-dir /caminho/para/binarios --output evidence.json com etcd, etcdctl e etcdutl 3.6.15. O exercício usa clusters sequenciais, até três processos simultâneos no mesmo host e apenas dados sintéticos. Precisa de portas locais e termina os processos que iniciou. Não aceita endpoints ou diretórios de dados existentes. Compara as observações dos dez grupos com as previsões da primeira secção e explica cada diferença. O valor bump=1000 só cobre este conjunto limitado. O consumidor está em memória; não foram ensaiados checkpoints duráveis, informers Kubernetes ou efeitos externos de negócio.
"""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)
A origem recuperada contém rule e obsolete; a cache antiga contém rule e ghost. O exercício demonstra que substituir o conjunto remove ghost, recupera obsolete e corrige rule, enquanto uma fusão poderia conservar estado indevido.
Armadilhas comuns
Fazer merge sem remover ausências; aceitar uma página como conjunto completo; abrir watch apenas para o futuro; confiar em qualquer cluster_id devolvido; assumir o formato JSON sem o observar.
Tópicos relacionados: Restauro: ponto recuperado e integridade · Snapshot, identidade e ponto recuperado · Revisões, watches e retoma controlada
A reconstrução precisa de origem autorizada, conjunto completo e uma ligação verificável às alterações seguintes. A comparação deve incluir valores, inclusões e exclusões dentro do mesmo âmbito.
Referência: Disaster recovery · DR recovery 2026-09; PostgreSQL 18, etcd 3.6 and selected AWS/Azure behavior