Newer
Older
hard-panel / panel / backend / app / digest.py
"""Прореживание точек истории с сохранением крайностей.

Общий код для всех мест, где точки приходится выбрасывать: ответ
`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])]