"""Прореживание точек истории с сохранением крайностей.
Общий код для всех мест, где точки приходится выбрасывать: ответ
`GET /servers/{id}/metrics` (страница сервера) и MCP-тул `server_history`.
Точка, оставшаяся от группы, несёт дайджест `agg` всей группы — min/max по
метрикам и n/exp замеров (формат и смысл — в docstring `metrics_retention`);
без этого прореживание съедало бы ровно те всплески и простои, ради которых
дайджест и хранится.
Прореживание в БД (по возрасту точки) живёт отдельно — в
`metrics_retention.thin_once`, потому что считает бакеты по календарю; здесь же
окно ответа режется на слоты равной длины, чтобы ось графика осталась
пропорциональной времени.
"""
from datetime import datetime, timezone
def _parse_ts(ts: str) -> datetime:
"""ts строки всегда с офсетом (`ingest._iso`), но наивный не должен ронять ответ."""
parsed = datetime.fromisoformat(ts)
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
def thin_points(points: list[dict], max_points: int) -> list[dict]:
"""Проредить точки (по возрастанию ts) до max_points, сливая их дайджесты.
Окно от первой точки до последней делится на `max_points` слотов равной
длительности, и в каждом остаётся последняя попавшая в него точка — самая
свежая, так что свежий край окна не теряется. В её `agg` собираются
крайности всех точек слота.
Режем по времени, а не по числу точек: при нарезке по счёту редкие ярусы
хранения (часовая точка старше недели) раздуваются до ширины сырых суток —
неделя ужималась в первые проценты ширины графика. Со слотами равной длины
ось почти пропорциональна времени; неточность остаётся только там, где
данные в окне неоднородны (самая старая точка может сдвинуться вправо на
длину слота).
Точку с ts не нашего формата пропускаем: на ось её всё равно не поставить,
а ронять из-за одной строки весь ответ незачем (ts пишет только ingest, и
только через `_iso`, так что такой строке взяться неоткуда).
"""
if max_points < 1 or len(points) <= max_points:
return points
anchor: datetime | None = None
stamped: list[tuple[float, dict]] = []
for point in points:
try:
ts = _parse_ts(point["ts"])
except ValueError:
continue
if anchor is None:
anchor = ts
stamped.append(((ts - anchor).total_seconds(), point))
if not stamped:
return []
span = stamped[-1][0] # окно меряем от первой разобранной точки (ts по возрастанию)
slot = span / max_points if span > 0 else 0.0
slots: dict[int, list[dict]] = {}
for offset, point in stamped:
index = min(int(offset / slot), max_points - 1) if slot > 0 else 0
slots.setdefault(index, []).append(point)
out = []
for index in sorted(slots):
group = slots[index]
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])]