Compare commits
No commits in common. "9d5e1e90ae8393ffdf21bfe473d4a5be4e43a753" and "6605940a81c5519b52ce420244a241fd8bb704db" have entirely different histories.
9d5e1e90ae
...
6605940a81
2 changed files with 11 additions and 91 deletions
36
app.py
36
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]:
|
async def _fetch_emby_session_history(days: int, server: str) -> tuple[list[dict], int]:
|
||||||
cfg = SERVICES.get("grafana")
|
cfg = SERVICES.get("grafana")
|
||||||
if not cfg:
|
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})"
|
query = f"max by ({labels}) ({selector})"
|
||||||
end = int(time.time())
|
end = int(time.time())
|
||||||
start = end - days * 86400
|
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"
|
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:
|
try:
|
||||||
async with httpx.AsyncClient(timeout=90) as client:
|
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(
|
response = await client.get(
|
||||||
url,
|
url,
|
||||||
params={"query": query, "start": chunk_start, "end": chunk_end, "step": step},
|
params={"query": query, "start": start, "end": end, "step": step},
|
||||||
headers=request_data["headers"], cookies=request_data["cookies"],
|
headers=request_data["headers"], cookies=request_data["cookies"],
|
||||||
)
|
)
|
||||||
response.raise_for_status()
|
response.raise_for_status()
|
||||||
payload = response.json()
|
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
|
|
||||||
except (httpx.HTTPError, ValueError) as exc:
|
except (httpx.HTTPError, ValueError) as exc:
|
||||||
log.warning("Emby sharing history query failed: %s", type(exc).__name__)
|
log.warning("Emby sharing history query failed: %s", type(exc).__name__)
|
||||||
raise HTTPException(502, "Emby session history is temporarily unavailable")
|
raise HTTPException(502, "Emby session history is temporarily unavailable")
|
||||||
|
if payload.get("status") != "success":
|
||||||
series = []
|
raise HTTPException(502, "Prometheus rejected the Emby session history query")
|
||||||
for item in series_by_metric.values():
|
return payload.get("data", {}).get("result", []), step
|
||||||
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
|
|
||||||
|
|
||||||
|
|
||||||
@app.get("/emby/account-sharing")
|
@app.get("/emby/account-sharing")
|
||||||
|
|
|
||||||
|
|
@ -216,60 +216,6 @@ def test_emby_account_sharing_endpoint_is_read_only_and_filterable(monkeypatch):
|
||||||
assert payload["users"][0]["username"] == "Alice"
|
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():
|
def test_capabilities_is_live_machine_readable_safety_map():
|
||||||
with TestClient(app.app) as client:
|
with TestClient(app.app) as client:
|
||||||
response = client.get(
|
response = client.get(
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue