1"""Periodic system metrics sampling (DB + event bus) and request counters."""23from __future__ import annotations45import asyncio6import logging7import time89from .config import Settings10from .db import Database11from .events import bus12from .hardware import GB, sample_telemetry_async1314log = logging.getLogger("llm_api.metrics")151617class MetricsCollector:18 def __init__(self, settings: Settings, db: Database, manager):19 self.settings = settings20 self.db = db21 self.manager = manager22 self.last = None23 self._task: asyncio.Task | None = None24 self.started_at = time.time()25 self.window: list[dict] = [] # recent throughput samples (tps) for dashboard2627 async def start(self) -> None:28 self._task = asyncio.create_task(self._loop(), name="metrics")2930 async def stop(self) -> None:31 if self._task:32 self._task.cancel()3334 async def sample(self) -> dict:35 tel = await sample_telemetry_async(self.settings.models_dir)36 worker = sum(lm.handle.memory_bytes() for lm in self.manager.loaded.values()) / GB37 cur = self.manager.current_model()38 d = tel.to_dict()39 d["worker_rss_gb"] = round(worker, 2)40 d["loaded_model"] = cur.model["id"] if cur else None41 d["loaded_models"] = [lm.model["id"] for lm in self.manager.loaded.values()]42 d["manager"] = self.manager.stats43 d["app_uptime_seconds"] = round(time.time() - self.started_at)44 self.last = d45 return d4647 async def _loop(self) -> None:48 tick = 049 while True:50 try:51 d = await self.sample()52 bus.publish("metrics", d)53 tick += 154 interval = max(3, self.settings.metrics_interval_seconds)55 if tick % max(1, interval // 3) == 0:56 await self.db.execute(57 "INSERT OR REPLACE INTO system_metrics(ts, mem_used_gb, mem_available_gb, mem_pressure, swap_used_gb, "58 "cpu_percent, gpu_percent, disk_free_gb, thermal, worker_rss_gb, loaded_model) VALUES(?,?,?,?,?,?,?,?,?,?,?)",59 (d["ts"], d["mem_used_gb"], d["mem_available_gb"], d["mem_pressure_percent"], d["swap_used_gb"],60 d["cpu_percent"], d["gpu_percent"], d["disk_free_gb"], d["thermal_state"], d["worker_rss_gb"],61 d["loaded_model"]))62 if tick % 200 == 0:63 cutoff = time.time() - self.settings.metrics_retention_days * 8640064 await self.db.execute("DELETE FROM system_metrics WHERE ts < ?", (cutoff,))65 await self.db.execute("DELETE FROM inference_requests WHERE created_at < ?", (cutoff,))66 if d["swap_used_gb"] > 4 or d["mem_pressure_level"] == "critical":67 bus.publish("alert", {"level": "warning", "message":68 f"Memory pressure {d['mem_pressure_level']} — swap {d['swap_used_gb']} GB. "69 "The loaded model exceeds the recommended operating envelope."})70 await asyncio.sleep(3)71 except asyncio.CancelledError:72 return73 except Exception:74 log.exception("metrics loop error")75 await asyncio.sleep(5)7677 async def history(self, minutes: int = 60) -> list[dict]:78 since = time.time() - minutes * 6079 return await self.db.fetchall("SELECT * FROM system_metrics WHERE ts >= ? ORDER BY ts", (since,))8081 async def request_stats(self, hours: int = 24) -> dict:82 since = time.time() - hours * 360083 totals = await self.db.fetchone(84 "SELECT COUNT(*) AS requests, COALESCE(SUM(prompt_tokens),0) AS prompt_tokens, "85 "COALESCE(SUM(completion_tokens),0) AS completion_tokens, AVG(tps) AS avg_tps, AVG(ttft_ms) AS avg_ttft_ms, "86 "SUM(CASE WHEN status>=400 THEN 1 ELSE 0 END) AS errors FROM inference_requests WHERE created_at >= ?", (since,))87 per_model = await self.db.fetchall(88 "SELECT model_id, COUNT(*) AS requests, COALESCE(SUM(completion_tokens),0) AS completion_tokens, "89 "AVG(tps) AS avg_tps, AVG(ttft_ms) AS avg_ttft_ms FROM inference_requests WHERE created_at >= ? "90 "GROUP BY model_id ORDER BY requests DESC", (since,))91 buckets = await self.db.fetchall(92 "SELECT CAST(created_at / 3600 AS INTEGER) * 3600 AS hour, COUNT(*) AS requests, "93 "COALESCE(SUM(completion_tokens),0) AS tokens FROM inference_requests WHERE created_at >= ? GROUP BY hour ORDER BY hour",94 (since,))95 all_time = await self.db.fetchone(96 "SELECT COUNT(*) AS requests, COALESCE(SUM(completion_tokens),0) AS completion_tokens, "97 "COALESCE(SUM(prompt_tokens),0) AS prompt_tokens FROM inference_requests")98 return {"window_hours": hours, "totals": totals, "per_model": per_model, "hourly": buckets, "all_time": all_time}99