Files
clientflow_backend/app/persistence.py

515 lines
19 KiB
Python

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