"""POST /api/v1/ingest — приём пакета метрик от hard-monitor.
Панель — мозг: здесь считается скорость сети (дельта сырых счётчиков),
обновляется карточка сервера и пишется точка метрик. Из потока пакетов
тут же генерируются ивенты: восстановление после offline, пороги
(cpu/ram/swap/disk/load, гистерезис), diff docker-контейнеров.
"""
import json
import re
from datetime import datetime, timezone
from fastapi import APIRouter, Header, HTTPException, Request
from app import events as ev
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
# Пороги (гистерезис: открытие high → закрытие back). Отдельные правила
# с UI-настройкой — позже; сейчас константы панели.
THRESHOLDS = {
"cpu": {"high": 90.0, "back": 80.0, "severity": "warning", "unit": "%"},
"ram": {"high": 90.0, "back": 80.0, "severity": "warning", "unit": "%"},
"swap": {"high": 60.0, "back": 40.0, "severity": "warning", "unit": "%"},
"disk": {"high": 90.0, "back": 85.0, "severity": "critical", "unit": "%"},
"load": {"high": 2.0, "back": 1.0, "severity": "warning", "unit": ""}, # load1 / cores
}
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)
# предыдущая точка — нужны сырые счётчики для дельты сети и diff docker
cursor = await db.execute(
"SELECT ts, net_json, docker_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]] = {}
prev_docker: list[dict] = []
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 = {}
try:
prev_docker = json.loads(prev["docker_json"])
except (ValueError, TypeError, json.JSONDecodeError):
prev_docker = []
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),
),
)
# --- Ивенты этапа 4: восстановление, пороги, diff docker -------------------
await ev.close_if_open(db, server_id, "server_online", "данные снова идут", {})
await _check_thresholds(db, server_id, payload)
await _diff_docker(db, server_id, prev_docker, payload.docker)
await db.commit()
return {"status": "ok"}
# --- Генерация ивентовиз пакета ------------------------------------------------
def _pct(used: int, total: int) -> float:
return used / total * 100.0 if total else 0.0
async def _check_thresholds(db, server_id: int, payload: IngestPayload) -> None:
"""Пороги cpu/ram/swap/disk/load с гистерезисом.
Состояние пары (…_high / …_recovered) — журнал: переживает рестарт
панели, повторно не спамит (opening только при переходе).
"""
swap = payload.memory.swap
items: list[tuple[str, float, dict]] = [ # (metric, value, data-дополнение)
("cpu", payload.cpu.percent, {}),
("ram", _pct(payload.memory.ram.used, payload.memory.ram.total), {}),
]
if swap.total:
items.append(("swap", _pct(swap.used, swap.total), {}))
if payload.memory.ram.total == 0:
items.pop(1) # агент не прислал память — не «0%»: не трогаем открытый порог
items += [("disk", d.percent, {"mount": d.mount}) for d in payload.disks]
if payload.cpu.count:
items.append(("load", payload.cpu.load[0] / payload.cpu.count, {}))
for metric, value, extra in items:
rule = THRESHOLDS[metric]
open_type = f"{metric}_high"
was_open = await ev.state_open_ext(db, server_id, open_type, mount=extra.get("mount"))
if value < rule["back"]:
if was_open:
await ev.record_event(
db,
type=f"{metric}_recovered",
severity="info",
message=_msg(metric, value, extra, rule["back"]),
server_id=server_id,
data={"value": round(value, 2), "threshold": rule["back"], **extra},
)
elif value >= rule["high"] and not was_open:
await ev.record_event(
db,
type=open_type,
severity=rule["severity"],
message=_msg(metric, value, extra, rule["high"]),
server_id=server_id,
data={"value": round(value, 2), "threshold": rule["high"], **extra},
)
def _msg(metric: str, value: float, extra: dict, at: float) -> str:
"""Компактное техное сообщение журнала: `cpu 95% (≥90)`, `disk / at 98% (≥90)`."""
if metric == "load":
return f"load1 {value:.1f}×cores ({'<' if value < at else '≥'}{at:.1f}×)"
unit = "%"
where = f" {extra['mount']}" if "mount" in extra else ""
cmp = "<" if value < at else "≥"
return f"{metric}{where} {value:.0f}{unit} ({cmp}{at:.0f}{unit})"
_EXIT_RE = re.compile(r"Exited \((\d*)\)")
async def _diff_docker(db, server_id: int, prev: list[dict], docker) -> None:
"""Diff контейнеров с предыдущим пакетом: added/removed/exited/started.
Список контейнеров — полная выборка с хоста, а не поток: сравнение
делаем на стороне панели (агент тупой, принцип «панель — мозг»).
"""
prev_by_name = {c["name"]: c for c in prev}
curr_by_name = {c.name: c.model_dump() for c in docker}
for name, c in curr_by_name.items():
was = prev_by_name.get(name)
if was is None:
await ev.record_event(db, type="container_added", severity="info",
message=f"{name} (added)", server_id=server_id,
data={"name": name, "image": c.get("image", "")})
elif _is_exited(was["status"]) and not _is_exited(c["status"]):
await ev.record_event(db, type="container_started", severity="info",
message=f"{name} (started)", server_id=server_id,
data={"name": name})
for name, was in prev_by_name.items():
if name not in curr_by_name:
await ev.record_event(db, type="container_removed", severity="info",
message=f"{name} (removed)", server_id=server_id,
data={"name": name, "image": was.get("image", "")})
elif not _is_exited(was["status"]) and _is_exited(curr_by_name[name]["status"]):
exit_code = _EXIT_RE.search(curr_by_name[name]["status"])
await ev.record_event(db, type="container_exited", severity="warning",
message=f"{name} (exited {'code ' + exit_code.group(1) if exit_code else '—'})",
server_id=server_id,
data={"name": name, **({"exit_code": int(exit_code.group(1))} if exit_code else {})})
def _is_exited(status: str) -> bool:
return status.startswith("Exited") or status.startswith("Dead")