Files
call-center/tests/test_track8_event_bus.py
T

420 lines
16 KiB
Python

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 == []