Add dry-run sqlite retention prune for tool/trace bloat.

Keeps recent days of chat tool messages, tool_log, traces, and sent outbound; admin API and CLI support apply + optional VACUUM.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-11 00:22:26 +08:00
parent cc9a516cb4
commit 48e5872045
4 changed files with 619 additions and 0 deletions

View file

@ -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),

View file

@ -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())

View file

@ -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",
]

View file

@ -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()