import asyncio from inspect import getsource from pathlib import Path import pytest import app.analyzer as analyzer import app.communication_service as communication_service import app.persistence as persistence from app.action_mapper import map_action_decision from app.schemas import ActionDecision, AnalyzeRequest, UsageInfo def _decision(code: str): decision = ActionDecision(action_code=code, confidence=.91) return decision, map_action_decision(decision), UsageInfo(provider="test"), "test" @pytest.mark.parametrize("code,task_id,ignored", [ ("SEND_QUOTE", "task-quote", False), ("MARK_NO_INTEREST", "task-negative", False), ("IGNORE_SPAM", None, True), ("IGNORE_BOUNCE", None, True), ]) def test_inbound_history_precedes_decision_and_communication_is_enriched( monkeypatch, code, task_id, ignored, ): calls = [] monkeypatch.setattr(analyzer, "save_inbound_message", lambda **kw: (calls.append("message") or ("message-1", "cw-42"))) async def decide(_request): calls.append("decision") return _decision(code) monkeypatch.setattr(analyzer, "decide_action", decide) monkeypatch.setattr(analyzer, "save_action_run", lambda **kw: (calls.append("run") or ("run-1", kw["message_id"]))) monkeypatch.setattr(analyzer, "create_task_from_action_result", lambda **kw: (calls.append("task") or task_id)) monkeypatch.setattr(communication_service, "upsert_chatwoot_inbound_communication", lambda **kw: calls.append(("communication", kw))) enriched = [] monkeypatch.setattr(communication_service, "enrich_chatwoot_communication_from_task", lambda **kw: enriched.append(kw)) response = asyncio.run(analyzer.analyze( AnalyzeRequest(last_customer_message="hello", source="chatwoot", conversation_id="conv"), raw_event_id="raw-1", source_event_id="cw-42", )) assert calls.index("message") < calls.index("decision") assert response.message_id == "message-1" assert response.task_id == task_id assert enriched[0]["source_message_id"] == "cw-42" assert enriched[0]["classification"] == code assert enriched[0]["ignored"] is ignored assert enriched[0]["task_id"] == task_id def test_canonical_identity_and_retry_guards_are_present(): message_source = getsource(persistence.save_inbound_message) run_source = getsource(persistence.save_action_run) assert "source_event_id = source_event_id or raw.get" in message_source assert "source_system=:source_system AND source_event_id=:source_event_id" in message_source assert "pg_advisory_xact_lock" in getsource(persistence._lock_source_identity) assert "existing_run" in run_source assert "raw_event_id=CAST(:raw_event_id AS UUID)" in run_source migration = Path("migrations/010_chatwoot_ingestion_identity.sql").read_text() assert "UNIQUE INDEX" in migration assert "messages(source_system, source_event_id)" in migration def test_existing_opportunity_and_customer_links_enrich_communication(): source = getsource(communication_service.enrich_chatwoot_communication_from_task) assert "t.opportunity_id::text" in source assert "NULLIF(t.customer_id, '')" in source assert "opportunity_id=opportunity_id" in source def test_auditor_recognizes_all_legacy_message_relations_and_states(): source = Path("scripts/audit_chatwoot_ingestion_gap.py").read_text() assert "m.source_event_id = re.source_event_id" in source assert "m.id = re.message_id" in source assert "m.raw_event_id = re.id" in source for state in ("MESSAGE_PRESENT_NO_COMMUNICATION", "TRUE_PROCESSED_WITHOUT_MESSAGE", "PROCESSING_ERROR", "IGNORED", "OK"): assert state in source class _Scalar: rowcount = 2 def scalar(self): return 2 class _Connection: def __init__(self): self.sql = [] def execute(self, statement, params=None): self.sql.append(str(statement)) return _Scalar() class _Begin: def __init__(self, connection): self.connection = connection def __enter__(self): return self.connection def __exit__(self, *_args): return False class _Engine: def __init__(self): self.connection = _Connection() def begin(self): return _Begin(self.connection) def test_backfill_dry_run_and_apply(monkeypatch): import scripts.backfill_chatwoot_ingestion_visibility as backfill dry_engine = _Engine() monkeypatch.setattr(backfill, "engine", dry_engine) dry = backfill.backfill(apply=False, hours=24) assert dry["message_candidates"] == 2 assert dry["messages_updated"] == 0 assert dry["communications_upserted"] == 0 assert not any("UPDATE messages" in sql or "INSERT INTO communications" in sql for sql in dry_engine.connection.sql) apply_engine = _Engine() monkeypatch.setattr(backfill, "engine", apply_engine) applied = backfill.backfill(apply=True, hours=24, source_event_id="cw-42") assert applied["messages_updated"] == 2 assert applied["communications_upserted"] == 2 assert any("UPDATE messages" in sql for sql in apply_engine.connection.sql) assert any("INSERT INTO communications" in sql for sql in apply_engine.connection.sql) source = Path("scripts/backfill_chatwoot_ingestion_visibility.py").read_text() assert "--messages-only" in source and "--communications-only" in source assert "INSERT INTO tasks" not in source and "INSERT INTO opportunities" not in source