from __future__ import annotations import argparse import sys import time from pathlib import Path import httpx from sqlalchemy import create_engine, select from sqlalchemy.orm import Session ROOT = Path(__file__).resolve().parents[1] if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) from services.shared.core import utc_now_iso from services.shared.db import DATABASE_URL, _normalize_database_url from services.shared.sql_models import EventInboxRow, EventOutboxRow def _engine_for(database_url: str | None): url = _normalize_database_url(database_url or DATABASE_URL) return create_engine(url, future=True) def _auth_headers(client: httpx.Client, auth_mode: str) -> dict[str, str]: if auth_mode == "legacy_headers": return {"X-User": "admin", "X-Role": "admin"} response = client.post( "/proxy/auth/auth/login", json={"username": "admin", "password": "admin123"}, timeout=10, ) response.raise_for_status() token = response.json()["access_token"] return {"Authorization": f"Bearer {token}"} def run_smoke_check( *, base_url: str, database_url: str | None = None, auth_mode: str = "legacy_headers", timeout_seconds: int = 15, ) -> dict: engine = _engine_for(database_url) with httpx.Client(base_url=base_url, timeout=10) as client: headers = _auth_headers(client, auth_mode) event_bus_health = client.get("/proxy/event-bus/health", headers=headers) if event_bus_health.status_code != 200: raise RuntimeError("event-bus-service health check failed") created = client.post( "/proxy/interaction/interactions", headers=headers, json={ "channel": "voice", "subject": f"Track8 smoke {utc_now_iso()}", "customer_id": None, "queue_id": "track8_smoke_queue", "priority": 3, }, ) created.raise_for_status() interaction_id = created.json()["interaction_id"] deadline = time.time() + timeout_seconds found: dict[str, str] = {} with Session(engine) as session: while time.time() < deadline: outbox_row = session.execute( select(EventOutboxRow).where( EventOutboxRow.entity_id == interaction_id, EventOutboxRow.event_type == "interaction.created", ) ).scalar_one_or_none() if outbox_row: found["event_id"] = outbox_row.event_id found["outbox_status"] = outbox_row.status audit_seen = session.execute( select(EventInboxRow).where( EventInboxRow.consumer_name == "audit-service", EventInboxRow.event_id == outbox_row.event_id, ) ).scalar_one_or_none() reporting_seen = session.execute( select(EventInboxRow).where( EventInboxRow.consumer_name == "reporting-service", EventInboxRow.event_id == outbox_row.event_id, ) ).scalar_one_or_none() if outbox_row.status == "published" and audit_seen and reporting_seen: return { "passed": True, "interaction_id": interaction_id, "event_id": outbox_row.event_id, "outbox_status": outbox_row.status, } session.expire_all() time.sleep(0.5) return { "passed": False, "interaction_id": interaction_id, "event_id": found.get("event_id"), "outbox_status": found.get("outbox_status"), } def main() -> None: parser = argparse.ArgumentParser(description="Run a lightweight Track 8 event-bus smoke check.") parser.add_argument("--base-url", default="http://localhost:8080") parser.add_argument("--database-url", default=None) parser.add_argument("--auth-mode", choices=["legacy_headers", "bearer"], default="legacy_headers") parser.add_argument("--timeout-seconds", type=int, default=15) args = parser.parse_args() result = run_smoke_check( base_url=args.base_url, database_url=args.database_url, auth_mode=args.auth_mode, timeout_seconds=args.timeout_seconds, ) print(result) if not result["passed"]: raise SystemExit(1) if __name__ == "__main__": main()