From d15fcf71290973114516766223a549f23ed9b28a Mon Sep 17 00:00:00 2001 From: Yera All Date: Sat, 4 Apr 2026 04:10:09 +0500 Subject: [PATCH] Add demo data prune job --- .gitlab-ci.yml | 25 ++ scripts/prune_demo_interactions.py | 378 +++++++++++++++++++++++++++++ 2 files changed, 403 insertions(+) create mode 100644 scripts/prune_demo_interactions.py diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index 1c14c1d..a601d2c 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -249,3 +249,28 @@ hotfix-analyst-assets: when: manual - when: never allow_failure: false + +prune-demo-data-production: + stage: deploy + tags: + - call-center-prod + needs: + - deploy-production-source + variables: + DOCKER_HOST: "" + DOCKER_TLS_VERIFY: "" + DOCKER_CERT_PATH: "" + DOCKER_TLS_CERTDIR: "" + DEPLOY_DIR: /home/gitlab-runner/deploy/call-center + KEEP_LATEST_INTERACTIONS: "6" + script: + - command -v python3 >/dev/null + - test -d "$DEPLOY_DIR" + - test -f "$DEPLOY_DIR/.env.production" + - mkdir -p "$DEPLOY_DIR/.db_backups" + - python3 scripts/prune_demo_interactions.py --env-file "$DEPLOY_DIR/.env.production" --base-dir "$DEPLOY_DIR" --backup-dir "$DEPLOY_DIR/.db_backups" --keep-latest-interactions "$KEEP_LATEST_INTERACTIONS" + rules: + - if: '$CI_COMMIT_BRANCH == $CI_DEFAULT_BRANCH' + when: manual + - when: never + allow_failure: false diff --git a/scripts/prune_demo_interactions.py b/scripts/prune_demo_interactions.py new file mode 100644 index 0000000..5605e01 --- /dev/null +++ b/scripts/prune_demo_interactions.py @@ -0,0 +1,378 @@ +#!/usr/bin/env python3 + +from __future__ import annotations + +import argparse +import sqlite3 +from collections.abc import Iterable +from dataclasses import dataclass +from datetime import datetime +from pathlib import Path + + +DEFAULT_KEEP_LATEST_INTERACTIONS = 6 + + +@dataclass +class PrunePlan: + keep_interaction_ids: list[str] + delete_interaction_ids: list[str] + delete_thread_ids: list[str] + delete_whatsapp_thread_ids: list[str] + delete_call_ids: list[str] + delete_ai_session_ids: list[str] + delete_voice_session_ids: list[str] + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser(description="Prune demo interactions from a SQLite call-center database.") + parser.add_argument("--env-file", type=Path, default=None, help="Path to a deployment env file.") + parser.add_argument("--base-dir", type=Path, default=None, help="Base directory for resolving relative sqlite paths.") + parser.add_argument("--db-path", type=Path, default=None, help="Override sqlite database path directly.") + parser.add_argument( + "--backup-dir", + type=Path, + default=None, + help="Directory where a SQLite backup should be written before deleting rows.", + ) + parser.add_argument( + "--keep-latest-interactions", + type=int, + default=DEFAULT_KEEP_LATEST_INTERACTIONS, + help="How many latest interactions should remain after pruning.", + ) + parser.add_argument("--dry-run", action="store_true", help="Print the plan without changing the database.") + return parser.parse_args() + + +def load_env_file(path: Path | None) -> dict[str, str]: + if not path or not path.exists(): + return {} + values: dict[str, str] = {} + for raw_line in path.read_text(encoding="utf-8").splitlines(): + line = raw_line.strip() + if not line or line.startswith("#") or "=" not in line: + continue + key, value = line.split("=", 1) + values[key.strip()] = value.strip().strip('"').strip("'") + return values + + +def resolve_db_path(args: argparse.Namespace) -> Path: + if args.db_path: + return args.db_path.resolve() + + env_values = load_env_file(args.env_file) + base_dir = args.base_dir.resolve() if args.base_dir else (args.env_file.resolve().parent if args.env_file else Path.cwd()) + database_url = env_values.get("DATABASE_URL", "").strip() + + if database_url.startswith("sqlite:///"): + raw_path = database_url[len("sqlite:///") :] + if raw_path.startswith("/"): + return Path(raw_path).resolve() + return (base_dir / raw_path).resolve() + + cc_data_dir = env_values.get("CC_DATA_DIR", ".data_local").strip() or ".data_local" + return (base_dir / cc_data_dir / "mvp_cc.db").resolve() + + +def current_timestamp() -> str: + return datetime.now().strftime("%Y%m%d_%H%M%S") + + +def table_exists(conn: sqlite3.Connection, table: str) -> bool: + row = conn.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' AND name = ?", + (table,), + ).fetchone() + return bool(row) + + +def column_exists(conn: sqlite3.Connection, table: str, column: str) -> bool: + if not table_exists(conn, table): + return False + rows = conn.execute(f"PRAGMA table_info({table})").fetchall() + return any(row[1] == column for row in rows) + + +def fetch_distinct_values(conn: sqlite3.Connection, table: str, column: str, interaction_ids: list[str]) -> list[str]: + if not interaction_ids or not column_exists(conn, table, column) or not column_exists(conn, table, "interaction_id"): + return [] + placeholders = ",".join("?" for _ in interaction_ids) + rows = conn.execute( + f""" + SELECT DISTINCT {column} + FROM {table} + WHERE interaction_id IN ({placeholders}) AND {column} IS NOT NULL AND {column} != '' + """, + interaction_ids, + ).fetchall() + return [str(row[0]) for row in rows if row[0]] + + +def fetch_related_session_ids( + conn: sqlite3.Connection, + *, + interaction_ids: list[str], + thread_ids: list[str], + call_ids: list[str], +) -> list[str]: + if not table_exists(conn, "ai_sessions"): + return [] + clauses: list[str] = [] + params: list[str] = [] + if interaction_ids and column_exists(conn, "ai_sessions", "interaction_id"): + placeholders = ",".join("?" for _ in interaction_ids) + clauses.append(f"interaction_id IN ({placeholders})") + params.extend(interaction_ids) + if thread_ids and column_exists(conn, "ai_sessions", "thread_id"): + placeholders = ",".join("?" for _ in thread_ids) + clauses.append(f"thread_id IN ({placeholders})") + params.extend(thread_ids) + if call_ids and column_exists(conn, "ai_sessions", "call_id"): + placeholders = ",".join("?" for _ in call_ids) + clauses.append(f"call_id IN ({placeholders})") + params.extend(call_ids) + if not clauses: + return [] + rows = conn.execute( + f"SELECT DISTINCT session_id FROM ai_sessions WHERE {' OR '.join(clauses)}", + params, + ).fetchall() + return [str(row[0]) for row in rows if row[0]] + + +def fetch_related_voice_session_ids(conn: sqlite3.Connection, *, interaction_ids: list[str], call_ids: list[str]) -> list[str]: + if not table_exists(conn, "voice_ai_sessions"): + return [] + clauses: list[str] = [] + params: list[str] = [] + if interaction_ids and column_exists(conn, "voice_ai_sessions", "interaction_id"): + placeholders = ",".join("?" for _ in interaction_ids) + clauses.append(f"interaction_id IN ({placeholders})") + params.extend(interaction_ids) + if call_ids and column_exists(conn, "voice_ai_sessions", "call_id"): + placeholders = ",".join("?" for _ in call_ids) + clauses.append(f"call_id IN ({placeholders})") + params.extend(call_ids) + if not clauses: + return [] + rows = conn.execute( + f"SELECT DISTINCT session_id FROM voice_ai_sessions WHERE {' OR '.join(clauses)}", + params, + ).fetchall() + return [str(row[0]) for row in rows if row[0]] + + +def build_plan(conn: sqlite3.Connection, keep_latest_interactions: int) -> PrunePlan: + keep_limit = max(keep_latest_interactions, 0) + keep_rows = conn.execute( + """ + SELECT interaction_id + FROM interactions + ORDER BY updated_at DESC, created_at DESC, id DESC + LIMIT ? + """, + (keep_limit,), + ).fetchall() + keep_interaction_ids = [str(row[0]) for row in keep_rows if row[0]] + + if keep_interaction_ids: + placeholders = ",".join("?" for _ in keep_interaction_ids) + delete_rows = conn.execute( + f""" + SELECT interaction_id + FROM interactions + WHERE interaction_id NOT IN ({placeholders}) + ORDER BY updated_at DESC, created_at DESC, id DESC + """, + keep_interaction_ids, + ).fetchall() + else: + delete_rows = conn.execute( + """ + SELECT interaction_id + FROM interactions + ORDER BY updated_at DESC, created_at DESC, id DESC + """ + ).fetchall() + + delete_interaction_ids = [str(row[0]) for row in delete_rows if row[0]] + delete_thread_ids = fetch_distinct_values(conn, "telegram_threads", "thread_id", delete_interaction_ids) + delete_whatsapp_thread_ids = fetch_distinct_values(conn, "whatsapp_threads", "thread_id", delete_interaction_ids) + delete_call_ids = fetch_distinct_values(conn, "asterisk_call_links", "call_id", delete_interaction_ids) + delete_ai_session_ids = fetch_related_session_ids( + conn, + interaction_ids=delete_interaction_ids, + thread_ids=delete_thread_ids + delete_whatsapp_thread_ids, + call_ids=delete_call_ids, + ) + delete_voice_session_ids = fetch_related_voice_session_ids( + conn, + interaction_ids=delete_interaction_ids, + call_ids=delete_call_ids, + ) + return PrunePlan( + keep_interaction_ids=keep_interaction_ids, + delete_interaction_ids=delete_interaction_ids, + delete_thread_ids=delete_thread_ids, + delete_whatsapp_thread_ids=delete_whatsapp_thread_ids, + delete_call_ids=delete_call_ids, + delete_ai_session_ids=delete_ai_session_ids, + delete_voice_session_ids=delete_voice_session_ids, + ) + + +def delete_by_values(conn: sqlite3.Connection, table: str, column: str, values: Iterable[str]) -> int: + value_list = [str(value) for value in values if str(value)] + if not value_list or not column_exists(conn, table, column): + return 0 + placeholders = ",".join("?" for _ in value_list) + cursor = conn.execute(f"DELETE FROM {table} WHERE {column} IN ({placeholders})", value_list) + return cursor.rowcount if cursor.rowcount is not None else 0 + + +def count_rows(conn: sqlite3.Connection, table: str) -> int: + if not table_exists(conn, table): + return 0 + row = conn.execute(f"SELECT COUNT(*) FROM {table}").fetchone() + return int(row[0] or 0) + + +def backup_database(db_path: Path, backup_dir: Path) -> Path: + backup_dir.mkdir(parents=True, exist_ok=True) + backup_path = backup_dir / f"{db_path.stem}_backup_{current_timestamp()}{db_path.suffix or '.db'}" + source = sqlite3.connect(db_path) + target = sqlite3.connect(backup_path) + try: + source.backup(target) + finally: + target.close() + source.close() + return backup_path + + +def prune_database(conn: sqlite3.Connection, plan: PrunePlan) -> dict[str, int]: + deleted: dict[str, int] = {} + + interaction_tables = [ + "interaction_timelines", + "telegram_messages", + "telegram_threads", + "whatsapp_messages", + "whatsapp_threads", + "asterisk_call_links", + "voice_events", + "call_recordings", + "voice_transcript_segments", + "ivr_sessions", + "ai_sessions", + "voice_ai_sessions", + "ai_turns", + "asterisk_call_action_log", + "asterisk_event_log", + ] + for table in interaction_tables: + deleted[table] = delete_by_values(conn, table, "interaction_id", plan.delete_interaction_ids) + + deleted["telegram_messages_by_thread"] = delete_by_values(conn, "telegram_messages", "thread_id", plan.delete_thread_ids) + deleted["whatsapp_messages_by_thread"] = delete_by_values(conn, "whatsapp_messages", "thread_id", plan.delete_whatsapp_thread_ids) + + call_tables = [ + "asterisk_call_links", + "voice_events", + "call_recordings", + "voice_transcript_segments", + "voice_ai_sessions", + "ai_sessions", + "asterisk_call_action_log", + "asterisk_event_log", + ] + for table in call_tables: + key = f"{table}_by_call" + deleted[key] = delete_by_values(conn, table, "call_id", plan.delete_call_ids) + + deleted["ai_jobs_by_thread"] = delete_by_values(conn, "ai_jobs", "thread_id", plan.delete_thread_ids + plan.delete_whatsapp_thread_ids) + deleted["ai_jobs_by_session"] = delete_by_values(conn, "ai_jobs", "session_id", plan.delete_ai_session_ids) + deleted["ai_turns_by_thread"] = delete_by_values(conn, "ai_turns", "thread_id", plan.delete_thread_ids + plan.delete_whatsapp_thread_ids) + deleted["ai_turns_by_session"] = delete_by_values( + conn, + "ai_turns", + "session_id", + plan.delete_ai_session_ids + plan.delete_voice_session_ids, + ) + deleted["voice_transcript_segments_by_session"] = delete_by_values( + conn, + "voice_transcript_segments", + "session_id", + plan.delete_voice_session_ids, + ) + deleted["interactions"] = delete_by_values(conn, "interactions", "interaction_id", plan.delete_interaction_ids) + return deleted + + +def main() -> int: + args = parse_args() + db_path = resolve_db_path(args) + if not db_path.exists(): + raise SystemExit(f"Database file not found: {db_path}") + + backup_dir = (args.backup_dir or (db_path.parent.parent / ".db_backups")).resolve() + conn = sqlite3.connect(db_path) + conn.row_factory = sqlite3.Row + try: + before_counts = { + "interactions": count_rows(conn, "interactions"), + "telegram_threads": count_rows(conn, "telegram_threads"), + "telegram_messages": count_rows(conn, "telegram_messages"), + "whatsapp_threads": count_rows(conn, "whatsapp_threads"), + "whatsapp_messages": count_rows(conn, "whatsapp_messages"), + "asterisk_call_links": count_rows(conn, "asterisk_call_links"), + } + plan = build_plan(conn, args.keep_latest_interactions) + print(f"Database: {db_path}") + print(f"Interactions before prune: {before_counts['interactions']}") + print(f"Keeping latest interactions: {len(plan.keep_interaction_ids)}") + print(f"Deleting interactions: {len(plan.delete_interaction_ids)}") + + if args.dry_run or not plan.delete_interaction_ids: + if not plan.delete_interaction_ids: + print("Nothing to prune.") + else: + print("Dry run only. No rows deleted.") + return 0 + + backup_path = backup_database(db_path, backup_dir) + print(f"Backup created: {backup_path}") + + conn.execute("PRAGMA foreign_keys = OFF") + conn.execute("BEGIN") + deleted = prune_database(conn, plan) + conn.commit() + + after_counts = { + "interactions": count_rows(conn, "interactions"), + "telegram_threads": count_rows(conn, "telegram_threads"), + "telegram_messages": count_rows(conn, "telegram_messages"), + "whatsapp_threads": count_rows(conn, "whatsapp_threads"), + "whatsapp_messages": count_rows(conn, "whatsapp_messages"), + "asterisk_call_links": count_rows(conn, "asterisk_call_links"), + } + + print("Deleted rows:") + for table, rowcount in deleted.items(): + if rowcount: + print(f" {table}: {rowcount}") + print("Row counts after prune:") + for table, value in after_counts.items(): + print(f" {table}: {value}") + return 0 + except Exception: + conn.rollback() + raise + finally: + conn.close() + + +if __name__ == "__main__": + raise SystemExit(main())