"""Единый журнал ивентов 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)