m5_smoke.sh 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395
  1. #!/usr/bin/env bash
  2. # Live M5 smoke test. Run from repo root:
  3. # bash scripts/m5_smoke.sh
  4. #
  5. # Walks through the 6 scenarios in M5_VERIFICATION.md:
  6. #
  7. # Step 2 — happy path: 1 alert via WS → 2 deliveries (Alice fcm + telegram)
  8. # Step 3 — 5 alerts, mixed severity + 30% dedupe → 10 deliveries
  9. # Step 4 — bad API key (correct format, wrong secret) → 0 deliveries,
  10. # ws_messages_total{result="bad_signature"} ticks
  11. # Step 5 — per-IP concurrency cap (35 conns, cap=32) → 32 ok + 3 rejected,
  12. # connection_rejected_total{transport="ws"} = 3
  13. # Step 6 — live tail: subscribe then send → 1 frame arrives,
  14. # tail_subscribers = 1 while connected
  15. # Step 7 — live tail company filter: subscribe to globex-002 only →
  16. # acme alerts filtered out, globex alert passes
  17. #
  18. # The script assumes the loadgen-ws binary is built at /tmp/loadgen-ws
  19. # (run `cd loadgen && go build -o /tmp/loadgen-ws ./cmd/ws`). The
  20. # per-IP cap test binary and tail driver are auto-built on first run;
  21. # they live in /tmp and are reused on subsequent runs.
  22. #
  23. # Exit code is the number of failed checks.
  24. set -e
  25. cd "$(dirname "$0")/.."
  26. PG="docker exec -i broad-announce-postgres-1 psql -U ba -d ba -A -t"
  27. INGESTD_METRICS=http://localhost:8800/metrics
  28. SRC_ACME=acme-001:prom-prod:s3cret-acme
  29. SRC_GLOBEX=globex-002:grafana:s3cret-globex
  30. WS_INGEST=ws://localhost:8800/v1/ingest/ws
  31. WS_TAIL=ws://localhost:8800/v1/tail/ws
  32. TAIL_TOKEN=tail-dev-token-please-change-in-prod
  33. fails=0
  34. pass() { echo " ✅ $*"; }
  35. fail() { echo " ❌ $*"; fails=$((fails+1)); }
  36. reset_state() {
  37. $PG -c "UPDATE individuals SET telegram_chat_id = NULL, telegram_user_id = NULL, telegram_invite_code = 'acme-bob-002' WHERE id = 'ind-acme-002';" >/dev/null
  38. $PG -c "UPDATE individuals SET telegram_chat_id = NULL, telegram_user_id = NULL WHERE id = 'ind-acme-003';" >/dev/null
  39. $PG -c "UPDATE subscriptions SET min_severity = 'critical' WHERE individual_id = 'ind-acme-002' AND source_id = 'prom-prod';" >/dev/null
  40. $PG -c "TRUNCATE deliveries;" >/dev/null
  41. curl -sS -X POST http://localhost:8830/admin/reset >/dev/null
  42. docker compose restart telegramd >/dev/null
  43. for i in 1 2 3 4 5 6 7 8 9 10; do
  44. if curl -sS http://localhost:8822/health 2>/dev/null | grep -q '"status":"ok"'; then
  45. sleep 1
  46. break
  47. fi
  48. sleep 1
  49. done
  50. }
  51. ws_counter() {
  52. # $1 = result label
  53. curl -sS "$INGESTD_METRICS" 2>/dev/null | \
  54. grep -E "^ba_ingestd_ws_messages_total\{result=\"$1\"" | \
  55. awk '{print $NF}' | awk -F. '{print $1+0; exit}' | head -1
  56. }
  57. ws_conn() {
  58. # $1 = state label
  59. curl -sS "$INGESTD_METRICS" 2>/dev/null | \
  60. grep -E "^ba_ingestd_ws_connections_total\{state=\"$1\"" | \
  61. awk '{print $NF}' | awk -F. '{print $1+0; exit}' | head -1
  62. }
  63. ws_rej() {
  64. curl -sS "$INGESTD_METRICS" 2>/dev/null | \
  65. grep -E '^ba_ingestd_connection_rejected_total\{[^}]*transport="ws"[^}]*\}' | \
  66. awk '{print $NF}' | awk -F. '{print $1+0; exit}' | head -1
  67. }
  68. tail_gauge() {
  69. curl -sS "$INGESTD_METRICS" 2>/dev/null | \
  70. grep -E '^ba_ingestd_tail_subscribers\{' | \
  71. awk '{print $NF}' | awk -F. '{print $1+0; exit}' | head -1
  72. }
  73. # ── Setup: build loadgen-ws + per-IP-cap test binary + tail driver ──
  74. mkdir -p /tmp/m5_smoke
  75. if [[ ! -x /tmp/loadgen-ws ]]; then
  76. echo "▸ Building /tmp/loadgen-ws"
  77. (cd loadgen && CGO_ENABLED=0 go build -o /tmp/loadgen-ws ./cmd/ws)
  78. fi
  79. if [[ ! -x /tmp/m5-tail-test ]]; then
  80. echo "▸ Building /tmp/m5-tail-test"
  81. CGO_ENABLED=0 go build -o /tmp/m5-tail-test ./loadgen/cmd/m5drivers/tail
  82. fi
  83. if [[ ! -x /tmp/m5-perip-test ]]; then
  84. echo "▸ Building /tmp/m5-perip-test"
  85. cat > /tmp/m5_smoke/perip_test.go <<'GO'
  86. package main
  87. import (
  88. "flag"
  89. "fmt"
  90. "sync"
  91. "sync/atomic"
  92. "git3.techno-world.net/lrosales/broad-announce/internal/wsclient"
  93. )
  94. func main() {
  95. wsURL := flag.String("ws", "ws://localhost:8800/v1/ingest/ws", "ws endpoint")
  96. apiKey := flag.String("key", "acme-001:prom-prod:s3cret-acme", "api key")
  97. cap := flag.Int("cap", 35, "number of concurrent WS conns to attempt (cap=32 should reject 3)")
  98. flag.Parse()
  99. var ok, rejected atomic.Int64
  100. var wg sync.WaitGroup
  101. for i := 0; i < *cap; i++ {
  102. wg.Add(1)
  103. go func() {
  104. defer wg.Done()
  105. c, err := wsclient.Connect(wsclient.Config{URL: *wsURL, APIKey: *apiKey, DialTimeout: 5e9})
  106. if err != nil {
  107. rejected.Add(1)
  108. return
  109. }
  110. ok.Add(1)
  111. defer c.Close()
  112. }()
  113. }
  114. wg.Wait()
  115. fmt.Printf("ok=%d rejected=%d (cap=32 → ok<=32, rejected>=3 if cap=35)\n", ok.Load(), rejected.Load())
  116. }
  117. GO
  118. mkdir -p ./scripts/m5_smoke_tmp
  119. cp /tmp/m5_smoke/perip_test.go ./scripts/m5_smoke_tmp/main.go
  120. CGO_ENABLED=0 go build -o /tmp/m5-perip-test ./scripts/m5_smoke_tmp
  121. rm -rf ./scripts/m5_smoke_tmp
  122. fi
  123. # ─────────────────────────────────────────────────────────────────
  124. echo "── M5 smoke — WebSocket ingest + live tail + per-IP cap ──"
  125. echo ""
  126. reset_state
  127. # ── Step 2: happy path ────────────────────────────────────────
  128. echo "── Step 2: 1 alert via WS → 2 deliveries (Alice fcm + telegram) ──"
  129. acme_recv_before=$(ws_counter received)
  130. acme_acc_before=$(ws_counter accepted)
  131. deliveries_before=$($PG -c "SELECT count(*) FROM deliveries WHERE channel IN ('fcm','telegram');" | tr -d ' \n')
  132. acme_recv_before=${acme_recv_before:-0}
  133. acme_acc_before=${acme_acc_before:-0}
  134. deliveries_before=${deliveries_before:-0}
  135. /tmp/loadgen-ws --target "$WS_INGEST" --api-key "$SRC_ACME" --count 1 --rate 1 2>&1 | tail -3
  136. sleep 2
  137. acme_recv_after=$(ws_counter received)
  138. acme_acc_after=$(ws_counter accepted)
  139. deliveries_after=$($PG -c "SELECT count(*) FROM deliveries WHERE channel IN ('fcm','telegram');" | tr -d ' \n')
  140. acme_recv_after=${acme_recv_after:-0}
  141. acme_acc_after=${acme_acc_after:-0}
  142. deliveries_after=${deliveries_after:-0}
  143. recv_delta=$((acme_recv_after - acme_recv_before))
  144. acc_delta=$((acme_acc_after - acme_acc_before))
  145. del_delta=$((deliveries_after - deliveries_before))
  146. if [[ $recv_delta -ge 1 ]]; then
  147. pass "ws_messages_total{result=\"received\"} +$recv_delta"
  148. else
  149. fail "ws_messages_total{result=\"received\"} delta was $recv_delta (expected ≥1)"
  150. fi
  151. if [[ $acc_delta -ge 1 ]]; then
  152. pass "ws_messages_total{result=\"accepted\"} +$acc_delta"
  153. else
  154. fail "ws_messages_total{result=\"accepted\"} delta was $acc_delta (expected ≥1)"
  155. fi
  156. if [[ $del_delta -ge 2 ]]; then
  157. pass "deliveries +$del_delta (Alice fcm + telegram)"
  158. else
  159. fail "deliveries delta was $del_delta (expected ≥2)"
  160. fi
  161. # ── Step 3: 5 alerts, mixed severity + 30% dedupe ───────────
  162. echo ""
  163. echo "── Step 3: 5 alerts, mixed severity + 30% dedupe → 10 deliveries ──"
  164. reset_state
  165. acme_recv_before=$(ws_counter received)
  166. acme_acc_before=$(ws_counter accepted)
  167. deliveries_before=$($PG -c "SELECT count(*) FROM deliveries WHERE channel IN ('fcm','telegram');" | tr -d ' \n')
  168. # Send 5 alerts at a steady rate. Default profile is
  169. # normal which uses 30% dedupe (i.e. ~2 of 5 alerts share
  170. # a dedupe_key), so we expect 3 new accepts and 2 dupes.
  171. /tmp/loadgen-ws --target "$WS_INGEST" --api-key "$SRC_ACME" --count 5 --rate 5 2>&1 | tail -3
  172. sleep 3
  173. acme_recv_after=$(ws_counter received)
  174. acme_acc_after=$(ws_counter accepted)
  175. deliveries_after=$($PG -c "SELECT count(*) FROM deliveries WHERE channel IN ('fcm','telegram');" | tr -d ' \n')
  176. recv_delta=$((acme_recv_after - acme_recv_before))
  177. acc_delta=$((acme_acc_after - acme_acc_before))
  178. del_delta=$((deliveries_after - deliveries_before))
  179. if [[ $recv_delta -eq 5 ]]; then
  180. pass "ws_messages_total{result=\"received\"} +$recv_delta"
  181. else
  182. fail "ws_messages_total{result=\"received\"} delta was $recv_delta (expected 5)"
  183. fi
  184. if [[ $acc_delta -ge 3 ]]; then
  185. pass "ws_messages_total{result=\"accepted\"} +$acc_delta (≥3 of 5 are new; rest deduped)"
  186. else
  187. fail "ws_messages_total{result=\"accepted\"} delta was $acc_delta (expected ≥3)"
  188. fi
  189. if [[ $del_delta -ge 8 ]]; then
  190. pass "deliveries +$del_delta (≥1 per non-deduped alert × 2 channels)"
  191. else
  192. fail "deliveries delta was $del_delta (expected ≥8)"
  193. fi
  194. # ── Step 4: bad API key ─────────────────────────────────────
  195. echo ""
  196. echo "── Step 4: bad API key (wrong secret) → 0 deliveries, bad_signature counter ticks ──"
  197. reset_state
  198. bad_before=$(ws_counter bad_signature)
  199. deliveries_before=$($PG -c "SELECT count(*) FROM deliveries WHERE channel IN ('fcm','telegram');" | tr -d ' \n')
  200. # Use the right company:source (so auth passes) but wrong HMAC.
  201. cat > /tmp/m5_smoke/badkey.go <<'GO'
  202. package main
  203. import (
  204. "crypto/hmac"
  205. "crypto/sha256"
  206. "encoding/hex"
  207. "encoding/json"
  208. "fmt"
  209. "time"
  210. "github.com/gorilla/websocket"
  211. )
  212. func main() {
  213. conn, _, err := websocket.DefaultDialer.Dial("ws://localhost:8800/v1/ingest/ws", nil)
  214. if err != nil { fmt.Println("dial err:", err); return }
  215. defer conn.Close()
  216. type Auth struct {
  217. APIKey string `json:"api_key"`
  218. }
  219. if err := conn.WriteJSON(Auth{APIKey: "acme-001:prom-prod:s3cret-acme"}); err != nil {
  220. fmt.Println("auth err:", err); return
  221. }
  222. conn.SetReadDeadline(time.Now().Add(2*time.Second))
  223. _, authReply, err := conn.ReadMessage()
  224. if err != nil { fmt.Println("auth reply err:", err); return }
  225. fmt.Println("auth reply:", string(authReply))
  226. body := map[string]any{"company_id":"acme-001","source_id":"prom-prod","severity":"info","title":"smoke-badkey","labels":map[string]string{"instance":"smoke-1"}}
  227. bb, _ := json.Marshal(body)
  228. // Sign with a wrong key — should fail HMAC verify.
  229. mac := hmac.New(sha256.New, []byte("WRONG-SECRET"))
  230. mac.Write(bb)
  231. hexMac := hex.EncodeToString(mac.Sum(nil))
  232. sig := "t=1700000000,v1=" + hexMac
  233. env := map[string]any{
  234. "alert": json.RawMessage(bb),
  235. "auth": sig,
  236. }
  237. if err := conn.WriteJSON(env); err != nil { fmt.Println("write err:", err); return }
  238. conn.SetReadDeadline(time.Now().Add(2*time.Second))
  239. _, ack2, err := conn.ReadMessage()
  240. if err != nil { fmt.Println("ack2 err:", err); return }
  241. fmt.Println("ack2:", string(ack2))
  242. }
  243. GO
  244. mkdir -p ./scripts/m5_smoke_tmp
  245. cp /tmp/m5_smoke/badkey.go ./scripts/m5_smoke_tmp/main.go
  246. CGO_ENABLED=0 go build -o /tmp/m5-badkey ./scripts/m5_smoke_tmp
  247. rm -rf ./scripts/m5_smoke_tmp
  248. /tmp/m5-badkey 2>&1 | tail -5
  249. sleep 2
  250. bad_after=$(ws_counter bad_signature)
  251. deliveries_after=$($PG -c "SELECT count(*) FROM deliveries WHERE channel IN ('fcm','telegram');" | tr -d ' \n')
  252. bad_delta=$((bad_after - bad_before))
  253. del_delta=$((deliveries_after - deliveries_before))
  254. if [[ $bad_delta -ge 1 ]]; then
  255. pass "ws_messages_total{result=\"bad_signature\"} +$bad_delta"
  256. else
  257. fail "ws_messages_total{result=\"bad_signature\"} delta was $bad_delta (expected ≥1)"
  258. fi
  259. if [[ $del_delta -eq 0 ]]; then
  260. pass "no deliveries created"
  261. else
  262. fail "deliveries delta was $del_delta (expected 0)"
  263. fi
  264. # ── Step 5: per-IP concurrency cap ─────────────────────────
  265. echo ""
  266. echo "── Step 5: per-IP cap (35 conns, cap=32) → 32 ok + 3 rejected ──"
  267. # Drain any residual per-IP count from prior tests. The per-IP
  268. # counter decrements on TCP close; some kernels buffer the
  269. # FIN/ACK and we can race the next test. A 5-second wait is
  270. # enough for the janitor to forget any stragglers (idle > 1s
  271. # in the test) and for the kernel to reap closed sockets.
  272. sleep 5
  273. rej_before=$(ws_rej)
  274. perip_out=$(/tmp/m5-perip-test --cap 35 2>&1 | tail -1)
  275. echo " perip_test: $perip_out"
  276. sleep 1
  277. rej_after=$(ws_rej)
  278. rej_delta=$((rej_after - rej_before))
  279. if echo "$perip_out" | grep -qE "ok=3[0-2] rejected=[2-3]"; then
  280. pass "perip_test reported 30-32 ok, 2-3 rejected (cap=32 held)"
  281. else
  282. fail "perip_test output unexpected: $perip_out"
  283. fi
  284. if [[ $rej_delta -ge 2 && $rej_delta -le 3 ]]; then
  285. pass "connection_rejected_total{transport=\"ws\"} +$rej_delta"
  286. else
  287. fail "connection_rejected_total{transport=\"ws\"} delta was $rej_delta (expected 2-3)"
  288. fi
  289. sleep 1
  290. # ── Step 6: live tail, no filter ───────────────────────────
  291. echo ""
  292. echo "── Step 6: live tail (no filter) — subscribe then send 1 alert → 1 frame arrives ──"
  293. tail_before=$(tail_gauge)
  294. /tmp/m5-tail-test --target "$WS_TAIL" --token "$TAIL_TOKEN" --duration 6s 2>&1 > /tmp/m5_tail1.log &
  295. TAIL_PID=$!
  296. sleep 2 # let the tail subscribe
  297. tail_mid=$(tail_gauge)
  298. if [[ $tail_mid -ge 1 ]]; then
  299. pass "tail_subscribers = $tail_mid while connected"
  300. else
  301. fail "tail_subscribers was $tail_mid (expected ≥1)"
  302. fi
  303. /tmp/loadgen-ws --target "$WS_INGEST" --api-key "$SRC_ACME" --count 1 --rate 1 2>&1 | tail -1
  304. wait $TAIL_PID 2>/dev/null
  305. if grep -q "^FRAME:" /tmp/m5_tail1.log; then
  306. pass "tail received at least 1 frame"
  307. else
  308. fail "tail received 0 frames (log: $(cat /tmp/m5_tail1.log))"
  309. fi
  310. tail_after=$(tail_gauge)
  311. if [[ $tail_after -eq 0 ]]; then
  312. pass "tail_subscribers back to 0 after disconnect"
  313. else
  314. fail "tail_subscribers was $tail_after after disconnect (expected 0)"
  315. fi
  316. # ── Step 7: live tail, company filter ─────────────────────
  317. echo ""
  318. echo "── Step 7: live tail (company=globex-002) — acme filtered, globex passes ──"
  319. /tmp/m5-tail-test --target "$WS_TAIL" --token "$TAIL_TOKEN" --company globex-002 --duration 6s 2>&1 > /tmp/m5_tail2.log &
  320. TAIL_PID=$!
  321. sleep 2
  322. # Send acme — should NOT reach the tail
  323. /tmp/loadgen-ws --target "$WS_INGEST" --api-key "$SRC_ACME" --count 1 --rate 1 2>&1 | tail -1
  324. sleep 1
  325. # Send globex — should reach the tail
  326. /tmp/loadgen-ws --target "$WS_INGEST" --api-key "$SRC_GLOBEX" --count 1 --rate 1 2>&1 | tail -1
  327. wait $TAIL_PID 2>/dev/null
  328. # Count frames per company
  329. acme_frames=$(grep -c "company_id\":\"acme-001" /tmp/m5_tail2.log || true)
  330. globex_frames=$(grep -c "company_id\":\"globex-002" /tmp/m5_tail2.log || true)
  331. if [[ $acme_frames -eq 0 ]]; then
  332. pass "0 acme frames reached globex-only tail (filter works)"
  333. else
  334. fail "got $acme_frames acme frames on globex-only tail (expected 0)"
  335. fi
  336. if [[ $globex_frames -ge 1 ]]; then
  337. pass "$globex_frames globex frame(s) reached the tail"
  338. else
  339. fail "got $globex_frames globex frames (expected ≥1)"
  340. fi
  341. # ── Summary ─────────────────────────────────────────────────
  342. echo ""
  343. if [[ $fails -eq 0 ]]; then
  344. echo "🟢 M5 smoke PASS — all checks green"
  345. exit 0
  346. else
  347. echo "🔴 M5 smoke FAIL — $fails check(s) failed"
  348. exit $fails
  349. fi