← Monitoring and Observability: measure and investigate
11 / 12 · 60 MIN

Executed retry and recovery

Run five local experiments and distinguish observed recovery, rejection and evidence limits.

Design an experiment that can fail

A recovery test needs an identity that can be followed before and after failure. The complete runner below uses Python 3.13 and Collector contrib 0.162.0, a local HTTP backend and synthetic markers. Each variant has its own temporary directory and loopback ports. The code validates configuration, starts a child process, waits for its readiness message and sends a span with a known identity. The backend records path, status and trace ID. Only processes created by the runner are terminated. There are five experiments and seven Collector starts because the persistent and memory variants restart. The configuration file is JSON, accepted by the YAML loader used in this experiment. Set DR_OTELCOL to the absolute binary path and save the code as run.py. Do not replace destinations with production endpoints to repeat the exercise.

Temporary failure and two restart variants

In the transient variant, the backend returns 503 until at least two attempts are observed. It then starts accepting and the same marker arrives with a 200 response. In the persistent variant, file_storage uses a preserved directory and fsync=true. The runner observes export failure, abruptly terminates its Collector and restarts the same configuration. The pending marker reappears at the backend. This demonstrates recovery after this process failure; it does not test volume loss or host power interruption. In the memory variant, the code observes two seconds after restart without finding the old marker. It then sends a new marker and confirms arrival. This control shows that the path works again. The negative window must remain explicit in the report; it is neither proof of absence for all future time nor a service RPO.

Rejections an HTTP counter can hide

The bad and partial variants use another marker and return HTTP 400 and HTTP 200 with one rejected span in partialSuccess respectively. One attempt was observed in each two-second window; a later probe confirms that the path accepts new data. Compare this record with the response received by the sender: the receiver accepted the request before the backend decision. Consequently, the fictional Vela dashboard cannot label every receiver HTTP 200 as complete delivery. Keep rejection counts and the boundary to which each response belongs. If there are 120 attempts for 100 distinct markers and 95 of these have observed acceptance, the per-marker proportion is 95/100. Do not use 95/120 to answer the same question. This calculation establishes neither durable storage nor deduplication at a vendor.

Hand results and gaps to operations

The reproducible result contains thirty checks, binary version and hash, experiment count and negative observation window. The code checks its own results and removes temporary directories on completion. If it fails, retain the error and investigate before adapting the expectation. Do not change an assertion merely to obtain a green result. In a fictional Lotus handover, attach authorized configuration, volume identity, submitted markers, rejections and destination receipt evidence. If a replica starts with an empty volume, distinguish process readiness from recovery of old data. Saturation, host failure, deferred authentication and real-backend integration experiments remain outstanding. The code ran synthetic data on loopback; it did not instrument applications, validate propagation between services or constitute independent specialist review. The next lesson explores those boundaries through a decision guide.

"""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))
IN PRACTICE

Marker 02 survived restart with the same directory; marker 03 did not appear in the two-second window without persistence.

Common pitfalls

Confusing receiver acceptance with final delivery; generalizing SIGKILL to volume loss; counting retries as new operations.

Related topics: Collector and queues · Telemetry incidents · Distributed traceability

Take this idea with you

Reproduce the failure and control first; report what you observed, the window and conditions not yet exercised.

Create account

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