diff --git a/README.md b/README.md
index d0a464f..aa96e1d 100644
--- a/README.md
+++ b/README.md
@@ -45,6 +45,7 @@ Known secret response fields such as Dockhand's `hawserToken` and `webhookSecret
| `/docker/inspect/{host}/{container}` | GET | Sanitized image, runtime, resources, mounts and state |
| `/docker/restart/{host}/{container}` | POST | Restart container; supports `dry_run=true` |
| `/config/reload` | POST | Reload YAML configuration and credential cache |
+| `/media/verify` | POST | ffprobe duration check on a NAS media file to catch truncated Sonarr/Radarr imports; flags `suspect` if deviation from `expected_minutes` exceeds `tolerance_pct` |
## VM lifecycle and inventory
@@ -97,6 +98,14 @@ Integration Compose definition: `tests/compose.integration.yaml` (binds only to
## Changelog
+### 2.4.8 — 17.09.2026
+
+- Fixed `POST /media/verify` false-positive rejection: filenames containing a legitimate apostrophe (e.g. "La'An" in a Star Trek episode title) were blocked as "unsafe characters". Replaced the naive single-quote wrapping + apostrophe blocklist with proper `shlex.quote()` escaping for the remote ffprobe command; only newline/NUL byte injection is still rejected.
+
+### 2.4.7 — 17.09.2026
+
+- Added `POST /media/verify`: probes a media file's real duration via `ffprobe` (running inside the existing `bazarrUHD` container on `arrapps`, which already mounts `/data` and ships ffmpeg — no new service needed) and flags it as `suspect` when it deviates from an `expected_minutes` value beyond `tolerance_pct`. Catches truncated Sonarr/Radarr imports (e.g. a re-grab that lands as 1.8 GB instead of the expected 8 GB for a UHD episode) after the file is already on the NAS.
+
### 2.3.2 — 22.07.2026
- Make Vaultwarden refresh durable: persistent named cache volume, protected runtime credentials, automatic API-key re-login and atomic cache writes.
diff --git a/app.py b/app.py
index 2b6da2e..7ddb58a 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, hashlib, shlex
from datetime import datetime, timezone
import httpx, yaml
from typing import Literal
@@ -9,15 +9,20 @@ from pydantic import BaseModel, Field
from fastapi import FastAPI, Request, HTTPException, Depends, Query
from fastapi.responses import JSONResponse, RedirectResponse, Response, HTMLResponse
from contextlib import asynccontextmanager
+from contextvars import ContextVar
log = logging.getLogger("butler")
-VERSION = "2.3.5"
+VERSION = "2.4.9"
API_DIR = os.environ.get("API_KEY_DIR", "/data/api")
VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache")
BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "")
CONFIG_PATH = os.environ.get("BUTLER_CONFIG", "/data/butler.yaml")
UI_PATH = os.environ.get("BUTLER_UI_PATH", os.path.join(os.path.dirname(__file__), "ui.html"))
+MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get(
+ "MEDIA_HANDOFF_ALLOWED_NETWORKS",
+ "10.2.1.119/32,10.5.85.12/32",
+)
# --- Config loading ---
@@ -46,6 +51,7 @@ _load_config()
# --- Audit log ---
_audit_log: list[dict] = []
+_audit_actor: ContextVar[str] = ContextVar("audit_actor", default="System/API")
MAX_AUDIT = 500
AUDIT_DB_PATH = os.environ.get("AUDIT_DB_PATH", "/data/state/audit.sqlite3")
@@ -69,8 +75,12 @@ def _init_audit_db() -> bool:
method TEXT NOT NULL,
status INTEGER NOT NULL,
detail TEXT NOT NULL,
- dry_run INTEGER NOT NULL
+ dry_run INTEGER NOT NULL,
+ actor TEXT NOT NULL DEFAULT 'Legacy/API'
)""")
+ columns = {row[1] for row in db.execute("PRAGMA table_info(audit)")}
+ if "actor" not in columns:
+ db.execute("ALTER TABLE audit ADD COLUMN actor TEXT NOT NULL DEFAULT 'Legacy/API'")
return True
except (OSError, sqlite3.Error):
return False
@@ -83,6 +93,7 @@ def _audit(endpoint: str, method: str, status: int, detail: str = "", dry_run: b
"status": status,
"detail": _redact_audit_detail(detail),
"dry_run": dry_run,
+ "actor": _audit_actor.get(),
}
_audit_log.append(entry)
if len(_audit_log) > MAX_AUDIT:
@@ -91,8 +102,8 @@ def _audit(endpoint: str, method: str, status: int, detail: str = "", dry_run: b
if _init_audit_db():
with sqlite3.connect(AUDIT_DB_PATH) as db:
db.execute(
- "INSERT INTO audit (ts, endpoint, method, status, detail, dry_run) VALUES (?, ?, ?, ?, ?, ?)",
- (entry["ts"], entry["endpoint"], entry["method"], entry["status"], entry["detail"], int(entry["dry_run"])),
+ "INSERT INTO audit (ts, endpoint, method, status, detail, dry_run, actor) VALUES (?, ?, ?, ?, ?, ?, ?)",
+ (entry["ts"], entry["endpoint"], entry["method"], entry["status"], entry["detail"], int(entry["dry_run"]), entry["actor"]),
)
db.execute("DELETE FROM audit WHERE id NOT IN (SELECT id FROM audit ORDER BY id DESC LIMIT ?)", (MAX_AUDIT,))
except (OSError, sqlite3.Error):
@@ -246,6 +257,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 +285,300 @@ 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 _clean_audit_actor(value: str) -> str:
+ cleaned = re.sub(r"[^\w .@/\-]", "", str(value or ""))[:40].strip()
+ return cleaned or "KI/API"
+
+
+@app.middleware("http")
+async def audit_actor_context(request: Request, call_next):
+ if request.headers.get("authorization", "").startswith("Bearer "):
+ actor = _clean_audit_actor(request.headers.get("x-butler-actor", "KI/API"))
+ elif _ui_session(request):
+ actor = "Weboberfläche"
+ else:
+ actor = "System/Öffentlich"
+ token = _audit_actor.set(actor)
+ try:
+ return await call_next(request)
+ finally:
+ _audit_actor.reset(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,
+ }
+
+
+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:
+ 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 = _emby_history_step(days)
+ 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
+ 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
+
+
+@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 +647,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")
@@ -433,7 +753,7 @@ async def capabilities(_=Depends(_verify)):
mode = "read_only" if method.upper() == "GET" else "mutation"
serialized = {"parameters": operation.get("parameters", []), "requestBody": operation.get("requestBody", {})}
dry_run = _schema_contains_property(serialized, "dry_run", components)
- critical = path.startswith(("/network/wireguard", "/caddy/"))
+ critical = path.startswith(("/network/wireguard", "/network/media-tunnel", "/caddy/"))
destructive = method.upper() == "DELETE" or any(
marker in path for marker in ("/destroy/", "/cleanup/", "/break-lock/", "/restore/")
)
@@ -568,7 +888,7 @@ async def audit(_=Depends(_verify), limit: int = Query(50, le=MAX_AUDIT)):
with sqlite3.connect(AUDIT_DB_PATH) as db:
db.row_factory = sqlite3.Row
rows = db.execute(
- "SELECT ts, endpoint, method, status, detail, dry_run FROM audit ORDER BY id DESC LIMIT ?",
+ "SELECT ts, endpoint, method, status, detail, dry_run, actor FROM audit ORDER BY id DESC LIMIT ?",
(limit,),
).fetchall()
return [dict(row) | {"dry_run": bool(row["dry_run"])} for row in rows]
@@ -1226,6 +1546,102 @@ async def vault_reload(_=Depends(_verify)):
return {"reloaded": True, "items": len(_vault_cache)}
+# --- Media file integrity verification ---
+
+MEDIA_VERIFY_PROBES = {
+ "arrapps": ("bazarrUHD", "sascha"),
+ "arr-chris": (None, "chris"),
+ "arr-chris-live": (None, "chris"),
+}
+
+
+class MediaVerifyRequest(BaseModel):
+ host: str = Field(..., max_length=64)
+ path: str = Field(..., max_length=1000)
+ expected_minutes: float | None = Field(None, ge=0, le=1440)
+ tolerance_pct: float = Field(20.0, gt=0, le=100)
+
+
+@app.post("/media/verify")
+async def media_verify(payload: MediaVerifyRequest, _=Depends(_verify)):
+ """Probe a media file's real duration via ffprobe to catch truncated imports.
+
+ Runs `ffprobe` inside an existing container (bazarrUHD on arrapps, which already
+ mounts /data and ships ffmpeg) so no new service is required. Compares the
+ measured duration against an expected runtime (minutes) supplied by the caller
+ (e.g. Sonarr/Radarr's runtime field) within tolerance_pct.
+ """
+ if not re.fullmatch(r"[A-Za-z0-9_.-]+", payload.host):
+ raise HTTPException(400, "Invalid host name")
+ if payload.host not in MEDIA_VERIFY_PROBES:
+ raise HTTPException(400, f"No ffprobe container configured for host {payload.host}")
+ if not re.fullmatch(r"/data/[^\x00]+\.(mkv|mp4|avi|m4v|ts)", payload.path):
+ raise HTTPException(400, "Path must be an absolute /data media file")
+ if ".." in payload.path or "\n" in payload.path or "\x00" in payload.path:
+ raise HTTPException(400, "Path contains unsafe characters")
+
+ container, _default_user = MEDIA_VERIFY_PROBES[payload.host]
+ if not container:
+ raise HTTPException(400, f"Host {payload.host} has no configured ffprobe container yet")
+
+ target = _find_inventory_host(payload.host)
+ if not target:
+ raise HTTPException(404, f"Host {payload.host} not found in inventory")
+
+ # expected_minutes == 0 or None: no reference runtime available, skip verification gracefully
+ if not payload.expected_minutes or payload.expected_minutes <= 0:
+ return {
+ "host": payload.host,
+ "path": payload.path,
+ "duration_seconds": None,
+ "duration_minutes": None,
+ "expected_minutes": payload.expected_minutes,
+ "verified": False,
+ "suspect": False,
+ "skipped": True,
+ "skip_reason": "no_reference_runtime",
+ "deviation_pct": None,
+ }
+
+ probe_cmd = (
+ f"docker exec {container} ffprobe -v error "
+ f"-show_entries format=duration -of default=noprint_wrappers=1:nokey=1 {shlex.quote(payload.path)}"
+ )
+ rc, out, err = await asyncio.to_thread(
+ _ssh, f'{target["user"]}@{target["ip"]}', f"sudo -n {probe_cmd} || {probe_cmd}", 30
+ )
+ if rc != 0:
+ _audit("/media/verify", "POST", 502, f"host={payload.host} path={payload.path}")
+ raise HTTPException(502, (err or out).strip()[:500] or "ffprobe failed")
+
+ raw_duration = out.strip()
+ try:
+ duration_seconds = float(raw_duration)
+ except ValueError:
+ _audit("/media/verify", "POST", 502, f"host={payload.host} unparsable duration")
+ raise HTTPException(502, "ffprobe returned no parsable duration; file is likely corrupt")
+
+ duration_minutes = duration_seconds / 60
+ result = {
+ "host": payload.host,
+ "path": payload.path,
+ "duration_seconds": round(duration_seconds, 1),
+ "duration_minutes": round(duration_minutes, 2),
+ "expected_minutes": payload.expected_minutes,
+ "verified": True,
+ "suspect": False,
+ }
+ if payload.expected_minutes:
+ deviation_pct = abs(duration_minutes - payload.expected_minutes) / payload.expected_minutes * 100
+ result["deviation_pct"] = round(deviation_pct, 1)
+ result["suspect"] = deviation_pct > payload.tolerance_pct
+ _audit(
+ "/media/verify", "POST", 200,
+ f"host={payload.host} dur={duration_minutes:.1f}m suspect={result['suspect']}",
+ )
+ return result
+
+
# --- VPS reverse-proxy and DNS management ---
VPS_SSH = "root@46.225.230.72"
@@ -1254,6 +1670,13 @@ class ProxyRouteRequest(BaseModel):
dns_token: str | None = None
+class DnsRrsetUpsertRequest(BaseModel):
+ record_type: Literal["A", "AAAA", "CNAME"] = "A"
+ value: str
+ ttl: int = Field(default=300, ge=60, le=86400)
+ comment: str = Field(default="Managed by Homelab Butler", max_length=200)
+
+
def _remote_python(script: str, timeout: int = 30) -> tuple[int, str, str]:
encoded = base64.b64encode(script.encode()).decode()
return _ssh(VPS_SSH, f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"', timeout=timeout)
@@ -1407,6 +1830,34 @@ SPEEDTEST_REPO_FILES = (
"streamscope/static/assets/longterm-metrics.js",
)
+GUCK_ADMIN_REPO_FILES = (
+ "guck-admin/Dockerfile",
+ "guck-admin/compose.yaml",
+ "guck-admin/requirements.txt",
+ "guck-admin/src/app.py",
+ "guck-admin/src/control.py",
+ "guck-admin/src/sharing_watchdog.py",
+ "guck-admin/src/templates/bandwidth.html",
+ "guck-admin/src/templates/base.html",
+ "guck-admin/src/templates/dashboard.html",
+ "guck-admin/src/templates/history.html",
+ "guck-admin/src/templates/sessions.html",
+ "guck-admin/src/templates/settings.html",
+ "guck-admin/src/templates/sharing.html",
+ "guck-admin/src/templates/users.html",
+ "guck-admin/static/admin.css",
+ "guck-admin/static/admin.js",
+ "guck-admin/static/icon.svg",
+ "guck-admin/static/manifest.webmanifest",
+ "guck-admin/static/sw.js",
+ "guck-admin/static/world.svg",
+)
+
+FORGEJO_DEPLOY_FILES = {
+ "sascha/speedtest": frozenset(SPEEDTEST_REPO_FILES),
+ "sascha/guck-vps": frozenset(GUCK_ADMIN_REPO_FILES),
+}
+
class SpeedtestDeployRequest(BaseModel):
stats_password: str
@@ -1414,7 +1865,7 @@ class SpeedtestDeployRequest(BaseModel):
async def _fetch_forgejo_text(repo: str, path: str) -> str:
- if repo != "sascha/speedtest" or path not in SPEEDTEST_REPO_FILES:
+ if path not in FORGEJO_DEPLOY_FILES.get(repo, frozenset()):
raise ValueError("unsupported Forgejo file")
cfg = SERVICES.get("forgejo", {})
base_url = cfg.get("url")
@@ -1518,6 +1969,143 @@ async def vps_speedtest_deploy(req: SpeedtestDeployRequest, _=Depends(_verify)):
return result
+class GuckAdminDeployRequest(BaseModel):
+ dry_run: bool = True
+
+
+def _validate_guck_admin_bundle(files: dict[str, str]) -> None:
+ if set(files) != set(GUCK_ADMIN_REPO_FILES):
+ raise ValueError("guck-admin source bundle is incomplete")
+ compose = files["guck-admin/compose.yaml"]
+ control = files["guck-admin/src/control.py"]
+ compose_required = (
+ "network_mode: host",
+ "NET_ADMIN",
+ "/app-config/guck-admin/data:/data",
+ "GUCK_LIMIT: /host/guck-limit.sh",
+ )
+ if any(item not in compose for item in compose_required):
+ raise ValueError("guck-admin compose is missing a required security or persistence setting")
+ policy_required = (
+ "CREATE TABLE IF NOT EXISTS custom_networks",
+ "def sync_custom_networks",
+ "2a00:8c40:f000::/36",
+ "45.58.235.0/24",
+ )
+ if any(item not in control for item in policy_required):
+ raise ValueError("guck-admin custom VPN policy is incomplete")
+
+
+def _guck_admin_remote_python(target: str, script: str, timeout: int = 60):
+ encoded = base64.b64encode(script.encode()).decode()
+ command = f"sudo -n python3 -c \"import base64;exec(base64.b64decode('{encoded}'))\""
+ return _ssh(target, command, timeout=timeout)
+
+
+def _deploy_guck_admin_compose(files: dict[str, str], dry_run: bool = True) -> dict:
+ _validate_guck_admin_bundle(files)
+ inventory = _find_inventory_host("guck-vps")
+ if not inventory:
+ raise RuntimeError("guck-vps is missing from Butler inventory")
+ target = f'{inventory["user"]}@{inventory["ip"]}'
+ if dry_run:
+ return {
+ "status": "validated",
+ "dry_run": True,
+ "host": "guck-vps",
+ "files": len(files),
+ "policy_networks": ["Mozilla Firefox VPN IPv6", "Fastly VPN IPv4"],
+ }
+ relative_files = {path.removeprefix("guck-admin/"): content for path, content in files.items()}
+ deploy_script = f"""from pathlib import Path
+import os, shutil
+stack = Path('/app-config/guck-admin')
+backup = Path('/app-config/deployment-backups/guck-admin-rollback')
+files = {relative_files!r}
+if backup.exists():
+ shutil.rmtree(backup)
+backup.mkdir(parents=True, exist_ok=True)
+stack.mkdir(parents=True, exist_ok=True)
+for relative, content in files.items():
+ target = stack / relative
+ old = backup / relative
+ if target.exists():
+ old.parent.mkdir(parents=True, exist_ok=True)
+ shutil.copy2(target, old)
+ target.parent.mkdir(parents=True, exist_ok=True)
+ temporary = target.with_name(target.name + '.butler-new')
+ temporary.write_text(content)
+ os.replace(temporary, target)
+"""
+ rollback_script = f"""from pathlib import Path
+import os, shutil
+stack = Path('/app-config/guck-admin')
+backup = Path('/app-config/deployment-backups/guck-admin-rollback')
+files = {tuple(relative_files)!r}
+for relative in files:
+ target = stack / relative
+ old = backup / relative
+ if old.exists():
+ target.parent.mkdir(parents=True, exist_ok=True)
+ shutil.copy2(old, target)
+ elif target.exists():
+ target.unlink()
+"""
+ rc, _out, err = _guck_admin_remote_python(target, deploy_script, timeout=90)
+ if rc != 0:
+ raise RuntimeError(f"guck-admin file deployment failed: {err[-300:]}")
+
+ def rollback():
+ _guck_admin_remote_python(target, rollback_script, timeout=90)
+ _ssh(target, "cd /app-config/guck-admin && sudo -n docker compose up -d --build --remove-orphans", timeout=600)
+
+ preflight = "cd /app-config/guck-admin && sudo -n docker compose config -q && sudo -n docker compose build --pull"
+ rc, _out, err = _ssh(target, preflight, timeout=600)
+ if rc != 0:
+ rollback()
+ raise RuntimeError(f"guck-admin build preflight failed: {err[-500:]}")
+ deploy = "cd /app-config/guck-admin && sudo -n docker compose up -d --remove-orphans"
+ rc, out, err = _ssh(target, deploy, timeout=240)
+ if rc != 0:
+ rollback()
+ raise RuntimeError(f"guck-admin deployment failed: {(err or out)[-500:]}")
+ health = "for i in $(seq 1 45); do curl -fsS --max-time 3 http://127.0.0.1:9090/health >/dev/null && exit 0; sleep 2; done; exit 1"
+ rc, _out, err = _ssh(target, health, timeout=105)
+ if rc != 0:
+ rollback()
+ raise RuntimeError(f"guck-admin health check failed and rollback was attempted: {err[-300:]}")
+ policy = "curl -fsS -X POST --max-time 120 http://127.0.0.1:9090/actions/limiter/refresh >/dev/null && sudo -n ipset test vpn-v6 2a00:8c40:f02d:a34c::1 && sudo -n ipset test vpn-v4 45.58.235.7 && sudo -n tc class show dev ens3 | python3 -c \"import sys; s=sys.stdin.read(); raise SystemExit(0 if '1:300' in s else 1)\""
+ rc, out, err = _ssh(target, policy, timeout=180)
+ if rc != 0:
+ rollback()
+ raise RuntimeError(f"guck-admin policy verification failed and rollback was attempted: {(err or out)[-500:]}")
+ return {
+ "status": "deployed",
+ "health": "ok",
+ "policy": "verified",
+ "host": "guck-vps",
+ "custom_networks": ["2a00:8c40:f000::/36", "45.58.235.0/24"],
+ "ipv4": True,
+ "ipv6": True,
+ }
+
+
+@app.post("/vps/guck-admin/deploy")
+async def vps_guck_admin_deploy(req: GuckAdminDeployRequest, _=Depends(_verify)):
+ try:
+ contents = await asyncio.gather(*(
+ _fetch_forgejo_text("sascha/guck-vps", path) for path in GUCK_ADMIN_REPO_FILES
+ ))
+ files = dict(zip(GUCK_ADMIN_REPO_FILES, contents))
+ result = await asyncio.to_thread(_deploy_guck_admin_compose, files, req.dry_run)
+ except ValueError as exc:
+ raise HTTPException(400, str(exc)) from exc
+ except Exception as exc:
+ raise HTTPException(502, f"guck-admin deployment failed: {str(exc)[-500:]}") from exc
+ _audit("/vps/guck-admin/deploy", "POST", 200, f"dry_run={req.dry_run}; Git-managed custom VPN policy")
+ return result
+
+
# --- VM Lifecycle Endpoints ---
import subprocess as _sp
@@ -1542,6 +2130,91 @@ def _ssh(host, cmd, timeout=600):
return 124, "", f"SSH command timed out after {timeout} seconds"
+PAPERLESS_GATEWAY = AUTOMATION1
+PAPERLESS_SSH = "sascha@10.5.1.120"
+PAPERLESS_CONSUME_DIR = "/app-config/paperless/consume"
+PAPERLESS_MAX_IMPORT_BYTES = 50 * 1024 * 1024
+
+
+def _ssh_bytes(host: str, cmd: str, payload: bytes, timeout: int = 120):
+ """Send a binary payload to a fixed remote command over Butler-managed SSH."""
+ try:
+ result = _sp.run(
+ ["ssh", "-o", "ConnectTimeout=10", "-o", "StrictHostKeyChecking=accept-new",
+ "-o", "UserKnownHostsFile=/tmp/butler_known_hosts", host, cmd],
+ input=payload, capture_output=True, timeout=timeout,
+ )
+ return result.returncode, result.stdout.decode(errors="replace"), result.stderr.decode(errors="replace")
+ except _sp.TimeoutExpired:
+ return 124, "", f"SSH upload timed out after {timeout} seconds"
+
+
+def _paperless_import_filename(filename: str, digest: str) -> str:
+ base_name = os.path.basename(filename or "document.pdf")
+ stem = re.sub(r"[^A-Za-z0-9._-]+", "_", base_name.rsplit(".", 1)[0]).strip("._-")
+ if not stem:
+ stem = "document"
+ return f"{stem[:100]}-{digest[:12]}.pdf"
+
+
+def _paperless_gateway_command(command: str) -> str:
+ """Route through automation1, whose deployment key is authorized on Homelab VMs."""
+ if "'" in command:
+ raise ValueError("Paperless remote command contains an unsafe quote")
+ return (
+ "ssh -o ConnectTimeout=10 -o StrictHostKeyChecking=accept-new "
+ f"-o UserKnownHostsFile=/tmp/paperless_known_hosts {PAPERLESS_SSH} '{command}'"
+ )
+
+
+@app.post("/paperless/import")
+async def paperless_import(request: Request, filename: str = Query(..., min_length=1, max_length=180), _=Depends(_verify)):
+ """Queue a PDF in Paperless without exposing Paperless or SSH to the caller."""
+ payload = await request.body()
+ if not payload or not payload.startswith(b"%PDF-"):
+ raise HTTPException(400, "Only valid PDF documents are accepted")
+ if len(payload) > PAPERLESS_MAX_IMPORT_BYTES:
+ raise HTTPException(413, "PDF exceeds the 50 MiB import limit")
+ digest = hashlib.sha256(payload).hexdigest()
+ import_name = _paperless_import_filename(filename, digest)
+ remote_path = f"{PAPERLESS_CONSUME_DIR}/{import_name}"
+ command = (
+ f"sudo -n install -d -o sascha -g sascha -m 0755 {PAPERLESS_CONSUME_DIR} && "
+ f"tmp=$(mktemp /tmp/paperless-import.XXXXXX) && "
+ f"cat > \"$tmp\" && sudo -n install -o sascha -g sascha -m 0644 \"$tmp\" {remote_path} && rm -f \"$tmp\""
+ )
+ rc, out, err = await asyncio.to_thread(
+ _ssh_bytes, PAPERLESS_GATEWAY, _paperless_gateway_command(command), payload, 180
+ )
+ if rc != 0:
+ raise HTTPException(502, (err or out).strip()[-500:] or "Paperless import transfer failed")
+ _audit("/paperless/import", "POST", 202, f"sha256={digest}; bytes={len(payload)}")
+ return {"status": "queued", "filename": import_name, "sha256": digest, "bytes": len(payload)}
+
+
+@app.get("/paperless/import/status")
+async def paperless_import_status(filename: str = Query(..., min_length=1, max_length=180), _=Depends(_verify)):
+ if not re.fullmatch(r"[A-Za-z0-9._-]+\.pdf", filename):
+ raise HTTPException(400, "Invalid import filename")
+ remote_path = f"{PAPERLESS_CONSUME_DIR}/{filename}"
+ command = (
+ f"if sudo -n test -f {remote_path}; then echo QUEUED; else echo CONSUMED; fi; "
+ f"logs=$(sudo -n docker logs --since 15m paperless-ngx 2>&1); "
+ f"printf \"%s\\n\" \"$logs\" | grep -F -- {filename} | tail -20 || true; "
+ f"task=$(printf \"%s\\n\" \"$logs\" | grep -F -- {filename} | "
+ f"grep -oE \"\\[[0-9a-f]{{8}}\\]\" | tail -1 | tr -d \"[]\"); "
+ f"if test -n \"$task\"; then printf \"%s\\n\" \"$logs\" | grep -F -- \"[$task]\" | tail -20; fi"
+ )
+ rc, out, err = await asyncio.to_thread(
+ _ssh, PAPERLESS_GATEWAY, _paperless_gateway_command(command), 45
+ )
+ if rc != 0:
+ raise HTTPException(502, (err or out).strip()[-500:] or "Paperless status check failed")
+ lines = out.splitlines()
+ state = lines[0].strip().lower() if lines else "unknown"
+ return {"status": state, "filename": filename, "recent_log": lines[1:]}
+
+
def _wireguard_status_command() -> str:
script = '''import json, subprocess
@@ -1682,6 +2355,512 @@ async def network_wireguard_remove_peer(host: str, req: WireGuardPeerRemoveReque
return {"host": host, **result}
+SASCHA_MEDIA_VPS_HOST = "pfannkuchen"
+SASCHA_MEDIA_EMBY_HOST = "emby-sascha"
+SASCHA_MEDIA_INTERFACE = "wg-media"
+SASCHA_MEDIA_VPS_ADDRESS = "10.11.13.1/32"
+SASCHA_MEDIA_EMBY_ADDRESS = "10.11.13.3/32"
+SASCHA_MEDIA_PORT = 51821
+SASCHA_MEDIA_MTU = 1340
+SASCHA_MEDIA_CONFIRMATION = "DEPLOY_DIRECT_SASCHA_MEDIA_TUNNEL"
+
+
+class SaschaMediaTunnelRequest(BaseModel):
+ dry_run: bool = True
+ confirmation: str | None = None
+
+
+class SaschaMediaEdgeOptimizeRequest(BaseModel):
+ dry_run: bool = True
+ confirmation: str | None = None
+
+
+def _sascha_media_edge_audit_command() -> str:
+ script = '''import json, re, subprocess
++from pathlib import Path
++
++def run(args):
++ proc = subprocess.run(args, text=True, capture_output=True, timeout=30)
++ return {"rc": proc.returncode, "stdout": proc.stdout.strip(), "stderr": proc.stderr.strip()[-300:]}
++
++result = {"hostname": "pfannkuchen", "caddy": {"container_running": False, "protocols": [], "protocols_explicit": False, "upstreams": [], "site_block": [], "emby_snippet": []}, "network": {}}
++inspect = run(["docker", "inspect", "caddy", "--format", "{{.State.Running}}"])
++result["caddy"]["container_running"] = inspect["rc"] == 0 and inspect["stdout"] == "true"
++adapt = run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"])
++if adapt["rc"] == 0:
++ try:
++ config = json.loads(adapt["stdout"])
++ servers = config.get("apps", {}).get("http", {}).get("servers", {})
++ explicit = []
++ def walk(value, matched=False):
++ if isinstance(value, dict):
++ current = matched
++ host = value.get("host")
++ if isinstance(host, list) and "tv.sascha-lutz.de" in host:
++ current = True
++ if current and isinstance(value.get("dial"), str):
++ result["caddy"]["upstreams"].append(value["dial"])
++ for child in value.values(): walk(child, current)
++ elif isinstance(value, list):
++ for child in value: walk(child, matched)
++ for server in servers.values():
++ protocols = server.get("protocols")
++ if isinstance(protocols, list): explicit.extend(protocols)
++ walk(server)
++ result["caddy"]["protocols_explicit"] = bool(explicit)
++ result["caddy"]["protocols"] = sorted(set(explicit)) if explicit else ["h1", "h2", "h3"]
++ result["caddy"]["upstreams"] = sorted(set(result["caddy"]["upstreams"]))
++ except Exception as exc:
++ result["caddy"]["adapt_error"] = str(exc)[:200]
++else:
++ result["caddy"]["adapt_error"] = adapt["stderr"] or "caddy adapt failed"
++
++path = Path("/app-config/caddy/Caddyfile")
++if path.exists():
++ lines = path.read_text(encoding="utf-8", errors="replace").splitlines()
++ collecting = False; depth = 0; selected = []
++ for line in lines:
++ if not collecting and re.match(r"^\\s*tv\\.sascha-lutz\\.de\\s*\\{", line): collecting = True
++ if collecting:
++ depth += line.count("{") - line.count("}")
++ if not re.search(r"(?i)(password|secret|token|private|basicauth|basic_auth|hash)", line): selected.append(line.strip())
++ else: selected.append("[REDACTED SENSITIVE DIRECTIVE]")
++ if depth == 0: break
++ result["caddy"]["site_block"] = selected
++ collecting = False; depth = 0; selected = []
++ for line in lines:
++ if not collecting and re.match(r"^\\s*\\(emby_config\\)\\s*\\{", line): collecting = True
++ if collecting:
++ depth += line.count("{") - line.count("}")
++ if not re.search(r"(?i)(password|secret|token|private|basicauth|basic_auth|hash)", line): selected.append(line.strip())
++ else: selected.append("[REDACTED SENSITIVE DIRECTIVE]")
++ if depth == 0: break
++ result["caddy"]["emby_snippet"] = selected
++
++result["network"]["wg_media_active"] = subprocess.run(["systemctl", "is-active", "--quiet", "wg-quick@wg-media"]).returncode == 0
++result["network"]["wg_media_enabled"] = subprocess.run(["systemctl", "is-enabled", "--quiet", "wg-quick@wg-media"]).returncode == 0
++result["network"]["udp_51821"] = "51821" in run(["ss", "-H", "-lun"]) ["stdout"]
++for name, target in (("legacy_route", "10.6.1.103"), ("direct_route", "10.11.13.3")):
++ result["network"][name] = run(["ip", "route", "get", target])["stdout"][:300]
++print(json.dumps(result))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+@app.get("/media/edge/sascha")
+async def sascha_media_edge_audit(_=Depends(_verify)):
+ """Return a redacted snapshot of the dedicated Hetzner edge for tv.sascha-lutz.de."""
+ inventory = await asyncio.to_thread(_find_inventory_host, SASCHA_MEDIA_VPS_HOST)
+ if not inventory:
+ raise HTTPException(404, "Hetzner media edge not found")
+ target = f'{inventory["user"]}@{inventory["ip"]}'
+ rc, out, err = await asyncio.to_thread(_ssh, target, _sascha_media_edge_audit_command(), 45)
+ if rc != 0:
+ raise HTTPException(502, (err or out).strip()[-500:] or "Sascha media edge audit failed")
+ try:
+ result = json.loads(out)
+ except json.JSONDecodeError as exc:
+ raise HTTPException(502, "Sascha media edge audit returned invalid JSON") from exc
+ _audit("/media/edge/sascha", "GET", 200, "redacted live media-edge snapshot")
+ return result
+
+
+def _sascha_media_edge_optimize_command() -> str:
+ script = '''import hashlib, json, os, re, shutil, subprocess, time
++from pathlib import Path
++path = Path("/app-config/caddy/Caddyfile")
++text = path.read_text(encoding="utf-8")
++old_pattern = re.compile(r"(?m)^(\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+)10\\.6\\.1\\.103:8096(\\s*)$")
++new_pattern = re.compile(r"(?m)^\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+10\\.11\\.13\\.3:8096\\s*$")
++old_count = len(old_pattern.findall(text)); new_count = len(new_pattern.findall(text))
++if old_count == 1:
++ text = old_pattern.sub(r"\\g<1>10.11.13.3:8096\\g<2>", text)
++elif not (old_count == 0 and new_count == 1):
++ raise RuntimeError("expected exactly one tv.sascha-lutz.de upstream")
++lines = text.splitlines(keepends=True)
++first = next((i for i, line in enumerate(lines) if line.strip() and not line.lstrip().startswith("#")), None)
++if first is None or lines[first].strip() != "{":
++ lines[0:0] = ["{\\n", " servers {\\n", " protocols h1 h2\\n", " }\\n", "}\\n", "\\n"]
++else:
++ depth = 0; global_end = None; servers_start = None; servers_end = None
++ for i in range(first, len(lines)):
++ stripped = lines[i].strip(); before = depth
++ if before == 1 and re.match(r"^servers(?:\\s+\\S+)?\\s*\\{$", stripped): servers_start = i
++ depth += lines[i].count("{") - lines[i].count("}")
++ if servers_start is not None and i > servers_start and depth == 1 and servers_end is None: servers_end = i
++ if i > first and depth == 0: global_end = i; break
++ if global_end is None: raise RuntimeError("unbalanced Caddy global options block")
++ if servers_start is None:
++ lines[global_end:global_end] = [" servers {\\n", " protocols h1 h2\\n", " }\\n"]
++ else:
++ if servers_end is None: raise RuntimeError("unbalanced Caddy servers block")
++ protocol_lines = [i for i in range(servers_start + 1, servers_end) if re.match(r"^\\s*protocols\\s+", lines[i])]
++ if len(protocol_lines) > 1: raise RuntimeError("multiple Caddy protocol directives")
++ if protocol_lines: lines[protocol_lines[0]] = re.sub(r"protocols\\s+.*", "protocols h1 h2", lines[protocol_lines[0]])
++ else: lines.insert(servers_start + 1, " protocols h1 h2\\n")
++text = "".join(lines)
++if re.search(r"(?im)^\\s*header(?:_down)?\\s+Alt-Svc", text): raise RuntimeError("manual Alt-Svc directive requires review")
++stamp = time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())
++backup = path.with_name("Caddyfile.pre-sascha-media-" + stamp)
++candidate = path.with_name("Caddyfile.sascha-media-candidate")
++shutil.copy2(path, backup); candidate.write_text(text, encoding="utf-8"); os.chmod(candidate, path.stat().st_mode)
++def run(args, timeout=30):
++ proc = subprocess.run(args, text=True, capture_output=True, timeout=timeout)
++ if proc.returncode != 0: raise RuntimeError((proc.stderr or proc.stdout).strip()[-500:] or "command failed")
++ return proc.stdout.strip()
++try:
++ run(["docker", "cp", str(candidate), "caddy:/tmp/Caddyfile.sascha-media-candidate"])
++ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/tmp/Caddyfile.sascha-media-candidate"])
++ with path.open("w", encoding="utf-8") as handle:
++ handle.write(text); handle.flush(); os.fsync(handle.fileno())
++ host_hash = hashlib.sha256(path.read_bytes()).hexdigest()
++ container_hash = run(["docker", "exec", "caddy", "sha256sum", "/etc/caddy/Caddyfile"]).split()[0]
++ if host_hash != container_hash: raise RuntimeError("host/container Caddyfile hash mismatch")
++ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/etc/caddy/Caddyfile"])
++ run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"])
++ adapted = json.loads(run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"]))
++ protocols = []
++ for server in adapted.get("apps", {}).get("http", {}).get("servers", {}).values(): protocols.extend(server.get("protocols", []))
++ if sorted(set(protocols)) != ["h1", "h2"]: raise RuntimeError("Caddy did not load h1+h2-only protocols")
++ if not run(["curl", "-fsS", "--max-time", "15", "http://10.11.13.3:8096/System/Ping"]): raise RuntimeError("Emby direct-path ping was empty")
++except Exception:
++ with path.open("w", encoding="utf-8") as handle:
++ handle.write(backup.read_text(encoding="utf-8")); handle.flush(); os.fsync(handle.fileno())
++ subprocess.run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"], text=True, capture_output=True, timeout=30)
++ raise
++finally:
++ candidate.unlink(missing_ok=True)
++ subprocess.run(["docker", "exec", "caddy", "rm", "-f", "/tmp/Caddyfile.sascha-media-candidate"], text=True, capture_output=True)
++print(json.dumps({"status": "optimized", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "backup": str(backup), "sha256": host_hash}))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+def _optimize_sascha_media_edge() -> dict:
+ return _ssh_json(_media_host_target(SASCHA_MEDIA_VPS_HOST), _sascha_media_edge_optimize_command(), 90)
+
+
+@app.post("/media/edge/sascha/optimize")
+async def optimize_sascha_media_edge(req: SaschaMediaEdgeOptimizeRequest, _=Depends(_verify)):
+ """Switch only tv.sascha-lutz.de to wg-media and disable HTTP/3 on its dedicated Hetzner edge."""
+ plan = {"status": "would_optimize", "hostname": "tv.sascha-lutz.de", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "zero_downtime_reload": True}
+ if req.dry_run:
+ _audit("/media/edge/sascha/optimize", "POST", 200, "dry_run=True", True)
+ return plan
+ if req.confirmation != "OPTIMIZE_TV_SASCHA_LUTZ_DE":
+ raise HTTPException(400, "confirmation must be OPTIMIZE_TV_SASCHA_LUTZ_DE")
+ try:
+ result = await asyncio.to_thread(_optimize_sascha_media_edge)
+ except Exception as exc:
+ _audit("/media/edge/sascha/optimize", "POST", 502, "optimization failed; Caddy rollback attempted")
+ raise HTTPException(502, str(exc)[-500:]) from exc
+ _audit("/media/edge/sascha/optimize", "POST", 200, "direct upstream and h1+h2 enabled")
+ return result
+
+
+def _media_tunnel_key_command() -> str:
+ script = '''import json, os, subprocess
++from pathlib import Path
++root = Path("/app-config/wireguard-media")
++private = root / "private.key"
++public = root / "public.key"
++root.mkdir(parents=True, exist_ok=True)
++os.chmod(root, 0o700)
++created = not private.exists()
++if created:
++ key = subprocess.run(["wg", "genkey"], check=True, text=True, capture_output=True).stdout.strip()
++ private.write_text(key + "\\n")
++ os.chmod(private, 0o600)
++if not public.exists() or created:
++ pub = subprocess.run(["wg", "pubkey"], input=private.read_text(), check=True, text=True, capture_output=True).stdout.strip()
++ public.write_text(pub + "\\n")
++ os.chmod(public, 0o644)
++print(json.dumps({"public_key": public.read_text().strip(), "created": created}))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+def _media_tunnel_install_command(role: Literal["vps", "emby"], peer_public_key: str) -> str:
+ if not re.fullmatch(r"[A-Za-z0-9+/]{43}=", peer_public_key):
+ raise ValueError("Invalid WireGuard public key")
+ settings = {
+ "vps": {
+ "address": SASCHA_MEDIA_VPS_ADDRESS,
+ "peer": SASCHA_MEDIA_EMBY_ADDRESS,
+ "endpoint": None,
+ "keepalive": None,
+ "listen": SASCHA_MEDIA_PORT,
+ },
+ "emby": {
+ "address": SASCHA_MEDIA_EMBY_ADDRESS,
+ "peer": SASCHA_MEDIA_VPS_ADDRESS,
+ "endpoint": f"46.225.230.72:{SASCHA_MEDIA_PORT}",
+ "keepalive": 15,
+ "listen": None,
+ },
+ }[role]
+ script = f'''import json, os, shutil, subprocess, time
++from pathlib import Path
++root = Path("/app-config/wireguard-media")
++config = root / "wg-media.conf"
++etc = Path("/etc/wireguard/wg-media.conf")
++backup_dir = root / "backups"
++backup_dir.mkdir(parents=True, exist_ok=True)
++private = (root / "private.key").read_text().strip()
++listen = {settings['listen']!r}
++existed = config.exists()
++backup = None
++if existed:
++ backup = backup_dir / ("wg-media.conf." + time.strftime("%Y%m%dT%H%M%SZ", time.gmtime()))
++ shutil.copy2(config, backup)
++lines = ["[Interface]", "Address = {settings['address']}", "MTU = {SASCHA_MEDIA_MTU}", "PrivateKey = " + private]
++if listen is not None: lines.append("ListenPort = " + str(listen))
++if {role!r} == "vps":
++ lines.extend(["PostUp = iptables -C INPUT -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT 2>/dev/null || iptables -I INPUT 1 -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT", "PreDown = iptables -D INPUT -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT 2>/dev/null || true"])
++lines.extend(["", "[Peer]", "PublicKey = {peer_public_key}", "AllowedIPs = {settings['peer']}"])
++if {settings['endpoint']!r}: lines.append("Endpoint = " + {settings['endpoint']!r})
++if {settings['keepalive']!r}: lines.append("PersistentKeepalive = " + str({settings['keepalive']!r}))
++candidate = root / "wgmtest.conf"
++candidate.write_text("\\n".join(lines) + "\\n")
++os.chmod(candidate, 0o600)
++check = subprocess.run(["wg-quick", "strip", str(candidate)], text=True, capture_output=True)
++if check.returncode != 0:
++ candidate.unlink(missing_ok=True)
++ raise RuntimeError(check.stderr.strip() or "wg-quick validation failed")
++os.replace(candidate, config)
++os.chmod(config, 0o600)
++etc.parent.mkdir(parents=True, exist_ok=True)
++if etc.is_symlink() or etc.exists():
++ if etc.is_symlink() and etc.resolve() == config.resolve(): pass
++ elif etc.exists():
++ etc_backup = backup_dir / ("etc-wg-media.conf." + time.strftime("%Y%m%dT%H%M%SZ", time.gmtime()))
++ shutil.move(etc, etc_backup)
++ etc.symlink_to(config)
++else: etc.symlink_to(config)
++proc = subprocess.run(["systemctl", "enable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30)
++if proc.returncode != 0:
++ if backup: shutil.copy2(backup, config)
++ else: config.unlink(missing_ok=True)
++ raise RuntimeError(proc.stderr.strip() or proc.stdout.strip() or "failed to start wg-media")
++print(json.dumps({{"role": {role!r}, "status": "configured", "address": {settings['address']!r}, "backup": str(backup) if backup else None, "existed": existed}}))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+def _media_tunnel_verify_command(source: str, destination: str, check_emby: bool = False) -> str:
+ extra = ''
+ if check_emby:
+ extra = '''\nhttp = subprocess.run(["curl", "-sS", "--max-time", "10", "-o", "/dev/null", "-w", "%{http_code}", "http://10.11.13.3:8096/System/Ping"], text=True, capture_output=True)\nresult["emby_http"] = http.stdout.strip() if http.returncode == 0 else "000"'''
+ script = f'''import json, subprocess, time
++result = {{"ping": False, "handshake": False}}
++for _ in range(10):
++ ping = subprocess.run(["ping", "-c", "1", "-W", "2", "-I", {source!r}, {destination!r}], text=True, capture_output=True)
++ hand = subprocess.run(["sudo", "-n", "wg", "show", "wg-media", "latest-handshakes"], text=True, capture_output=True)
++ now = int(time.time())
++ stamps = []
++ for line in hand.stdout.splitlines():
++ try: stamps.append(int(line.split()[-1]))
++ except Exception: pass
++ result["ping"] = ping.returncode == 0
++ result["handshake"] = any(stamp > 0 and now - stamp < 60 for stamp in stamps)
++ if result["ping"] and result["handshake"]: break
++ time.sleep(2)
++{extra}
++print(json.dumps(result))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+def _media_tunnel_rollback_command(remove_keys: bool) -> str:
+ script = f'''import json, shutil, subprocess
++from pathlib import Path
++root = Path("/app-config/wireguard-media")
++config = root / "wg-media.conf"
++backups = sorted((root / "backups").glob("wg-media.conf.*")) if (root / "backups").exists() else []
++subprocess.run(["systemctl", "disable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30)
++if backups:
++ shutil.copy2(backups[-1], config)
++ subprocess.run(["systemctl", "enable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30)
++ status = "restored"
++else:
++ config.unlink(missing_ok=True)
++ Path("/etc/wireguard/wg-media.conf").unlink(missing_ok=True)
++ if {remove_keys!r}:
++ (root / "private.key").unlink(missing_ok=True); (root / "public.key").unlink(missing_ok=True)
++ status = "removed"
++print(json.dumps({{"status": status}}))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+def _media_tunnel_key_cleanup_command() -> str:
+ script = '''import json
++from pathlib import Path
++root = Path("/app-config/wireguard-media")
++if not (root / "wg-media.conf").exists():
++ (root / "private.key").unlink(missing_ok=True)
++ (root / "public.key").unlink(missing_ok=True)
++ status = "new_keys_removed"
++else:
++ status = "kept_for_existing_config"
++print(json.dumps({"status": status}))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+def _media_host_target(name: str) -> str:
+ inventory = _find_inventory_host(name)
+ if not inventory:
+ raise RuntimeError(f"Inventory host missing: {name}")
+ return f'{inventory["user"]}@{inventory["ip"]}'
+
+
+def _ssh_json(target: str, command: str, timeout: int = 45) -> dict:
+ rc, out, err = _ssh(target, command, timeout)
+ if rc != 0:
+ raise RuntimeError((err or out).strip()[-500:] or f"remote command failed on {target}")
+ try:
+ return json.loads(out)
+ except json.JSONDecodeError as exc:
+ raise RuntimeError(f"remote command returned invalid JSON on {target}") from exc
+
+
+def _deploy_sascha_media_tunnel() -> dict:
+ vps = _media_host_target(SASCHA_MEDIA_VPS_HOST)
+ emby = _media_host_target(SASCHA_MEDIA_EMBY_HOST)
+ vps_key = _ssh_json(vps, _media_tunnel_key_command())
+ emby_key = _ssh_json(emby, _media_tunnel_key_command())
+ configured = []
+ try:
+ vps_install = _ssh_json(vps, _media_tunnel_install_command("vps", emby_key["public_key"]), 60)
+ configured.append((vps, bool(vps_key.get("created"))))
+ emby_install = _ssh_json(emby, _media_tunnel_install_command("emby", vps_key["public_key"]), 60)
+ configured.append((emby, bool(emby_key.get("created"))))
+ vps_check = _ssh_json(vps, _media_tunnel_verify_command("10.11.13.1", "10.11.13.3", True), 35)
+ emby_check = _ssh_json(emby, _media_tunnel_verify_command("10.11.13.3", "10.11.13.1"), 35)
+ if not (vps_check.get("ping") and vps_check.get("handshake") and emby_check.get("ping") and emby_check.get("handshake")):
+ raise RuntimeError("direct media tunnel verification failed")
+ if vps_check.get("emby_http") not in {"200", "401"}:
+ raise RuntimeError("Emby did not answer through the direct media tunnel")
+ return {
+ "status": "deployed", "interface": SASCHA_MEDIA_INTERFACE,
+ "vps_address": SASCHA_MEDIA_VPS_ADDRESS, "emby_address": SASCHA_MEDIA_EMBY_ADDRESS,
+ "listen_port": SASCHA_MEDIA_PORT, "mtu": SASCHA_MEDIA_MTU,
+ "handshake": True, "ping_vps_to_emby": True, "ping_emby_to_vps": True,
+ "emby_http": vps_check.get("emby_http"),
+ "rollback_backups": [vps_install.get("backup"), emby_install.get("backup")],
+ }
+ except Exception:
+ for target, created in reversed(configured):
+ try: _ssh_json(target, _media_tunnel_rollback_command(created), 60)
+ except Exception: pass
+ installed_targets = {target for target, _created in configured}
+ for target, key in ((vps, vps_key), (emby, emby_key)):
+ if target not in installed_targets and key.get("created"):
+ try: _ssh_json(target, _media_tunnel_key_cleanup_command(), 30)
+ except Exception: pass
+ raise
+
+
+@app.post("/network/media-tunnel/sascha")
+async def deploy_sascha_media_tunnel(req: SaschaMediaTunnelRequest, _=Depends(_verify)):
+ """Deploy the fixed direct Hetzner-to-emby-sascha WireGuard media tunnel."""
+ plan = {
+ "status": "would_deploy", "interface": SASCHA_MEDIA_INTERFACE,
+ "vps_address": SASCHA_MEDIA_VPS_ADDRESS, "emby_address": SASCHA_MEDIA_EMBY_ADDRESS,
+ "listen_port": SASCHA_MEDIA_PORT, "mtu": SASCHA_MEDIA_MTU,
+ "allowed_ips": [SASCHA_MEDIA_VPS_ADDRESS, SASCHA_MEDIA_EMBY_ADDRESS],
+ "keeps_legacy_node6_path": True,
+ }
+ if req.dry_run:
+ _audit("/network/media-tunnel/sascha", "POST", 200, "dry_run=True", True)
+ return plan
+ if req.confirmation != SASCHA_MEDIA_CONFIRMATION:
+ raise HTTPException(400, f"confirmation must be {SASCHA_MEDIA_CONFIRMATION}")
+ try:
+ result = await asyncio.to_thread(_deploy_sascha_media_tunnel)
+ except Exception as exc:
+ _audit("/network/media-tunnel/sascha", "POST", 502, "deployment failed; rollback attempted")
+ raise HTTPException(502, str(exc)[-500:]) from exc
+ _audit("/network/media-tunnel/sascha", "POST", 200, "direct media tunnel deployed")
+ return result
+
+
+def _media_benchmark_server_command() -> str:
+ server = '''import socket
++payload = b"\\0" * (1024 * 1024)
++with socket.socket() as listener:
++ listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
++ listener.bind(("0.0.0.0", 5209)); listener.listen(4); listener.settimeout(80)
++ for _ in range(2):
++ conn, _addr = listener.accept()
++ with conn:
++ conn.settimeout(30)
++ for _ in range(256): conn.sendall(payload)
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(server.encode()).decode()
+ command = f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+ return (
+ "sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true; "
+ "sudo -n systemd-run --unit=butler-media-benchmark --collect --property=RuntimeMaxSec=90 "
+ f"/bin/sh -c {json.dumps(command)}"
+ )
+
+
+def _media_benchmark_client_command() -> str:
+ script = '''import json, socket, time
++results = {}
++for name, host in (("legacy_node6", "10.6.1.103"), ("direct_wg_media", "10.11.13.3")):
++ total = 0; started = time.monotonic()
++ with socket.create_connection((host, 5209), timeout=10) as conn:
++ conn.settimeout(40)
++ while True:
++ chunk = conn.recv(1024 * 1024)
++ if not chunk: break
++ total += len(chunk)
++ elapsed = time.monotonic() - started
++ results[name] = {"bytes": total, "seconds": round(elapsed, 3), "mbit_s": round(total * 8 / elapsed / 1000000, 1)}
++print(json.dumps(results))
++'''.replace("\n+", "\n")
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+def _benchmark_sascha_media_paths() -> dict:
+ vps = _media_host_target(SASCHA_MEDIA_VPS_HOST)
+ emby = _media_host_target(SASCHA_MEDIA_EMBY_HOST)
+ rc, out, err = _ssh(emby, _media_benchmark_server_command(), 30)
+ if rc != 0:
+ raise RuntimeError((err or out).strip()[-500:] or "failed to start benchmark server")
+ time.sleep(2)
+ try:
+ result = _ssh_json(vps, _media_benchmark_client_command(), 90)
+ finally:
+ _ssh(emby, "sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true", 20)
+ if any(item.get("bytes") != 256 * 1024 * 1024 for item in result.values()):
+ raise RuntimeError("benchmark transferred an unexpected byte count")
+ return result
+
+
+@app.post("/network/media-tunnel/sascha/benchmark")
+async def benchmark_sascha_media_paths(_=Depends(_verify)):
+ """Compare the legacy node6 route with the direct WireGuard media path using fixed transient TCP streams."""
+ try:
+ result = await asyncio.to_thread(_benchmark_sascha_media_paths)
+ except Exception as exc:
+ _audit("/network/media-tunnel/sascha/benchmark", "POST", 502, "benchmark failed")
+ raise HTTPException(502, str(exc)[-500:]) from exc
+ _audit("/network/media-tunnel/sascha/benchmark", "POST", 200, "fixed 256 MiB TCP comparison")
+ return {"status": "completed", "direction": "emby-sascha_to_hetzner", "results": result}
+
+
SYSCTL_AUDIT_KEYS = (
"net.core.default_qdisc",
"net.core.rmem_default",
@@ -2164,6 +3343,53 @@ async def _hetzner_zone_and_rrsets(zone_name: str):
return zone, rr_response.json().get("rrsets", []), headers
+@app.put("/dns/rrset/{zone_name}/{record_name}")
+async def dns_rrset_upsert(zone_name: str, record_name: str, req: DnsRrsetUpsertRequest, _=Depends(_verify)):
+ zone_name = zone_name.strip().lower().rstrip(".")
+ hostname = _validate_managed_hostname(f"{record_name}.{zone_name}")
+ record_name = hostname[: -(len(zone_name) + 1)]
+ record_type = req.record_type.upper()
+ value = req.value.strip().rstrip("." if record_type == "CNAME" else "")
+ try:
+ if record_type == "A" and ipaddress.ip_address(value).version != 4:
+ raise ValueError
+ if record_type == "AAAA" and ipaddress.ip_address(value).version != 6:
+ raise ValueError
+ if record_type == "CNAME":
+ _validate_managed_hostname(value)
+ except ValueError as exc:
+ raise HTTPException(400, f"Invalid {record_type} record value") from exc
+ zone, rrsets, headers = await _hetzner_zone_and_rrsets(zone_name)
+ existing = next((item for item in rrsets if item.get("name") == record_name and item.get("type") == record_type), None)
+ payload = {
+ "name": record_name,
+ "type": record_type,
+ "ttl": req.ttl,
+ "records": [{"value": value, "comment": req.comment}],
+ }
+ async with httpx.AsyncClient(timeout=30) as client:
+ if existing:
+ response = await client.put(
+ f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets/{record_name}/{record_type}",
+ headers=headers,
+ json=payload,
+ )
+ else:
+ response = await client.post(
+ f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets",
+ headers=headers,
+ json=payload,
+ )
+ if response.status_code not in {200, 201}:
+ raise HTTPException(502, f"Hetzner RRSet upsert failed: HTTP {response.status_code}")
+ verify = await client.get(f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets", headers=headers)
+ current = [item for item in verify.json().get("rrsets", []) if item.get("name") == record_name and item.get("type") == record_type] if verify.status_code == 200 else []
+ if not current or not any(record.get("value") == value for record in current[0].get("records", [])):
+ raise HTTPException(502, "RRSet read-back does not contain requested value")
+ _audit(f"/dns/rrset/{zone_name}/{record_name}", "PUT", 200, f"type={record_type} value={value}")
+ return {"zone": zone_name, "zone_id": zone.get("id"), "name": record_name, "created": existing is None, "rrset": current[0]}
+
+
@app.get("/dns/rrset/{zone_name}/{record_name}")
async def dns_rrset_get(zone_name: str, record_name: str, _=Depends(_verify)):
zone_name = zone_name.strip().lower().rstrip(".")
@@ -2829,6 +4055,205 @@ async def tts_health(_=Depends(_verify)):
return results
+def _sab_history_command() -> str:
+ container_script = r'''import json, re, subprocess, urllib.parse, urllib.request
+text = open('/config/sabnzbd.ini', encoding='utf-8', errors='replace').read()
+match = re.search(r'^api_key\s*=\s*(\S+)', text, re.M)
+port_match = re.search(r'^port\s*=\s*(\d+)', text, re.M)
+api_key = match.group(1) if match else ''
+port = port_match.group(1) if port_match else '7777'
+params = urllib.parse.urlencode({'mode': 'history', 'limit': 100, 'output': 'json', 'apikey': api_key})
+with urllib.request.urlopen('http://127.0.0.1:' + port + '/api?' + params, timeout=20) as response:
+ data = json.load(response)
+slots = data.get('history', {}).get('slots', [])
+allowed = ('nzo_id', 'name', 'category', 'status', 'script', 'script_line', 'fail_message', 'completed', 'storage', 'path')
+print(json.dumps([{key: item.get(key) for key in allowed} for item in slots]))
+'''
+ container_encoded = base64.b64encode(container_script.encode()).decode()
+ host_script = f'''import subprocess, sys
+command = ["sudo", "-n", "docker", "exec", "sabnzbd", "python3", "-c", "import base64;exec(base64.b64decode('{container_encoded}'))"]
+proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
+if proc.returncode != 0:
+ command = command[2:]
+ proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
+if proc.returncode != 0:
+ print((proc.stderr or proc.stdout)[-500:], file=sys.stderr)
+ raise SystemExit(proc.returncode)
+print(proc.stdout)
+'''
+ encoded = base64.b64encode(host_script.encode()).decode()
+ return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+@app.get("/media/handoff/sab-history")
+async def media_handoff_sab_history(_=Depends(_verify)):
+ """Return sanitized SAB history without exposing the SAB API key."""
+ inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
+ if not inventory:
+ raise HTTPException(404, "Host sabnzbd not found")
+ target = f'{inventory["user"]}@{inventory["ip"]}'
+ rc, out, err = await asyncio.to_thread(_ssh, target, _sab_history_command(), 40)
+ if rc != 0:
+ raise HTTPException(502, (err or out).strip()[-500:] or "SAB history failed")
+ try:
+ return {"history": json.loads(out)}
+ except json.JSONDecodeError as exc:
+ raise HTTPException(502, "invalid SAB history response") from exc
+
+
+def _media_handoff_diagnostics_command() -> str:
+ script = r'''import hashlib, json, os, subprocess, urllib.error, urllib.request
+
+
+def run(args):
+ proc = subprocess.run(args, capture_output=True, text=True, timeout=20)
+ return proc.returncode, proc.stdout.strip(), proc.stderr.strip()
+
+
+def docker_exec(command):
+ for prefix in (["sudo", "-n", "docker"], ["docker"]):
+ rc, out, err = run(prefix + ["exec", "sabnzbd", "sh", "-lc", command])
+ if rc == 0 or "not found" not in (err + out).lower():
+ return rc, out, err
+ return rc, out, err
+
+inspect_rc, inspect_out, _ = run(["sudo", "-n", "docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
+if inspect_rc != 0:
+ inspect_rc, inspect_out, _ = run(["docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
+py_rc, py_out, _ = docker_exec("command -v python3")
+curl_rc, curl_out, _ = docker_exec("command -v curl")
+stat_rc, stat_out, _ = docker_exec("test -f /usenet/scripts/movetdarr.sh && stat -c '%a %s' /usenet/scripts/movetdarr.sh && sha256sum /usenet/scripts/movetdarr.sh")
+log_rc, log_out, _ = docker_exec("tail -n 200 /usenet/scripts/postprocess.log 2>/dev/null || true")
+source_rc, source_out, _ = docker_exec("find /usenet/complete -mindepth 2 -maxdepth 2 -type d ! -name '_UNPACK_*' -print 2>/dev/null | sort | tail -100")
+target_rc, target_out, _ = docker_exec("find /tdarr/complete -mindepth 2 -maxdepth 2 -type d -print 2>/dev/null | sort | tail -100")
+probe_code = 0
+probe_body = ""
+probe = urllib.request.Request(
+ "http://10.5.85.2:8888/media/handoff",
+ method="POST",
+ headers={"Content-Type": "application/json"},
+ data=b'{"action":"status","jobId":"diagnostic-probe"}',
+)
+try:
+ with urllib.request.urlopen(probe, timeout=10) as response:
+ probe_code = response.status
+ probe_body = response.read(500).decode(errors="replace")
+except urllib.error.HTTPError as exc:
+ probe_code = exc.code
+ probe_body = exc.read(500).decode(errors="replace")
+except Exception as exc:
+ probe_body = type(exc).__name__
+script_info = {"exists": stat_rc == 0, "executable": False, "sha256": None, "mode": None, "size": None}
+if stat_rc == 0:
+ lines = stat_out.splitlines()
+ if lines:
+ parts = lines[0].split()
+ if len(parts) >= 2:
+ script_info.update(mode=parts[0], size=int(parts[1]), executable=any(ch in parts[0][-3:] for ch in "1357"))
+ if len(lines) > 1:
+ script_info["sha256"] = lines[1].split()[0]
+print(json.dumps({
+ "container_running": inspect_rc == 0 and inspect_out == "true",
+ "tools": {"python3": py_rc == 0 and bool(py_out), "curl": curl_rc == 0 and bool(curl_out)},
+ "script": script_info,
+ "caller_probe": {"http_status": probe_code, "body": probe_body},
+ "source_directories": source_out.splitlines() if source_rc == 0 else [],
+ "target_directories": target_out.splitlines() if target_rc == 0 else [],
+ "recent_log": log_out.splitlines()[-200:],
+}))
+'''
+ encoded = base64.b64encode(script.encode()).decode()
+ return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
+
+
+@app.get("/media/handoff/diagnostics")
+async def media_handoff_diagnostics(_=Depends(_verify)):
+ """Read-only diagnostics for the fixed SABnzbd handoff script and caller path."""
+ inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
+ if not inventory:
+ raise HTTPException(404, "Host sabnzbd not found")
+ target = f'{inventory["user"]}@{inventory["ip"]}'
+ rc, out, err = await asyncio.to_thread(_ssh, target, _media_handoff_diagnostics_command(), 30)
+ if rc != 0:
+ raise HTTPException(502, (err or out).strip()[-500:] or "media handoff diagnostics failed")
+ try:
+ return json.loads(out)
+ except json.JSONDecodeError as exc:
+ raise HTTPException(502, "invalid media handoff diagnostic response") from exc
+
+
+class MediaHandoffPayload(BaseModel):
+ action: Literal["start", "moved", "status", "fail"]
+ category: str | None = Field(None, max_length=32)
+ directory: str | None = Field(None, max_length=1000)
+ release: str | None = Field(None, max_length=500)
+ cleanName: str | None = Field(None, max_length=500)
+ expectedFiles: int | None = Field(None, ge=1, le=100)
+ jobId: str | None = Field(None, max_length=80)
+ reason: str | None = Field(None, max_length=300)
+
+
+def _media_handoff_caller_allowed(request: Request) -> bool:
+ auth = request.headers.get("authorization", "")
+ if BUTLER_TOKEN and secrets.compare_digest(auth, f"Bearer {BUTLER_TOKEN}"):
+ return True
+ try:
+ caller = ipaddress.ip_address(request.client.host if request.client else "")
+ networks = [
+ ipaddress.ip_network(value.strip(), strict=False)
+ for value in MEDIA_HANDOFF_ALLOWED_NETWORKS.split(",")
+ if value.strip()
+ ]
+ except ValueError:
+ return False
+ return any(caller in network for network in networks)
+
+
+@app.post("/media/handoff")
+async def media_handoff(payload: MediaHandoffPayload, request: Request):
+ """Narrow SABnzbd-to-n8n bridge; no generic unauthenticated proxy access."""
+ if not _media_handoff_caller_allowed(request):
+ caller = request.client.host if request.client else "unknown"
+ _audit("/media/handoff", "POST", 403, f"caller={caller} action={payload.action}")
+ raise HTTPException(403, "Media handoff caller is not allowed")
+
+ # 16.09.2026: serienen/videoen ergaenzt. Die englischen Arr-Instanzen
+ # sonarrEN (Port 8991, Root /data/FHD/serienen) und radarrEN (Port 7880,
+ # Root /data/FHD/videoen) nutzen eigene SAB-Kategorien. Ohne sie brach der
+ # Handoff mit 422 ab und Releases blieben in /usenet/complete liegen.
+ allowed_categories = {"serien4k", "serien", "serienen", "video4k", "video", "videoen"}
+ if payload.action == "start":
+ if payload.category not in allowed_categories:
+ raise HTTPException(422, "Unsupported media category")
+ expected_prefix = f"/usenet/complete/{payload.category}/"
+ if not payload.directory or not payload.directory.startswith(expected_prefix):
+ raise HTTPException(422, "Invalid media handoff directory")
+ if not payload.release:
+ raise HTTPException(422, "Release name is required")
+ else:
+ if not payload.jobId or not re.fullmatch(r"[a-z0-9-]{8,80}", payload.jobId):
+ raise HTTPException(422, "Valid jobId is required")
+
+ cfg = SERVICES.get("n8n")
+ if not cfg or not cfg.get("url"):
+ raise HTTPException(503, "n8n service is not configured")
+ target = f"{cfg['url'].rstrip('/')}/webhook/media-handoff"
+ body = payload.model_dump(exclude_none=True) if hasattr(payload, "model_dump") else payload.dict(exclude_none=True)
+ try:
+ async with httpx.AsyncClient(verify=False, timeout=15) as client:
+ response = await client.post(target, json=body, headers={"Content-Type": "application/json"})
+ except httpx.HTTPError as exc:
+ _audit("/media/handoff", "POST", 502, f"action={payload.action} error={type(exc).__name__}")
+ raise HTTPException(502, "n8n media handoff is unavailable") from exc
+
+ try:
+ result = response.json()
+ except Exception:
+ result = {"ok": False, "error": "invalid_n8n_response"}
+ _audit("/media/handoff", "POST", response.status_code, f"action={payload.action}")
+ return JSONResponse(content=result, status_code=response.status_code)
+
+
@app.api_route("/{service}/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"])
async def proxy(service: str, path: str, request: Request, _=Depends(_verify)):
SKIP_SERVICES = {"vm", "inventory", "ansible", "debug", "tts", "status", "audit", "config"}
@@ -2838,7 +4263,7 @@ async def proxy(service: str, path: str, request: Request, _=Depends(_verify)):
if not cfg:
raise HTTPException(404, f"Unknown service: {service}. Available: {list(SERVICES.keys())}")
- base_url = cfg["url"]
+ base_url = cfg.get("url")
auth_type = cfg["auth"]
headers = dict(request.headers)
cookies = {}
diff --git a/tests/test_app.py b/tests/test_app.py
index 82356fd..45011c2 100644
--- a/tests/test_app.py
+++ b/tests/test_app.py
@@ -1,6 +1,7 @@
import os
import asyncio
import json
+import shlex
import time
from datetime import datetime, timedelta, timezone
@@ -43,7 +44,207 @@ def test_health_exposes_current_version():
with TestClient(app.app) as client:
response = client.get("/health")
assert response.status_code == 200
- assert response.json()["version"] == app.VERSION == "2.3.5"
+ assert response.json()["version"] == app.VERSION == "2.4.1"
+
+
+def test_media_handoff_proxies_strict_category_contract(monkeypatch):
+ captured = {}
+
+ class FakeResponse:
+ status_code = 200
+
+ def json(self):
+ return {"ok": True, "jobId": "test-job-1234", "state": "registered"}
+
+ class FakeClient:
+ def __init__(self, **kwargs):
+ captured["client"] = kwargs
+
+ async def __aenter__(self):
+ return self
+
+ async def __aexit__(self, *_args):
+ return None
+
+ async def post(self, url, json, headers):
+ captured.update(url=url, json=json, headers=headers)
+ return FakeResponse()
+
+ monkeypatch.setattr(app.httpx, "AsyncClient", FakeClient)
+ with TestClient(app.app) as client:
+ monkeypatch.setattr(app, "SERVICES", {"n8n": {"url": "http://n8n:5678", "auth": "n8n"}})
+ response = client.post(
+ "/media/handoff",
+ headers={"Authorization": "Bearer test-token"},
+ json={
+ "action": "start",
+ "category": "serien4k",
+ "directory": "/usenet/complete/serien4k/Show.S01E01",
+ "release": "Show.S01E01-GRP",
+ "cleanName": "Show S01E01",
+ "expectedFiles": 2,
+ },
+ )
+ assert response.status_code == 200
+ assert response.json()["state"] == "registered"
+ assert captured["url"] == "http://n8n:5678/webhook/media-handoff"
+ assert captured["json"]["category"] == "serien4k"
+ assert captured["json"]["expectedFiles"] == 2
+
+
+def test_media_handoff_rejects_untrusted_caller_before_proxy(monkeypatch):
+ monkeypatch.setattr(app.httpx, "AsyncClient", lambda **_kwargs: (_ for _ in ()).throw(AssertionError("must not proxy")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/handoff",
+ json={"action": "status", "jobId": "test-job-1234"},
+ )
+ assert response.status_code == 403
+
+
+def test_media_handoff_rejects_wrong_category_path(monkeypatch):
+ monkeypatch.setattr(app.httpx, "AsyncClient", lambda **_kwargs: (_ for _ in ()).throw(AssertionError("must not proxy")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/handoff",
+ headers={"Authorization": "Bearer test-token"},
+ json={
+ "action": "start",
+ "category": "video4k",
+ "directory": "/usenet/complete/serien4k/Wrong",
+ "release": "Wrong",
+ },
+ )
+ assert response.status_code == 422
+
+
+def test_media_handoff_accepts_english_arr_categories(monkeypatch):
+ """serienen/videoen muessen durchgehen (sonarrEN 8991 / radarrEN 7880).
+
+ 16.09.2026: fehlten in allowed_categories -> HTTP 422 -> movetdarr.sh
+ brach mit "n8n-Handoff konnte nicht registriert werden" ab und liess
+ Mutiny.2026 x2 tagelang in /usenet/complete/videoen liegen.
+ """
+ seen = []
+
+ class FakeResponse:
+ status_code = 200
+
+ def json(self):
+ return {"ok": True, "jobId": "test-job-5678", "state": "registered"}
+
+ class FakeClient:
+ def __init__(self, **_kwargs):
+ pass
+
+ async def __aenter__(self):
+ return self
+
+ async def __aexit__(self, *_args):
+ return None
+
+ async def post(self, url, json, headers):
+ seen.append(json)
+ return FakeResponse()
+
+ monkeypatch.setattr(app.httpx, "AsyncClient", FakeClient)
+ with TestClient(app.app) as client:
+ monkeypatch.setattr(app, "SERVICES", {"n8n": {"url": "http://n8n:5678", "auth": "n8n"}})
+ for category in ("serienen", "videoen"):
+ response = client.post(
+ "/media/handoff",
+ headers={"Authorization": "Bearer test-token"},
+ json={
+ "action": "start",
+ "category": category,
+ "directory": f"/usenet/complete/{category}/Release.2026",
+ "release": "Release.2026-GRP",
+ "cleanName": "Release 2026",
+ "expectedFiles": 1,
+ },
+ )
+ assert response.status_code == 200, (category, response.text)
+ assert response.json()["state"] == "registered", category
+
+ assert [item["category"] for item in seen] == ["serienen", "videoen"]
+
+
+def test_media_handoff_still_rejects_unknown_category(monkeypatch):
+ """Fail-closed bleibt: eine frei erfundene Kategorie wird nicht geproxyt."""
+ monkeypatch.setattr(
+ app.httpx,
+ "AsyncClient",
+ lambda **_kwargs: (_ for _ in ()).throw(AssertionError("must not proxy")),
+ )
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/handoff",
+ headers={"Authorization": "Bearer test-token"},
+ json={
+ "action": "start",
+ "category": "hoerbuecher",
+ "directory": "/usenet/complete/hoerbuecher/Buch",
+ "release": "Buch",
+ },
+ )
+ assert response.status_code == 422
+
+
+def test_media_handoff_english_category_path_must_match(monkeypatch):
+ """Pfadpruefung gilt auch fuer die neuen Kategorien."""
+ monkeypatch.setattr(
+ app.httpx,
+ "AsyncClient",
+ lambda **_kwargs: (_ for _ in ()).throw(AssertionError("must not proxy")),
+ )
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/handoff",
+ headers={"Authorization": "Bearer test-token"},
+ json={
+ "action": "start",
+ "category": "videoen",
+ "directory": "/usenet/complete/video4k/Wrong",
+ "release": "Wrong",
+ },
+ )
+ assert response.status_code == 422
+
+
+def test_paperless_import_queues_pdf_through_butler(monkeypatch):
+ captured = {}
+
+ def fake_upload(host, command, payload, timeout):
+ captured.update(host=host, command=command, payload=payload, timeout=timeout)
+ return 0, "", ""
+
+ monkeypatch.setattr(app, "_ssh_bytes", fake_upload)
+ pdf = b"%PDF-1.7\nvalid-test-payload"
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/paperless/import?filename=Hausordnung%20Stand%2009.05.2022.pdf",
+ content=pdf,
+ headers={"Authorization": "Bearer test-token", "Content-Type": "application/pdf"},
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["status"] == "queued"
+ assert body["filename"].startswith("Hausordnung_Stand_09.05.2022-")
+ assert captured["host"] == app.PAPERLESS_GATEWAY
+ assert app.PAPERLESS_SSH in captured["command"]
+ assert captured["payload"] == pdf
+ assert app.PAPERLESS_CONSUME_DIR in captured["command"]
+
+
+def test_paperless_import_rejects_non_pdf(monkeypatch):
+ monkeypatch.setattr(app, "_ssh_bytes", lambda *args: (_ for _ in ()).throw(AssertionError("must not upload")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/paperless/import?filename=bad.pdf",
+ content=b"not a pdf",
+ headers={"Authorization": "Bearer test-token"},
+ )
+ assert response.status_code == 400
def test_ui_serves_self_contained_operator_console():
@@ -54,6 +255,11 @@ def test_ui_serves_self_contained_operator_console():
assert 'id="operations-grid"' in response.text
assert 'id="doctor-form"' in response.text
assert 'id="preflight-form"' in response.text
+ assert 'id="sharing-form"' in response.text
+ assert 'id="sharing-events"' in response.text
+ assert "/emby/account-sharing" in response.text
+ assert 'id="login"' not in response.text
+ assert "Butler-Token" not in response.text
assert "localStorage" not in response.text
@@ -86,6 +292,185 @@ def test_ui_session_login_uses_httponly_cookie_and_csrf():
assert allowed.status_code == 200
+def test_ui_anonymous_session_is_automatic_and_strictly_read_only():
+ app._ui_sessions.clear()
+ with TestClient(app.app) as client:
+ session = client.get("/ui/session")
+ assert session.status_code == 200
+ assert session.json()["authenticated"] is True
+ assert session.json()["read_only"] is True
+ assert "HttpOnly" in session.headers.get("set-cookie", "")
+
+ capabilities = client.get("/capabilities")
+ assert capabilities.status_code == 200
+
+ csrf = client.cookies.get("butler_csrf")
+ mutation = client.post("/config/reload", headers={"X-CSRF-Token": csrf})
+ assert mutation.status_code == 403
+ assert "read-only" in mutation.json()["detail"].lower()
+
+
+def test_emby_network_identity_normalizes_ipv4_and_ipv6_privacy_addresses():
+ ipv4 = app._emby_network_identity("203.0.113.9:443")
+ assert ipv4["ip"] == "203.0.113.9"
+ assert ipv4["network"] == "203.0.113.9/32"
+ assert ipv4["identity"] == "203.0.113.9/32"
+
+ first = app._emby_network_identity("2003:abcd:1234:5678::1")
+ privacy_peer = app._emby_network_identity("[2003:abcd:1234:5678:ffff::99]:443")
+ sibling_subnet = app._emby_network_identity("2003:abcd:1234:9999::1")
+ assert first["network"] == privacy_peer["network"] == "2003:abcd:1234:5678::/64"
+ assert first["parent"] == sibling_subnet["parent"] == "2003:abcd:1234::/48"
+ assert first["identity"] == sibling_subnet["identity"] == "2003:abcd:1234::/48"
+
+
+def test_emby_sharing_flags_concurrent_distinct_networks_but_not_sibling_ipv6_subnets():
+ def series(endpoint, values, city):
+ return {
+ "metric": {
+ "job": "emby-sascha", "username": "Alice", "remoteEndPoint": endpoint,
+ "city": city, "region": "Test", "countryCode": "DE",
+ "latitude": "48.1", "longitude": "11.5",
+ },
+ "values": [[timestamp, "1"] for timestamp in values],
+ }
+
+ payload = [
+ series("2606:4700:1234:1000::1", [100, 160, 220, 280, 340, 400], "Home"),
+ series("2606:4700:1234:2000::2", [100, 160, 220, 280, 340, 400], "Home privacy subnet"),
+ series("2001:4860:9999:1000::1", [100, 160, 220, 280, 340, 400], "Away"),
+ ]
+
+ result = app._analyze_emby_sharing(payload, step_seconds=60)
+
+ assert result["summary"]["concurrent_events"] == 1
+ event = result["events"][0]
+ assert event["type"] == "concurrent_networks"
+ assert event["username"] == "Alice"
+ assert len(event["evidence"]) == 2
+ assert {item["identity"] for item in event["evidence"]} == {
+ "2606:4700:1234::/48", "2001:4860:9999::/48"
+ }
+
+
+def test_emby_sharing_ignores_short_overlap_inside_prometheus_staleness_window():
+ def series(endpoint):
+ return {"metric": {"job": "emby-sascha", "username": "Alice", "remoteEndPoint": endpoint,
+ "city": "Munich", "region": "Bavaria", "countryCode": "DE",
+ "latitude": "48.1", "longitude": "11.5"},
+ "values": [[100, "1"], [160, "1"]]}
+
+ result = app._analyze_emby_sharing([series("8.8.8.8"), series("1.1.1.1")], step_seconds=60)
+
+ assert result["summary"]["concurrent_events"] == 0
+
+
+def test_emby_sharing_flags_geographically_impossible_network_change():
+ payload = [
+ {
+ "metric": {"job": "emby-chris", "username": "Bob", "remoteEndPoint": "8.8.8.8",
+ "city": "Berlin", "region": "Berlin", "countryCode": "DE",
+ "latitude": "52.5200", "longitude": "13.4050"},
+ "values": [[100, "1"], [160, "1"]],
+ },
+ {
+ "metric": {"job": "emby-chris", "username": "Bob", "remoteEndPoint": "1.1.1.1",
+ "city": "New York", "region": "New York", "countryCode": "US",
+ "latitude": "40.7128", "longitude": "-74.0060"},
+ "values": [[400, "1"], [460, "1"]],
+ },
+ ]
+
+ result = app._analyze_emby_sharing(payload, step_seconds=60)
+
+ assert result["summary"]["concurrent_events"] == 0
+ assert result["summary"]["impossible_travel_events"] == 1
+ event = result["events"][0]
+ assert event["type"] == "impossible_travel"
+ assert event["distance_km"] > 6000
+ assert event["required_speed_kmh"] > 1000
+
+
+def test_emby_account_sharing_endpoint_is_read_only_and_filterable(monkeypatch):
+ async def history(days, server):
+ assert days == 7
+ assert server == "all"
+ return ([{
+ "metric": {"job": "emby-sascha", "username": "Alice", "remoteEndPoint": "8.8.8.8",
+ "city": "Munich", "region": "Bavaria", "countryCode": "DE",
+ "latitude": "48.1", "longitude": "11.5"},
+ "values": [[100, "1"], [160, "1"]],
+ }], 60)
+
+ monkeypatch.setattr(app, "_fetch_emby_session_history", history)
+ with TestClient(app.app) as client:
+ response = client.get(
+ "/emby/account-sharing?days=7&username=alice",
+ headers={"Authorization": "Bearer test-token"},
+ )
+
+ assert response.status_code == 200
+ payload = response.json()
+ assert payload["policy"]["mode"] == "conservative"
+ assert payload["policy"]["ipv6_detection_identity"] == "/48"
+ assert payload["summary"]["users_analyzed"] == 1
+ 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(
@@ -101,6 +486,9 @@ def test_capabilities_is_live_machine_readable_safety_map():
assert removal["mode"] == "mutation"
assert removal["dry_run"] is True
assert removal["critical"] is True
+ tunnel = by_operation[("POST", "/network/media-tunnel/sascha")]
+ assert tunnel["dry_run"] is True
+ assert tunnel["critical"] is True
assert by_operation[("DELETE", "/vm/destroy/{vmid}")]["dry_run"] is True
assert all(item["path"] != "/{service}/{path}" for item in payload["operations"])
assert payload["model_contract"]["instruction"].startswith("Prefer read_only")
@@ -136,6 +524,39 @@ def test_audit_persists_and_redacts_secrets(tmp_path, monkeypatch):
assert entry["detail"].count("[REDACTED]") == 3
+def test_audit_records_ai_and_web_ui_actor(tmp_path, monkeypatch):
+ db = tmp_path / "audit.sqlite3"
+ monkeypatch.setattr(app, "AUDIT_DB_PATH", str(db))
+
+ async def empty_status():
+ return {}
+
+ monkeypatch.setattr(app, "_collect_service_status", empty_status)
+ app._ui_sessions.clear()
+ with TestClient(app.app) as client:
+ ai = client.get(
+ "/status",
+ headers={"Authorization": "Bearer test-token", "X-Butler-Actor": "Trulla"},
+ )
+ assert ai.status_code == 200
+
+ client.headers.pop("Authorization", None)
+ client.get("/ui/session")
+ web = client.get("/status")
+ assert web.status_code == 200
+
+ entries = client.get("/audit").json()
+
+ assert [item["actor"] for item in entries[:2]] == ["Weboberfläche", "Trulla"]
+
+
+def test_ui_audit_has_actor_column_and_filter():
+ with TestClient(app.app) as client:
+ response = client.get("/ui")
+ assert '
Akteur | ' in response.text
+ assert 'id="audit-actor"' in response.text
+
+
def test_wireguard_status_returns_redacted_live_state(monkeypatch):
payload = {
"interface": "wg0",
@@ -481,6 +902,115 @@ def test_uptime_monitor_remove_rejects_invalid_expected_name_before_ssh(monkeypa
assert response.status_code == 400
+def test_media_verify_returns_duration_and_not_suspect_when_within_tolerance(monkeypatch):
+ monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
+ calls = []
+
+ def fake_ssh(host, command, timeout=30):
+ calls.append((host, command, timeout))
+ return 0, "2967.355000\n", ""
+
+ monkeypatch.setattr(app, "_ssh", fake_ssh)
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/verify",
+ headers={"Authorization": "Bearer test-token"},
+ json={"host": "arrapps", "path": "/data/UHD/serien/Show/ep.mkv", "expected_minutes": 49.5},
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["duration_minutes"] == 49.46
+ assert body["suspect"] is False
+ assert calls[0][0] == "sascha@10.2.1.100"
+ assert "bazarrUHD" in calls[0][1]
+ assert "ffprobe" in calls[0][1]
+
+
+def test_media_verify_flags_truncated_file_as_suspect(monkeypatch):
+ monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
+ monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (0, "1836.0\n", ""))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/verify",
+ headers={"Authorization": "Bearer test-token"},
+ json={"host": "arrapps", "path": "/data/UHD/serien/Show/ep.mkv", "expected_minutes": 49.5},
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["suspect"] is True
+ assert body["deviation_pct"] > 20
+
+
+def test_media_verify_rejects_path_outside_data_before_ssh(monkeypatch):
+ monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/verify",
+ headers={"Authorization": "Bearer test-token"},
+ json={"host": "arrapps", "path": "/etc/passwd"},
+ )
+ assert response.status_code == 400
+
+
+def test_media_verify_allows_apostrophe_in_filename_and_quotes_it_safely(monkeypatch):
+ monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
+ calls = []
+
+ def fake_ssh(host, command, timeout=30):
+ calls.append((host, command, timeout))
+ return 0, "2967.0\n", ""
+
+ monkeypatch.setattr(app, "_ssh", fake_ssh)
+ path = "/data/FHD/serien/Star Trek - Strange New Worlds (2022)/Season 04/Once La'An a Time.mkv"
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/verify",
+ headers={"Authorization": "Bearer test-token"},
+ json={"host": "arrapps", "path": path, "expected_minutes": 49.5},
+ )
+ assert response.status_code == 200
+ # shlex.quote must produce a command the remote shell parses as ONE argument,
+ # i.e. no unescaped apostrophe breaks out of quoting.
+ executed_cmd = calls[0][1]
+ quoted = shlex.quote(path)
+ assert quoted in executed_cmd
+ assert shlex.split(executed_cmd)[-1] == path
+
+
+def test_media_verify_rejects_shell_metacharacters_before_ssh(monkeypatch):
+ monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/verify",
+ headers={"Authorization": "Bearer test-token"},
+ json={"host": "arrapps", "path": "/data/UHD/serien/a\x00.mkv"},
+ )
+ assert response.status_code == 400
+
+
+def test_media_verify_rejects_unknown_host_before_ssh(monkeypatch):
+ monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/verify",
+ headers={"Authorization": "Bearer test-token"},
+ json={"host": "unknown-host", "path": "/data/UHD/serien/ep.mkv"},
+ )
+ assert response.status_code == 400
+
+
+def test_media_verify_returns_502_on_unparsable_ffprobe_output(monkeypatch):
+ monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
+ monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (0, "N/A\n", ""))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/verify",
+ headers={"Authorization": "Bearer test-token"},
+ json={"host": "arrapps", "path": "/data/UHD/serien/ep.mkv", "expected_minutes": 49.5},
+ )
+ assert response.status_code == 502
+
+
def test_invalid_log_target_is_rejected_before_ssh():
with TestClient(app.app) as client:
response = client.get(
@@ -950,3 +1480,124 @@ def test_speedtest_deploy_requires_strong_secrets_and_uses_full_git_app(monkeypa
assert deploy[1]["compose.yaml"] == "content:compose.yaml"
assert deploy[2] == "correct-horse-battery-staple"
assert deploy[3] == "streamscope-session-secret-with-entropy"
+
+
+def test_sascha_media_edge_audit_uses_fixed_hetzner_host(monkeypatch):
+ payload = {
+ "hostname": "pfannkuchen",
+ "caddy": {"container_running": True, "protocols": ["h1", "h2", "h3"], "upstreams": ["10.6.1.103:8096"]},
+ "network": {"wg_media": False, "udp_51821": False},
+ }
+ calls = []
+ monkeypatch.setattr(app, "_find_inventory_host", lambda name: {"name": name, "user": "root", "ip": "46.225.230.72"})
+ monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=30: (calls.append((host, command, timeout)) or (0, json.dumps(payload), "")))
+ with TestClient(app.app) as client:
+ response = client.get("/media/edge/sascha", headers={"Authorization": "Bearer test-token"})
+ assert response.status_code == 200
+ assert response.json()["caddy"]["upstreams"] == ["10.6.1.103:8096"]
+ assert calls[0][0] == "root@46.225.230.72"
+ assert "PrivateKey" not in response.text
+
+
+def test_sascha_media_tunnel_defaults_to_side_effect_free_dry_run(monkeypatch):
+ monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("dry-run must not SSH")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/network/media-tunnel/sascha",
+ headers={"Authorization": "Bearer test-token"},
+ json={},
+ )
+ assert response.status_code == 200
+ body = response.json()
+ assert body["status"] == "would_deploy"
+ assert body["vps_address"] == "10.11.13.1/32"
+ assert body["emby_address"] == "10.11.13.3/32"
+ assert body["listen_port"] == 51821
+
+
+def test_sascha_media_tunnel_requires_explicit_confirmation(monkeypatch):
+ monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("unconfirmed request must not SSH")))
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/network/media-tunnel/sascha",
+ headers={"Authorization": "Bearer test-token"},
+ json={"dry_run": False},
+ )
+ assert response.status_code == 400
+
+
+def test_sascha_media_tunnel_apply_returns_redacted_result(monkeypatch):
+ result = {
+ "status": "deployed",
+ "interface": "wg-media",
+ "vps_address": "10.11.13.1/32",
+ "emby_address": "10.11.13.3/32",
+ "listen_port": 51821,
+ "handshake": True,
+ "ping_vps_to_emby": True,
+ "ping_emby_to_vps": True,
+ }
+ monkeypatch.setattr(app, "_deploy_sascha_media_tunnel", lambda: result)
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/network/media-tunnel/sascha",
+ headers={"Authorization": "Bearer test-token"},
+ json={"dry_run": False, "confirmation": "DEPLOY_DIRECT_SASCHA_MEDIA_TUNNEL"},
+ )
+ assert response.status_code == 200
+ assert response.json()["handshake"] is True
+ assert "private" not in response.text.lower()
+
+
+def test_sascha_media_tunnel_candidate_uses_valid_wireguard_interface_name():
+ import base64
+ import re
+
+ command = app._media_tunnel_install_command("vps", "A" * 43 + "=")
+ encoded = re.search(r"b64decode\('([^']+)'\)", command).group(1)
+ script = base64.b64decode(encoded).decode()
+ assert 'root / "wgmtest.conf"' in script
+ assert len("wgmtest") <= 15
+
+
+def test_sascha_media_path_benchmark_returns_both_fixed_routes(monkeypatch):
+ monkeypatch.setattr(app, "_benchmark_sascha_media_paths", lambda: {
+ "legacy_node6": {"bytes": 268435456, "seconds": 4.0, "mbit_s": 536.9},
+ "direct_wg_media": {"bytes": 268435456, "seconds": 3.0, "mbit_s": 715.8},
+ })
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/network/media-tunnel/sascha/benchmark",
+ headers={"Authorization": "Bearer test-token"},
+ )
+ assert response.status_code == 200
+ assert response.json()["direction"] == "emby-sascha_to_hetzner"
+ assert set(response.json()["results"]) == {"legacy_node6", "direct_wg_media"}
+
+
+def test_sascha_media_edge_optimize_is_dry_run_by_default(monkeypatch):
+ monkeypatch.setattr(app, "_optimize_sascha_media_edge", lambda: (_ for _ in ()).throw(AssertionError("dry-run must not mutate")))
+ with TestClient(app.app) as client:
+ response = client.post("/media/edge/sascha/optimize", headers={"Authorization": "Bearer test-token"}, json={})
+ assert response.status_code == 200
+ assert response.json()["protocols"] == ["h1", "h2"]
+ assert response.json()["upstream"] == "10.11.13.3:8096"
+
+
+def test_sascha_media_edge_optimize_requires_confirmation(monkeypatch):
+ monkeypatch.setattr(app, "_optimize_sascha_media_edge", lambda: (_ for _ in ()).throw(AssertionError("must not mutate")))
+ with TestClient(app.app) as client:
+ response = client.post("/media/edge/sascha/optimize", headers={"Authorization": "Bearer test-token"}, json={"dry_run": False})
+ assert response.status_code == 400
+
+
+def test_sascha_media_edge_optimize_applies_fixed_safe_result(monkeypatch):
+ monkeypatch.setattr(app, "_optimize_sascha_media_edge", lambda: {"status": "optimized", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "backup": "/app-config/caddy/backup", "sha256": "a" * 64})
+ with TestClient(app.app) as client:
+ response = client.post(
+ "/media/edge/sascha/optimize",
+ headers={"Authorization": "Bearer test-token"},
+ json={"dry_run": False, "confirmation": "OPTIMIZE_TV_SASCHA_LUTZ_DE"},
+ )
+ assert response.status_code == 200
+ assert response.json()["status"] == "optimized"
diff --git a/tests/test_bw_manager_deploy.py b/tests/test_bw_manager_deploy.py
index cf5a339..670435b 100644
--- a/tests/test_bw_manager_deploy.py
+++ b/tests/test_bw_manager_deploy.py
@@ -18,6 +18,8 @@ def load_app(monkeypatch):
def test_bw_manager_deploy_rejects_bundle_without_and_gate(monkeypatch):
app = load_app(monkeypatch)
+ if not hasattr(app, "BW_MANAGER_REPO_FILES"):
+ pytest.skip("BW-Manager was intentionally retired on 13.08.2026")
files = {path: "placeholder" for path in app.BW_MANAGER_REPO_FILES}
files["compose.yaml"] = "services:\n bw-manager:\n build: ./src\n"
files["src/app.py"] = "def old_policy(): pass\n"
diff --git a/tests/test_guck_admin_deploy.py b/tests/test_guck_admin_deploy.py
new file mode 100644
index 0000000..74d1f79
--- /dev/null
+++ b/tests/test_guck_admin_deploy.py
@@ -0,0 +1,54 @@
+import app
+import pytest
+from fastapi.testclient import TestClient
+
+
+def valid_bundle():
+ files={path:'placeholder' for path in app.GUCK_ADMIN_REPO_FILES}
+ files['guck-admin/compose.yaml']='''services:\n guck-admin:\n build: .\n network_mode: host\n cap_add: [NET_ADMIN]\n environment:\n GUCK_LIMIT: /host/guck-limit.sh\n volumes:\n - /app-config/guck-admin/data:/data\n control-monitor:\n build: .\n network_mode: host\n cap_add: [NET_ADMIN]\n'''
+ files['guck-admin/src/control.py']='''CREATE TABLE IF NOT EXISTS custom_networks\n2a00:8c40:f000::/36\n45.58.235.0/24\ndef sync_custom_networks(): pass\n'''
+ return files
+
+
+def test_guck_admin_bundle_requires_persistent_custom_network_policy():
+ files=valid_bundle()
+ files['guck-admin/src/control.py']='def old_control(): pass\n'
+ with pytest.raises(ValueError,match='custom VPN'):
+ app._validate_guck_admin_bundle(files)
+
+
+def test_guck_admin_deploy_dry_run_has_no_remote_side_effect(monkeypatch):
+ calls=[]
+ monkeypatch.setattr(app,'_find_inventory_host',lambda host:{'user':'debian','ip':'141.94.237.199'})
+ monkeypatch.setattr(app,'_ssh',lambda *args,**kwargs:calls.append(args))
+ result=app._deploy_guck_admin_compose(valid_bundle(),dry_run=True)
+ assert result['status']=='validated'
+ assert result['host']=='guck-vps'
+ assert calls==[]
+
+
+def test_guck_admin_deploy_uses_inventory_target_and_verifies_policy(monkeypatch):
+ calls=[]
+ monkeypatch.setattr(app,'_find_inventory_host',lambda host:{'user':'debian','ip':'141.94.237.199'})
+ def fake_ssh(host,command,timeout=600):
+ calls.append((host,command,timeout))
+ if 'ipset test vpn-v6' in command:
+ return 0,'policy ok',''
+ return 0,'ok',''
+ monkeypatch.setattr(app,'_ssh',fake_ssh)
+ result=app._deploy_guck_admin_compose(valid_bundle(),dry_run=False)
+ assert result['status']=='deployed'
+ assert result['policy']=='verified'
+ assert all(host=='debian@141.94.237.199' for host,_,_ in calls)
+ commands='\n'.join(command for _,command,_ in calls)
+ assert 'docker compose build' in commands
+ assert '127.0.0.1:9090/health' in commands
+ assert '/actions/limiter/refresh' in commands
+ assert 'ipset test vpn-v6 2a00:8c40:f02d:a34c::1' in commands
+ assert 'ipset test vpn-v4 45.58.235.7' in commands
+
+
+def test_guck_admin_deploy_endpoint_requires_auth(monkeypatch):
+ monkeypatch.setattr(app,'BUTLER_TOKEN','test-token')
+ response=TestClient(app.app).post('/vps/guck-admin/deploy',json={'dry_run':True})
+ assert response.status_code in (401,403)
diff --git a/ui.html b/ui.html
index 9ad3d27..7b8c45b 100644
--- a/ui.html
+++ b/ui.html
@@ -11,47 +11,33 @@
.shell{display:grid;grid-template-columns:248px minmax(0,1fr);min-height:100vh}.sidebar{position:sticky;top:0;height:100vh;padding:22px 16px;border-right:1px solid var(--border2);background:rgba(15,16,17,.92);backdrop-filter:blur(18px);z-index:20}.brand{display:flex;align-items:center;gap:11px;padding:0 8px 24px}.logo{width:34px;height:34px;border:1px solid rgba(113,112,255,.35);border-radius:10px;display:grid;place-items:center;background:linear-gradient(145deg,rgba(113,112,255,.22),rgba(255,255,255,.03));box-shadow:inset 0 0 16px rgba(113,112,255,.12)}.brand strong{font-size:14px;font-weight:590}.brand small{display:block;color:var(--muted);font-size:11px;margin-top:2px}.nav-label{font-size:10px;text-transform:uppercase;letter-spacing:.1em;color:var(--dim);padding:15px 10px 7px}.nav button{width:100%;display:flex;align-items:center;gap:10px;border:0;background:transparent;color:var(--muted);padding:9px 10px;border-radius:7px;text-align:left;font-size:13px;font-weight:510}.nav button:hover,.nav button.active{background:rgba(255,255,255,.05);color:var(--text)}.nav .icon{width:18px;text-align:center;color:var(--dim)}.sidebar-foot{position:absolute;left:16px;right:16px;bottom:18px}.session-pill{display:flex;align-items:center;justify-content:space-between;border:1px solid var(--border);border-radius:8px;padding:9px 10px;color:var(--muted);font-size:11px;background:rgba(255,255,255,.02)}.dot{width:7px;height:7px;border-radius:50%;background:var(--green);box-shadow:0 0 10px rgba(16,185,129,.65)}
main{min-width:0}.topbar{height:68px;position:sticky;top:0;z-index:15;display:flex;align-items:center;justify-content:space-between;padding:0 30px;border-bottom:1px solid var(--border2);background:rgba(8,9,10,.78);backdrop-filter:blur(18px)}.topbar h1{font-size:15px;margin:0;font-weight:510}.top-actions{display:flex;gap:8px}.btn{border:1px solid var(--border);background:rgba(255,255,255,.03);color:var(--secondary);border-radius:7px;padding:8px 12px;font-size:12px;font-weight:510}.btn:hover{background:rgba(255,255,255,.07);color:var(--text)}.btn.primary{background:var(--accent2);border-color:transparent;color:white}.btn.danger{color:#fca5a5}.mobile-menu{display:none}.content{max-width:1440px;margin:0 auto;padding:30px}.view{display:none}.view.active{display:block}.eyebrow{font-size:11px;color:var(--accent);text-transform:uppercase;letter-spacing:.11em;font-weight:590}.hero{display:flex;justify-content:space-between;align-items:flex-end;gap:24px;margin:4px 0 26px}.hero h2{font-size:31px;letter-spacing:-.7px;font-weight:510;margin:8px 0 6px}.hero p{margin:0;color:var(--muted);font-size:14px;line-height:1.55}.updated{font:11px ui-monospace,SFMono-Regular,Menlo,monospace;color:var(--dim)}
.metrics{display:grid;grid-template-columns:repeat(4,minmax(0,1fr));gap:12px;margin-bottom:18px}.metric,.panel,.operation{border:1px solid var(--border);background:rgba(255,255,255,.025);border-radius:var(--radius)}.metric{padding:17px}.metric-label{font-size:11px;color:var(--muted);margin-bottom:12px}.metric-value{font-size:25px;letter-spacing:-.45px;font-weight:510}.metric-meta{font-size:11px;color:var(--dim);margin-top:7px}.metric.good .metric-value{color:#a7f3d0}.metric.warn .metric-value{color:#fcd34d}.metric.bad .metric-value{color:#fca5a5}.grid-2{display:grid;grid-template-columns:minmax(0,1.4fr) minmax(300px,.6fr);gap:14px}.panel{padding:18px;min-width:0}.panel-head{display:flex;align-items:center;justify-content:space-between;margin-bottom:15px}.panel-title{font-size:13px;font-weight:590}.panel-sub{font-size:11px;color:var(--dim)}.empty{border:1px dashed var(--border);border-radius:8px;color:var(--dim);padding:26px;text-align:center;font-size:12px}.finding{display:grid;grid-template-columns:9px 1fr auto;gap:10px;align-items:start;padding:11px 0;border-bottom:1px solid var(--border2)}.finding:last-child{border-bottom:0}.finding-dot{width:7px;height:7px;border-radius:50%;background:var(--yellow);margin-top:5px}.finding.critical .finding-dot{background:var(--red)}.finding strong{font-size:12px}.finding p{font-size:11px;color:var(--muted);margin:4px 0 0}.badge{display:inline-flex;align-items:center;border:1px solid var(--border);border-radius:999px;padding:3px 7px;font:10px ui-monospace,SFMono-Regular,Menlo,monospace;color:var(--muted);white-space:nowrap}.badge.get,.badge.read_only{color:#a7f3d0;border-color:rgba(16,185,129,.25);background:rgba(16,185,129,.06)}.badge.post,.badge.put,.badge.patch{color:#c4b5fd;border-color:rgba(113,112,255,.3);background:rgba(113,112,255,.07)}.badge.delete,.badge.destructive{color:#fca5a5;border-color:rgba(239,68,68,.25);background:rgba(239,68,68,.06)}
- .form-row{display:grid;grid-template-columns:1fr auto;gap:9px}.form-row.triple{grid-template-columns:160px 1fr auto}.input{width:100%;border:1px solid var(--border);background:rgba(255,255,255,.025);color:var(--text);border-radius:7px;padding:10px 12px;outline:0;font-size:13px}.input:focus{border-color:rgba(113,112,255,.6);box-shadow:0 0 0 3px rgba(113,112,255,.1)}select.input{appearance:none}.result{margin-top:14px;min-height:110px}.layer-grid{display:grid;grid-template-columns:repeat(auto-fit,minmax(190px,1fr));gap:9px}.layer{border:1px solid var(--border2);background:rgba(255,255,255,.02);border-radius:8px;padding:12px}.layer h4{font-size:11px;margin:0 0 8px;text-transform:uppercase;color:var(--muted);letter-spacing:.07em}.layer pre,.json{white-space:pre-wrap;word-break:break-word;margin:0;color:var(--secondary);font:11px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace}.safe-banner{display:flex;align-items:center;gap:12px;border-radius:9px;padding:14px;border:1px solid rgba(16,185,129,.25);background:rgba(16,185,129,.06)}.safe-banner.blocked{border-color:rgba(239,68,68,.25);background:rgba(239,68,68,.06)}.safe-icon{font-size:21px}
+ .form-row{display:grid;grid-template-columns:1fr auto;gap:9px}.form-row.triple{grid-template-columns:160px 1fr auto}.sharing-form{grid-template-columns:140px 170px minmax(180px,1fr) auto}.input{width:100%;border:1px solid var(--border);background:rgba(255,255,255,.025);color:var(--text);border-radius:7px;padding:10px 12px;outline:0;font-size:13px}.input:focus{border-color:rgba(113,112,255,.6);box-shadow:0 0 0 3px rgba(113,112,255,.1)}select.input{appearance:none}.result{margin-top:14px;min-height:110px}.layer-grid{display:grid;grid-template-columns:repeat(auto-fit,minmax(190px,1fr));gap:9px}.layer{border:1px solid var(--border2);background:rgba(255,255,255,.02);border-radius:8px;padding:12px}.layer h4{font-size:11px;margin:0 0 8px;text-transform:uppercase;color:var(--muted);letter-spacing:.07em}.layer pre,.json{white-space:pre-wrap;word-break:break-word;margin:0;color:var(--secondary);font:11px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace}.safe-banner{display:flex;align-items:center;gap:12px;border-radius:9px;padding:14px;border:1px solid rgba(16,185,129,.25);background:rgba(16,185,129,.06)}.safe-banner.blocked{border-color:rgba(239,68,68,.25);background:rgba(239,68,68,.06)}.safe-icon{font-size:21px}
.toolbar{display:flex;align-items:center;gap:9px;flex-wrap:wrap;margin-bottom:15px}.toolbar .input{max-width:360px}.chips{display:flex;gap:6px;flex-wrap:wrap}.chip{border:1px solid var(--border);background:transparent;color:var(--muted);border-radius:999px;padding:6px 10px;font-size:11px}.chip.active,.chip:hover{color:var(--text);background:rgba(255,255,255,.05)}.operations{display:grid;grid-template-columns:repeat(3,minmax(0,1fr));gap:10px}.operation{padding:14px;display:flex;flex-direction:column;gap:10px;min-height:145px;transition:.16s ease}.operation:hover{border-color:rgba(113,112,255,.35);transform:translateY(-1px);background:rgba(255,255,255,.035)}.operation-top{display:flex;justify-content:space-between;gap:8px}.operation-path{font:11px/1.45 ui-monospace,SFMono-Regular,Menlo,monospace;color:var(--secondary);word-break:break-all}.operation h3{font-size:12px;font-weight:590;margin:0}.operation p{font-size:11px;color:var(--muted);line-height:1.45;margin:0;flex:1}.operation-badges{display:flex;gap:5px;flex-wrap:wrap}.table-wrap{overflow:auto;border:1px solid var(--border);border-radius:9px}table{width:100%;border-collapse:collapse;min-width:760px}th,td{text-align:left;padding:10px 12px;border-bottom:1px solid var(--border2);font-size:11px}th{color:var(--dim);text-transform:uppercase;letter-spacing:.07em;font-size:9px;background:rgba(255,255,255,.02)}td{color:var(--secondary)}td.mono{font-family:ui-monospace,SFMono-Regular,Menlo,monospace}.status-code.ok{color:#a7f3d0}.status-code.err{color:#fca5a5}
.login{position:fixed;inset:0;z-index:100;background:radial-gradient(circle at 50% 15%,rgba(113,112,255,.18),transparent 30%),#08090a;display:grid;place-items:center;padding:20px}.login-card{width:min(420px,100%);border:1px solid var(--border);background:#0f1011;border-radius:14px;padding:28px;box-shadow:0 28px 90px rgba(0,0,0,.5)}.login-card .logo{margin-bottom:22px}.login-card h1{font-size:25px;letter-spacing:-.5px;font-weight:510;margin:0 0 8px}.login-card p{font-size:13px;color:var(--muted);line-height:1.5;margin:0 0 20px}.login-card form{display:grid;gap:10px}.error{color:#fca5a5;font-size:11px;min-height:16px}.toast{position:fixed;right:20px;bottom:20px;z-index:120;background:#191a1b;border:1px solid var(--border);border-radius:9px;padding:11px 14px;font-size:12px;color:var(--secondary);box-shadow:0 12px 40px rgba(0,0,0,.35);transform:translateY(20px);opacity:0;pointer-events:none;transition:.2s}.toast.show{transform:none;opacity:1}.spinner{width:15px;height:15px;border:2px solid rgba(255,255,255,.15);border-top-color:var(--accent);border-radius:50%;animation:spin .7s linear infinite;display:inline-block}@keyframes spin{to{transform:rotate(360deg)}}
@media(max-width:1050px){.operations{grid-template-columns:repeat(2,minmax(0,1fr))}.metrics{grid-template-columns:repeat(2,minmax(0,1fr))}.grid-2{grid-template-columns:1fr}}
- @media(max-width:760px){.shell{display:block}.sidebar{position:fixed;transform:translateX(-100%);transition:.2s;width:260px}.sidebar.open{transform:none}.mobile-menu{display:inline-flex}.topbar{padding:0 15px}.content{padding:20px 14px}.hero{align-items:flex-start;flex-direction:column}.hero h2{font-size:25px}.operations{grid-template-columns:1fr}.metrics{grid-template-columns:repeat(2,1fr)}.form-row,.form-row.triple{grid-template-columns:1fr}.btn,.input{min-height:44px}.top-actions .desktop-only{display:none}}
+ @media(max-width:760px){.shell{display:block}.sidebar{position:fixed;transform:translateX(-100%);transition:.2s;width:260px}.sidebar.open{transform:none}.mobile-menu{display:inline-flex}.topbar{padding:0 15px}.content{padding:20px 14px}.hero{align-items:flex-start;flex-direction:column}.hero h2{font-size:25px}.operations{grid-template-columns:1fr}.metrics{grid-template-columns:repeat(2,1fr)}.form-row,.form-row.triple,.sharing-form{grid-template-columns:1fr}.btn,.input{min-height:44px}.top-actions .desktop-only{display:none}}
@media(max-width:420px){.metrics{grid-template-columns:1fr}.content{padding:18px 10px}.panel{padding:14px}.metric{padding:14px}}
-
-
- 🍳
- Homelab Control Plane
- Pfannkuchen Butler
- Geschützte Operations- und Diagnosekonsole. Der Token bleibt nur für den Login im Arbeitsspeicher und wird danach durch eine HttpOnly-Sitzung ersetzt.
-
-
-
-
-
+
- Übersicht
+
Live Zustand
Dein Homelab auf einen Blick.
Deterministisch aus Butler-Collectoren – keine geschätzten Zustände.
Noch nicht geladen
@@ -69,14 +55,29 @@
Gib einen bekannten Host oder Service ein.
+
+ Konservative Anomalieerkennung
Emby Account Sharing
Zeitgleiche öffentliche Netze und geografisch unmögliche Wechsel – über beide Emby-Server.
read-only
+
+
+ IPv4 wird exakt verglichen. IPv6 wird als /64 angezeigt und für die Haushaltserkennung konservativ auf /48 zusammengefasst.
+
+ Analysierte Benutzer
—
Auffällige Benutzer
—
Zeitgleiche Netze
—
Unmögliche Wechsel
—
+ Analyse noch nicht gestartet.
Benutzer-Risiko
nur begründete Treffer
+
+
Read-only Safety Gate
Wartung & Drift
Blocker erkennen, bevor eine Änderung Schaden anrichtet.
Noch kein Preflight ausgeführt.
- Persistente Historie
Audit
Redigierte Butler-Aktionen aus SQLite.
- Letzte Aktionen
| Zeit | Methode | Endpoint | Status | Dry-run | Detail |
|---|
| Wird geladen … |
+ KI- und API-Protokoll
Butler-Aktivitäten
Wer hat wann welche Butler-Funktion aufgerufen? Sensible Werte bleiben redigiert.
+ Letzte Aktivitäten
| Zeit | Akteur | Methode | Endpoint | Status | Dry-run | Detail |
|---|
| Wird geladen … |