ai-agent / src /ai_agent /utils /shutdown.py
katospiegel's picture
Deploy develop: FastAPI+React frontend, multi-stage Docker (ai_agent serve)
07c2476 verified
Raw
History Blame Contribute Delete
6.21 kB
# utils/shutdown.py
"""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"))
# ---------------------------------------------------------------------------
# Cache DB helpers
# ---------------------------------------------------------------------------
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 # noqa: PLC0415
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 # noqa: PLC0415
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.")
# ---------------------------------------------------------------------------
# Log file rotation helper
# ---------------------------------------------------------------------------
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 "",
)
# ---------------------------------------------------------------------------
# Background cleanup thread
# ---------------------------------------------------------------------------
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
# ---------------------------------------------------------------------------
# Public API
# ---------------------------------------------------------------------------
_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
# Stop any previously running background thread before restarting.
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, # won't block process exit
)
_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)