Files

131 lines
4.5 KiB
Python

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()