diff --git a/.env-sample b/.env-sample index 867ee46..dfd0045 100644 --- a/.env-sample +++ b/.env-sample @@ -36,6 +36,7 @@ DATABASE_URL=username:password@hostname:port # For PostgreSQL DATABASE_PATH=data/comet.db # Only relevant for SQLite DATABASE_BATCH_SIZE=20000 # The batch size for the database import and export operations DATABASE_READ_REPLICA_URLS='' # Optional JSON array of PostgreSQL read-only URLs, e.g. '["user:pass@replica-1/db", "user:pass@replica-2/db"]' +DATABASE_STARTUP_CLEANUP_INTERVAL=3600 # Minimum seconds between heavy startup cleanup sweeps (0=every start, -1=disable) # ============================== # # Cache Settings (Seconds) # diff --git a/CHANGELOG.md b/CHANGELOG.md index ee34921..01e96a9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,7 @@ ### Features * add optional PostgreSQL read replica routing with transparent primary fallback +* add `DATABASE_STARTUP_CLEANUP_INTERVAL` to throttle heavy startup cleanup sweeps across workers ## [2.31.0](https://github.com/g0ldyy/comet/compare/v2.30.0...v2.31.0) (2025-12-08) diff --git a/comet/core/database.py b/comet/core/database.py index a2eefad..b83350f 100644 --- a/comet/core/database.py +++ b/comet/core/database.py @@ -6,6 +6,8 @@ import traceback from comet.core.logger import logger from comet.core.models import database, settings +STARTUP_CLEANUP_LOCK_ID = 0xC0FFEE + DATABASE_VERSION = "1.0" @@ -28,6 +30,15 @@ async def setup_database(): """ ) + await database.execute( + """ + CREATE TABLE IF NOT EXISTS db_maintenance ( + id INTEGER PRIMARY KEY CHECK (id = 1), + last_startup_cleanup REAL + ) + """ + ) + current_version = await database.fetch_val( """ SELECT version FROM db_version WHERE id = 1 @@ -593,13 +604,41 @@ async def setup_database(): await database.execute("PRAGMA auto_vacuum=OFF") await database.execute("DELETE FROM ongoing_searches") + await database.execute("DELETE FROM active_connections") + await database.execute("DELETE FROM metrics_cache") + + await _run_startup_cleanup() + + except Exception as e: + logger.error(f"Error setting up the database: {e}") + logger.exception(traceback.format_exc()) + + +async def _run_startup_cleanup(): + interval = settings.DATABASE_STARTUP_CLEANUP_INTERVAL + if interval is None or interval < 0: + return + + current_time = time.time() + should_run = True if interval == 0 else await _should_run_startup_cleanup(current_time, interval) + if not should_run: + logger.log("DATABASE", "Startup cleanup skipped (recent run)") + return + + lock_acquired = await _try_acquire_startup_cleanup_lock() + if not lock_acquired: + logger.log("DATABASE", "Startup cleanup already running elsewhere; skipping") + return + + try: + logger.log("DATABASE", "Running startup cleanup sweep") await database.execute( """ DELETE FROM first_searches WHERE timestamp + :cache_ttl < :current_time; """, - {"cache_ttl": settings.TORRENT_CACHE_TTL, "current_time": time.time()}, + {"cache_ttl": settings.TORRENT_CACHE_TTL, "current_time": current_time}, ) await database.execute( @@ -607,7 +646,7 @@ async def setup_database(): DELETE FROM metadata_cache WHERE timestamp + :cache_ttl < :current_time; """, - {"cache_ttl": settings.METADATA_CACHE_TTL, "current_time": time.time()}, + {"cache_ttl": settings.METADATA_CACHE_TTL, "current_time": current_time}, ) if settings.TORRENT_CACHE_TTL >= 0: @@ -616,7 +655,7 @@ async def setup_database(): DELETE FROM torrents WHERE timestamp + :cache_ttl < :current_time; """, - {"cache_ttl": settings.TORRENT_CACHE_TTL, "current_time": time.time()}, + {"cache_ttl": settings.TORRENT_CACHE_TTL, "current_time": current_time}, ) await database.execute( @@ -624,18 +663,54 @@ async def setup_database(): DELETE FROM debrid_availability WHERE timestamp + :cache_ttl < :current_time; """, - {"cache_ttl": settings.DEBRID_CACHE_TTL, "current_time": time.time()}, + {"cache_ttl": settings.DEBRID_CACHE_TTL, "current_time": current_time}, ) await database.execute("DELETE FROM download_links_cache") - await database.execute("DELETE FROM active_connections") + await database.execute( + """ + INSERT INTO db_maintenance (id, last_startup_cleanup) + VALUES (1, :timestamp) + ON CONFLICT (id) DO UPDATE SET last_startup_cleanup = :timestamp + """, + {"timestamp": current_time}, + ) + finally: + if lock_acquired: + await _release_startup_cleanup_lock() - await database.execute("DELETE FROM metrics_cache") - except Exception as e: - logger.error(f"Error setting up the database: {e}") - logger.exception(traceback.format_exc()) +async def _should_run_startup_cleanup(current_time: float, interval: int) -> bool: + row = await database.fetch_one( + "SELECT last_startup_cleanup FROM db_maintenance WHERE id = 1" + ) + if not row or row["last_startup_cleanup"] is None: + return True + + last_run = float(row["last_startup_cleanup"]) + return (current_time - last_run) >= interval + + +async def _try_acquire_startup_cleanup_lock() -> bool: + if settings.DATABASE_TYPE != "postgresql": + return True + + result = await database.fetch_val( + "SELECT pg_try_advisory_lock(:lock_id)", + {"lock_id": STARTUP_CLEANUP_LOCK_ID}, + ) + return bool(result) + + +async def _release_startup_cleanup_lock(): + if settings.DATABASE_TYPE != "postgresql": + return + + await database.fetch_val( + "SELECT pg_advisory_unlock(:lock_id)", + {"lock_id": STARTUP_CLEANUP_LOCK_ID}, + ) async def cleanup_expired_locks(): diff --git a/comet/core/models.py b/comet/core/models.py index 8b40e0d..25fd2a8 100644 --- a/comet/core/models.py +++ b/comet/core/models.py @@ -35,6 +35,7 @@ class AppSettings(BaseSettings): DATABASE_PATH: Optional[str] = "data/comet.db" DATABASE_BATCH_SIZE: Optional[int] = 20000 DATABASE_READ_REPLICA_URLS: List[str] = Field(default_factory=list) + DATABASE_STARTUP_CLEANUP_INTERVAL: Optional[int] = 3600 METADATA_CACHE_TTL: Optional[int] = 2592000 # 30 days TORRENT_CACHE_TTL: Optional[int] = 1296000 # 15 days LIVE_TORRENT_CACHE_TTL: Optional[int] = 1296000 # 15 days