Compare commits

..

No commits in common. "main" and "fix/sascha-media-tunnel-validation-20260905" have entirely different histories.

3 changed files with 7 additions and 735 deletions

View file

@ -45,7 +45,6 @@ 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
@ -98,14 +97,6 @@ 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.

474
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, shlex
import os, json, asyncio, logging, time, base64, re, subprocess, ipaddress, secrets, sqlite3, math, hashlib
from datetime import datetime, timezone
import httpx, yaml
from typing import Literal
@ -12,17 +12,14 @@ from contextlib import asynccontextmanager
from contextvars import ContextVar
log = logging.getLogger("butler")
VERSION = "2.4.9"
VERSION = "2.3.8"
API_DIR = os.environ.get("API_KEY_DIR", "/data/api")
VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache")
BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "")
CONFIG_PATH = os.environ.get("BUTLER_CONFIG", "/data/butler.yaml")
UI_PATH = os.environ.get("BUTLER_UI_PATH", os.path.join(os.path.dirname(__file__), "ui.html"))
MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get(
"MEDIA_HANDOFF_ALLOWED_NETWORKS",
"10.2.1.119/32,10.5.85.12/32",
)
MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get("MEDIA_HANDOFF_ALLOWED_NETWORKS", "10.2.1.119/32")
# --- Config loading ---
@ -1546,102 +1543,6 @@ 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"
@ -1670,13 +1571,6 @@ 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)
@ -2370,11 +2264,6 @@ class SaschaMediaTunnelRequest(BaseModel):
confirmation: str | None = None
class SaschaMediaEdgeOptimizeRequest(BaseModel):
dry_run: bool = True
confirmation: str | None = None
def _sascha_media_edge_audit_command() -> str:
script = '''import json, re, subprocess
+from pathlib import Path
@ -2383,7 +2272,7 @@ def _sascha_media_edge_audit_command() -> str:
+ proc = subprocess.run(args, text=True, capture_output=True, timeout=30)
+ return {"rc": proc.returncode, "stdout": proc.stdout.strip(), "stderr": proc.stderr.strip()[-300:]}
+
+result = {"hostname": "pfannkuchen", "caddy": {"container_running": False, "protocols": [], "protocols_explicit": False, "upstreams": [], "site_block": [], "emby_snippet": []}, "network": {}}
+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"])
@ -2427,15 +2316,6 @@ def _sascha_media_edge_audit_command() -> str:
+ else: selected.append("[REDACTED SENSITIVE DIRECTIVE]")
+ if depth == 0: break
+ result["caddy"]["site_block"] = selected
+ collecting = False; depth = 0; selected = []
+ for line in lines:
+ if not collecting and re.match(r"^\\s*\\(emby_config\\)\\s*\\{", line): collecting = True
+ if collecting:
+ depth += line.count("{") - line.count("}")
+ if not re.search(r"(?i)(password|secret|token|private|basicauth|basic_auth|hash)", line): selected.append(line.strip())
+ else: selected.append("[REDACTED SENSITIVE DIRECTIVE]")
+ if depth == 0: break
+ result["caddy"]["emby_snippet"] = selected
+
+result["network"]["wg_media_active"] = subprocess.run(["systemctl", "is-active", "--quiet", "wg-quick@wg-media"]).returncode == 0
+result["network"]["wg_media_enabled"] = subprocess.run(["systemctl", "is-enabled", "--quiet", "wg-quick@wg-media"]).returncode == 0
@ -2466,100 +2346,6 @@ async def sascha_media_edge_audit(_=Depends(_verify)):
return result
def _sascha_media_edge_optimize_command() -> str:
script = '''import hashlib, json, os, re, shutil, subprocess, time
+from pathlib import Path
+path = Path("/app-config/caddy/Caddyfile")
+text = path.read_text(encoding="utf-8")
+old_pattern = re.compile(r"(?m)^(\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+)10\\.6\\.1\\.103:8096(\\s*)$")
+new_pattern = re.compile(r"(?m)^\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+10\\.11\\.13\\.3:8096\\s*$")
+old_count = len(old_pattern.findall(text)); new_count = len(new_pattern.findall(text))
+if old_count == 1:
+ text = old_pattern.sub(r"\\g<1>10.11.13.3:8096\\g<2>", text)
+elif not (old_count == 0 and new_count == 1):
+ raise RuntimeError("expected exactly one tv.sascha-lutz.de upstream")
+lines = text.splitlines(keepends=True)
+first = next((i for i, line in enumerate(lines) if line.strip() and not line.lstrip().startswith("#")), None)
+if first is None or lines[first].strip() != "{":
+ lines[0:0] = ["{\\n", " servers {\\n", " protocols h1 h2\\n", " }\\n", "}\\n", "\\n"]
+else:
+ depth = 0; global_end = None; servers_start = None; servers_end = None
+ for i in range(first, len(lines)):
+ stripped = lines[i].strip(); before = depth
+ if before == 1 and re.match(r"^servers(?:\\s+\\S+)?\\s*\\{$", stripped): servers_start = i
+ depth += lines[i].count("{") - lines[i].count("}")
+ if servers_start is not None and i > servers_start and depth == 1 and servers_end is None: servers_end = i
+ if i > first and depth == 0: global_end = i; break
+ if global_end is None: raise RuntimeError("unbalanced Caddy global options block")
+ if servers_start is None:
+ lines[global_end:global_end] = [" servers {\\n", " protocols h1 h2\\n", " }\\n"]
+ else:
+ if servers_end is None: raise RuntimeError("unbalanced Caddy servers block")
+ protocol_lines = [i for i in range(servers_start + 1, servers_end) if re.match(r"^\\s*protocols\\s+", lines[i])]
+ if len(protocol_lines) > 1: raise RuntimeError("multiple Caddy protocol directives")
+ if protocol_lines: lines[protocol_lines[0]] = re.sub(r"protocols\\s+.*", "protocols h1 h2", lines[protocol_lines[0]])
+ else: lines.insert(servers_start + 1, " protocols h1 h2\\n")
+text = "".join(lines)
+if re.search(r"(?im)^\\s*header(?:_down)?\\s+Alt-Svc", text): raise RuntimeError("manual Alt-Svc directive requires review")
+stamp = time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())
+backup = path.with_name("Caddyfile.pre-sascha-media-" + stamp)
+candidate = path.with_name("Caddyfile.sascha-media-candidate")
+shutil.copy2(path, backup); candidate.write_text(text, encoding="utf-8"); os.chmod(candidate, path.stat().st_mode)
+def run(args, timeout=30):
+ proc = subprocess.run(args, text=True, capture_output=True, timeout=timeout)
+ if proc.returncode != 0: raise RuntimeError((proc.stderr or proc.stdout).strip()[-500:] or "command failed")
+ return proc.stdout.strip()
+try:
+ run(["docker", "cp", str(candidate), "caddy:/tmp/Caddyfile.sascha-media-candidate"])
+ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/tmp/Caddyfile.sascha-media-candidate"])
+ with path.open("w", encoding="utf-8") as handle:
+ handle.write(text); handle.flush(); os.fsync(handle.fileno())
+ host_hash = hashlib.sha256(path.read_bytes()).hexdigest()
+ container_hash = run(["docker", "exec", "caddy", "sha256sum", "/etc/caddy/Caddyfile"]).split()[0]
+ if host_hash != container_hash: raise RuntimeError("host/container Caddyfile hash mismatch")
+ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/etc/caddy/Caddyfile"])
+ run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"])
+ adapted = json.loads(run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"]))
+ protocols = []
+ for server in adapted.get("apps", {}).get("http", {}).get("servers", {}).values(): protocols.extend(server.get("protocols", []))
+ if sorted(set(protocols)) != ["h1", "h2"]: raise RuntimeError("Caddy did not load h1+h2-only protocols")
+ if not run(["curl", "-fsS", "--max-time", "15", "http://10.11.13.3:8096/System/Ping"]): raise RuntimeError("Emby direct-path ping was empty")
+except Exception:
+ with path.open("w", encoding="utf-8") as handle:
+ handle.write(backup.read_text(encoding="utf-8")); handle.flush(); os.fsync(handle.fileno())
+ subprocess.run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"], text=True, capture_output=True, timeout=30)
+ raise
+finally:
+ candidate.unlink(missing_ok=True)
+ subprocess.run(["docker", "exec", "caddy", "rm", "-f", "/tmp/Caddyfile.sascha-media-candidate"], text=True, capture_output=True)
+print(json.dumps({"status": "optimized", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "backup": str(backup), "sha256": host_hash}))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _optimize_sascha_media_edge() -> dict:
return _ssh_json(_media_host_target(SASCHA_MEDIA_VPS_HOST), _sascha_media_edge_optimize_command(), 90)
@app.post("/media/edge/sascha/optimize")
async def optimize_sascha_media_edge(req: SaschaMediaEdgeOptimizeRequest, _=Depends(_verify)):
"""Switch only tv.sascha-lutz.de to wg-media and disable HTTP/3 on its dedicated Hetzner edge."""
plan = {"status": "would_optimize", "hostname": "tv.sascha-lutz.de", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "zero_downtime_reload": True}
if req.dry_run:
_audit("/media/edge/sascha/optimize", "POST", 200, "dry_run=True", True)
return plan
if req.confirmation != "OPTIMIZE_TV_SASCHA_LUTZ_DE":
raise HTTPException(400, "confirmation must be OPTIMIZE_TV_SASCHA_LUTZ_DE")
try:
result = await asyncio.to_thread(_optimize_sascha_media_edge)
except Exception as exc:
_audit("/media/edge/sascha/optimize", "POST", 502, "optimization failed; Caddy rollback attempted")
raise HTTPException(502, str(exc)[-500:]) from exc
_audit("/media/edge/sascha/optimize", "POST", 200, "direct upstream and h1+h2 enabled")
return result
def _media_tunnel_key_command() -> str:
script = '''import json, os, subprocess
+from pathlib import Path
@ -2623,7 +2409,7 @@ def _media_tunnel_install_command(role: Literal["vps", "emby"], peer_public_key:
+lines.extend(["", "[Peer]", "PublicKey = {peer_public_key}", "AllowedIPs = {settings['peer']}"])
+if {settings['endpoint']!r}: lines.append("Endpoint = " + {settings['endpoint']!r})
+if {settings['keepalive']!r}: lines.append("PersistentKeepalive = " + str({settings['keepalive']!r}))
+candidate = root / "wgmtest.conf"
+candidate = root / "wg-media-candidate.conf"
+candidate.write_text("\\n".join(lines) + "\\n")
+os.chmod(candidate, 0o600)
+check = subprocess.run(["wg-quick", "strip", str(candidate)], text=True, capture_output=True)
@ -2793,74 +2579,6 @@ async def deploy_sascha_media_tunnel(req: SaschaMediaTunnelRequest, _=Depends(_v
return result
def _media_benchmark_server_command() -> str:
server = '''import socket
+payload = b"\\0" * (1024 * 1024)
+with socket.socket() as listener:
+ listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
+ listener.bind(("0.0.0.0", 5209)); listener.listen(4); listener.settimeout(80)
+ for _ in range(2):
+ conn, _addr = listener.accept()
+ with conn:
+ conn.settimeout(30)
+ for _ in range(256): conn.sendall(payload)
+'''.replace("\n+", "\n")
encoded = base64.b64encode(server.encode()).decode()
command = f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
return (
"sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true; "
"sudo -n systemd-run --unit=butler-media-benchmark --collect --property=RuntimeMaxSec=90 "
f"/bin/sh -c {json.dumps(command)}"
)
def _media_benchmark_client_command() -> str:
script = '''import json, socket, time
+results = {}
+for name, host in (("legacy_node6", "10.6.1.103"), ("direct_wg_media", "10.11.13.3")):
+ total = 0; started = time.monotonic()
+ with socket.create_connection((host, 5209), timeout=10) as conn:
+ conn.settimeout(40)
+ while True:
+ chunk = conn.recv(1024 * 1024)
+ if not chunk: break
+ total += len(chunk)
+ elapsed = time.monotonic() - started
+ results[name] = {"bytes": total, "seconds": round(elapsed, 3), "mbit_s": round(total * 8 / elapsed / 1000000, 1)}
+print(json.dumps(results))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _benchmark_sascha_media_paths() -> dict:
vps = _media_host_target(SASCHA_MEDIA_VPS_HOST)
emby = _media_host_target(SASCHA_MEDIA_EMBY_HOST)
rc, out, err = _ssh(emby, _media_benchmark_server_command(), 30)
if rc != 0:
raise RuntimeError((err or out).strip()[-500:] or "failed to start benchmark server")
time.sleep(2)
try:
result = _ssh_json(vps, _media_benchmark_client_command(), 90)
finally:
_ssh(emby, "sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true", 20)
if any(item.get("bytes") != 256 * 1024 * 1024 for item in result.values()):
raise RuntimeError("benchmark transferred an unexpected byte count")
return result
@app.post("/network/media-tunnel/sascha/benchmark")
async def benchmark_sascha_media_paths(_=Depends(_verify)):
"""Compare the legacy node6 route with the direct WireGuard media path using fixed transient TCP streams."""
try:
result = await asyncio.to_thread(_benchmark_sascha_media_paths)
except Exception as exc:
_audit("/network/media-tunnel/sascha/benchmark", "POST", 502, "benchmark failed")
raise HTTPException(502, str(exc)[-500:]) from exc
_audit("/network/media-tunnel/sascha/benchmark", "POST", 200, "fixed 256 MiB TCP comparison")
return {"status": "completed", "direction": "emby-sascha_to_hetzner", "results": result}
SYSCTL_AUDIT_KEYS = (
"net.core.default_qdisc",
"net.core.rmem_default",
@ -3343,53 +3061,6 @@ 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(".")
@ -4055,133 +3726,6 @@ 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)
@ -4213,15 +3757,9 @@ 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")
# 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"}
allowed_categories = {"serien4k", "serien", "video4k", "video"}
if payload.action == "start":
if payload.category not in allowed_categories:
raise HTTPException(422, "Unsupported media category")

View file

@ -1,7 +1,6 @@
import os
import asyncio
import json
import shlex
import time
from datetime import datetime, timedelta, timezone
@ -44,7 +43,7 @@ def test_health_exposes_current_version():
with TestClient(app.app) as client:
response = client.get("/health")
assert response.status_code == 200
assert response.json()["version"] == app.VERSION == "2.4.1"
assert response.json()["version"] == app.VERSION == "2.3.8"
def test_media_handoff_proxies_strict_category_contract(monkeypatch):
@ -118,99 +117,6 @@ 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 = {}
@ -902,115 +808,6 @@ 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(
@ -1547,57 +1344,3 @@ def test_sascha_media_tunnel_apply_returns_redacted_result(monkeypatch):
assert response.status_code == 200
assert response.json()["handshake"] is True
assert "private" not in response.text.lower()
def test_sascha_media_tunnel_candidate_uses_valid_wireguard_interface_name():
import base64
import re
command = app._media_tunnel_install_command("vps", "A" * 43 + "=")
encoded = re.search(r"b64decode\('([^']+)'\)", command).group(1)
script = base64.b64decode(encoded).decode()
assert 'root / "wgmtest.conf"' in script
assert len("wgmtest") <= 15
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"