bot_host / analytics.py
ItsBounvy's picture
Upload 31 files
66916b2 verified
Raw History Blame Contribute Delete
21.3 kB
"""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/<slug>/<rest>`` -> ``/proxy/<slug>/...``
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/<slug>/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/<slug>"
return f"/proxy/<slug>/{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/<slug>/<rest>...``; 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",
]