| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224 |
- #!/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()
|