From d23fdc7896a258ffbd5b89777feb5ef22c1a34c6 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 21:28:56 +0200 Subject: [PATCH] Add Emby account-sharing analysis and read-only UI access (app.py) --- app.py | 280 +++++++++++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 271 insertions(+), 9 deletions(-) diff --git a/app.py b/app.py index 2b6da2e..bac1a05 100644 --- a/app.py +++ b/app.py @@ -1,7 +1,7 @@ """Homelab Butler v2.1 – Unified API proxy for Pfannkuchen homelab. Reads service config from butler.yaml, credentials from Vaultwarden cache with flat-file fallback.""" -import os, json, asyncio, logging, time, base64, re, subprocess, ipaddress, secrets, sqlite3 +import os, json, asyncio, logging, time, base64, re, subprocess, ipaddress, secrets, sqlite3, math from datetime import datetime, timezone import httpx, yaml from typing import Literal @@ -246,6 +246,24 @@ def _ui_session(request: Request) -> dict | None: return None return session + +def _issue_ui_session(read_only: bool) -> JSONResponse: + session_id = secrets.token_urlsafe(32) + csrf = secrets.token_urlsafe(24) + _ui_sessions[session_id] = { + "csrf": csrf, + "expires": time.time() + UI_SESSION_TTL, + "read_only": read_only, + } + response = JSONResponse({ + "authenticated": True, + "read_only": read_only, + "expires_in": UI_SESSION_TTL, + }) + response.set_cookie("butler_session", session_id, max_age=UI_SESSION_TTL, httponly=True, samesite="strict", path="/") + response.set_cookie("butler_csrf", csrf, max_age=UI_SESSION_TTL, httponly=False, samesite="strict", path="/") + return response + def _verify(request: Request): if not BUTLER_TOKEN: return @@ -256,10 +274,253 @@ def _verify(request: Request): if not session: raise HTTPException(401, "Invalid token") if request.method not in {"GET", "HEAD", "OPTIONS"}: + if session.get("read_only", False): + raise HTTPException(403, "Anonymous UI session is read-only") csrf = request.headers.get("x-csrf-token", "") if not csrf or not secrets.compare_digest(csrf, session["csrf"]): raise HTTPException(403, "Invalid CSRF token") +def _emby_network_identity(endpoint: str) -> dict: + """Normalize an Emby endpoint without treating IPv6 privacy addresses as new households.""" + raw = str(endpoint or "").strip() + if raw.startswith("[") and "]" in raw: + raw = raw[1:raw.index("]")] + try: + address = ipaddress.ip_address(raw) + except ValueError: + if raw.count(":") == 1: + raw = raw.rsplit(":", 1)[0] + try: + address = ipaddress.ip_address(raw) + except ValueError as exc: + raise ValueError("Invalid Emby remote endpoint") from exc + if address.version == 4: + network = ipaddress.ip_network(f"{address}/32", strict=False) + parent = network + else: + network = ipaddress.ip_network(f"{address}/64", strict=False) + parent = ipaddress.ip_network(f"{address}/48", strict=False) + return { + "ip": str(address), + "version": address.version, + "network": str(network), + "parent": str(parent), + "identity": str(parent if address.version == 6 else network), + "public": address.is_global, + } + + +def _emby_location(metric: dict) -> dict: + def coordinate(name): + try: + return float(metric.get(name, 0)) + except (TypeError, ValueError): + return 0.0 + return { + "city": metric.get("city", ""), + "region": metric.get("region", ""), + "country": metric.get("countryCode", ""), + "latitude": coordinate("latitude"), + "longitude": coordinate("longitude"), + } + + +def _analyze_emby_sharing(series: list[dict], step_seconds: int) -> dict: + observations: dict[str, dict[int, dict[str, dict]]] = {} + tracks: dict[str, dict[str, dict]] = {} + identities_by_user: dict[str, set[str]] = {} + servers_by_user: dict[str, set[str]] = {} + for item in series: + metric = item.get("metric", {}) + username = str(metric.get("username", "")).strip() + if not username: + continue + try: + network = _emby_network_identity(metric.get("remoteEndPoint", "")) + except ValueError: + continue + if not network["public"]: + continue + evidence = { + **network, + "server": metric.get("job", ""), + "location": _emby_location(metric), + } + identities_by_user.setdefault(username, set()).add(network["identity"]) + servers_by_user.setdefault(username, set()).add(str(metric.get("job", ""))) + track = tracks.setdefault(username, {}).setdefault(network["identity"], {"evidence": evidence, "timestamps": []}) + for value in item.get("values", []): + if not isinstance(value, list) or len(value) < 2 or str(value[1]).lower() in {"0", "nan"}: + continue + timestamp = int(float(value[0])) + track["timestamps"].append(timestamp) + observations.setdefault(username, {}).setdefault(timestamp, {}).setdefault(network["identity"], evidence) + + raw_events = [] + for username, timeline in observations.items(): + buckets = [] + for timestamp in sorted(timeline): + evidence = timeline[timestamp] + if len(evidence) >= 2: + buckets.append((timestamp, tuple(sorted(evidence)), evidence)) + current = None + for timestamp, identity_key, evidence in buckets: + if current and current["identity_key"] == identity_key and timestamp - current["end_ts"] <= step_seconds * 2: + current["end_ts"] = timestamp + current["samples"] += 1 + continue + if current and current["samples"] >= 2 and current["end_ts"] - current["start_ts"] + step_seconds >= 360: + raw_events.append(current) + current = { + "username": username, "identity_key": identity_key, "start_ts": timestamp, + "end_ts": timestamp, "samples": 1, "evidence": list(evidence.values()), + } + if current and current["samples"] >= 2 and current["end_ts"] - current["start_ts"] + step_seconds >= 360: + raw_events.append(current) + + events = [{ + "type": "concurrent_networks", + "severity": "high", + "username": item["username"], + "start": datetime.fromtimestamp(item["start_ts"], timezone.utc).isoformat(), + "end": datetime.fromtimestamp(item["end_ts"], timezone.utc).isoformat(), + "duration_seconds": item["end_ts"] - item["start_ts"] + step_seconds, + "samples": item["samples"], + "evidence": item["evidence"], + "reason": "Zeitgleiche Nutzung desselben Emby-Benutzers aus unterschiedlichen öffentlichen Netzen", + } for item in raw_events] + + def distance_km(first: dict, second: dict) -> float: + lat1, lon1 = first["latitude"], first["longitude"] + lat2, lon2 = second["latitude"], second["longitude"] + if not all((-90 <= lat <= 90 and -180 <= lon <= 180) for lat, lon in ((lat1, lon1), (lat2, lon2))): + return 0.0 + phi1, phi2 = math.radians(lat1), math.radians(lat2) + dphi, dlambda = math.radians(lat2 - lat1), math.radians(lon2 - lon1) + value = math.sin(dphi / 2) ** 2 + math.cos(phi1) * math.cos(phi2) * math.sin(dlambda / 2) ** 2 + return 6371.0 * 2 * math.atan2(math.sqrt(value), math.sqrt(max(0.0, 1 - value))) + + travel_events = [] + for username, user_tracks in tracks.items(): + intervals = [] + for identity, track in user_tracks.items(): + current = None + for timestamp in sorted(set(track["timestamps"])): + if current and timestamp - current["end"] <= step_seconds * 2: + current["end"] = timestamp + current["samples"] += 1 + else: + if current and current["samples"] >= 2: + intervals.append(current) + current = {"identity": identity, "start": timestamp, "end": timestamp, "samples": 1, "evidence": track["evidence"]} + if current and current["samples"] >= 2: + intervals.append(current) + intervals.sort(key=lambda item: item["start"]) + for previous, current in zip(intervals, intervals[1:]): + if previous["identity"] == current["identity"] or current["start"] <= previous["end"]: + continue + distance = distance_km(previous["evidence"]["location"], current["evidence"]["location"]) + gap_hours = max((current["start"] - previous["end"]) / 3600, 1 / 60) + speed = distance / gap_hours + if distance < 300 or speed <= 1000: + continue + travel_events.append({ + "type": "impossible_travel", "severity": "medium", "username": username, + "start": datetime.fromtimestamp(previous["end"], timezone.utc).isoformat(), + "end": datetime.fromtimestamp(current["start"], timezone.utc).isoformat(), + "duration_seconds": current["start"] - previous["end"], + "samples": previous["samples"] + current["samples"], + "distance_km": round(distance, 1), "required_speed_kmh": round(speed, 1), + "evidence": [previous["evidence"], current["evidence"]], + "reason": "Geografischer Wechsel zwischen öffentlichen Netzen wäre in der verfügbaren Zeit nicht plausibel", + }) + events.extend(travel_events) + events.sort(key=lambda item: item["start"], reverse=True) + flagged = {item["username"] for item in events} + users = [{ + "username": username, + "risk": "high" if any(item["username"] == username and item["type"] == "concurrent_networks" for item in events) else ("medium" if username in flagged else "none"), + "events": sum(item["username"] == username for item in events), + "network_identities": len(identities_by_user.get(username, set())), + "servers": sorted(servers_by_user.get(username, set())), + } for username in sorted(observations)] + return { + "summary": { + "users_analyzed": len(observations), + "flagged_users": len(flagged), + "concurrent_events": len(raw_events), + "impossible_travel_events": len(travel_events), + }, + "users": users, + "events": events, + } + + +async def _fetch_emby_session_history(days: int, server: str) -> tuple[list[dict], int]: + cfg = SERVICES.get("grafana") + if not cfg: + raise HTTPException(503, "Grafana service is not configured") + request_data = _service_auth(cfg) + datasource_uid = os.environ.get("EMBY_PROMETHEUS_UID", "bdpu4276997nkc") + labels = "job,username,remoteEndPoint,city,region,countryCode,latitude,longitude" + selector = 'emby_sessions{username!=""}' + if server != "all": + selector = f'emby_sessions{{username!="",job="{server}"}}' + query = f"max by ({labels}) ({selector})" + end = int(time.time()) + start = end - days * 86400 + 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" + try: + async with httpx.AsyncClient(timeout=90) as client: + 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") + 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") +async def emby_account_sharing( + days: int = Query(30, ge=1, le=90), + username: str | None = Query(None, min_length=1, max_length=100), + server: Literal["all", "emby-sascha", "emby-chris"] = "all", + _=Depends(_verify), +): + """Conservative read-only analysis of concurrent networks and geographically impossible changes.""" + series, step = await _fetch_emby_session_history(days, server) + if username: + wanted = username.casefold() + series = [item for item in series if str(item.get("metric", {}).get("username", "")).casefold() == wanted] + result = _analyze_emby_sharing(series, step) + return { + "generated": datetime.now(timezone.utc).isoformat(), + "period": {"days": days, "server": server, "step_seconds": step, "series": len(series)}, + "policy": { + "mode": "conservative", + "ipv4_detection_identity": "/32", + "ipv6_display_network": "/64", + "ipv6_detection_identity": "/48", + "minimum_samples": 2, + "minimum_concurrent_seconds": 360, + "prometheus_staleness_guard": True, + "impossible_travel_minimum_km": 300, + "impossible_travel_speed_kmh": 1000, + "private_networks_excluded": True, + "automatic_enforcement": False, + }, + **result, + } + + def _get_key(cfg): vault_key = cfg.get("vault_key") if vault_key and vault_key in _vault_cache: @@ -328,18 +589,19 @@ async def ui(): async def ui_login(payload: UiLoginRequest): if not BUTLER_TOKEN or not secrets.compare_digest(payload.token, BUTLER_TOKEN): raise HTTPException(401, "Invalid token") - session_id = secrets.token_urlsafe(32) - csrf = secrets.token_urlsafe(24) - _ui_sessions[session_id] = {"csrf": csrf, "expires": time.time() + UI_SESSION_TTL} - response = JSONResponse({"authenticated": True, "expires_in": UI_SESSION_TTL}) - response.set_cookie("butler_session", session_id, max_age=UI_SESSION_TTL, httponly=True, samesite="strict", path="/") - response.set_cookie("butler_csrf", csrf, max_age=UI_SESSION_TTL, httponly=False, samesite="strict", path="/") - return response + return _issue_ui_session(read_only=False) @app.get("/ui/session") async def ui_session(request: Request): - return {"authenticated": _ui_session(request) is not None} + session = _ui_session(request) + if session: + return { + "authenticated": True, + "read_only": session.get("read_only", False), + "expires_in": max(0, int(session["expires"] - time.time())), + } + return _issue_ui_session(read_only=True) @app.post("/ui/logout")