diff --git a/panel/backend/app/api/ingest.py b/panel/backend/app/api/ingest.py index aa0390f..a4b72de 100644 --- a/panel/backend/app/api/ingest.py +++ b/panel/backend/app/api/ingest.py @@ -93,12 +93,21 @@ now = _now() packet_ts = _iso(payload.ts) if payload.ts else _iso(now) - # предыдущая точка — нужны сырые счётчики для дельты сети и diff docker + # предыдущий пакет — сырые счётчики для дельты сети и diff docker; тяжёлые + # детали живут в server_state (в историю метрик они не пишутся) cursor = await db.execute( - "SELECT ts, net_json, docker_json FROM metrics WHERE server_id = ? ORDER BY id DESC LIMIT 1", + "SELECT ts, net_json, docker_json FROM server_state WHERE server_id = ?", (server_id,), ) prev = await cursor.fetchone() + if prev is None: + # апгрейд на живом деплое: состояние ещё не заведено — берём последнюю + # старую строку метрик (до этого релиза в ней лежали те же поля) + cursor = await db.execute( + "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] = [] @@ -133,13 +142,12 @@ ), ) - # точка метрик; net_json — сырые счётчики, из них будет считаться дельта + # точка метрик: только числа и диски (тяжёлые детали пакета — в server_state) await db.execute( """INSERT INTO metrics (server_id, ts, cpu, load1, load5, load15, ram_used, ram_total, swap_used, swap_total, uptime, - net_in_mbs, net_out_mbs, disks_json, net_json, processes_json, - docker_json, extra_json) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", + net_in_mbs, net_out_mbs, disks_json) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", ( server_id, packet_ts, @@ -153,9 +161,25 @@ net_in_mbs, net_out_mbs, j([d.model_dump() for d in payload.disks]), + ), + ) + + # состояние сервера: последний пакет «как есть» — дельта сети, diff docker, + # вкладки «Процессы»/Docker на странице сервера + await db.execute( + """INSERT INTO server_state (server_id, ts, net_json, docker_json, + processes_json, extra_json) VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT(server_id) DO UPDATE SET + ts = excluded.ts, net_json = excluded.net_json, + docker_json = excluded.docker_json, + processes_json = excluded.processes_json, + extra_json = excluded.extra_json""", + ( + server_id, + packet_ts, j([n.model_dump() for n in payload.net]), - j([p.model_dump() for p in payload.processes]), j([c.model_dump() for c in payload.docker]), + j([p.model_dump() for p in payload.processes]), j(payload.extra), ), ) diff --git a/panel/backend/app/api/servers.py b/panel/backend/app/api/servers.py index 32463aa..66811ce 100644 --- a/panel/backend/app/api/servers.py +++ b/panel/backend/app/api/servers.py @@ -114,7 +114,10 @@ raise HTTPException(status_code=404, detail="server not found") item = _server_dict(server) cursor = await db.execute( - "SELECT * FROM metrics WHERE server_id = ? ORDER BY id DESC LIMIT 1", + """SELECT m.*, st.processes_json AS st_processes, st.docker_json AS st_docker, + st.extra_json AS st_extra + FROM metrics m LEFT JOIN server_state st ON st.server_id = m.server_id + WHERE m.server_id = ? ORDER BY m.id DESC LIMIT 1""", (server_id,), ) metrics = await cursor.fetchone() @@ -129,9 +132,10 @@ "net_in_mbs": metrics["net_in_mbs"], "net_out_mbs": metrics["net_out_mbs"], "disks": json.loads(metrics["disks_json"]), - "processes": json.loads(metrics["processes_json"]), - "docker": json.loads(metrics["docker_json"]), - "extra": json.loads(metrics["extra_json"]), + # тяжёлые детали — из состояния последнего пакета (в истории их нет) + "processes": json.loads(metrics["st_processes"] or "[]"), + "docker": json.loads(metrics["st_docker"] or "[]"), + "extra": json.loads(metrics["st_extra"] or "{}"), } else: item["metrics"] = None diff --git a/panel/backend/app/db.py b/panel/backend/app/db.py index d116710..d09a5ac 100644 --- a/panel/backend/app/db.py +++ b/panel/backend/app/db.py @@ -6,15 +6,25 @@ from app.config import get_settings _conn: aiosqlite.Connection | None = None +# У старой БД auto_vacuum=0: чтобы он заработал, нужен разовый VACUUM (см. vacuum_if_needed) +_vacuum_pending = False async def init_db() -> None: """Открыть БД, включить WAL, применить schema.sql.""" - global _conn + global _conn, _vacuum_pending settings = get_settings() settings.database_path.parent.mkdir(parents=True, exist_ok=True) _conn = await aiosqlite.connect(settings.database_path) _conn.row_factory = aiosqlite.Row + # auto_vacuum=INCREMENTAL: удалённые страницы возвращаются файлу (иначе БД + # навсегда остаётся в пиковом размере). На пустой БД действует сразу; на + # существующей — только после VACUUM, поэтому запоминаем флаг. + cursor = await _conn.execute("PRAGMA auto_vacuum") + (auto_vacuum,) = await cursor.fetchone() + _vacuum_pending = auto_vacuum == 0 + if _vacuum_pending: + await _conn.execute("PRAGMA auto_vacuum=INCREMENTAL") await _conn.execute("PRAGMA journal_mode=WAL") await _conn.execute("PRAGMA foreign_keys=ON") schema = (Path(__file__).with_name("schema.sql")).read_text(encoding="utf-8") @@ -38,6 +48,23 @@ raise +async def vacuum_if_needed() -> bool: + """Разовый VACUUM, если у БД был auto_vacuum=0 (см. init_db). Возвращает «делали ли».""" + global _vacuum_pending + if not _vacuum_pending or _conn is None: + return False + _vacuum_pending = False + await _conn.execute("VACUUM") + await _conn.commit() + return True + + +async def incremental_vacuum(pages: int = 1000) -> None: + """Вернуть файлу до pages свободных страниц (работает при auto_vacuum=INCREMENTAL).""" + if _conn is not None: + await _conn.execute(f"PRAGMA incremental_vacuum({int(pages)})") + + async def close_db() -> None: global _conn if _conn is not None: diff --git a/panel/backend/app/main.py b/panel/backend/app/main.py index c46aacd..7e8c490 100644 --- a/panel/backend/app/main.py +++ b/panel/backend/app/main.py @@ -11,9 +11,10 @@ from app.security import require_admin, require_mcp from app.api import auth_routes, events, ingest, mcp_tokens, servers, services, shares -from app.db import close_db, init_db +from app.db import close_db, init_db, vacuum_if_needed from app.events import offline_loop from app.mcp import mcp, mcp_endpoint +from app.metrics_retention import retention_loop from app.services_probe import loop as services_probe_loop from app.shares_probe import loop as shares_probe_loop @@ -25,6 +26,10 @@ @asynccontextmanager async def lifespan(_: FastAPI): await init_db() + # на старых БД один раз перестраиваем файл, включая auto_vacuum (иначе + # удалённые страницы не вернутся файлу) — до старта фоновых петель + if await vacuum_if_needed(): + print("database compacted once (auto_vacuum enabled)", flush=True) # фоновый пробер сетевых хранилищ (share_samples) probe_task = asyncio.create_task(shares_probe_loop()) # фоновый пробер health-эндпоинтов сервисов (service_samples) @@ -33,6 +38,8 @@ gc_task = asyncio.create_task(auth_gc_loop()) # watchdog offline (нет пакетов дольше interval * multiplier) + ретеншн ивентов watchdog_task = asyncio.create_task(offline_loop()) + # разряжение истории метрик по возрасту точки + годовой горизонт + retention_task = asyncio.create_task(retention_loop()) # lifespan смонтированных sub-app не вызывается — MCP session manager # стартуем здесь, иначе /mcp отвечает 500 ("Task group is not initialized") async with mcp.session_manager.run(): @@ -41,6 +48,7 @@ health_task.cancel() gc_task.cancel() watchdog_task.cancel() + retention_task.cancel() await close_db() diff --git a/panel/backend/app/mcp.py b/panel/backend/app/mcp.py index dd4f366..688c026 100644 --- a/panel/backend/app/mcp.py +++ b/panel/backend/app/mcp.py @@ -122,8 +122,10 @@ ## Правила - ID серверов и сервисов бери из обзорных тулов — угадывание запрещено. -- `server_history(hours)`: разумно 1 (детали), 6, 24 (сутки) или 168 (неделя; - глубже недели данных нет — история метрик хранится 7 дней). +- `server_history(hours)`: разумно 1 (детали), 6, 24 (сутки), 168 (неделя) или + 720 (месяц). История разряжается по возрасту точки: первые сутки — как + прислал агент, 1–3 сут — точка в минуту, 3–7 сут — точка в 5 минут, дальше — + точка в час; старше года точки удаляются. - Числа передавай как есть из тулов (проценты, МБ/с, ms) — не пересчитывай единицы самостоятельно. - Если сервер offline — не паникуй и не повторяй вопросы тулами; скажи @@ -264,7 +266,9 @@ """История метрик сервера за `hours` часов, прореженная до `max_points`. Возвращает точки: ts, cpu, ram%, swap, load, сеть MB/s, диски. - История метрик хранится 7 дней — глубже точки удаляются. + История разряжается по возрасту точки: первые сутки — сырой интервал + агента, 1–3 сут — точка в минуту, 3–7 сут — точка в 5 минут, дальше — + точка в час; старше года точки удаляются. """ since = (datetime.now(timezone.utc) - timedelta(hours=hours)).isoformat() db = get_db() diff --git a/panel/backend/app/metrics_retention.py b/panel/backend/app/metrics_retention.py new file mode 100644 index 0000000..ecb1a68 --- /dev/null +++ b/panel/backend/app/metrics_retention.py @@ -0,0 +1,111 @@ +"""Разряжение истории метрик по возрасту точки + годовой горизонт. + +Политика (границы — по возрасту, а не по календарным суткам: иначе на стыке +суток скачок плотности): + + 0 – 24 ч всё (сырой интервал агента) + 24 ч – 3 сут 1 точка/мин + 3 – 7 сут 1 точка/5 мин + 7 сут – 365 сут 1 точка/час + старше 365 сут удаляем + +В каждом бакете остаётся реальная (последняя) точка — никаких средних, только +прореживание. Строки берутся парой (id, ts) без JSON. Формат ts — ISO-8601 UTC +фиксированной длины, поэтому бакет считается срезом строки: 'YYYY-MM-DDTHH:MM' +(минута), 'YYYY-MM-DDTHH:M' (пятнадцатиминутный бакет M//5), 'YYYY-MM-DDTHH' +(час) — парсить даты не нужно, лексикографика корректна. + +Замечание про соседей ingest'а: тяжёлые JSON последнего пакета живут в +server_state, поэтому прореживание истории ничего в приёме пакетов не ломает. + +Первый прогон на старой БД — самый дорогой: в окне «7 сут – год» может лежать +~1 млн точек на сервер (интервал 30 с) — это десятки МБ памяти на окно и +минуты работы; дальше окна уже разряжены и проходы дешёвые. +""" + +import asyncio +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 минут +DELETE_CHUNK = 500 + + +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] + + +async def _thin_server_tier(db, server_id: int, older: str, younger: str, width: str) -> int: + """Проредить окно (older, younger] одного сервера до одного ряда на бакет. + + Возвращает число удалённых строк. + """ + cursor = await db.execute( + "SELECT id, ts FROM metrics WHERE server_id = ? AND ts < ? AND ts >= ? ORDER BY 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"]] + 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) + deleted += cursor.rowcount or 0 + return deleted + + +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()] + + deleted = 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) + + # горизонт: всё, что старше года, не нужно вовсе + 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} + + +async def retention_loop() -> None: + """Проход сразу при старте (деплой начинает худеть немедленно), далее по таймеру.""" + while True: + try: + result = await thin_once() + if result["deleted"]: + print(f"metrics retention: удалено {result['deleted']} точек", 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 7d12e96..c144956 100644 --- a/panel/backend/app/schema.sql +++ b/panel/backend/app/schema.sql @@ -32,7 +32,10 @@ net_in_mbs REAL NOT NULL DEFAULT 0, -- суммарный вход, МБ/с (panel считает по дельте) net_out_mbs REAL NOT NULL DEFAULT 0, -- суммарный выход, МБ/с disks_json TEXT NOT NULL DEFAULT '[]', - net_json TEXT NOT NULL DEFAULT '[]', -- сырые счётчики интерфейсов — для следующей дельты + -- колонки-наследие: тяжёлые детали пакета теперь живут в server_state + -- (нужны только от ПОСЛЕДНЕГО пакета, в истории — мёртвый груз ~89% JSON); + -- в схеме оставлены ради старых БД, новые строки пишут дефолты + net_json TEXT NOT NULL DEFAULT '[]', processes_json TEXT NOT NULL DEFAULT '[]', docker_json TEXT NOT NULL DEFAULT '[]', extra_json TEXT NOT NULL DEFAULT '{}' -- temps, smart и прочее, что panel пока не разбирает @@ -40,6 +43,17 @@ CREATE INDEX IF NOT EXISTS idx_metrics_server_ts ON metrics(server_id, ts); +-- Последний пакет сервера в «сыром» виде: только он нужен для дельты сети, +-- diff контейнеров и вкладок «Процессы»/Docker. Одна строка на сервер. +CREATE TABLE IF NOT EXISTS server_state ( + server_id INTEGER PRIMARY KEY REFERENCES servers(id) ON DELETE CASCADE, + ts TEXT NOT NULL, -- ts последнего пакета (для дельты) + net_json TEXT NOT NULL DEFAULT '[]', -- сырые счётчики интерфейсов + docker_json TEXT NOT NULL DEFAULT '[]', -- список контейнеров (для diff) + processes_json TEXT NOT NULL DEFAULT '[]', + extra_json TEXT NOT NULL DEFAULT '{}' +); + CREATE TABLE IF NOT EXISTS events ( id INTEGER PRIMARY KEY AUTOINCREMENT, server_id INTEGER REFERENCES servers(id) ON DELETE CASCADE,