From 89017f71a6958ade3a6a3dc1b691fb080b3b1c42 Mon Sep 17 00:00:00 2001 From: plx Date: Fri, 14 Aug 2026 12:30:01 +0000 Subject: [PATCH] perf: optimize historical communication fallback --- app/communication_service.py | 116 +++++++++++++++++- app/db.py | 7 ++ ...istorical_communication_lookup_indexes.sql | 9 ++ ...ical_communication_fallback_performance.py | 99 +++++++++++++++ 4 files changed, 225 insertions(+), 6 deletions(-) create mode 100644 migrations/008_historical_communication_lookup_indexes.sql create mode 100644 tests/test_historical_communication_fallback_performance.py diff --git a/app/communication_service.py b/app/communication_service.py index aa04029..e7989a8 100644 --- a/app/communication_service.py +++ b/app/communication_service.py @@ -214,12 +214,114 @@ def _list_message_backed_chatwoot_items_for_opportunity(opportunity_id: str, lim if not str(opportunity_id or "").strip(): return [] with engine.begin() as conn: - rows = conn.execute(text(""" + candidate_sql = """ + SELECT m.id AS message_id + FROM opp JOIN messages m + ON m.source_system = 'chatwoot' + AND NULLIF(opp.conversation_id, '') IS NOT NULL + AND m.conversation_id = opp.conversation_id + UNION + SELECT m.id + FROM opp JOIN raw_events re + ON NULLIF(opp.conversation_id, '') IS NOT NULL + AND re.conversation_id = opp.conversation_id + JOIN messages m ON m.source_system = 'chatwoot' AND m.raw_event_id = re.id + UNION + SELECT m.id + FROM opp JOIN raw_events re + ON NULLIF(opp.conversation_id, '') IS NOT NULL + AND re.conversation_id = opp.conversation_id + JOIN messages m ON m.source_system = 'chatwoot' AND m.id = re.message_id + UNION + SELECT m.id + FROM opp JOIN tasks t ON t.opportunity_id = opp.id + JOIN messages m ON m.source_system = 'chatwoot' AND m.id = t.message_id + UNION + SELECT m.id + FROM opp JOIN tasks t ON t.opportunity_id = opp.id + JOIN messages m ON m.source_system = 'chatwoot' AND m.raw_event_id = t.raw_event_id + UNION + SELECT m.id + FROM opp JOIN tasks t ON t.opportunity_id = opp.id + JOIN raw_events re ON re.id = t.raw_event_id + JOIN messages m ON m.source_system = 'chatwoot' AND m.id = re.message_id + UNION + SELECT m.id + FROM opp JOIN tasks t ON t.opportunity_id = opp.id + JOIN action_runs ar ON ar.id = t.action_run_id + JOIN messages m ON m.source_system = 'chatwoot' AND m.id = ar.message_id + UNION + SELECT m.id + FROM opp JOIN tasks t ON t.opportunity_id = opp.id + JOIN action_runs ar ON ar.id = t.action_run_id + JOIN messages m ON m.source_system = 'chatwoot' AND m.raw_event_id = ar.raw_event_id + UNION + SELECT m.id + FROM opp JOIN tasks t ON t.opportunity_id = opp.id + JOIN action_runs ar ON ar.id = t.action_run_id + JOIN raw_events re ON re.id = ar.raw_event_id + JOIN messages m ON m.source_system = 'chatwoot' AND m.id = re.message_id + """ + params = {"opportunity_id": opportunity_id, "limit": max(1, int(limit or 20))} + possible = conn.execute(text(f""" + WITH opp AS ( + SELECT id, conversation_id FROM opportunities + WHERE id = CAST(:opportunity_id AS UUID) LIMIT 1 + ), candidate_message_ids AS ({candidate_sql}) + SELECT EXISTS (SELECT 1 FROM candidate_message_ids LIMIT 1) + """), params).scalar() + if not possible: + return [] + + detail_sql = """ WITH opp AS ( SELECT id, conversation_id, contact_id FROM opportunities WHERE id = CAST(:opportunity_id AS UUID) LIMIT 1 + ), candidate_message_ids AS ( + __CANDIDATE_SQL__ + ), candidate_messages AS ( + SELECT m.* FROM messages m + JOIN candidate_message_ids candidate ON candidate.message_id = m.id + ), message_raw_pairs AS ( + SELECT m.id AS message_id, re.id AS raw_event_id + FROM candidate_messages m JOIN raw_events re ON re.id = m.raw_event_id + UNION + SELECT m.id, re.id + FROM candidate_messages m JOIN raw_events re ON re.message_id = m.id + ), raw_context AS ( + SELECT m.id AS message_id, pair.raw_event_id + FROM candidate_messages m + LEFT JOIN message_raw_pairs pair ON pair.message_id = m.id + ), action_pairs AS ( + SELECT context.message_id, context.raw_event_id, ar.id AS action_run_id + FROM raw_context context JOIN action_runs ar ON ar.message_id = context.message_id + UNION + SELECT context.message_id, context.raw_event_id, ar.id + FROM raw_context context JOIN action_runs ar ON ar.raw_event_id = context.raw_event_id + ), action_context AS ( + SELECT context.message_id, context.raw_event_id, pair.action_run_id + FROM raw_context context + LEFT JOIN action_pairs pair + ON pair.message_id = context.message_id + AND pair.raw_event_id IS NOT DISTINCT FROM context.raw_event_id + ), task_pairs AS ( + SELECT context.message_id, context.raw_event_id, context.action_run_id, t.id AS task_id + FROM action_context context JOIN tasks t ON t.message_id = context.message_id + UNION + SELECT context.message_id, context.raw_event_id, context.action_run_id, t.id + FROM action_context context JOIN tasks t ON t.raw_event_id = context.raw_event_id + UNION + SELECT context.message_id, context.raw_event_id, context.action_run_id, t.id + FROM action_context context JOIN tasks t ON t.action_run_id = context.action_run_id + ), full_context AS ( + SELECT context.message_id, context.raw_event_id, context.action_run_id, pair.task_id + FROM action_context context + LEFT JOIN task_pairs pair + ON pair.message_id = context.message_id + AND pair.raw_event_id IS NOT DISTINCT FROM context.raw_event_id + AND pair.action_run_id IS NOT DISTINCT FROM context.action_run_id ), ranked AS ( SELECT DISTINCT ON (m.id) m.id::text AS id, @@ -271,10 +373,11 @@ def _list_message_backed_chatwoot_items_for_opportunity(opportunity_id: str, lim m.created_at, m.created_at AS updated_at FROM opp - JOIN messages m ON m.source_system = 'chatwoot' - LEFT JOIN raw_events re ON re.id = m.raw_event_id OR re.message_id = m.id - LEFT JOIN action_runs ar ON ar.message_id = m.id OR ar.raw_event_id = re.id - LEFT JOIN tasks t ON t.message_id = m.id OR t.raw_event_id = re.id OR t.action_run_id = ar.id + JOIN full_context context ON TRUE + JOIN candidate_messages m ON m.id = context.message_id + LEFT JOIN raw_events re ON re.id = context.raw_event_id + LEFT JOIN action_runs ar ON ar.id = context.action_run_id + LEFT JOIN tasks t ON t.id = context.task_id WHERE t.opportunity_id = opp.id OR ( @@ -287,7 +390,8 @@ def _list_message_backed_chatwoot_items_for_opportunity(opportunity_id: str, lim FROM ranked ORDER BY created_at DESC LIMIT :limit - """), {"opportunity_id": opportunity_id, "limit": max(1, int(limit or 20))}).mappings().all() + """.replace("__CANDIDATE_SQL__", candidate_sql) + rows = conn.execute(text(detail_sql), params).mappings().all() return [dict(row) for row in rows] diff --git a/app/db.py b/app/db.py index ab8d30f..1b71524 100644 --- a/app/db.py +++ b/app/db.py @@ -295,6 +295,13 @@ def ensure_core_schema() -> None: conn.execute(text("CREATE INDEX IF NOT EXISTS idx_raw_events_created_at ON raw_events(created_at DESC)")) 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 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)")) + conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_message ON tasks(message_id)")) + conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_raw_event ON tasks(raw_event_id)")) + conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_action_run ON tasks(action_run_id)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_action_runs_created_at ON action_runs(created_at DESC)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_task_events_task ON task_events(task_id, created_at DESC)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_business_events_created_at ON business_events(created_at DESC)")) diff --git a/migrations/008_historical_communication_lookup_indexes.sql b/migrations/008_historical_communication_lookup_indexes.sql new file mode 100644 index 0000000..d9a5d56 --- /dev/null +++ b/migrations/008_historical_communication_lookup_indexes.sql @@ -0,0 +1,9 @@ +-- Index the individual relationship paths used by the historical Chatwoot +-- fallback. Conversation, opportunity and primary-key indexes already exist. +CREATE INDEX IF NOT EXISTS idx_messages_raw_event ON messages(raw_event_id); +CREATE INDEX IF NOT EXISTS idx_raw_events_message ON raw_events(message_id); +CREATE INDEX IF NOT EXISTS idx_action_runs_message ON action_runs(message_id); +CREATE INDEX IF NOT EXISTS idx_action_runs_raw_event ON action_runs(raw_event_id); +CREATE INDEX IF NOT EXISTS idx_tasks_message ON tasks(message_id); +CREATE INDEX IF NOT EXISTS idx_tasks_raw_event ON tasks(raw_event_id); +CREATE INDEX IF NOT EXISTS idx_tasks_action_run ON tasks(action_run_id); diff --git a/tests/test_historical_communication_fallback_performance.py b/tests/test_historical_communication_fallback_performance.py new file mode 100644 index 0000000..a450c96 --- /dev/null +++ b/tests/test_historical_communication_fallback_performance.py @@ -0,0 +1,99 @@ +from contextlib import nullcontext +from pathlib import Path + +import app.communication_service as service + + +class _Result: + def __init__(self, *, scalar_value=None, rows=()): + self.scalar_value = scalar_value + self.rows = list(rows) + + def scalar(self): + return self.scalar_value + + def mappings(self): + return self + + def all(self): + return self.rows + + +class _Connection: + def __init__(self, *, possible, rows=()): + self.possible = possible + self.rows = rows + self.statements = [] + + def execute(self, statement, params=None): + self.statements.append((str(statement), params)) + if len(self.statements) == 1: + return _Result(scalar_value=self.possible) + return _Result(rows=self.rows) + + +class _Engine: + def __init__(self, connection): + self.connection = connection + + def begin(self): + return nullcontext(self.connection) + + +def test_no_historical_candidates_returns_after_one_statement(monkeypatch): + connection = _Connection(possible=False) + monkeypatch.setattr(service, "engine", _Engine(connection)) + + assert service._list_message_backed_chatwoot_items_for_opportunity( + "dc3c020e-33eb-4067-b7da-ab684533a425" + ) == [] + assert len(connection.statements) == 1 + assert "SELECT EXISTS" in connection.statements[0][0] + + +def test_positive_fallback_uses_two_statements_and_preserves_rows(monkeypatch): + expected = [{"id": "message-1", "source_system": "chatwoot", "created_at": "2026-08-14"}] + connection = _Connection(possible=True, rows=expected) + monkeypatch.setattr(service, "engine", _Engine(connection)) + + result = service._list_message_backed_chatwoot_items_for_opportunity( + "dc3c020e-33eb-4067-b7da-ab684533a425", limit=12 + ) + + assert result == expected + assert len(connection.statements) == 2 + assert connection.statements[1][1]["limit"] == 12 + + +def test_relationship_paths_are_unioned_without_or_join_conditions(monkeypatch): + connection = _Connection(possible=True) + monkeypatch.setattr(service, "engine", _Engine(connection)) + service._list_message_backed_chatwoot_items_for_opportunity( + "dc3c020e-33eb-4067-b7da-ab684533a425" + ) + sql = connection.statements[1][0] + + assert "re.id = m.raw_event_id OR re.message_id = m.id" not in sql + assert "ar.message_id = m.id OR ar.raw_event_id = re.id" not in sql + assert "t.message_id = m.id OR t.raw_event_id = re.id" not in sql + assert sql.count(" UNION\n") >= 10 + assert "SELECT DISTINCT ON (m.id)" in sql + assert "ORDER BY m.id, COALESCE(t.created_at, m.created_at) DESC" in sql + assert "ORDER BY created_at DESC" in sql + + +def test_only_missing_relationship_indexes_are_added(): + migration = Path("migrations/008_historical_communication_lookup_indexes.sql").read_text() + expected = { + "messages(raw_event_id)", + "raw_events(message_id)", + "action_runs(message_id)", + "action_runs(raw_event_id)", + "tasks(message_id)", + "tasks(raw_event_id)", + "tasks(action_run_id)", + } + assert all(columns in migration for columns in expected) + assert "messages(conversation_id)" not in migration + assert "raw_events(conversation_id)" not in migration + assert "tasks(opportunity_id)" not in migration