diff --git a/interfaces/admin/routes.py b/interfaces/admin/routes.py index 2836a80a..300819e4 100644 --- a/interfaces/admin/routes.py +++ b/interfaces/admin/routes.py @@ -1104,6 +1104,55 @@ def build_admin_router() -> APIRouter: "max_age_days": int(max_age_days), } + @router.post("/admin/api/runtime/sqlite-retention/prune") + def api_runtime_sqlite_retention_prune( + payload: dict[str, Any] | None = Body(default=None), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + """Dry-run or apply retention prune for tool/trace/outbound noise. + + Body: + keep_days (default 30), dry_run (default true), vacuum (default false), + include_scheduled_runs (default false), include_outbound (default true). + Apply requires dry_run=false and confirm=true. + """ + from svc.persistence.sqlite_retention import prune_sqlite_retention + + body = payload or {} + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + try: + keep_days = int(body.get("keep_days", 30)) + except Exception: + keep_days = 30 + dry_run = body.get("dry_run", True) + if isinstance(dry_run, str): + dry_run = dry_run.strip().lower() not in {"0", "false", "no", "off"} + else: + dry_run = bool(dry_run) + vacuum = bool(body.get("vacuum", False)) + include_scheduled_runs = bool(body.get("include_scheduled_runs", False)) + include_outbound = body.get("include_outbound", True) + if isinstance(include_outbound, str): + include_outbound = include_outbound.strip().lower() not in {"0", "false", "no", "off"} + else: + include_outbound = bool(include_outbound) + if not dry_run and not bool(body.get("confirm", False)): + return { + "ok": False, + "error": "confirm_required", + "hint": "Set dry_run=false and confirm=true to delete. Prefer dry_run first.", + } + return prune_sqlite_retention( + store, + keep_days=keep_days, + dry_run=dry_run, + vacuum=vacuum and not dry_run, + include_outbound_sent=include_outbound, + include_scheduled_runs=include_scheduled_runs, + ) + @router.get("/admin/api/tenants") def api_tenants( scope: str | None = Query(default=None), diff --git a/runtime/operations/scripts/prune_sqlite_retention.py b/runtime/operations/scripts/prune_sqlite_retention.py new file mode 100644 index 00000000..f8e9ed5e --- /dev/null +++ b/runtime/operations/scripts/prune_sqlite_retention.py @@ -0,0 +1,71 @@ +"""Prune aged tool/trace noise from the assistant SQLite (or Postgres) DB. + +Dry-run (default):: + + python runtime/operations/scripts/prune_sqlite_retention.py --keep-days 30 + +Apply delete (keeps last N days of tool messages / tool_log / traces / sent outbound):: + + python runtime/operations/scripts/prune_sqlite_retention.py --keep-days 30 --apply --yes + +Also reclaim file space on SQLite (needs ~DB-size free disk; stop writers if possible):: + + python runtime/operations/scripts/prune_sqlite_retention.py --keep-days 30 --apply --yes --vacuum + +Point at a specific file:: + + set OPS_ASSISTANT_DB_PATH=D:\\path\\ai_ops.sqlite + python runtime/operations/scripts/prune_sqlite_retention.py --keep-days 14 --dry-run +""" + +from __future__ import annotations + +import argparse +import json +import os +import sys + + +def main() -> int: + p = argparse.ArgumentParser(description=__doc__) + p.add_argument("--keep-days", type=int, default=30, help="Retain rows newer than this many days.") + p.add_argument("--dry-run", action="store_true", default=False, help="Force dry-run (default).") + p.add_argument("--apply", action="store_true", help="Actually delete (requires --yes).") + p.add_argument("--yes", action="store_true", help="Confirm destructive delete.") + p.add_argument("--vacuum", action="store_true", help="SQLite VACUUM after delete.") + p.add_argument("--include-scheduled-runs", action="store_true", help="Also prune old scheduled_job_run.") + p.add_argument("--no-outbound", action="store_true", help="Skip channel_outbound_message prune.") + p.add_argument("--db", default="", help="Override OPS_ASSISTANT_DB_PATH for this run.") + args = p.parse_args() + + if str(args.db or "").strip(): + os.environ["OPS_ASSISTANT_DB_PATH"] = str(args.db).strip() + os.environ["AIA_ASSISTANT_DB_BACKEND"] = "sqlite" + + from svc.persistence.assistant_store import get_assistant_store, reset_assistant_store_singleton + from svc.persistence.sqlite_retention import prune_sqlite_retention + + reset_assistant_store_singleton() + store = get_assistant_store() + + dry_run = not bool(args.apply) + if args.dry_run: + dry_run = True + if args.apply and not args.yes: + print("Refusing --apply without --yes", file=sys.stderr) + return 2 + + out = prune_sqlite_retention( + store, + keep_days=int(args.keep_days), + dry_run=dry_run, + vacuum=bool(args.vacuum) and not dry_run, + include_outbound_sent=not bool(args.no_outbound), + include_scheduled_runs=bool(args.include_scheduled_runs), + ) + print(json.dumps(out, ensure_ascii=False, indent=2)) + return 0 if out.get("ok") else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/svc/persistence/sqlite_retention.py b/svc/persistence/sqlite_retention.py new file mode 100644 index 00000000..4950274c --- /dev/null +++ b/svc/persistence/sqlite_retention.py @@ -0,0 +1,358 @@ +from __future__ import annotations + +import sqlite3 +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any + + +DEFAULT_KEEP_DAYS = 30 +MAX_KEEP_DAYS = 3650 +MIN_KEEP_DAYS = 1 + + +def _utcnow() -> datetime: + return datetime.now(timezone.utc) + + +def cutoff_iso(*, keep_days: int, now: datetime | None = None) -> str: + days = max(MIN_KEEP_DAYS, min(int(keep_days), MAX_KEEP_DAYS)) + base = now or _utcnow() + return (base - timedelta(days=days)).isoformat() + + +def _is_postgres_store(store: Any) -> bool: + if bool(getattr(store, "_use_pg", False)): + return True + url = str(getattr(store, "_postgres_url", "") or getattr(store, "db_path", "") or "") + return url.startswith("postgres://") or url.startswith("postgresql://") + + +def _table_exists(conn: Any, name: str, *, postgres: bool) -> bool: + if postgres: + try: + row = conn.execute("SELECT to_regclass(?)", (f"public.{name}",)).fetchone() + if not row: + return False + val = row[0] if not isinstance(row, dict) else next(iter(row.values())) + return bool(val) + except Exception: + return False + row = conn.execute( + "SELECT 1 FROM sqlite_master WHERE type='table' AND name=? LIMIT 1", + (name,), + ).fetchone() + return bool(row) + + +def _length_expr(col: str, *, postgres: bool) -> str: + if postgres: + return f"COALESCE(LENGTH(CAST({col} AS text)), 0)" + return f"COALESCE(LENGTH(COALESCE({col}, '')), 0)" + + +def _count_and_bytes( + conn: Any, + *, + sql: str, + params: tuple[Any, ...], +) -> tuple[int, int]: + row = conn.execute(sql, params).fetchone() + if not row: + return 0, 0 + if isinstance(row, dict): + vals = list(row.values()) + return int(vals[0] or 0), int(vals[1] or 0) + return int(row[0] or 0), int(row[1] or 0) + + +def plan_sqlite_retention( + store: Any, + *, + keep_days: int = DEFAULT_KEEP_DAYS, + include_tool_messages: bool = True, + include_tool_log: bool = True, + include_trace_events: bool = True, + include_outbound_sent: bool = True, + include_scheduled_runs: bool = False, + now: datetime | None = None, +) -> dict[str, Any]: + """Dry-run friendly plan: counts + approximate payload bytes to delete.""" + cutoff = cutoff_iso(keep_days=keep_days, now=now) + postgres = _is_postgres_store(store) + targets: dict[str, Any] = {} + total_rows = 0 + total_bytes = 0 + + with store._connect() as conn: + if include_tool_messages and _table_exists(conn, "chat_message", postgres=postgres): + le = _length_expr("content", postgres=postgres) + n, b = _count_and_bytes( + conn, + sql=( + f"SELECT COUNT(*), COALESCE(SUM({le}), 0) FROM chat_message " + "WHERE lower(role)='tool' AND timestamp < ?" + ), + params=(cutoff,), + ) + targets["chat_message_tool"] = {"rows": n, "bytes": b, "cutoff": cutoff} + total_rows += n + total_bytes += b + + if include_tool_log and _table_exists(conn, "tool_log", postgres=postgres): + le = ( + f"{_length_expr('args', postgres=postgres)} + {_length_expr('result', postgres=postgres)}" + ) + n, b = _count_and_bytes( + conn, + sql=( + f"SELECT COUNT(*), COALESCE(SUM({le}), 0) FROM tool_log " + "WHERE timestamp < ?" + ), + params=(cutoff,), + ) + targets["tool_log"] = {"rows": n, "bytes": b, "cutoff": cutoff} + total_rows += n + total_bytes += b + + if include_trace_events and _table_exists(conn, "trace_event", postgres=postgres): + le = _length_expr("payload", postgres=postgres) + n, b = _count_and_bytes( + conn, + sql=( + f"SELECT COUNT(*), COALESCE(SUM({le}), 0) FROM trace_event " + "WHERE timestamp < ?" + ), + params=(cutoff,), + ) + targets["trace_event"] = {"rows": n, "bytes": b, "cutoff": cutoff} + total_rows += n + total_bytes += b + + if include_outbound_sent and _table_exists( + conn, "channel_outbound_message", postgres=postgres + ): + le = _length_expr("text", postgres=postgres) + n, b = _count_and_bytes( + conn, + sql=( + f"SELECT COUNT(*), COALESCE(SUM({le}), 0) FROM channel_outbound_message " + "WHERE created_at < ? AND lower(status) IN ('sent', 'failed', 'dead')" + ), + params=(cutoff,), + ) + targets["channel_outbound_message"] = {"rows": n, "bytes": b, "cutoff": cutoff} + total_rows += n + total_bytes += b + + if include_scheduled_runs and _table_exists( + conn, "scheduled_job_run", postgres=postgres + ): + le = ( + f"{_length_expr('reply_text', postgres=postgres)} + " + f"{_length_expr('error', postgres=postgres)}" + ) + n, b = _count_and_bytes( + conn, + sql=( + f"SELECT COUNT(*), COALESCE(SUM({le}), 0) FROM scheduled_job_run " + "WHERE created_at < ? AND lower(status) NOT IN ('queued', 'running')" + ), + params=(cutoff,), + ) + targets["scheduled_job_run"] = {"rows": n, "bytes": b, "cutoff": cutoff} + total_rows += n + total_bytes += b + + db_size_bytes = None + if not postgres: + try: + db_size_bytes = int(Path(str(store.db_path)).stat().st_size) + except Exception: + db_size_bytes = None + + return { + "ok": True, + "dry_run": True, + "keep_days": max(MIN_KEEP_DAYS, min(int(keep_days), MAX_KEEP_DAYS)), + "cutoff": cutoff, + "backend": "postgresql" if postgres else "sqlite", + "db_path": str(getattr(store, "db_path", "") or ""), + "db_size_bytes": db_size_bytes, + "targets": targets, + "total_rows": total_rows, + "total_bytes": total_bytes, + "total_mb": round(total_bytes / 1024.0 / 1024.0, 2), + } + + +def apply_sqlite_retention( + store: Any, + *, + keep_days: int = DEFAULT_KEEP_DAYS, + include_tool_messages: bool = True, + include_tool_log: bool = True, + include_trace_events: bool = True, + include_outbound_sent: bool = True, + include_scheduled_runs: bool = False, + vacuum: bool = False, + now: datetime | None = None, +) -> dict[str, Any]: + """Delete aged noisy rows. Optional VACUUM (SQLite only; needs free disk ≈ DB size).""" + plan = plan_sqlite_retention( + store, + keep_days=keep_days, + include_tool_messages=include_tool_messages, + include_tool_log=include_tool_log, + include_trace_events=include_trace_events, + include_outbound_sent=include_outbound_sent, + include_scheduled_runs=include_scheduled_runs, + now=now, + ) + cutoff = str(plan["cutoff"]) + deleted: dict[str, int] = {} + postgres = _is_postgres_store(store) + + with store._connect() as conn: + if include_tool_messages and _table_exists(conn, "chat_message", postgres=postgres): + cur = conn.execute( + "DELETE FROM chat_message WHERE lower(role)='tool' AND timestamp < ?", + (cutoff,), + ) + deleted["chat_message_tool"] = int(getattr(cur, "rowcount", 0) or 0) + + if include_tool_log and _table_exists(conn, "tool_log", postgres=postgres): + cur = conn.execute("DELETE FROM tool_log WHERE timestamp < ?", (cutoff,)) + deleted["tool_log"] = int(getattr(cur, "rowcount", 0) or 0) + + if include_trace_events and _table_exists(conn, "trace_event", postgres=postgres): + cur = conn.execute("DELETE FROM trace_event WHERE timestamp < ?", (cutoff,)) + deleted["trace_event"] = int(getattr(cur, "rowcount", 0) or 0) + + if include_outbound_sent and _table_exists( + conn, "channel_outbound_message", postgres=postgres + ): + cur = conn.execute( + "DELETE FROM channel_outbound_message " + "WHERE created_at < ? AND lower(status) IN ('sent', 'failed', 'dead')", + (cutoff,), + ) + deleted["channel_outbound_message"] = int(getattr(cur, "rowcount", 0) or 0) + + if include_scheduled_runs and _table_exists( + conn, "scheduled_job_run", postgres=postgres + ): + cur = conn.execute( + "DELETE FROM scheduled_job_run " + "WHERE created_at < ? AND lower(status) NOT IN ('queued', 'running')", + (cutoff,), + ) + deleted["scheduled_job_run"] = int(getattr(cur, "rowcount", 0) or 0) + + vacuum_result: dict[str, Any] | None = None + if vacuum: + if postgres: + vacuum_result = { + "ok": False, + "skipped": True, + "reason": "postgresql_use_manual_vacuum", + } + else: + vacuum_result = _vacuum_sqlite(store) + + db_size_after = None + if not postgres: + try: + db_size_after = int(Path(str(store.db_path)).stat().st_size) + except Exception: + db_size_after = None + + return { + "ok": True, + "dry_run": False, + "keep_days": plan["keep_days"], + "cutoff": cutoff, + "backend": plan["backend"], + "db_path": plan["db_path"], + "plan": { + "total_rows": plan["total_rows"], + "total_bytes": plan["total_bytes"], + "total_mb": plan["total_mb"], + "targets": plan["targets"], + }, + "deleted": deleted, + "deleted_rows": int(sum(deleted.values())), + "db_size_bytes_before": plan.get("db_size_bytes"), + "db_size_bytes_after": db_size_after, + "vacuum": vacuum_result, + } + + +def _vacuum_sqlite(store: Any) -> dict[str, Any]: + path = str(getattr(store, "db_path", "") or "").strip() + if not path or path.startswith("postgres"): + return {"ok": False, "error": "not_sqlite"} + try: + dispose = getattr(store, "_dispose_engines", None) + if callable(dispose): + try: + dispose() + except Exception: + pass + conn = sqlite3.connect(path, timeout=120.0) + try: + conn.execute("VACUUM") + conn.commit() + finally: + conn.close() + size = int(Path(path).stat().st_size) + return {"ok": True, "db_size_bytes": size} + except Exception as exc: + return {"ok": False, "error": f"{type(exc).__name__}: {exc}"} + + +def prune_sqlite_retention( + store: Any, + *, + keep_days: int = DEFAULT_KEEP_DAYS, + dry_run: bool = True, + vacuum: bool = False, + include_tool_messages: bool = True, + include_tool_log: bool = True, + include_trace_events: bool = True, + include_outbound_sent: bool = True, + include_scheduled_runs: bool = False, + now: datetime | None = None, +) -> dict[str, Any]: + """Entry point used by admin API / CLI.""" + if dry_run: + return plan_sqlite_retention( + store, + keep_days=keep_days, + include_tool_messages=include_tool_messages, + include_tool_log=include_tool_log, + include_trace_events=include_trace_events, + include_outbound_sent=include_outbound_sent, + include_scheduled_runs=include_scheduled_runs, + now=now, + ) + return apply_sqlite_retention( + store, + keep_days=keep_days, + include_tool_messages=include_tool_messages, + include_tool_log=include_tool_log, + include_trace_events=include_trace_events, + include_outbound_sent=include_outbound_sent, + include_scheduled_runs=include_scheduled_runs, + vacuum=vacuum, + now=now, + ) + + +__all__ = [ + "DEFAULT_KEEP_DAYS", + "apply_sqlite_retention", + "cutoff_iso", + "plan_sqlite_retention", + "prune_sqlite_retention", +] diff --git a/tests/test_sqlite_retention.py b/tests/test_sqlite_retention.py new file mode 100644 index 00000000..f021dd10 --- /dev/null +++ b/tests/test_sqlite_retention.py @@ -0,0 +1,141 @@ +from __future__ import annotations + +import os +import tempfile +import unittest +from datetime import datetime, timedelta, timezone +from pathlib import Path + +from svc.persistence.sqlite_retention import prune_sqlite_retention +from svc.persistence.sqlite_store import SqliteStore + + +class SqliteRetentionTests(unittest.TestCase): + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory(ignore_cleanup_errors=True) + self.db = Path(self._tmp.name) / "ret.sqlite" + os.environ["OPS_ASSISTANT_DB_PATH"] = str(self.db) + os.environ["AIA_ASSISTANT_DB_BACKEND"] = "sqlite" + self.store = SqliteStore(str(self.db)) + self.sess = self.store.create_session("retention-test") + self.session_id = str(self.sess.id) + self.now = datetime(2026, 8, 10, 12, 0, 0, tzinfo=timezone.utc) + self.old = (self.now - timedelta(days=40)).isoformat() + self.fresh = (self.now - timedelta(days=2)).isoformat() + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _seed(self) -> None: + with self.store._connect() as conn: + conn.execute( + "INSERT INTO chat_message(session_id, role, content, timestamp) VALUES (?,?,?,?)", + (self.session_id, "tool", "OLD_TOOL_PAYLOAD" * 100, self.old), + ) + conn.execute( + "INSERT INTO chat_message(session_id, role, content, timestamp) VALUES (?,?,?,?)", + (self.session_id, "tool", "FRESH_TOOL", self.fresh), + ) + conn.execute( + "INSERT INTO chat_message(session_id, role, content, timestamp) VALUES (?,?,?,?)", + (self.session_id, "assistant", "keep me", self.old), + ) + conn.execute( + "INSERT INTO tool_log(session_id, tool_name, specialist, args, result, timestamp) " + "VALUES (?,?,?,?,?,?)", + (self.session_id, "t1", "", "{}", '"OLD_RESULT"', self.old), + ) + conn.execute( + "INSERT INTO tool_log(session_id, tool_name, specialist, args, result, timestamp) " + "VALUES (?,?,?,?,?,?)", + (self.session_id, "t2", "", "{}", '"FRESH"', self.fresh), + ) + conn.execute( + "INSERT INTO trace_event(session_id, trace_id, span_id, parent_span_id, event_type, payload, timestamp) " + "VALUES (?,?,?,?,?,?,?)", + (self.session_id, "tr1", "sp1", "", "x", "{}", self.old), + ) + conn.execute( + "INSERT INTO channel_outbound_message" + "(id, tenant_id, channel, account_id, chat_id, text, status, source, created_at, error) " + "VALUES (?,?,?,?,?,?,?,?,?,?)", + ("ob-old", "", "whatsapp", "wa", "c", "hi", "sent", "", self.old, ""), + ) + conn.execute( + "INSERT INTO channel_outbound_message" + "(id, tenant_id, channel, account_id, chat_id, text, status, source, created_at, error) " + "VALUES (?,?,?,?,?,?,?,?,?,?)", + ("ob-pending", "", "whatsapp", "wa", "c", "hi", "pending", "", self.old, ""), + ) + + def test_dry_run_then_apply(self) -> None: + self._seed() + plan = prune_sqlite_retention(self.store, keep_days=30, dry_run=True, now=self.now) + self.assertTrue(plan.get("ok")) + self.assertTrue(plan.get("dry_run")) + self.assertGreaterEqual(int(plan.get("total_rows") or 0), 3) + self.assertGreater(int((plan.get("targets") or {}).get("chat_message_tool", {}).get("rows") or 0), 0) + + # Nothing deleted yet. + with self.store._connect() as conn: + n_tool = conn.execute( + "SELECT COUNT(*) FROM chat_message WHERE role='tool'" + ).fetchone()[0] + self.assertEqual(int(n_tool), 2) + + out = prune_sqlite_retention( + self.store, + keep_days=30, + dry_run=False, + vacuum=False, + now=self.now, + ) + self.assertTrue(out.get("ok")) + self.assertFalse(out.get("dry_run")) + self.assertGreaterEqual(int(out.get("deleted_rows") or 0), 3) + + with self.store._connect() as conn: + roles = [ + r[0] + for r in conn.execute("SELECT role FROM chat_message ORDER BY id").fetchall() + ] + self.assertIn("assistant", roles) + self.assertIn("tool", roles) # fresh tool kept + self.assertEqual(roles.count("tool"), 1) + tool_logs = conn.execute("SELECT COUNT(*) FROM tool_log").fetchone()[0] + self.assertEqual(int(tool_logs), 1) + traces = conn.execute("SELECT COUNT(*) FROM trace_event").fetchone()[0] + self.assertEqual(int(traces), 0) + outbound = { + r[0]: r[1] + for r in conn.execute( + "SELECT id, status FROM channel_outbound_message" + ).fetchall() + } + self.assertNotIn("ob-old", outbound) + self.assertIn("ob-pending", outbound) + + def test_vacuum_shrinks_after_large_delete(self) -> None: + self._seed() + # Inflate then delete + vacuum. + with self.store._connect() as conn: + conn.execute( + "INSERT INTO chat_message(session_id, role, content, timestamp) VALUES (?,?,?,?)", + (self.session_id, "tool", "X" * 200_000, self.old), + ) + before = self.db.stat().st_size + out = prune_sqlite_retention( + self.store, + keep_days=30, + dry_run=False, + vacuum=True, + now=self.now, + ) + self.assertTrue(out.get("ok")) + self.assertTrue((out.get("vacuum") or {}).get("ok")) + after = self.db.stat().st_size + self.assertLess(after, before) + + +if __name__ == "__main__": + unittest.main()