from datetime import datetime, timezone from typing import Any, Dict, Optional, Tuple import json from sqlalchemy import text from app.config import settings from app.db import engine from app.schemas import ( ActionDecision, ActionResult, AnalyzeRequest, CurrentState, UsageInfo, ) def _json(value) -> str: if hasattr(value, "model_dump"): value = value.model_dump() return json.dumps(value or {}, ensure_ascii=False) def _uuid_or_none(value): value = str(value or "").strip() 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 _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, 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, action_decision: ActionDecision, action_result: ActionResult, usage: UsageInfo, needs_review: bool, model: str, decision_source: str, 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" 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, ) 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 ( message_id, raw_event_id, conversation_id, contact_id, source_system, model, provider, openrouter_generation_id, decision_source, action_decision, action_result, prompt_tokens, completion_tokens, total_tokens, cost, usage, needs_review ) VALUES ( CAST(:message_id AS UUID), CAST(:raw_event_id AS UUID), :conversation_id, :contact_id, :source_system, :model, :provider, :openrouter_generation_id, :decision_source, CAST(:action_decision AS JSONB), CAST(:action_result AS JSONB), :prompt_tokens, :completion_tokens, :total_tokens, :cost, CAST(:usage AS JSONB), :needs_review ) RETURNING id::text """), { "message_id": message_id, "raw_event_id": raw_event_id, "conversation_id": conversation_id, "contact_id": request.contact_id, "source_system": request.source or "manual", "model": model, "provider": usage.provider, "openrouter_generation_id": usage.id, "decision_source": decision_source, "action_decision": _json(action_decision), "action_result": _json(action_result), "prompt_tokens": usage.prompt_tokens, "completion_tokens": usage.completion_tokens, "total_tokens": usage.total_tokens, "cost": usage.cost, "usage": _json(usage), "needs_review": needs_review, }).fetchone() action_run_id = run_row[0] if raw_event_id: conn.execute(text(""" UPDATE raw_events SET message_id = CAST(:message_id AS UUID), action_run_id = CAST(:action_run_id AS UUID) WHERE id = CAST(:raw_event_id AS UUID) """), { "message_id": message_id, "action_run_id": action_run_id, "raw_event_id": raw_event_id, }) return action_run_id, message_id def save_raw_event( source_system: str, event_type: str | None, source_event_id: str | None, conversation_id: str | None, contact_id: str | None, payload: dict, ) -> Dict[str, Any]: with engine.begin() as conn: row = conn.execute(text(""" INSERT INTO raw_events ( source_system, event_type, source_event_id, conversation_id, contact_id, payload ) VALUES ( :source_system, :event_type, :source_event_id, :conversation_id, :contact_id, CAST(:payload AS JSONB) ) ON CONFLICT (source_system, source_event_id) WHERE source_event_id IS NOT NULL DO UPDATE SET payload = EXCLUDED.payload, event_type = EXCLUDED.event_type, conversation_id = COALESCE(EXCLUDED.conversation_id, raw_events.conversation_id), contact_id = COALESCE(EXCLUDED.contact_id, raw_events.contact_id) RETURNING id::text, processed, ignored, message_id::text, action_run_id::text, (xmax = 0) AS inserted """), { "source_system": source_system, "event_type": event_type, "source_event_id": source_event_id, "conversation_id": conversation_id, "contact_id": contact_id, "payload": _json(payload), }).mappings().first() return dict(row or {}) def mark_raw_event_processed( raw_event_id: str, action_run_id: str | None = None, message_id: str | None = None, ignored: bool = False, error: str | None = None, ) -> None: with engine.begin() as conn: conn.execute(text(""" UPDATE raw_events SET processed = TRUE, ignored = :ignored, processing_error = :error, action_run_id = COALESCE(CAST(:action_run_id AS UUID), action_run_id), message_id = COALESCE(CAST(:message_id AS UUID), message_id), processed_at = now() WHERE id = CAST(:raw_event_id AS UUID) """), { "raw_event_id": raw_event_id, "ignored": ignored, "error": error, "action_run_id": _uuid_or_none(action_run_id), "message_id": _uuid_or_none(message_id), }) def mark_raw_event_error(raw_event_id: str, error: str) -> None: """Regista erro de processamento sem marcar o evento como processado. Isto evita filas silenciosas: o evento deixa de ficar em processed=false/ignored=false/processing_error=null, mas continua elegível para recovery explícito com scripts administrativos. """ with engine.begin() as conn: conn.execute(text(""" UPDATE raw_events SET processed = FALSE, ignored = FALSE, processing_error = :error, processed_at = now() WHERE id = CAST(:raw_event_id AS UUID) """), { "raw_event_id": raw_event_id, "error": str(error or "processing_exception")[:1000], }) def get_state_for_conversation(conversation_id: str | None) -> CurrentState: if not conversation_id: return CurrentState() with engine.begin() as conn: row = conn.execute(text(""" SELECT t.action_code, t.route, t.status, t.action, t.note, t.created_at FROM tasks t WHERE t.conversation_id = :conversation_id ORDER BY t.created_at DESC LIMIT 1 """), { "conversation_id": conversation_id, }).mappings().first() if not row: return CurrentState() return CurrentState( last_action_code=row.get("action_code") or "desconhecido", last_route=row.get("route") or "desconhecido", last_task_status=row.get("status") or "desconhecido", metadata={ "last_action": row.get("action"), "last_note": row.get("note"), "last_created_at": str(row.get("created_at")), }, ) def get_recent_chatwoot_public_context( conversation_id: str | None, *, current_source_event_id: str | None = None, max_messages: int = 2, ) -> list[str]: """Devolve as últimas mensagens públicas Chatwoot para contexto LLM. Usado quando o webhook não traz `conversation.messages`. Lê raw_events já guardados e exclui a mensagem atual para evitar que o LLM confunda histórico com o email recebido agora. """ if not conversation_id: return [] with engine.begin() as conn: rows = conn.execute(text(""" SELECT source_event_id, payload #>> '{message_type}' AS message_type, payload #>> '{content}' AS content, created_at FROM raw_events WHERE source_system = 'chatwoot' AND event_type = 'message_created' AND conversation_id = :conversation_id AND COALESCE(payload #>> '{content}', '') <> '' AND (CAST(:current_source_event_id AS TEXT) IS NULL OR source_event_id <> CAST(:current_source_event_id AS TEXT)) AND COALESCE(payload #>> '{private}', 'false') <> 'true' ORDER BY created_at DESC LIMIT :limit """), { "conversation_id": str(conversation_id), "current_source_event_id": str(current_source_event_id or "") or None, "limit": max(1, int(max_messages or 2)), }).mappings().all() lines: list[str] = [] for row in reversed(rows): role = "BLIF" if str(row.get("message_type") or "").lower() == "outgoing" else "Cliente" content = str(row.get("content") or "") content = " ".join(content.split()) if len(content) > 360: content = content[:359].rstrip() + "…" if content: lines.append(f"- {role}: {content}") return lines