"""POST /api/v1/ingest — приём пакета метрик от hard-monitor.
Панель — мозг: здесь считается скорость сети (дельта сырых счётчиков),
обновляется карточка сервера и пишется точка метрик. Ивенты (пороги,
diff дисков/docker, offline) появятся на этапе 4 — тоже здесь.
"""
import json
from datetime import datetime, timezone
from fastapi import APIRouter, Header, HTTPException, Request
from app.db import get_db, j
from app.models import IngestPayload
from app.security import hash_key
router = APIRouter(prefix="/api/v1")
MIB = 1024 * 1024
# «Онлайн», если последний пакет был не позже чем interval * множитель + запас
OFFLINE_GRACE_SEC = 60
def _now() -> datetime:
return datetime.now(timezone.utc)
def _iso(dt: datetime) -> str:
return dt.astimezone(timezone.utc).isoformat()
def _client_ip(request: Request, x_forwarded_for: str) -> str:
if x_forwarded_for:
return x_forwarded_for.split(",")[0].strip()
return request.client.host if request.client else ""
def _calc_net_rates(
payload: IngestPayload, prev_ts: datetime | None, prev_counters: dict[str, tuple[int, int]], now: datetime
) -> tuple[float, float]:
"""МБ/с по дельте счётчиков между пакетами.
Счётчики psutil живут с загрузки ОС и не сбрасываются рестартом агента,
поэтому дельта на стороне panel всегда валидна. Отрицательная дельта
(обнуление после ребута) → 0.
"""
dt = (now - prev_ts).total_seconds() if prev_ts else 0
if dt <= 0 or not prev_counters:
return 0.0, 0.0
in_total = 0.0
out_total = 0.0
for iface in payload.net:
if iface.iface == "lo" or iface.iface not in prev_counters:
continue
prev_sent, prev_recv = prev_counters[iface.iface]
if iface.bytes_recv >= prev_recv:
in_total += (iface.bytes_recv - prev_recv) / dt / MIB
if iface.bytes_sent >= prev_sent:
out_total += (iface.bytes_sent - prev_sent) / dt / MIB
return round(in_total, 3), round(out_total, 3)
@router.post("/ingest")
async def ingest(
payload: IngestPayload,
request: Request,
x_server_key: str = Header(alias="X-Server-Key"),
x_forwarded_for: str = Header(default="", alias="X-Forwarded-For"),
) -> dict:
db = get_db()
cursor = await db.execute(
"SELECT id FROM servers WHERE key_hash = ?", (hash_key(x_server_key),)
)
server = await cursor.fetchone()
if server is None:
raise HTTPException(status_code=401, detail="unknown server key")
server_id = server["id"]
now = _now()
packet_ts = _iso(payload.ts) if payload.ts else _iso(now)
# предыдущая точка — нужны сырые счётчики для дельты сети
cursor = await db.execute(
"SELECT ts, net_json FROM metrics WHERE server_id = ? ORDER BY id DESC LIMIT 1",
(server_id,),
)
prev = await cursor.fetchone()
prev_ts = None
prev_counters: dict[str, tuple[int, int]] = {}
if prev is not None:
try:
prev_ts = datetime.fromisoformat(prev["ts"])
for iface in json.loads(prev["net_json"]):
prev_counters[iface["iface"]] = (iface["bytes_sent"], iface["bytes_recv"])
except (ValueError, KeyError, TypeError, json.JSONDecodeError):
prev_ts = None
prev_counters = {}
net_in_mbs, net_out_mbs = _calc_net_rates(payload, prev_ts, prev_counters, now)
# обновить карточку сервера (факты из пакета)
await db.execute(
"""UPDATE servers SET hostname = ?, os = ?, kernel = ?, ips_json = ?,
source_ip = ?, interval = ?, last_seen = ? WHERE id = ?""",
(
payload.hostname,
payload.os,
payload.kernel,
j(payload.ips),
_client_ip(request, x_forwarded_for),
payload.interval,
_iso(now),
server_id,
),
)
# точка метрик; net_json — сырые счётчики, из них будет считаться дельта
await db.execute(
"""INSERT INTO metrics (server_id, ts, cpu, load1, load5, load15,
ram_used, ram_total, swap_used, swap_total, uptime,
net_in_mbs, net_out_mbs, disks_json, net_json, processes_json,
docker_json, extra_json)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(
server_id,
packet_ts,
payload.cpu.percent,
*(payload.cpu.load + [0.0] * (3 - len(payload.cpu.load))),
payload.memory.ram.used,
payload.memory.ram.total,
payload.memory.swap.used,
payload.memory.swap.total,
payload.uptime,
net_in_mbs,
net_out_mbs,
j([d.model_dump() for d in payload.disks]),
j([n.model_dump() for n in payload.net]),
j([p.model_dump() for p in payload.processes]),
j([c.model_dump() for c in payload.docker]),
j(payload.extra),
),
)
await db.commit()
return {"status": "ok"}