Newer
Older
hard-panel / panel / backend / app / events.py
"""Единый журнал ивентов 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)