← Monitorização e Observabilidade: medir e investigar
11 / 12 · 60 MIN

Retry e recuperação executados

Executa cinco experiências locais e distingue recuperação observada, rejeição e limites da evidência.

Desenhar uma experiência que possa falhar

Um teste de recuperação precisa de uma identidade que possa ser seguida antes e depois da falha. O runner completo abaixo usa Python 3.13 e Collector contrib 0.162.0, um backend HTTP local e marcadores sintéticos. Cada variante tem uma pasta temporária própria e portas loopback. O código valida a configuração, inicia um subprocesso, espera pela mensagem de prontidão e envia um span com identidade conhecida. O backend regista caminho, estado e trace ID. Apenas os processos criados pelo runner são terminados. São cinco experiências e sete arranques do Collector, porque as variantes com e sem persistência reiniciam. O ficheiro de configuração é JSON, aceite pelo carregador YAML utilizado neste ensaio. Define a variável DR_OTELCOL com o caminho absoluto do binário e guarda o código num ficheiro run.py. Não substituas os destinos por produção para repetir o exercício.

Falha temporária e dois tipos de reinício

Na variante transient, o backend responde 503 até observar pelo menos duas tentativas. Depois passa a aceitar e o mesmo marcador chega com resposta 200. Na variante persistent, file_storage usa uma pasta preservada e fsync=true. O runner observa a falha de exportação, termina abruptamente o seu Collector e reinicia a mesma configuração. O marcador pendente reaparece no backend. Isto demonstra recuperação após esta falha de processo; não testa perda do volume ou corte de energia no host. Na variante memory, o código observa dois segundos após reinício sem encontrar o marcador antigo. Em seguida envia um marcador novo e confirma a chegada. Este controlo mostra que o caminho voltou a funcionar. A janela negativa deve permanecer explícita no relatório; não é uma prova de ausência para todo o futuro nem um RPO do serviço.

Rejeições que um contador HTTP pode esconder

As variantes bad e partial usam outro marcador e devolvem, respetivamente, HTTP 400 e HTTP 200 com um span rejeitado em partialSuccess. Em cada variante foi observada uma tentativa na janela de dois segundos; uma sonda posterior confirma que o caminho aceita novos dados. Compara este registo com a resposta recebida pelo emissor: o receiver aceitou o pedido antes da decisão do backend. Portanto, o dashboard fictício Vela não pode chamar entrega completa a todos os HTTP 200 do receiver. Guarda a contagem de rejeições e a fronteira a que cada resposta pertence. Se há 120 tentativas para 100 marcadores distintos e 95 destes têm aceitação observada, a proporção por marcador é 95/100. Não uses 95/120 para responder à mesma pergunta. Este cálculo não demonstra armazenamento durável nem desduplicação num fornecedor.

Entregar resultados e lacunas ao RUN

O resultado reproduzível reúne trinta verificações, versão e hash do binário, número de experiências e janela de observação negativa. O código verifica os próprios resultados e remove os diretórios temporários ao terminar. Se falhar, conserva a mensagem de erro e investiga antes de adaptar a expectativa. Não alteres uma asserção só para obter verde. Num handover fictício Lótus, anexa configuração autorizada, identidade dos volumes, marcadores submetidos, rejeições e evidência de receção por destino. Se uma réplica arranca com volume vazio, distingue prontidão de processo de recuperação dos dados antigos. Ainda faltam ensaios de saturação, falha do host, autenticação diferida e integração com o backend real. O código executou dados sintéticos em loopback; não instrumentou aplicações, não validou propagação entre serviços e não representa revisão independente por especialista. A aula seguinte trabalha essas fronteiras através de um guião de decisão.

"""Original DR loopback experiment. Python 3.13; Collector contrib 0.162.0.
Run: DR_OTELCOL=/absolute/path/to/otelcol-contrib python3 run.py
Kills only Collector child processes created here. Uses temporary synthetic data.
No production credentials, application instrumentation or external backend.
"""
import hashlib
import http.server
import json
import os
from pathlib import Path
import socket
import subprocess
import tempfile
import threading
import time
import urllib.request

BINARY = Path(os.environ['DR_OTELCOL']).resolve()
VERSION = subprocess.check_output([str(BINARY), '--version'], text=True).strip()
assert VERSION == 'otelcol-contrib version 0.162.0', VERSION
checks = []


def check(name, condition):
    assert condition, name
    checks.append(name)


def wait_for(predicate, timeout=8):
    until = time.monotonic() + timeout
    while time.monotonic() < until:
        if predicate():
            return
        time.sleep(.02)
    raise AssertionError('Timed out waiting for experiment condition')


def port():
    with socket.socket() as sock:
        sock.bind(('127.0.0.1', 0))
        return sock.getsockname()[1]


class Backend(http.server.ThreadingHTTPServer):
    daemon_threads = True

    def __init__(self):
        super().__init__(('127.0.0.1', 0), Handler)
        self.mode = 'outage'
        self.records = []
        self.lock = threading.Lock()

    def snapshot(self):
        with self.lock:
            return list(self.records)


class Handler(http.server.BaseHTTPRequestHandler):
    def log_message(self, *_):
        pass

    def do_POST(self):
        payload = json.loads(self.rfile.read(int(self.headers['Content-Length'])))
        spans = [span for resource in payload.get('resourceSpans', [])
                 for scope in resource.get('scopeSpans', [])
                 for span in scope.get('spans', [])]
        with self.server.lock:
            mode = self.server.mode
            status = 503 if mode == 'outage' else 400 if mode == 'bad' else 200
            self.server.records.append({'status': status, 'path': self.path,
                                        'ids': [s['traceId'] for s in spans]})
        body = {'message': 'synthetic backend unavailable'} if status == 503 else (
            {'message': 'synthetic permanent rejection'} if status == 400 else (
                {'partialSuccess': {'rejectedSpans': str(len(spans)),
                                    'errorMessage': 'synthetic rejection'}}
                if mode == 'partial' else {}))
        data = json.dumps(body).encode()
        self.send_response(status)
        self.send_header('Content-Type', 'application/json')
        self.send_header('Content-Length', str(len(data)))
        self.end_headers()
        self.wfile.write(data)


def payload(marker):
    return {'resourceSpans': [{'resource': {'attributes': [
        {'key': 'service.name', 'value': {'stringValue': 'dr-synthetic-recovery'}}]},
        'scopeSpans': [{'spans': [{'traceId': marker, 'spanId': '0000000000000001',
                                  'name': 'synthetic-operation', 'kind': 1,
                                  'startTimeUnixNano': '1791028800000000000',
                                  'endTimeUnixNano': '1791028800001000000'}]}]}]}


def run_case(root, name, mode, persistent=False):
    work = root / name
    work.mkdir()
    backend = Backend()
    backend.mode = mode
    thread = threading.Thread(target=backend.serve_forever, daemon=True)
    thread.start()
    receiver, metrics = port(), port()
    queue = {'enabled': True, 'queue_size': 10, 'num_consumers': 1}
    config = {
        'receivers': {'otlp/lab': {'protocols': {'http': {'endpoint': f'127.0.0.1:{receiver}'}}}},
        'exporters': {'otlp_http/lab': {
            'endpoint': f'http://127.0.0.1:{backend.server_port}',
            'encoding': 'json', 'compression': 'none', 'timeout': '1s',
            'sending_queue': queue,
            'retry_on_failure': {'initial_interval': '100ms', 'max_interval': '200ms',
                                 'max_elapsed_time': '30s'}}},
        'service': {'telemetry': {'metrics': {'readers': [{'pull': {'exporter': {
            'prometheus': {'host': '127.0.0.1', 'port': metrics}}}}]}},
            'pipelines': {'traces': {'receivers': ['otlp/lab'], 'exporters': ['otlp_http/lab']}}}}
    if persistent:
        queue['storage'] = 'file_storage'
        config['extensions'] = {'file_storage': {'directory': str(work / 'storage'),
                                                'create_directory': True, 'fsync': True}}
        config['service']['extensions'] = ['file_storage']
    cfg = work / 'config.json'
    cfg.write_text(json.dumps(config))
    result = subprocess.run([str(BINARY), 'validate', '--config', str(cfg)], capture_output=True)
    check(name + ': configuration validates', result.returncode == 0)
    process = None
    logs = None

    def start(suffix):
        nonlocal process, logs
        logfile = work / (suffix + '.log')
        logs = logfile.open('wb')
        process = subprocess.Popen([str(BINARY), '--config', str(cfg)], stdout=logs, stderr=logs)
        try:
            wait_for(lambda: b'Everything is ready' in logfile.read_bytes())
        except Exception:
            raise AssertionError(logfile.read_text())

    def stop(abrupt=False):
        nonlocal process, logs
        if process and process.poll() is None:
            process.kill() if abrupt else process.terminate()
            try:
                process.wait(timeout=5)
            except subprocess.TimeoutExpired:
                process.kill()
                process.wait(timeout=5)
        if logs:
            logs.close()

    def send(marker):
        request = urllib.request.Request(f'http://127.0.0.1:{receiver}/v1/traces',
            data=json.dumps(payload(marker)).encode(), headers={'Content-Type': 'application/json'})
        with urllib.request.urlopen(request, timeout=3) as response:
            body = json.loads(response.read())
            check(name + ': receiver accepts ' + marker[-2:],
                  response.status == 200 and not body.get('partialSuccess', {}).get('rejectedSpans'))

    marker = {'transient': '1', 'persistent': '2', 'memory': '3', 'bad': '4', 'partial': '5'}[name].zfill(32)
    try:
        start('first')
        send(marker)
        wait_for(lambda: len(backend.snapshot()) >= 1)
        check(name + ': base endpoint adds trace path',
              all(r['path'] == '/v1/traces' for r in backend.snapshot()))
        if name == 'transient':
            wait_for(lambda: len(backend.snapshot()) >= 2)
            check('transient: repeated 503 preserves trace identity',
                  all(r['status'] == 503 and r['ids'] == [marker] for r in backend.snapshot()))
            backend.mode = 'success'
            wait_for(lambda: any(r['status'] == 200 and marker in r['ids'] for r in backend.snapshot()))
            check('transient: marker reaches backend after recovery', True)
        elif name in ('persistent', 'memory'):
            check(name + ': pending marker met unavailable backend',
                  backend.snapshot()[0]['status'] == 503)
            stop(abrupt=True)
            boundary = len(backend.snapshot())
            backend.mode = 'success'
            start('second')
            if persistent:
                wait_for(lambda: any(marker in r['ids'] and r['status'] == 200
                                     for r in backend.snapshot()[boundary:]))
                check('persistent: pending marker recovered after SIGKILL and restart', True)
                check('persistent: storage file exists', any((work / 'storage').iterdir()))
            else:
                # Bounded negative observation, not proof that arrival is impossible forever.
                time.sleep(2)
                check('memory: old marker absent in two-second post-restart window',
                      not any(marker in r['ids'] for r in backend.snapshot()[boundary:]))
                fresh = '6'.zfill(32)
                send(fresh)
                wait_for(lambda: any(fresh in r['ids'] for r in backend.snapshot()[boundary:]))
                check('memory: new probe confirms recovered path', True)
        else:
            time.sleep(2)
            check(name + ': one attempt observed in two-second window', len(backend.snapshot()) == 1)
            backend.mode = 'success'
            fresh = ('7' if name == 'bad' else '8').zfill(32)
            send(fresh)
            wait_for(lambda: any(fresh in r['ids'] for r in backend.snapshot()))
            check(name + ': later valid probe delivered', True)
    finally:
        stop()
        backend.shutdown()
        backend.server_close()
        thread.join(timeout=3)


with tempfile.TemporaryDirectory(prefix='dr-otel-recovery-') as directory:
    root = Path(directory)
    run_case(root, 'transient', 'outage')
    run_case(root, 'persistent', 'outage', persistent=True)
    run_case(root, 'memory', 'outage')
    run_case(root, 'bad', 'bad')
    run_case(root, 'partial', 'partial')

print(json.dumps({'version': VERSION, 'binarySha256': hashlib.sha256(BINARY.read_bytes()).hexdigest(),
                  'checks_passed': len(checks), 'checks': checks,
                  'experiments': 5, 'collector_starts': 7, 'negative_window_seconds': 2,
                  'scope': 'Synthetic loopback HTTP backend and owned Collector children only. '
                           'No vendor backend, application propagation, host-power failure, '
                           'queue saturation or production acceptance tested.'}, indent=2))
NA PRÁTICA

O marcador 02 sobreviveu ao reinício com a mesma pasta; o marcador 03 não apareceu na janela de dois segundos sem persistência.

Armadilhas comuns

Confundir aceitação do receiver com entrega final; generalizar SIGKILL a perda do volume; contar retries como operações novas.

Tópicos relacionados: Collector e filas · Incidentes de telemetria · Rastreabilidade distribuída

Leva esta ideia contigo

Reproduz primeiro a falha e o controlo; relata o que observaste, a janela e as condições ainda não exercitadas.

Criar conta

Referência: Collector release and original DR recovery experiments · Observability 2026-09; selected OpenTelemetry, Prometheus and Dynatrace Classic concepts