# ============================================ # Projet : API-KA # Fichier : src/api/reqstats.py # Node : m3u96b # Author : Simon-Pierre Boucher # Contact : contact@spboucher.ai # Date : 2026-08-17 # ============================================ """Journalisation légère des requêtes API vers la table ``api_requests``. Conçu pour ne JAMAIS ralentir le chemin de requête : le middleware appelle ``record()`` (un simple ``deque.append``, O(1), sans I/O) et une tâche de fond (démarrée dans le lifespan de l'app) vide le tampon en lot toutes les ``FLUSH_INTERVAL`` secondes via un thread, hors boucle d'événements. Rétention : ``RETENTION_DAYS`` jours — une purge s'exécute au plus une fois par jour lors d'un flush. Alimente exclusivement la page /stats. """ from __future__ import annotations import asyncio import datetime import threading from collections import deque from src.config import SERVICES from src.utils.logger import get_logger FLUSH_INTERVAL = 5.0 # secondes entre deux vidages du tampon MAX_BUFFER = 10_000 # garde-fou mémoire : au-delà, les plus anciens sont perdus RETENTION_DAYS = 90 _buffer: deque[dict] = deque(maxlen=MAX_BUFFER) _lock = threading.Lock() _last_purge: datetime.date | None = None # Chemins statiques connus, conservés tels quels dans la colonne ``endpoint``. _KNOWN_PATHS = frozenset( { "/", "/health", "/contact", "/stats", "/docs", "/redoc", "/openapi.json", "/favicon.svg", "/apple-touch-icon.png", "/og.png", "/api/v1/runs", "/api/stats", "/api/stats/dashboard", "/api/stats/report", "/api/stats/ecosystem-report", } ) _SERVICE_SUBROUTES = frozenset({"latest", "stats"}) # Familles de routes API repliées sous « /* » (auth, agent, iOS, # monitoring) : cardinalité bornée sans les noyer dans « (autre) ». _API_FAMILIES = ("/api/auth", "/api/agent", "/api/ios", "/api/monitoring") def normalize_endpoint(path: str) -> str: """Replie un chemin de requête vers un endpoint à cardinalité bornée. Les paramètres de route (dates) sont remplacés par des gabarits, les assets statiques /ka/* sont agrégés et les chemins inconnus (scans de bots, 404) sont regroupés sous « (autre) ». """ path = path.rstrip("/") or "/" if path in _KNOWN_PATHS: return path if path.startswith("/ka/"): return "/ka/*" if path.startswith("/api/v1/"): parts = path.split("/") # ['', 'api', 'v1', service, ...] service = parts[3] if len(parts) > 3 else "" if service in SERVICES: if len(parts) == 4: return f"/api/v1/{service}" sub = parts[4] if sub == "date": return f"/api/v1/{service}/date/{{date}}" if sub in _SERVICE_SUBROUTES and len(parts) == 5: return f"/api/v1/{service}/{sub}" return "/api/v1/(autre)" for family in _API_FAMILIES: if path == family or path.startswith(family + "/"): return f"{family}/*" return "(autre)" def record( ts: datetime.datetime, method: str, path: str, status: int, duration_ms: float ) -> None: """Empile une requête dans le tampon mémoire (non bloquant, jamais d'I/O).""" entry = { "ts": ts, "method": method, "endpoint": normalize_endpoint(path), "status": int(status), "duration_ms": round(float(duration_ms), 2), } with _lock: _buffer.append(entry) def _drain() -> list[dict]: """Vide le tampon et retourne son contenu.""" with _lock: entries = list(_buffer) _buffer.clear() return entries def flush() -> int: """Insère en lot le contenu du tampon dans ``api_requests`` (synchrone). Appelée depuis un thread par la tâche de fond ; purge les lignes plus vieilles que ``RETENTION_DAYS`` jours au plus une fois par jour. """ global _last_purge entries = _drain() if not entries: _maybe_purge() return 0 try: from src.database.db import session_scope from src.database.models import ApiRequest with session_scope() as session: session.bulk_insert_mappings(ApiRequest.__mapper__, entries) except Exception: # pragma: no cover — la stat ne doit jamais casser l'API get_logger("apika.reqstats").exception("Échec du flush api_requests") return 0 _maybe_purge() return len(entries) def _maybe_purge() -> None: """Purge les requêtes plus vieilles que RETENTION_DAYS (1 fois/jour max).""" global _last_purge today = datetime.datetime.now(tz=datetime.UTC).date() if _last_purge == today: return _last_purge = today try: from sqlalchemy import delete from src.database.db import session_scope from src.database.models import ApiRequest cutoff = datetime.datetime.now(tz=datetime.UTC) - datetime.timedelta( days=RETENTION_DAYS ) with session_scope() as session: session.execute(delete(ApiRequest).where(ApiRequest.ts < cutoff)) except Exception: # pragma: no cover get_logger("apika.reqstats").exception("Échec de la purge api_requests") async def flusher_task() -> None: """Tâche de fond : flush périodique du tampon, hors event loop (thread).""" try: while True: await asyncio.sleep(FLUSH_INTERVAL) await asyncio.to_thread(flush) except asyncio.CancelledError: # Dernier flush au shutdown pour ne rien perdre. await asyncio.to_thread(flush) raise