diff --git a/app/persistence.py b/app/persistence.py index 3df7462..4cc5120 100644 --- a/app/persistence.py +++ b/app/persistence.py @@ -1,3 +1,4 @@ +from datetime import datetime, timezone from typing import Any, Dict, Optional, Tuple import json @@ -32,6 +33,95 @@ def _lock_source_identity(conn: Any, source_system: str, source_event_id: str | }) +def _message_created_at(value: Any) -> datetime | None: + """Normalize a source message timestamp without substituting processing time.""" + if value in (None, ""): + return None + if isinstance(value, datetime): + return value if value.tzinfo else value.replace(tzinfo=timezone.utc) + if isinstance(value, (int, float)) or str(value).strip().replace(".", "", 1).isdigit(): + try: + return datetime.fromtimestamp(float(value), tz=timezone.utc) + except (OverflowError, TypeError, ValueError): + return None + try: + parsed = datetime.fromisoformat(str(value).strip().replace("Z", "+00:00")) + except ValueError: + return None + return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc) + + +def save_factual_chatwoot_message( + *, + raw_event_id: str, + source_event_id: str | None, + conversation_id: str | None, + contact_id: str | None, + direction: str, + raw_body: str, + clean_body: str | None = None, + source_created_at: Any = None, + metadata: Dict[str, Any] | None = None, +) -> tuple[str, bool]: + """Persist one factual public Chatwoot message, without workflow effects. + + Returns ``(message_id, inserted)``. Canonical source identity and the + transaction advisory lock make webhook delivery and backfill idempotent. + """ + if not settings.clientflow_persist: + return "", False + normalized_direction = str(direction or "").strip().lower() + if normalized_direction not in {"inbound", "outbound"}: + raise ValueError("direction must be inbound or outbound") + created_at = _message_created_at(source_created_at) + with engine.begin() as conn: + _lock_source_identity(conn, "chatwoot", source_event_id) + existing = None + if source_event_id: + existing = conn.execute(text(""" + SELECT id::text FROM messages + WHERE source_system='chatwoot' AND source_event_id=:source_event_id + LIMIT 1 + """), {"source_event_id": source_event_id}).first() + if not existing: + 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() + inserted = existing is None + if inserted: + 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, created_at + ) VALUES ( + CAST(:raw_event_id AS UUID), 'chatwoot', :source_event_id, + :conversation_id, :contact_id, :direction, :raw_body, :clean_body, + NULL, CAST(:metadata AS JSONB), + COALESCE(CAST(:created_at AS TIMESTAMPTZ), now()) + ) RETURNING id::text + """), { + "raw_event_id": raw_event_id, + "source_event_id": source_event_id, + "conversation_id": conversation_id, + "contact_id": contact_id, + "direction": normalized_direction, + "raw_body": raw_body, + "clean_body": clean_body if clean_body is not None else raw_body, + "metadata": _json(metadata), + "created_at": created_at, + }).first() + message_id = str(row[0]) + else: + message_id = str(existing[0]) + conn.execute(text(""" + UPDATE raw_events SET message_id=CAST(:message_id AS UUID) + WHERE id=CAST(:raw_event_id AS UUID) AND message_id IS NULL + """), {"message_id": message_id, "raw_event_id": raw_event_id}) + return message_id, inserted + + def save_inbound_message( *, request: AnalyzeRequest, diff --git a/app/webhooks_chatwoot.py b/app/webhooks_chatwoot.py index 24d4ddf..4ad8270 100644 --- a/app/webhooks_chatwoot.py +++ b/app/webhooks_chatwoot.py @@ -17,6 +17,7 @@ from app.persistence import ( get_state_for_conversation, mark_raw_event_error, mark_raw_event_processed, + save_factual_chatwoot_message, save_raw_event, ) from app.schemas import AnalyzeRequest @@ -32,6 +33,12 @@ def as_str(value: Any) -> Optional[str]: return str(value) +def as_bool(value: Any) -> bool: + if isinstance(value, str): + return value.strip().lower() in {"1", "true", "yes"} + return bool(value) + + def get_nested(data: Dict[str, Any], *keys: str) -> Any: current: Any = data @@ -103,6 +110,11 @@ def message_timestamp(message: Dict[str, Any]) -> Any: ) +def source_message_timestamp(payload: Dict[str, Any]) -> Any: + message = payload.get("message") if isinstance(payload.get("message"), dict) else payload + return message.get("created_at") or payload.get("created_at") or message.get("timestamp") or payload.get("timestamp") + + def extract_message_content_for_context(message: Dict[str, Any]) -> str: try: content = extract_chatwoot_content(message, message) @@ -261,9 +273,21 @@ def extract_chatwoot_event(payload: Dict[str, Any]) -> Dict[str, Any]: message_type_text = str(message_type).lower() is_outgoing = message_type_text in {"outgoing", "outbound", "1"} + # The sender of an outgoing message is normally the agent, not the customer. + # Retain a contact only when Chatwoot supplied an explicit conversation/contact + # identity rather than falling back to the sender object. + if ( + is_outgoing + and payload.get("contact_id") is None + and not isinstance(payload.get("contact"), dict) + and not isinstance(conversation.get("contact"), dict) + and not isinstance(get_nested(conversation, "contact_inbox", "contact"), dict) + ): + contact_id = None + is_private = bool( - message.get("private") - or payload.get("private") + as_bool(message.get("private")) + or as_bool(payload.get("private")) or message.get("content_type") == "input_select" ) @@ -292,6 +316,7 @@ def extract_chatwoot_event(payload: Dict[str, Any]) -> Dict[str, Any]: "is_private": is_private, "sender_name": as_str(sender_name), "sender_type": as_str(sender_type), + "created_at": source_message_timestamp(payload), } @@ -375,6 +400,25 @@ async def process_saved_chatwoot_raw_event(raw_event_id: str, payload: Dict[str, } if extracted["is_outgoing"]: + factual_message_id = None + if extracted.get("event_type") == "message_created" and not extracted.get("is_private"): + factual_message_id, _ = save_factual_chatwoot_message( + raw_event_id=raw_event_id, + source_event_id=extracted.get("source_event_id"), + conversation_id=extracted.get("conversation_id"), + contact_id=extracted.get("contact_id"), + direction="outbound", + raw_body=extracted["content"], + clean_body=extracted["content"], + source_created_at=extracted.get("created_at"), + metadata={ + "message_type": extracted.get("message_type"), + "public": True, + "private": False, + "sender_name": extracted.get("sender_name"), + "sender_type": extracted.get("sender_type"), + }, + ) auto_complete_result = auto_complete_task_from_outgoing_message( conversation_id=extracted["conversation_id"], contact_id=extracted["contact_id"], @@ -400,6 +444,7 @@ async def process_saved_chatwoot_raw_event(raw_event_id: str, payload: Dict[str, "status": "outgoing_processed" if completed else "ignored", "reason": auto_complete_result.get("status"), "raw_event_id": raw_event_id, + "message_id": factual_message_id, "conversation_id": extracted["conversation_id"], "contact_id": extracted["contact_id"], "auto_complete": auto_complete_result, diff --git a/scripts/backfill_chatwoot_outbound_messages.py b/scripts/backfill_chatwoot_outbound_messages.py new file mode 100644 index 0000000..4d4793c --- /dev/null +++ b/scripts/backfill_chatwoot_outbound_messages.py @@ -0,0 +1,93 @@ +#!/usr/bin/env python3 +"""Backfill factual public Chatwoot outbound messages; 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 +from app.persistence import save_factual_chatwoot_message +from app.webhooks_chatwoot import extract_chatwoot_event + + +def _candidate_rows() -> list[dict]: + with engine.connect() as conn: + rows = conn.execute(text(""" + SELECT re.id::text, re.source_event_id, re.conversation_id, + re.contact_id, re.payload, re.created_at + FROM raw_events re + WHERE re.source_system = 'chatwoot' + AND re.event_type = 'message_created' + AND lower(COALESCE(re.payload->>'message_type', + re.payload #>> '{message,message_type}', '')) + IN ('outgoing', 'outbound', '1') + AND lower(COALESCE(re.payload->>'private', + re.payload #>> '{message,private}', 'false')) + NOT IN ('true', '1', 'yes') + ORDER BY re.created_at, re.id + """)).mappings().all() + conn.rollback() + return [dict(row) for row in rows] + + +def backfill(*, apply: bool = False) -> dict[str, int]: + result = {"candidates": 0, "inserted": 0, "already_present": 0, "skipped": 0, "errors": 0} + for row in _candidate_rows(): + payload = row.get("payload") if isinstance(row.get("payload"), dict) else {} + extracted = extract_chatwoot_event(payload) + if not extracted.get("is_outgoing") or extracted.get("is_private") or not extracted.get("content"): + result["skipped"] += 1 + continue + result["candidates"] += 1 + with engine.connect() as conn: + present = conn.execute(text(""" + SELECT 1 FROM messages + WHERE (source_system='chatwoot' AND source_event_id=:source_event_id) + OR raw_event_id=CAST(:raw_event_id AS UUID) + LIMIT 1 + """), {"source_event_id": row.get("source_event_id"), "raw_event_id": row["id"]}).first() + conn.rollback() + if present: + result["already_present"] += 1 + continue + if not apply: + continue + try: + _, inserted = save_factual_chatwoot_message( + raw_event_id=row["id"], + source_event_id=extracted.get("source_event_id") or row.get("source_event_id"), + conversation_id=extracted.get("conversation_id") or row.get("conversation_id"), + contact_id=extracted.get("contact_id") or row.get("contact_id"), + direction="outbound", + raw_body=extracted["content"], + clean_body=extracted["content"], + source_created_at=extracted.get("created_at") or row.get("created_at"), + metadata={"message_type": extracted.get("message_type"), "public": True, + "private": False, "backfilled": True, + "sender_name": extracted.get("sender_name"), + "sender_type": extracted.get("sender_type")}, + ) + result["inserted" if inserted else "already_present"] += 1 + except Exception: + result["errors"] += 1 + return result + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--apply", action="store_true", help="Insert missing factual messages") + args = parser.parse_args() + result = backfill(apply=args.apply) + print(" ".join(f"{key}={value}" for key, value in result.items()), f"apply={args.apply}") + return 1 if result["errors"] else 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_chatwoot_outbound_factual_timeline.py b/tests/test_chatwoot_outbound_factual_timeline.py new file mode 100644 index 0000000..d9af836 --- /dev/null +++ b/tests/test_chatwoot_outbound_factual_timeline.py @@ -0,0 +1,200 @@ +import asyncio +from datetime import datetime, timezone + +import app.webhooks_chatwoot as webhook +import app.persistence as persistence +from app.persistence import _message_created_at + + +def _payload(*, private=False, message_id=712, content="Resposta pública"): + return { + "event": "message_created", + "id": message_id, + "message_type": "outgoing", + "private": private, + "content": content, + "created_at": 1_700_000_000, + "conversation": {"id": 2450}, + "sender": {"id": 9, "name": "Operador", "type": "user"}, + } + + +def test_chatwoot_timestamp_preserves_epoch_and_iso(): + assert _message_created_at(1_700_000_000) == datetime.fromtimestamp(1_700_000_000, tz=timezone.utc) + assert _message_created_at("2026-08-15T10:20:30Z") == datetime(2026, 8, 15, 10, 20, 30, tzinfo=timezone.utc) + + +def test_outgoing_public_persists_fact_before_existing_outgoing_processing(monkeypatch): + calls = [] + monkeypatch.setattr(webhook, "save_factual_chatwoot_message", lambda **kw: (calls.append(("message", kw)) or ("msg-1", True))) + monkeypatch.setattr(webhook, "auto_complete_task_from_outgoing_message", lambda **kw: (calls.append(("outgoing", kw)) or {"status": "no_matching_pending_task"})) + monkeypatch.setattr(webhook, "mark_raw_event_processed", lambda **kw: calls.append(("processed", kw))) + + result = asyncio.run(webhook.process_saved_chatwoot_raw_event("00000000-0000-0000-0000-000000000001", _payload())) + + assert [call[0] for call in calls] == ["message", "outgoing", "processed"] + assert calls[0][1]["direction"] == "outbound" + assert calls[0][1]["source_event_id"] == "712" + assert calls[0][1]["contact_id"] is None + assert calls[0][1]["source_created_at"] == 1_700_000_000 + assert result["message_id"] == "msg-1" + + +def test_outgoing_public_does_not_enter_inbound_analysis(monkeypatch): + monkeypatch.setattr(webhook, "save_factual_chatwoot_message", lambda **kw: ("msg-1", True)) + monkeypatch.setattr(webhook, "auto_complete_task_from_outgoing_message", lambda **kw: {"status": "no_matching_pending_task"}) + monkeypatch.setattr(webhook, "mark_raw_event_processed", lambda **kw: None) + + async def forbidden(*args, **kwargs): + raise AssertionError("outgoing factual messages must not enter classification") + + monkeypatch.setattr(webhook, "analyze", forbidden) + asyncio.run(webhook.process_saved_chatwoot_raw_event("00000000-0000-0000-0000-000000000001", _payload())) + + +def test_outgoing_private_note_is_not_persisted_as_public_message(monkeypatch): + monkeypatch.setattr(webhook, "save_factual_chatwoot_message", lambda **kw: (_ for _ in ()).throw(AssertionError("private note persisted"))) + monkeypatch.setattr(webhook, "auto_complete_task_from_outgoing_message", lambda **kw: {"status": "ignored_private_note"}) + monkeypatch.setattr(webhook, "mark_raw_event_processed", lambda **kw: None) + + result = asyncio.run(webhook.process_saved_chatwoot_raw_event("00000000-0000-0000-0000-000000000001", _payload(private=True))) + assert result["message_id"] is None + + +def test_string_false_private_value_is_public(): + assert webhook.extract_chatwoot_event(_payload(private="false"))["is_private"] is False + + +def test_non_message_created_outgoing_event_is_not_persisted(monkeypatch): + payload = _payload() + payload["event"] = "conversation_updated" + monkeypatch.setattr(webhook, "save_factual_chatwoot_message", lambda **kw: (_ for _ in ()).throw(AssertionError("wrong event persisted"))) + monkeypatch.setattr(webhook, "auto_complete_task_from_outgoing_message", lambda **kw: {"status": "no_matching_pending_task"}) + monkeypatch.setattr(webhook, "mark_raw_event_processed", lambda **kw: None) + asyncio.run(webhook.process_saved_chatwoot_raw_event("00000000-0000-0000-0000-000000000001", payload)) + + +class _Rows: + def __init__(self, row=None): + self.row = row + + def first(self): + return self.row + + +class _PersistenceConnection: + def __init__(self): + self.message_id = None + self.inserts = 0 + + def execute(self, statement, params=None): + sql = str(statement) + if "pg_advisory_xact_lock" in sql: + return _Rows() + if "SELECT id::text FROM messages" in sql: + return _Rows((self.message_id,)) if self.message_id else _Rows() + if "INSERT INTO messages" in sql: + self.inserts += 1 + self.message_id = "00000000-0000-0000-0000-000000000099" + return _Rows((self.message_id,)) + if "UPDATE raw_events SET message_id" in sql: + return _Rows() + raise AssertionError(sql) + + +class _Begin: + def __init__(self, connection): + self.connection = connection + + def __enter__(self): + return self.connection + + def __exit__(self, *_args): + return False + + +class _PersistenceEngine: + def __init__(self): + self.connection = _PersistenceConnection() + + def begin(self): + return _Begin(self.connection) + + +def test_factual_persistence_is_idempotent_by_source_identity(monkeypatch): + fake = _PersistenceEngine() + monkeypatch.setattr(persistence, "engine", fake) + monkeypatch.setattr(persistence.settings, "clientflow_persist", True) + kwargs = { + "raw_event_id": "00000000-0000-0000-0000-000000000001", + "source_event_id": "712", + "conversation_id": "2450", + "contact_id": None, + "direction": "outbound", + "raw_body": "Resposta", + "source_created_at": 1_700_000_000, + } + first = persistence.save_factual_chatwoot_message(**kwargs) + second = persistence.save_factual_chatwoot_message(**kwargs) + assert first[1] is True + assert second == (first[0], False) + assert fake.connection.inserts == 1 + + +class _BackfillPresence: + def __init__(self, present=False): + self.present = present + + def execute(self, statement, params=None): + return _Rows((1,)) if self.present else _Rows() + + def rollback(self): + return None + + def __enter__(self): + return self + + def __exit__(self, *_args): + return False + + +class _BackfillEngine: + def __init__(self, presence): + self.presence = presence + + def connect(self): + return self.presence + + +def test_outbound_backfill_dry_run_apply_and_repeat(monkeypatch): + import scripts.backfill_chatwoot_outbound_messages as backfill + + row = { + "id": "00000000-0000-0000-0000-000000000001", + "source_event_id": "712", + "conversation_id": "2450", + "contact_id": None, + "payload": _payload(), + "created_at": datetime(2026, 8, 15, tzinfo=timezone.utc), + } + presence = _BackfillPresence(False) + monkeypatch.setattr(backfill, "_candidate_rows", lambda: [row]) + monkeypatch.setattr(backfill, "engine", _BackfillEngine(presence)) + writes = [] + + def save(**kwargs): + writes.append(kwargs) + presence.present = True + return "msg-1", True + + monkeypatch.setattr(backfill, "save_factual_chatwoot_message", save) + + assert backfill.backfill(apply=False) == { + "candidates": 1, "inserted": 0, "already_present": 0, "skipped": 0, "errors": 0, + } + assert writes == [] + assert backfill.backfill(apply=True)["inserted"] == 1 + repeated = backfill.backfill(apply=True) + assert repeated["inserted"] == 0 + assert repeated["already_present"] == 1 + assert len(writes) == 1