"""Фоновый пробер 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, j
from app.events import record_event
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"}
# самосвидетельство §5 спеки: белые списки и лимиты
REPORT_MODES = {"maintenance", "readonly", "draining"}
REPORT_MAX_GAUGES = 20
def load_json(body: bytes | None) -> dict | list | None:
"""JSON из тела health — всё прочее → None (не-JSON тело норма)."""
if not body:
return None
try:
return json.loads(body.decode("utf-8", errors="replace"))
except (ValueError, UnicodeError):
return None
def classify(code: int, payload: dict | list | None) -> tuple[str, str]:
"""(state, message) по коду и разобранному телу; payload = load_json(body)."""
if code < 200 or code >= 300:
return "down", f"http {code}"
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", ""
def parse_report(payload: dict | list | None) -> dict:
"""Самосвидетельство §5: берём ТОЛЬКО известные поля.
Сырое тело не храним и не логируем: эндпоинт без авторизации, неизвестные
ключи могут содержать ПД. Неизвестное значение mode игнорируется — mode
не статус лайвности.
"""
if not isinstance(payload, dict):
return {}
report: dict = {}
for key in ("service", "version", "build", "environment"):
value = payload.get(key)
if isinstance(value, str) and (trimmed := value.strip()[:80]):
report[key] = trimmed
if isinstance(payload.get("uptime"), (int, float)) and not isinstance(payload["uptime"], bool):
report["uptime"] = int(payload["uptime"])
since = payload.get("since")
if isinstance(since, str):
try:
uptime = (datetime.now(timezone.utc) - datetime.fromisoformat(since)).total_seconds()
report["uptime"] = int(max(0, uptime))
except ValueError:
pass
mode = payload.get("mode")
if isinstance(mode, str) and (clean := mode.strip().lower()) in REPORT_MODES:
report["mode"] = clean
gauges = payload.get("gauges")
if isinstance(gauges, dict):
picked = {}
for name, value in list(gauges.items())[:REPORT_MAX_GAUGES]:
# bool — подмножество int: "cpu_warn": true считается сломанным, пропускаем
if isinstance(value, (int, float)) and not isinstance(value, bool) and isinstance(name, str):
picked[name[:40]] = round(float(value), 4)
if picked:
report["gauges"] = picked
return report
async def probe_once(client: httpx.AsyncClient, url: str) -> dict:
"""Один GET: {state, code, latency_ms, message, report}."""
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()
payload = load_json(body)
state, message = classify(code, payload)
report = parse_report(payload) # есть и на down (http-код): версия полезна
except httpx.HTTPError as exc: # таймаут, DNS, отказ соединения
report = {}
text = str(exc)
message = (text.split(";")[0][:200] if text else type(exc).__name__)
except Exception as exc: # noqa: BLE001 — любая ошибка = вниз, цикл не роняем
report = {}
message = str(exc)[:200]
latency_ms = round((time.monotonic() - started) * 1000, 1)
return {"state": state, "code": code, "latency_ms": latency_ms, "message": message,
"report": report}
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()
prev_states = await _last_states(db, [row["id"] for row in rows])
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], "report": {}}
else:
probe = result
await db.execute(
"INSERT INTO service_samples (service_id, ts, state, code, latency_ms, message, report)"
" VALUES (?, ?, ?, ?, ?, ?, ?)",
(row["id"], now, probe["state"], probe["code"], probe["latency_ms"],
probe["message"], j(probe["report"])),
)
# журнал ивентов: открытие/закрытие инцидент-периода
await _transition_events(
prev_states.get(row["id"]), probe, service_id=row["id"], db=db,
)
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 _last_states(db, service_ids: list[int]) -> dict[int, str]:
"""Предыдущее состояние каждого сервиса (до вставки новых точек)."""
states: dict[int, str] = {}
for sid in service_ids:
cursor = await db.execute(
"SELECT state FROM service_samples WHERE service_id = ? ORDER BY id DESC LIMIT 1",
(sid,),
)
row = await cursor.fetchone()
if row is not None:
states[sid] = row["state"]
return states
async def _transition_events(prev_state: str | None, probe: dict, *, service_id: int, db) -> None:
"""Ивенты на краях периодов (журнал + Synapse): down/degraded/recovered.
Период открывается первым не-up замером и закрывается первым up —
спама на каждую не-up проб нет (переходы фиксируются по смене).
"""
new_state = probe["state"]
if prev_state == new_state:
return
# имя сервиса нужно в сообщении/конверте — один дешёвый запрос
cursor = await db.execute("SELECT id, name, url FROM services WHERE id = ?", (service_id,))
row = await cursor.fetchone()
svc = dict(row) if row is not None else {"id": service_id, "name": "?", "url": ""}
if prev_state in (None, "up", "pending") and new_state in ("down", "degraded"):
# первый замер не-up (в т.ч. на новом сервисе) — период открылся
await record_event(
db, type=f"service_{new_state}",
severity="critical" if new_state == "down" else "warning",
message=_service_message(new_state, probe, svc),
data=_service_data(new_state, probe, svc),
)
elif prev_state == "degraded" and new_state == "down":
# эскалация внутри периода
await record_event(
db, type="service_down", severity="critical",
message=_service_message(new_state, probe, svc),
data=_service_data(new_state, probe, svc),
)
elif prev_state in ("down", "degraded") and new_state == "up":
await record_event(
db, type="service_recovered", severity="info",
message=f"{svc['name']} снова в порядке",
data={"service_id": svc["id"], "name": svc["name"]},
)
def _service_message(state: str, probe: dict, service: dict) -> str:
where = f" ({probe['message']})" if probe["message"] else ""
return f"{service['name']} — {state}{where}"
def _service_data(state: str, probe: dict, service: dict) -> dict:
data = {
"type": f"service_{state}",
"service_id": service["id"],
"name": service["name"],
"url": service["url"],
"state": state,
}
if probe.get("code") is not None:
data["code"] = probe["code"]
if probe.get("message"):
data["message"] = probe["message"][:200]
return data
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()
cursor = await db.execute(
"SELECT state FROM service_samples WHERE service_id = ? ORDER BY id DESC LIMIT 1",
(service_id,),
)
prev_row = await cursor.fetchone()
await db.execute(
"INSERT INTO service_samples (service_id, ts, state, code, latency_ms, message, report)"
" VALUES (?, ?, ?, ?, ?, ?, ?)",
(
service_id,
datetime.now(timezone.utc).isoformat(),
probe["state"], probe["code"], probe["latency_ms"], probe["message"],
j(probe["report"]),
),
)
await _transition_events(prev_row["state"] if prev_row else None, probe,
service_id=service_id, db=db)
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)