Newer
Older
hard-panel / panel / backend / app / services_probe.py
"""Фоновый пробер health-эндпоинтов сервисов.

Панель сама опрашивает настроенные в разделе «Сервисы» URL (GET) каждые
GHARD_HEALTH_INTERVAL секунд, все сервисы — параллельно. Формат отклика
health-эндпоинтов у gnexus-сервисов пока не договорён — работает
прогрессивный минимум (см. docs/health-endpoint-spec.md):

    2xx → up; тело-JSON со строковым status: ok-набор → up,
    degraded-набор → degraded, всё прочее → down;
    не-2xx / таймаут / отказ сети → down (точка пишется в любом случае,
    latency замеряется).

Старше 7 дней точки удаляются (политика хранения).
"""

import asyncio
import json
import time
from datetime import datetime, timedelta, timezone

import httpx

from app.config import get_settings
from app.db import get_db

RETENTION_DAYS = 7

# тело health в 64 KB не уместится — больше и читать не надо
MAX_BODY = 64 * 1024

# значение body.status (в lowercase) → состояние (docs/health-endpoint-spec)
OK_VALUES = {"ok", "healthy", "pass", "passing", "up"}
DEGRADED_VALUES = {"degraded", "warn", "warning", "partial", "busy"}


def classify(code: int, body: bytes | None) -> tuple[str, str]:
    """(state, message) по коду и телу отклика. code задан, если дошёл до HTTP."""
    if code < 200 or code >= 300:
        return "down", f"http {code}"
    if body:
        try:
            payload = json.loads(body.decode("utf-8", errors="replace"))
        except (ValueError, UnicodeError):
            return "up", ""
        if isinstance(payload, dict):
            status = payload.get("status")
            if isinstance(status, str):
                value = status.strip().lower()
                if value in OK_VALUES:
                    return "up", value
                if value in DEGRADED_VALUES:
                    return "degraded", value
                # не-словарный или незнакомый статус — счёт не в пользу сервиса
                return "down", value
    return "up", ""


async def probe_once(client: httpx.AsyncClient, url: str) -> dict:
    """Один GET: {state, code, latency_ms, message}."""
    started = time.monotonic()
    code: int | None = None
    message = ""
    state = "down"
    try:
        response = await client.send(client.build_request("GET", url), stream=True)
        code = response.status_code
        chunks: list[bytes] = []
        # health-тела маленькие: как только набрали MAX_BODY — дальше не читаем
        async for chunk in response.aiter_bytes(4096):
            chunks.append(chunk)
            if sum(map(len, chunks)) >= MAX_BODY:
                break
        body = b"".join(chunks)[:MAX_BODY]
        await response.aclose()
        state, message = classify(code, body)
    except httpx.HTTPError as exc:  # таймаут, DNS, отказ соединения
        text = str(exc)
        message = (text.split(";")[0][:200] if text else type(exc).__name__)
    except Exception as exc:  # noqa: BLE001 — любая ошибка = вниз, цикл не роняем
        message = str(exc)[:200]
    latency_ms = round((time.monotonic() - started) * 1000, 1)
    return {"state": state, "code": code, "latency_ms": latency_ms, "message": message}


async def sample_once() -> int:
    """Одна волна проб по всем сервисам. Возвращает число точек."""
    db = get_db()
    cursor = await db.execute("SELECT id, url FROM services")
    rows = await cursor.fetchall()
    if not rows:
        return 0
    settings = get_settings()
    async with httpx.AsyncClient(
        timeout=settings.health_timeout, follow_redirects=True, verify=False
    ) as client:
        results = await asyncio.gather(
            *(probe_once(client, row["url"]) for row in rows), return_exceptions=True
        )
    now = datetime.now(timezone.utc).isoformat()
    points = 0
    for row, result in zip(rows, results):
        if isinstance(result, Exception):  # проб по отдельному url не уронил волну
            probe = {"state": "down", "code": None, "latency_ms": 0.0, "message": str(result)[:200]}
        else:
            probe = result
        await db.execute(
            "INSERT INTO service_samples (service_id, ts, state, code, latency_ms, message)"
            " VALUES (?, ?, ?, ?, ?, ?)",
            (row["id"], now, probe["state"], probe["code"], probe["latency_ms"], probe["message"]),
        )
        points += 1
    cutoff = (datetime.now(timezone.utc) - timedelta(days=RETENTION_DAYS)).isoformat()
    await db.execute("DELETE FROM service_samples WHERE ts < ?", (cutoff,))
    await db.commit()
    return points


async def probe_service_once(service_id: int, url: str) -> dict:
    """Одиночная проба одного сервиса (POST/PATCH в API вызывает сразу)."""
    settings = get_settings()
    async with httpx.AsyncClient(
        timeout=settings.health_timeout, follow_redirects=True, verify=False
    ) as client:
        probe = await probe_once(client, url)
    db = get_db()
    await db.execute(
        "INSERT INTO service_samples (service_id, ts, state, code, latency_ms, message)"
        " VALUES (?, ?, ?, ?, ?, ?)",
        (
            service_id,
            datetime.now(timezone.utc).isoformat(),
            probe["state"], probe["code"], probe["latency_ms"], probe["message"],
        ),
    )
    await db.commit()
    return probe


async def loop() -> None:
    interval = get_settings().health_interval
    while True:
        try:
            await sample_once()
        except Exception as exc:  # не роняем цикл из-за одной ошибки
            print(f"services probe error: {exc}", flush=True)
        await asyncio.sleep(interval)