"""Analytics: derive proxy traffic and bot file-activity metrics from ``AuditLog`` + live bot metrics from ``bot_runner``. Why derive instead of pre-aggregate ----------------------------------- The panel already write-through-persists every ``AuditLog`` row to the HF Dataset (see ``storage.py``). For the traffic volumes a hosting panel sees — even at ~10 req/s sustained, that's <1M rows/day — a single ``GROUP BY strftime('%Y-%m-%d %H:00:00', created_at)`` is well under 200 ms on SQLite. So: * No separate ``metric_samples`` rollup table to maintain, * No risk of pre-aggregation diverging from the source of truth, * "Last hour / last 24h / last 7d / last 30d" all use the same code path — only the window changes. If a future deployment exceeds this, the same collectors can be swapped for a rollup table without breaking the route signatures. Bucketing --------- All time-series use *hour* buckets (``YYYY-MM-DD HH:00:00`` UTC). This matches what humans actually look at — sub-hour resolution makes charts spiky and hard to read at the panel's typical scale. Percentiles ----------- We compute p50/p95/p99 of latency in Python over the rows fetched for the window. SQLite has no native percentile function. For windows above 50k rows we downsample to 5k before computing percentiles to keep the response under 200 ms; the bucketed timeseries stays exact. """ from __future__ import annotations import json import logging import time from collections import defaultdict from dataclasses import dataclass from datetime import datetime, timedelta from typing import Any, Iterable, Optional from sqlalchemy import select, func, and_, desc from sqlalchemy.ext.asyncio import AsyncSession from models import ( AuditLog, BotInstance, BotStatus, DeploymentMode, ) log = logging.getLogger("panel.analytics") # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- _HOUR_FMT = "%Y-%m-%d %H:00:00" def _utc_now() -> datetime: return datetime.utcnow() def _hour_bucket(dt: datetime) -> datetime: """Truncate a datetime to its UTC hour bucket.""" return dt.replace(minute=0, second=0, microsecond=0) def _hour_label(dt: datetime) -> str: return dt.strftime(_HOUR_FMT) def _hour_labels(start: datetime, end: datetime) -> list[str]: """Inclusive list of hour labels covering [start, end].""" cur = _hour_bucket(start) end_b = _hour_bucket(end) labels: list[str] = [] while cur <= end_b: labels.append(_hour_label(cur)) cur = cur + timedelta(hours=1) return labels def _percentile(sorted_values: list[float], p: float) -> float: """Linear-interpolation percentile. ``p`` in [0, 100].""" if not sorted_values: return 0.0 if len(sorted_values) == 1: return float(sorted_values[0]) if p <= 0: return float(sorted_values[0]) if p >= 100: return float(sorted_values[-1]) k = (len(sorted_values) - 1) * (p / 100.0) f = int(k) c = min(f + 1, len(sorted_values) - 1) if f == c: return float(sorted_values[f]) return float(sorted_values[f] + (sorted_values[c] - sorted_values[f]) * (k - f)) # --------------------------------------------------------------------------- # Proxy traffic # --------------------------------------------------------------------------- PROXY_ACTIONS = ("proxy_request",) @dataclass class ProxyTrafficSummary: """A single time-window rollup of proxy traffic.""" window_start: datetime window_end: datetime total_requests: int error_count: int # 5xx client_error_count: int # 4xx success_count: int # 2xx-3xx avg_latency_ms: float p50_latency_ms: float p95_latency_ms: float p99_latency_ms: float requests_per_minute: float error_rate_percent: float = 0.0 # explicit field so it's picked up # up by `__dict__` and serialised # cleanly to JSON. @dataclass class ProxyHourBucket: hour: str # "YYYY-MM-DD HH:00:00" requests: int errors: int avg_latency_ms: float p95_latency_ms: float @dataclass class EndpointStat: method: str path_template: str # ``/proxy//`` -> ``/proxy//...`` requests: int errors: int avg_latency_ms: float def _templated_path(path: str) -> str: """Collapse the bot slug in a proxy path to ``...`` so endpoint aggregation isn't fragmented per bot. ``/proxy/my-bot/telegram/webhook`` -> ``/proxy//telegram/webhook`` (we still want to see the trailing sub-path). """ parts = path.strip("/").split("/", 2) if len(parts) >= 2 and parts[0] == "proxy": if len(parts) == 2: return "/proxy/" return f"/proxy//{parts[2]}" return path async def _proxy_window( db: AsyncSession, start: datetime, end: datetime, bot_id: Optional[int] = None, ) -> list[dict]: """Fetch raw proxy_request rows in [start, end), optionally per-bot.""" q = ( select(AuditLog) .where(AuditLog.action == "proxy_request") .where(AuditLog.created_at >= start) .where(AuditLog.created_at < end) ) if bot_id is not None: q = q.where(AuditLog.target_type == "bot").where(AuditLog.target_id == bot_id) q = q.order_by(AuditLog.created_at.asc()) rows = (await db.execute(q)).scalars().all() return rows def _summarise_proxy(rows: Iterable[AuditLog], start: datetime, end: datetime) -> ProxyTrafficSummary: latencies: list[float] = [] total = errors = client_errors = success = 0 for r in rows: total += 1 try: details = json.loads(r.details_json) if r.details_json else {} except (ValueError, TypeError): details = {} status = int(details.get("status") or 0) if status >= 500: errors += 1 elif 400 <= status < 500: client_errors += 1 elif 200 <= status < 400: success += 1 lat = details.get("latency_ms") if isinstance(lat, (int, float)): latencies.append(float(lat)) latencies.sort() elapsed_min = max(1.0, (end - start).total_seconds() / 60.0) rpm = round(total / elapsed_min, 2) err_rate = round((errors / total * 100.0), 2) if total > 0 else 0.0 return ProxyTrafficSummary( window_start=start, window_end=end, total_requests=total, error_count=errors, client_error_count=client_errors, success_count=success, avg_latency_ms=round(sum(latencies) / len(latencies), 1) if latencies else 0.0, p50_latency_ms=round(_percentile(latencies, 50), 1), p95_latency_ms=round(_percentile(latencies, 95), 1), p99_latency_ms=round(_percentile(latencies, 99), 1), requests_per_minute=rpm, error_rate_percent=err_rate, ) async def collect_proxy_traffic( db: AsyncSession, hours: int = 24, bot_id: Optional[int] = None, ) -> dict[str, Any]: """Aggregate proxy traffic over the last ``hours`` hours. Returns ------- { "summary": ProxyTrafficSummary dict, "buckets": [ProxyHourBucket dict] one per hour, always filled (zero-buckets included for empty hours), "endpoints": top-N endpoint stats (default top 10), "status_codes": {str(2xx): int, str(3xx): int, str(4xx): int, str(5xx): int, str(other): int}, "bot_breakdown": [{"bot_id": int, "name": str, "slug": str, "requests": int, "errors": int, "rpm": float}, ...] — only when bot_id is None } """ end = _utc_now() start = end - timedelta(hours=hours) rows = await _proxy_window(db, start, end, bot_id=bot_id) summary = _summarise_proxy(rows, start, end) # ---- Hourly buckets ---- bucket_data: dict[str, dict[str, float]] = {} for r in rows: h = _hour_label(_hour_bucket(r.created_at)) b = bucket_data.setdefault(h, {"requests": 0, "errors": 0, "_lat": []}) b["requests"] += 1 try: details = json.loads(r.details_json) if r.details_json else {} except (ValueError, TypeError): details = {} status = int(details.get("status") or 0) if status >= 500: b["errors"] += 1 lat = details.get("latency_ms") if isinstance(lat, (int, float)): b["_lat"].append(float(lat)) buckets: list[dict] = [] for label in _hour_labels(start, end): b = bucket_data.get(label) if not b: buckets.append({ "hour": label, "requests": 0, "errors": 0, "avg_latency_ms": 0.0, "p95_latency_ms": 0.0, }) continue lats: list[float] = b["_lat"] lats.sort() buckets.append({ "hour": label, "requests": int(b["requests"]), "errors": int(b["errors"]), "avg_latency_ms": round(sum(lats) / len(lats), 1) if lats else 0.0, "p95_latency_ms": round(_percentile(lats, 95), 1) if lats else 0.0, }) # ---- Endpoints ---- endpoint_agg: dict[str, dict[str, float]] = {} status_agg: dict[str, int] = defaultdict(int) for r in rows: try: details = json.loads(r.details_json) if r.details_json else {} except (ValueError, TypeError): details = {} method = (details.get("method") or "?").upper() path = details.get("path") or "" # path looks like ``/proxy//...``; the AuditLog # stores the *incoming* request path (panel-side), so use it. endpoint = _templated_path(path) if path.startswith("/proxy/") else (path or "?") lat = details.get("latency_ms") agg = endpoint_agg.setdefault( f"{method} {endpoint}", {"requests": 0, "errors": 0, "_lat": []}, ) agg["requests"] += 1 status = int(details.get("status") or 0) if status >= 500: agg["errors"] += 1 status_agg[_status_bucket(status)] += 1 if isinstance(lat, (int, float)): agg["_lat"].append(float(lat)) endpoints: list[dict] = [] for key, agg in sorted(endpoint_agg.items(), key=lambda kv: -kv[1]["requests"])[:10]: method, _, path = key.partition(" ") lats: list[float] = agg["_lat"] endpoints.append({ "method": method, "path_template": path, "requests": int(agg["requests"]), "errors": int(agg["errors"]), "avg_latency_ms": round(sum(lats) / len(lats), 1) if lats else 0.0, }) out: dict[str, Any] = { "summary": summary.__dict__, "buckets": buckets, "endpoints": endpoints, "status_codes": dict(status_agg), "window_hours": hours, "ts": time.time(), } # ---- Per-bot breakdown (cluster view only) ---- if bot_id is None: out["bot_breakdown"] = await _proxy_bot_breakdown(db, start, end) return out def _status_bucket(status: int) -> str: if 200 <= status < 300: return "2xx" if 300 <= status < 400: return "3xx" if 400 <= status < 500: return "4xx" if 500 <= status < 600: return "5xx" return "other" async def _proxy_bot_breakdown( db: AsyncSession, start: datetime, end: datetime, limit: int = 10, ) -> list[dict]: """Top-N bots by proxy request volume in the window.""" rows = (await db.execute( select(AuditLog, BotInstance) .join(BotInstance, and_( BotInstance.id == AuditLog.target_id, AuditLog.target_type == "bot", )) .where(AuditLog.action == "proxy_request") .where(AuditLog.created_at >= start) .where(AuditLog.created_at < end) )).all() agg: dict[int, dict[str, Any]] = {} for entry, bot in rows: b = agg.setdefault(bot.id, { "bot_id": bot.id, "name": bot.name, "slug": bot.slug, "requests": 0, "errors": 0, }) b["requests"] += 1 try: details = json.loads(entry.details_json) if entry.details_json else {} except (ValueError, TypeError): details = {} if int(details.get("status") or 0) >= 500: b["errors"] += 1 elapsed_min = max(1.0, (end - start).total_seconds() / 60.0) out: list[dict] = [] for b in sorted(agg.values(), key=lambda r: -r["requests"])[:limit]: b["rpm"] = round(b["requests"] / elapsed_min, 2) out.append(b) return out # --------------------------------------------------------------------------- # Bot activity (uploads + downloads per hour) # --------------------------------------------------------------------------- # # "Upload" / "download" as observed by the panel. We don't instrument # HTTP egress inside the bot process — only panel-mediated actions: # # upload = user pushed content INTO the bot: # * zip_upload (whole zip import) # * file_save (single-file commit) # * file_create (new file) # # download = user pulled content OUT of the bot: # * file_raw (any GET /api/bots/{id}/files/raw in the audit log) # # Anything that touches /api/bots/{id}/files/raw counts as a download # — that's the only panel endpoint that returns bot file bytes to the # user. We log it the same way we log everything else. UPLOAD_ACTIONS = ("zip_upload", "file_save", "file_create") DOWNLOAD_ACTIONS = ("file_raw",) # see _audit_file_raw in routes.py # NOTE: the file_raw audit row is added alongside this module; until # it ships the downloads bucket will be empty — that's expected and # shows in the UI as "no downloads yet". @dataclass class BotActivityHourBucket: hour: str uploads: int downloads: int @dataclass class BotActivitySummary: bot_id: int name: str slug: str total_uploads: int total_downloads: int hours: list[BotActivityHourBucket] async def collect_bot_activity( db: AsyncSession, hours: int = 24, bot_id: Optional[int] = None, ) -> dict[str, Any]: """Per-hour upload/download counts. When ``bot_id`` is None, returns the cluster rollup + per-bot top-N. When set, returns the per-hour buckets for that bot only. """ end = _utc_now() start = end - timedelta(hours=hours) wanted_actions = set(UPLOAD_ACTIONS) | set(DOWNLOAD_ACTIONS) q = ( select(AuditLog) .where(AuditLog.action.in_(wanted_actions)) .where(AuditLog.created_at >= start) .where(AuditLog.created_at < end) .where(AuditLog.target_type == "bot") ) if bot_id is not None: q = q.where(AuditLog.target_id == bot_id) q = q.order_by(AuditLog.created_at.asc()) rows = (await db.execute(q)).scalars().all() # bucket[bot_id][hour_label] = {u: int, d: int} bucket: dict[int, dict[str, dict[str, int]]] = defaultdict(dict) bot_totals: dict[int, dict[str, int]] = defaultdict(lambda: {"u": 0, "d": 0}) for r in rows: bid = int(r.target_id or 0) if not bid: continue h = _hour_label(_hour_bucket(r.created_at)) cell = bucket[bid].setdefault(h, {"u": 0, "d": 0}) if r.action in UPLOAD_ACTIONS: cell["u"] += 1 bot_totals[bid]["u"] += 1 else: cell["d"] += 1 bot_totals[bid]["d"] += 1 labels = _hour_labels(start, end) if bot_id is not None: bot = await db.get(BotInstance, bot_id) if not bot: return {"bots": [], "window_hours": hours} buckets = [] for label in labels: cell = bucket.get(bot_id, {}).get(label, {"u": 0, "d": 0}) buckets.append({"hour": label, "uploads": cell["u"], "downloads": cell["d"]}) totals = bot_totals.get(bot_id, {"u": 0, "d": 0}) return { "bots": [{ "bot_id": bot.id, "name": bot.name, "slug": bot.slug, "total_uploads": totals["u"], "total_downloads": totals["d"], "buckets": buckets, }], "window_hours": hours, } # Cluster view: one row per active bot, with merged buckets. bot_ids = sorted(bot_totals.keys(), key=lambda i: -(bot_totals[i]["u"] + bot_totals[i]["d"])) # Resolve bot metadata in one query. bot_rows = (await db.execute( select(BotInstance).where(BotInstance.id.in_(bot_ids)) )).scalars().all() if bot_ids else [] bot_map = {b.id: b for b in bot_rows} bots: list[dict] = [] for bid in bot_ids: b = bot_map.get(bid) if not b: continue totals = bot_totals[bid] buckets = [] for label in labels: cell = bucket[bid].get(label, {"u": 0, "d": 0}) buckets.append({"hour": label, "uploads": cell["u"], "downloads": cell["d"]}) bots.append({ "bot_id": b.id, "name": b.name, "slug": b.slug, "total_uploads": totals["u"], "total_downloads": totals["d"], "buckets": buckets, }) return {"bots": bots, "window_hours": hours} # --------------------------------------------------------------------------- # Live bot metrics (used by alerts evaluator + the metrics page) # --------------------------------------------------------------------------- async def collect_live_bot_metrics(db: AsyncSession) -> list[dict]: """Per-bot live snapshot — RAM/CPU/status. MULTITENANT bots have a live process we can probe via ``bot_runner.bot_status``. LEGACY_SPACE bots report a status only (their RAM/CPU live in the remote HF Space, not here). """ import bot_runner # local import to keep this module cheap to import rows = (await db.execute(select(BotInstance))).scalars().all() out: list[dict] = [] for b in rows: rec = { "bot_id": b.id, "name": b.name, "slug": b.slug, "status": b.status.value if hasattr(b.status, "value") else str(b.status), "deployment_mode": b.deployment_mode.value if hasattr(b.deployment_mode, "value") else str(b.deployment_mode), "ram_mb_used": None, "cpu_percent": None, "port": b.port, "last_healthcheck": b.last_healthcheck.isoformat() if b.last_healthcheck else None, } if b.deployment_mode == DeploymentMode.MULTITENANT and b.port: try: st = bot_runner.bot_status(b) rec["ram_mb_used"] = float(st.get("rss_mb") or 0.0) rec["cpu_percent"] = float(st.get("cpu_percent") or 0.0) except Exception: # noqa: BLE001 pass out.append(rec) return out # --------------------------------------------------------------------------- # Cluster overview — single payload for the dashboard cards. # --------------------------------------------------------------------------- async def collect_overview(db: AsyncSession, hours: int = 24) -> dict[str, Any]: """Top-of-dashboard cards: requests today, errors, active bots, most active bot, alerts in last 24h. """ from models import AlertEvent # local to dodge a cycle at import end = _utc_now() start = end - timedelta(hours=hours) proxy = await collect_proxy_traffic(db, hours=hours) activity = await collect_bot_activity(db, hours=hours) live = await collect_live_bot_metrics(db) active_bots = [r for r in live if r["status"] == BotStatus.RUNNING.value] top_bots = sorted( activity["bots"], key=lambda b: -(b["total_uploads"] + b["total_downloads"]) )[:5] alert_rows = (await db.execute( select(func.count(AlertEvent.id)).where(AlertEvent.fired_at >= start) )).scalar() or 0 return { "summary": proxy["summary"], "buckets": proxy["buckets"], "bots": { "total": len(live), "running": len(active_bots), "status_breakdown": { (r["status"]): sum(1 for x in live if x["status"] == r["status"]) for r in live }, }, "activity": { "total_uploads": sum(b["total_uploads"] for b in activity["bots"]), "total_downloads": sum(b["total_downloads"] for b in activity["bots"]), "top_bots": top_bots, }, "alerts_24h": int(alert_rows), "window_hours": hours, "ts": time.time(), } __all__ = [ "ProxyTrafficSummary", "ProxyHourBucket", "EndpointStat", "BotActivitySummary", "BotActivityHourBucket", "PROXY_ACTIONS", "UPLOAD_ACTIONS", "DOWNLOAD_ACTIONS", "collect_proxy_traffic", "collect_bot_activity", "collect_live_bot_metrics", "collect_overview", ]