deerflow-code/offline-backend-20260512/backend/scripts/cleanup_sqlite_after_mysql_migration.py
2026-09-07 18:24:55 +08:00

265 lines
10 KiB
Python

"""Drop the stale business tables from the SQLite file after the MySQL switch.
Context
-------
In SQLite mode the LangGraph checkpointer **and** the application share one
file: ``.deer-flow/data/deerflow.db``. After ``database.backend`` is switched
to ``mysql`` (and ``migrate_business_sqlite_to_mysql.py`` has copied the data
over), the business tables in ``deerflow.db`` are dead weight — the app now
reads them from MySQL — while the LangGraph checkpoint tables in the *same*
file are still actively used (``checkpointer.type`` stays ``sqlite``).
This script drops only the stale business tables and ``VACUUM``s to reclaim
the freed pages.
Safety design
-------------
* **Explicit DROP allowlist** — only the known business tables are dropped.
LangGraph tables (``checkpoints`` / ``writes`` / ``store`` /
``store_migrations``) and any unrecognised table are left untouched.
* **MySQL pre-check** — before dropping, it connects to MySQL and verifies
every business table exists there. A missing table aborts the run (the
migration looks incomplete). Skippable only via an explicit flag.
* **Dry-run by default** — nothing is changed unless ``--apply`` is passed.
* **Automatic backup** — ``--apply`` first copies ``deerflow.db`` to a
timestamped ``.bak`` file (disable with ``--no-backup``).
IMPORTANT: stop the backend first. ``DROP TABLE`` / ``VACUUM`` need an
exclusive lock on the SQLite file; with the server running they will fail
with "database is locked" (a safe failure — nothing is changed).
Usage
-----
# from backend/ , with the backend's interpreter (has the MySQL driver)
PYTHONPATH=. uv run python scripts/cleanup_sqlite_after_mysql_migration.py # dry-run
PYTHONPATH=. uv run python scripts/cleanup_sqlite_after_mysql_migration.py --apply # do it
"""
from __future__ import annotations
import argparse
import asyncio
import os
import re
import shutil
import sqlite3
import sys
from datetime import datetime
from pathlib import Path
# Business tables migrated to MySQL — identical to BUSINESS_TABLES in
# migrate_business_sqlite_to_mysql.py. These are the DROP allowlist.
BUSINESS_TABLES: list[str] = [
"users",
"threads_meta",
"runs",
"run_events",
"feedback",
"agents",
"skills",
"scheduled_tasks",
"scheduled_task_subscriptions",
"scheduled_task_runs",
"scheduled_task_delivery_profiles",
"scheduler_threads",
"llm_call_metrics",
"notifications",
"recommended_questions",
"tags",
"tag_assignments",
"thread_shares",
]
# Newer business tables that may exist depending on deployment age. Dropped if
# present; absence is fine. They are NOT checked against MySQL (a fresh MySQL
# install creates them via create_all / Alembic, but an older SQLite file may
# never have had them).
OPTIONAL_BUSINESS_TABLES: list[str] = [
"ai_writing_sessions",
"article_types",
"tool_call_metrics",
]
# SQLite-side Alembic version tracker. After the MySQL switch the SQLite file
# only serves the checkpointer, which does not use Alembic — so this row is
# stale. Dropped, but never MySQL-checked.
ALEMBIC_TABLE = "alembic_version"
# LangGraph infrastructure tables — these stay on SQLite and must NEVER be
# dropped. The checkpointer/store keep using them.
KEEP_TABLES: set[str] = {"checkpoints", "writes", "store", "store_migrations"}
def _repo_root_from_script() -> Path:
# scripts/ -> backend/ -> offline-backend-20260512/
return Path(__file__).resolve().parents[2]
def _load_mysql_url(explicit: str | None) -> str | None:
if explicit:
return explicit
env = os.environ.get("MYSQL_DATABASE_URL")
if env:
return env
# Fall back to the .env next to the backend directory.
for candidate in (
_repo_root_from_script() / ".env",
Path(__file__).resolve().parents[1] / ".env",
):
if candidate.is_file():
for line in candidate.read_text(encoding="utf-8").splitlines():
if line.startswith("MYSQL_DATABASE_URL="):
return line.split("=", 1)[1].strip().strip('"').strip("'")
return None
def _normalize_async_url(url: str) -> str:
"""Force the URL onto the asyncmy async driver SQLAlchemy understands."""
if url.startswith(("mysql+asyncmy://", "mysql+aiomysql://")):
return url
return re.sub(r"^mysql(\+\w+)?://", "mysql+asyncmy://", url)
def _default_sqlite_path() -> Path:
return _repo_root_from_script() / "backend" / ".deer-flow" / "data" / "deerflow.db"
async def _mysql_existing_tables(url: str) -> set[str]:
"""Return the set of table names present in the MySQL business database."""
from sqlalchemy import text
from sqlalchemy.engine import make_url
from sqlalchemy.ext.asyncio import create_async_engine
async_url = _normalize_async_url(url)
db_name = make_url(async_url).database
engine = create_async_engine(async_url)
try:
async with engine.connect() as conn:
rows = await conn.execute(
text("SELECT table_name FROM information_schema.tables WHERE table_schema = :s"),
{"s": db_name},
)
return {r[0] for r in rows}
finally:
await engine.dispose()
def _sqlite_tables(con: sqlite3.Connection) -> list[str]:
return [r[0] for r in con.execute("SELECT name FROM sqlite_master WHERE type='table' ORDER BY name")]
def _row_count(con: sqlite3.Connection, table: str) -> int:
return con.execute(f'SELECT COUNT(*) FROM "{table}"').fetchone()[0]
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
parser.add_argument("--sqlite", default=str(_default_sqlite_path()), help="Path to deerflow.db")
parser.add_argument("--mysql-url", default=None, help="MySQL URL (default: $MYSQL_DATABASE_URL or .env)")
parser.add_argument("--apply", action="store_true", help="Actually drop tables + VACUUM (default: dry-run)")
parser.add_argument("--no-backup", action="store_true", help="Skip the pre-apply backup copy")
parser.add_argument("--skip-mysql-check", action="store_true", help="Do NOT verify tables exist in MySQL (unsafe)")
parser.add_argument("--drop-alembic", action="store_true", help="Also drop the stale SQLite alembic_version table")
args = parser.parse_args()
sqlite_path = Path(args.sqlite).resolve()
if not sqlite_path.is_file():
print(f"ERROR: SQLite file not found: {sqlite_path}")
return 1
size_before = sqlite_path.stat().st_size
print(f"SQLite file : {sqlite_path}")
print(f"Size : {size_before / 1024 / 1024:.1f} MB")
print(f"Mode : {'APPLY' if args.apply else 'DRY-RUN'}\n")
con = sqlite3.connect(sqlite_path)
try:
present = set(_sqlite_tables(con))
drop_business = [t for t in BUSINESS_TABLES if t in present]
drop_optional = [t for t in OPTIONAL_BUSINESS_TABLES if t in present]
to_drop = drop_business + drop_optional
if args.drop_alembic and ALEMBIC_TABLE in present:
to_drop.append(ALEMBIC_TABLE)
keep = sorted(present & KEEP_TABLES)
known = set(BUSINESS_TABLES) | set(OPTIONAL_BUSINESS_TABLES) | KEEP_TABLES | {ALEMBIC_TABLE}
unknown = sorted(present - known)
missing_business = [t for t in BUSINESS_TABLES if t not in present]
if missing_business:
print(f"NOTE: business tables already absent from SQLite: {', '.join(missing_business)}\n")
print("KEEP (LangGraph checkpointer / store - never dropped):")
for t in keep:
print(f" + {t:<34}{_row_count(con, t):>10} rows")
if unknown:
print("\nUNKNOWN (not recognised - left untouched, review manually):")
for t in unknown:
print(f" ? {t:<34}{_row_count(con, t):>10} rows")
print("\nDROP (stale business tables):")
for t in to_drop:
print(f" - {t:<34}{_row_count(con, t):>10} rows")
if not to_drop:
print(" (nothing to drop - already clean)")
return 0
# --- MySQL safety pre-check -------------------------------------
if args.skip_mysql_check:
print("\nWARNING: --skip-mysql-check set - NOT verifying the MySQL side.")
else:
mysql_url = _load_mysql_url(args.mysql_url)
if not mysql_url:
print("\nERROR: no MySQL URL ($MYSQL_DATABASE_URL / .env / --mysql-url).")
print(" Pass --skip-mysql-check to bypass (only if you are certain).")
return 1
print("\nVerifying business tables exist in MySQL ...")
try:
mysql_tables = asyncio.run(_mysql_existing_tables(mysql_url))
except Exception as exc: # noqa: BLE001
print(f"ERROR: could not query MySQL: {exc}")
print(" Run this on the backend host (MySQL access is host-restricted),")
print(" or pass --skip-mysql-check to bypass (only if you are certain).")
return 1
missing_in_mysql = [t for t in drop_business if t not in mysql_tables]
if missing_in_mysql:
print(f"ABORT: these business tables are MISSING in MySQL: {', '.join(missing_in_mysql)}")
print(" The migration looks incomplete - not dropping anything.")
return 1
print(f"OK - all {len(drop_business)} business tables present in MySQL.")
if not args.apply:
print("\nDry-run only. Re-run with --apply to drop the tables and VACUUM.")
return 0
# --- apply ------------------------------------------------------
if not args.no_backup:
backup = sqlite_path.with_name(
f"{sqlite_path.name}.bak-{datetime.now():%Y%m%d-%H%M%S}"
)
print(f"\nBacking up -> {backup}")
shutil.copy2(sqlite_path, backup)
print("Dropping tables ...")
for t in to_drop:
con.execute(f'DROP TABLE IF EXISTS "{t}"')
con.commit()
print("VACUUM (reclaiming space) ...")
con.isolation_level = None # VACUUM cannot run inside a transaction
con.execute("VACUUM")
finally:
con.close()
size_after = sqlite_path.stat().st_size
print(
f"\nDone. Size {size_before / 1024 / 1024:.1f} MB "
f"-> {size_after / 1024 / 1024:.1f} MB "
f"(reclaimed {(size_before - size_after) / 1024 / 1024:.1f} MB)"
)
return 0
if __name__ == "__main__":
sys.exit(main())