from __future__ import annotations import json from pathlib import Path from uuid import uuid4 from fastapi.testclient import TestClient from sqlalchemy import create_engine, select from sqlalchemy.orm import Session import gateway.app as gateway_module from scripts import event_bus_smoke, track8_check from services.audit_service import app as audit_module from services.event_bus_service.app import app as event_bus_app from services.interaction_service.app import app as interaction_app from services.ivr_service.app import app as ivr_app from services.recording_service.app import app as recording_app from services.reporting_service import app as reporting_module from services.shared.event_bus import append_outbox_event, build_envelope, mark_outbox_failed from services.shared.sql_models import ( AuditEventRow, Base, EventInboxRow, EventOutboxRow, ReportingEventLogRow, ReportingEventRow, ) from services.supervisor_service.app import app as supervisor_app def _admin_headers() -> dict[str, str]: return {"X-User": "admin", "X-Role": "admin"} def _supervisor_headers() -> dict[str, str]: return {"X-User": "supervisor", "X-Role": "supervisor"} def _operator_headers() -> dict[str, str]: return {"X-User": "operator", "X-Role": "operator"} def _write_audio_fixture(path: Path) -> Path: path.parent.mkdir(parents=True, exist_ok=True) path.write_bytes(b"RIFF" + b"\x00" * 32) return path def _build_valid_flow(queue_id: str, resolved_queue_id: str) -> dict: return { "name": f"Track8 Flow {queue_id}", "description": "Track8 test flow", "queue_id": queue_id, "entry_node_id": "root", "flow_json": { "nodes": [ { "node_id": "root", "prompt_text": "Press 1 for support", "prompt_audio_key": "ivr/track8-root", "is_terminal": False, "invalid_target_node_id": "root", "options": [{"digit": "1", "target_node_id": "support"}], }, { "node_id": "support", "prompt_text": "Routing to support", "prompt_audio_key": "ivr/track8-support", "is_terminal": True, "outcome_code": "support_route", "resolved_queue_id": resolved_queue_id, "resolved_queue_code": "support", "options": [], }, ] }, "is_active": True, } def test_append_outbox_event_persists_full_envelope(): from services.shared.db import get_session session = get_session() try: entity_id = f"int_{uuid4().hex[:12]}" row = append_outbox_event( session, event_type="interaction.created", producer_service="interaction-service", entity_type="interaction", entity_id=entity_id, payload={"interaction_id": entity_id, "queue_id": "q_track8"}, correlation_id="corr_track8", ) session.commit() stored = session.execute( select(EventOutboxRow).where(EventOutboxRow.event_id == row.event_id) ).scalar_one() envelope = json.loads(stored.payload_json) assert stored.status == "pending" assert stored.routing_key == "interaction.created" assert envelope["event_id"] == stored.event_id assert envelope["correlation_id"] == "corr_track8" assert envelope["payload"]["interaction_id"] == entity_id finally: session.close() def test_interaction_supervisor_recording_and_ivr_write_outbox(tmp_path: Path, monkeypatch) -> None: monkeypatch.setenv("EVENT_BUS_ENABLED", "1") interaction_client = TestClient(interaction_app) supervisor_client = TestClient(supervisor_app) recording_client = TestClient(recording_app) ivr_client = TestClient(ivr_app) interaction = interaction_client.post( "/interactions", json={ "channel": "voice", "subject": "Track8 Outbox Interaction", "queue_id": f"q_{uuid4().hex[:8]}", "priority": 3, }, headers={"X-User": "operator", "X-Role": "operator"}, ) assert interaction.status_code == 200 interaction_id = interaction.json()["interaction_id"] assigned = interaction_client.patch( f"/interactions/{interaction_id}/assign", json={"assignee": "operator_a"}, headers=_admin_headers(), ) assert assigned.status_code == 200 escalated = interaction_client.post( f"/interactions/{interaction_id}/escalate", json={"target_queue_id": "line2"}, headers=_supervisor_headers(), ) assert escalated.status_code == 200 closed = interaction_client.patch( f"/interactions/{interaction_id}/status", json={"status": "closed"}, headers={"X-User": "operator", "X-Role": "operator"}, ) assert closed.status_code == 200 agent_state = supervisor_client.post( "/supervisor/agent-states", json={"agent_id": f"agent_{uuid4().hex[:6]}", "state": "READY", "queue_id": "line2"}, headers=_supervisor_headers(), ) assert agent_state.status_code == 200 monkeypatch.setenv("CC_RECORDINGS_DIR", str(tmp_path / "recordings")) source_path = _write_audio_fixture(tmp_path / "fixtures" / "track8.wav") registered = recording_client.post( "/recordings/register", headers=_supervisor_headers(), json={ "call_id": f"call_{uuid4().hex[:8]}", "interaction_id": interaction_id, "source_path": str(source_path), "file_name": "track8.wav", "mime_type": "audio/wav", "duration_seconds": 2, }, ) assert registered.status_code == 200 queue_id = f"ivr_q_{uuid4().hex[:6]}" flow = ivr_client.post("/ivr/flows", headers=_admin_headers(), json=_build_valid_flow(queue_id, "support_line")) assert flow.status_code == 200 started = ivr_client.post( "/ivr/sessions/start", headers=_supervisor_headers(), json={"call_id": f"call_{uuid4().hex[:6]}", "queue_id": queue_id, "interaction_id": interaction_id}, ) assert started.status_code == 200 session_id = started.json()["session"]["session_id"] completed = ivr_client.post( f"/ivr/sessions/{session_id}/dtmf", headers=_supervisor_headers(), json={"digit": "1"}, ) assert completed.status_code == 200 assert completed.json()["completed"] is True from services.shared.db import get_session session = get_session() try: event_types = { row.event_type for row in session.execute( select(EventOutboxRow).where( EventOutboxRow.event_type.in_( [ "interaction.created", "interaction.assigned", "interaction.escalated", "interaction.closed", "agent.state.changed", "call.recording.ready", "ivr.completed", ] ) ) ).scalars() } assert "interaction.created" in event_types assert "interaction.assigned" in event_types assert "interaction.escalated" in event_types assert "interaction.closed" in event_types assert "agent.state.changed" in event_types assert "call.recording.ready" in event_types assert "ivr.completed" in event_types finally: session.close() def test_audit_and_reporting_consumers_are_idempotent() -> None: envelope = build_envelope( event_type="ivr.completed", producer="ivr-service", entity_type="ivr_session", entity_id=f"ivs_{uuid4().hex[:12]}", payload={ "interaction_id": f"int_{uuid4().hex[:12]}", "queue_id": "line2", "resolved_queue_id": "support_line", "channel": "voice", "outcome_code": "support_route", }, ).model_dump() audit_module._handle_event(envelope) audit_module._handle_event(envelope) reporting_module._handle_event(envelope) reporting_module._handle_event(envelope) from services.shared.db import get_session session = get_session() try: audit_rows = session.execute( select(AuditEventRow).where(AuditEventRow.action == "ivr.completed") ).scalars().all() audit_inbox = session.execute( select(EventInboxRow).where( EventInboxRow.consumer_name == "audit-service", EventInboxRow.event_id == envelope["event_id"], ) ).scalars().all() reporting_log = session.execute( select(ReportingEventLogRow).where(ReportingEventLogRow.event_id == envelope["event_id"]) ).scalars().all() reporting_events = session.execute( select(ReportingEventRow).order_by(ReportingEventRow.id.desc()) ).scalars().all() reporting_inbox = session.execute( select(EventInboxRow).where( EventInboxRow.consumer_name == "reporting-service", EventInboxRow.event_id == envelope["event_id"], ) ).scalars().all() assert len(audit_inbox) == 1 assert len(reporting_inbox) == 1 assert len([row for row in audit_rows if envelope["event_id"] in row.metadata_json]) == 1 assert len(reporting_log) == 1 assert any(row.queue_id == "support_line" for row in reporting_events) finally: session.close() def test_event_bus_service_retry_endpoint_and_gateway_contracts(monkeypatch) -> None: monkeypatch.setenv("EVENT_BUS_ENABLED", "0") from services.shared.db import get_session session = get_session() try: row = append_outbox_event( session, event_type="interaction.created", producer_service="interaction-service", entity_type="interaction", entity_id=f"int_{uuid4().hex[:12]}", payload={"interaction_id": f"int_{uuid4().hex[:12]}"}, ) mark_outbox_failed(session, row, "broker down") event_id = row.event_id session.commit() finally: session.close() client = TestClient(event_bus_app) retried = client.post(f"/bus/outbox/{event_id}/retry", headers=_admin_headers()) assert retried.status_code == 200 assert retried.json()["status"] == "pending" denied = client.get("/bus/outbox", headers=_operator_headers()) assert denied.status_code == 403 gateway_client = TestClient(gateway_module.app) contracts = gateway_client.get("/contracts") assert contracts.status_code == 200 body = contracts.json() assert "stage_8" in body assert "/bus/outbox" in body["stage_8"] def test_event_bus_smoke_and_track8_check_with_temp_database(tmp_path: Path, monkeypatch) -> None: db_path = tmp_path / "track8.db" engine = create_engine(f"sqlite:///{db_path.as_posix()}", future=True) Base.metadata.create_all(engine) class DummyResponse: def __init__(self, status_code: int, payload: dict | None = None): self.status_code = status_code self._payload = payload or {} def raise_for_status(self) -> None: if self.status_code >= 400: raise RuntimeError("HTTP error") def json(self) -> dict: return self._payload class DummyClient: def __init__(self, *args, **kwargs): self.base_url = kwargs.get("base_url") def __enter__(self): return self def __exit__(self, exc_type, exc, tb): return None def get(self, path: str, headers: dict | None = None, timeout: int | None = None): if path == "/proxy/event-bus/health": return DummyResponse(200, {"status": "ok"}) raise AssertionError(f"Unexpected GET path {path}") def post(self, path: str, headers: dict | None = None, json: dict | None = None, timeout: int | None = None): if path == "/proxy/interaction/interactions": interaction_id = f"int_{uuid4().hex[:12]}" with Session(engine) as session: outbox = EventOutboxRow( event_id=f"evt_{uuid4().hex[:12]}", event_type="interaction.created", event_version=1, producer_service="interaction-service", entity_type="interaction", entity_id=interaction_id, correlation_id=None, routing_key="interaction.created", payload_json=json and __import__("json").dumps( { "event_id": f"evt_payload_{uuid4().hex[:12]}", "event_type": "interaction.created", "event_version": 1, "occurred_at": "2026-03-02T12:00:00+00:00", "producer": "interaction-service", "entity_type": "interaction", "entity_id": interaction_id, "routing_key": "interaction.created", "payload": {"interaction_id": interaction_id}, } ) or "{}", status="published", attempt_count=0, last_error=None, available_at="2026-03-02T12:00:00+00:00", published_at="2026-03-02T12:00:01+00:00", created_at="2026-03-02T12:00:00+00:00", updated_at="2026-03-02T12:00:01+00:00", ) session.add(outbox) session.add( EventInboxRow( consumer_name="audit-service", event_id=outbox.event_id, event_type=outbox.event_type, processed_at="2026-03-02T12:00:02+00:00", status="processed", notes=None, ) ) session.add( EventInboxRow( consumer_name="reporting-service", event_id=outbox.event_id, event_type=outbox.event_type, processed_at="2026-03-02T12:00:02+00:00", status="processed", notes=None, ) ) session.commit() return DummyResponse(200, {"interaction_id": interaction_id}) raise AssertionError(f"Unexpected POST path {path}") monkeypatch.setattr(event_bus_smoke, "_engine_for", lambda database_url=None: engine) monkeypatch.setattr(event_bus_smoke.httpx, "Client", DummyClient) smoke = event_bus_smoke.run_smoke_check(base_url="http://example.local", auth_mode="legacy_headers") assert smoke["passed"] is True issues = track8_check.evaluate_track8( database_url=f"sqlite:///{db_path.as_posix()}", require_bus_enabled=False, base_url=None, ) assert issues == []