diff --git a/README.md b/README.md index 5d4ce9e..9a57e9d 100644 --- a/README.md +++ b/README.md @@ -245,6 +245,10 @@ | `GHARD_PANEL_PORT` | `8000` | host-порт (docker-compose) | | `GHARD_OFFLINE_MULTIPLIER` | `3` | offline = нет пакетов дольше N × interval | | `GHARD_SHARE_INTERVAL` | `60` | период опроса сетевых хранилищ, сек | +| `GHARD_SYNAPSE_URL` | *(пусто)* | URL Synapse (Gnexus Synapse Integration); пусто — отчёт о событиях отключён | +| `GHARD_SYNAPSE_API_KEY` | *(пусто)* | ключ `syn_*` источника hard-panel (выдаётся один раз в админке Synapse) | +| `GHARD_SYNAPSE_DEFAULT_SOURCE` | `hard-panel` | имя источника, зарегистрированное в Synapse | +| `GHARD_SYNAPSE_TIMEOUT` | `10` | таймаут HTTP Synapse, сек | ## API (кратко) @@ -262,6 +266,8 @@ - `GET /api/v1/services/{id}/incidents?limit=50` — журнал моментов недоступности (периоды down+degraded из истории проб; у идущего периода `end=null`) +- `GET /api/v1/events?limit=100` — единый журнал событий (все серверы и + сервисы, по убыванию; 1..1000) - `GET/PATCH /api/v1/me` — профиль + язык (GET), override языка (PATCH; только cookie-сессия — у Bearer-админа нет личности) @@ -293,6 +299,50 @@ «продолжается»). Там же — редактирование (имя/URL; смена URL даёт немедленный повторный замер) и удаление сервиса. +## Журнал событий и интеграция с Synapse + +Панель — источник живности экосистемы: всё, что она замечает, пишется в +единый журнал (`events`, хранение 90 дней; читается `GET /api/v1/events`) и — +если настроен Synapse — уходит конвертом v1 в +[Gnexus Synapse](https://git.gnexus.space/root/gn-synapse) (хаб всех исходящих +уведомлений экосистемы, конвенция в handbook +[10-platform/notifications.md](https://git.gnexus.space/root/gnexus-handbook/raw/master/10-platform/notifications.md)). +Журнал же служит состоянием: пара «открытие/закрытие» (например +`server_offline`/`server_online`) переживает рестарт панели — повторных +оповещений нет. + +Что фиксируется: + +| Событие | Когда | severity журнала | Конверт: subject × action | priority | +|---|---|---|---|---| +| `cpu_high` / `cpu_recovered` | cpu ≥90% / <80% (гистерезис) | warning | `resource` × `over_threshold` / `back_normal` | high / low | +| `ram_high` / `ram_recovered` | ram ≥90% / <80% | warning | как cpu | high / low | +| `swap_high` / `swap_recovered` | swap ≥60% / <40% | warning | как cpu | high / low | +| `disk_high` / `disk_recovered` | диск (mount) ≥90% / <85% | critical | как cpu | high / low | +| `load_high` / `load_recovered` | load1 ≥2×ядер / <1× | warning | как cpu | high / low | +| `server_offline` / `server_online` | нет пакетов дольше interval × N / данные снова идут | warning / info | `server` × `offline` / `online` | high / low | +| `container_added`/`removed`/`started` | появился / исчез / снова работает | info | `container` × `added` / `removed` / `started` | low | +| `container_exited` | контейнер завершился (код в data) | warning | `container` × `exited` | normal | +| `service_down` / `service_degraded` | первый замер вниз из up / деградация | critical / warning | `service` × `down` / `degraded` | high / normal | +| `service_recovered` | сервис снова в порядке | info | `service` × `recovered` | low | + +Повторяющиеся открытия (порог/офлайн/сервис — всё, что может «пылить») +получают `dedup_key` вида `{type}-{id}-{UTC-день}` (дедупликация в Synapse, +окно 24 ч); `ttl_seconds` = 24 ч — устаревший алерт не маршрутизируется. +Конверты не содержат адресатов и каналов — «куда доставить», по конвенции +экосистемы, решает Synapse по своим правилам. Приоритеты high/critical +отправляются `send()` (ошибка видна в логе панели), остальные — +fire-and-forget `emit()`: сбой доставки никогда не ломает приём метрик и +пробы. + +Подключение: сначала зарегистрируйте источник `hard-panel` в админке Synapse +(тройки `subject × action` регистрируются один раз там же), ключ `syn_*` +показывается один раз — в `.env`/docker env (`GHARD_SYNAPSE_URL`, +`GHARD_SYNAPSE_API_KEY`) и в gnexus-creds. Пустой `GHARD_SYNAPSE_API_KEY` = +интеграция отключена (панель молча пишет только журнал). SDK — +[gn-synapse-client-py](https://git.gnexus.space/root/gn-synapse-client-py) +(`gnexus-synapse`, async-клиент `AsyncSynapseClient`, один инстанс на процесс). + ## Язык интерфейса Три языка: en, uk, ru. По умолчанию панель следует языку аккаунта gnexus-auth diff --git a/panel/backend/app/api/events.py b/panel/backend/app/api/events.py new file mode 100644 index 0000000..2f4b094 --- /dev/null +++ b/panel/backend/app/api/events.py @@ -0,0 +1,40 @@ +"""GET /api/v1/events — единый журнал ивентов (генерация в app/events.py +и точках генерации: ingest, services_probe, offline-loop). Журнал читают +UI и будущие потребители; в Synapse ивенты уходят автоматически. +""" + +from fastapi import APIRouter, Depends, Query + +from app.db import get_db +from app.security import require_admin + +router = APIRouter(prefix="/api/v1", dependencies=[Depends(require_admin)]) + + +@router.get("/events") +async def list_events(limit: int = Query(default=100, ge=1, le=1000)) -> list[dict]: + "Свежие ивенты по убыванию (все серверы)." + db = get_db() + cursor = await db.execute( + """SELECT e.id, e.server_id, e.type, e.severity, e.message, e.data_json, e.ts, + s.name AS server_name + FROM events e LEFT JOIN servers s ON s.id = e.server_id + ORDER BY e.id DESC LIMIT ?""", + (limit,), + ) + rows = await cursor.fetchall() + import json as _json + + return [ + { + "id": row["id"], + "server_id": row["server_id"], + "server_name": row["server_name"], + "type": row["type"], + "severity": row["severity"], + "message": row["message"], + "data": _json.loads(row["data_json"]), + "ts": row["ts"], + } + for row in rows + ] \ No newline at end of file diff --git a/panel/backend/app/api/ingest.py b/panel/backend/app/api/ingest.py index 98b4ed9..aa0390f 100644 --- a/panel/backend/app/api/ingest.py +++ b/panel/backend/app/api/ingest.py @@ -1,15 +1,18 @@ """POST /api/v1/ingest — приём пакета метрик от hard-monitor. Панель — мозг: здесь считается скорость сети (дельта сырых счётчиков), -обновляется карточка сервера и пишется точка метрик. Ивенты (пороги, -diff дисков/docker, offline) появятся на этапе 4 — тоже здесь. +обновляется карточка сервера и пишется точка метрик. Из потока пакетов +тут же генерируются ивенты: восстановление после 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 @@ -20,6 +23,16 @@ # «Онлайн», если последний пакет был не позже чем 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) @@ -80,14 +93,15 @@ now = _now() packet_ts = _iso(payload.ts) if payload.ts else _iso(now) - # предыдущая точка — нужны сырые счётчики для дельты сети + # предыдущая точка — нужны сырые счётчики для дельты сети и diff docker cursor = await db.execute( - "SELECT ts, net_json FROM metrics WHERE server_id = ? ORDER BY id DESC LIMIT 1", + "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"]) @@ -96,6 +110,10 @@ 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) @@ -141,5 +159,108 @@ 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"} \ No newline at end of file + 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") \ No newline at end of file diff --git a/panel/backend/app/config.py b/panel/backend/app/config.py index c11a55a..8e2b1f6 100644 --- a/panel/backend/app/config.py +++ b/panel/backend/app/config.py @@ -31,6 +31,12 @@ health_interval: int = 30 # таймаут одного health-запроса (сек) health_timeout: float = 5.0 + # Gnexus Synapse (этап 4, репортер ивентов): пустой api_key = интеграция + # выключена. default_source — зарегистрированное в админке Synapse имя + synapse_url: str = "" + synapse_api_key: str = "" + synapse_default_source: str = "hard-panel" + synapse_timeout: float = 10.0 @lru_cache diff --git a/panel/backend/app/events.py b/panel/backend/app/events.py new file mode 100644 index 0000000..3c17f97 --- /dev/null +++ b/panel/backend/app/events.py @@ -0,0 +1,246 @@ +"""Единый журнал ивентов hard-panel + репортер Gnexus Synapse (этап 4). + +Панель — мозг: все ивенты генерируются здесь из потока метрик/проб и пишутся +в таблицу events (retention 90 дней). Если задан GHARD_SYNAPSE_API_KEY — +каждый ивент уходит конвертом v1 в Synapse (handbook +10-platform/notifications.md): клиент gn-synapse-client-py, один инстанс +на процесс, отправка fire-and-forget — свою очередь ретраев не строим, +доставку ретраит Synapse по своим правилам маршрутизации. +""" + +import asyncio +import json +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone + +from app.config import get_settings +from app.db import get_db + +EVENTS_RETENTION_DAYS = 90 +# ttl конверта: алерты актуальности — сутки, дальше не маршрутизируются +TTL_SECONDS = 24 * 3600 + + +# --- Маппинг ивент → конверт Synapse ----------------------------------------- + +@dataclass(frozen=True) +class Sink: + subject: str + action: str + priority: str # low | normal | high | critical (шкала Synapse) + dedup: bool # повторяющиеся открытия — dedup по типу+сущность+UTC-день + + +THRESHOLDS = ("cpu", "ram", "swap", "disk", "load") + +THRESHOLD_SINKS: dict[str, Sink] = {} +for _m in THRESHOLDS: + THRESHOLD_SINKS[f"{_m}_high"] = Sink("resource", "over_threshold", "high", True) + THRESHOLD_SINKS[f"{_m}_recovered"] = Sink("resource", "back_normal", "low", False) + +SINKS: dict[str, Sink] = { + "server_offline": Sink("server", "offline", "high", True), + "server_online": Sink("server", "online", "low", False), + "container_added": Sink("container", "added", "low", False), + "container_removed": Sink("container", "removed", "low", False), + "container_started": Sink("container", "started", "low", False), + "container_exited": Sink("container", "exited", "normal", False), + "service_down": Sink("service", "down", "high", True), + "service_degraded": Sink("service", "degraded", "normal", True), + "service_recovered": Sink("service", "recovered", "low", False), + **THRESHOLD_SINKS, +} + +# Пары state machine «открыт / закрыт»: state_open смотрит последнее из двух +# событий в журнале (сам журнал — состояние, переживает рестарт панели). +STATE_PAIRS: dict[str, str] = { + "server_offline": "server_online", + **{f"{_m}_high": f"{_m}_recovered" for _m in THRESHOLDS}, +} + + +# --- Клиент Synapse ----------------------------------------------------------- + +_client = None +_warned_no_config = False + + +def synapse_client(): + global _client, _warned_no_config + if _client is not None: + return _client + settings = get_settings() + if not settings.synapse_api_key or not settings.synapse_url: + if not _warned_no_config: + print("synapse reporter disabled (GHARD_SYNAPSE_API_KEY / SYNAPSE_URL unset)", flush=True) + _warned_no_config = True + return None + from gnexus_synapse import AsyncSynapseClient + + _client = AsyncSynapseClient( + settings.synapse_url, + settings.synapse_api_key, + timeout=settings.synapse_timeout, + default_source=settings.synapse_default_source or "hard-panel", + ) + return _client + + +def _ship(coro) -> None: + """Fire-and-forget: падения Synapse не роняют вызвавший код (лог + всё).""" + task = asyncio.create_task(coro) + task.add_done_callback(_ship_done) + + +def _ship_done(task: asyncio.Task) -> None: + if not task.cancelled() and task.exception() is not None: + print(f"synapse report failed: {task.exception()}", flush=True) + + +def _utc_date() -> str: + return datetime.now(timezone.utc).strftime("%Y%m%d") + + +# --- Журнал ------------------------------------------------------------------- + +async def record_event( + db, + *, + type: str, + severity: str, + message: str, + server_id: int | None = None, + data: dict | None = None, +) -> None: + """Записать ивент в журнал и (если настроен) уйти конвертом в Synapse. + + severity — для журнала/UI (info|warning|critical); приоритет конверта + задан по типу (SINKS). dedup_key повторяющихся открытий — тип + сущность + + число UTC-дня (best-effort, окно dedup Synapse 24 ч). + """ + now = datetime.now(timezone.utc) + await db.execute( + """INSERT INTO events (server_id, type, severity, message, data_json, ts) + VALUES (?, ?, ?, ?, ?, ?)""", + (server_id, type, severity, message, json.dumps(data or {}, ensure_ascii=False), now.isoformat()), + ) + await db.commit() + + sink = SINKS.get(type) + if sink is None: + return + client = synapse_client() + if client is None: + return + + payload = dict(data or {}) + payload["type"] = type + priority = sink.priority + if server_id is not None: + payload.setdefault("server_id", server_id) + + dedup = None + if sink.dedup: + dedup = f"{type}-{server_id or payload.get('service_id')}-{_utc_date()}" + + # инциденты (high/critical) — send (ошибки видны в логе), остальное — emit: + # клиент сам молча логнет, бизнес-код ничего не ловит + if priority in ("high", "critical"): + _ship( + client.send( + subject=sink.subject, + action=sink.action, + priority=priority, + payload=payload, + dedup_key=dedup, + ttl_seconds=TTL_SECONDS, + ) + ) + else: + _ship( + client.emit( + subject=sink.subject, + action=sink.action, + priority=priority, + payload=payload, + dedup_key=dedup, + ttl_seconds=TTL_SECONDS, + ) + ) + + +# --- State machine открытия/закрытия (по журналу) ------------------------------ + +async def state_open(db, server_id: int, open_type: str) -> bool: + """Пара (…_high, …_recovered): открыта, если последнее из двух — открытие.""" + cursor = await db.execute( + "SELECT type FROM events WHERE server_id = ? AND type IN (?, ?) ORDER BY id DESC LIMIT 1", + (server_id, open_type, STATE_PAIRS[open_type]), + ) + row = await cursor.fetchone() + return row is not None and row["type"] == open_type + + +async def state_open_ext(db, server_id: int, open_type: str, *, mount: str | None = None) -> bool: + """Вариант с доп. ключом сущности (disk-маунты: пара на каждый маунт).""" + if mount is None: + return await state_open(db, server_id, open_type) + cursor = await db.execute( + """SELECT type FROM events + WHERE server_id = ? AND type IN (?, ?) + AND json_extract(data_json, '$.mount') = ? + ORDER BY id DESC LIMIT 1""", + (server_id, open_type, STATE_PAIRS[open_type], mount), + ) + row = await cursor.fetchone() + return row is not None and row["type"] == open_type + + +async def close_if_open(db, server_id: int, close_type: str, message: str, data: dict) -> None: + """Закрыть пару, если открыта (приход пакета → recovered, online).""" + open_type = {v: k for k, v in STATE_PAIRS.items()}[close_type] + if await state_open(db, server_id, open_type): + await record_event( + db, type=close_type, severity="info", message=message, server_id=server_id, data=data + ) + + +# --- Watchdog offline --------------------------------------------------------- + +async def offline_loop() -> None: + """Нет пакета дольше interval * offline_multiplier → server_offline. + + Журнал — state: повторные прогоны не спамят (пара offline/online). + Тут же ретеншн ивентов (EVENTS_RETENTION_DAYS). + """ + while True: + try: + db = get_db() + settings = get_settings() + now = datetime.now(timezone.utc) + cursor = await db.execute( + "SELECT id, name, hostname, interval, last_seen FROM servers WHERE last_seen IS NOT NULL" + ) + for row in await cursor.fetchall(): + if await state_open(db, row["id"], "server_offline"): + continue + interval = max(row["interval"] or 30, 15) + limit = max(180.0, interval * settings.offline_multiplier) + silent = (now - datetime.fromisoformat(row["last_seen"])).total_seconds() + # grace: короткие моргания агента ивента не дают + if silent > limit: + minutes = int(silent // 60) + await record_event( + db, + type="server_offline", + severity="warning", + message=f"{row['name']} — нет данных {minutes} мин", + server_id=row["id"], + data={"hostname": row["hostname"], "silent_minutes": minutes}, + ) + cutoff = (now - timedelta(days=EVENTS_RETENTION_DAYS)).isoformat() + await db.execute("DELETE FROM events WHERE ts < ?", (cutoff,)) + await db.commit() + except Exception as exc: # не роняем цикл из-за одной ошибки + print(f"offline watchdog error: {exc}", flush=True) + await asyncio.sleep(15) \ No newline at end of file diff --git a/panel/backend/app/main.py b/panel/backend/app/main.py index 643b8eb..8197f7b 100644 --- a/panel/backend/app/main.py +++ b/panel/backend/app/main.py @@ -10,8 +10,9 @@ from app.security import require_admin -from app.api import auth_routes, ingest, servers, services, shares +from app.api import auth_routes, events, ingest, servers, services, shares from app.db import close_db, init_db +from app.events import offline_loop from app.mcp import mcp, mcp_asgi from app.services_probe import loop as services_probe_loop from app.shares_probe import loop as shares_probe_loop @@ -30,6 +31,8 @@ health_task = asyncio.create_task(services_probe_loop()) # подчистка expired-сессий и oauth state (каждые 10 минут) gc_task = asyncio.create_task(auth_gc_loop()) + # watchdog offline (нет пакетов дольше interval * multiplier) + ретеншн ивентов + watchdog_task = asyncio.create_task(offline_loop()) # lifespan смонтированных sub-app не вызывается — MCP session manager # стартуем здесь, иначе /mcp отвечает 500 ("Task group is not initialized") async with mcp.session_manager.run(): @@ -37,6 +40,7 @@ probe_task.cancel() health_task.cancel() gc_task.cancel() + watchdog_task.cancel() await close_db() @@ -70,6 +74,7 @@ app.include_router(servers.router) app.include_router(shares.router) app.include_router(services.router) +app.include_router(events.router) # MCP для ИИ-агентов (streamable HTTP, stateless) — POST/GET/DELETE /mcp. # Доступ как у REST API (require_admin): Bearer-токен, если задан; diff --git a/panel/backend/app/services_probe.py b/panel/backend/app/services_probe.py index f3c8c7f..0ee881d 100644 --- a/panel/backend/app/services_probe.py +++ b/panel/backend/app/services_probe.py @@ -22,6 +22,7 @@ from app.config import get_settings from app.db import get_db, j +from app.events import record_event RETENTION_DAYS = 7 @@ -149,6 +150,7 @@ *(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 не уронил волну @@ -162,6 +164,10 @@ (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,)) @@ -169,6 +175,77 @@ 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() @@ -177,6 +254,11 @@ ) 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 (?, ?, ?, ?, ?, ?, ?)", @@ -187,6 +269,8 @@ 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 diff --git a/panel/backend/pyproject.toml b/panel/backend/pyproject.toml index cac5722..ce1aa6f 100644 --- a/panel/backend/pyproject.toml +++ b/panel/backend/pyproject.toml @@ -11,6 +11,8 @@ "pydantic-settings>=2.3", "mcp>=1.10,<2", "gnexus-gauth @ git+https://git.gnexus.space/git/root/gnexus-auth-client-py.git", + # v0.1.2 — на сервере теги не запушены, пин по коммиту (тег v0.1.2) + "gnexus-synapse @ git+https://git.gnexus.space/git/root/gn-synapse-client-py.git@1834128d3c27c2cf937b3a01f1be6277288e275b", ] [build-system]