"""Фоновый пробер 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)