SPB Git forge

spb/llm-api

Public
0commits 0branches 0releases
0 Bsize
maindefault branch
—last push
4.9 KB · 99 lines python
Raw Blame History
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