Add Emby account-sharing analysis and read-only UI access (app.py)
This commit is contained in:
parent
63978bff18
commit
d23fdc7896
1 changed files with 271 additions and 9 deletions
280
app.py
280
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")
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue