fix: preserve Chatwoot message visibility and identity
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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)"))
|
||||
|
||||
@@ -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 (
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
|
||||
4
migrations/010_chatwoot_ingestion_identity.sql
Normal file
4
migrations/010_chatwoot_ingestion_identity.sql
Normal file
@@ -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;
|
||||
@@ -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:
|
||||
|
||||
132
scripts/backfill_chatwoot_ingestion_visibility.py
Normal file
132
scripts/backfill_chatwoot_ingestion_visibility.py
Normal file
@@ -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())
|
||||
141
tests/test_chatwoot_ingestion_consistency.py
Normal file
141
tests/test_chatwoot_ingestion_consistency.py
Normal file
@@ -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
|
||||
Reference in New Issue
Block a user