From f9664d5e36b057b544f8c18ed49a88c76660d7d3 Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:06:07 +0200 Subject: [PATCH 1/2] feat: benchmark direct media path (app.py) --- app.py | 81 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 79 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index 1f63730..b93e8d6 100644 --- a/app.py +++ b/app.py @@ -12,7 +12,7 @@ from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") -VERSION = "2.3.9" +VERSION = "2.4.0" API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -2272,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": []}, "network": {}} ++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"]) @@ -2316,6 +2316,15 @@ 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 @@ -2579,6 +2588,74 @@ 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", -- 2.49.1 From fa01dd77fb86f421fbdedcc3065b84135b1b167d Mon Sep 17 00:00:00 2001 From: sascha Date: Sat, 5 Sep 2026 13:06:08 +0200 Subject: [PATCH 2/2] feat: benchmark direct media path (tests/test_app.py) --- tests/test_app.py | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/tests/test_app.py b/tests/test_app.py index 4552a81..a98a293 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -43,7 +43,7 @@ def test_health_exposes_current_version(): with TestClient(app.app) as client: response = client.get("/health") assert response.status_code == 200 - assert response.json()["version"] == app.VERSION == "2.3.9" + assert response.json()["version"] == app.VERSION == "2.4.0" def test_media_handoff_proxies_strict_category_contract(monkeypatch): @@ -1355,3 +1355,18 @@ def test_sascha_media_tunnel_candidate_uses_valid_wireguard_interface_name(): 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"} -- 2.49.1