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))
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
Reproduce the failure and control first; report what you observed, the window and conditions not yet exercised.
Reference: Collector release and original DR recovery experiments · Observability 2026-09; selected OpenTelemetry, Prometheus and Dynatrace Classic concepts