diff --git a/app.py b/app.py index b23d8f2..5171de7 100644 --- a/app.py +++ b/app.py @@ -484,12 +484,6 @@ def _analyze_emby_sharing(series: list[dict], step_seconds: int) -> dict: } -def _emby_history_step(days: int) -> int: - """Keep query_range below Prometheus' 11,000-points-per-series limit.""" - duration_seconds = days * 86400 - return max(60, math.ceil((duration_seconds / 10_500) / 60) * 60) - - async def _fetch_emby_session_history(days: int, server: str) -> tuple[list[dict], int]: cfg = SERVICES.get("grafana") if not cfg: @@ -503,43 +497,23 @@ async def _fetch_emby_session_history(days: int, server: str) -> tuple[list[dict query = f"max by ({labels}) ({selector})" end = int(time.time()) start = end - days * 86400 - step = _emby_history_step(days) + step = max(60, math.ceil(((end - start) / 30000) / 60) * 60) url = f"{request_data['base_url'].rstrip('/')}/api/datasources/proxy/uid/{datasource_uid}/api/v1/query_range" - series_by_metric: dict[str, dict] = {} - chunk_seconds = 30 * 86400 try: async with httpx.AsyncClient(timeout=90) as client: - chunk_start = start - while chunk_start < end: - chunk_end = min(chunk_start + chunk_seconds, end) - response = await client.get( - url, - params={"query": query, "start": chunk_start, "end": chunk_end, "step": step}, - headers=request_data["headers"], cookies=request_data["cookies"], - ) - response.raise_for_status() - payload = response.json() - if payload.get("status") != "success": - raise HTTPException(502, "Prometheus rejected the Emby session history query") - for item in payload.get("data", {}).get("result", []): - metric = item.get("metric", {}) - key = json.dumps(metric, sort_keys=True, separators=(",", ":")) - merged = series_by_metric.setdefault(key, {"metric": metric, "values": []}) - merged["values"].extend(item.get("values", [])) - chunk_start = chunk_end + response = await client.get( + url, + params={"query": query, "start": start, "end": end, "step": step}, + headers=request_data["headers"], cookies=request_data["cookies"], + ) + response.raise_for_status() + payload = response.json() except (httpx.HTTPError, ValueError) as exc: log.warning("Emby sharing history query failed: %s", type(exc).__name__) raise HTTPException(502, "Emby session history is temporarily unavailable") - - series = [] - for item in series_by_metric.values(): - values_by_timestamp = { - float(value[0]): value for value in item["values"] - if isinstance(value, (list, tuple)) and len(value) >= 2 - } - item["values"] = [values_by_timestamp[ts] for ts in sorted(values_by_timestamp)] - series.append(item) - return series, step + if payload.get("status") != "success": + raise HTTPException(502, "Prometheus rejected the Emby session history query") + return payload.get("data", {}).get("result", []), step @app.get("/emby/account-sharing") diff --git a/tests/test_app.py b/tests/test_app.py index fe4a29b..5c6b2d7 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -216,60 +216,6 @@ def test_emby_account_sharing_endpoint_is_read_only_and_filterable(monkeypatch): assert payload["users"][0]["username"] == "Alice" -def test_emby_history_step_stays_below_prometheus_resolution_limit(): - assert app._emby_history_step(7) == 60 - for days in (7, 30, 90): - step = app._emby_history_step(days) - assert step % 60 == 0 - assert (days * 86400) / step <= 10_500 - - -def test_emby_history_fetch_chunks_90_days_and_merges_equal_series(monkeypatch): - calls = [] - - class FakeResponse: - def __init__(self, params): - self.params = params - - def raise_for_status(self): - return None - - def json(self): - start = self.params["start"] - end = self.params["end"] - return {"status": "success", "data": {"result": [{ - "metric": {"job": "emby-sascha", "username": "Alice", "remoteEndPoint": "8.8.8.8"}, - "values": [[start, "1"], [end, "1"]], - }]}} - - class FakeClient: - def __init__(self, **_kwargs): - pass - - async def __aenter__(self): - return self - - async def __aexit__(self, *_args): - return None - - async def get(self, _url, params, headers, cookies): - calls.append(params) - return FakeResponse(params) - - monkeypatch.setattr(app, "SERVICES", {"grafana": {"url": "http://grafana", "auth": "none"}}) - monkeypatch.setattr(app.httpx, "AsyncClient", FakeClient) - monkeypatch.setattr(app.time, "time", lambda: 10_000_000) - - series, step = asyncio.run(app._fetch_emby_session_history(90, "all")) - - assert len(calls) == 3 - assert all(call["end"] - call["start"] <= 30 * 86400 for call in calls) - assert step == 780 - assert len(series) == 1 - timestamps = [value[0] for value in series[0]["values"]] - assert timestamps == sorted(set(timestamps)) - - def test_capabilities_is_live_machine_readable_safety_map(): with TestClient(app.app) as client: response = client.get(