#!/usr/bin/env python3 """ m10_smoke.py — M10 soak test (5k/s, 10 min, zero DLQ, runaway-source isolation) Usage: # Local docker-compose: docker compose --profile loadgen-m10 up -d python3 scripts/m10_smoke.py # Remote parres: ssh root@192.168.44.94 "cd /root/broad-announce && docker compose --profile loadgen-m10 up -d && python3 scripts/m10_smoke.py" Exit code 0 = all green. Exit code 1 = assertion failed. Step 1 — pre-flight Step 2 — 5k/s soak (10 min, ramp 30s) Step 3 — runaway-source fault injection (60s) Step 4 — DLQ invariant Step 5 — teardown (profile down) """ import subprocess import sys import time import urllib.request import urllib.parse import json import base64 sys.path.insert(0, __file__.rsplit("/", 1)[0]) import m10_lib as lib PROM = "http://localhost:9090" GRAFANA = "http://localhost:3001" SOAK_DURATION_MIN = 10 # minutes SOAK_RAMP_SEC = 30 # ramp-up seconds RUNAWAY_DURATION_SEC = 60 # runaway fault injection duration CLUSTER_TARGET = 5000 # alerts/sec cluster-wide target P99_THRESHOLD = 5.0 # seconds — p99 must be under this DLQ_EXPECTED = 0 # zero DLQ is the invariant RATE_TOLERANCE = 0.05 # ±5% def curl_json(url: str) -> dict | None: try: with urllib.request.urlopen(url, timeout=10) as r: return json.loads(r.read()) except Exception: return None def pass_(msg: str): print(f" ✅ {msg}") def warn_(msg: str): print(f" ⚠️ {msg}") def fail_(msg: str): print(f" ❌ {msg}") sys.exit(1) def step1_preflight() -> dict: """Verify all services are up and DLQ baseline is clean.""" print("Step 1 — pre-flight") targets = lib.get_prometheus_targets() required = ["ingestd", "routerd", "deliverd-fcm", "deliverd-telegram", "admind", "archiverd", "prometheus", "loadgen-http-1", "loadgen-http-2", "loadgen-http-3"] all_up = True for svc in required: status = targets.get(svc, "0") if status == "1": pass_(f"{svc} is up") else: warn_(f"{svc} is {'not scraped' if status == '0' else status}") all_up = False # Check DLQ baseline dlq_now = lib.assert_dlq_count_equals(0, window_seconds=60) pass_(f"DLQ baseline clean: {dlq_now} rows") # Check circuit breaker is closed cb_state = lib.get_circuit_breaker_state("nats") cb_names = {0: "CLOSED", 1: "HALF-OPEN", 2: "OPEN"} if cb_state == 0: pass_(f"circuit breaker CLOSED (nats)") elif cb_state < 0: warn_(f"circuit breaker gauge not initialized yet (expected on cold start)") else: warn_(f"circuit breaker is {cb_names.get(cb_state, cb_state)} — may recover under load") if not all_up: fail_("not all services are up — fix before running M10 smoke") return targets def step2_start_loadgen() -> subprocess.Popen: """Start the 3-instance loadgen cluster.""" print("\nStep 2 — starting 3-instance loadgen cluster (5k/s)") proc = subprocess.Popen( ["docker", "compose", "--profile", "loadgen-m10", "up", "-d"], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) code = proc.wait() if code != 0: fail_(f"docker compose --profile loadgen-m10 up -d failed (exit {code})") pass_("3 loadgen-http instances started") # Wait for ramp-up print(f" waiting {SOAK_RAMP_SEC}s for ramp-up to complete...", flush=True) time.sleep(SOAK_RAMP_SEC) pass_(f"ramp-up complete — now targeting {CLUSTER_TARGET}/s") return proc def step2_monitor_soak() -> dict: """ Monitor the soak: sample p99 + rate + DLQ every 30s. Fails fast on any breach. Returns dict of samples for the log. """ print(f"\n Monitoring soak for {SOAK_DURATION_MIN} minutes...") samples = [] start = time.time() deadline = start + SOAK_DURATION_MIN * 60 sample_interval = 30 # seconds between samples while time.time() < deadline: time.sleep(sample_interval) elapsed = int(time.time() - start) // 60 try: rate = lib.assert_rate_near(CLUSTER_TARGET, tolerance=RATE_TOLERANCE, window_seconds=30) p99 = lib.assert_ingest_p99_under(P99_THRESHOLD, window_seconds=60) dlq = lib.assert_dlq_count_equals(DLQ_EXPECTED, window_seconds=60) print(f" [{elapsed}m] rate={rate:.0f}/s p99={p99:.3f}s dlq={dlq}") samples.append({"elapsed_min": elapsed, "rate": rate, "p99": p99, "dlq": dlq}) except AssertionError as e: fail_(f"soak breach at {elapsed}m: {e}") return samples def step3_runaway_test() -> None: """ Runaway-source fault injection. - Start a 4th loadgen instance firing at 10× per-source cap for the same company. - Verify p99 for the OTHER sources (acme-002, acme-003) stays under threshold. - The runaway (acme-001) can be anything. """ print(f"\nStep 3 — runaway-source fault injection ({RUNAWAY_DURATION_SEC}s)") print(" (not yet implemented — requires --rate override on loadgen-http-1)") warn_(f"runaway-source test skipped (W4 enhancement pending)") # TODO: spin up a 4th instance at 10× cap targeting acme-001 # Expected: per_source_p99("acme-002") < 5s and per_source_p99("acme-003") < 5s def step4_dlq_invariant(samples: list[dict]) -> int: """Assert zero DLQ rows for the entire soak window.""" print("\nStep 4 — DLQ invariant check") # The soak samples already checked DLQ delta, but do a final absolute check. # Query the full soak window. soak_seconds = SOAK_DURATION_MIN * 60 dlq = lib.assert_dlq_count_equals(DLQ_EXPECTED, window_seconds=soak_seconds) pass_(f"DLQ count over {SOAK_DURATION_MIN}min soak: {dlq} (expected 0)") return dlq def step5_teardown(proc) -> None: """Bring down the loadgen cluster.""" print("\nStep 5 — teardown") r = subprocess.run( ["docker", "compose", "--profile", "loadgen-m10", "down", "-v"], capture_output=True, ) if r.returncode == 0: pass_("loadgen cluster torn down") else: warn_(f"teardown returned {r.returncode}: {r.stderr.decode().strip()}") def print_summary(samples: list[dict], dlq_final: int) -> None: """Print a summary table of the soak run.""" print("\n=== M10 Soak Summary ===") print(f"Duration: {SOAK_DURATION_MIN} min") print(f"Target: {CLUSTER_TARGET}/s (±{RATE_TOLERANCE*100:.0f}%)") print(f"p99 threshold: {P99_THRESHOLD}s") print() if samples: print(f"{'Time':>6} {'Rate/s':>8} {'p99(s)':>7} {'DLQ':>4}") print("-" * 35) for s in samples: print(f"{s['elapsed_min']:>5}m {s['rate']:>8.0f} {s['p99']:>7.3f} {s['dlq']:>4}") print() print(f"Final DLQ count: {dlq_final} (expected 0)") print() print("🎉 M10 smoke: all checks complete.") def main(): print(f"M10 Soak Test — target {CLUSTER_TARGET}/s for {SOAK_DURATION_MIN} min") print(f"p99 threshold: {P99_THRESHOLD}s | DLQ expected: {DLQ_EXPECTED}") print() try: step1_preflight() proc = step2_start_loadgen() samples = step2_monitor_soak() step3_runaway_test() dlq_final = step4_dlq_invariant(samples) step5_teardown(proc) print_summary(samples, dlq_final) except AssertionError as e: print(f"\n💥 M10 smoke FAILED: {e}") sys.exit(1) except KeyboardInterrupt: print("\nInterrupted.") sys.exit(1) if __name__ == "__main__": main()