diff --git a/README.md b/README.md index f0b6f9a..ce94ec7 100644 --- a/README.md +++ b/README.md @@ -273,6 +273,14 @@ (`GET /api/v1/servers/{id}/metrics`), MCP-тул `server_history` — там же, при своём прореживании он сливает дайджесты, а не выбрасывает их. +Ответ истории тоже прореживается, и тоже без потери краёв: `limit` — это +сколько точек вернуть максимум, а если в окне их больше, они прореживаются +**равномерно по всему периоду** (`points[0]` — начало запрошенного интервала, +`points[-1]` — сейчас), дайджесты выброшенных точек сливаются в оставшуюся. +Обрезка по свежему краю была бы хуже вдвойне: запрос за неделю отдавал бы +только последние часы — сырые сутки вытесняли всю старую историю вместе с +дайджестами, и полосы на длинных периодах просто не доезжали до графика. + Отдельно от истории живёт таблица `server_state` — «состояние последнего пакета» на сервер: сырые счётчики сети (нужны, чтобы посчитать МБ/с от предыдущего пакета), список docker-контейнеров (нужен для diff'а ивентов), diff --git a/panel/backend/app/api/servers.py b/panel/backend/app/api/servers.py index 95ed8ea..15f76da 100644 --- a/panel/backend/app/api/servers.py +++ b/panel/backend/app/api/servers.py @@ -11,11 +11,18 @@ from app.config import get_settings from app.db import get_db, j +from app.digest import thin_points from app.models import ServerCreate, ServerUpdate from app.security import generate_server_key, hash_key, require_admin router = APIRouter(prefix="/api/v1", dependencies=[Depends(require_admin)]) +# сколько строк истории читаем из БД под запрос. Окно 7 сут после разряжения +# держит ~10 тыс. строк при интервале 15 с (сырые сутки + ярусы 1/мин и +# 1/5мин), так что запас есть; агент с невероятно частым интервалом упрётся в +# потолок и потеряет самый старый край окна — а не свежий, как было раньше. +FETCH_CAP = 30_000 + def _now_iso() -> str: return datetime.now(timezone.utc).isoformat() @@ -179,9 +186,14 @@ ) -> list[dict]: """История метрик (графики): время, cpu, ram, сеть. По возрастанию ts. + `limit` — сколько точек вернуть максимум. Если в окне их больше, они + равномерно прореживаются по всему окну (а не обрезаются по свежему краю: + иначе запрос за неделю отдавал бы только последние часы — сырые сутки + вытесняют всю старую историю); дайджесты выброшенных точек сливаются в + оставшуюся, поэтому пики и простои не теряются. + У схлопнутых бакетов (старше суток) приходит `agg` — дайджест интервала - (min/max по метрикам, n/exp замеров): пики и простои не теряются при - разряжении. + (min/max по метрикам, n/exp замеров). """ db = get_db() conditions = ["server_id = ?"] @@ -192,7 +204,7 @@ if until is not None: conditions.append("ts <= ?") params.append(until.astimezone(timezone.utc).isoformat()) - params.append(limit) + params.append(FETCH_CAP) 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, @@ -202,7 +214,7 @@ params, ) rows = await cursor.fetchall() - return [ + points = [ { "ts": row["ts"], "cpu": row["cpu"], @@ -216,4 +228,7 @@ "agg": json.loads(row["agg_json"] or "{}") or None, } for row in reversed(rows) - ] \ No newline at end of file + ] + if len(points) > limit: + points = thin_points(points, limit) + return points \ No newline at end of file diff --git a/panel/backend/app/digest.py b/panel/backend/app/digest.py new file mode 100644 index 0000000..4ac9e22 --- /dev/null +++ b/panel/backend/app/digest.py @@ -0,0 +1,57 @@ +"""Прореживание точек истории с сохранением крайностей. + +Общий код для всех мест, где точки приходится выбрасывать: ответ +`GET /servers/{id}/metrics` (страница сервера) и MCP-тул `server_history`. +Точка, оставшаяся от группы, несёт дайджест `agg` всей группы — min/max по +метрикам и n/exp замеров (формат и смысл — в docstring `metrics_retention`); +без этого прореживание съедало бы ровно те всплески и простои, ради которых +дайджест и хранится. + +Прореживание в БД (по возрасту точки) живёт отдельно — в +`metrics_retention.thin_once`, потому что считает бакеты по календарю, а не по +числу точек; здесь же группы нарезаются равномерно по окну ответа. +""" + + +def thin_points(points: list[dict], max_points: int) -> list[dict]: + """Проредить точки (по возрастанию ts) до max_points, сливая их дайджесты. + + В каждой группе остаётся последняя точка — самая свежая, поэтому свежий + край окна не теряется; в её `agg` собираются крайности всей группы. + """ + 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) + if lo >= hi: + break # вырожденные слоты возможны только при max_points > len(points) + group = points[lo:hi] + point = dict(group[-1]) + digest: dict = {} + for member in group: + merge_digest(digest, member.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 in ("n", "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])] diff --git a/panel/backend/app/mcp.py b/panel/backend/app/mcp.py index 371616f..d3b4ba5 100644 --- a/panel/backend/app/mcp.py +++ b/panel/backend/app/mcp.py @@ -36,6 +36,7 @@ ) from app.api.shares import list_shares as _list_shares from app.db import get_db +from app.digest import thin_points from app.models import ServerCreate as _ServerCreate AGENT_INSTRUCTIONS = """\ @@ -125,7 +126,9 @@ - `server_history(hours)`: разумно 1 (детали), 6, 24 (сутки), 168 (неделя) или 720 (месяц). История разряжается по возрасту точки: первые сутки — как прислал агент, 1–3 сут — точка в минуту, 3–7 сут — точка в 5 минут, дальше — - точка в час; старше года точки удаляются. Вопросы вида «был ли всплеск или + точка в час; старше года точки удаляются. Окно покрывается целиком (точки + прореживаются по нему равномерно, первая — начало периода, последняя — + сейчас). Вопросы вида «был ли всплеск или простой» решай по `agg` точки (min/max за интервал и n/exp замеров), а не по значению самой точки: точка — это последний замер интервала, а пик мог быть внутри. @@ -273,6 +276,10 @@ агента, 1–3 сут — точка в минуту, 3–7 сут — точка в 5 минут, дальше — точка в час; старше года точки удаляются. + Если точек в окне больше `max_points`, они прореживаются равномерно по + всему окну (а не обрезаются по свежему краю) — окно всегда покрыто + целиком: `points[0]` — начало запрошенного периода, `points[-1]` — сейчас. + У схлопнутых точек (старше суток) есть `agg` — дайджест интервала: `n`/`exp` (сколько замеров пришло из ожидаемых — n < exp значит простой), min/max по cpu/ram/swap/load/сеть и по маунтам диска. Пики и провалы @@ -285,7 +292,7 @@ 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""", + ORDER BY ts DESC LIMIT 30000""", (server_id, since), ) rows = await cursor.fetchall() @@ -309,51 +316,10 @@ for row in reversed(rows) ] if len(points) > max_points: - points = _thin_points(points, 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-чеки) -----------------------------------------------------