m10_smoke.py 7.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224
  1. #!/usr/bin/env python3
  2. """
  3. m10_smoke.py — M10 soak test (5k/s, 10 min, zero DLQ, runaway-source isolation)
  4. Usage:
  5. # Local docker-compose:
  6. docker compose --profile loadgen-m10 up -d
  7. python3 scripts/m10_smoke.py
  8. # Remote parres:
  9. ssh root@192.168.44.94 "cd /root/broad-announce && docker compose --profile loadgen-m10 up -d && python3 scripts/m10_smoke.py"
  10. Exit code 0 = all green. Exit code 1 = assertion failed.
  11. Step 1 — pre-flight
  12. Step 2 — 5k/s soak (10 min, ramp 30s)
  13. Step 3 — runaway-source fault injection (60s)
  14. Step 4 — DLQ invariant
  15. Step 5 — teardown (profile down)
  16. """
  17. import subprocess
  18. import sys
  19. import time
  20. import urllib.request
  21. import urllib.parse
  22. import json
  23. import base64
  24. sys.path.insert(0, __file__.rsplit("/", 1)[0])
  25. import m10_lib as lib
  26. PROM = "http://localhost:9090"
  27. GRAFANA = "http://localhost:3001"
  28. SOAK_DURATION_MIN = 10 # minutes
  29. SOAK_RAMP_SEC = 30 # ramp-up seconds
  30. RUNAWAY_DURATION_SEC = 60 # runaway fault injection duration
  31. CLUSTER_TARGET = 5000 # alerts/sec cluster-wide target
  32. P99_THRESHOLD = 5.0 # seconds — p99 must be under this
  33. DLQ_EXPECTED = 0 # zero DLQ is the invariant
  34. RATE_TOLERANCE = 0.05 # ±5%
  35. def curl_json(url: str) -> dict | None:
  36. try:
  37. with urllib.request.urlopen(url, timeout=10) as r:
  38. return json.loads(r.read())
  39. except Exception:
  40. return None
  41. def pass_(msg: str):
  42. print(f" ✅ {msg}")
  43. def warn_(msg: str):
  44. print(f" ⚠️ {msg}")
  45. def fail_(msg: str):
  46. print(f" ❌ {msg}")
  47. sys.exit(1)
  48. def step1_preflight() -> dict:
  49. """Verify all services are up and DLQ baseline is clean."""
  50. print("Step 1 — pre-flight")
  51. targets = lib.get_prometheus_targets()
  52. required = ["ingestd", "routerd", "deliverd-fcm", "deliverd-telegram",
  53. "admind", "archiverd", "prometheus",
  54. "loadgen-http-1", "loadgen-http-2", "loadgen-http-3"]
  55. all_up = True
  56. for svc in required:
  57. status = targets.get(svc, "0")
  58. if status == "1":
  59. pass_(f"{svc} is up")
  60. else:
  61. warn_(f"{svc} is {'not scraped' if status == '0' else status}")
  62. all_up = False
  63. # Check DLQ baseline
  64. dlq_now = lib.assert_dlq_count_equals(0, window_seconds=60)
  65. pass_(f"DLQ baseline clean: {dlq_now} rows")
  66. # Check circuit breaker is closed
  67. cb_state = lib.get_circuit_breaker_state("nats")
  68. cb_names = {0: "CLOSED", 1: "HALF-OPEN", 2: "OPEN"}
  69. if cb_state == 0:
  70. pass_(f"circuit breaker CLOSED (nats)")
  71. elif cb_state < 0:
  72. warn_(f"circuit breaker gauge not initialized yet (expected on cold start)")
  73. else:
  74. warn_(f"circuit breaker is {cb_names.get(cb_state, cb_state)} — may recover under load")
  75. if not all_up:
  76. fail_("not all services are up — fix before running M10 smoke")
  77. return targets
  78. def step2_start_loadgen() -> subprocess.Popen:
  79. """Start the 3-instance loadgen cluster."""
  80. print("\nStep 2 — starting 3-instance loadgen cluster (5k/s)")
  81. proc = subprocess.Popen(
  82. ["docker", "compose", "--profile", "loadgen-m10", "up", "-d"],
  83. stdout=subprocess.DEVNULL,
  84. stderr=subprocess.DEVNULL,
  85. )
  86. code = proc.wait()
  87. if code != 0:
  88. fail_(f"docker compose --profile loadgen-m10 up -d failed (exit {code})")
  89. pass_("3 loadgen-http instances started")
  90. # Wait for ramp-up
  91. print(f" waiting {SOAK_RAMP_SEC}s for ramp-up to complete...", flush=True)
  92. time.sleep(SOAK_RAMP_SEC)
  93. pass_(f"ramp-up complete — now targeting {CLUSTER_TARGET}/s")
  94. return proc
  95. def step2_monitor_soak() -> dict:
  96. """
  97. Monitor the soak: sample p99 + rate + DLQ every 30s.
  98. Fails fast on any breach.
  99. Returns dict of samples for the log.
  100. """
  101. print(f"\n Monitoring soak for {SOAK_DURATION_MIN} minutes...")
  102. samples = []
  103. start = time.time()
  104. deadline = start + SOAK_DURATION_MIN * 60
  105. sample_interval = 30 # seconds between samples
  106. while time.time() < deadline:
  107. time.sleep(sample_interval)
  108. elapsed = int(time.time() - start) // 60
  109. try:
  110. rate = lib.assert_rate_near(CLUSTER_TARGET, tolerance=RATE_TOLERANCE, window_seconds=30)
  111. p99 = lib.assert_ingest_p99_under(P99_THRESHOLD, window_seconds=60)
  112. dlq = lib.assert_dlq_count_equals(DLQ_EXPECTED, window_seconds=60)
  113. print(f" [{elapsed}m] rate={rate:.0f}/s p99={p99:.3f}s dlq={dlq}")
  114. samples.append({"elapsed_min": elapsed, "rate": rate, "p99": p99, "dlq": dlq})
  115. except AssertionError as e:
  116. fail_(f"soak breach at {elapsed}m: {e}")
  117. return samples
  118. def step3_runaway_test() -> None:
  119. """
  120. Runaway-source fault injection.
  121. - Start a 4th loadgen instance firing at 10× per-source cap for the same company.
  122. - Verify p99 for the OTHER sources (acme-002, acme-003) stays under threshold.
  123. - The runaway (acme-001) can be anything.
  124. """
  125. print(f"\nStep 3 — runaway-source fault injection ({RUNAWAY_DURATION_SEC}s)")
  126. print(" (not yet implemented — requires --rate override on loadgen-http-1)")
  127. warn_(f"runaway-source test skipped (W4 enhancement pending)")
  128. # TODO: spin up a 4th instance at 10× cap targeting acme-001
  129. # Expected: per_source_p99("acme-002") < 5s and per_source_p99("acme-003") < 5s
  130. def step4_dlq_invariant(samples: list[dict]) -> int:
  131. """Assert zero DLQ rows for the entire soak window."""
  132. print("\nStep 4 — DLQ invariant check")
  133. # The soak samples already checked DLQ delta, but do a final absolute check.
  134. # Query the full soak window.
  135. soak_seconds = SOAK_DURATION_MIN * 60
  136. dlq = lib.assert_dlq_count_equals(DLQ_EXPECTED, window_seconds=soak_seconds)
  137. pass_(f"DLQ count over {SOAK_DURATION_MIN}min soak: {dlq} (expected 0)")
  138. return dlq
  139. def step5_teardown(proc) -> None:
  140. """Bring down the loadgen cluster."""
  141. print("\nStep 5 — teardown")
  142. r = subprocess.run(
  143. ["docker", "compose", "--profile", "loadgen-m10", "down", "-v"],
  144. capture_output=True,
  145. )
  146. if r.returncode == 0:
  147. pass_("loadgen cluster torn down")
  148. else:
  149. warn_(f"teardown returned {r.returncode}: {r.stderr.decode().strip()}")
  150. def print_summary(samples: list[dict], dlq_final: int) -> None:
  151. """Print a summary table of the soak run."""
  152. print("\n=== M10 Soak Summary ===")
  153. print(f"Duration: {SOAK_DURATION_MIN} min")
  154. print(f"Target: {CLUSTER_TARGET}/s (±{RATE_TOLERANCE*100:.0f}%)")
  155. print(f"p99 threshold: {P99_THRESHOLD}s")
  156. print()
  157. if samples:
  158. print(f"{'Time':>6} {'Rate/s':>8} {'p99(s)':>7} {'DLQ':>4}")
  159. print("-" * 35)
  160. for s in samples:
  161. print(f"{s['elapsed_min']:>5}m {s['rate']:>8.0f} {s['p99']:>7.3f} {s['dlq']:>4}")
  162. print()
  163. print(f"Final DLQ count: {dlq_final} (expected 0)")
  164. print()
  165. print("🎉 M10 smoke: all checks complete.")
  166. def main():
  167. print(f"M10 Soak Test — target {CLUSTER_TARGET}/s for {SOAK_DURATION_MIN} min")
  168. print(f"p99 threshold: {P99_THRESHOLD}s | DLQ expected: {DLQ_EXPECTED}")
  169. print()
  170. try:
  171. step1_preflight()
  172. proc = step2_start_loadgen()
  173. samples = step2_monitor_soak()
  174. step3_runaway_test()
  175. dlq_final = step4_dlq_invariant(samples)
  176. step5_teardown(proc)
  177. print_summary(samples, dlq_final)
  178. except AssertionError as e:
  179. print(f"\n💥 M10 smoke FAILED: {e}")
  180. sys.exit(1)
  181. except KeyboardInterrupt:
  182. print("\nInterrupted.")
  183. sys.exit(1)
  184. if __name__ == "__main__":
  185. main()