Compare commits
No commits in common. "main" and "fix/paperless-import-result-logs" have entirely different histories.
main
...
fix/paperl
3 changed files with 4 additions and 1270 deletions
|
|
@ -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/inspect/{host}/{container}` | GET | Sanitized image, runtime, resources, mounts and state |
|
||||||
| `/docker/restart/{host}/{container}` | POST | Restart container; supports `dry_run=true` |
|
| `/docker/restart/{host}/{container}` | POST | Restart container; supports `dry_run=true` |
|
||||||
| `/config/reload` | POST | Reload YAML configuration and credential cache |
|
| `/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
|
## VM lifecycle and inventory
|
||||||
|
|
||||||
|
|
@ -98,14 +97,6 @@ Integration Compose definition: `tests/compose.integration.yaml` (binds only to
|
||||||
|
|
||||||
## Changelog
|
## 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
|
### 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.
|
- Make Vaultwarden refresh durable: persistent named cache volume, protected runtime credentials, automatic API-key re-login and atomic cache writes.
|
||||||
|
|
|
||||||
865
app.py
865
app.py
|
|
@ -1,7 +1,7 @@
|
||||||
"""Homelab Butler v2.1 – Unified API proxy for Pfannkuchen homelab.
|
"""Homelab Butler v2.1 – Unified API proxy for Pfannkuchen homelab.
|
||||||
Reads service config from butler.yaml, credentials from Vaultwarden cache with flat-file fallback."""
|
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
|
from datetime import datetime, timezone
|
||||||
import httpx, yaml
|
import httpx, yaml
|
||||||
from typing import Literal
|
from typing import Literal
|
||||||
|
|
@ -12,17 +12,13 @@ from contextlib import asynccontextmanager
|
||||||
from contextvars import ContextVar
|
from contextvars import ContextVar
|
||||||
|
|
||||||
log = logging.getLogger("butler")
|
log = logging.getLogger("butler")
|
||||||
VERSION = "2.4.9"
|
VERSION = "2.3.5"
|
||||||
|
|
||||||
API_DIR = os.environ.get("API_KEY_DIR", "/data/api")
|
API_DIR = os.environ.get("API_KEY_DIR", "/data/api")
|
||||||
VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache")
|
VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache")
|
||||||
BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "")
|
BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "")
|
||||||
CONFIG_PATH = os.environ.get("BUTLER_CONFIG", "/data/butler.yaml")
|
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"))
|
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",
|
|
||||||
)
|
|
||||||
|
|
||||||
# --- Config loading ---
|
# --- Config loading ---
|
||||||
|
|
||||||
|
|
@ -753,7 +749,7 @@ async def capabilities(_=Depends(_verify)):
|
||||||
mode = "read_only" if method.upper() == "GET" else "mutation"
|
mode = "read_only" if method.upper() == "GET" else "mutation"
|
||||||
serialized = {"parameters": operation.get("parameters", []), "requestBody": operation.get("requestBody", {})}
|
serialized = {"parameters": operation.get("parameters", []), "requestBody": operation.get("requestBody", {})}
|
||||||
dry_run = _schema_contains_property(serialized, "dry_run", components)
|
dry_run = _schema_contains_property(serialized, "dry_run", components)
|
||||||
critical = path.startswith(("/network/wireguard", "/network/media-tunnel", "/caddy/"))
|
critical = path.startswith(("/network/wireguard", "/caddy/"))
|
||||||
destructive = method.upper() == "DELETE" or any(
|
destructive = method.upper() == "DELETE" or any(
|
||||||
marker in path for marker in ("/destroy/", "/cleanup/", "/break-lock/", "/restore/")
|
marker in path for marker in ("/destroy/", "/cleanup/", "/break-lock/", "/restore/")
|
||||||
)
|
)
|
||||||
|
|
@ -1546,102 +1542,6 @@ async def vault_reload(_=Depends(_verify)):
|
||||||
return {"reloaded": True, "items": len(_vault_cache)}
|
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 reverse-proxy and DNS management ---
|
||||||
|
|
||||||
VPS_SSH = "root@46.225.230.72"
|
VPS_SSH = "root@46.225.230.72"
|
||||||
|
|
@ -1670,13 +1570,6 @@ class ProxyRouteRequest(BaseModel):
|
||||||
dns_token: str | None = None
|
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]:
|
def _remote_python(script: str, timeout: int = 30) -> tuple[int, str, str]:
|
||||||
encoded = base64.b64encode(script.encode()).decode()
|
encoded = base64.b64encode(script.encode()).decode()
|
||||||
return _ssh(VPS_SSH, f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"', timeout=timeout)
|
return _ssh(VPS_SSH, f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"', timeout=timeout)
|
||||||
|
|
@ -2355,512 +2248,6 @@ async def network_wireguard_remove_peer(host: str, req: WireGuardPeerRemoveReque
|
||||||
return {"host": host, **result}
|
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 = (
|
SYSCTL_AUDIT_KEYS = (
|
||||||
"net.core.default_qdisc",
|
"net.core.default_qdisc",
|
||||||
"net.core.rmem_default",
|
"net.core.rmem_default",
|
||||||
|
|
@ -3343,53 +2730,6 @@ async def _hetzner_zone_and_rrsets(zone_name: str):
|
||||||
return zone, rr_response.json().get("rrsets", []), headers
|
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}")
|
@app.get("/dns/rrset/{zone_name}/{record_name}")
|
||||||
async def dns_rrset_get(zone_name: str, record_name: str, _=Depends(_verify)):
|
async def dns_rrset_get(zone_name: str, record_name: str, _=Depends(_verify)):
|
||||||
zone_name = zone_name.strip().lower().rstrip(".")
|
zone_name = zone_name.strip().lower().rstrip(".")
|
||||||
|
|
@ -4055,205 +3395,6 @@ async def tts_health(_=Depends(_verify)):
|
||||||
return results
|
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)
|
|
||||||
|
|
||||||
|
|
||||||
def _media_handoff_caller_allowed(request: Request) -> bool:
|
|
||||||
auth = request.headers.get("authorization", "")
|
|
||||||
if BUTLER_TOKEN and secrets.compare_digest(auth, f"Bearer {BUTLER_TOKEN}"):
|
|
||||||
return True
|
|
||||||
try:
|
|
||||||
caller = ipaddress.ip_address(request.client.host if request.client else "")
|
|
||||||
networks = [
|
|
||||||
ipaddress.ip_network(value.strip(), strict=False)
|
|
||||||
for value in MEDIA_HANDOFF_ALLOWED_NETWORKS.split(",")
|
|
||||||
if value.strip()
|
|
||||||
]
|
|
||||||
except ValueError:
|
|
||||||
return False
|
|
||||||
return any(caller in network for network in networks)
|
|
||||||
|
|
||||||
|
|
||||||
@app.post("/media/handoff")
|
|
||||||
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"}
|
|
||||||
if payload.action == "start":
|
|
||||||
if payload.category not in allowed_categories:
|
|
||||||
raise HTTPException(422, "Unsupported media category")
|
|
||||||
expected_prefix = f"/usenet/complete/{payload.category}/"
|
|
||||||
if not payload.directory or not payload.directory.startswith(expected_prefix):
|
|
||||||
raise HTTPException(422, "Invalid media handoff directory")
|
|
||||||
if not payload.release:
|
|
||||||
raise HTTPException(422, "Release name is required")
|
|
||||||
else:
|
|
||||||
if not payload.jobId or not re.fullmatch(r"[a-z0-9-]{8,80}", payload.jobId):
|
|
||||||
raise HTTPException(422, "Valid jobId is required")
|
|
||||||
|
|
||||||
cfg = SERVICES.get("n8n")
|
|
||||||
if not cfg or not cfg.get("url"):
|
|
||||||
raise HTTPException(503, "n8n service is not configured")
|
|
||||||
target = f"{cfg['url'].rstrip('/')}/webhook/media-handoff"
|
|
||||||
body = payload.model_dump(exclude_none=True) if hasattr(payload, "model_dump") else payload.dict(exclude_none=True)
|
|
||||||
try:
|
|
||||||
async with httpx.AsyncClient(verify=False, timeout=15) as client:
|
|
||||||
response = await client.post(target, json=body, headers={"Content-Type": "application/json"})
|
|
||||||
except httpx.HTTPError as exc:
|
|
||||||
_audit("/media/handoff", "POST", 502, f"action={payload.action} error={type(exc).__name__}")
|
|
||||||
raise HTTPException(502, "n8n media handoff is unavailable") from exc
|
|
||||||
|
|
||||||
try:
|
|
||||||
result = response.json()
|
|
||||||
except Exception:
|
|
||||||
result = {"ok": False, "error": "invalid_n8n_response"}
|
|
||||||
_audit("/media/handoff", "POST", response.status_code, f"action={payload.action}")
|
|
||||||
return JSONResponse(content=result, status_code=response.status_code)
|
|
||||||
|
|
||||||
|
|
||||||
@app.api_route("/{service}/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"])
|
@app.api_route("/{service}/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"])
|
||||||
async def proxy(service: str, path: str, request: Request, _=Depends(_verify)):
|
async def proxy(service: str, path: str, request: Request, _=Depends(_verify)):
|
||||||
SKIP_SERVICES = {"vm", "inventory", "ansible", "debug", "tts", "status", "audit", "config"}
|
SKIP_SERVICES = {"vm", "inventory", "ansible", "debug", "tts", "status", "audit", "config"}
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,6 @@
|
||||||
import os
|
import os
|
||||||
import asyncio
|
import asyncio
|
||||||
import json
|
import json
|
||||||
import shlex
|
|
||||||
import time
|
import time
|
||||||
from datetime import datetime, timedelta, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
|
|
||||||
|
|
@ -44,171 +43,7 @@ def test_health_exposes_current_version():
|
||||||
with TestClient(app.app) as client:
|
with TestClient(app.app) as client:
|
||||||
response = client.get("/health")
|
response = client.get("/health")
|
||||||
assert response.status_code == 200
|
assert response.status_code == 200
|
||||||
assert response.json()["version"] == app.VERSION == "2.4.1"
|
assert response.json()["version"] == app.VERSION == "2.3.5"
|
||||||
|
|
||||||
|
|
||||||
def test_media_handoff_proxies_strict_category_contract(monkeypatch):
|
|
||||||
captured = {}
|
|
||||||
|
|
||||||
class FakeResponse:
|
|
||||||
status_code = 200
|
|
||||||
|
|
||||||
def json(self):
|
|
||||||
return {"ok": True, "jobId": "test-job-1234", "state": "registered"}
|
|
||||||
|
|
||||||
class FakeClient:
|
|
||||||
def __init__(self, **kwargs):
|
|
||||||
captured["client"] = kwargs
|
|
||||||
|
|
||||||
async def __aenter__(self):
|
|
||||||
return self
|
|
||||||
|
|
||||||
async def __aexit__(self, *_args):
|
|
||||||
return None
|
|
||||||
|
|
||||||
async def post(self, url, json, headers):
|
|
||||||
captured.update(url=url, json=json, headers=headers)
|
|
||||||
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"}})
|
|
||||||
response = client.post(
|
|
||||||
"/media/handoff",
|
|
||||||
headers={"Authorization": "Bearer test-token"},
|
|
||||||
json={
|
|
||||||
"action": "start",
|
|
||||||
"category": "serien4k",
|
|
||||||
"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):
|
|
||||||
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",
|
|
||||||
json={"action": "status", "jobId": "test-job-1234"},
|
|
||||||
)
|
|
||||||
assert response.status_code == 403
|
|
||||||
|
|
||||||
|
|
||||||
def test_media_handoff_rejects_wrong_category_path(monkeypatch):
|
|
||||||
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": "video4k",
|
|
||||||
"directory": "/usenet/complete/serien4k/Wrong",
|
|
||||||
"release": "Wrong",
|
|
||||||
},
|
|
||||||
)
|
|
||||||
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):
|
def test_paperless_import_queues_pdf_through_butler(monkeypatch):
|
||||||
|
|
@ -486,9 +321,6 @@ def test_capabilities_is_live_machine_readable_safety_map():
|
||||||
assert removal["mode"] == "mutation"
|
assert removal["mode"] == "mutation"
|
||||||
assert removal["dry_run"] is True
|
assert removal["dry_run"] is True
|
||||||
assert removal["critical"] 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 by_operation[("DELETE", "/vm/destroy/{vmid}")]["dry_run"] is True
|
||||||
assert all(item["path"] != "/{service}/{path}" for item in payload["operations"])
|
assert all(item["path"] != "/{service}/{path}" for item in payload["operations"])
|
||||||
assert payload["model_contract"]["instruction"].startswith("Prefer read_only")
|
assert payload["model_contract"]["instruction"].startswith("Prefer read_only")
|
||||||
|
|
@ -902,115 +734,6 @@ def test_uptime_monitor_remove_rejects_invalid_expected_name_before_ssh(monkeypa
|
||||||
assert response.status_code == 400
|
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():
|
def test_invalid_log_target_is_rejected_before_ssh():
|
||||||
with TestClient(app.app) as client:
|
with TestClient(app.app) as client:
|
||||||
response = client.get(
|
response = client.get(
|
||||||
|
|
@ -1480,124 +1203,3 @@ def test_speedtest_deploy_requires_strong_secrets_and_uses_full_git_app(monkeypa
|
||||||
assert deploy[1]["compose.yaml"] == "content:compose.yaml"
|
assert deploy[1]["compose.yaml"] == "content:compose.yaml"
|
||||||
assert deploy[2] == "correct-horse-battery-staple"
|
assert deploy[2] == "correct-horse-battery-staple"
|
||||||
assert deploy[3] == "streamscope-session-secret-with-entropy"
|
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"
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue