From 650caedacfa4bb8ad7d9625c1b3ae80395101eff Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 12:58:08 +0200 Subject: [PATCH 01/20] feat: add Sascha direct media tunnel support (app.py) --- app.py | 313 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 311 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index 650d4f5..22aa55d 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.3.6" +VERSION = "2.3.7" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -750,7 +750,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/") ) @@ -2249,6 +2249,315 @@ 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 + + +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": []}, "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 ++ ++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 _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 / "wg-media.conf.candidate" ++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_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 + 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 + + SYSCTL_AUDIT_KEYS = ( "net.core.default_qdisc", "net.core.rmem_default", From 32885c48b9da0c56bda4349a3dcf3e99eab7639c Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 12:58:08 +0200 Subject: [PATCH 02/20] feat: add Sascha direct media tunnel support (tests/test_app.py) --- tests/test_app.py | 72 ++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 71 insertions(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index 80304be..5f52007 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -43,7 +43,7 @@ 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.6" + assert response.json()["version"] == app.VERSION == "2.3.7" def test_media_handoff_proxies_strict_category_contract(monkeypatch): @@ -392,6 +392,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") @@ -1274,3 +1277,70 @@ 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() From ebdb77144ad8070fa31bfbefd92359a0d7ec58c4 Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:02:45 +0200 Subject: [PATCH 03/20] fix: validate wg candidate safely (app.py) --- app.py | 25 +++++++++++++++++++++++-- 1 file changed, 23 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index 22aa55d..613089b 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.3.7" +VERSION = "2.3.8" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -2409,7 +2409,7 @@ def _media_tunnel_install_command(role: Literal["vps", "emby"], peer_public_key: +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 / "wg-media.conf.candidate" ++candidate = root / "wg-media-candidate.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) @@ -2485,6 +2485,22 @@ def _media_tunnel_rollback_command(remove_keys: bool) -> str: 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: @@ -2531,6 +2547,11 @@ def _deploy_sascha_media_tunnel() -> dict: 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 From 08c7a5a59df9cf9a118e1c0543d825702fc7e643 Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:02:46 +0200 Subject: [PATCH 04/20] fix: validate wg candidate safely (tests/test_app.py) --- tests/test_app.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index 5f52007..c1bb9c8 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -43,7 +43,7 @@ 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.7" + assert response.json()["version"] == app.VERSION == "2.3.8" def test_media_handoff_proxies_strict_category_contract(monkeypatch): From a9f6c86c7cc25d3384d94070cd280e500217ed9e Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:04:05 +0200 Subject: [PATCH 05/20] fix: valid WireGuard interface candidate name (app.py) --- app.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index 613089b..1f63730 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.3.8" +VERSION = "2.3.9" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -2409,7 +2409,7 @@ def _media_tunnel_install_command(role: Literal["vps", "emby"], peer_public_key: +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 / "wg-media-candidate.conf" ++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) From b0fd3c1e361e8866b6084becdae0935fcde3af5e Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:04:06 +0200 Subject: [PATCH 06/20] fix: valid WireGuard interface candidate name (tests/test_app.py) --- tests/test_app.py | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index c1bb9c8..4552a81 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -43,7 +43,7 @@ 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.8" + assert response.json()["version"] == app.VERSION == "2.3.9" def test_media_handoff_proxies_strict_category_contract(monkeypatch): @@ -1344,3 +1344,14 @@ def test_sascha_media_tunnel_apply_returns_redacted_result(monkeypatch): 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 From f9664d5e36b057b544f8c18ed49a88c76660d7d3 Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:06:07 +0200 Subject: [PATCH 07/20] feat: benchmark direct media path (app.py) --- app.py | 81 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 79 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index 1f63730..b93e8d6 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.3.9" +VERSION = "2.4.0" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -2272,7 +2272,7 @@ def _sascha_media_edge_audit_command() -> str: + 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": []}, "network": {}} ++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"]) @@ -2316,6 +2316,15 @@ def _sascha_media_edge_audit_command() -> str: + 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 @@ -2579,6 +2588,74 @@ async def deploy_sascha_media_tunnel(req: SaschaMediaTunnelRequest, _=Depends(_v 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", From fa01dd77fb86f421fbdedcc3065b84135b1b167d Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:06:08 +0200 Subject: [PATCH 08/20] feat: benchmark direct media path (tests/test_app.py) --- tests/test_app.py | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index 4552a81..a98a293 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -43,7 +43,7 @@ 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.9" + assert response.json()["version"] == app.VERSION == "2.4.0" def test_media_handoff_proxies_strict_category_contract(monkeypatch): @@ -1355,3 +1355,18 @@ def test_sascha_media_tunnel_candidate_uses_valid_wireguard_interface_name(): 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"} From 07d3e019e28e9a0b509a202c5dc19ef8335919eb Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:13:18 +0200 Subject: [PATCH 09/20] feat: optimize Sascha Emby edge (app.py) --- app.py | 101 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 100 insertions(+), 1 deletion(-) diff --git a/app.py b/app.py index b93e8d6..d5a06ef 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.4.0" +VERSION = "2.4.1" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -2264,6 +2264,11 @@ class SaschaMediaTunnelRequest(BaseModel): 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 @@ -2355,6 +2360,100 @@ async def sascha_media_edge_audit(_=Depends(_verify)): 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 From 3c3bc0fcda6de493c067090a3d420c188aa54dff Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:13:19 +0200 Subject: [PATCH 10/20] feat: optimize Sascha Emby edge (tests/test_app.py) --- tests/test_app.py | 30 +++++++++++++++++++++++++++++- 1 file changed, 29 insertions(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index a98a293..3178eb8 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -43,7 +43,7 @@ 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.4.0" + assert response.json()["version"] == app.VERSION == "2.4.1" def test_media_handoff_proxies_strict_category_contract(monkeypatch): @@ -1370,3 +1370,31 @@ def test_sascha_media_path_benchmark_returns_both_fixed_routes(monkeypatch): 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" From 8892ea45973db61729ab22b329cb33538d3e850e Mon Sep 17 00:00:00 2001 From: sascha Date: Mon, 7 Sep 2026 08:58:05 +0200 Subject: [PATCH 11/20] feat: add authenticated DNS RRSet upsert endpoint --- app.py | 190 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 188 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index d5a06ef..b950c4d 100644 --- a/app.py +++ b/app.py @@ -12,14 +12,17 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.4.1" +VERSION = "2.4.6" 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") +MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get( + "MEDIA_HANDOFF_ALLOWED_NETWORKS", + "10.2.1.119/32,10.5.85.12/32", +) # --- Config loading --- @@ -1571,6 +1574,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) @@ -3237,6 +3247,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(".") @@ -3902,6 +3959,133 @@ 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) @@ -3933,6 +4117,8 @@ def _media_handoff_caller_allowed(request: Request) -> bool: 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") allowed_categories = {"serien4k", "serien", "video4k", "video"} From f3b9a0c33a2a135978bcb2148bf1de066e93e4bb Mon Sep 17 00:00:00 2001 From: Trulla Date: Wed, 16 Sep 2026 12:11:31 +0200 Subject: [PATCH 12/20] Allow serienen/videoen in media handoff allowed_categories kannte nur serien4k|serien|video4k|video. Die englischen Arr-Instanzen sonarrEN (Port 8991, Root /data/FHD/serienen) und radarrEN (Port 7880, Root /data/FHD/videoen) nutzen eigene SAB-Kategorien. Folge: POST /media/handoff antwortete 422 "Unsupported media category", movetdarr.sh brach mit "n8n-Handoff konnte nicht registriert werden - kein Move" ab und Mutiny.2026 x2 (~21 GB) lagen seit 08./09.09.2026 unangetastet in /usenet/complete/videoen. radarrEN hat den Film weiter monitored, hasFile=false, Queue leer. Die Pfadpruefung (expected_prefix) bleibt unveraendert und gilt auch fuer die neuen Kategorien. 3 neue Tests: serienen+videoen werden geproxyt, erfundene Kategorie bleibt 422, Pfad-Mismatch bei videoen bleibt 422. 6/6 media_handoff-Tests gruen. test_health_exposes_current_version schlug schon vor dieser Aenderung fehl (Test erwartet 2.4.1, VERSION ist 2.4.6). --- app.py | 6 ++- tests/test_app.py | 93 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 98 insertions(+), 1 deletion(-) diff --git a/app.py b/app.py index b950c4d..cd59803 100644 --- a/app.py +++ b/app.py @@ -4121,7 +4121,11 @@ async def media_handoff(payload: MediaHandoffPayload, request: Request): _audit("/media/handoff", "POST", 403, f"caller={caller} action={payload.action}") raise HTTPException(403, "Media handoff caller is not allowed") - allowed_categories = {"serien4k", "serien", "video4k", "video"} + # 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") diff --git a/tests/test_app.py b/tests/test_app.py index 3178eb8..adf6f9c 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -117,6 +117,99 @@ def test_media_handoff_rejects_wrong_category_path(monkeypatch): 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 = {} From cadcf8766f1d0fdd2f42a9a03447a03879c09556 Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 17 Sep 2026 15:50:16 +0200 Subject: [PATCH 13/20] Update app.py for media verify endpoint --- app.py | 83 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 82 insertions(+), 1 deletion(-) diff --git a/app.py b/app.py index cd59803..7549309 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.4.6" +VERSION = "2.4.7" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -1546,6 +1546,87 @@ 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, gt=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 "'" 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") + + probe_cmd = ( + f"docker exec {container} ffprobe -v error " + f"-show_entries format=duration -of default=noprint_wrappers=1:nokey=1 '{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" From 8230e03745fd5745f84f781f31d1966aaa5eb55e Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 17 Sep 2026 15:50:17 +0200 Subject: [PATCH 14/20] Update tests/test_app.py for media verify endpoint --- tests/test_app.py | 84 +++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 84 insertions(+) diff --git a/tests/test_app.py b/tests/test_app.py index adf6f9c..ae9cde0 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -901,6 +901,90 @@ 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_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'; rm -rf /'.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"}, + ) + assert response.status_code == 502 + + def test_invalid_log_target_is_rejected_before_ssh(): with TestClient(app.app) as client: response = client.get( From 1346f5f8a62e4d7a05784eac13c2c8a122f7f530 Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 17 Sep 2026 15:50:17 +0200 Subject: [PATCH 15/20] Update README.md for media verify endpoint --- README.md | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/README.md b/README.md index d0a464f..2b4cc1e 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,10 @@ Integration Compose definition: `tests/compose.integration.yaml` (binds only to ## Changelog +### 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. From 4d48fd40c083cf012c2e68785fda40919b1f97fa Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 17 Sep 2026 21:54:49 +0200 Subject: [PATCH 16/20] Fix media/verify apostrophe false-positive in app.py --- app.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/app.py b/app.py index 7549309..44508ca 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, math, hashlib +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 @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.4.7" +VERSION = "2.4.8" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -1577,7 +1577,7 @@ async def media_verify(payload: MediaVerifyRequest, _=Depends(_verify)): 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 "'" in payload.path: + 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] @@ -1590,7 +1590,7 @@ async def media_verify(payload: MediaVerifyRequest, _=Depends(_verify)): probe_cmd = ( f"docker exec {container} ffprobe -v error " - f"-show_entries format=duration -of default=noprint_wrappers=1:nokey=1 '{payload.path}'" + 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 From f5c650fe3d18d6cec681aebc9e7675b4e1ac153c Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 17 Sep 2026 21:54:50 +0200 Subject: [PATCH 17/20] Fix media/verify apostrophe false-positive in tests/test_app.py --- tests/test_app.py | 28 +++++++++++++++++++++++++++- 1 file changed, 27 insertions(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index ae9cde0..1bc0857 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 @@ -951,13 +952,38 @@ def test_media_verify_rejects_path_outside_data_before_ssh(monkeypatch): 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'; rm -rf /'.mkv"}, + json={"host": "arrapps", "path": "/data/UHD/serien/a\x00.mkv"}, ) assert response.status_code == 400 From 211eab032f7eb166bbb90326d528c32021e567b9 Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 17 Sep 2026 21:54:50 +0200 Subject: [PATCH 18/20] Fix media/verify apostrophe false-positive in README.md --- README.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/README.md b/README.md index 2b4cc1e..aa96e1d 100644 --- a/README.md +++ b/README.md @@ -98,6 +98,10 @@ 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. From 686272a7c81447badaae32fad4e72f382adc5843 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 18 Sep 2026 07:19:06 +0200 Subject: [PATCH 19/20] Handle runtime=0 gracefully in app.py --- app.py | 19 +++++++++++++++++-- 1 file changed, 17 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index 44508ca..7ddb58a 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.4.8" +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") @@ -1558,7 +1558,7 @@ MEDIA_VERIFY_PROBES = { class MediaVerifyRequest(BaseModel): host: str = Field(..., max_length=64) path: str = Field(..., max_length=1000) - expected_minutes: float | None = Field(None, gt=0, le=1440) + expected_minutes: float | None = Field(None, ge=0, le=1440) tolerance_pct: float = Field(20.0, gt=0, le=100) @@ -1588,6 +1588,21 @@ async def media_verify(payload: MediaVerifyRequest, _=Depends(_verify)): 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)}" From de37b51829e8030cff722cf8bc7a52c9a709f61d Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 18 Sep 2026 07:19:07 +0200 Subject: [PATCH 20/20] Handle runtime=0 gracefully in tests/test_app.py --- tests/test_app.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index 1bc0857..45011c2 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -1006,7 +1006,7 @@ def test_media_verify_returns_502_on_unparsable_ffprobe_output(monkeypatch): response = client.post( "/media/verify", headers={"Authorization": "Bearer test-token"}, - json={"host": "arrapps", "path": "/data/UHD/serien/ep.mkv"}, + json={"host": "arrapps", "path": "/data/UHD/serien/ep.mkv", "expected_minutes": 49.5}, ) assert response.status_code == 502