|
|
| """Periodic cleanup and shutdown hooks.
|
|
|
| Two mechanisms keep the app tidy at runtime and on exit:
|
|
|
| **Background cleanup thread** (started by :func:`register`):
|
| - Sweeps expired cache rows every ``CLEANUP_INTERVAL_SECONDS`` (default 7200).
|
| - Purges old log files on the same interval (files older than
|
| ``LOG_RETENTION_DAYS`` days, default 7).
|
|
|
| **atexit hook** (also registered by :func:`register`):
|
| - Runs a final cache sweep, then VACUUM and closes the connection cleanly
|
| (triggers a WAL checkpoint). Handles the case where the process exits
|
| before the next background interval fires.
|
|
|
| Call :func:`register` once at startup (see ``cli.py``).
|
| """
|
|
|
| from __future__ import annotations
|
|
|
| import atexit
|
| import logging
|
| import os
|
| import threading
|
| import time
|
| from pathlib import Path
|
|
|
| log = logging.getLogger("ai_agent.shutdown")
|
|
|
| LOG_RETENTION_DAYS: int = int(os.getenv("LOG_RETENTION_DAYS", "7"))
|
| CLEANUP_INTERVAL_SECONDS: int = int(os.getenv("CLEANUP_INTERVAL_SECONDS", "7200"))
|
|
|
|
|
|
|
|
|
|
|
|
|
| def _sweep_cache_db() -> None:
|
| """Delete expired rows from the cache DB (lightweight, runs periodically)."""
|
| from ai_agent.utils.cache_db import get_cache_db_or_none
|
|
|
| db = get_cache_db_or_none()
|
| if db is None:
|
| return
|
| try:
|
| deleted = db.sweep_expired()
|
| if deleted:
|
| log.debug("Cache sweep: removed %d expired row(s).", deleted)
|
| except Exception:
|
| log.exception("Periodic cache sweep failed.")
|
|
|
|
|
| def _vacuum_and_close_cache_db() -> None:
|
| """Final shutdown: VACUUM and close the cache DB (runs via atexit)."""
|
| from ai_agent.utils.cache_db import get_cache_db_or_none
|
|
|
| db = get_cache_db_or_none()
|
| if db is None:
|
| return
|
| try:
|
| db.vacuum_and_close()
|
| log.info("Cache DB shutdown: VACUUM complete.")
|
| except Exception:
|
| log.exception("Cache DB shutdown cleanup failed.")
|
|
|
|
|
|
|
|
|
|
|
|
|
| def _purge_old_logs() -> None:
|
| """Delete log files older than LOG_RETENTION_DAYS inside LOG_DIR.
|
|
|
| Age is determined by each file's modification time (``st_mtime``), which
|
| reflects when data was last written. Parsing the date from the filename
|
| prefix is deliberately avoided: ``TimedRotatingFileHandler`` keeps the
|
| startup date in the base name when it rotates, so a filename-based age
|
| would be wrong for rotated files.
|
| """
|
| log_dir = Path(os.getenv("LOG_DIR", "logs"))
|
| if not log_dir.is_dir():
|
| return
|
|
|
| cutoff = time.time() - LOG_RETENTION_DAYS * 86_400
|
| removed = 0
|
| errors = 0
|
| for entry in log_dir.iterdir():
|
| if not entry.is_file():
|
| continue
|
| if not (entry.name.startswith("app_") and ".log" in entry.name):
|
| continue
|
| try:
|
| if entry.stat().st_mtime < cutoff:
|
| entry.unlink()
|
| removed += 1
|
| except Exception:
|
| log.exception("Failed to delete old log file: %s", entry)
|
| errors += 1
|
|
|
| if removed or errors:
|
| log.info(
|
| "Log cleanup: removed %d file(s) older than %d day(s)%s.",
|
| removed,
|
| LOG_RETENTION_DAYS,
|
| f", {errors} error(s)" if errors else "",
|
| )
|
|
|
|
|
|
|
|
|
|
|
|
|
| def _cleanup_loop(interval: int, stop_event: threading.Event) -> None:
|
| """Run periodic sweeps until *stop_event* is set or the process exits.
|
|
|
| The first sweep runs immediately on startup so stale data is removed
|
| without waiting for the first interval to elapse.
|
| """
|
| while True:
|
| _sweep_cache_db()
|
| _purge_old_logs()
|
| if stop_event.wait(timeout=interval):
|
| break
|
|
|
|
|
|
|
|
|
|
|
|
|
| _stop_event: threading.Event | None = None
|
| _cleanup_thread: threading.Thread | None = None
|
|
|
|
|
| def _stop_background_cleanup_and_close_cache_db() -> None:
|
| """Stop the background cleanup thread before closing the shared cache DB."""
|
| global _stop_event, _cleanup_thread
|
|
|
| if _stop_event is not None:
|
| _stop_event.set()
|
|
|
| thread = _cleanup_thread
|
| if (
|
| thread is not None
|
| and thread.is_alive()
|
| and thread is not threading.current_thread()
|
| ):
|
| thread.join(timeout=5.0)
|
|
|
| _vacuum_and_close_cache_db()
|
|
|
|
|
| def register() -> None:
|
| """Start the background cleanup thread and register the atexit hook.
|
|
|
| Safe to call multiple times (idempotent).
|
| """
|
| global _stop_event, _cleanup_thread
|
|
|
|
|
| if _stop_event is not None:
|
| _stop_event.set()
|
|
|
| previous_thread = _cleanup_thread
|
| if (
|
| previous_thread is not None
|
| and previous_thread.is_alive()
|
| and previous_thread is not threading.current_thread()
|
| ):
|
| previous_thread.join(timeout=5.0)
|
|
|
| _stop_event = threading.Event()
|
| _cleanup_thread = threading.Thread(
|
| target=_cleanup_loop,
|
| args=(CLEANUP_INTERVAL_SECONDS, _stop_event),
|
| name="cache-log-cleanup",
|
| daemon=True,
|
| )
|
| _cleanup_thread.start()
|
| log.info(
|
| "Background cleanup started (interval: %ds, log retention: %dd).",
|
| CLEANUP_INTERVAL_SECONDS,
|
| LOG_RETENTION_DAYS,
|
| )
|
|
|
| atexit.unregister(_stop_background_cleanup_and_close_cache_db)
|
| atexit.register(_stop_background_cleanup_and_close_cache_db)
|
|
|