"""Periodic system metrics sampling (DB + event bus) and request counters.""" from __future__ import annotations import asyncio import logging import time from .config import Settings from .db import Database from .events import bus from .hardware import GB, sample_telemetry_async log = logging.getLogger("llm_api.metrics") class MetricsCollector: def __init__(self, settings: Settings, db: Database, manager): self.settings = settings self.db = db self.manager = manager self.last = None self._task: asyncio.Task | None = None self.started_at = time.time() self.window: list[dict] = [] # recent throughput samples (tps) for dashboard async def start(self) -> None: self._task = asyncio.create_task(self._loop(), name="metrics") async def stop(self) -> None: if self._task: self._task.cancel() async def sample(self) -> dict: tel = await sample_telemetry_async(self.settings.models_dir) worker = sum(lm.handle.memory_bytes() for lm in self.manager.loaded.values()) / GB cur = self.manager.current_model() d = tel.to_dict() d["worker_rss_gb"] = round(worker, 2) d["loaded_model"] = cur.model["id"] if cur else None d["loaded_models"] = [lm.model["id"] for lm in self.manager.loaded.values()] d["manager"] = self.manager.stats d["app_uptime_seconds"] = round(time.time() - self.started_at) self.last = d return d async def _loop(self) -> None: tick = 0 while True: try: d = await self.sample() bus.publish("metrics", d) tick += 1 interval = max(3, self.settings.metrics_interval_seconds) if tick % max(1, interval // 3) == 0: await self.db.execute( "INSERT OR REPLACE INTO system_metrics(ts, mem_used_gb, mem_available_gb, mem_pressure, swap_used_gb, " "cpu_percent, gpu_percent, disk_free_gb, thermal, worker_rss_gb, loaded_model) VALUES(?,?,?,?,?,?,?,?,?,?,?)", (d["ts"], d["mem_used_gb"], d["mem_available_gb"], d["mem_pressure_percent"], d["swap_used_gb"], d["cpu_percent"], d["gpu_percent"], d["disk_free_gb"], d["thermal_state"], d["worker_rss_gb"], d["loaded_model"])) if tick % 200 == 0: cutoff = time.time() - self.settings.metrics_retention_days * 86400 await self.db.execute("DELETE FROM system_metrics WHERE ts < ?", (cutoff,)) await self.db.execute("DELETE FROM inference_requests WHERE created_at < ?", (cutoff,)) if d["swap_used_gb"] > 4 or d["mem_pressure_level"] == "critical": bus.publish("alert", {"level": "warning", "message": f"Memory pressure {d['mem_pressure_level']} — swap {d['swap_used_gb']} GB. " "The loaded model exceeds the recommended operating envelope."}) await asyncio.sleep(3) except asyncio.CancelledError: return except Exception: log.exception("metrics loop error") await asyncio.sleep(5) async def history(self, minutes: int = 60) -> list[dict]: since = time.time() - minutes * 60 return await self.db.fetchall("SELECT * FROM system_metrics WHERE ts >= ? ORDER BY ts", (since,)) async def request_stats(self, hours: int = 24) -> dict: since = time.time() - hours * 3600 totals = await self.db.fetchone( "SELECT COUNT(*) AS requests, COALESCE(SUM(prompt_tokens),0) AS prompt_tokens, " "COALESCE(SUM(completion_tokens),0) AS completion_tokens, AVG(tps) AS avg_tps, AVG(ttft_ms) AS avg_ttft_ms, " "SUM(CASE WHEN status>=400 THEN 1 ELSE 0 END) AS errors FROM inference_requests WHERE created_at >= ?", (since,)) per_model = await self.db.fetchall( "SELECT model_id, COUNT(*) AS requests, COALESCE(SUM(completion_tokens),0) AS completion_tokens, " "AVG(tps) AS avg_tps, AVG(ttft_ms) AS avg_ttft_ms FROM inference_requests WHERE created_at >= ? " "GROUP BY model_id ORDER BY requests DESC", (since,)) buckets = await self.db.fetchall( "SELECT CAST(created_at / 3600 AS INTEGER) * 3600 AS hour, COUNT(*) AS requests, " "COALESCE(SUM(completion_tokens),0) AS tokens FROM inference_requests WHERE created_at >= ? GROUP BY hour ORDER BY hour", (since,)) all_time = await self.db.fetchone( "SELECT COUNT(*) AS requests, COALESCE(SUM(completion_tokens),0) AS completion_tokens, " "COALESCE(SUM(prompt_tokens),0) AS prompt_tokens FROM inference_requests") return {"window_hours": hours, "totals": totals, "per_model": per_model, "hourly": buckets, "all_time": all_time}