from __future__ import annotations import argparse import asyncio import os from pathlib import Path import shutil import subprocess import sys import time from typing import Any import httpx ROOT = Path(__file__).resolve().parents[1] DATA_ROOT = ROOT / ".data_smoke" SERVICE_SPECS = [ {"name": "auth", "module": "services.auth_service.app:app", "port": 18001}, {"name": "audit", "module": "services.audit_service.app:app", "port": 18002}, {"name": "customer", "module": "services.customer_service.app:app", "port": 18003}, {"name": "interaction", "module": "services.interaction_service.app:app", "port": 18004}, {"name": "routing", "module": "services.routing_service.app:app", "port": 18005}, {"name": "voice", "module": "services.voice_adapter_service.app:app", "port": 18006}, {"name": "telegram", "module": "services.telegram_adapter_service.app:app", "port": 18007}, {"name": "kb", "module": "services.kb_service.app:app", "port": 18008}, {"name": "reporting", "module": "services.reporting_service.app:app", "port": 18009}, {"name": "supervisor", "module": "services.supervisor_service.app:app", "port": 18010}, {"name": "gateway", "module": "gateway.app:app", "port": 18080}, ] async def wait_for_health(base_url: str, retries: int = 80, delay: float = 0.25) -> None: async with httpx.AsyncClient(timeout=2) as client: for _ in range(retries): try: response = await client.get(f"{base_url}/health") if response.status_code == 200: return except Exception: pass await asyncio.sleep(delay) raise RuntimeError(f"Service not healthy: {base_url}/health") def start_services(database_url: str | None, data_dir: Path) -> list[subprocess.Popen]: data_dir.mkdir(parents=True, exist_ok=True) processes: list[subprocess.Popen] = [] for spec in SERVICE_SPECS: env = os.environ.copy() env["CC_DATA_DIR"] = str(data_dir) if database_url: env["DATABASE_URL"] = database_url if spec["name"] == "gateway": env["AUTH_SERVICE_URL"] = "http://127.0.0.1:18001" env["AUDIT_SERVICE_URL"] = "http://127.0.0.1:18002" env["CUSTOMER_SERVICE_URL"] = "http://127.0.0.1:18003" env["INTERACTION_SERVICE_URL"] = "http://127.0.0.1:18004" env["ROUTING_SERVICE_URL"] = "http://127.0.0.1:18005" env["VOICE_ADAPTER_SERVICE_URL"] = "http://127.0.0.1:18006" env["TELEGRAM_ADAPTER_SERVICE_URL"] = "http://127.0.0.1:18007" env["KB_SERVICE_URL"] = "http://127.0.0.1:18008" env["REPORTING_SERVICE_URL"] = "http://127.0.0.1:18009" env["SUPERVISOR_SERVICE_URL"] = "http://127.0.0.1:18010" cmd = [ sys.executable, "-m", "uvicorn", spec["module"], "--host", "127.0.0.1", "--port", str(spec["port"]), ] proc = subprocess.Popen( cmd, cwd=str(ROOT), env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) processes.append(proc) return processes def stop_services(processes: list[subprocess.Popen]) -> None: for proc in processes: if proc.poll() is not None: continue try: proc.terminate() except Exception: pass for proc in processes: try: proc.wait(timeout=2) except Exception: try: proc.kill() except Exception: pass async def run_gate_checks() -> dict[str, Any]: for spec in SERVICE_SPECS: await wait_for_health(f"http://127.0.0.1:{spec['port']}") base = "http://127.0.0.1:18080" admin = {"X-User": "admin", "X-Role": "admin"} operator = {"X-User": "operator", "X-Role": "operator"} supervisor = {"X-User": "supervisor", "X-Role": "supervisor"} async with httpx.AsyncClient(base_url=base, timeout=10) as client: health = await client.get("/health") health.raise_for_status() login = await client.post("/proxy/auth/auth/login", json={"username": "admin", "password": "admin123"}) login.raise_for_status() queue = await client.post( "/proxy/routing/queues", headers=admin, json={ "name": "Main Queue", "description": "Smoke queue", "rules": [ {"channel": "voice", "priority": 3, "strategy": "round_robin", "sla_seconds": 30}, {"channel": "telegram", "priority": 3, "strategy": "round_robin", "sla_seconds": 45}, ], }, ) queue.raise_for_status() queue_id = queue.json()["queue_id"] customer = await client.post( "/proxy/customer/customers", json={ "display_name": "Gate User", "phones": ["+77010000011"], "preferred_phone": "+77010000011", "tags": ["smoke"], }, ) customer.raise_for_status() customer_id = customer.json()["customer_id"] interaction = await client.post( "/proxy/interaction/interactions", headers=operator, json={ "channel": "voice", "subject": "Gate voice check", "customer_id": customer_id, "queue_id": queue_id, "priority": 3, }, ) interaction.raise_for_status() interaction_id = interaction.json()["interaction_id"] assign = await client.patch( f"/proxy/interaction/interactions/{interaction_id}/assign", headers=supervisor, json={"assignee": "operator_a"}, ) assign.raise_for_status() escalate = await client.post( f"/proxy/interaction/interactions/{interaction_id}/escalate", headers=operator, json={"target_queue_id": "line2"}, ) escalate.raise_for_status() close = await client.patch( f"/proxy/interaction/interactions/{interaction_id}/status", headers=operator, json={"status": "closed"}, ) close.raise_for_status() voice_event = await client.post( "/proxy/voice/integrations/voice/events", headers=operator, json={ "event_type": "call.started", "call_id": "smoke_call_1", "interaction_id": interaction_id, "payload": {"source": "live_smoke"}, }, ) voice_event.raise_for_status() tg_event = await client.post( "/proxy/telegram/integrations/telegram/webhook", json={ "chat_id": "smoke_chat", "text": "Smoke message", "customer_external_id": None, "payload": {"source": "live_smoke"}, }, ) tg_event.raise_for_status() interactions = await client.get("/proxy/interaction/interactions") interactions.raise_for_status() rows = interactions.json() return { "gateway": health.json(), "queue_id": queue_id, "customer_id": customer_id, "interaction_id": interaction_id, "interactions_total": len(rows), } async def main() -> None: parser = argparse.ArgumentParser(description="Live Gate 1/2 smoke") parser.add_argument( "--database-url", default=os.getenv("DATABASE_URL", ""), help="Optional DB URL override for started services", ) args = parser.parse_args() db_url = args.database_url.strip() or None run_data_dir = DATA_ROOT / f"run_{int(time.time() * 1000)}" processes = start_services(db_url, run_data_dir) try: result = await run_gate_checks() print("Gate 1/2 live smoke passed") print(result) if db_url: print(f"DB mode: {db_url}") finally: stop_services(processes) try: shutil.rmtree(run_data_dir) except Exception: pass if __name__ == "__main__": asyncio.run(main())