94 lines
3.8 KiB
Python
94 lines
3.8 KiB
Python
#!/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())
|