Compare commits

..

34 commits

Author SHA1 Message Date
f891b4859f Merge pull request 'Handle expected_minutes=0 gracefully (defense-in-depth)' (#61) from fix/media-verify-runtime-zero-handling into main 2026-09-18 07:19:45 +02:00
de37b51829 Handle runtime=0 gracefully in tests/test_app.py 2026-09-18 07:19:07 +02:00
686272a7c8 Handle runtime=0 gracefully in app.py 2026-09-18 07:19:06 +02:00
af19a60853 Merge pull request 'Fix /media/verify apostrophe false-positive' (#60) from fix/media-verify-apostrophe-escaping into main 2026-09-17 21:55:06 +02:00
211eab032f Fix media/verify apostrophe false-positive in README.md 2026-09-17 21:54:50 +02:00
f5c650fe3d Fix media/verify apostrophe false-positive in tests/test_app.py 2026-09-17 21:54:50 +02:00
4d48fd40c0 Fix media/verify apostrophe false-positive in app.py 2026-09-17 21:54:49 +02:00
e1dafd99b7 Merge pull request 'Add POST /media/verify - ffprobe-based truncated import detection' (#59) from feature/media-verify-endpoint into main 2026-09-17 15:50:29 +02:00
1346f5f8a6 Update README.md for media verify endpoint 2026-09-17 15:50:17 +02:00
8230e03745 Update tests/test_app.py for media verify endpoint 2026-09-17 15:50:17 +02:00
cadcf8766f Update app.py for media verify endpoint 2026-09-17 15:50:16 +02:00
453de2b037 Merge pull request 'Allow serienen/videoen in media handoff' (#58) from fix/handoff-en-categories-20260916 into main 2026-09-16 12:12:15 +02:00
Trulla
f3b9a0c33a 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).
2026-09-16 12:11:31 +02:00
f24809a22a Merge pull request 'feat: authenticated DNS RRSet upsert' (#57) from feat/dns-rrset-upsert-20260907 into main 2026-09-07 09:04:19 +02:00
8892ea4597 feat: add authenticated DNS RRSet upsert endpoint 2026-09-07 08:58:05 +02:00
161faf6c65 Merge pull request 'feat: Sascha-Emby auf Direkt-WG und h1+h2 umstellen' (#56) from feat/optimize-tv-sascha-edge-20260905 into main 2026-09-05 13:13:21 +02:00
3c3bc0fcda feat: optimize Sascha Emby edge (tests/test_app.py) 2026-09-05 13:13:19 +02:00
07d3e019e2 feat: optimize Sascha Emby edge (app.py) 2026-09-05 13:13:18 +02:00
36d122d0bb Merge pull request 'feat: Media-Pfade node6 gegen Direkt-WireGuard messen' (#55) from feat/sascha-media-path-benchmark-20260905 into main 2026-09-05 13:06:09 +02:00
fa01dd77fb feat: benchmark direct media path (tests/test_app.py) 2026-09-05 13:06:08 +02:00
f9664d5e36 feat: benchmark direct media path (app.py) 2026-09-05 13:06:07 +02:00
6c92c7023a Merge pull request 'fix: gültiger WireGuard-Kandidatenname' (#54) from fix/wg-candidate-interface-name-20260905 into main 2026-09-05 13:04:08 +02:00
b0fd3c1e36 fix: valid WireGuard interface candidate name (tests/test_app.py) 2026-09-05 13:04:06 +02:00
a9f6c86c7c fix: valid WireGuard interface candidate name (app.py) 2026-09-05 13:04:05 +02:00
02762cc651 Merge pull request 'fix: WireGuard-Kandidat korrekt validieren' (#53) from fix/sascha-media-tunnel-validation-20260905 into main 2026-09-05 13:02:48 +02:00
08c7a5a59d fix: validate wg candidate safely (tests/test_app.py) 2026-09-05 13:02:46 +02:00
ebdb77144a fix: validate wg candidate safely (app.py) 2026-09-05 13:02:45 +02:00
b693ce7467 Merge pull request 'feat: direkter Media-Tunnel für tv.sascha-lutz.de' (#52) from feat/sascha-direct-media-tunnel-20260905 into main 2026-09-05 12:58:10 +02:00
32885c48b9 feat: add Sascha direct media tunnel support (tests/test_app.py) 2026-09-05 12:58:08 +02:00
650caedacf feat: add Sascha direct media tunnel support (app.py) 2026-09-05 12:58:08 +02:00
d36a3ea615 Merge pull request 'fix: wait for every file in season packs' (#51) from fix/media-handoff-season-packs-20260905 into main 2026-09-05 08:20:50 +02:00
05785444d5 fix: wait for every file in season packs 2026-09-05 08:20:48 +02:00
742d3ebfa7 fix: wait for every file in season packs 2026-09-05 08:20:47 +02:00
7ddb88b54e Merge pull request 'feat: add restricted media handoff bridge' (#50) from feat/media-handoff-orchestrator-20260905 into main 2026-09-05 07:59:54 +02:00
3 changed files with 1137 additions and 6 deletions

View file

@ -45,6 +45,7 @@ Known secret response fields such as Dockhand's `hawserToken` and `webhookSecret
| `/docker/inspect/{host}/{container}` | GET | Sanitized image, runtime, resources, mounts and state |
| `/docker/restart/{host}/{container}` | POST | Restart container; supports `dry_run=true` |
| `/config/reload` | POST | Reload YAML configuration and credential cache |
| `/media/verify` | POST | ffprobe duration check on a NAS media file to catch truncated Sonarr/Radarr imports; flags `suspect` if deviation from `expected_minutes` exceeds `tolerance_pct` |
## VM lifecycle and inventory
@ -97,6 +98,14 @@ Integration Compose definition: `tests/compose.integration.yaml` (binds only to
## Changelog
### 2.4.8 — 17.09.2026
- Fixed `POST /media/verify` false-positive rejection: filenames containing a legitimate apostrophe (e.g. "La'An" in a Star Trek episode title) were blocked as "unsafe characters". Replaced the naive single-quote wrapping + apostrophe blocklist with proper `shlex.quote()` escaping for the remote ffprobe command; only newline/NUL byte injection is still rejected.
### 2.4.7 — 17.09.2026
- Added `POST /media/verify`: probes a media file's real duration via `ffprobe` (running inside the existing `bazarrUHD` container on `arrapps`, which already mounts `/data` and ships ffmpeg — no new service needed) and flags it as `suspect` when it deviates from an `expected_minutes` value beyond `tolerance_pct`. Catches truncated Sonarr/Radarr imports (e.g. a re-grab that lands as 1.8 GB instead of the expected 8 GB for a UHD episode) after the file is already on the NAS.
### 2.3.2 — 22.07.2026
- Make Vaultwarden refresh durable: persistent named cache volume, protected runtime credentials, automatic API-key re-login and atomic cache writes.

803
app.py
View file

@ -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,14 +12,17 @@ from contextlib import asynccontextmanager
from contextvars import ContextVar
log = logging.getLogger("butler")
VERSION = "2.3.6"
VERSION = "2.4.9"
API_DIR = os.environ.get("API_KEY_DIR", "/data/api")
VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache")
BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "")
CONFIG_PATH = os.environ.get("BUTLER_CONFIG", "/data/butler.yaml")
UI_PATH = os.environ.get("BUTLER_UI_PATH", os.path.join(os.path.dirname(__file__), "ui.html"))
MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get("MEDIA_HANDOFF_ALLOWED_NETWORKS", "10.2.1.119/32")
MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get(
"MEDIA_HANDOFF_ALLOWED_NETWORKS",
"10.2.1.119/32,10.5.85.12/32",
)
# --- Config loading ---
@ -750,7 +753,7 @@ async def capabilities(_=Depends(_verify)):
mode = "read_only" if method.upper() == "GET" else "mutation"
serialized = {"parameters": operation.get("parameters", []), "requestBody": operation.get("requestBody", {})}
dry_run = _schema_contains_property(serialized, "dry_run", components)
critical = path.startswith(("/network/wireguard", "/caddy/"))
critical = path.startswith(("/network/wireguard", "/network/media-tunnel", "/caddy/"))
destructive = method.upper() == "DELETE" or any(
marker in path for marker in ("/destroy/", "/cleanup/", "/break-lock/", "/restore/")
)
@ -1543,6 +1546,102 @@ async def vault_reload(_=Depends(_verify)):
return {"reloaded": True, "items": len(_vault_cache)}
# --- Media file integrity verification ---
MEDIA_VERIFY_PROBES = {
"arrapps": ("bazarrUHD", "sascha"),
"arr-chris": (None, "chris"),
"arr-chris-live": (None, "chris"),
}
class MediaVerifyRequest(BaseModel):
host: str = Field(..., max_length=64)
path: str = Field(..., max_length=1000)
expected_minutes: float | None = Field(None, ge=0, le=1440)
tolerance_pct: float = Field(20.0, gt=0, le=100)
@app.post("/media/verify")
async def media_verify(payload: MediaVerifyRequest, _=Depends(_verify)):
"""Probe a media file's real duration via ffprobe to catch truncated imports.
Runs `ffprobe` inside an existing container (bazarrUHD on arrapps, which already
mounts /data and ships ffmpeg) so no new service is required. Compares the
measured duration against an expected runtime (minutes) supplied by the caller
(e.g. Sonarr/Radarr's runtime field) within tolerance_pct.
"""
if not re.fullmatch(r"[A-Za-z0-9_.-]+", payload.host):
raise HTTPException(400, "Invalid host name")
if payload.host not in MEDIA_VERIFY_PROBES:
raise HTTPException(400, f"No ffprobe container configured for host {payload.host}")
if not re.fullmatch(r"/data/[^\x00]+\.(mkv|mp4|avi|m4v|ts)", payload.path):
raise HTTPException(400, "Path must be an absolute /data media file")
if ".." in payload.path or "\n" in payload.path or "\x00" in payload.path:
raise HTTPException(400, "Path contains unsafe characters")
container, _default_user = MEDIA_VERIFY_PROBES[payload.host]
if not container:
raise HTTPException(400, f"Host {payload.host} has no configured ffprobe container yet")
target = _find_inventory_host(payload.host)
if not target:
raise HTTPException(404, f"Host {payload.host} not found in inventory")
# expected_minutes == 0 or None: no reference runtime available, skip verification gracefully
if not payload.expected_minutes or payload.expected_minutes <= 0:
return {
"host": payload.host,
"path": payload.path,
"duration_seconds": None,
"duration_minutes": None,
"expected_minutes": payload.expected_minutes,
"verified": False,
"suspect": False,
"skipped": True,
"skip_reason": "no_reference_runtime",
"deviation_pct": None,
}
probe_cmd = (
f"docker exec {container} ffprobe -v error "
f"-show_entries format=duration -of default=noprint_wrappers=1:nokey=1 {shlex.quote(payload.path)}"
)
rc, out, err = await asyncio.to_thread(
_ssh, f'{target["user"]}@{target["ip"]}', f"sudo -n {probe_cmd} || {probe_cmd}", 30
)
if rc != 0:
_audit("/media/verify", "POST", 502, f"host={payload.host} path={payload.path}")
raise HTTPException(502, (err or out).strip()[:500] or "ffprobe failed")
raw_duration = out.strip()
try:
duration_seconds = float(raw_duration)
except ValueError:
_audit("/media/verify", "POST", 502, f"host={payload.host} unparsable duration")
raise HTTPException(502, "ffprobe returned no parsable duration; file is likely corrupt")
duration_minutes = duration_seconds / 60
result = {
"host": payload.host,
"path": payload.path,
"duration_seconds": round(duration_seconds, 1),
"duration_minutes": round(duration_minutes, 2),
"expected_minutes": payload.expected_minutes,
"verified": True,
"suspect": False,
}
if payload.expected_minutes:
deviation_pct = abs(duration_minutes - payload.expected_minutes) / payload.expected_minutes * 100
result["deviation_pct"] = round(deviation_pct, 1)
result["suspect"] = deviation_pct > payload.tolerance_pct
_audit(
"/media/verify", "POST", 200,
f"host={payload.host} dur={duration_minutes:.1f}m suspect={result['suspect']}",
)
return result
# --- VPS reverse-proxy and DNS management ---
VPS_SSH = "root@46.225.230.72"
@ -1571,6 +1670,13 @@ class ProxyRouteRequest(BaseModel):
dns_token: str | None = None
class DnsRrsetUpsertRequest(BaseModel):
record_type: Literal["A", "AAAA", "CNAME"] = "A"
value: str
ttl: int = Field(default=300, ge=60, le=86400)
comment: str = Field(default="Managed by Homelab Butler", max_length=200)
def _remote_python(script: str, timeout: int = 30) -> tuple[int, str, str]:
encoded = base64.b64encode(script.encode()).decode()
return _ssh(VPS_SSH, f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"', timeout=timeout)
@ -2249,6 +2355,512 @@ async def network_wireguard_remove_peer(host: str, req: WireGuardPeerRemoveReque
return {"host": host, **result}
SASCHA_MEDIA_VPS_HOST = "pfannkuchen"
SASCHA_MEDIA_EMBY_HOST = "emby-sascha"
SASCHA_MEDIA_INTERFACE = "wg-media"
SASCHA_MEDIA_VPS_ADDRESS = "10.11.13.1/32"
SASCHA_MEDIA_EMBY_ADDRESS = "10.11.13.3/32"
SASCHA_MEDIA_PORT = 51821
SASCHA_MEDIA_MTU = 1340
SASCHA_MEDIA_CONFIRMATION = "DEPLOY_DIRECT_SASCHA_MEDIA_TUNNEL"
class SaschaMediaTunnelRequest(BaseModel):
dry_run: bool = True
confirmation: str | None = None
class SaschaMediaEdgeOptimizeRequest(BaseModel):
dry_run: bool = True
confirmation: str | None = None
def _sascha_media_edge_audit_command() -> str:
script = '''import json, re, subprocess
+from pathlib import Path
+
+def run(args):
+ proc = subprocess.run(args, text=True, capture_output=True, timeout=30)
+ return {"rc": proc.returncode, "stdout": proc.stdout.strip(), "stderr": proc.stderr.strip()[-300:]}
+
+result = {"hostname": "pfannkuchen", "caddy": {"container_running": False, "protocols": [], "protocols_explicit": False, "upstreams": [], "site_block": [], "emby_snippet": []}, "network": {}}
+inspect = run(["docker", "inspect", "caddy", "--format", "{{.State.Running}}"])
+result["caddy"]["container_running"] = inspect["rc"] == 0 and inspect["stdout"] == "true"
+adapt = run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"])
+if adapt["rc"] == 0:
+ try:
+ config = json.loads(adapt["stdout"])
+ servers = config.get("apps", {}).get("http", {}).get("servers", {})
+ explicit = []
+ def walk(value, matched=False):
+ if isinstance(value, dict):
+ current = matched
+ host = value.get("host")
+ if isinstance(host, list) and "tv.sascha-lutz.de" in host:
+ current = True
+ if current and isinstance(value.get("dial"), str):
+ result["caddy"]["upstreams"].append(value["dial"])
+ for child in value.values(): walk(child, current)
+ elif isinstance(value, list):
+ for child in value: walk(child, matched)
+ for server in servers.values():
+ protocols = server.get("protocols")
+ if isinstance(protocols, list): explicit.extend(protocols)
+ walk(server)
+ result["caddy"]["protocols_explicit"] = bool(explicit)
+ result["caddy"]["protocols"] = sorted(set(explicit)) if explicit else ["h1", "h2", "h3"]
+ result["caddy"]["upstreams"] = sorted(set(result["caddy"]["upstreams"]))
+ except Exception as exc:
+ result["caddy"]["adapt_error"] = str(exc)[:200]
+else:
+ result["caddy"]["adapt_error"] = adapt["stderr"] or "caddy adapt failed"
+
+path = Path("/app-config/caddy/Caddyfile")
+if path.exists():
+ lines = path.read_text(encoding="utf-8", errors="replace").splitlines()
+ collecting = False; depth = 0; selected = []
+ for line in lines:
+ if not collecting and re.match(r"^\\s*tv\\.sascha-lutz\\.de\\s*\\{", line): collecting = True
+ if collecting:
+ depth += line.count("{") - line.count("}")
+ if not re.search(r"(?i)(password|secret|token|private|basicauth|basic_auth|hash)", line): selected.append(line.strip())
+ else: selected.append("[REDACTED SENSITIVE DIRECTIVE]")
+ if depth == 0: break
+ result["caddy"]["site_block"] = selected
+ collecting = False; depth = 0; selected = []
+ for line in lines:
+ if not collecting and re.match(r"^\\s*\\(emby_config\\)\\s*\\{", line): collecting = True
+ if collecting:
+ depth += line.count("{") - line.count("}")
+ if not re.search(r"(?i)(password|secret|token|private|basicauth|basic_auth|hash)", line): selected.append(line.strip())
+ else: selected.append("[REDACTED SENSITIVE DIRECTIVE]")
+ if depth == 0: break
+ result["caddy"]["emby_snippet"] = selected
+
+result["network"]["wg_media_active"] = subprocess.run(["systemctl", "is-active", "--quiet", "wg-quick@wg-media"]).returncode == 0
+result["network"]["wg_media_enabled"] = subprocess.run(["systemctl", "is-enabled", "--quiet", "wg-quick@wg-media"]).returncode == 0
+result["network"]["udp_51821"] = "51821" in run(["ss", "-H", "-lun"]) ["stdout"]
+for name, target in (("legacy_route", "10.6.1.103"), ("direct_route", "10.11.13.3")):
+ result["network"][name] = run(["ip", "route", "get", target])["stdout"][:300]
+print(json.dumps(result))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
@app.get("/media/edge/sascha")
async def sascha_media_edge_audit(_=Depends(_verify)):
"""Return a redacted snapshot of the dedicated Hetzner edge for tv.sascha-lutz.de."""
inventory = await asyncio.to_thread(_find_inventory_host, SASCHA_MEDIA_VPS_HOST)
if not inventory:
raise HTTPException(404, "Hetzner media edge not found")
target = f'{inventory["user"]}@{inventory["ip"]}'
rc, out, err = await asyncio.to_thread(_ssh, target, _sascha_media_edge_audit_command(), 45)
if rc != 0:
raise HTTPException(502, (err or out).strip()[-500:] or "Sascha media edge audit failed")
try:
result = json.loads(out)
except json.JSONDecodeError as exc:
raise HTTPException(502, "Sascha media edge audit returned invalid JSON") from exc
_audit("/media/edge/sascha", "GET", 200, "redacted live media-edge snapshot")
return result
def _sascha_media_edge_optimize_command() -> str:
script = '''import hashlib, json, os, re, shutil, subprocess, time
+from pathlib import Path
+path = Path("/app-config/caddy/Caddyfile")
+text = path.read_text(encoding="utf-8")
+old_pattern = re.compile(r"(?m)^(\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+)10\\.6\\.1\\.103:8096(\\s*)$")
+new_pattern = re.compile(r"(?m)^\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+10\\.11\\.13\\.3:8096\\s*$")
+old_count = len(old_pattern.findall(text)); new_count = len(new_pattern.findall(text))
+if old_count == 1:
+ text = old_pattern.sub(r"\\g<1>10.11.13.3:8096\\g<2>", text)
+elif not (old_count == 0 and new_count == 1):
+ raise RuntimeError("expected exactly one tv.sascha-lutz.de upstream")
+lines = text.splitlines(keepends=True)
+first = next((i for i, line in enumerate(lines) if line.strip() and not line.lstrip().startswith("#")), None)
+if first is None or lines[first].strip() != "{":
+ lines[0:0] = ["{\\n", " servers {\\n", " protocols h1 h2\\n", " }\\n", "}\\n", "\\n"]
+else:
+ depth = 0; global_end = None; servers_start = None; servers_end = None
+ for i in range(first, len(lines)):
+ stripped = lines[i].strip(); before = depth
+ if before == 1 and re.match(r"^servers(?:\\s+\\S+)?\\s*\\{$", stripped): servers_start = i
+ depth += lines[i].count("{") - lines[i].count("}")
+ if servers_start is not None and i > servers_start and depth == 1 and servers_end is None: servers_end = i
+ if i > first and depth == 0: global_end = i; break
+ if global_end is None: raise RuntimeError("unbalanced Caddy global options block")
+ if servers_start is None:
+ lines[global_end:global_end] = [" servers {\\n", " protocols h1 h2\\n", " }\\n"]
+ else:
+ if servers_end is None: raise RuntimeError("unbalanced Caddy servers block")
+ protocol_lines = [i for i in range(servers_start + 1, servers_end) if re.match(r"^\\s*protocols\\s+", lines[i])]
+ if len(protocol_lines) > 1: raise RuntimeError("multiple Caddy protocol directives")
+ if protocol_lines: lines[protocol_lines[0]] = re.sub(r"protocols\\s+.*", "protocols h1 h2", lines[protocol_lines[0]])
+ else: lines.insert(servers_start + 1, " protocols h1 h2\\n")
+text = "".join(lines)
+if re.search(r"(?im)^\\s*header(?:_down)?\\s+Alt-Svc", text): raise RuntimeError("manual Alt-Svc directive requires review")
+stamp = time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())
+backup = path.with_name("Caddyfile.pre-sascha-media-" + stamp)
+candidate = path.with_name("Caddyfile.sascha-media-candidate")
+shutil.copy2(path, backup); candidate.write_text(text, encoding="utf-8"); os.chmod(candidate, path.stat().st_mode)
+def run(args, timeout=30):
+ proc = subprocess.run(args, text=True, capture_output=True, timeout=timeout)
+ if proc.returncode != 0: raise RuntimeError((proc.stderr or proc.stdout).strip()[-500:] or "command failed")
+ return proc.stdout.strip()
+try:
+ run(["docker", "cp", str(candidate), "caddy:/tmp/Caddyfile.sascha-media-candidate"])
+ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/tmp/Caddyfile.sascha-media-candidate"])
+ with path.open("w", encoding="utf-8") as handle:
+ handle.write(text); handle.flush(); os.fsync(handle.fileno())
+ host_hash = hashlib.sha256(path.read_bytes()).hexdigest()
+ container_hash = run(["docker", "exec", "caddy", "sha256sum", "/etc/caddy/Caddyfile"]).split()[0]
+ if host_hash != container_hash: raise RuntimeError("host/container Caddyfile hash mismatch")
+ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/etc/caddy/Caddyfile"])
+ run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"])
+ adapted = json.loads(run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"]))
+ protocols = []
+ for server in adapted.get("apps", {}).get("http", {}).get("servers", {}).values(): protocols.extend(server.get("protocols", []))
+ if sorted(set(protocols)) != ["h1", "h2"]: raise RuntimeError("Caddy did not load h1+h2-only protocols")
+ if not run(["curl", "-fsS", "--max-time", "15", "http://10.11.13.3:8096/System/Ping"]): raise RuntimeError("Emby direct-path ping was empty")
+except Exception:
+ with path.open("w", encoding="utf-8") as handle:
+ handle.write(backup.read_text(encoding="utf-8")); handle.flush(); os.fsync(handle.fileno())
+ subprocess.run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"], text=True, capture_output=True, timeout=30)
+ raise
+finally:
+ candidate.unlink(missing_ok=True)
+ subprocess.run(["docker", "exec", "caddy", "rm", "-f", "/tmp/Caddyfile.sascha-media-candidate"], text=True, capture_output=True)
+print(json.dumps({"status": "optimized", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "backup": str(backup), "sha256": host_hash}))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _optimize_sascha_media_edge() -> dict:
return _ssh_json(_media_host_target(SASCHA_MEDIA_VPS_HOST), _sascha_media_edge_optimize_command(), 90)
@app.post("/media/edge/sascha/optimize")
async def optimize_sascha_media_edge(req: SaschaMediaEdgeOptimizeRequest, _=Depends(_verify)):
"""Switch only tv.sascha-lutz.de to wg-media and disable HTTP/3 on its dedicated Hetzner edge."""
plan = {"status": "would_optimize", "hostname": "tv.sascha-lutz.de", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "zero_downtime_reload": True}
if req.dry_run:
_audit("/media/edge/sascha/optimize", "POST", 200, "dry_run=True", True)
return plan
if req.confirmation != "OPTIMIZE_TV_SASCHA_LUTZ_DE":
raise HTTPException(400, "confirmation must be OPTIMIZE_TV_SASCHA_LUTZ_DE")
try:
result = await asyncio.to_thread(_optimize_sascha_media_edge)
except Exception as exc:
_audit("/media/edge/sascha/optimize", "POST", 502, "optimization failed; Caddy rollback attempted")
raise HTTPException(502, str(exc)[-500:]) from exc
_audit("/media/edge/sascha/optimize", "POST", 200, "direct upstream and h1+h2 enabled")
return result
def _media_tunnel_key_command() -> str:
script = '''import json, os, subprocess
+from pathlib import Path
+root = Path("/app-config/wireguard-media")
+private = root / "private.key"
+public = root / "public.key"
+root.mkdir(parents=True, exist_ok=True)
+os.chmod(root, 0o700)
+created = not private.exists()
+if created:
+ key = subprocess.run(["wg", "genkey"], check=True, text=True, capture_output=True).stdout.strip()
+ private.write_text(key + "\\n")
+ os.chmod(private, 0o600)
+if not public.exists() or created:
+ pub = subprocess.run(["wg", "pubkey"], input=private.read_text(), check=True, text=True, capture_output=True).stdout.strip()
+ public.write_text(pub + "\\n")
+ os.chmod(public, 0o644)
+print(json.dumps({"public_key": public.read_text().strip(), "created": created}))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _media_tunnel_install_command(role: Literal["vps", "emby"], peer_public_key: str) -> str:
if not re.fullmatch(r"[A-Za-z0-9+/]{43}=", peer_public_key):
raise ValueError("Invalid WireGuard public key")
settings = {
"vps": {
"address": SASCHA_MEDIA_VPS_ADDRESS,
"peer": SASCHA_MEDIA_EMBY_ADDRESS,
"endpoint": None,
"keepalive": None,
"listen": SASCHA_MEDIA_PORT,
},
"emby": {
"address": SASCHA_MEDIA_EMBY_ADDRESS,
"peer": SASCHA_MEDIA_VPS_ADDRESS,
"endpoint": f"46.225.230.72:{SASCHA_MEDIA_PORT}",
"keepalive": 15,
"listen": None,
},
}[role]
script = f'''import json, os, shutil, subprocess, time
+from pathlib import Path
+root = Path("/app-config/wireguard-media")
+config = root / "wg-media.conf"
+etc = Path("/etc/wireguard/wg-media.conf")
+backup_dir = root / "backups"
+backup_dir.mkdir(parents=True, exist_ok=True)
+private = (root / "private.key").read_text().strip()
+listen = {settings['listen']!r}
+existed = config.exists()
+backup = None
+if existed:
+ backup = backup_dir / ("wg-media.conf." + time.strftime("%Y%m%dT%H%M%SZ", time.gmtime()))
+ shutil.copy2(config, backup)
+lines = ["[Interface]", "Address = {settings['address']}", "MTU = {SASCHA_MEDIA_MTU}", "PrivateKey = " + private]
+if listen is not None: lines.append("ListenPort = " + str(listen))
+if {role!r} == "vps":
+ lines.extend(["PostUp = iptables -C INPUT -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT 2>/dev/null || iptables -I INPUT 1 -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT", "PreDown = iptables -D INPUT -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT 2>/dev/null || true"])
+lines.extend(["", "[Peer]", "PublicKey = {peer_public_key}", "AllowedIPs = {settings['peer']}"])
+if {settings['endpoint']!r}: lines.append("Endpoint = " + {settings['endpoint']!r})
+if {settings['keepalive']!r}: lines.append("PersistentKeepalive = " + str({settings['keepalive']!r}))
+candidate = root / "wgmtest.conf"
+candidate.write_text("\\n".join(lines) + "\\n")
+os.chmod(candidate, 0o600)
+check = subprocess.run(["wg-quick", "strip", str(candidate)], text=True, capture_output=True)
+if check.returncode != 0:
+ candidate.unlink(missing_ok=True)
+ raise RuntimeError(check.stderr.strip() or "wg-quick validation failed")
+os.replace(candidate, config)
+os.chmod(config, 0o600)
+etc.parent.mkdir(parents=True, exist_ok=True)
+if etc.is_symlink() or etc.exists():
+ if etc.is_symlink() and etc.resolve() == config.resolve(): pass
+ elif etc.exists():
+ etc_backup = backup_dir / ("etc-wg-media.conf." + time.strftime("%Y%m%dT%H%M%SZ", time.gmtime()))
+ shutil.move(etc, etc_backup)
+ etc.symlink_to(config)
+else: etc.symlink_to(config)
+proc = subprocess.run(["systemctl", "enable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30)
+if proc.returncode != 0:
+ if backup: shutil.copy2(backup, config)
+ else: config.unlink(missing_ok=True)
+ raise RuntimeError(proc.stderr.strip() or proc.stdout.strip() or "failed to start wg-media")
+print(json.dumps({{"role": {role!r}, "status": "configured", "address": {settings['address']!r}, "backup": str(backup) if backup else None, "existed": existed}}))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _media_tunnel_verify_command(source: str, destination: str, check_emby: bool = False) -> str:
extra = ''
if check_emby:
extra = '''\nhttp = subprocess.run(["curl", "-sS", "--max-time", "10", "-o", "/dev/null", "-w", "%{http_code}", "http://10.11.13.3:8096/System/Ping"], text=True, capture_output=True)\nresult["emby_http"] = http.stdout.strip() if http.returncode == 0 else "000"'''
script = f'''import json, subprocess, time
+result = {{"ping": False, "handshake": False}}
+for _ in range(10):
+ ping = subprocess.run(["ping", "-c", "1", "-W", "2", "-I", {source!r}, {destination!r}], text=True, capture_output=True)
+ hand = subprocess.run(["sudo", "-n", "wg", "show", "wg-media", "latest-handshakes"], text=True, capture_output=True)
+ now = int(time.time())
+ stamps = []
+ for line in hand.stdout.splitlines():
+ try: stamps.append(int(line.split()[-1]))
+ except Exception: pass
+ result["ping"] = ping.returncode == 0
+ result["handshake"] = any(stamp > 0 and now - stamp < 60 for stamp in stamps)
+ if result["ping"] and result["handshake"]: break
+ time.sleep(2)
+{extra}
+print(json.dumps(result))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _media_tunnel_rollback_command(remove_keys: bool) -> str:
script = f'''import json, shutil, subprocess
+from pathlib import Path
+root = Path("/app-config/wireguard-media")
+config = root / "wg-media.conf"
+backups = sorted((root / "backups").glob("wg-media.conf.*")) if (root / "backups").exists() else []
+subprocess.run(["systemctl", "disable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30)
+if backups:
+ shutil.copy2(backups[-1], config)
+ subprocess.run(["systemctl", "enable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30)
+ status = "restored"
+else:
+ config.unlink(missing_ok=True)
+ Path("/etc/wireguard/wg-media.conf").unlink(missing_ok=True)
+ if {remove_keys!r}:
+ (root / "private.key").unlink(missing_ok=True); (root / "public.key").unlink(missing_ok=True)
+ status = "removed"
+print(json.dumps({{"status": status}}))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _media_tunnel_key_cleanup_command() -> str:
script = '''import json
+from pathlib import Path
+root = Path("/app-config/wireguard-media")
+if not (root / "wg-media.conf").exists():
+ (root / "private.key").unlink(missing_ok=True)
+ (root / "public.key").unlink(missing_ok=True)
+ status = "new_keys_removed"
+else:
+ status = "kept_for_existing_config"
+print(json.dumps({"status": status}))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _media_host_target(name: str) -> str:
inventory = _find_inventory_host(name)
if not inventory:
raise RuntimeError(f"Inventory host missing: {name}")
return f'{inventory["user"]}@{inventory["ip"]}'
def _ssh_json(target: str, command: str, timeout: int = 45) -> dict:
rc, out, err = _ssh(target, command, timeout)
if rc != 0:
raise RuntimeError((err or out).strip()[-500:] or f"remote command failed on {target}")
try:
return json.loads(out)
except json.JSONDecodeError as exc:
raise RuntimeError(f"remote command returned invalid JSON on {target}") from exc
def _deploy_sascha_media_tunnel() -> dict:
vps = _media_host_target(SASCHA_MEDIA_VPS_HOST)
emby = _media_host_target(SASCHA_MEDIA_EMBY_HOST)
vps_key = _ssh_json(vps, _media_tunnel_key_command())
emby_key = _ssh_json(emby, _media_tunnel_key_command())
configured = []
try:
vps_install = _ssh_json(vps, _media_tunnel_install_command("vps", emby_key["public_key"]), 60)
configured.append((vps, bool(vps_key.get("created"))))
emby_install = _ssh_json(emby, _media_tunnel_install_command("emby", vps_key["public_key"]), 60)
configured.append((emby, bool(emby_key.get("created"))))
vps_check = _ssh_json(vps, _media_tunnel_verify_command("10.11.13.1", "10.11.13.3", True), 35)
emby_check = _ssh_json(emby, _media_tunnel_verify_command("10.11.13.3", "10.11.13.1"), 35)
if not (vps_check.get("ping") and vps_check.get("handshake") and emby_check.get("ping") and emby_check.get("handshake")):
raise RuntimeError("direct media tunnel verification failed")
if vps_check.get("emby_http") not in {"200", "401"}:
raise RuntimeError("Emby did not answer through the direct media tunnel")
return {
"status": "deployed", "interface": SASCHA_MEDIA_INTERFACE,
"vps_address": SASCHA_MEDIA_VPS_ADDRESS, "emby_address": SASCHA_MEDIA_EMBY_ADDRESS,
"listen_port": SASCHA_MEDIA_PORT, "mtu": SASCHA_MEDIA_MTU,
"handshake": True, "ping_vps_to_emby": True, "ping_emby_to_vps": True,
"emby_http": vps_check.get("emby_http"),
"rollback_backups": [vps_install.get("backup"), emby_install.get("backup")],
}
except Exception:
for target, created in reversed(configured):
try: _ssh_json(target, _media_tunnel_rollback_command(created), 60)
except Exception: pass
installed_targets = {target for target, _created in configured}
for target, key in ((vps, vps_key), (emby, emby_key)):
if target not in installed_targets and key.get("created"):
try: _ssh_json(target, _media_tunnel_key_cleanup_command(), 30)
except Exception: pass
raise
@app.post("/network/media-tunnel/sascha")
async def deploy_sascha_media_tunnel(req: SaschaMediaTunnelRequest, _=Depends(_verify)):
"""Deploy the fixed direct Hetzner-to-emby-sascha WireGuard media tunnel."""
plan = {
"status": "would_deploy", "interface": SASCHA_MEDIA_INTERFACE,
"vps_address": SASCHA_MEDIA_VPS_ADDRESS, "emby_address": SASCHA_MEDIA_EMBY_ADDRESS,
"listen_port": SASCHA_MEDIA_PORT, "mtu": SASCHA_MEDIA_MTU,
"allowed_ips": [SASCHA_MEDIA_VPS_ADDRESS, SASCHA_MEDIA_EMBY_ADDRESS],
"keeps_legacy_node6_path": True,
}
if req.dry_run:
_audit("/network/media-tunnel/sascha", "POST", 200, "dry_run=True", True)
return plan
if req.confirmation != SASCHA_MEDIA_CONFIRMATION:
raise HTTPException(400, f"confirmation must be {SASCHA_MEDIA_CONFIRMATION}")
try:
result = await asyncio.to_thread(_deploy_sascha_media_tunnel)
except Exception as exc:
_audit("/network/media-tunnel/sascha", "POST", 502, "deployment failed; rollback attempted")
raise HTTPException(502, str(exc)[-500:]) from exc
_audit("/network/media-tunnel/sascha", "POST", 200, "direct media tunnel deployed")
return result
def _media_benchmark_server_command() -> str:
server = '''import socket
+payload = b"\\0" * (1024 * 1024)
+with socket.socket() as listener:
+ listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
+ listener.bind(("0.0.0.0", 5209)); listener.listen(4); listener.settimeout(80)
+ for _ in range(2):
+ conn, _addr = listener.accept()
+ with conn:
+ conn.settimeout(30)
+ for _ in range(256): conn.sendall(payload)
+'''.replace("\n+", "\n")
encoded = base64.b64encode(server.encode()).decode()
command = f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
return (
"sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true; "
"sudo -n systemd-run --unit=butler-media-benchmark --collect --property=RuntimeMaxSec=90 "
f"/bin/sh -c {json.dumps(command)}"
)
def _media_benchmark_client_command() -> str:
script = '''import json, socket, time
+results = {}
+for name, host in (("legacy_node6", "10.6.1.103"), ("direct_wg_media", "10.11.13.3")):
+ total = 0; started = time.monotonic()
+ with socket.create_connection((host, 5209), timeout=10) as conn:
+ conn.settimeout(40)
+ while True:
+ chunk = conn.recv(1024 * 1024)
+ if not chunk: break
+ total += len(chunk)
+ elapsed = time.monotonic() - started
+ results[name] = {"bytes": total, "seconds": round(elapsed, 3), "mbit_s": round(total * 8 / elapsed / 1000000, 1)}
+print(json.dumps(results))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _benchmark_sascha_media_paths() -> dict:
vps = _media_host_target(SASCHA_MEDIA_VPS_HOST)
emby = _media_host_target(SASCHA_MEDIA_EMBY_HOST)
rc, out, err = _ssh(emby, _media_benchmark_server_command(), 30)
if rc != 0:
raise RuntimeError((err or out).strip()[-500:] or "failed to start benchmark server")
time.sleep(2)
try:
result = _ssh_json(vps, _media_benchmark_client_command(), 90)
finally:
_ssh(emby, "sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true", 20)
if any(item.get("bytes") != 256 * 1024 * 1024 for item in result.values()):
raise RuntimeError("benchmark transferred an unexpected byte count")
return result
@app.post("/network/media-tunnel/sascha/benchmark")
async def benchmark_sascha_media_paths(_=Depends(_verify)):
"""Compare the legacy node6 route with the direct WireGuard media path using fixed transient TCP streams."""
try:
result = await asyncio.to_thread(_benchmark_sascha_media_paths)
except Exception as exc:
_audit("/network/media-tunnel/sascha/benchmark", "POST", 502, "benchmark failed")
raise HTTPException(502, str(exc)[-500:]) from exc
_audit("/network/media-tunnel/sascha/benchmark", "POST", 200, "fixed 256 MiB TCP comparison")
return {"status": "completed", "direction": "emby-sascha_to_hetzner", "results": result}
SYSCTL_AUDIT_KEYS = (
"net.core.default_qdisc",
"net.core.rmem_default",
@ -2731,6 +3343,53 @@ async def _hetzner_zone_and_rrsets(zone_name: str):
return zone, rr_response.json().get("rrsets", []), headers
@app.put("/dns/rrset/{zone_name}/{record_name}")
async def dns_rrset_upsert(zone_name: str, record_name: str, req: DnsRrsetUpsertRequest, _=Depends(_verify)):
zone_name = zone_name.strip().lower().rstrip(".")
hostname = _validate_managed_hostname(f"{record_name}.{zone_name}")
record_name = hostname[: -(len(zone_name) + 1)]
record_type = req.record_type.upper()
value = req.value.strip().rstrip("." if record_type == "CNAME" else "")
try:
if record_type == "A" and ipaddress.ip_address(value).version != 4:
raise ValueError
if record_type == "AAAA" and ipaddress.ip_address(value).version != 6:
raise ValueError
if record_type == "CNAME":
_validate_managed_hostname(value)
except ValueError as exc:
raise HTTPException(400, f"Invalid {record_type} record value") from exc
zone, rrsets, headers = await _hetzner_zone_and_rrsets(zone_name)
existing = next((item for item in rrsets if item.get("name") == record_name and item.get("type") == record_type), None)
payload = {
"name": record_name,
"type": record_type,
"ttl": req.ttl,
"records": [{"value": value, "comment": req.comment}],
}
async with httpx.AsyncClient(timeout=30) as client:
if existing:
response = await client.put(
f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets/{record_name}/{record_type}",
headers=headers,
json=payload,
)
else:
response = await client.post(
f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets",
headers=headers,
json=payload,
)
if response.status_code not in {200, 201}:
raise HTTPException(502, f"Hetzner RRSet upsert failed: HTTP {response.status_code}")
verify = await client.get(f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets", headers=headers)
current = [item for item in verify.json().get("rrsets", []) if item.get("name") == record_name and item.get("type") == record_type] if verify.status_code == 200 else []
if not current or not any(record.get("value") == value for record in current[0].get("records", [])):
raise HTTPException(502, "RRSet read-back does not contain requested value")
_audit(f"/dns/rrset/{zone_name}/{record_name}", "PUT", 200, f"type={record_type} value={value}")
return {"zone": zone_name, "zone_id": zone.get("id"), "name": record_name, "created": existing is None, "rrset": current[0]}
@app.get("/dns/rrset/{zone_name}/{record_name}")
async def dns_rrset_get(zone_name: str, record_name: str, _=Depends(_verify)):
zone_name = zone_name.strip().lower().rstrip(".")
@ -3396,12 +4055,140 @@ async def tts_health(_=Depends(_verify)):
return results
def _sab_history_command() -> str:
container_script = r'''import json, re, subprocess, urllib.parse, urllib.request
text = open('/config/sabnzbd.ini', encoding='utf-8', errors='replace').read()
match = re.search(r'^api_key\s*=\s*(\S+)', text, re.M)
port_match = re.search(r'^port\s*=\s*(\d+)', text, re.M)
api_key = match.group(1) if match else ''
port = port_match.group(1) if port_match else '7777'
params = urllib.parse.urlencode({'mode': 'history', 'limit': 100, 'output': 'json', 'apikey': api_key})
with urllib.request.urlopen('http://127.0.0.1:' + port + '/api?' + params, timeout=20) as response:
data = json.load(response)
slots = data.get('history', {}).get('slots', [])
allowed = ('nzo_id', 'name', 'category', 'status', 'script', 'script_line', 'fail_message', 'completed', 'storage', 'path')
print(json.dumps([{key: item.get(key) for key in allowed} for item in slots]))
'''
container_encoded = base64.b64encode(container_script.encode()).decode()
host_script = f'''import subprocess, sys
command = ["sudo", "-n", "docker", "exec", "sabnzbd", "python3", "-c", "import base64;exec(base64.b64decode('{container_encoded}'))"]
proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
if proc.returncode != 0:
command = command[2:]
proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
if proc.returncode != 0:
print((proc.stderr or proc.stdout)[-500:], file=sys.stderr)
raise SystemExit(proc.returncode)
print(proc.stdout)
'''
encoded = base64.b64encode(host_script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
@app.get("/media/handoff/sab-history")
async def media_handoff_sab_history(_=Depends(_verify)):
"""Return sanitized SAB history without exposing the SAB API key."""
inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
if not inventory:
raise HTTPException(404, "Host sabnzbd not found")
target = f'{inventory["user"]}@{inventory["ip"]}'
rc, out, err = await asyncio.to_thread(_ssh, target, _sab_history_command(), 40)
if rc != 0:
raise HTTPException(502, (err or out).strip()[-500:] or "SAB history failed")
try:
return {"history": json.loads(out)}
except json.JSONDecodeError as exc:
raise HTTPException(502, "invalid SAB history response") from exc
def _media_handoff_diagnostics_command() -> str:
script = r'''import hashlib, json, os, subprocess, urllib.error, urllib.request
def run(args):
proc = subprocess.run(args, capture_output=True, text=True, timeout=20)
return proc.returncode, proc.stdout.strip(), proc.stderr.strip()
def docker_exec(command):
for prefix in (["sudo", "-n", "docker"], ["docker"]):
rc, out, err = run(prefix + ["exec", "sabnzbd", "sh", "-lc", command])
if rc == 0 or "not found" not in (err + out).lower():
return rc, out, err
return rc, out, err
inspect_rc, inspect_out, _ = run(["sudo", "-n", "docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
if inspect_rc != 0:
inspect_rc, inspect_out, _ = run(["docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
py_rc, py_out, _ = docker_exec("command -v python3")
curl_rc, curl_out, _ = docker_exec("command -v curl")
stat_rc, stat_out, _ = docker_exec("test -f /usenet/scripts/movetdarr.sh && stat -c '%a %s' /usenet/scripts/movetdarr.sh && sha256sum /usenet/scripts/movetdarr.sh")
log_rc, log_out, _ = docker_exec("tail -n 200 /usenet/scripts/postprocess.log 2>/dev/null || true")
source_rc, source_out, _ = docker_exec("find /usenet/complete -mindepth 2 -maxdepth 2 -type d ! -name '_UNPACK_*' -print 2>/dev/null | sort | tail -100")
target_rc, target_out, _ = docker_exec("find /tdarr/complete -mindepth 2 -maxdepth 2 -type d -print 2>/dev/null | sort | tail -100")
probe_code = 0
probe_body = ""
probe = urllib.request.Request(
"http://10.5.85.2:8888/media/handoff",
method="POST",
headers={"Content-Type": "application/json"},
data=b'{"action":"status","jobId":"diagnostic-probe"}',
)
try:
with urllib.request.urlopen(probe, timeout=10) as response:
probe_code = response.status
probe_body = response.read(500).decode(errors="replace")
except urllib.error.HTTPError as exc:
probe_code = exc.code
probe_body = exc.read(500).decode(errors="replace")
except Exception as exc:
probe_body = type(exc).__name__
script_info = {"exists": stat_rc == 0, "executable": False, "sha256": None, "mode": None, "size": None}
if stat_rc == 0:
lines = stat_out.splitlines()
if lines:
parts = lines[0].split()
if len(parts) >= 2:
script_info.update(mode=parts[0], size=int(parts[1]), executable=any(ch in parts[0][-3:] for ch in "1357"))
if len(lines) > 1:
script_info["sha256"] = lines[1].split()[0]
print(json.dumps({
"container_running": inspect_rc == 0 and inspect_out == "true",
"tools": {"python3": py_rc == 0 and bool(py_out), "curl": curl_rc == 0 and bool(curl_out)},
"script": script_info,
"caller_probe": {"http_status": probe_code, "body": probe_body},
"source_directories": source_out.splitlines() if source_rc == 0 else [],
"target_directories": target_out.splitlines() if target_rc == 0 else [],
"recent_log": log_out.splitlines()[-200:],
}))
'''
encoded = base64.b64encode(script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
@app.get("/media/handoff/diagnostics")
async def media_handoff_diagnostics(_=Depends(_verify)):
"""Read-only diagnostics for the fixed SABnzbd handoff script and caller path."""
inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
if not inventory:
raise HTTPException(404, "Host sabnzbd not found")
target = f'{inventory["user"]}@{inventory["ip"]}'
rc, out, err = await asyncio.to_thread(_ssh, target, _media_handoff_diagnostics_command(), 30)
if rc != 0:
raise HTTPException(502, (err or out).strip()[-500:] or "media handoff diagnostics failed")
try:
return json.loads(out)
except json.JSONDecodeError as exc:
raise HTTPException(502, "invalid media handoff diagnostic response") from exc
class MediaHandoffPayload(BaseModel):
action: Literal["start", "moved", "status", "fail"]
category: str | None = Field(None, max_length=32)
directory: str | None = Field(None, max_length=1000)
release: str | None = Field(None, max_length=500)
cleanName: str | None = Field(None, max_length=500)
expectedFiles: int | None = Field(None, ge=1, le=100)
jobId: str | None = Field(None, max_length=80)
reason: str | None = Field(None, max_length=300)
@ -3426,9 +4213,15 @@ 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"}
# 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")

View file

@ -1,6 +1,7 @@
import os
import asyncio
import json
import shlex
import time
from datetime import datetime, timedelta, timezone
@ -43,7 +44,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.4.1"
def test_media_handoff_proxies_strict_category_contract(monkeypatch):
@ -81,12 +82,14 @@ def test_media_handoff_proxies_strict_category_contract(monkeypatch):
"directory": "/usenet/complete/serien4k/Show.S01E01",
"release": "Show.S01E01-GRP",
"cleanName": "Show S01E01",
"expectedFiles": 2,
},
)
assert response.status_code == 200
assert response.json()["state"] == "registered"
assert captured["url"] == "http://n8n:5678/webhook/media-handoff"
assert captured["json"]["category"] == "serien4k"
assert captured["json"]["expectedFiles"] == 2
def test_media_handoff_rejects_untrusted_caller_before_proxy(monkeypatch):
@ -115,6 +118,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 = {}
@ -390,6 +486,9 @@ def test_capabilities_is_live_machine_readable_safety_map():
assert removal["mode"] == "mutation"
assert removal["dry_run"] is True
assert removal["critical"] is True
tunnel = by_operation[("POST", "/network/media-tunnel/sascha")]
assert tunnel["dry_run"] is True
assert tunnel["critical"] is True
assert by_operation[("DELETE", "/vm/destroy/{vmid}")]["dry_run"] is True
assert all(item["path"] != "/{service}/{path}" for item in payload["operations"])
assert payload["model_contract"]["instruction"].startswith("Prefer read_only")
@ -803,6 +902,115 @@ def test_uptime_monitor_remove_rejects_invalid_expected_name_before_ssh(monkeypa
assert response.status_code == 400
def test_media_verify_returns_duration_and_not_suspect_when_within_tolerance(monkeypatch):
monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
calls = []
def fake_ssh(host, command, timeout=30):
calls.append((host, command, timeout))
return 0, "2967.355000\n", ""
monkeypatch.setattr(app, "_ssh", fake_ssh)
with TestClient(app.app) as client:
response = client.post(
"/media/verify",
headers={"Authorization": "Bearer test-token"},
json={"host": "arrapps", "path": "/data/UHD/serien/Show/ep.mkv", "expected_minutes": 49.5},
)
assert response.status_code == 200
body = response.json()
assert body["duration_minutes"] == 49.46
assert body["suspect"] is False
assert calls[0][0] == "sascha@10.2.1.100"
assert "bazarrUHD" in calls[0][1]
assert "ffprobe" in calls[0][1]
def test_media_verify_flags_truncated_file_as_suspect(monkeypatch):
monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (0, "1836.0\n", ""))
with TestClient(app.app) as client:
response = client.post(
"/media/verify",
headers={"Authorization": "Bearer test-token"},
json={"host": "arrapps", "path": "/data/UHD/serien/Show/ep.mkv", "expected_minutes": 49.5},
)
assert response.status_code == 200
body = response.json()
assert body["suspect"] is True
assert body["deviation_pct"] > 20
def test_media_verify_rejects_path_outside_data_before_ssh(monkeypatch):
monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH")))
with TestClient(app.app) as client:
response = client.post(
"/media/verify",
headers={"Authorization": "Bearer test-token"},
json={"host": "arrapps", "path": "/etc/passwd"},
)
assert response.status_code == 400
def test_media_verify_allows_apostrophe_in_filename_and_quotes_it_safely(monkeypatch):
monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
calls = []
def fake_ssh(host, command, timeout=30):
calls.append((host, command, timeout))
return 0, "2967.0\n", ""
monkeypatch.setattr(app, "_ssh", fake_ssh)
path = "/data/FHD/serien/Star Trek - Strange New Worlds (2022)/Season 04/Once La'An a Time.mkv"
with TestClient(app.app) as client:
response = client.post(
"/media/verify",
headers={"Authorization": "Bearer test-token"},
json={"host": "arrapps", "path": path, "expected_minutes": 49.5},
)
assert response.status_code == 200
# shlex.quote must produce a command the remote shell parses as ONE argument,
# i.e. no unescaped apostrophe breaks out of quoting.
executed_cmd = calls[0][1]
quoted = shlex.quote(path)
assert quoted in executed_cmd
assert shlex.split(executed_cmd)[-1] == path
def test_media_verify_rejects_shell_metacharacters_before_ssh(monkeypatch):
monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH")))
with TestClient(app.app) as client:
response = client.post(
"/media/verify",
headers={"Authorization": "Bearer test-token"},
json={"host": "arrapps", "path": "/data/UHD/serien/a\x00.mkv"},
)
assert response.status_code == 400
def test_media_verify_rejects_unknown_host_before_ssh(monkeypatch):
monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH")))
with TestClient(app.app) as client:
response = client.post(
"/media/verify",
headers={"Authorization": "Bearer test-token"},
json={"host": "unknown-host", "path": "/data/UHD/serien/ep.mkv"},
)
assert response.status_code == 400
def test_media_verify_returns_502_on_unparsable_ffprobe_output(monkeypatch):
monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "sascha", "ip": "10.2.1.100"})
monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (0, "N/A\n", ""))
with TestClient(app.app) as client:
response = client.post(
"/media/verify",
headers={"Authorization": "Bearer test-token"},
json={"host": "arrapps", "path": "/data/UHD/serien/ep.mkv", "expected_minutes": 49.5},
)
assert response.status_code == 502
def test_invalid_log_target_is_rejected_before_ssh():
with TestClient(app.app) as client:
response = client.get(
@ -1272,3 +1480,124 @@ def test_speedtest_deploy_requires_strong_secrets_and_uses_full_git_app(monkeypa
assert deploy[1]["compose.yaml"] == "content:compose.yaml"
assert deploy[2] == "correct-horse-battery-staple"
assert deploy[3] == "streamscope-session-secret-with-entropy"
def test_sascha_media_edge_audit_uses_fixed_hetzner_host(monkeypatch):
payload = {
"hostname": "pfannkuchen",
"caddy": {"container_running": True, "protocols": ["h1", "h2", "h3"], "upstreams": ["10.6.1.103:8096"]},
"network": {"wg_media": False, "udp_51821": False},
}
calls = []
monkeypatch.setattr(app, "_find_inventory_host", lambda name: {"name": name, "user": "root", "ip": "46.225.230.72"})
monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=30: (calls.append((host, command, timeout)) or (0, json.dumps(payload), "")))
with TestClient(app.app) as client:
response = client.get("/media/edge/sascha", headers={"Authorization": "Bearer test-token"})
assert response.status_code == 200
assert response.json()["caddy"]["upstreams"] == ["10.6.1.103:8096"]
assert calls[0][0] == "root@46.225.230.72"
assert "PrivateKey" not in response.text
def test_sascha_media_tunnel_defaults_to_side_effect_free_dry_run(monkeypatch):
monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("dry-run must not SSH")))
with TestClient(app.app) as client:
response = client.post(
"/network/media-tunnel/sascha",
headers={"Authorization": "Bearer test-token"},
json={},
)
assert response.status_code == 200
body = response.json()
assert body["status"] == "would_deploy"
assert body["vps_address"] == "10.11.13.1/32"
assert body["emby_address"] == "10.11.13.3/32"
assert body["listen_port"] == 51821
def test_sascha_media_tunnel_requires_explicit_confirmation(monkeypatch):
monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("unconfirmed request must not SSH")))
with TestClient(app.app) as client:
response = client.post(
"/network/media-tunnel/sascha",
headers={"Authorization": "Bearer test-token"},
json={"dry_run": False},
)
assert response.status_code == 400
def test_sascha_media_tunnel_apply_returns_redacted_result(monkeypatch):
result = {
"status": "deployed",
"interface": "wg-media",
"vps_address": "10.11.13.1/32",
"emby_address": "10.11.13.3/32",
"listen_port": 51821,
"handshake": True,
"ping_vps_to_emby": True,
"ping_emby_to_vps": True,
}
monkeypatch.setattr(app, "_deploy_sascha_media_tunnel", lambda: result)
with TestClient(app.app) as client:
response = client.post(
"/network/media-tunnel/sascha",
headers={"Authorization": "Bearer test-token"},
json={"dry_run": False, "confirmation": "DEPLOY_DIRECT_SASCHA_MEDIA_TUNNEL"},
)
assert response.status_code == 200
assert response.json()["handshake"] is True
assert "private" not in response.text.lower()
def test_sascha_media_tunnel_candidate_uses_valid_wireguard_interface_name():
import base64
import re
command = app._media_tunnel_install_command("vps", "A" * 43 + "=")
encoded = re.search(r"b64decode\('([^']+)'\)", command).group(1)
script = base64.b64decode(encoded).decode()
assert 'root / "wgmtest.conf"' in script
assert len("wgmtest") <= 15
def test_sascha_media_path_benchmark_returns_both_fixed_routes(monkeypatch):
monkeypatch.setattr(app, "_benchmark_sascha_media_paths", lambda: {
"legacy_node6": {"bytes": 268435456, "seconds": 4.0, "mbit_s": 536.9},
"direct_wg_media": {"bytes": 268435456, "seconds": 3.0, "mbit_s": 715.8},
})
with TestClient(app.app) as client:
response = client.post(
"/network/media-tunnel/sascha/benchmark",
headers={"Authorization": "Bearer test-token"},
)
assert response.status_code == 200
assert response.json()["direction"] == "emby-sascha_to_hetzner"
assert set(response.json()["results"]) == {"legacy_node6", "direct_wg_media"}
def test_sascha_media_edge_optimize_is_dry_run_by_default(monkeypatch):
monkeypatch.setattr(app, "_optimize_sascha_media_edge", lambda: (_ for _ in ()).throw(AssertionError("dry-run must not mutate")))
with TestClient(app.app) as client:
response = client.post("/media/edge/sascha/optimize", headers={"Authorization": "Bearer test-token"}, json={})
assert response.status_code == 200
assert response.json()["protocols"] == ["h1", "h2"]
assert response.json()["upstream"] == "10.11.13.3:8096"
def test_sascha_media_edge_optimize_requires_confirmation(monkeypatch):
monkeypatch.setattr(app, "_optimize_sascha_media_edge", lambda: (_ for _ in ()).throw(AssertionError("must not mutate")))
with TestClient(app.app) as client:
response = client.post("/media/edge/sascha/optimize", headers={"Authorization": "Bearer test-token"}, json={"dry_run": False})
assert response.status_code == 400
def test_sascha_media_edge_optimize_applies_fixed_safe_result(monkeypatch):
monkeypatch.setattr(app, "_optimize_sascha_media_edge", lambda: {"status": "optimized", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "backup": "/app-config/caddy/backup", "sha256": "a" * 64})
with TestClient(app.app) as client:
response = client.post(
"/media/edge/sascha/optimize",
headers={"Authorization": "Bearer test-token"},
json={"dry_run": False, "confirmation": "OPTIMIZE_TV_SASCHA_LUTZ_DE"},
)
assert response.status_code == 200
assert response.json()["status"] == "optimized"