Newer
Older
hard-panel / panel / backend / app / metrics_retention.py
"""Разряжение истории метрик по возрасту точки + годовой горизонт.

Политика (границы — по возрасту, а не по календарным суткам: иначе на стыке
суток скачок плотности):

    0 – 24 ч          всё (сырой интервал агента)
    24 ч – 3 сут     1 точка/мин
    3 – 7 сут        1 точка/5 мин
    7 сут – 365 сут  1 точка/час
    старше 365 сут   удаляем

В каждом бакете остаётся реальная (последняя) точка — никаких средних, только
прореживание. НО чтобы прореживание не теряло крайности (CPU под 98% две
минуты, диск под 98%, провал до нуля), вместе с выжившей точкой сохраняется
дайджест бакета в `metrics.agg_json`:

    {"n": 114, "exp": 120,              # сколько точек пришло / сколько ожидалось
     "cpu": [3.2, 98.4], "ram": [21, 88], "swap": [0, 12],
     "load": [0.1, 3.4], "net": [1.2, 90.9], "out": [0.4, 60.6],
     "disks": {"/": [41, 98]}}          # min/max по каждой метрике и маунту

График по нему рисует полосу min–max (пики и провалы видно на всей глубине), а
`n` против `exp` показывает, что в бакете был простой: пакетов пришло меньше,
чем должно было (сам факт падения пишется ещё и в журнал событий, он живёт год).
Дайджест пишется только там, где бакет реально схлопывается (n > 1), поэтому
сырые точки последних суток остаются как есть, а повторный прогон ничего не
перезаписывает.

Строки бакета читаются без JSON (кроме дисков) в порядке ts, id; бакеты копятся
в словаре (не больше ~8600 бакетов на годовой ярус), удаление — по границам
бакета через индекс `idx_metrics_server_ts`, без списка id в памяти. Формат ts —
ISO-8601 UTC фиксированной длины, поэтому ключ бакета считается срезом строки;
границы бакета (для DELETE) собираются через datetime — один раз на схлопнутый
бакет, это даёт корректный перенос через час/сутки.

Первый прогон на старой БД — самый дорогой: в окне «7 сут – год» может лежать
~1 млн точек на сервер (интервал 30 с) — это минуты работы; дальше окна уже
разряжены и проходы дешёвые.

Замечание про соседей ingest'а: тяжёлые JSON последнего пакета живут в
server_state, поэтому прореживание истории ничего в приёме пакетов не ломает.
"""

import asyncio
import json
from datetime import datetime, timedelta, timezone

from app.db import get_db, incremental_vacuum

# (возраст старше, возраст не старше, ширина бакета) — в порядке ярусов
TIERS: list[tuple[timedelta, timedelta, str]] = [
    (timedelta(days=3), timedelta(days=1), "minute"),
    (timedelta(days=7), timedelta(days=3), "5min"),
    (timedelta(days=365), timedelta(days=7), "hour"),
]
HORIZON = timedelta(days=365)

LOOP_INTERVAL_S = 900  # 15 минут

# ширина бакета в секундах — из неё считается «сколько замеров ожидалось» (exp)
BUCKET_WIDTH_S = {"minute": 60, "5min": 300, "hour": 3600}


def _bucket(ts: str, width: str) -> str:
    """Ключ бакета по строке ISO-8601 (`2026-10-05T08:59:31.759686+00:00`)."""
    if width == "hour":
        return ts[:13]
    if width == "5min":
        return ts[:14] + str(int(ts[14:16]) // 5)
    return ts[:16]


def _bucket_bounds(key: str, width: str) -> tuple[str, str]:
    """ISO-границы бакета (включительно) — удаляем по ним, не держа id в памяти."""
    if width == "hour":
        start = datetime.fromisoformat(key + ":00:00+00:00")
        step = timedelta(hours=1)
    elif width == "5min":
        start = datetime.fromisoformat(f"{key[:14]}{int(key[14:]) * 5:02d}:00+00:00")
        step = timedelta(minutes=5)
    else:  # minute
        start = datetime.fromisoformat(key + ":00+00:00")
        step = timedelta(minutes=1)
    end = start + step - timedelta(microseconds=1)
    return start.isoformat(), end.isoformat()


class _Bucket:
    """Накопитель одного бакета: крайности по метрикам + число точек."""

    __slots__ = ("n", "survivor", "lo", "hi", "disks_lo", "disks_hi", "has_ram", "has_swap")

    def __init__(self, row) -> None:
        self.n = 0
        self.survivor = 0
        self.lo: dict[str, float] = {}
        self.hi: dict[str, float] = {}
        self.disks_lo: dict[str, float] = {}
        self.disks_hi: dict[str, float] = {}
        self.has_ram = False
        self.has_swap = False
        self.add(row)

    def add(self, row) -> None:
        self.n += 1
        if row["id"] > self.survivor:
            self.survivor = row["id"]
        self._put("cpu", row["cpu"])
        self._put("load", row["load1"])
        self._put("net", row["net_in_mbs"])
        self._put("out", row["net_out_mbs"])
        if row["ram_total"]:
            self.has_ram = True
            self._put("ram", row["ram_used"] / row["ram_total"] * 100.0)
        if row["swap_total"]:
            self.has_swap = True
            self._put("swap", row["swap_used"] / row["swap_total"] * 100.0)
        try:
            for disk in json.loads(row["disks_json"]):
                mount, pct = disk["mount"], float(disk["percent"])
                self._disk_put(mount, pct)
        except (ValueError, KeyError, TypeError, json.JSONDecodeError):
            pass

    def _put(self, name: str, value: float) -> None:
        if value is None:
            return
        if name not in self.lo or value < self.lo[name]:
            self.lo[name] = value
        if name not in self.hi or value > self.hi[name]:
            self.hi[name] = value

    def _disk_put(self, mount: str, pct: float) -> None:
        if mount not in self.disks_lo or pct < self.disks_lo[mount]:
            self.disks_lo[mount] = pct
        if mount not in self.disks_hi or pct > self.disks_hi[mount]:
            self.disks_hi[mount] = pct

    def digest(self, exp: int) -> str:
        """Компактный дайджест бакета: крайности (1 знак у процентов, 2 у чисел)."""
        pct = {"cpu", "ram", "swap"}
        out: dict = {"n": self.n, "exp": exp}
        for name in ("cpu", "ram", "swap", "load", "net", "out"):
            if name in self.lo:
                r = 1 if name in pct else 2
                out[name] = [round(self.lo[name], r), round(self.hi[name], r)]
        if self.disks_lo:
            out["disks"] = {
                mount: [round(self.disks_lo[mount], 1), round(self.disks_hi[mount], 1)]
                for mount in self.disks_lo
            }
        return json.dumps(out, ensure_ascii=False, separators=(",", ":"))


async def _thin_server_tier(db, server_id: int, older: str, younger: str, width: str,
                            interval: int) -> tuple[int, int]:
    """Схлопнуть бакеты окна (older, younger] одного сервера до одной точки.

    Возвращает (удалено строк, схлопнуто бакетов). В выжившую точку пишется
    дайджест бакета — крайности за интервал не теряются.
    """
    cursor = await db.execute(
        """SELECT id, ts, cpu, load1, ram_used, ram_total, swap_used, swap_total,
                  net_in_mbs, net_out_mbs, disks_json
           FROM metrics WHERE server_id = ? AND ts < ? AND ts >= ? ORDER BY ts, id""",
        (server_id, younger, older),
    )
    buckets: dict[str, _Bucket] = {}
    async for row in cursor:
        key = _bucket(row["ts"], width)
        bucket = buckets.get(key)
        if bucket is None:
            buckets[key] = _Bucket(row)
        else:
            bucket.add(row)
    await cursor.close()
    if not buckets:
        return 0, 0

    exp = max(1, round(BUCKET_WIDTH_S[width] / max(interval, 1)))
    deleted = 0
    flushed = 0
    for key, bucket in buckets.items():
        if bucket.n <= 1:
            continue  # схлопывать нечего (и повторный прогон ничего не переписывает)
        await db.execute(
            "UPDATE metrics SET agg_json = ? WHERE id = ?", (bucket.digest(exp), bucket.survivor)
        )
        start, end = _bucket_bounds(key, width)
        cursor = await db.execute(
            """DELETE FROM metrics WHERE server_id = ? AND ts >= ? AND ts <= ? AND id != ?""",
            (server_id, start, end, bucket.survivor),
        )
        deleted += cursor.rowcount or 0
        flushed += 1
    return deleted, flushed


async def thin_once() -> dict:
    """Один проход: схлопнуть ярусы и отрезать всё старше горизонта."""
    db = get_db()
    now = datetime.now(timezone.utc)

    cursor = await db.execute("SELECT id, interval FROM servers")
    servers = [(row["id"], row["interval"] or 30) for row in await cursor.fetchall()]

    deleted = 0
    buckets = 0
    for older, younger, width in TIERS:
        older_iso = (now - older).isoformat()
        younger_iso = (now - younger).isoformat()
        for server_id, interval in servers:
            gone, flushed = await _thin_server_tier(
                db, server_id, older_iso, younger_iso, width, interval
            )
            deleted += gone
            buckets += flushed

    # горизонт: всё, что старше года, не нужно вовсе
    cursor = await db.execute("DELETE FROM metrics WHERE ts < ?", ((now - HORIZON).isoformat(),))
    deleted += cursor.rowcount or 0
    await db.commit()
    if deleted:
        # вернуть файлу освободившиеся страницы (нужен auto_vacuum=INCREMENTAL)
        await incremental_vacuum(2000)
    return {"deleted": deleted, "buckets": buckets}


async def retention_loop() -> None:
    """Проход сразу при старте (деплой начинает худеть немедленно), далее по таймеру."""
    while True:
        try:
            result = await thin_once()
            if result["deleted"]:
                print(
                    f"metrics retention: удалено {result['deleted']} точек, "
                    f"схлопнуто бакетов {result['buckets']}",
                    flush=True,
                )
        except Exception as exc:  # не роняем цикл из-за одной ошибки
            print(f"metrics retention error: {exc}", flush=True)
        await asyncio.sleep(LOOP_INTERVAL_S)