diff --git a/app.py b/app.py index 22aa55d..650d4f5 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.6" 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", "/network/media-tunnel", "/caddy/")) + critical = path.startswith(("/network/wireguard", "/caddy/")) destructive = method.upper() == "DELETE" or any( marker in path for marker in ("/destroy/", "/cleanup/", "/break-lock/", "/restore/") ) @@ -2249,315 +2249,6 @@ 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 5f52007..80304be 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.6" def test_media_handoff_proxies_strict_category_contract(monkeypatch): @@ -392,9 +392,6 @@ 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") @@ -1277,70 +1274,3 @@ 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()