Files
clientflow_backend/scripts/backfill_chatwoot_ingestion_visibility.py

133 lines
7.1 KiB
Python

#!/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())