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", 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()