"""Разряжение истории метрик по возрасту точки + годовой горизонт.
Политика (границы — по возрасту, а не по календарным суткам: иначе на стыке
суток скачок плотности):
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)