From c4d55d37ed98faa7dc6756d6962ab174dbee437a Mon Sep 17 00:00:00 2001 From: plx Date: Fri, 14 Aug 2026 14:08:23 +0000 Subject: [PATCH] fix: preserve Chatwoot message visibility and identity --- app/analyzer.py | 55 ++++--- app/communication_service.py | 70 +++++++++ app/db.py | 5 + app/persistence.py | 143 +++++++++++++----- app/task_service.py | 12 ++ app/webhooks_chatwoot.py | 11 +- .../010_chatwoot_ingestion_identity.sql | 4 + scripts/audit_chatwoot_ingestion_gap.py | 29 ++-- .../backfill_chatwoot_ingestion_visibility.py | 132 ++++++++++++++++ tests/test_chatwoot_ingestion_consistency.py | 141 +++++++++++++++++ 10 files changed, 531 insertions(+), 71 deletions(-) create mode 100644 migrations/010_chatwoot_ingestion_identity.sql create mode 100644 scripts/backfill_chatwoot_ingestion_visibility.py create mode 100644 tests/test_chatwoot_ingestion_consistency.py diff --git a/app/analyzer.py b/app/analyzer.py index a860edd..fcefb82 100644 --- a/app/analyzer.py +++ b/app/analyzer.py @@ -4,13 +4,32 @@ from app.action_decider import decide_action from app.action_mapper import map_action_decision from app.config import settings from app.persistence import save_action_run +from app.persistence import save_inbound_message from app.schemas import ActionDecision, AnalyzeRequest, AnalyzeResponse, UsageInfo from app.task_service import create_task_from_action_result -async def analyze(request: AnalyzeRequest, raw_event_id: str | None = None) -> AnalyzeResponse: +async def analyze( + request: AnalyzeRequest, + raw_event_id: str | None = None, + source_event_id: str | None = None, +) -> AnalyzeResponse: needs_review = False + # The inbound fact is durable before any LLM/policy decision. This keeps + # ignored classifications and decision failures visible and replay-safe. + message_id, source_event_id = save_inbound_message( + request=request, raw_event_id=raw_event_id, source_event_id=source_event_id, + raw_body=request.last_customer_message, clean_body=request.last_customer_message, + ) + if (request.source or "").lower() == "chatwoot" and source_event_id: + from app.communication_service import upsert_chatwoot_inbound_communication + upsert_chatwoot_inbound_communication( + source_message_id=source_event_id, conversation_id=request.conversation_id, + contact_id=request.contact_id, body=request.last_customer_message, + metadata={"raw_event_id": raw_event_id, "message_id": message_id}, + ) + try: decision, action_result, usage, decision_source = await decide_action(request) except Exception as exc: @@ -32,24 +51,11 @@ async def analyze(request: AnalyzeRequest, raw_event_id: str | None = None) -> A ) decision_source = "fallback" - if str(action_result.action_code or "").upper() == "IGNORE_BOUNCE": - # v4.8.5: NDR/bounce emails are already readable in Chatwoot/Thunderbird. - # They must not create ClientFlow work, opportunities or fiscal customers. - return AnalyzeResponse( - app="ClientFlow", - model=settings.openrouter_model, - action_decision=decision, - action_result=action_result, - usage=usage, - needs_review=False, - action_run_id=None, - message_id=None, - task_id=None, - ) - if not action_result.safe_to_post: needs_review = True + is_ignored_bounce = str(action_result.action_code or "").upper() == "IGNORE_BOUNCE" + action_run_id, message_id = save_action_run( request=request, action_decision=decision, @@ -61,6 +67,8 @@ async def analyze(request: AnalyzeRequest, raw_event_id: str | None = None) -> A raw_body=request.last_customer_message, clean_body=request.last_customer_message, raw_event_id=raw_event_id, + source_event_id=source_event_id, + message_id=message_id, ) normalized_message_for_idem = re.sub(r"\s+", " ", str(request.last_customer_message or "").strip().casefold()) @@ -74,7 +82,7 @@ async def analyze(request: AnalyzeRequest, raw_event_id: str | None = None) -> A conversation_id=request.conversation_id, contact_id=request.contact_id, source_system=request.source or "manual", - source_event_id=message_id, + source_event_id=source_event_id or message_id, metadata={ "created_from_analyzer": True, "content_fingerprint": message_fingerprint, @@ -89,6 +97,19 @@ async def analyze(request: AnalyzeRequest, raw_event_id: str | None = None) -> A }, ) + if (request.source or "").lower() == "chatwoot" and source_event_id: + from app.communication_service import enrich_chatwoot_communication_from_task + action_code = str(action_result.action_code or "").upper() + enrich_chatwoot_communication_from_task( + source_message_id=source_event_id, classification=action_code, + confidence=float(getattr(decision, "confidence", 0.0) or 0.0), + ignored=is_ignored_bounce or action_code in {"IGNORE_SPAM", "SPAM"}, + conversation_id=request.conversation_id, contact_id=request.contact_id, + body=request.last_customer_message, task_id=task_id, + metadata={"raw_event_id": raw_event_id, "message_id": message_id, + "action_run_id": action_run_id, "decision_source": decision_source}, + ) + return AnalyzeResponse( app="ClientFlow", model=settings.openrouter_model, diff --git a/app/communication_service.py b/app/communication_service.py index e7989a8..3b708f6 100644 --- a/app/communication_service.py +++ b/app/communication_service.py @@ -477,6 +477,76 @@ def link_communication_to_opportunity(communication_id: str, opportunity_id: Opt """), {"id": communication_id, "opportunity_id": opportunity_id or ""}) +def upsert_chatwoot_inbound_communication( + *, source_message_id: str, conversation_id: Optional[str], contact_id: Optional[str], + body: str, classification: Optional[str] = None, confidence: Optional[float] = None, + status: str = "new", customer_id: Optional[str] = None, + opportunity_id: Optional[str] = None, task_id: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None, +) -> Optional[str]: + """Materialize one inbound Chatwoot message without schema or work creation.""" + if not str(source_message_id or "").strip(): + return None + with engine.begin() as conn: + row = conn.execute(text(""" + INSERT INTO communications ( + source_system, source_message_id, conversation_id, contact_id, + direction, body, classification, confidence, status, + customer_id, opportunity_id, task_id, metadata + ) VALUES ( + 'chatwoot', :source_message_id, :conversation_id, :contact_id, + 'inbound', :body, :classification, :confidence, :status, + CAST(:customer_id AS UUID), CAST(:opportunity_id AS UUID), + CAST(:task_id AS UUID), CAST(:metadata AS JSONB) + ) + ON CONFLICT (source_system, source_message_id) WHERE source_message_id IS NOT NULL + DO UPDATE SET + conversation_id=COALESCE(NULLIF(EXCLUDED.conversation_id,''), communications.conversation_id), + contact_id=COALESCE(NULLIF(EXCLUDED.contact_id,''), communications.contact_id), + body=COALESCE(NULLIF(EXCLUDED.body,''), communications.body), + classification=COALESCE(EXCLUDED.classification, communications.classification), + confidence=COALESCE(EXCLUDED.confidence, communications.confidence), + status=CASE WHEN EXCLUDED.classification IS NULL THEN communications.status ELSE EXCLUDED.status END, + customer_id=COALESCE(EXCLUDED.customer_id, communications.customer_id), + opportunity_id=COALESCE(EXCLUDED.opportunity_id, communications.opportunity_id), + task_id=COALESCE(EXCLUDED.task_id, communications.task_id), + metadata=COALESCE(communications.metadata,'{}'::jsonb) || EXCLUDED.metadata, + updated_at=now() + RETURNING id::text + """), {"source_message_id": str(source_message_id), + "conversation_id": conversation_id or "", "contact_id": contact_id or "", + "body": body or "", "classification": classification, + "confidence": confidence, "status": status or "new", + "customer_id": customer_id or None, "opportunity_id": opportunity_id or None, + "task_id": task_id or None, + "metadata": json.dumps(metadata or {}, ensure_ascii=False, default=str)}).first() + return str(row[0]) if row else None + + +def enrich_chatwoot_communication_from_task( + *, source_message_id: str, classification: str, confidence: Optional[float], + ignored: bool, conversation_id: Optional[str], contact_id: Optional[str], body: str, + task_id: Optional[str], metadata: Optional[Dict[str, Any]] = None, +) -> Optional[str]: + customer_id = opportunity_id = None + if task_id: + with engine.begin() as conn: + linked = conn.execute(text(""" + SELECT NULLIF(t.customer_id, '') AS customer_id, t.opportunity_id::text + FROM tasks t WHERE t.id=CAST(:task_id AS UUID) + """), {"task_id": task_id}).mappings().first() + if linked: + customer_id = linked.get("customer_id") + opportunity_id = linked.get("opportunity_id") + return upsert_chatwoot_inbound_communication( + source_message_id=source_message_id, conversation_id=conversation_id, + contact_id=contact_id, body=body, classification=classification, + confidence=confidence, status="ignored" if ignored else "classified", + customer_id=customer_id, opportunity_id=opportunity_id, task_id=task_id, + metadata=metadata, + ) + + def record_outbound_communication( diff --git a/app/db.py b/app/db.py index 1b71524..ab34d3f 100644 --- a/app/db.py +++ b/app/db.py @@ -296,6 +296,11 @@ def ensure_core_schema() -> None: conn.execute(text("CREATE INDEX IF NOT EXISTS idx_raw_events_conversation ON raw_events(conversation_id)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_messages_conversation ON messages(conversation_id)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_messages_raw_event ON messages(raw_event_id)")) + conn.execute(text(""" + CREATE UNIQUE INDEX IF NOT EXISTS ux_messages_source_event + ON messages(source_system, source_event_id) + WHERE source_event_id IS NOT NULL + """)) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_raw_events_message ON raw_events(message_id)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_action_runs_message ON action_runs(message_id)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_action_runs_raw_event ON action_runs(raw_event_id)")) diff --git a/app/persistence.py b/app/persistence.py index b027a9e..3df7462 100644 --- a/app/persistence.py +++ b/app/persistence.py @@ -25,6 +25,92 @@ def _uuid_or_none(value): return value or None +def _lock_source_identity(conn: Any, source_system: str, source_event_id: str | None) -> None: + if source_event_id: + conn.execute(text("SELECT pg_advisory_xact_lock(hashtext(:identity))"), { + "identity": f"message:{source_system}:{source_event_id}", + }) + + +def save_inbound_message( + *, + request: AnalyzeRequest, + raw_event_id: Optional[str] = None, + source_event_id: Optional[str] = None, + raw_body: Optional[str] = None, + clean_body: Optional[str] = None, +) -> Tuple[str, Optional[str]]: + """Persist the canonical inbound message before classification. + + Chatwoot identity comes from its raw event when callers do not pass it. + The advisory lock makes retries safe even while older databases are waiting + for the canonical identity index migration. + """ + if not settings.clientflow_persist: + return "", source_event_id + source_system = request.source or "manual" + conversation_id = request.conversation_id or "manual" + with engine.begin() as conn: + if raw_event_id: + raw = conn.execute(text(""" + SELECT source_system, source_event_id + FROM raw_events WHERE id = CAST(:raw_event_id AS UUID) + """), {"raw_event_id": raw_event_id}).mappings().first() + if raw: + source_system = str(raw.get("source_system") or source_system) + source_event_id = source_event_id or raw.get("source_event_id") + _lock_source_identity(conn, source_system, source_event_id) + existing = None + if source_event_id: + existing = conn.execute(text(""" + SELECT id::text FROM messages + WHERE source_system=:source_system AND source_event_id=:source_event_id + LIMIT 1 + """), {"source_system": source_system, "source_event_id": source_event_id}).first() + if not existing and raw_event_id: + existing = conn.execute(text(""" + SELECT id::text FROM messages WHERE raw_event_id=CAST(:raw_event_id AS UUID) LIMIT 1 + """), {"raw_event_id": raw_event_id}).first() + if existing: + message_id = str(existing[0]) + conn.execute(text(""" + UPDATE messages SET + source_event_id=COALESCE(source_event_id, :source_event_id), + raw_event_id=COALESCE(raw_event_id, CAST(:raw_event_id AS UUID)), + conversation_id=COALESCE(NULLIF(conversation_id,''), :conversation_id), + contact_id=COALESCE(NULLIF(contact_id,''), :contact_id), + raw_body=COALESCE(NULLIF(raw_body,''), :raw_body), + clean_body=COALESCE(NULLIF(clean_body,''), :clean_body) + WHERE id=CAST(:message_id AS UUID) + """), {"message_id": message_id, "source_event_id": source_event_id, + "raw_event_id": raw_event_id, "conversation_id": conversation_id, + "contact_id": request.contact_id, "raw_body": raw_body or request.last_customer_message, + "clean_body": clean_body or request.last_customer_message}) + else: + row = conn.execute(text(""" + INSERT INTO messages ( + raw_event_id, source_system, source_event_id, conversation_id, + contact_id, direction, raw_body, clean_body, previous_context, metadata + ) VALUES ( + CAST(:raw_event_id AS UUID), :source_system, :source_event_id, + :conversation_id, :contact_id, 'inbound', :raw_body, :clean_body, + :previous_context, CAST(:metadata AS JSONB) + ) RETURNING id::text + """), {"raw_event_id": raw_event_id, "source_system": source_system, + "source_event_id": source_event_id, "conversation_id": conversation_id, + "contact_id": request.contact_id, "raw_body": raw_body or request.last_customer_message, + "clean_body": clean_body or request.last_customer_message, + "previous_context": request.previous_context, + "metadata": _json({"current_state": request.current_state.model_dump()})}).first() + message_id = str(row[0]) + if raw_event_id: + conn.execute(text("""UPDATE raw_events SET message_id=CAST(:message_id AS UUID) + WHERE id=CAST(:raw_event_id AS UUID)"""), { + "message_id": message_id, "raw_event_id": raw_event_id, + }) + return message_id, source_event_id + + def save_action_run( *, request: AnalyzeRequest, @@ -37,53 +123,30 @@ def save_action_run( raw_body: Optional[str] = None, clean_body: Optional[str] = None, raw_event_id: Optional[str] = None, + source_event_id: Optional[str] = None, + message_id: Optional[str] = None, ) -> Tuple[str, str]: if not settings.clientflow_persist: return "", "" conversation_id = request.conversation_id or "manual" - with engine.begin() as conn: - message_row = conn.execute(text(""" - INSERT INTO messages ( - raw_event_id, - source_system, - source_event_id, - conversation_id, - contact_id, - direction, - raw_body, - clean_body, - previous_context, - metadata - ) - VALUES ( - CAST(:raw_event_id AS UUID), - :source_system, - NULL, - :conversation_id, - :contact_id, - 'inbound', - :raw_body, - :clean_body, - :previous_context, - CAST(:metadata AS JSONB) - ) - RETURNING id::text - """), { - "raw_event_id": raw_event_id, - "source_system": request.source or "manual", - "conversation_id": conversation_id, - "contact_id": request.contact_id, - "raw_body": raw_body or request.last_customer_message, - "clean_body": clean_body or request.last_customer_message, - "previous_context": request.previous_context, - "metadata": _json({ - "current_state": request.current_state.model_dump(), - }), - }).fetchone() + if not message_id: + message_id, source_event_id = save_inbound_message( + request=request, raw_event_id=raw_event_id, source_event_id=source_event_id, + raw_body=raw_body, clean_body=clean_body, + ) - message_id = message_row[0] + with engine.begin() as conn: + _lock_source_identity(conn, request.source or "manual", source_event_id) + existing_run = conn.execute(text(""" + SELECT id::text FROM action_runs + WHERE message_id=CAST(:message_id AS UUID) + OR (CAST(:raw_event_id AS UUID) IS NOT NULL AND raw_event_id=CAST(:raw_event_id AS UUID)) + ORDER BY created_at LIMIT 1 + """), {"message_id": message_id, "raw_event_id": raw_event_id}).first() + if existing_run: + return str(existing_run[0]), str(message_id) run_row = conn.execute(text(""" INSERT INTO action_runs ( diff --git a/app/task_service.py b/app/task_service.py index a9ab557..fc73abd 100644 --- a/app/task_service.py +++ b/app/task_service.py @@ -191,6 +191,18 @@ def create_task_from_action_result( action_code = data.get("action_code") or "REVIEW_MANUALLY" if str(action_code or "").strip().upper() == "IGNORE_BOUNCE": return None + if raw_event_id: + # A retry may re-run classification after the action_run was already + # persisted. Raw event identity wins over a potentially different LLM + # answer and prevents a second work item for the same inbound fact. + with engine.begin() as conn: + existing = conn.execute(text(""" + SELECT id::text FROM tasks + WHERE raw_event_id=CAST(:raw_event_id AS UUID) + ORDER BY created_at LIMIT 1 + """), {"raw_event_id": raw_event_id}).first() + if existing: + return str(existing[0]) config = get_action_config(action_code) route = data.get("route") or config["route"] diff --git a/app/webhooks_chatwoot.py b/app/webhooks_chatwoot.py index 9ff130f..24d4ddf 100644 --- a/app/webhooks_chatwoot.py +++ b/app/webhooks_chatwoot.py @@ -416,7 +416,11 @@ async def process_saved_chatwoot_raw_event(raw_event_id: str, payload: Dict[str, contact_id=extracted["contact_id"], ) - response = await analyze(analyze_request, raw_event_id=raw_event_id) + response = await analyze( + analyze_request, + raw_event_id=raw_event_id, + source_event_id=extracted.get("source_event_id"), + ) chatwoot_note_result = { "status": "skipped", @@ -429,11 +433,14 @@ async def process_saved_chatwoot_raw_event(raw_event_id: str, payload: Dict[str, action_result=response.action_result, ) + ignored_action = str(response.action_result.action_code or "").upper() in { + "IGNORE_BOUNCE", "IGNORE_SPAM", "SPAM", + } mark_raw_event_processed( raw_event_id=raw_event_id, action_run_id=response.action_run_id, message_id=response.message_id, - ignored=False, + ignored=ignored_action, error=None, ) diff --git a/migrations/010_chatwoot_ingestion_identity.sql b/migrations/010_chatwoot_ingestion_identity.sql new file mode 100644 index 0000000..38c493d --- /dev/null +++ b/migrations/010_chatwoot_ingestion_identity.sql @@ -0,0 +1,4 @@ +-- Canonical message identity for replay-safe Chatwoot ingestion. +CREATE UNIQUE INDEX IF NOT EXISTS ux_messages_source_event +ON messages(source_system, source_event_id) +WHERE source_event_id IS NOT NULL; diff --git a/scripts/audit_chatwoot_ingestion_gap.py b/scripts/audit_chatwoot_ingestion_gap.py index ce20808..d785f2b 100755 --- a/scripts/audit_chatwoot_ingestion_gap.py +++ b/scripts/audit_chatwoot_ingestion_gap.py @@ -18,6 +18,9 @@ from app.db import engine OUT_CSV = Path("/tmp/clientflow_chatwoot_ingestion_gap.csv") OUT_MD = Path("/tmp/clientflow_chatwoot_ingestion_gap.md") +# Historical report label retained for downstream parsers; pending inbound is +# now grouped under the required PROCESSING_ERROR state. +LEGACY_PENDING_STATE = "PENDING_INCOMING_RAW_EVENT" def clean(value: Any, limit: int = 240) -> str: @@ -45,13 +48,15 @@ def main() -> int: re.payload #>> '{sender,email}' AS sender_email, re.payload #>> '{conversation,meta,sender,email}' AS meta_sender_email, left(coalesce(re.payload #>> '{content}', ''), 600) AS content, - m.id IS NOT NULL AS has_message, + EXISTS ( + SELECT 1 FROM messages m + WHERE (m.source_system = 'chatwoot' AND m.source_event_id = re.source_event_id) + OR m.id = re.message_id + OR m.raw_event_id = re.id + ) AS has_message, c.id IS NOT NULL AS has_communication, o.id IS NOT NULL AS has_opportunity FROM raw_events re - LEFT JOIN messages m - ON m.source_system = 'chatwoot' - AND m.source_event_id = re.source_event_id LEFT JOIN communications c ON c.source_system = 'chatwoot' AND c.source_message_id = re.source_event_id @@ -76,16 +81,16 @@ def main() -> int: status = "NON_INCOMING_OR_UNKNOWN" if not args.include_ok: continue - elif not row.get("processed") and not row.get("ignored") and not row.get("processing_error"): - status = "PENDING_INCOMING_RAW_EVENT" elif row.get("processing_error"): - status = "INCOMING_PROCESSING_ERROR" + status = "PROCESSING_ERROR" elif row.get("ignored"): - status = "INCOMING_IGNORED" - elif not row.get("has_message") and not row.get("has_communication"): - status = "PROCESSED_BUT_NOT_VISIBLE" - elif not row.get("has_opportunity"): - status = "VISIBLE_BUT_UNLINKED_CONVERSATION" + status = "IGNORED" + elif row.get("has_message") and not row.get("has_communication"): + status = "MESSAGE_PRESENT_NO_COMMUNICATION" + elif row.get("processed") and not row.get("has_message"): + status = "TRUE_PROCESSED_WITHOUT_MESSAGE" + elif not row.get("processed"): + status = "PROCESSING_ERROR" else: status = "OK" if not args.include_ok: diff --git a/scripts/backfill_chatwoot_ingestion_visibility.py b/scripts/backfill_chatwoot_ingestion_visibility.py new file mode 100644 index 0000000..995e9f6 --- /dev/null +++ b/scripts/backfill_chatwoot_ingestion_visibility.py @@ -0,0 +1,132 @@ +#!/usr/bin/env python3 +"""Repair Chatwoot message identity and communication visibility; dry-run by default.""" +from __future__ import annotations + +import argparse +import sys +from pathlib import Path + +from sqlalchemy import text + +PROJECT_ROOT = Path(__file__).resolve().parents[1] +if str(PROJECT_ROOT) not in sys.path: + sys.path.insert(0, str(PROJECT_ROOT)) + +from app.db import engine + + +def backfill(*, apply: bool = False, hours: int = 168, + source_event_id: str | None = None, messages: bool = True, + communications: bool = True) -> dict[str, int]: + params = {"hours": max(1, int(hours)), "source_event_id": source_event_id or None} + result = {"message_candidates": 0, "messages_updated": 0, + "communication_candidates": 0, "communications_upserted": 0} + with engine.begin() as conn: + if messages: + predicate = """ + m.source_system='chatwoot' AND m.source_event_id IS NULL AND re.source_system='chatwoot' + AND re.source_event_id IS NOT NULL AND m.raw_event_id=re.id + AND re.created_at >= now() - (:hours * interval '1 hour') + AND (CAST(:source_event_id AS TEXT) IS NULL OR re.source_event_id=CAST(:source_event_id AS TEXT)) + """ + result["message_candidates"] = int(conn.execute(text( + "SELECT COUNT(*) FROM messages m JOIN raw_events re ON " + predicate + ), params).scalar() or 0) + if apply and result["message_candidates"]: + updated = conn.execute(text(""" + UPDATE messages m SET source_event_id=re.source_event_id + FROM raw_events re WHERE """ + predicate + """ + AND NOT EXISTS ( + SELECT 1 FROM messages existing + WHERE existing.source_system=m.source_system + AND existing.source_event_id=re.source_event_id + AND existing.id<>m.id + ) + """), params) + result["messages_updated"] = int(updated.rowcount or 0) + + if communications: + candidates = """ + WITH chatwoot_messages AS ( + SELECT DISTINCT ON (m.id) + m.id, COALESCE(m.source_event_id, re.source_event_id) AS source_message_id, + COALESCE(m.conversation_id,re.conversation_id) AS conversation_id, + COALESCE(m.contact_id,re.contact_id) AS contact_id, + COALESCE(m.clean_body,m.raw_body,re.payload #>> '{content}',re.payload #>> '{message,content}') AS body, + re.id AS raw_event_id, re.ignored, m.created_at + FROM messages m + LEFT JOIN raw_events re ON re.id=m.raw_event_id OR re.message_id=m.id + WHERE m.source_system='chatwoot' + AND COALESCE(m.source_event_id,re.source_event_id) IS NOT NULL + AND m.created_at >= now() - (:hours * interval '1 hour') + AND (CAST(:source_event_id AS TEXT) IS NULL OR COALESCE(m.source_event_id,re.source_event_id)=CAST(:source_event_id AS TEXT)) + ORDER BY m.id, re.created_at DESC NULLS LAST + ), enriched AS ( + SELECT cm.*, + ar.action_result ->> 'action_code' AS classification, + task.id AS task_id, task.opportunity_id, NULLIF(task.customer_id,'') AS customer_id + FROM chatwoot_messages cm + LEFT JOIN LATERAL ( + SELECT a.* FROM action_runs a + WHERE a.message_id=cm.id OR a.raw_event_id=cm.raw_event_id + ORDER BY a.created_at DESC LIMIT 1 + ) ar ON TRUE + LEFT JOIN LATERAL ( + SELECT t.* FROM tasks t + WHERE t.message_id=cm.id OR t.raw_event_id=cm.raw_event_id OR t.action_run_id=ar.id + ORDER BY t.created_at DESC LIMIT 1 + ) task ON TRUE + ) + """ + result["communication_candidates"] = int(conn.execute(text(candidates + """ + SELECT COUNT(*) FROM enriched e WHERE NOT EXISTS ( + SELECT 1 FROM communications c + WHERE c.source_system='chatwoot' AND c.source_message_id=e.source_message_id) + """), params).scalar() or 0) + if apply: + inserted = conn.execute(text(candidates + """ + INSERT INTO communications ( + source_system,source_message_id,conversation_id,contact_id,direction,body, + classification,status,customer_id,opportunity_id,task_id,metadata,created_at,updated_at) + SELECT 'chatwoot',source_message_id,conversation_id,contact_id,'inbound',body, + classification, + CASE WHEN ignored OR classification IN ('IGNORE_BOUNCE','IGNORE_SPAM','SPAM') + THEN 'ignored' WHEN classification IS NULL THEN 'new' ELSE 'classified' END, + CASE WHEN customer_id ~* '^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$' + THEN CAST(customer_id AS UUID) ELSE NULL END, + opportunity_id,task_id, + jsonb_build_object('backfilled',true,'message_id',id::text,'raw_event_id',raw_event_id::text), + created_at,now() + FROM enriched + ON CONFLICT (source_system,source_message_id) WHERE source_message_id IS NOT NULL + DO UPDATE SET + classification=COALESCE(EXCLUDED.classification,communications.classification), + customer_id=COALESCE(EXCLUDED.customer_id,communications.customer_id), + opportunity_id=COALESCE(EXCLUDED.opportunity_id,communications.opportunity_id), + task_id=COALESCE(EXCLUDED.task_id,communications.task_id), + metadata=COALESCE(communications.metadata,'{}'::jsonb)||EXCLUDED.metadata, + updated_at=now() + """), params) + result["communications_upserted"] = int(inserted.rowcount or 0) + return result + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--apply", action="store_true") + parser.add_argument("--hours", type=int, default=168) + parser.add_argument("--source-event-id") + mode = parser.add_mutually_exclusive_group() + mode.add_argument("--messages-only", action="store_true") + mode.add_argument("--communications-only", action="store_true") + args = parser.parse_args() + outcome = backfill(apply=args.apply, hours=args.hours, + source_event_id=args.source_event_id, + messages=not args.communications_only, + communications=not args.messages_only) + print(" ".join(f"{key}={value}" for key, value in outcome.items()), f"apply={args.apply}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_chatwoot_ingestion_consistency.py b/tests/test_chatwoot_ingestion_consistency.py new file mode 100644 index 0000000..18ab2dc --- /dev/null +++ b/tests/test_chatwoot_ingestion_consistency.py @@ -0,0 +1,141 @@ +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