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