diff --git a/panel/backend/app/api/servers.py b/panel/backend/app/api/servers.py index 66811ce..95ed8ea 100644 --- a/panel/backend/app/api/servers.py +++ b/panel/backend/app/api/servers.py @@ -177,7 +177,12 @@ until: datetime | None = Query(default=None), limit: int = Query(default=500, ge=1, le=5000), ) -> list[dict]: - """История метрик (графики): время, cpu, ram, сеть. По возрастанию ts.""" + """История метрик (графики): время, cpu, ram, сеть. По возрастанию ts. + + У схлопнутых бакетов (старше суток) приходит `agg` — дайджест интервала + (min/max по метрикам, n/exp замеров): пики и простои не теряются при + разряжении. + """ db = get_db() conditions = ["server_id = ?"] params: list = [server_id] @@ -190,7 +195,8 @@ params.append(limit) cursor = await db.execute( f"""SELECT ts, cpu, load1, load5, load15, ram_used, ram_total, - swap_used, swap_total, uptime, net_in_mbs, net_out_mbs, disks_json + swap_used, swap_total, uptime, net_in_mbs, net_out_mbs, + disks_json, agg_json FROM metrics WHERE {' AND '.join(conditions)} ORDER BY ts DESC LIMIT ?""", params, @@ -207,6 +213,7 @@ "net_in_mbs": row["net_in_mbs"], "net_out_mbs": row["net_out_mbs"], "disks": json.loads(row["disks_json"]), + "agg": json.loads(row["agg_json"] or "{}") or None, } for row in reversed(rows) ] \ No newline at end of file diff --git a/panel/backend/app/db.py b/panel/backend/app/db.py index d09a5ac..b09408e 100644 --- a/panel/backend/app/db.py +++ b/panel/backend/app/db.py @@ -36,6 +36,7 @@ await _try_alter("ALTER TABLE services ADD COLUMN site_title TEXT") await _try_alter("ALTER TABLE services ADD COLUMN icon TEXT") await _try_alter("ALTER TABLE services ADD COLUMN meta_fetched_at TEXT") + await _try_alter("ALTER TABLE metrics ADD COLUMN agg_json TEXT NOT NULL DEFAULT '{}'") await _conn.commit() diff --git a/panel/backend/app/mcp.py b/panel/backend/app/mcp.py index 688c026..371616f 100644 --- a/panel/backend/app/mcp.py +++ b/panel/backend/app/mcp.py @@ -59,7 +59,7 @@ откликов + журнал моментов недоступности. - **Сетевые хранилища** — примонтированные NAS-шары и их заполнение. - **Журнал событий** — всё, что панель зафиксировала: offline/пороги/ - контейнеры/сервисы (хранение 90 дней). + контейнеры/сервисы (хранение — год). ## Язык Отвечай пользователю на его языке. Вывод тулов технический (статусы @@ -125,7 +125,10 @@ - `server_history(hours)`: разумно 1 (детали), 6, 24 (сутки), 168 (неделя) или 720 (месяц). История разряжается по возрасту точки: первые сутки — как прислал агент, 1–3 сут — точка в минуту, 3–7 сут — точка в 5 минут, дальше — - точка в час; старше года точки удаляются. + точка в час; старше года точки удаляются. Вопросы вида «был ли всплеск или + простой» решай по `agg` точки (min/max за интервал и n/exp замеров), а не по + значению самой точки: точка — это последний замер интервала, а пик мог быть + внутри. - Числа передавай как есть из тулов (проценты, МБ/с, ms) — не пересчитывай единицы самостоятельно. - Если сервер offline — не паникуй и не повторяй вопросы тулами; скажи @@ -269,12 +272,18 @@ История разряжается по возрасту точки: первые сутки — сырой интервал агента, 1–3 сут — точка в минуту, 3–7 сут — точка в 5 минут, дальше — точка в час; старше года точки удаляются. + + У схлопнутых точек (старше суток) есть `agg` — дайджест интервала: + `n`/`exp` (сколько замеров пришло из ожидаемых — n < exp значит простой), + min/max по cpu/ram/swap/load/сеть и по маунтам диска. Пики и провалы + внутри интервала видны по нему, даже если сама точка их не показывает. """ since = (datetime.now(timezone.utc) - timedelta(hours=hours)).isoformat() db = get_db() cursor = await db.execute( """SELECT ts, cpu, load1, load5, load15, ram_used, ram_total, - swap_used, swap_total, uptime, net_in_mbs, net_out_mbs, disks_json + swap_used, swap_total, uptime, net_in_mbs, net_out_mbs, + disks_json, agg_json FROM metrics WHERE server_id = ? AND ts >= ? ORDER BY ts DESC LIMIT 5000""", (server_id, since), @@ -295,15 +304,56 @@ "net_in_mbs": row["net_in_mbs"], "net_out_mbs": row["net_out_mbs"], "disks": json.loads(row["disks_json"]), + "agg": json.loads(row["agg_json"] or "{}") or None, } for row in reversed(rows) ] if len(points) > max_points: - step = len(points) / max_points - points = [points[int(i * step)] for i in range(max_points)] + points = _thin_points(points, max_points) return _dumps({"server_id": server_id, "hours": hours, "points": points}) +def _thin_points(points: list[dict], max_points: int) -> list[dict]: + """Проредить точки до max_points, сливая дайджесты выкинутых в оставшуюся. + + Иначе прореживание в самом туле съедало бы ровно те крайности, ради + которых дайджест и хранится. + """ + step = len(points) / max_points + out = [] + for i in range(max_points): + lo = int(i * step) + hi = int((i + 1) * step) if i + 1 < max_points else len(points) + point = dict(points[lo]) + digest: dict = {} + for skipped in points[lo:hi]: + _merge_digest(digest, skipped.get("agg")) + if digest: + point["agg"] = digest + out.append(point) + return out + + +def _merge_digest(acc: dict, agg: dict | None) -> dict: + """Слить дайджесты: n/exp складываются, min/max расширяются.""" + if not agg: + return acc + for key, value in agg.items(): + if key == "n" or key == "exp": + acc[key] = acc.get(key, 0) + value + elif key == "disks": + disks = acc.setdefault("disks", {}) + for mount, pair in value.items(): + disks[mount] = _merge_pair(disks.get(mount), pair) + else: + acc[key] = _merge_pair(acc.get(key), value) + return acc + + +def _merge_pair(cur: list | None, pair: list) -> list: + return list(pair) if cur is None else [min(cur[0], pair[0]), max(cur[1], pair[1])] + + # --- Сервисы (health-чеки) ----------------------------------------------------- @@ -349,7 +399,7 @@ @mcp.tool() async def events_recent(limit: int = 30, severity: str = "") -> str: - """Хвост единого журнала событий (90 дней хранения). + """Хвост единого журнала событий (хранение — год). severity — необязательный фильтр: info | warning | critical. Возвращает: id, ts, type, severity, message, data (payload события). diff --git a/panel/backend/app/metrics_retention.py b/panel/backend/app/metrics_retention.py index ecb1a68..73e033d 100644 --- a/panel/backend/app/metrics_retention.py +++ b/panel/backend/app/metrics_retention.py @@ -10,20 +10,39 @@ старше 365 сут удаляем В каждом бакете остаётся реальная (последняя) точка — никаких средних, только -прореживание. Строки берутся парой (id, ts) без JSON. Формат ts — ISO-8601 UTC -фиксированной длины, поэтому бакет считается срезом строки: 'YYYY-MM-DDTHH:MM' -(минута), 'YYYY-MM-DDTHH:M' (пятнадцатиминутный бакет M//5), 'YYYY-MM-DDTHH' -(час) — парсить даты не нужно, лексикографика корректна. +прореживание. НО чтобы прореживание не теряло крайности (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, поэтому прореживание истории ничего в приёме пакетов не ломает. - -Первый прогон на старой БД — самый дорогой: в окне «7 сут – год» может лежать -~1 млн точек на сервер (интервал 30 с) — это десятки МБ памяти на окно и -минуты работы; дальше окна уже разряжены и проходы дешёвые. """ import asyncio +import json from datetime import datetime, timedelta, timezone from app.db import get_db, incremental_vacuum @@ -37,7 +56,9 @@ HORIZON = timedelta(days=365) LOOP_INTERVAL_S = 900 # 15 минут -DELETE_CHUNK = 500 + +# ширина бакета в секундах — из неё считается «сколько замеров ожидалось» (exp) +BUCKET_WIDTH_S = {"minute": 60, "5min": 300, "hour": 3600} def _bucket(ts: str, width: str) -> str: @@ -49,45 +70,151 @@ return ts[:16] -async def _thin_server_tier(db, server_id: int, older: str, younger: str, width: str) -> int: - """Проредить окно (older, younger] одного сервера до одного ряда на бакет. +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 FROM metrics WHERE server_id = ? AND ts < ? AND ts >= ? ORDER BY id", + """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), ) - rows = await cursor.fetchall() - if not rows: - return 0 - last_kept: dict[str, int] = {} # бакет → id последней (максимальный id) точки - for row in rows: - last_kept[_bucket(row["ts"], width)] = row["id"] - doomed = [row["id"] for row in rows if last_kept[_bucket(row["ts"], width)] != row["id"]] + 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 - for start in range(0, len(doomed), DELETE_CHUNK): - chunk = doomed[start : start + DELETE_CHUNK] - marks = ",".join("?" * len(chunk)) - cursor = await db.execute(f"DELETE FROM metrics WHERE id IN ({marks})", chunk) + 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 - return deleted + flushed += 1 + return deleted, flushed async def thin_once() -> dict: - """Один проход: проредить ярусы и отрезать всё старше горизонта.""" + """Один проход: схлопнуть ярусы и отрезать всё старше горизонта.""" db = get_db() now = datetime.now(timezone.utc) - cursor = await db.execute("SELECT id FROM servers") - server_ids = [row["id"] for row in await cursor.fetchall()] + 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 in server_ids: - deleted += await _thin_server_tier(db, server_id, older_iso, younger_iso, width) + 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(),)) @@ -96,7 +223,7 @@ if deleted: # вернуть файлу освободившиеся страницы (нужен auto_vacuum=INCREMENTAL) await incremental_vacuum(2000) - return {"deleted": deleted} + return {"deleted": deleted, "buckets": buckets} async def retention_loop() -> None: @@ -105,7 +232,11 @@ try: result = await thin_once() if result["deleted"]: - print(f"metrics retention: удалено {result['deleted']} точек", flush=True) + 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) diff --git a/panel/backend/app/schema.sql b/panel/backend/app/schema.sql index c144956..31ee7c8 100644 --- a/panel/backend/app/schema.sql +++ b/panel/backend/app/schema.sql @@ -32,6 +32,9 @@ net_in_mbs REAL NOT NULL DEFAULT 0, -- суммарный вход, МБ/с (panel считает по дельте) net_out_mbs REAL NOT NULL DEFAULT 0, -- суммарный выход, МБ/с disks_json TEXT NOT NULL DEFAULT '[]', + -- дайджест схлопнутого бакета (разряжение): n/exp и min/max по метрикам — + -- крайности (спайки, провалы) и простои не теряются при прореживании + agg_json TEXT NOT NULL DEFAULT '{}', -- колонки-наследие: тяжёлые детали пакета теперь живут в server_state -- (нужны только от ПОСЛЕДНЕГО пакета, в истории — мёртвый груз ~89% JSON); -- в схеме оставлены ради старых БД, новые строки пишут дефолты