From cb54a7977a9fa52562bc9279539c34814b5e4e52 Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 13 Aug 2026 13:09:07 +0200 Subject: [PATCH 01/72] fix: make bridge secrets readable by service uid --- app.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/app.py b/app.py index 6a36206..7aa0419 100644 --- a/app.py +++ b/app.py @@ -1914,7 +1914,8 @@ for filename, value in files.items(): if path.is_dir(): path.rmdir() path.write_text(value) - os.chmod(path, 0o600) + os.chown(path, 10001, 10001) + os.chmod(path, 0o400) """.format(files_json=json.dumps(files)) encoded = base64.b64encode(installer.encode()).decode() command = f"sudo python3 -c {__import__('shlex').quote(f'import base64;exec(base64.b64decode({encoded!r}))')}" From dedde311f1f099dffcfd629d8e3ddda193bb96c7 Mon Sep 17 00:00:00 2001 From: sascha Date: Thu, 13 Aug 2026 17:55:26 +0200 Subject: [PATCH 02/72] Remove obsolete BW Manager deployment endpoint --- app.py | 111 --------------------------------------------------------- 1 file changed, 111 deletions(-) diff --git a/app.py b/app.py index 7aa0419..0e2999c 100644 --- a/app.py +++ b/app.py @@ -1165,117 +1165,6 @@ async def vps_speedtest_deploy(req: SpeedtestDeployRequest, _=Depends(_verify)): return result -BW_MANAGER_REPO_FILES = ( - ".env.example", ".gitignore", "README.md", "compose.yaml", - "src/.dockerignore", "src/Dockerfile", "src/app.py", - "src/remote_policy.py", "src/requirements.txt", - "src/templates/base.html", "src/templates/history.html", - "src/templates/index.html", "src/templates/users.html", -) - - -async def _fetch_bw_manager_text(path: str) -> str: - if path not in BW_MANAGER_REPO_FILES: - raise ValueError("unsupported BW Manager file") - cfg = SERVICES.get("forgejo", {}) - base_url, token = cfg.get("url"), _get_key(cfg) - if not base_url or not token: - raise RuntimeError("Forgejo service configuration is unavailable") - url = f"{base_url}/api/v1/repos/sascha/bw-manager/contents/{path}" - async with httpx.AsyncClient(timeout=30) as client: - response = await client.get( - url, params={"ref": "main"}, - headers={"Authorization": f"token {token}"}, - ) - response.raise_for_status() - return base64.b64decode(response.json()["content"]).decode() - - -def _deploy_bw_manager_compose(files: dict[str, str]) -> dict: - if set(files) != set(BW_MANAGER_REPO_FILES): - raise ValueError("BW Manager source bundle is incomplete") - if "build: ./src" not in files["compose.yaml"]: - raise ValueError("BW Manager compose contract is invalid") - if "build_gated_targets" not in files["src/app.py"]: - raise ValueError("BW Manager candidate lacks the user/network AND gate") - rc, working_dir, err = _ssh( - VPS_SSH, - "docker inspect -f '{{ index .Config.Labels \"com.docker.compose.project.working_dir\" }}' bw-manager", - timeout=30, - ) - working_dir = working_dir.strip() - if rc != 0 or not re.fullmatch(r"/app-config/[A-Za-z0-9_./-]+", working_dir): - raise RuntimeError(f"cannot determine safe BW Manager working directory: {(err or working_dir)[-300:]}") - timestamp = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") - candidate = f"/app-config/deployment-candidates/bw-manager-{timestamp}" - backup = f"/app-config/deployment-backups/bw-manager-{timestamp}" - image = "bw-manager-bw-manager" - rollback_image = f"{image}:rollback-{timestamp}" - init_script = f"""from pathlib import Path -import shutil -candidate = Path({candidate!r}) -if candidate.exists(): shutil.rmtree(candidate) -candidate.mkdir(parents=True) -live_env = Path({working_dir!r}) / '.env' -if live_env.exists(): shutil.copy2(live_env, candidate / '.env') -""" - rc, _out, err = _remote_python(init_script) - if rc != 0: - raise RuntimeError(f"BW Manager candidate initialization failed: {err[-300:]}") - # Stage one file per SSH call. Sending the complete repository in one - # command exceeds Linux's argv limit once app.py and templates are encoded. - for relative, content in files.items(): - file_script = f"""from pathlib import Path -target = Path({candidate!r}) / {relative!r} -target.parent.mkdir(parents=True, exist_ok=True) -target.write_text({content!r}) -""" - rc, _out, err = _remote_python(file_script) - if rc != 0: - raise RuntimeError(f"BW Manager staging failed for {relative}: {err[-300:]}") - rc, _out, err = _ssh(VPS_SSH, f"cd {candidate} && docker compose config -q && docker compose build --pull", timeout=600) - if rc != 0: - raise RuntimeError(f"BW Manager candidate build failed: {err[-500:]}") - deploy_script = f"""from pathlib import Path -import shutil -live, backup, candidate = Path({working_dir!r}), Path({backup!r}), Path({candidate!r}) -backup.parent.mkdir(parents=True, exist_ok=True) -if backup.exists(): shutil.rmtree(backup) -shutil.copytree(live, backup) -for relative in {BW_MANAGER_REPO_FILES!r}: - source, target = candidate / relative, live / relative - target.parent.mkdir(parents=True, exist_ok=True) - shutil.copy2(source, target) -""" - rc, _out, err = _remote_python(deploy_script) - if rc != 0: - raise RuntimeError(f"BW Manager live file switch failed: {err[-300:]}") - _ssh(VPS_SSH, f"docker image tag {image} {rollback_image}", timeout=60) - rollback = ( - f"rm -rf {working_dir} && cp -a {backup} {working_dir} && " - f"docker image tag {rollback_image} {image} && cd {working_dir} && " - "docker compose up -d --no-build" - ) - rc, out, err = _ssh(VPS_SSH, f"cd {working_dir} && docker compose up -d --build --remove-orphans", timeout=600) - if rc != 0: - _ssh(VPS_SSH, rollback, timeout=180) - raise RuntimeError(f"BW Manager deployment failed: {(err or out)[-500:]}") - health = "for i in $(seq 1 45); do curl -fsS --max-time 3 http://127.0.0.1:8870/api/status >/dev/null && exit 0; sleep 2; done; exit 1" - rc, _out, err = _ssh(VPS_SSH, health, timeout=105) - if rc != 0: - _ssh(VPS_SSH, rollback, timeout=180) - raise RuntimeError(f"BW Manager health failed; rollback attempted: {err[-300:]}") - return {"status": "deployed", "health": "ok", "working_dir": working_dir, "backup": backup} - - -@app.post("/vps/bw-manager/deploy") -async def vps_bw_manager_deploy(_=Depends(_verify)): - contents = await asyncio.gather(*(_fetch_bw_manager_text(path) for path in BW_MANAGER_REPO_FILES)) - result = await asyncio.to_thread(_deploy_bw_manager_compose, dict(zip(BW_MANAGER_REPO_FILES, contents))) - _audit("/vps/bw-manager/deploy", "POST", 200, "Git-managed BW Manager deployment") - return result - - # --- VM Lifecycle Endpoints --- import subprocess as _sp From fde0b2d469b1e99c34675650177f4732e1c140c9 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 06:51:02 +0200 Subject: [PATCH 03/72] Ansible-NFS-Aktion: app.py aktualisieren --- app.py | 31 ++++++++++++++++++++----------- 1 file changed, 20 insertions(+), 11 deletions(-) diff --git a/app.py b/app.py index 0e2999c..187d584 100644 --- a/app.py +++ b/app.py @@ -1571,25 +1571,34 @@ async def ansible_run(request: Request, _=Depends(_verify)): if not hostname: return JSONResponse({"error": "limit/hostname required"}, status_code=400) action = body.get("action", "setup") - if action not in {"setup", "tune", "pvetune"}: - return JSONResponse({"error": "action must be setup, tune or pvetune"}, status_code=400) + if action not in {"setup", "tune", "pvetune", "nfs"}: + return JSONResponse({"error": "action must be setup, tune, pvetune or nfs"}, status_code=400) if not re.fullmatch(r"[a-zA-Z0-9_.:-]+", hostname): return JSONResponse({"error": "invalid hostname/limit"}, status_code=400) - if action in {"tune", "pvetune"}: - approved_files = ( - "roles/sysctl/defaults/main.yml", - "roles/sysctl/tasks/main.yml", - "group_vars/vps/sysctl.yml", - "sysctl-proxmox.yaml", - "roles/sysctl_proxmox/tasks/main.yml", - ) + if action in {"tune", "pvetune", "nfs"}: + if action == "nfs": + approved_files = ( + "roles/nfs_stability/tasks/main.yml", + "nfs-stability.yml", + "pfannkuchen.sh", + ) + prepare = "mkdir -p roles/nfs_stability/tasks && " + else: + approved_files = ( + "roles/sysctl/defaults/main.yml", + "roles/sysctl/tasks/main.yml", + "group_vars/vps/sysctl.yml", + "sysctl-proxmox.yaml", + "roles/sysctl_proxmox/tasks/main.yml", + ) + prepare = "" file_sync = " && ".join( f"git show origin/master:{path} > {path}" for path in approved_files ) command = ( "cd /app-config/ansible && " "git fetch origin master && " - f"{file_sync} && " + f"{prepare}{file_sync} && " f"bash pfannkuchen.sh {action} {hostname}" ) else: From ad863eac15abc1278c60d3410eaa0344a7b7a229 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 06:51:03 +0200 Subject: [PATCH 04/72] Ansible-NFS-Aktion: tests/test_app.py aktualisieren --- tests/test_app.py | 25 +++++++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/tests/test_app.py b/tests/test_app.py index 3a3be69..34932ce 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -235,6 +235,31 @@ def test_ansible_run_supports_safe_tune_action_and_syncs_approved_files(monkeypa assert len(calls) == 1 +def test_ansible_run_supports_scoped_nfs_action(monkeypatch): + calls = [] + + def fake_ssh(host, command, timeout=600): + calls.append((host, command, timeout)) + return 0, "changed=1 failed=0", "" + + monkeypatch.setattr(app, "_ssh", fake_ssh) + with TestClient(app.app) as client: + response = client.post( + "/ansible/run", + headers={"Authorization": "Bearer test-token"}, + json={"hostname": "arrapps", "action": "nfs"}, + ) + assert response.status_code == 200 + assert response.json()["action"] == "nfs" + command = calls[0][1] + assert "mkdir -p roles/nfs_stability/tasks" in command + assert "git show origin/master:roles/nfs_stability/tasks/main.yml" in command + assert "git show origin/master:nfs-stability.yml" in command + assert "bash pfannkuchen.sh nfs arrapps" in command + assert "git pull --ff-only" not in command + assert len(calls) == 1 + + def test_ansible_run_rejects_unknown_action_and_shell_metacharacters(monkeypatch): monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH"))) with TestClient(app.app) as client: From f7993383d1697faa3688af23d8b9f33b8bc0c3c2 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 08:29:36 +0200 Subject: [PATCH 05/72] Veralteten Semaphore-Proxy entfernen --- butler.yaml | 9 --------- 1 file changed, 9 deletions(-) diff --git a/butler.yaml b/butler.yaml index ebd0c61..d0b0928 100644 --- a/butler.yaml +++ b/butler.yaml @@ -105,15 +105,6 @@ services: description: "Git server (Gitea fork)" health_path: "/api/healthz" - semaphore: - url: "http://10.4.1.116:3010" - auth: bearer - key_file: semaphore - vault_key: semaphore_token - description: "Ansible UI/API" - health_path: "/api/projects" - health_expected: [200] - fileflows: url: "http://10.2.1.104:8268" auth: none From 83433644962a09089eb2799f0dc786ef045b5702 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 09:52:52 +0200 Subject: [PATCH 06/72] Read-only Host-Forensik: app.py --- app.py | 56 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 56 insertions(+) diff --git a/app.py b/app.py index 187d584..00c017e 100644 --- a/app.py +++ b/app.py @@ -1258,6 +1258,62 @@ async def system_sysctl_audit(host: str, _=Depends(_verify)): raise HTTPException(502, "sysctl audit returned invalid JSON") from exc return {"host": host, **result} + +def _host_forensics_command(since_hours: int) -> str: + script = f'''import glob, json, os, subprocess + +def run(command): + proc = subprocess.run(command, shell=True, text=True, capture_output=True, timeout=30) + return {{"rc": proc.returncode, "stdout": proc.stdout.strip()[-12000:], "stderr": proc.stderr.strip()[-1000:]}} + +checks = {{ + "hostname": run("hostnamectl --static 2>/dev/null || hostname"), + "uptime": run("uptime"), + "disk": run("df -hT / /var/lib/docker 2>/dev/null || df -hT /"), + "failed_units": run("systemctl --failed --no-legend --no-pager"), + "docker_binary": run("command -v docker || true"), + "docker_packages": run("dpkg-query -W -f='${{Package}}|${{Status}}|${{Version}}\\n' 'docker*' 'containerd*' 2>/dev/null || true"), + "docker_units": run("systemctl is-active docker containerd 2>/dev/null; systemctl is-enabled docker containerd 2>/dev/null"), + "docker_containers": run("docker ps -a --format '{{{{.Names}}}}|{{{{.Image}}}}|{{{{.Status}}}}' 2>/dev/null || true"), + "docker_images": run("docker image ls --format '{{{{.Repository}}}}:{{{{.Tag}}}}|{{{{.ID}}}}|{{{{.Size}}}}' 2>/dev/null || true"), + "docker_volumes": run("docker volume ls --format '{{{{.Name}}}}' 2>/dev/null || true"), + "docker_disk_usage": run("docker system df 2>/dev/null || true"), + "iptables_docker_refs": run("iptables-save 2>/dev/null | grep -ci docker || true"), + "nft_docker_refs": run("nft list ruleset 2>/dev/null | grep -ci docker || true"), + "forward_policy": run("iptables -S FORWARD 2>/dev/null | head -40"), + "lvm": run("lvs -o lv_name,lv_size,data_percent,metadata_percent --units g --noheadings 2>/dev/null || true"), + "qemu_configs": run("ls -l /etc/pve/nodes/$(hostname)/qemu-server 2>/dev/null || true"), + "recent_system_files": run("find /etc/systemd/system /etc/docker /etc/network -type f -mmin -{since_hours * 60} -printf '%TY-%Tm-%Td %TH:%TM:%TS %p\\n' 2>/dev/null | sort"), + "recent_iso_builder_files": run("find /app-config/ansible/iso-builder -type f -mmin -{since_hours * 60} -printf '%TY-%Tm-%Td %TH:%TM:%TS %p\\n' 2>/dev/null | sort"), + "ansible_git_status": run("git -C /app-config/ansible status --short 2>/dev/null || true"), + "iso_builder_hashes": run("sha256sum /app-config/ansible/iso-builder/* 2>/dev/null || true"), +}} +print(json.dumps(checks)) +''' + encoded = base64.b64encode(script.encode()).decode() + return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' + + +@app.get("/system/forensics/{host}") +async def system_forensics(host: str, since_hours: int = Query(48, ge=1, le=168), _=Depends(_verify)): + """Read-only host residue audit for failed deployments and package/network drift.""" + if not re.fullmatch(r"[a-z0-9][a-z0-9-]{0,62}", host): + raise HTTPException(400, "Invalid host name") + inventory = await asyncio.to_thread(_find_inventory_host, host) + if not inventory: + raise HTTPException(404, f"Host {host} not found") + target = f'{inventory["user"]}@{inventory["ip"]}' + rc, out, err = await asyncio.to_thread(_ssh, target, _host_forensics_command(since_hours), 60) + if rc != 0: + raise HTTPException(502, (err or out).strip()[-500:] or "host forensics failed") + try: + result = json.loads(out) + except json.JSONDecodeError as exc: + raise HTTPException(502, "host forensics returned invalid JSON") from exc + _audit(f"/system/forensics/{host}", "GET", 200, f"since_hours={since_hours}") + return {"host": host, "since_hours": since_hours, "checks": result} + + def _pve_auth(): pv = _parse_kv("proxmox") return f"PVEAPIToken={pv.get('tokenid','')}={pv.get('secret','')}" From a5d3b0c7ceebd0e8a9faac8e1e48e4603994886c Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 09:52:52 +0200 Subject: [PATCH 07/72] Read-only Host-Forensik: tests/test_app.py --- tests/test_app.py | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/tests/test_app.py b/tests/test_app.py index 34932ce..8c32eb7 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -174,6 +174,33 @@ def test_sysctl_audit_rejects_unknown_host_without_ssh(monkeypatch): assert response.status_code == 404 +def test_host_forensics_is_read_only_and_uses_inventory(monkeypatch): + payload = {"docker_binary": {"rc": 0, "stdout": "/usr/bin/docker", "stderr": ""}} + calls = [] + monkeypatch.setattr(app, "_find_inventory_host", lambda name: {"name": name, "user": "root", "ip": "10.5.85.13"}) + monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=600: (calls.append((host, command, timeout)) or (0, __import__("json").dumps(payload), ""))) + with TestClient(app.app) as client: + response = client.get("/system/forensics/node3?since_hours=24", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 200 + assert response.json()["checks"] == payload + assert calls[0][0] == "root@10.5.85.13" + assert calls[0][2] == 60 + assert "base64.b64decode" in calls[0][1] + command = app._host_forensics_command(24) + for destructive in ("systemctl restart", "systemctl stop", "systemctl disable", "docker rm", "docker system prune", "iptables -F", "rm -rf"): + assert destructive not in command + + +def test_host_forensics_rejects_unknown_host_and_invalid_window(monkeypatch): + monkeypatch.setattr(app, "_find_inventory_host", lambda _name: None) + monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not SSH"))) + with TestClient(app.app) as client: + missing = client.get("/system/forensics/not-there", headers={"Authorization": "Bearer test-token"}) + bad_window = client.get("/system/forensics/node3?since_hours=999", headers={"Authorization": "Bearer test-token"}) + assert missing.status_code == 404 + assert bad_window.status_code == 422 + + def test_invalid_log_target_is_rejected_before_ssh(): with TestClient(app.app) as client: response = client.get( From 46f2293078dbf5bba960a7b7046606f02c85d019 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 09:55:08 +0200 Subject: [PATCH 08/72] Forensik um Firewallregeln und redaktierten ISO-Diff erweitern --- app.py | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/app.py b/app.py index 00c017e..128f5c2 100644 --- a/app.py +++ b/app.py @@ -1266,6 +1266,14 @@ def run(command): proc = subprocess.run(command, shell=True, text=True, capture_output=True, timeout=30) return {{"rc": proc.returncode, "stdout": proc.stdout.strip()[-12000:], "stderr": proc.stderr.strip()[-1000:]}} +def redacted_git_diff(): + proc = subprocess.run(["git", "-C", "/app-config/ansible", "diff", "--", "iso-builder/build-iso.sh", "iso-builder/preseed.cfg.tpl", "pfannkuchen.ini"], text=True, capture_output=True, timeout=30) + sensitive = ("password", "passwd", "secret", "token", "private", "credential", "ssh-rsa", "ssh-ed25519") + lines = [] + for line in proc.stdout.splitlines(): + lines.append("[REDACTED SENSITIVE DIFF LINE]" if any(word in line.lower() for word in sensitive) else line) + return {{"rc": proc.returncode, "stdout": "\\n".join(lines)[-12000:], "stderr": proc.stderr.strip()[-1000:]}} + checks = {{ "hostname": run("hostnamectl --static 2>/dev/null || hostname"), "uptime": run("uptime"), @@ -1279,6 +1287,7 @@ checks = {{ "docker_volumes": run("docker volume ls --format '{{{{.Name}}}}' 2>/dev/null || true"), "docker_disk_usage": run("docker system df 2>/dev/null || true"), "iptables_docker_refs": run("iptables-save 2>/dev/null | grep -ci docker || true"), + "iptables_docker_rules": run("iptables-save 2>/dev/null | grep -i docker || true"), "nft_docker_refs": run("nft list ruleset 2>/dev/null | grep -ci docker || true"), "forward_policy": run("iptables -S FORWARD 2>/dev/null | head -40"), "lvm": run("lvs -o lv_name,lv_size,data_percent,metadata_percent --units g --noheadings 2>/dev/null || true"), @@ -1286,6 +1295,8 @@ checks = {{ "recent_system_files": run("find /etc/systemd/system /etc/docker /etc/network -type f -mmin -{since_hours * 60} -printf '%TY-%Tm-%Td %TH:%TM:%TS %p\\n' 2>/dev/null | sort"), "recent_iso_builder_files": run("find /app-config/ansible/iso-builder -type f -mmin -{since_hours * 60} -printf '%TY-%Tm-%Td %TH:%TM:%TS %p\\n' 2>/dev/null | sort"), "ansible_git_status": run("git -C /app-config/ansible status --short 2>/dev/null || true"), + "minecraft_inventory": run("grep -in 'minecraft' /app-config/ansible/pfannkuchen.ini 2>/dev/null || true"), + "iso_builder_diff_redacted": redacted_git_diff(), "iso_builder_hashes": run("sha256sum /app-config/ansible/iso-builder/* 2>/dev/null || true"), }} print(json.dumps(checks)) From 1f27750ac1416e691b77520f8ef01f58b0c4d98c Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:01:27 +0200 Subject: [PATCH 09/72] Sichere Qwen-Bereinigung: app.py --- app.py | 133 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 133 insertions(+) diff --git a/app.py b/app.py index 128f5c2..a931ba0 100644 --- a/app.py +++ b/app.py @@ -1325,6 +1325,139 @@ async def system_forensics(host: str, since_hours: int = Query(48, ge=1, le=168) return {"host": host, "since_hours": since_hours, "checks": result} +def _docker_residue_cleanup_command(dry_run: bool) -> str: + script = f'''import json, os, shlex, shutil, subprocess + +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()}} + +def docker_rule_count(): + proc = run(["iptables-save"]) + return sum(1 for line in proc["stdout"].splitlines() if "docker" in line.lower()) + +result = {{"dry_run": {str(dry_run)}, "before_rule_count": docker_rule_count(), "removed_rules": [], "removed_chains": [], "removed_links": [], "removed_paths": [], "errors": []}} +docker_binary = shutil.which("docker") +unit_state = run(["systemctl", "is-active", "docker", "containerd"])["stdout"].splitlines() +if docker_binary or any(state == "active" for state in unit_state): + result["error"] = "Docker or containerd is still installed/active; refusing residue cleanup" + print(json.dumps(result)); raise SystemExit(2) +if result["dry_run"]: + result["would_remove_paths"] = [path for path in ("/var/lib/docker", "/var/lib/containerd", "/etc/docker") if os.path.exists(path)] + print(json.dumps(result)); raise SystemExit(0) +for table in ("filter", "nat"): + saved = run(["iptables-save", "-t", table]) + rules = [] + for line in saved["stdout"].splitlines(): + if line.startswith("-A ") and "docker" in line.lower(): + rules.append(line) + for line in rules: + args = ["iptables", "-t", table] + shlex.split(line) + args[3] = "-D" + removed = run(args) + if removed["rc"] == 0: + result["removed_rules"].append(table + ":" + line) + else: + result["errors"].append(table + ":" + line + ":" + removed["stderr"]) +for table, chains in (("filter", ("DOCKER-USER", "DOCKER-FORWARD", "DOCKER-BRIDGE", "DOCKER-CT", "DOCKER-INTERNAL", "DOCKER")), ("nat", ("DOCKER",))): + for chain in chains: + run(["iptables", "-t", table, "-F", chain]) + deleted = run(["iptables", "-t", table, "-X", chain]) + if deleted["rc"] == 0: + result["removed_chains"].append(table + ":" + chain) +for link in ("docker0", "docker_gwbridge"): + exists = run(["ip", "link", "show", link]) + if exists["rc"] == 0: + deleted = run(["ip", "link", "delete", link]) + if deleted["rc"] == 0: result["removed_links"].append(link) + else: result["errors"].append(link + ":" + deleted["stderr"]) +for path in ("/var/lib/docker", "/var/lib/containerd", "/etc/docker"): + if os.path.exists(path): + shutil.rmtree(path) + result["removed_paths"].append(path) +result["after_rule_count"] = docker_rule_count() +result["forward_rules"] = run(["iptables", "-S", "FORWARD"])["stdout"].splitlines() +print(json.dumps(result)) +if result["errors"] or result["after_rule_count"] != 0: raise SystemExit(1) +''' + encoded = base64.b64encode(script.encode()).decode() + return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' + + +@app.post("/system/cleanup/docker-residue/{host}") +async def cleanup_docker_residue(host: str, dry_run: bool = Query(True), _=Depends(_verify)): + """Remove only stale Docker firewall/data residue after Docker itself is absent.""" + if not re.fullmatch(r"[a-z0-9][a-z0-9-]{0,62}", host): + raise HTTPException(400, "Invalid host name") + inventory = await asyncio.to_thread(_find_inventory_host, host) + if not inventory: + raise HTTPException(404, f"Host {host} not found") + target = f'{inventory["user"]}@{inventory["ip"]}' + rc, out, err = await asyncio.to_thread(_ssh, target, _docker_residue_cleanup_command(dry_run), 90) + try: + result = json.loads(out) + except json.JSONDecodeError as exc: + raise HTTPException(502, (err or out).strip()[-500:] or "cleanup returned invalid JSON") from exc + if rc != 0: + raise HTTPException(409 if result.get("error") else 502, result) + _audit(f"/system/cleanup/docker-residue/{host}", "POST", 200, f"dry_run={dry_run}", dry_run=dry_run) + return {"host": host, **result} + + +def _iso_builder_restore_command(dry_run: bool) -> str: + script = f'''import glob, hashlib, json, os, subprocess, tempfile +repo = "/app-config/ansible" +paths = ("iso-builder/build-iso.sh", "iso-builder/preseed.cfg.tpl") +result = {{"dry_run": {str(dry_run)}, "restored": [], "removed_outputs": [], "validation": {{}}}} +def run(args): + proc = subprocess.run(args, cwd=repo, text=True, capture_output=True, timeout=60) + return {{"rc": proc.returncode, "stdout": proc.stdout.strip(), "stderr": proc.stderr.strip()}} +fetch = run(["git", "fetch", "origin", "master"]) +if fetch["rc"] != 0: + result["error"] = "git fetch failed"; result["detail"] = fetch["stderr"][-500:]; print(json.dumps(result)); raise SystemExit(1) +outputs = sorted(glob.glob(os.path.join(repo, "iso-builder/output/debian-13-minecraft*.iso"))) +result["would_remove_outputs"] = outputs +for path in paths: + blob = subprocess.run(["git", "show", "origin/master:" + path], cwd=repo, capture_output=True, timeout=30) + if blob.returncode != 0: + result["error"] = "missing canonical file " + path; print(json.dumps(result)); raise SystemExit(1) + current = open(os.path.join(repo, path), "rb").read() if os.path.exists(os.path.join(repo, path)) else b"" + result.setdefault("hashes", {{}})[path] = {{"live_before": hashlib.sha256(current).hexdigest(), "canonical": hashlib.sha256(blob.stdout).hexdigest()}} + if not result["dry_run"]: + destination = os.path.join(repo, path) + fd, temporary = tempfile.mkstemp(dir=os.path.dirname(destination)) + with os.fdopen(fd, "wb") as handle: handle.write(blob.stdout) + os.chmod(temporary, 0o755 if path.endswith(".sh") else 0o644) + os.replace(temporary, destination) + result["restored"].append(path) +if not result["dry_run"]: + for output in outputs: + os.remove(output); result["removed_outputs"].append(output) + syntax = run(["bash", "-n", "iso-builder/build-iso.sh"]) + diff = run(["git", "diff", "--quiet", "origin/master", "--", *paths]) + result["validation"] = {{"bash_syntax_rc": syntax["rc"], "canonical_diff_rc": diff["rc"]}} + if syntax["rc"] != 0 or diff["rc"] != 0: + result["error"] = "post-restore validation failed"; print(json.dumps(result)); raise SystemExit(1) +print(json.dumps(result)) +''' + encoded = base64.b64encode(script.encode()).decode() + return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' + + +@app.post("/system/restore/iso-builder") +async def restore_iso_builder(dry_run: bool = Query(True), _=Depends(_verify)): + """Restore only the canonical ISO-builder files and remove generated Minecraft ISOs.""" + rc, out, err = await asyncio.to_thread(_ssh, AUTOMATION1, _iso_builder_restore_command(dry_run), 120) + try: + result = json.loads(out) + except json.JSONDecodeError as exc: + raise HTTPException(502, (err or out).strip()[-500:] or "restore returned invalid JSON") from exc + if rc != 0: + raise HTTPException(502, result) + _audit("/system/restore/iso-builder", "POST", 200, f"dry_run={dry_run}", dry_run=dry_run) + return result + + def _pve_auth(): pv = _parse_kv("proxmox") return f"PVEAPIToken={pv.get('tokenid','')}={pv.get('secret','')}" From c62f6709d7117455d5686d33e876c71e48e42aff Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:01:27 +0200 Subject: [PATCH 10/72] Sichere Qwen-Bereinigung: tests/test_app.py --- tests/test_app.py | 39 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 39 insertions(+) diff --git a/tests/test_app.py b/tests/test_app.py index 8c32eb7..b5db980 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -201,6 +201,45 @@ def test_host_forensics_rejects_unknown_host_and_invalid_window(monkeypatch): assert bad_window.status_code == 422 +def test_docker_residue_cleanup_defaults_to_dry_run(monkeypatch): + payload = {"dry_run": True, "before_rule_count": 25, "removed_rules": [], "removed_chains": [], "removed_links": [], "removed_paths": [], "errors": []} + calls = [] + monkeypatch.setattr(app, "_find_inventory_host", lambda name: {"name": name, "user": "root", "ip": "10.5.85.13"}) + monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=600: (calls.append((host, command, timeout)) or (0, __import__("json").dumps(payload), ""))) + with TestClient(app.app) as client: + response = client.post("/system/cleanup/docker-residue/node3", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 200 + assert response.json()["dry_run"] is True + assert calls[0][0] == "root@10.5.85.13" + assert calls[0][2] == 90 + + +def test_docker_residue_cleanup_refuses_active_docker(monkeypatch): + payload = {"dry_run": False, "error": "Docker or containerd is still installed/active; refusing residue cleanup"} + monkeypatch.setattr(app, "_find_inventory_host", lambda name: {"name": name, "user": "root", "ip": "10.5.85.13"}) + monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (2, __import__("json").dumps(payload), "")) + with TestClient(app.app) as client: + response = client.post("/system/cleanup/docker-residue/node3?dry_run=false", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 409 + + +def test_iso_builder_restore_defaults_to_dry_run_and_has_no_free_target(monkeypatch): + payload = {"dry_run": True, "restored": [], "removed_outputs": [], "validation": {}} + calls = [] + monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=600: (calls.append((host, command, timeout)) or (0, __import__("json").dumps(payload), ""))) + with TestClient(app.app) as client: + response = client.post("/system/restore/iso-builder", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 200 + assert response.json()["dry_run"] is True + assert calls[0][0] == app.AUTOMATION1 + assert calls[0][2] == 120 + command = app._iso_builder_restore_command(True) + encoded = command.split("base64.b64decode('", 1)[1].split("')", 1)[0] + decoded = __import__("base64").b64decode(encoded).decode() + assert "origin/master" in decoded + assert "/app-config/ansible" in decoded + + def test_invalid_log_target_is_rejected_before_ssh(): with TestClient(app.app) as client: response = client.get( From a324b136dfaeff4a5058c067be0d515f7c79bb5f Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:30:30 +0200 Subject: [PATCH 11/72] Sichere Caddy- und DNS-Stilllegung: app.py --- app.py | 153 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 153 insertions(+) diff --git a/app.py b/app.py index a931ba0..7f6ba54 100644 --- a/app.py +++ b/app.py @@ -1458,6 +1458,159 @@ async def restore_iso_builder(dry_run: bool = Query(True), _=Depends(_verify)): return result +CADDY_HOST = "root@46.225.230.72" +CADDYFILE_PATH = "/app-config/caddy/Caddyfile" +HETZNER_DNS_API = "https://api.hetzner.cloud/v1" + + +def _validate_managed_hostname(hostname: str) -> str: + hostname = hostname.strip().lower().rstrip(".") + if not re.fullmatch(r"[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?(?:\.[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?)+", hostname): + raise HTTPException(400, "Invalid hostname") + if not hostname.endswith(".sascha-lutz.de"): + raise HTTPException(400, "Hostname is outside the managed zone") + return hostname + + +def _caddy_site_remove_command(hostname: str, dry_run: bool) -> str: + script = '''import datetime, hashlib, json, os, re, subprocess, sys +path = %r +hostname = %r +dry_run = %r + +def digest(data): + return hashlib.sha256(data.encode()).hexdigest() + +def run(args): + p = subprocess.run(args, text=True, capture_output=True, timeout=30) + return {"rc": p.returncode, "stdout": p.stdout.strip()[-2000:], "stderr": p.stderr.strip()[-2000:]} + +text = open(path, encoding="utf-8").read() +pattern = re.compile(r"(?m)^[ \\t]*" + re.escape(hostname) + r"[ \\t]*\\{") +match = pattern.search(text) +result = {"hostname": hostname, "dry_run": dry_run, "found": bool(match), "before_sha256": digest(text)} +if not match: + print(json.dumps(result)); raise SystemExit(0) +start = match.start(); depth = 0; end = None +for idx in range(match.end() - 1, len(text)): + if text[idx] == "{": depth += 1 + elif text[idx] == "}": + depth -= 1 + if depth == 0: + end = idx + 1 + while end < len(text) and text[end] in " \\t": end += 1 + while end < len(text) and text[end] == "\\n": end += 1 + break +if end is None: + result["error"] = "Unbalanced Caddy site block"; print(json.dumps(result)); raise SystemExit(2) +result["line_start"] = text.count("\\n", 0, start) + 1 +result["line_end"] = text.count("\\n", 0, end) + 1 +if dry_run: + print(json.dumps(result)); raise SystemExit(0) +new = text[:start] + text[end:] +backup = path + ".pre-outline-removal-" + datetime.datetime.now().strftime("%%Y%%m%%d-%%H%%M%%S") +open(backup, "w", encoding="utf-8").write(text) +with open(path, "w", encoding="utf-8") as f: + f.write(new); f.flush(); os.fsync(f.fileno()) +result["backup"] = backup +result["after_sha256"] = digest(new) +validation = run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/etc/caddy/Caddyfile"]) +result["validation"] = validation +if validation["rc"] != 0: + with open(path, "w", encoding="utf-8") as f: + f.write(text); f.flush(); os.fsync(f.fileno()) + result["rolled_back"] = True; print(json.dumps(result)); raise SystemExit(3) +host_sha = run(["sha256sum", path]) +container_sha = run(["docker", "exec", "caddy", "sha256sum", "/etc/caddy/Caddyfile"]) +result["host_container_hash_match"] = bool(host_sha["stdout"] and container_sha["stdout"] and host_sha["stdout"].split()[0] == container_sha["stdout"].split()[0]) +if not result["host_container_hash_match"]: + with open(path, "w", encoding="utf-8") as f: + f.write(text); f.flush(); os.fsync(f.fileno()) + result["rolled_back"] = True; print(json.dumps(result)); raise SystemExit(4) +reload = run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"]) +result["reload"] = reload +if reload["rc"] != 0: + with open(path, "w", encoding="utf-8") as f: + f.write(text); f.flush(); os.fsync(f.fileno()) + run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"]) + result["rolled_back"] = True; print(json.dumps(result)); raise SystemExit(5) +container_text = run(["docker", "exec", "caddy", "sh", "-c", "cat /etc/caddy/Caddyfile"]) +result["hostname_absent"] = hostname not in container_text["stdout"] +print(json.dumps(result)) +''' % (CADDYFILE_PATH, hostname, dry_run) + encoded = base64.b64encode(script.encode()).decode() + return f"python3 -c \"import base64;exec(base64.b64decode('{encoded}'))\"" + + +@app.delete("/caddy/site/{hostname}") +async def caddy_site_remove(hostname: str, _=Depends(_verify), dry_run: bool = Query(True)): + hostname = _validate_managed_hostname(hostname) + rc, out, err = _ssh(CADDY_HOST, _caddy_site_remove_command(hostname, dry_run), timeout=90) + try: + result = json.loads(out) + except Exception: + raise HTTPException(502, (err or out or "Caddy removal returned no JSON")[-1000:]) + if rc != 0: + raise HTTPException(502, result) + _audit(f"/caddy/site/{hostname}", "DELETE", 200, f"dry_run={dry_run}") + return result + + +async def _hetzner_zone_and_rrsets(zone_name: str): + token = _read("HETZNER_DNS_TOKEN") + if not token: + raise HTTPException(503, "HETZNER_DNS_TOKEN unavailable in Butler vault cache") + headers = {"Authorization": f"Bearer {token}"} + async with httpx.AsyncClient(timeout=30) as client: + zones_response = await client.get(f"{HETZNER_DNS_API}/zones", headers=headers, params={"name": zone_name}) + if zones_response.status_code != 200: + raise HTTPException(502, f"Hetzner zones lookup failed: HTTP {zones_response.status_code}") + zones = zones_response.json().get("zones", []) + zone = next((item for item in zones if item.get("name") == zone_name), None) + if not zone: + raise HTTPException(404, "DNS zone not found") + rr_response = await client.get(f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets", headers=headers) + if rr_response.status_code != 200: + raise HTTPException(502, f"Hetzner RRSet lookup failed: HTTP {rr_response.status_code}") + return zone, rr_response.json().get("rrsets", []), headers + + +@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(".") + hostname = _validate_managed_hostname(f"{record_name}.{zone_name}") + record_name = hostname[: -(len(zone_name) + 1)] + zone, rrsets, _headers = await _hetzner_zone_and_rrsets(zone_name) + selected = [r for r in rrsets if r.get("name") == record_name and r.get("type") in {"A", "AAAA", "CNAME"}] + return {"zone": zone_name, "zone_id": zone.get("id"), "name": record_name, "rrsets": selected} + + +@app.delete("/dns/rrset/{zone_name}/{record_name}") +async def dns_rrset_delete(zone_name: str, record_name: str, _=Depends(_verify), dry_run: bool = Query(True)): + zone_name = zone_name.strip().lower().rstrip(".") + hostname = _validate_managed_hostname(f"{record_name}.{zone_name}") + record_name = hostname[: -(len(zone_name) + 1)] + zone, rrsets, headers = await _hetzner_zone_and_rrsets(zone_name) + selected = [r for r in rrsets if r.get("name") == record_name and r.get("type") in {"A", "AAAA", "CNAME"}] + result = {"zone": zone_name, "zone_id": zone.get("id"), "name": record_name, "dry_run": dry_run, "rrsets": selected, "deleted": []} + if dry_run or not selected: + return result + async with httpx.AsyncClient(timeout=30) as client: + for rrset in selected: + url = f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets/{record_name}/{rrset['type']}" + response = await client.delete(url, headers=headers) + if response.status_code not in {200, 204}: + raise HTTPException(502, f"Hetzner RRSet delete failed for {rrset['type']}: HTTP {response.status_code}") + result["deleted"].append(rrset["type"]) + verify_response = await client.get(f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets", headers=headers) + remaining = verify_response.json().get("rrsets", []) if verify_response.status_code == 200 else selected + result["remaining"] = [r for r in remaining if r.get("name") == record_name and r.get("type") in {"A", "AAAA", "CNAME"}] + if result["remaining"]: + raise HTTPException(502, "RRSet read-back still contains deleted record") + _audit(f"/dns/rrset/{zone_name}/{record_name}", "DELETE", 200, "deleted=" + ",".join(result["deleted"])) + return result + + def _pve_auth(): pv = _parse_kv("proxmox") return f"PVEAPIToken={pv.get('tokenid','')}={pv.get('secret','')}" From 14c757bbd2156ddbf5efb4c6884acb7d6d154479 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:30:31 +0200 Subject: [PATCH 12/72] Sichere Caddy- und DNS-Stilllegung: tests/test_app.py --- tests/test_app.py | 37 +++++++++++++++++++++++++++++++++++++ 1 file changed, 37 insertions(+) diff --git a/tests/test_app.py b/tests/test_app.py index b5db980..a3ea996 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -240,6 +240,43 @@ def test_iso_builder_restore_defaults_to_dry_run_and_has_no_free_target(monkeypa assert "/app-config/ansible" in decoded +def test_caddy_site_remove_defaults_to_dry_run(monkeypatch): + payload = {"hostname": "wiki.sascha-lutz.de", "dry_run": True, "found": True, "line_start": 10, "line_end": 13} + calls = [] + monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=600: (calls.append((host, command, timeout)) or (0, __import__("json").dumps(payload), ""))) + with TestClient(app.app) as client: + response = client.delete("/caddy/site/wiki.sascha-lutz.de", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 200 + assert response.json()["dry_run"] is True + assert calls[0][0] == app.CADDY_HOST + assert calls[0][2] == 90 + + +def test_caddy_site_remove_rejects_unmanaged_hostname_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.delete("/caddy/site/example.com", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 400 + + +def test_dns_rrset_dry_run_returns_only_selected_records(monkeypatch): + async def fake_lookup(zone): + assert zone == "sascha-lutz.de" + return {"id": 96805, "name": zone}, [ + {"name": "wiki", "type": "A", "records": [{"value": "46.225.230.72"}]}, + {"name": "wiki", "type": "AAAA", "records": [{"value": "::1"}]}, + {"name": "git", "type": "A", "records": [{"value": "46.225.230.72"}]}, + ], {"Authorization": "Bearer hidden"} + monkeypatch.setattr(app, "_hetzner_zone_and_rrsets", fake_lookup) + with TestClient(app.app) as client: + response = client.delete("/dns/rrset/sascha-lutz.de/wiki", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 200 + body = response.json() + assert body["dry_run"] is True + assert [item["type"] for item in body["rrsets"]] == ["A", "AAAA"] + assert body["deleted"] == [] + + def test_invalid_log_target_is_rejected_before_ssh(): with TestClient(app.app) as client: response = client.get( From 3ba6001c83bd41d9de4d8757f0d7af7133a31cdb Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:32:46 +0200 Subject: [PATCH 13/72] DNS-Token bei Bedarf sicher aus Vault-Cache nachladen --- app.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/app.py b/app.py index 7f6ba54..c5ffbf8 100644 --- a/app.py +++ b/app.py @@ -1557,9 +1557,10 @@ async def caddy_site_remove(hostname: str, _=Depends(_verify), dry_run: bool = Q async def _hetzner_zone_and_rrsets(zone_name: str): - token = _read("HETZNER_DNS_TOKEN") - if not token: - raise HTTPException(503, "HETZNER_DNS_TOKEN unavailable in Butler vault cache") + try: + token = await asyncio.to_thread(_get_hetzner_dns_token) + except RuntimeError as exc: + raise HTTPException(503, str(exc)) headers = {"Authorization": f"Bearer {token}"} async with httpx.AsyncClient(timeout=30) as client: zones_response = await client.get(f"{HETZNER_DNS_API}/zones", headers=headers, params={"name": zone_name}) From 5d68ff90a8dd6de3cece6e9e84146f3e20d98e03 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:34:29 +0200 Subject: [PATCH 14/72] =?UTF-8?q?Hetzner-DNS-Vaultalias=20sicher=20aufl?= =?UTF-8?q?=C3=B6sen:=20app.py?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/app.py b/app.py index c5ffbf8..ee5ffdf 100644 --- a/app.py +++ b/app.py @@ -984,6 +984,9 @@ def _get_hetzner_dns_token() -> str: token = _read("HETZNER_DNS_TOKEN") if token: return token + aliases = [value for key, value in _vault_cache.items() if "hetzner" in key.lower() and "dns" in key.lower() and value.strip()] + if len(aliases) == 1: + return aliases[0].strip() rc, _out, err = _ssh( "sascha@10.4.1.116", "sudo bash /data/stacks/homelab-butler/vault-sync.sh", From 90021b5a4822df3c61d6cb400aebb4d4167e4ff9 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:34:29 +0200 Subject: [PATCH 15/72] =?UTF-8?q?Hetzner-DNS-Vaultalias=20sicher=20aufl?= =?UTF-8?q?=C3=B6sen:=20tests/test=5Fapp.py?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/test_app.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/tests/test_app.py b/tests/test_app.py index a3ea996..cf541dd 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -277,6 +277,13 @@ def test_dns_rrset_dry_run_returns_only_selected_records(monkeypatch): assert body["deleted"] == [] +def test_hetzner_dns_token_accepts_single_sanitized_vault_alias(monkeypatch): + monkeypatch.setattr(app, "_vault_cache", {"hetzner-dns-api": "secret-value", "other": "ignored"}) + monkeypatch.setattr(app, "_read", lambda _name: None) + monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not sync vault"))) + assert app._get_hetzner_dns_token() == "secret-value" + + def test_invalid_log_target_is_rejected_before_ssh(): with TestClient(app.app) as client: response = client.get( From 87e2f6d902c1d92035441b98f72805868c91694d Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:36:26 +0200 Subject: [PATCH 16/72] Hetzner-DNS-Token aus beschrifteter Secure Note normalisieren --- app.py | 27 +++++++++++++++++++++++---- 1 file changed, 23 insertions(+), 4 deletions(-) diff --git a/app.py b/app.py index ee5ffdf..db5ebbe 100644 --- a/app.py +++ b/app.py @@ -980,13 +980,32 @@ print(backup) return {"status": "reloaded", "backup": backup} +def _normalize_hetzner_dns_token(raw: str) -> str | None: + raw = (raw or "").strip() + if not raw: + return None + if raw.isascii() and not any(ch.isspace() for ch in raw) and re.fullmatch(r"[A-Za-z0-9._-]{24,}", raw): + return raw + candidates = [] + labelled = re.findall(r"(?is)(?:token|api[ -]?key)[^\n:=]{0,80}(?::|=|\n)\s*([A-Za-z0-9._-]{24,})", raw) + candidates.extend(value for value in labelled if value.isascii()) + for line in raw.splitlines(): + if "token" not in line.lower() and "api key" not in line.lower() and "api-key" not in line.lower(): + continue + value = re.split(r"[:=]", line, maxsplit=1)[-1].strip().strip("`'\"") + if value.isascii() and re.fullmatch(r"[A-Za-z0-9._-]{24,}", value): + candidates.append(value) + return candidates[0] if len(set(candidates)) == 1 else None + + def _get_hetzner_dns_token() -> str: - token = _read("HETZNER_DNS_TOKEN") + token = _normalize_hetzner_dns_token(_read("HETZNER_DNS_TOKEN") or "") if token: return token - aliases = [value for key, value in _vault_cache.items() if "hetzner" in key.lower() and "dns" in key.lower() and value.strip()] - if len(aliases) == 1: - return aliases[0].strip() + aliases = [_normalize_hetzner_dns_token(value) for key, value in _vault_cache.items() if "hetzner" in key.lower() and "dns" in key.lower()] + aliases = [value for value in aliases if value] + if len(set(aliases)) == 1: + return aliases[0] rc, _out, err = _ssh( "sascha@10.4.1.116", "sudo bash /data/stacks/homelab-butler/vault-sync.sh", From 98c9afaef22730ec9c6627488a5f4b977a2581cf Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:37:50 +0200 Subject: [PATCH 17/72] Eindeutigen langen DNS-Token aus Secure Note erkennen --- app.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/app.py b/app.py index db5ebbe..19c164e 100644 --- a/app.py +++ b/app.py @@ -989,6 +989,8 @@ def _normalize_hetzner_dns_token(raw: str) -> str | None: candidates = [] labelled = re.findall(r"(?is)(?:token|api[ -]?key)[^\n:=]{0,80}(?::|=|\n)\s*([A-Za-z0-9._-]{24,})", raw) candidates.extend(value for value in labelled if value.isascii()) + broad = re.findall(r"(? Date: Fri, 14 Aug 2026 10:40:17 +0200 Subject: [PATCH 18/72] Stillgelegten Outline-Proxy entfernen --- butler.yaml | 7 ------- 1 file changed, 7 deletions(-) diff --git a/butler.yaml b/butler.yaml index d0b0928..4624d82 100644 --- a/butler.yaml +++ b/butler.yaml @@ -45,13 +45,6 @@ services: description: "Media request management" health_path: "/api/v1/status" - outline: - url: "http://10.1.1.100:3000" - auth: bearer - key_file: outline - vault_key: outline_api_key - description: "Wiki" - n8n: url: "http://10.4.1.113:5678" auth: n8n From 2f2d498e7fb0a90cc5e1ced7b5cabdc4682b74d9 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:43:15 +0200 Subject: [PATCH 19/72] Outline-Monitore in Uptime Kuma read-only auditieren --- app.py | 31 ++++++++++++++++++++++++------- 1 file changed, 24 insertions(+), 7 deletions(-) diff --git a/app.py b/app.py index 19c164e..84a123a 100644 --- a/app.py +++ b/app.py @@ -1284,19 +1284,35 @@ async def system_sysctl_audit(host: str, _=Depends(_verify)): def _host_forensics_command(since_hours: int) -> str: - script = f'''import glob, json, os, subprocess + script = f'''import glob, json, os, re, subprocess def run(command): proc = subprocess.run(command, shell=True, text=True, capture_output=True, timeout=30) return {{"rc": proc.returncode, "stdout": proc.stdout.strip()[-12000:], "stderr": proc.stderr.strip()[-1000:]}} -def redacted_git_diff(): - proc = subprocess.run(["git", "-C", "/app-config/ansible", "diff", "--", "iso-builder/build-iso.sh", "iso-builder/preseed.cfg.tpl", "pfannkuchen.ini"], text=True, capture_output=True, timeout=30) - sensitive = ("password", "passwd", "secret", "token", "private", "credential", "ssh-rsa", "ssh-ed25519") +def safe_git_diff(paths): + proc = subprocess.run(["git", "-C", "/app-config/ansible", "diff", "--"] + paths, text=True, capture_output=True, timeout=30) + sensitive = re.compile(r"pass|secret|token|api[_-]?key|private[_-]?key", re.I) lines = [] for line in proc.stdout.splitlines(): - lines.append("[REDACTED SENSITIVE DIFF LINE]" if any(word in line.lower() for word in sensitive) else line) - return {{"rc": proc.returncode, "stdout": "\\n".join(lines)[-12000:], "stderr": proc.stderr.strip()[-1000:]}} + lines.append("[REDACTED SENSITIVE DIFF LINE]" if sensitive.search(line) else line) + return {{"rc": proc.returncode, "stdout": "\\n".join(lines)[-12000:], "stderr": proc.stderr.strip()[-4000:]}} + +def kuma_outline_monitors(): + path = "/app-config/kuma/kuma.db" + if not os.path.exists(path): + return {{"rc": 0, "stdout": "[]", "stderr": ""}} + try: + import sqlite3 + connection = sqlite3.connect("file:" + path + "?mode=ro", uri=True) + columns = [row[1] for row in connection.execute("pragma table_info(monitor)")] + wanted = [name for name in ("id", "name", "url", "hostname", "active") if name in columns] + rows = [dict(zip(wanted, row)) for row in connection.execute("select " + ",".join(wanted) + " from monitor")] + selected = [row for row in rows if "outline" in json.dumps(row).lower() or "wiki.sascha-lutz.de" in json.dumps(row).lower()] + connection.close() + return {{"rc": 0, "stdout": json.dumps(selected), "stderr": ""}} + except Exception as exc: + return {{"rc": 1, "stdout": "", "stderr": str(exc)}} checks = {{ "hostname": run("hostnamectl --static 2>/dev/null || hostname"), @@ -1320,8 +1336,9 @@ checks = {{ "recent_iso_builder_files": run("find /app-config/ansible/iso-builder -type f -mmin -{since_hours * 60} -printf '%TY-%Tm-%Td %TH:%TM:%TS %p\\n' 2>/dev/null | sort"), "ansible_git_status": run("git -C /app-config/ansible status --short 2>/dev/null || true"), "minecraft_inventory": run("grep -in 'minecraft' /app-config/ansible/pfannkuchen.ini 2>/dev/null || true"), - "iso_builder_diff_redacted": redacted_git_diff(), + "iso_builder_diff_redacted": safe_git_diff(["iso-builder/build-iso.sh", "iso-builder/preseed.cfg.tpl", "pfannkuchen.ini"]), "iso_builder_hashes": run("sha256sum /app-config/ansible/iso-builder/* 2>/dev/null || true"), + "kuma_outline_monitors": kuma_outline_monitors(), }} print(json.dumps(checks)) ''' From a37784c8969e36359d10ec914f9a79defbea2256 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:49:01 +0200 Subject: [PATCH 20/72] =?UTF-8?q?Sichere=20Uptime-Monitor-Entfernung=20erg?= =?UTF-8?q?=C3=A4nzen?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app.py | 78 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 78 insertions(+) diff --git a/app.py b/app.py index 84a123a..526598c 100644 --- a/app.py +++ b/app.py @@ -1499,6 +1499,84 @@ async def restore_iso_builder(dry_run: bool = Query(True), _=Depends(_verify)): return result +UPTIME_HOST = "sascha@10.5.85.5" +UPTIME_DB = "/app-config/kuma/kuma.db" + + +def _uptime_monitor_remove_command(monitor_id: int, expected_name: str, dry_run: bool) -> str: + script = '''import datetime, json, os, shutil, sqlite3, subprocess +monitor_id = %r +expected_name = %r +dry_run = %r +path = %r + +def docker(*args): + return subprocess.run(["sudo", "docker", *args], text=True, capture_output=True, timeout=60) + +connection = sqlite3.connect("file:" + path + "?mode=ro", uri=True) +connection.row_factory = sqlite3.Row +row = connection.execute("select id,name,url,active from monitor where id=?", (monitor_id,)).fetchone() +connection.close() +result = {"monitor_id": monitor_id, "expected_name": expected_name, "dry_run": dry_run, "found": dict(row) if row else None} +if not row: + print(json.dumps(result)); raise SystemExit(0) +if row["name"] != expected_name: + result["error"] = "Monitor name mismatch"; print(json.dumps(result)); raise SystemExit(2) +if dry_run: + print(json.dumps(result)); raise SystemExit(0) +stop = docker("stop", "kuma") +if stop.returncode != 0: + result["error"] = "Could not stop Kuma"; result["stderr"] = stop.stderr[-500:]; print(json.dumps(result)); raise SystemExit(3) +backup = path + ".pre-monitor-removal-" + datetime.datetime.now().strftime("%%Y%%m%%d-%%H%%M%%S") +try: + shutil.copy2(path, backup) + connection = sqlite3.connect(path) + connection.execute("pragma foreign_keys=off") + tables = [item[0] for item in connection.execute("select name from sqlite_master where type='table'")] + cleaned = [] + for table in tables: + if table == "monitor" or not table.replace("_", "").isalnum(): + continue + for fk in connection.execute('pragma foreign_key_list("' + table + '")'): + if fk[2] == "monitor" and fk[3].replace("_", "").isalnum(): + cursor = connection.execute('delete from "' + table + '" where "' + fk[3] + '"=?', (monitor_id,)) + if cursor.rowcount: + cleaned.append({"table": table, "rows": cursor.rowcount}) + deleted = connection.execute("delete from monitor where id=? and name=?", (monitor_id, expected_name)).rowcount + connection.commit(); connection.close() + result["backup"] = backup; result["dependencies_cleaned"] = cleaned; result["deleted"] = deleted +except Exception as exc: + shutil.copy2(backup, path) + result["error"] = str(exc) +finally: + start = docker("start", "kuma") + result["container_start_rc"] = start.returncode +if result.get("error") or result.get("deleted") != 1 or result["container_start_rc"] != 0: + print(json.dumps(result)); raise SystemExit(4) +connection = sqlite3.connect("file:" + path + "?mode=ro", uri=True) +result["remaining"] = connection.execute("select count(*) from monitor where id=?", (monitor_id,)).fetchone()[0] +connection.close() +print(json.dumps(result)) +''' % (monitor_id, expected_name, dry_run, UPTIME_DB) + encoded = base64.b64encode(script.encode()).decode() + return f"python3 -c \"import base64;exec(base64.b64decode('{encoded}'))\"" + + +@app.delete("/uptime/monitor/{monitor_id}") +async def uptime_monitor_remove(monitor_id: int, expected_name: str = Query(..., min_length=1, max_length=100), dry_run: bool = Query(True), _=Depends(_verify)): + if not re.fullmatch(r"[A-Za-z0-9 ._()-]+", expected_name): + raise HTTPException(400, "Invalid expected monitor name") + rc, out, err = _ssh(UPTIME_HOST, _uptime_monitor_remove_command(monitor_id, expected_name, dry_run), timeout=120) + try: + result = json.loads(out) + except Exception: + raise HTTPException(502, (err or out or "Uptime cleanup returned no JSON")[-1000:]) + if rc != 0: + raise HTTPException(409 if result.get("error") == "Monitor name mismatch" else 502, result) + _audit(f"/uptime/monitor/{monitor_id}", "DELETE", 200, f"dry_run={dry_run}") + return result + + CADDY_HOST = "root@46.225.230.72" CADDYFILE_PATH = "/app-config/caddy/Caddyfile" HETZNER_DNS_API = "https://api.hetzner.cloud/v1" From c40e869ef8731142d4152125962e46ffb86e82c7 Mon Sep 17 00:00:00 2001 From: sascha Date: Fri, 14 Aug 2026 10:49:02 +0200 Subject: [PATCH 21/72] Uptime-Monitor-Entfernung testen --- tests/test_app.py | 22 ++++++++++++++++++++-- 1 file changed, 20 insertions(+), 2 deletions(-) diff --git a/tests/test_app.py b/tests/test_app.py index cf541dd..0b30b71 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -1,5 +1,6 @@ import os import asyncio +import json import time from datetime import datetime, timedelta, timezone @@ -278,10 +279,27 @@ def test_dns_rrset_dry_run_returns_only_selected_records(monkeypatch): def test_hetzner_dns_token_accepts_single_sanitized_vault_alias(monkeypatch): - monkeypatch.setattr(app, "_vault_cache", {"hetzner-dns-api": "secret-value", "other": "ignored"}) + token = "a" * 48 + monkeypatch.setattr(app, "_vault_cache", {"hetzner-dns-api": f"Hetzner DNS API Token: {token}\nFür sascha-lutz.de", "other": "ignored"}) monkeypatch.setattr(app, "_read", lambda _name: None) monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not sync vault"))) - assert app._get_hetzner_dns_token() == "secret-value" + assert app._get_hetzner_dns_token() == token + + +def test_uptime_monitor_remove_defaults_to_dry_run(monkeypatch): + payload = {"monitor_id": 71, "expected_name": "Outline Wiki", "dry_run": True, "found": {"id": 71, "name": "Outline Wiki"}} + monkeypatch.setattr(app, "_ssh", lambda target, command, timeout=120: (0, json.dumps(payload), "")) + with TestClient(app.app) as client: + response = client.delete("/uptime/monitor/71", params={"expected_name": "Outline Wiki"}, headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 200 + assert response.json()["dry_run"] is True + + +def test_uptime_monitor_remove_rejects_invalid_expected_name_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.delete("/uptime/monitor/71", params={"expected_name": "Outline; rm -rf /"}, headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 400 def test_invalid_log_target_is_rejected_before_ssh(): From 67d59536e910e828748d96e93bdc784c2237cef9 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:06:03 +0200 Subject: [PATCH 22/72] Add redacted WireGuard status API (#37) --- app.py | 61 +++++++++++++++++++++++++++++++++++++++++++++++ tests/test_app.py | 46 +++++++++++++++++++++++++++++++++++ 2 files changed, 107 insertions(+) diff --git a/app.py b/app.py index 526598c..c43dec1 100644 --- a/app.py +++ b/app.py @@ -1213,6 +1213,67 @@ def _ssh(host, cmd, timeout=600): return 124, "", f"SSH command timed out after {timeout} seconds" +def _wireguard_status_command() -> str: + script = '''import json, subprocess + +def run(args): + proc = subprocess.run(args, text=True, capture_output=True, timeout=15) + if proc.returncode != 0: + raise RuntimeError((proc.stderr or proc.stdout).strip() or "command failed") + return proc.stdout.strip() + +result = {"interface": "wg0", "addresses": [], "listen_port": None, "service_active": False, "service_enabled": False, "routes": [], "peers": []} +result["service_active"] = subprocess.run(["systemctl", "is-active", "--quiet", "wg-quick@wg0"]).returncode == 0 +result["service_enabled"] = subprocess.run(["systemctl", "is-enabled", "--quiet", "wg-quick@wg0"]).returncode == 0 +try: + addr_data = json.loads(run(["ip", "-j", "address", "show", "dev", "wg0"])) + for item in addr_data: + for address in item.get("addr_info", []): + result["addresses"].append(address["local"] + "/" + str(address["prefixlen"])) + result["routes"] = json.loads(run(["ip", "-j", "route", "show", "dev", "wg0"])) + rows = run(["wg", "show", "wg0", "dump"]).splitlines() + if rows: + interface = rows[0].split("\\t") + result["listen_port"] = int(interface[2]) + for raw in rows[1:]: + fields = raw.split("\\t") + result["peers"].append({ + "public_key": fields[0], + "endpoint": None if fields[2] == "(none)" else fields[2], + "allowed_ips": [] if fields[3] == "(none)" else fields[3].split(","), + "latest_handshake": int(fields[4]), + "rx_bytes": int(fields[5]), + "tx_bytes": int(fields[6]), + "persistent_keepalive": int(fields[7]), + }) +except Exception as exc: + result["error"] = str(exc)[:300] +print(json.dumps(result)) +''' + encoded = base64.b64encode(script.encode()).decode() + return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' + + +@app.get("/network/wireguard/{host}") +async def network_wireguard_status(host: str, _=Depends(_verify)): + """Return redacted WireGuard state without private or preshared keys.""" + if not re.fullmatch(r"[a-z0-9][a-z0-9-]{0,62}", host): + raise HTTPException(400, "Invalid host name") + inventory = await asyncio.to_thread(_find_inventory_host, host) + if not inventory: + raise HTTPException(404, f"Host {host} not found") + target = f'{inventory["user"]}@{inventory["ip"]}' + rc, out, err = await asyncio.to_thread(_ssh, target, _wireguard_status_command(), 30) + if rc != 0: + raise HTTPException(502, (err or out).strip()[-500:] or "WireGuard status failed") + try: + result = json.loads(out) + except json.JSONDecodeError as exc: + raise HTTPException(502, "WireGuard status returned invalid JSON") from exc + _audit(f"/network/wireguard/{host}", "GET", 200, "redacted live status") + return {"host": host, **result} + + SYSCTL_AUDIT_KEYS = ( "net.core.default_qdisc", "net.core.rmem_default", diff --git a/tests/test_app.py b/tests/test_app.py index 0b30b71..eb3be0d 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -46,6 +46,52 @@ def test_health_exposes_current_version(): assert response.json()["version"] == app.VERSION == "2.3.5" +def test_wireguard_status_returns_redacted_live_state(monkeypatch): + payload = { + "interface": "wg0", + "addresses": ["10.11.12.1/32"], + "listen_port": 37888, + "service_active": True, + "service_enabled": True, + "routes": [{"dst": "10.11.12.3", "prefsrc": "10.11.12.1"}], + "peers": [{ + "public_key": "peer-public-key", + "endpoint": "203.0.113.9:51820", + "allowed_ips": ["10.11.12.3/32"], + "latest_handshake": 123, + "rx_bytes": 456, + "tx_bytes": 789, + "persistent_keepalive": 25, + }], + } + monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "debian", "ip": "141.94.237.199"}) + monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=30: (0, json.dumps(payload), "")) + + with TestClient(app.app) as client: + response = client.get( + "/network/wireguard/guck-vps", + headers={"Authorization": "Bearer test-token"}, + ) + + assert response.status_code == 200 + assert response.json()["host"] == "guck-vps" + assert response.json()["peers"][0]["allowed_ips"] == ["10.11.12.3/32"] + assert "private" not in response.text.lower() + + +def test_wireguard_status_rejects_unknown_host_without_ssh(monkeypatch): + monkeypatch.setattr(app, "_find_inventory_host", lambda host: None) + monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("SSH must not run"))) + + with TestClient(app.app) as client: + response = client.get( + "/network/wireguard/does-not-exist", + headers={"Authorization": "Bearer test-token"}, + ) + + assert response.status_code == 404 + + def test_tts_generate_returns_cloned_wav(monkeypatch): captured = {} From 28a39483501c31035eb02af5c46266a46faa293b Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:08:09 +0200 Subject: [PATCH 23/72] Fix WireGuard status for unprivileged hosts (#38) --- app.py | 4 ++-- tests/test_app.py | 12 ++++++++++++ 2 files changed, 14 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index c43dec1..37b94c6 100644 --- a/app.py +++ b/app.py @@ -1231,7 +1231,7 @@ try: for address in item.get("addr_info", []): result["addresses"].append(address["local"] + "/" + str(address["prefixlen"])) result["routes"] = json.loads(run(["ip", "-j", "route", "show", "dev", "wg0"])) - rows = run(["wg", "show", "wg0", "dump"]).splitlines() + rows = run(["sudo", "-n", "wg", "show", "wg0", "dump"]).splitlines() if rows: interface = rows[0].split("\\t") result["listen_port"] = int(interface[2]) @@ -1244,7 +1244,7 @@ try: "latest_handshake": int(fields[4]), "rx_bytes": int(fields[5]), "tx_bytes": int(fields[6]), - "persistent_keepalive": int(fields[7]), + "persistent_keepalive": 0 if fields[7] == "off" else int(fields[7]), }) except Exception as exc: result["error"] = str(exc)[:300] diff --git a/tests/test_app.py b/tests/test_app.py index eb3be0d..c49a918 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -92,6 +92,18 @@ def test_wireguard_status_rejects_unknown_host_without_ssh(monkeypatch): assert response.status_code == 404 +def test_wireguard_status_command_uses_sudo_and_accepts_off_keepalive(): + import base64 + import re + + command = app._wireguard_status_command() + encoded = re.search(r"b64decode\('([^']+)'\)", command).group(1) + script = base64.b64decode(encoded).decode() + + assert '["sudo", "-n", "wg", "show", "wg0", "dump"]' in script + assert '0 if fields[7] == "off" else int(fields[7])' in script + + def test_tts_generate_returns_cloned_wav(monkeypatch): captured = {} From f2c5fa705175bd33e47ab628c716bbf28c9ff5f8 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:10:46 +0200 Subject: [PATCH 24/72] Add scoped removal for obsolete OVH-Hetzner peer (#39) --- app.py | 79 +++++++++++++++++++++++++++++++++++++++++++++++ tests/test_app.py | 31 +++++++++++++++++++ 2 files changed, 110 insertions(+) diff --git a/app.py b/app.py index 37b94c6..5a99cdc 100644 --- a/app.py +++ b/app.py @@ -1274,6 +1274,85 @@ async def network_wireguard_status(host: str, _=Depends(_verify)): return {"host": host, **result} +class WireGuardPeerRemoveRequest(BaseModel): + public_key: str + expected_allowed_ip: str + dry_run: bool = True + + +def _wireguard_remove_peer_command(public_key: str, expected_allowed_ip: str, dry_run: bool) -> str: + script = f'''import json, os, re, shutil, subprocess, time +from pathlib import Path + +public_key = {public_key!r} +expected = {expected_allowed_ip!r} +dry_run = {dry_run!r} +config = Path("/etc/wireguard/wg0.conf") +text = config.read_text() +sections = re.split(r"(?=^\\[Peer\\]\\s*$)", text, flags=re.M) +matches = [] +for index, section in enumerate(sections): + key_match = re.search(r"^PublicKey\\s*=\\s*(\\S+)\\s*$", section, re.M) + allowed_match = re.search(r"^AllowedIPs\\s*=\\s*(.+?)\\s*$", section, re.M) + allowed = [item.strip() for item in allowed_match.group(1).split(",")] if allowed_match else [] + if key_match and key_match.group(1) == public_key and expected in allowed: + matches.append((index, allowed)) +if len(matches) != 1: + print(json.dumps({{"error": "expected exactly one matching peer", "matches": len(matches)}})); raise SystemExit(2) +index, allowed = matches[0] +result = {{"status": "would_remove" if dry_run else "removed", "allowed_ips": allowed, "removed_routes": [], "backup": None}} +if dry_run: + print(json.dumps(result)); raise SystemExit(0) +backup = config.with_name("wg0.conf.butler-" + time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())) +shutil.copy2(config, backup) +result["backup"] = str(backup) +new_text = "".join(section for number, section in enumerate(sections) if number != index) +tmp = config.with_name("wg0.conf.butler-tmp") +tmp.write_text(new_text) +os.chmod(tmp, config.stat().st_mode) +os.chown(tmp, config.stat().st_uid, config.stat().st_gid) +os.replace(tmp, config) +try: + subprocess.run(["wg", "set", "wg0", "peer", public_key, "remove"], check=True, text=True, capture_output=True) + for route in allowed: + proc = subprocess.run(["ip", "route", "del", route, "dev", "wg0"], text=True, capture_output=True) + if proc.returncode == 0: result["removed_routes"].append(route) + peers = subprocess.run(["wg", "show", "wg0", "peers"], check=True, text=True, capture_output=True).stdout.split() + if public_key in peers: raise RuntimeError("peer still active") +except Exception: + shutil.copy2(backup, config) + raise +print(json.dumps(result)) +''' + encoded = base64.b64encode(script.encode()).decode() + return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' + + +@app.delete("/network/wireguard/{host}/peer") +async def network_wireguard_remove_peer(host: str, req: WireGuardPeerRemoveRequest, _=Depends(_verify)): + allowed = {"guck-vps": "10.7.1.0/24", "pfannkuchen": "10.200.200.60/32"} + if host not in allowed: + raise HTTPException(403, "Peer removal is restricted to the obsolete OVH-Hetzner transit") + if req.expected_allowed_ip != allowed[host]: + raise HTTPException(400, "Unexpected AllowedIP for this host") + if not re.fullmatch(r"[A-Za-z0-9+/]{43}=", req.public_key): + raise HTTPException(400, "Invalid WireGuard public key") + inventory = await asyncio.to_thread(_find_inventory_host, host) + if not inventory: + raise HTTPException(404, f"Host {host} not found") + target = f'{inventory["user"]}@{inventory["ip"]}' + command = _wireguard_remove_peer_command(req.public_key, req.expected_allowed_ip, req.dry_run) + rc, out, err = await asyncio.to_thread(_ssh, target, command, 45) + if rc != 0: + raise HTTPException(502, (err or out).strip()[-500:] or "WireGuard peer removal failed") + try: + result = json.loads(out) + except json.JSONDecodeError as exc: + raise HTTPException(502, "WireGuard peer removal returned invalid JSON") from exc + _audit(f"/network/wireguard/{host}/peer", "DELETE", 200, f"dry_run={req.dry_run}") + return {"host": host, **result} + + SYSCTL_AUDIT_KEYS = ( "net.core.default_qdisc", "net.core.rmem_default", diff --git a/tests/test_app.py b/tests/test_app.py index c49a918..726d11d 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -104,6 +104,37 @@ def test_wireguard_status_command_uses_sudo_and_accepts_off_keepalive(): assert '0 if fields[7] == "off" else int(fields[7])' in script +def test_wireguard_peer_remove_is_scoped_and_audited(monkeypatch): + calls = [] + monkeypatch.setattr(app, "_find_inventory_host", lambda host: {"user": "debian", "ip": "141.94.237.199"}) + monkeypatch.setattr(app, "_ssh", lambda host, command, timeout=30: (calls.append((host, command, timeout)) or (0, json.dumps({"status": "removed", "removed_routes": ["10.7.1.0/24"]}), ""))) + + with TestClient(app.app) as client: + response = client.request( + "DELETE", + "/network/wireguard/guck-vps/peer", + headers={"Authorization": "Bearer test-token"}, + json={"public_key": "A" * 43 + "=", "expected_allowed_ip": "10.7.1.0/24", "dry_run": False}, + ) + + assert response.status_code == 200 + assert response.json()["status"] == "removed" + assert calls[0][0] == "debian@141.94.237.199" + assert calls[0][2] == 45 + + +def test_wireguard_peer_remove_rejects_non_allowlisted_host(monkeypatch): + monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("SSH must not run"))) + with TestClient(app.app) as client: + response = client.request( + "DELETE", + "/network/wireguard/node7/peer", + headers={"Authorization": "Bearer test-token"}, + json={"public_key": "A" * 43 + "=", "expected_allowed_ip": "10.7.1.0/24", "dry_run": True}, + ) + assert response.status_code == 403 + + def test_tts_generate_returns_cloned_wav(monkeypatch): captured = {} From c7dd2ea97230add71a156d918fe298cc2cc1d43a Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:18:29 +0200 Subject: [PATCH 25/72] Add live capabilities and safety map (#40) --- app.py | 76 +++++++++++++++++++++++++++++++++++++++++++++++ tests/test_app.py | 27 +++++++++++++++++ 2 files changed, 103 insertions(+) diff --git a/app.py b/app.py index 5a99cdc..aa94ed4 100644 --- a/app.py +++ b/app.py @@ -264,6 +264,7 @@ async def root(): "openapi": "/openapi.json", "services": svc_list, "endpoints": { + "capabilities": "GET /capabilities - live machine-readable operation and safety map", "proxy": "GET/POST/PUT/DELETE /{service}/{path} - proxy to backend with auto-auth", "vm_list": "GET /vm/list", "vm_create": "POST /vm/create {node, ip, hostname, cores?, memory?, disk?}", @@ -288,6 +289,80 @@ async def root(): async def health(): return {"status": "ok", "vault_items": len(_vault_cache), "services": len(SERVICES), "version": VERSION} + +def _schema_contains_property(node, property_name: str, components: dict, seen: set[str] | None = None) -> bool: + """Resolve local OpenAPI refs and look for a request property.""" + seen = seen or set() + if isinstance(node, list): + return any(_schema_contains_property(item, property_name, components, seen) for item in node) + if not isinstance(node, dict): + return False + if node.get("name") == property_name or property_name in node.get("properties", {}): + return True + ref = node.get("$ref", "") + if ref.startswith("#/components/schemas/"): + name = ref.rsplit("/", 1)[-1] + if name in seen: + return False + return _schema_contains_property(components.get(name, {}), property_name, components, seen | {name}) + return any( + _schema_contains_property(value, property_name, components, seen) + for key, value in node.items() + if key != "properties" + ) + + +@app.get("/capabilities") +async def capabilities(_=Depends(_verify)): + """Live operation catalog with safety metadata for AI agents.""" + schema = app.openapi() + components = schema.get("components", {}).get("schemas", {}) + operations = [] + for path, methods in schema.get("paths", {}).items(): + if path == "/{service}/{path}" or path in {"/", "/health", "/openapi.json", "/docs", "/redoc"}: + continue + for method, operation in methods.items(): + if method.upper() not in {"GET", "POST", "PUT", "PATCH", "DELETE"}: + continue + mode = "read_only" if method.upper() == "GET" else "mutation" + serialized = {"parameters": operation.get("parameters", []), "requestBody": operation.get("requestBody", {})} + dry_run = _schema_contains_property(serialized, "dry_run", components) + critical = path.startswith(("/network/wireguard", "/caddy/")) + destructive = method.upper() == "DELETE" or any( + marker in path for marker in ("/destroy/", "/cleanup/", "/break-lock/", "/restore/") + ) + operations.append({ + "method": method.upper(), + "path": path, + "summary": operation.get("summary", ""), + "description": operation.get("description", ""), + "mode": mode, + "dry_run": dry_run, + "critical": critical, + "destructive": destructive, + "confirmation_required": mode == "mutation", + }) + operations.sort(key=lambda item: (item["path"], item["method"])) + counts = { + "total": len(operations), + "read_only": sum(item["mode"] == "read_only" for item in operations), + "mutations": sum(item["mode"] == "mutation" for item in operations), + "destructive": sum(item["destructive"] for item in operations), + } + return { + "schema_version": 1, + "service": "homelab-butler", + "version": VERSION, + "generated": datetime.now(timezone.utc).isoformat(), + "counts": counts, + "operations": operations, + "proxy": {"path": "/{service}/{path}", "note": "Generic backend proxy; inspect /info services and OpenAPI before use"}, + "model_contract": { + "instruction": "Prefer read_only operations. Before every mutation inspect its schema, use dry_run when available, and obtain confirmation for critical or destructive actions.", + "source_of_truth": "/openapi.json", + }, + } + def _classify_http_status(status_code: int, expected: set[int]) -> str: """Return a deterministic service state suitable for small models.""" if status_code in expected: @@ -408,6 +483,7 @@ async def info(_=Depends(_verify)): for name, cfg in SERVICES.items() }, "endpoints": { + "capabilities": "/capabilities", "status": "/status", "overview": "/overview?details=false", "audit": "/audit", diff --git a/tests/test_app.py b/tests/test_app.py index 726d11d..a9a0ae1 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -46,6 +46,33 @@ def test_health_exposes_current_version(): assert response.json()["version"] == app.VERSION == "2.3.5" +def test_capabilities_is_live_machine_readable_safety_map(): + with TestClient(app.app) as client: + response = client.get( + "/capabilities", + headers={"Authorization": "Bearer test-token"}, + ) + assert response.status_code == 200 + payload = response.json() + by_operation = {(item["method"], item["path"]): item for item in payload["operations"]} + assert ("GET", "/network/wireguard/{host}") in by_operation + assert by_operation[("GET", "/network/wireguard/{host}")]["mode"] == "read_only" + removal = by_operation[("DELETE", "/network/wireguard/{host}/peer")] + assert removal["mode"] == "mutation" + assert removal["dry_run"] is True + assert removal["critical"] is True + assert by_operation[("DELETE", "/vm/destroy/{vmid}")]["dry_run"] is True + assert all(item["path"] != "/{service}/{path}" for item in payload["operations"]) + assert payload["model_contract"]["instruction"].startswith("Prefer read_only") + + +def test_info_advertises_capabilities_endpoint(): + with TestClient(app.app) as client: + response = client.get("/info", headers={"Authorization": "Bearer test-token"}) + assert response.status_code == 200 + assert response.json()["endpoints"]["capabilities"] == "/capabilities" + + def test_wireguard_status_returns_redacted_live_state(monkeypatch): payload = { "interface": "wg0", From 5058ddf0a11ba2acb2ccd15130c12dabaed8e5ce Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:36:22 +0200 Subject: [PATCH 26/72] Add Butler operational safety suite (#41) --- .gitignore | 1 + app.py | 195 +++++++++++++++++++++++++++++++++++++++++++++- butler.yaml | 4 + compose.yaml | 2 + tests/test_app.py | 111 +++++++++++++++++++++++++- 5 files changed, 308 insertions(+), 5 deletions(-) diff --git a/.gitignore b/.gitignore index d5b4207..47bca58 100644 --- a/.gitignore +++ b/.gitignore @@ -3,3 +3,4 @@ vault-sync.log __pycache__/ *.pyc +state/ diff --git a/app.py b/app.py index aa94ed4..91be060 100644 --- a/app.py +++ b/app.py @@ -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 +import os, json, asyncio, logging, time, base64, re, subprocess, ipaddress, secrets, sqlite3 from datetime import datetime, timezone import httpx, yaml from typing import Literal @@ -46,6 +46,33 @@ _load_config() _audit_log: list[dict] = [] MAX_AUDIT = 500 +AUDIT_DB_PATH = os.environ.get("AUDIT_DB_PATH", "/data/state/audit.sqlite3") + + +def _redact_audit_detail(detail: str) -> str: + return re.sub( + r"(?i)\b(token|password|api[_-]?key|secret)=([^\s]+)", + lambda match: f"{match.group(1)}=[REDACTED]", + detail, + )[:200] + + +def _init_audit_db() -> bool: + try: + os.makedirs(os.path.dirname(AUDIT_DB_PATH), exist_ok=True) + with sqlite3.connect(AUDIT_DB_PATH) as db: + db.execute("""CREATE TABLE IF NOT EXISTS audit ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ts TEXT NOT NULL, + endpoint TEXT NOT NULL, + method TEXT NOT NULL, + status INTEGER NOT NULL, + detail TEXT NOT NULL, + dry_run INTEGER NOT NULL + )""") + return True + except (OSError, sqlite3.Error): + return False def _audit(endpoint: str, method: str, status: int, detail: str = "", dry_run: bool = False): entry = { @@ -53,12 +80,22 @@ def _audit(endpoint: str, method: str, status: int, detail: str = "", dry_run: b "endpoint": endpoint, "method": method, "status": status, - "detail": detail[:200], + "detail": _redact_audit_detail(detail), "dry_run": dry_run, } _audit_log.append(entry) if len(_audit_log) > MAX_AUDIT: _audit_log.pop(0) + try: + if _init_audit_db(): + with sqlite3.connect(AUDIT_DB_PATH) as db: + db.execute( + "INSERT INTO audit (ts, endpoint, method, status, detail, dry_run) VALUES (?, ?, ?, ?, ?, ?)", + (entry["ts"], entry["endpoint"], entry["method"], entry["status"], entry["detail"], int(entry["dry_run"])), + ) + db.execute("DELETE FROM audit WHERE id NOT IN (SELECT id FROM audit ORDER BY id DESC LIMIT ?)", (MAX_AUDIT,)) + except (OSError, sqlite3.Error): + pass API_DIR = os.environ.get("API_KEY_DIR", "/data/api") VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache") @@ -92,6 +129,7 @@ async def _periodic_cache_reload(): async def lifespan(app: FastAPI): _load_config() _load_vault_cache() + _init_audit_db() task = asyncio.create_task(_periodic_cache_reload()) yield task.cancel() @@ -265,6 +303,9 @@ async def root(): "services": svc_list, "endpoints": { "capabilities": "GET /capabilities - live machine-readable operation and safety map", + "doctor": "GET /doctor/{target} - correlated service/host/backup/disk diagnosis", + "drift": "GET /drift - inventory coverage gaps", + "maintenance_preflight": "GET /maintenance/preflight?action=general&target=HOST - read-only safety gate", "proxy": "GET/POST/PUT/DELETE /{service}/{path} - proxy to backend with auto-auth", "vm_list": "GET /vm/list", "vm_create": "POST /vm/create {node, ip, hostname, cores?, memory?, disk?}", @@ -457,6 +498,17 @@ async def status(_=Depends(_verify)): @app.get("/audit") async def audit(_=Depends(_verify), limit: int = Query(50, le=MAX_AUDIT)): """Recent API calls (newest first).""" + try: + if os.path.exists(AUDIT_DB_PATH): + with sqlite3.connect(AUDIT_DB_PATH) as db: + db.row_factory = sqlite3.Row + rows = db.execute( + "SELECT ts, endpoint, method, status, detail, dry_run FROM audit ORDER BY id DESC LIMIT ?", + (limit,), + ).fetchall() + return [dict(row) | {"dry_run": bool(row["dry_run"])} for row in rows] + except (OSError, sqlite3.Error): + pass return list(reversed(_audit_log[-limit:])) @app.post("/config/reload") @@ -484,6 +536,9 @@ async def info(_=Depends(_verify)): }, "endpoints": { "capabilities": "/capabilities", + "doctor": "/doctor/{target}", + "drift": "/drift", + "maintenance_preflight": "/maintenance/preflight", "status": "/status", "overview": "/overview?details=false", "audit": "/audit", @@ -599,8 +654,11 @@ async def _collect_backup_status(concurrency: int = 10) -> dict: """Query VM backups concurrently; one slow host no longer blocks all others serially.""" semaphore = asyncio.Semaphore(concurrency) hosts = [host for host in await _get_inventory_hosts_async() if not host["name"].startswith("node")] + exempt_hosts = set(_config.get("backup", {}).get("exempt_hosts", [])) async def inspect_backup(host: dict): + if host["name"] in exempt_hosts: + return host["name"], {"state": "exempt", "ok": True, "reason": "backup policy exemption"} async with semaphore: rc, out, err = await asyncio.to_thread( _ssh, @@ -612,7 +670,7 @@ async def _collect_backup_status(concurrency: int = 10) -> dict: pairs = await asyncio.gather(*(inspect_backup(host) for host in hosts)) results = dict(pairs) - summary = {"total": len(results), "healthy": 0, "warning": 0, "critical": 0, "unknown": 0} + summary = {"total": len(results), "healthy": 0, "warning": 0, "critical": 0, "unknown": 0, "exempt": 0} for item in results.values(): summary[item["state"]] += 1 return {"summary": summary, "hosts": results} @@ -780,6 +838,135 @@ def _add_component(summary: dict, bucket: dict, state: str): bucket[normalized] += 1 +async def _collect_operational_snapshot() -> dict: + services, hosts, backups, disks = await asyncio.gather( + _collect_service_status(), _collect_health_all(), _collect_backup_status(), _collect_disk_usage() + ) + return {"services": services, "hosts": hosts, "backups": backups, "disks": disks} + + +@app.get("/doctor/{target}") +async def doctor(target: str, _=Depends(_verify)): + """Correlate service, host, backup and disk layers for one known target.""" + if not re.fullmatch(r"[A-Za-z0-9_.-]+", target): + raise HTTPException(400, "Invalid target") + snapshot = await _collect_operational_snapshot() + layers = {} + findings = [] + if target in snapshot["services"]: + service = snapshot["services"][target] + layers["service"] = service + if service.get("status") != "healthy": + findings.append({"severity": "critical" if service.get("status") in ("offline", "auth_failed", "misconfigured") else "warning", "code": "service_unhealthy", "message": service.get("message", "Service probe failed")}) + if target in snapshot["hosts"]: + host = snapshot["hosts"][target] + layers["host"] = host + if not host.get("reachable"): + findings.append({"severity": "critical", "code": "host_unreachable", "message": "Host is not reachable over SSH"}) + bad = [line for line in host.get("containers", []) if "unhealthy" in line.lower() or "restarting" in line.lower()] + if bad: + findings.append({"severity": "critical", "code": "container_unhealthy", "message": bad[0][:200]}) + backup = snapshot["backups"].get("hosts", {}).get(target) + if backup is not None: + layers["backup"] = backup + if backup.get("state") not in ("healthy", "exempt"): + findings.append({"severity": "critical" if backup.get("state") in ("critical", "unknown") else "warning", "code": "backup_unhealthy", "message": "Backup is stale or could not be verified"}) + disk = snapshot["disks"].get(target) + if disk is not None: + layers["disk"] = disk + pct = int(str(disk.get("pct", "0")).rstrip("%") or 0) + if pct >= 80: + findings.append({"severity": "critical" if pct >= 90 else "warning", "code": "disk_high", "message": f"Root filesystem usage is {pct}%"}) + if not layers: + raise HTTPException(404, "Target not found") + state = "critical" if any(item["severity"] == "critical" for item in findings) else "warning" if findings else "healthy" + return {"target": target, "state": state, "findings": findings, "layers": layers, "next_checks": [f"/logs/{target}/{{container}}", f"/system/forensics/{target}"] if "host" in layers else []} + + +@app.get("/drift") +async def drift(_=Depends(_verify)): + """Report coverage drift between inventory, backup and disk collectors.""" + snapshot = await _collect_operational_snapshot() + inventory = set(snapshot["hosts"]) + managed_hosts = {name for name in inventory if not name.startswith("node")} + backups = set(snapshot["backups"].get("hosts", {})) + disks = set(snapshot["disks"]) + findings = [] + for name in sorted(managed_hosts - backups): + findings.append({"severity": "warning", "code": "inventory_missing_backup", "target": name}) + for name in sorted(inventory - disks): + findings.append({"severity": "warning", "code": "inventory_missing_disk", "target": name}) + for name in sorted(backups - inventory): + findings.append({"severity": "warning", "code": "backup_without_inventory", "target": name}) + return { + "state": "warning" if findings else "healthy", + "findings": findings, + "coverage": {"inventory": len(inventory), "backups": len(backups), "disks": len(disks)}, + "model_contract": {"instruction": "Treat findings as coverage gaps, not proof that the target is offline."}, + } + + +async def _collect_active_backups(concurrency: int = 10) -> dict: + hosts = [host for host in await _get_inventory_hosts_async() if not host["name"].startswith("node")] + exempt = set(_config.get("backup", {}).get("exempt_hosts", [])) + semaphore = asyncio.Semaphore(concurrency) + + async def check(host: dict): + if host["name"] in exempt: + return host["name"], "exempt" + async with semaphore: + rc, out, _err = await asyncio.to_thread( + _ssh, + f'{host["user"]}@{host["ip"]}', + "sudo -n systemctl is-active borg-backup.service 2>/dev/null || true", + 10, + ) + state = out.strip().splitlines()[-1] if out.strip() else "unknown" + return host["name"], state if rc == 0 else "unknown" + + return dict(await asyncio.gather(*(check(host) for host in hosts))) + + +@app.get("/maintenance/preflight") +async def maintenance_preflight( + action: Literal["general", "docker", "network", "vm"] = Query("general"), + target: str | None = Query(None), + _=Depends(_verify), +): + """Read-only safety gate before maintenance or mutations.""" + if target and not re.fullmatch(r"[A-Za-z0-9_.-]+", target): + raise HTTPException(400, "Invalid target") + snapshot, active_backups = await asyncio.gather(_collect_operational_snapshot(), _collect_active_backups()) + if target and target not in snapshot["hosts"] and target not in snapshot["services"]: + raise HTTPException(404, "Target not found") + selected = {target} if target else set(snapshot["hosts"]) + blockers = [] + warnings = [] + for name in sorted(selected): + host = snapshot["hosts"].get(name) + if host and not host.get("reachable"): + blockers.append({"code": "host_unreachable", "target": name}) + if host and any("unhealthy" in line.lower() or "restarting" in line.lower() for line in host.get("containers", [])): + blockers.append({"code": "container_unhealthy", "target": name}) + if active_backups.get(name) in ("active", "activating"): + blockers.append({"code": "backup_active", "target": name}) + backup = snapshot["backups"].get("hosts", {}).get(name, {}) + if backup.get("state") not in (None, "healthy", "exempt"): + warnings.append({"code": "backup_unhealthy", "target": name}) + disk = snapshot["disks"].get(name, {}) + pct = int(str(disk.get("pct", "0")).rstrip("%") or 0) + if pct >= 80: + warnings.append({"code": "disk_high", "target": name, "pct": pct}) + return { + "safe": not blockers, + "action": action, + "target": target, + "blockers": blockers, + "warnings": warnings, + "model_contract": {"instruction": "Do not start the requested maintenance while safe is false."}, + } + + @app.get("/overview", response_model=OverviewResponse, response_model_exclude_none=True) async def overview(details: bool = Query(False), _=Depends(_verify)): """Compact deterministic homelab verdict designed for small language models.""" @@ -843,7 +1030,7 @@ async def overview(details: bool = Query(False), _=Depends(_verify)): for name in sorted(backups.get("hosts", {})): item = backups["hosts"][name] raw_state = item.get("state", "unknown") - state = "healthy" if raw_state == "healthy" else "critical" if raw_state in ("critical", "unknown") else "warning" + state = "healthy" if raw_state in ("healthy", "exempt") else "critical" if raw_state in ("critical", "unknown") else "warning" _add_component(summary, component_summary["backups"], state) if state != "healthy": findings.append({ diff --git a/butler.yaml b/butler.yaml index 4624d82..0f8ca9e 100644 --- a/butler.yaml +++ b/butler.yaml @@ -105,6 +105,10 @@ services: health_path: "/" timeout: 300 +backup: + # guck-vps contains Git-managed edge configuration and has no Borgmatic installation. + exempt_hosts: [guck-vps] + # VM lifecycle settings vm: automation_host: "sascha@10.5.85.5" diff --git a/compose.yaml b/compose.yaml index b593f71..e0d1f56 100644 --- a/compose.yaml +++ b/compose.yaml @@ -11,10 +11,12 @@ services: - /home/sascha/.ssh:/root/.ssh:ro - ./butler.yaml:/data/butler.yaml:ro - ./app.py:/app/app.py:ro + - ./state:/data/state environment: - API_KEY_DIR=/data/api - VAULT_CACHE_DIR=/data/vault-cache - BUTLER_TOKEN=${BUTLER_TOKEN} + - AUDIT_DB_PATH=/data/state/audit.sqlite3 volumes: vault-cache: diff --git a/tests/test_app.py b/tests/test_app.py index a9a0ae1..645220a 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -71,6 +71,28 @@ def test_info_advertises_capabilities_endpoint(): response = client.get("/info", headers={"Authorization": "Bearer test-token"}) assert response.status_code == 200 assert response.json()["endpoints"]["capabilities"] == "/capabilities" + assert response.json()["endpoints"]["doctor"] == "/doctor/{target}" + assert response.json()["endpoints"]["drift"] == "/drift" + assert response.json()["endpoints"]["maintenance_preflight"] == "/maintenance/preflight" + + +def test_audit_persists_and_redacts_secrets(tmp_path, monkeypatch): + db = tmp_path / "audit.sqlite3" + monkeypatch.setattr(app, "AUDIT_DB_PATH", str(db)) + app._init_audit_db() + app._audit("/danger", "POST", 200, "host=x token=abc password=hunter2 api_key=secret") + app._audit_log.clear() + + with TestClient(app.app) as client: + response = client.get("/audit", headers={"Authorization": "Bearer test-token"}) + + assert response.status_code == 200 + entry = response.json()[0] + assert entry["endpoint"] == "/danger" + assert "abc" not in entry["detail"] + assert "hunter2" not in entry["detail"] + assert "secret" not in entry["detail"] + assert entry["detail"].count("[REDACTED]") == 3 def test_wireguard_status_returns_redacted_live_state(monkeypatch): @@ -633,6 +655,7 @@ def test_backup_collection_runs_hosts_concurrently(monkeypatch): {"name": f"vm-{index}", "user": "sascha", "ip": f"10.1.1.{index}"} for index in range(1, 5) ]) + monkeypatch.setattr(app, "_config", {}) def fake_ssh(*_args, **_kwargs): nonlocal active, max_active @@ -647,7 +670,24 @@ def test_backup_collection_runs_hosts_concurrently(monkeypatch): assert max_active > 1 assert result["summary"] == { - "total": 4, "healthy": 4, "warning": 0, "critical": 0, "unknown": 0 + "total": 4, "healthy": 4, "warning": 0, "critical": 0, "unknown": 0, "exempt": 0 + } + + +def test_backup_policy_exempts_host_without_borg_call(monkeypatch): + monkeypatch.setattr(app, "_get_inventory_hosts", lambda: [ + {"name": "guck-vps", "user": "debian", "ip": "141.94.237.199"} + ]) + monkeypatch.setattr(app, "_config", {"backup": {"exempt_hosts": ["guck-vps"]}}) + monkeypatch.setattr(app, "_ssh", lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("must not call borg"))) + + result = asyncio.run(app._collect_backup_status()) + + assert result["summary"] == { + "total": 1, "healthy": 0, "warning": 0, "critical": 0, "unknown": 0, "exempt": 1 + } + assert result["hosts"]["guck-vps"] == { + "state": "exempt", "ok": True, "reason": "backup policy exemption" } @@ -661,6 +701,75 @@ def test_overview_openapi_has_stable_enums_and_schema(): assert finding_schema["properties"]["severity"]["enum"] == ["healthy", "warning", "critical"] +def test_doctor_correlates_host_layers(monkeypatch): + async def snapshot(): + return { + "services": {}, + "hosts": {"guck-vps": {"reachable": True, "containers": ["caddy: Up 3 days"]}}, + "backups": {"summary": {}, "hosts": {"guck-vps": {"state": "exempt", "ok": True}}}, + "disks": {"guck-vps": {"pct": "8%"}}, + } + + monkeypatch.setattr(app, "_collect_operational_snapshot", snapshot) + with TestClient(app.app) as client: + response = client.get("/doctor/guck-vps", headers={"Authorization": "Bearer test-token"}) + + assert response.status_code == 200 + result = response.json() + assert result["state"] == "healthy" + assert result["layers"]["host"]["reachable"] is True + assert result["layers"]["backup"]["state"] == "exempt" + assert result["layers"]["disk"]["pct"] == "8%" + assert result["findings"] == [] + + +def test_drift_reports_inventory_coverage_gaps(monkeypatch): + async def snapshot(): + return { + "services": {}, + "hosts": {"vm-a": {"reachable": True}, "vm-b": {"reachable": True}, "node1": {"reachable": True}}, + "backups": {"summary": {}, "hosts": {"vm-a": {"state": "healthy"}, "orphan": {"state": "healthy"}}}, + "disks": {"vm-a": {"pct": "10%"}, "node1": {"pct": "20%"}}, + } + + monkeypatch.setattr(app, "_collect_operational_snapshot", snapshot) + with TestClient(app.app) as client: + response = client.get("/drift", headers={"Authorization": "Bearer test-token"}) + + assert response.status_code == 200 + result = response.json() + assert result["state"] == "warning" + assert {item["code"] for item in result["findings"]} == { + "inventory_missing_backup", "inventory_missing_disk", "backup_without_inventory" + } + + +def test_maintenance_preflight_blocks_active_target_backup(monkeypatch): + async def snapshot(): + return { + "services": {}, + "hosts": {"emby-chris": {"reachable": True, "containers": ["emby: Up 2 days"]}}, + "backups": {"summary": {}, "hosts": {"emby-chris": {"state": "healthy"}}}, + "disks": {"emby-chris": {"pct": "30%"}}, + } + + async def active_backups(): + return {"emby-chris": "active"} + + monkeypatch.setattr(app, "_collect_operational_snapshot", snapshot) + monkeypatch.setattr(app, "_collect_active_backups", active_backups) + with TestClient(app.app) as client: + response = client.get( + "/maintenance/preflight?action=docker&target=emby-chris", + headers={"Authorization": "Bearer test-token"}, + ) + + assert response.status_code == 200 + result = response.json() + assert result["safe"] is False + assert result["blockers"] == [{"code": "backup_active", "target": "emby-chris"}] + + def test_overview_is_compact_deterministic_and_light_model_friendly(monkeypatch): async def service_data(): return { From a1a69567d4c851ea4fb5c78556b91084b02ccdae Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:56:23 +0200 Subject: [PATCH 27/72] feat(ui): update app.py --- app.py | 70 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 68 insertions(+), 2 deletions(-) diff --git a/app.py b/app.py index 91be060..2b6da2e 100644 --- a/app.py +++ b/app.py @@ -7,7 +7,7 @@ import httpx, yaml from typing import Literal from pydantic import BaseModel, Field from fastapi import FastAPI, Request, HTTPException, Depends, Query -from fastapi.responses import JSONResponse, RedirectResponse, Response +from fastapi.responses import JSONResponse, RedirectResponse, Response, HTMLResponse from contextlib import asynccontextmanager log = logging.getLogger("butler") @@ -17,6 +17,7 @@ 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")) # --- Config loading --- @@ -227,12 +228,37 @@ async def _dockhand_login(client): # --- Auth --- +_ui_sessions: dict[str, dict] = {} +UI_SESSION_TTL = 8 * 60 * 60 + + +class UiLoginRequest(BaseModel): + token: str + + +def _ui_session(request: Request) -> dict | None: + session_id = request.cookies.get("butler_session", "") + session = _ui_sessions.get(session_id) + if not session: + return None + if session["expires"] <= time.time(): + _ui_sessions.pop(session_id, None) + return None + return session + def _verify(request: Request): if not BUTLER_TOKEN: return auth = request.headers.get("authorization", "") - if auth != f"Bearer {BUTLER_TOKEN}": + if secrets.compare_digest(auth, f"Bearer {BUTLER_TOKEN}"): + return + session = _ui_session(request) + if not session: raise HTTPException(401, "Invalid token") + if request.method not in {"GET", "HEAD", "OPTIONS"}: + csrf = request.headers.get("x-csrf-token", "") + if not csrf or not secrets.compare_digest(csrf, session["csrf"]): + raise HTTPException(403, "Invalid CSRF token") def _get_key(cfg): vault_key = cfg.get("vault_key") @@ -290,6 +316,45 @@ def _inventory_hosts(text: str) -> list[dict]: # --- Routes --- +@app.get("/ui", response_class=HTMLResponse) +async def ui(): + try: + return HTMLResponse(open(UI_PATH, encoding="utf-8").read()) + except FileNotFoundError: + raise HTTPException(503, "Butler UI asset is missing") + + +@app.post("/ui/login") +async def ui_login(payload: UiLoginRequest): + if not BUTLER_TOKEN or not secrets.compare_digest(payload.token, BUTLER_TOKEN): + raise HTTPException(401, "Invalid token") + session_id = secrets.token_urlsafe(32) + csrf = secrets.token_urlsafe(24) + _ui_sessions[session_id] = {"csrf": csrf, "expires": time.time() + UI_SESSION_TTL} + response = JSONResponse({"authenticated": True, "expires_in": UI_SESSION_TTL}) + response.set_cookie("butler_session", session_id, max_age=UI_SESSION_TTL, httponly=True, samesite="strict", path="/") + response.set_cookie("butler_csrf", csrf, max_age=UI_SESSION_TTL, httponly=False, samesite="strict", path="/") + return response + + +@app.get("/ui/session") +async def ui_session(request: Request): + return {"authenticated": _ui_session(request) is not None} + + +@app.post("/ui/logout") +async def ui_logout(request: Request): + session = _ui_session(request) + if session: + csrf = request.headers.get("x-csrf-token", "") + if not csrf or not secrets.compare_digest(csrf, session["csrf"]): + raise HTTPException(403, "Invalid CSRF token") + _ui_sessions.pop(request.cookies.get("butler_session", ""), None) + response = JSONResponse({"authenticated": False}) + response.delete_cookie("butler_session", path="/") + response.delete_cookie("butler_csrf", path="/") + return response + @app.get("/") async def root(): """AI self-onboarding: returns all available endpoints and services.""" @@ -536,6 +601,7 @@ async def info(_=Depends(_verify)): }, "endpoints": { "capabilities": "/capabilities", + "ui": "/ui", "doctor": "/doctor/{target}", "drift": "/drift", "maintenance_preflight": "/maintenance/preflight", From 5723eebfb47d657265a553ac908c80065db87a30 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:56:24 +0200 Subject: [PATCH 28/72] feat(ui): update tests/test_app.py --- tests/test_app.py | 41 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/tests/test_app.py b/tests/test_app.py index 645220a..82356fd 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -46,6 +46,46 @@ def test_health_exposes_current_version(): assert response.json()["version"] == app.VERSION == "2.3.5" +def test_ui_serves_self_contained_operator_console(): + with TestClient(app.app) as client: + response = client.get("/ui") + assert response.status_code == 200 + assert "Pfannkuchen Butler" in response.text + assert 'id="operations-grid"' in response.text + assert 'id="doctor-form"' in response.text + assert 'id="preflight-form"' in response.text + assert "localStorage" not in response.text + + +def test_ui_asset_is_mounted_read_only_in_compose(): + from pathlib import Path + import yaml + compose = yaml.safe_load(Path(app.__file__).with_name("compose.yaml").read_text(encoding="utf-8")) + mounts = compose["services"]["homelab-butler"]["volumes"] + assert "./ui.html:/app/ui.html:ro" in mounts + + +def test_ui_session_login_uses_httponly_cookie_and_csrf(): + app._ui_sessions.clear() + with TestClient(app.app) as client: + denied = client.post("/ui/login", json={"token": "wrong"}) + assert denied.status_code == 401 + + login = client.post("/ui/login", json={"token": "test-token"}) + assert login.status_code == 200 + assert "HttpOnly" in login.headers.get("set-cookie", "") + csrf = client.cookies.get("butler_csrf") + assert csrf + + capabilities = client.get("/capabilities") + assert capabilities.status_code == 200 + + blocked = client.post("/config/reload") + assert blocked.status_code == 403 + allowed = client.post("/config/reload", headers={"X-CSRF-Token": csrf}) + assert allowed.status_code == 200 + + def test_capabilities_is_live_machine_readable_safety_map(): with TestClient(app.app) as client: response = client.get( @@ -74,6 +114,7 @@ def test_info_advertises_capabilities_endpoint(): assert response.json()["endpoints"]["doctor"] == "/doctor/{target}" assert response.json()["endpoints"]["drift"] == "/drift" assert response.json()["endpoints"]["maintenance_preflight"] == "/maintenance/preflight" + assert response.json()["endpoints"]["ui"] == "/ui" def test_audit_persists_and_redacts_secrets(tmp_path, monkeypatch): From 6d05c16f589250059cfe3bf9043c44f655daf69d Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:56:24 +0200 Subject: [PATCH 29/72] feat(ui): update compose.yaml --- compose.yaml | 1 + 1 file changed, 1 insertion(+) diff --git a/compose.yaml b/compose.yaml index e0d1f56..a07a30f 100644 --- a/compose.yaml +++ b/compose.yaml @@ -11,6 +11,7 @@ services: - /home/sascha/.ssh:/root/.ssh:ro - ./butler.yaml:/data/butler.yaml:ro - ./app.py:/app/app.py:ro + - ./ui.html:/app/ui.html:ro - ./state:/data/state environment: - API_KEY_DIR=/data/api From 94f6c6b241178c248eb323430c0a835e4b282f1d Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:56:25 +0200 Subject: [PATCH 30/72] feat(ui): add operator console asset --- ui.html | 124 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 124 insertions(+) create mode 100644 ui.html diff --git a/ui.html b/ui.html new file mode 100644 index 0000000..9ad3d27 --- /dev/null +++ b/ui.html @@ -0,0 +1,124 @@ + + + + + + + Pfannkuchen Butler + + + + + + +
+ + + + From d23fdc7896a258ffbd5b89777feb5ef22c1a34c6 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 21:28:56 +0200 Subject: [PATCH 31/72] Add Emby account-sharing analysis and read-only UI access (app.py) --- app.py | 280 +++++++++++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 271 insertions(+), 9 deletions(-) diff --git a/app.py b/app.py index 2b6da2e..bac1a05 100644 --- a/app.py +++ b/app.py @@ -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 +import os, json, asyncio, logging, time, base64, re, subprocess, ipaddress, secrets, sqlite3, math from datetime import datetime, timezone import httpx, yaml from typing import Literal @@ -246,6 +246,24 @@ def _ui_session(request: Request) -> dict | None: return None return session + +def _issue_ui_session(read_only: bool) -> JSONResponse: + session_id = secrets.token_urlsafe(32) + csrf = secrets.token_urlsafe(24) + _ui_sessions[session_id] = { + "csrf": csrf, + "expires": time.time() + UI_SESSION_TTL, + "read_only": read_only, + } + response = JSONResponse({ + "authenticated": True, + "read_only": read_only, + "expires_in": UI_SESSION_TTL, + }) + response.set_cookie("butler_session", session_id, max_age=UI_SESSION_TTL, httponly=True, samesite="strict", path="/") + response.set_cookie("butler_csrf", csrf, max_age=UI_SESSION_TTL, httponly=False, samesite="strict", path="/") + return response + def _verify(request: Request): if not BUTLER_TOKEN: return @@ -256,10 +274,253 @@ def _verify(request: Request): if not session: raise HTTPException(401, "Invalid token") if request.method not in {"GET", "HEAD", "OPTIONS"}: + if session.get("read_only", False): + raise HTTPException(403, "Anonymous UI session is read-only") csrf = request.headers.get("x-csrf-token", "") if not csrf or not secrets.compare_digest(csrf, session["csrf"]): raise HTTPException(403, "Invalid CSRF token") +def _emby_network_identity(endpoint: str) -> dict: + """Normalize an Emby endpoint without treating IPv6 privacy addresses as new households.""" + raw = str(endpoint or "").strip() + if raw.startswith("[") and "]" in raw: + raw = raw[1:raw.index("]")] + try: + address = ipaddress.ip_address(raw) + except ValueError: + if raw.count(":") == 1: + raw = raw.rsplit(":", 1)[0] + try: + address = ipaddress.ip_address(raw) + except ValueError as exc: + raise ValueError("Invalid Emby remote endpoint") from exc + if address.version == 4: + network = ipaddress.ip_network(f"{address}/32", strict=False) + parent = network + else: + network = ipaddress.ip_network(f"{address}/64", strict=False) + parent = ipaddress.ip_network(f"{address}/48", strict=False) + return { + "ip": str(address), + "version": address.version, + "network": str(network), + "parent": str(parent), + "identity": str(parent if address.version == 6 else network), + "public": address.is_global, + } + + +def _emby_location(metric: dict) -> dict: + def coordinate(name): + try: + return float(metric.get(name, 0)) + except (TypeError, ValueError): + return 0.0 + return { + "city": metric.get("city", ""), + "region": metric.get("region", ""), + "country": metric.get("countryCode", ""), + "latitude": coordinate("latitude"), + "longitude": coordinate("longitude"), + } + + +def _analyze_emby_sharing(series: list[dict], step_seconds: int) -> dict: + observations: dict[str, dict[int, dict[str, dict]]] = {} + tracks: dict[str, dict[str, dict]] = {} + identities_by_user: dict[str, set[str]] = {} + servers_by_user: dict[str, set[str]] = {} + for item in series: + metric = item.get("metric", {}) + username = str(metric.get("username", "")).strip() + if not username: + continue + try: + network = _emby_network_identity(metric.get("remoteEndPoint", "")) + except ValueError: + continue + if not network["public"]: + continue + evidence = { + **network, + "server": metric.get("job", ""), + "location": _emby_location(metric), + } + identities_by_user.setdefault(username, set()).add(network["identity"]) + servers_by_user.setdefault(username, set()).add(str(metric.get("job", ""))) + track = tracks.setdefault(username, {}).setdefault(network["identity"], {"evidence": evidence, "timestamps": []}) + for value in item.get("values", []): + if not isinstance(value, list) or len(value) < 2 or str(value[1]).lower() in {"0", "nan"}: + continue + timestamp = int(float(value[0])) + track["timestamps"].append(timestamp) + observations.setdefault(username, {}).setdefault(timestamp, {}).setdefault(network["identity"], evidence) + + raw_events = [] + for username, timeline in observations.items(): + buckets = [] + for timestamp in sorted(timeline): + evidence = timeline[timestamp] + if len(evidence) >= 2: + buckets.append((timestamp, tuple(sorted(evidence)), evidence)) + current = None + for timestamp, identity_key, evidence in buckets: + if current and current["identity_key"] == identity_key and timestamp - current["end_ts"] <= step_seconds * 2: + current["end_ts"] = timestamp + current["samples"] += 1 + continue + if current and current["samples"] >= 2 and current["end_ts"] - current["start_ts"] + step_seconds >= 360: + raw_events.append(current) + current = { + "username": username, "identity_key": identity_key, "start_ts": timestamp, + "end_ts": timestamp, "samples": 1, "evidence": list(evidence.values()), + } + if current and current["samples"] >= 2 and current["end_ts"] - current["start_ts"] + step_seconds >= 360: + raw_events.append(current) + + events = [{ + "type": "concurrent_networks", + "severity": "high", + "username": item["username"], + "start": datetime.fromtimestamp(item["start_ts"], timezone.utc).isoformat(), + "end": datetime.fromtimestamp(item["end_ts"], timezone.utc).isoformat(), + "duration_seconds": item["end_ts"] - item["start_ts"] + step_seconds, + "samples": item["samples"], + "evidence": item["evidence"], + "reason": "Zeitgleiche Nutzung desselben Emby-Benutzers aus unterschiedlichen öffentlichen Netzen", + } for item in raw_events] + + def distance_km(first: dict, second: dict) -> float: + lat1, lon1 = first["latitude"], first["longitude"] + lat2, lon2 = second["latitude"], second["longitude"] + if not all((-90 <= lat <= 90 and -180 <= lon <= 180) for lat, lon in ((lat1, lon1), (lat2, lon2))): + return 0.0 + phi1, phi2 = math.radians(lat1), math.radians(lat2) + dphi, dlambda = math.radians(lat2 - lat1), math.radians(lon2 - lon1) + value = math.sin(dphi / 2) ** 2 + math.cos(phi1) * math.cos(phi2) * math.sin(dlambda / 2) ** 2 + return 6371.0 * 2 * math.atan2(math.sqrt(value), math.sqrt(max(0.0, 1 - value))) + + travel_events = [] + for username, user_tracks in tracks.items(): + intervals = [] + for identity, track in user_tracks.items(): + current = None + for timestamp in sorted(set(track["timestamps"])): + if current and timestamp - current["end"] <= step_seconds * 2: + current["end"] = timestamp + current["samples"] += 1 + else: + if current and current["samples"] >= 2: + intervals.append(current) + current = {"identity": identity, "start": timestamp, "end": timestamp, "samples": 1, "evidence": track["evidence"]} + if current and current["samples"] >= 2: + intervals.append(current) + intervals.sort(key=lambda item: item["start"]) + for previous, current in zip(intervals, intervals[1:]): + if previous["identity"] == current["identity"] or current["start"] <= previous["end"]: + continue + distance = distance_km(previous["evidence"]["location"], current["evidence"]["location"]) + gap_hours = max((current["start"] - previous["end"]) / 3600, 1 / 60) + speed = distance / gap_hours + if distance < 300 or speed <= 1000: + continue + travel_events.append({ + "type": "impossible_travel", "severity": "medium", "username": username, + "start": datetime.fromtimestamp(previous["end"], timezone.utc).isoformat(), + "end": datetime.fromtimestamp(current["start"], timezone.utc).isoformat(), + "duration_seconds": current["start"] - previous["end"], + "samples": previous["samples"] + current["samples"], + "distance_km": round(distance, 1), "required_speed_kmh": round(speed, 1), + "evidence": [previous["evidence"], current["evidence"]], + "reason": "Geografischer Wechsel zwischen öffentlichen Netzen wäre in der verfügbaren Zeit nicht plausibel", + }) + events.extend(travel_events) + events.sort(key=lambda item: item["start"], reverse=True) + flagged = {item["username"] for item in events} + users = [{ + "username": username, + "risk": "high" if any(item["username"] == username and item["type"] == "concurrent_networks" for item in events) else ("medium" if username in flagged else "none"), + "events": sum(item["username"] == username for item in events), + "network_identities": len(identities_by_user.get(username, set())), + "servers": sorted(servers_by_user.get(username, set())), + } for username in sorted(observations)] + return { + "summary": { + "users_analyzed": len(observations), + "flagged_users": len(flagged), + "concurrent_events": len(raw_events), + "impossible_travel_events": len(travel_events), + }, + "users": users, + "events": events, + } + + +async def _fetch_emby_session_history(days: int, server: str) -> tuple[list[dict], int]: + cfg = SERVICES.get("grafana") + if not cfg: + raise HTTPException(503, "Grafana service is not configured") + request_data = _service_auth(cfg) + datasource_uid = os.environ.get("EMBY_PROMETHEUS_UID", "bdpu4276997nkc") + labels = "job,username,remoteEndPoint,city,region,countryCode,latitude,longitude" + selector = 'emby_sessions{username!=""}' + if server != "all": + selector = f'emby_sessions{{username!="",job="{server}"}}' + query = f"max by ({labels}) ({selector})" + end = int(time.time()) + start = end - days * 86400 + step = max(60, math.ceil(((end - start) / 30000) / 60) * 60) + url = f"{request_data['base_url'].rstrip('/')}/api/datasources/proxy/uid/{datasource_uid}/api/v1/query_range" + try: + async with httpx.AsyncClient(timeout=90) as client: + response = await client.get( + url, + params={"query": query, "start": start, "end": end, "step": step}, + headers=request_data["headers"], cookies=request_data["cookies"], + ) + response.raise_for_status() + payload = response.json() + except (httpx.HTTPError, ValueError) as exc: + log.warning("Emby sharing history query failed: %s", type(exc).__name__) + raise HTTPException(502, "Emby session history is temporarily unavailable") + if payload.get("status") != "success": + raise HTTPException(502, "Prometheus rejected the Emby session history query") + return payload.get("data", {}).get("result", []), step + + +@app.get("/emby/account-sharing") +async def emby_account_sharing( + days: int = Query(30, ge=1, le=90), + username: str | None = Query(None, min_length=1, max_length=100), + server: Literal["all", "emby-sascha", "emby-chris"] = "all", + _=Depends(_verify), +): + """Conservative read-only analysis of concurrent networks and geographically impossible changes.""" + series, step = await _fetch_emby_session_history(days, server) + if username: + wanted = username.casefold() + series = [item for item in series if str(item.get("metric", {}).get("username", "")).casefold() == wanted] + result = _analyze_emby_sharing(series, step) + return { + "generated": datetime.now(timezone.utc).isoformat(), + "period": {"days": days, "server": server, "step_seconds": step, "series": len(series)}, + "policy": { + "mode": "conservative", + "ipv4_detection_identity": "/32", + "ipv6_display_network": "/64", + "ipv6_detection_identity": "/48", + "minimum_samples": 2, + "minimum_concurrent_seconds": 360, + "prometheus_staleness_guard": True, + "impossible_travel_minimum_km": 300, + "impossible_travel_speed_kmh": 1000, + "private_networks_excluded": True, + "automatic_enforcement": False, + }, + **result, + } + + def _get_key(cfg): vault_key = cfg.get("vault_key") if vault_key and vault_key in _vault_cache: @@ -328,18 +589,19 @@ async def ui(): async def ui_login(payload: UiLoginRequest): if not BUTLER_TOKEN or not secrets.compare_digest(payload.token, BUTLER_TOKEN): raise HTTPException(401, "Invalid token") - session_id = secrets.token_urlsafe(32) - csrf = secrets.token_urlsafe(24) - _ui_sessions[session_id] = {"csrf": csrf, "expires": time.time() + UI_SESSION_TTL} - response = JSONResponse({"authenticated": True, "expires_in": UI_SESSION_TTL}) - response.set_cookie("butler_session", session_id, max_age=UI_SESSION_TTL, httponly=True, samesite="strict", path="/") - response.set_cookie("butler_csrf", csrf, max_age=UI_SESSION_TTL, httponly=False, samesite="strict", path="/") - return response + return _issue_ui_session(read_only=False) @app.get("/ui/session") async def ui_session(request: Request): - return {"authenticated": _ui_session(request) is not None} + session = _ui_session(request) + if session: + return { + "authenticated": True, + "read_only": session.get("read_only", False), + "expires_in": max(0, int(session["expires"] - time.time())), + } + return _issue_ui_session(read_only=True) @app.post("/ui/logout") From 1ae67e158b42200a3cd7926584d7eaeeb1ef4df4 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 21:28:57 +0200 Subject: [PATCH 32/72] Add Emby account-sharing analysis and read-only UI access (ui.html) --- ui.html | 54 +++++++++++++++++++++++++++--------------------------- 1 file changed, 27 insertions(+), 27 deletions(-) diff --git a/ui.html b/ui.html index 9ad3d27..f837d9f 100644 --- a/ui.html +++ b/ui.html @@ -11,47 +11,33 @@ .shell{display:grid;grid-template-columns:248px minmax(0,1fr);min-height:100vh}.sidebar{position:sticky;top:0;height:100vh;padding:22px 16px;border-right:1px solid var(--border2);background:rgba(15,16,17,.92);backdrop-filter:blur(18px);z-index:20}.brand{display:flex;align-items:center;gap:11px;padding:0 8px 24px}.logo{width:34px;height:34px;border:1px solid rgba(113,112,255,.35);border-radius:10px;display:grid;place-items:center;background:linear-gradient(145deg,rgba(113,112,255,.22),rgba(255,255,255,.03));box-shadow:inset 0 0 16px rgba(113,112,255,.12)}.brand strong{font-size:14px;font-weight:590}.brand small{display:block;color:var(--muted);font-size:11px;margin-top:2px}.nav-label{font-size:10px;text-transform:uppercase;letter-spacing:.1em;color:var(--dim);padding:15px 10px 7px}.nav button{width:100%;display:flex;align-items:center;gap:10px;border:0;background:transparent;color:var(--muted);padding:9px 10px;border-radius:7px;text-align:left;font-size:13px;font-weight:510}.nav button:hover,.nav button.active{background:rgba(255,255,255,.05);color:var(--text)}.nav .icon{width:18px;text-align:center;color:var(--dim)}.sidebar-foot{position:absolute;left:16px;right:16px;bottom:18px}.session-pill{display:flex;align-items:center;justify-content:space-between;border:1px solid var(--border);border-radius:8px;padding:9px 10px;color:var(--muted);font-size:11px;background:rgba(255,255,255,.02)}.dot{width:7px;height:7px;border-radius:50%;background:var(--green);box-shadow:0 0 10px rgba(16,185,129,.65)} main{min-width:0}.topbar{height:68px;position:sticky;top:0;z-index:15;display:flex;align-items:center;justify-content:space-between;padding:0 30px;border-bottom:1px solid var(--border2);background:rgba(8,9,10,.78);backdrop-filter:blur(18px)}.topbar h1{font-size:15px;margin:0;font-weight:510}.top-actions{display:flex;gap:8px}.btn{border:1px solid var(--border);background:rgba(255,255,255,.03);color:var(--secondary);border-radius:7px;padding:8px 12px;font-size:12px;font-weight:510}.btn:hover{background:rgba(255,255,255,.07);color:var(--text)}.btn.primary{background:var(--accent2);border-color:transparent;color:white}.btn.danger{color:#fca5a5}.mobile-menu{display:none}.content{max-width:1440px;margin:0 auto;padding:30px}.view{display:none}.view.active{display:block}.eyebrow{font-size:11px;color:var(--accent);text-transform:uppercase;letter-spacing:.11em;font-weight:590}.hero{display:flex;justify-content:space-between;align-items:flex-end;gap:24px;margin:4px 0 26px}.hero h2{font-size:31px;letter-spacing:-.7px;font-weight:510;margin:8px 0 6px}.hero p{margin:0;color:var(--muted);font-size:14px;line-height:1.55}.updated{font:11px ui-monospace,SFMono-Regular,Menlo,monospace;color:var(--dim)} .metrics{display:grid;grid-template-columns:repeat(4,minmax(0,1fr));gap:12px;margin-bottom:18px}.metric,.panel,.operation{border:1px solid var(--border);background:rgba(255,255,255,.025);border-radius:var(--radius)}.metric{padding:17px}.metric-label{font-size:11px;color:var(--muted);margin-bottom:12px}.metric-value{font-size:25px;letter-spacing:-.45px;font-weight:510}.metric-meta{font-size:11px;color:var(--dim);margin-top:7px}.metric.good .metric-value{color:#a7f3d0}.metric.warn .metric-value{color:#fcd34d}.metric.bad .metric-value{color:#fca5a5}.grid-2{display:grid;grid-template-columns:minmax(0,1.4fr) minmax(300px,.6fr);gap:14px}.panel{padding:18px;min-width:0}.panel-head{display:flex;align-items:center;justify-content:space-between;margin-bottom:15px}.panel-title{font-size:13px;font-weight:590}.panel-sub{font-size:11px;color:var(--dim)}.empty{border:1px dashed var(--border);border-radius:8px;color:var(--dim);padding:26px;text-align:center;font-size:12px}.finding{display:grid;grid-template-columns:9px 1fr auto;gap:10px;align-items:start;padding:11px 0;border-bottom:1px solid var(--border2)}.finding:last-child{border-bottom:0}.finding-dot{width:7px;height:7px;border-radius:50%;background:var(--yellow);margin-top:5px}.finding.critical .finding-dot{background:var(--red)}.finding strong{font-size:12px}.finding p{font-size:11px;color:var(--muted);margin:4px 0 0}.badge{display:inline-flex;align-items:center;border:1px solid var(--border);border-radius:999px;padding:3px 7px;font:10px ui-monospace,SFMono-Regular,Menlo,monospace;color:var(--muted);white-space:nowrap}.badge.get,.badge.read_only{color:#a7f3d0;border-color:rgba(16,185,129,.25);background:rgba(16,185,129,.06)}.badge.post,.badge.put,.badge.patch{color:#c4b5fd;border-color:rgba(113,112,255,.3);background:rgba(113,112,255,.07)}.badge.delete,.badge.destructive{color:#fca5a5;border-color:rgba(239,68,68,.25);background:rgba(239,68,68,.06)} - .form-row{display:grid;grid-template-columns:1fr auto;gap:9px}.form-row.triple{grid-template-columns:160px 1fr auto}.input{width:100%;border:1px solid var(--border);background:rgba(255,255,255,.025);color:var(--text);border-radius:7px;padding:10px 12px;outline:0;font-size:13px}.input:focus{border-color:rgba(113,112,255,.6);box-shadow:0 0 0 3px rgba(113,112,255,.1)}select.input{appearance:none}.result{margin-top:14px;min-height:110px}.layer-grid{display:grid;grid-template-columns:repeat(auto-fit,minmax(190px,1fr));gap:9px}.layer{border:1px solid var(--border2);background:rgba(255,255,255,.02);border-radius:8px;padding:12px}.layer h4{font-size:11px;margin:0 0 8px;text-transform:uppercase;color:var(--muted);letter-spacing:.07em}.layer pre,.json{white-space:pre-wrap;word-break:break-word;margin:0;color:var(--secondary);font:11px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace}.safe-banner{display:flex;align-items:center;gap:12px;border-radius:9px;padding:14px;border:1px solid rgba(16,185,129,.25);background:rgba(16,185,129,.06)}.safe-banner.blocked{border-color:rgba(239,68,68,.25);background:rgba(239,68,68,.06)}.safe-icon{font-size:21px} + .form-row{display:grid;grid-template-columns:1fr auto;gap:9px}.form-row.triple{grid-template-columns:160px 1fr auto}.sharing-form{grid-template-columns:140px 170px minmax(180px,1fr) auto}.input{width:100%;border:1px solid var(--border);background:rgba(255,255,255,.025);color:var(--text);border-radius:7px;padding:10px 12px;outline:0;font-size:13px}.input:focus{border-color:rgba(113,112,255,.6);box-shadow:0 0 0 3px rgba(113,112,255,.1)}select.input{appearance:none}.result{margin-top:14px;min-height:110px}.layer-grid{display:grid;grid-template-columns:repeat(auto-fit,minmax(190px,1fr));gap:9px}.layer{border:1px solid var(--border2);background:rgba(255,255,255,.02);border-radius:8px;padding:12px}.layer h4{font-size:11px;margin:0 0 8px;text-transform:uppercase;color:var(--muted);letter-spacing:.07em}.layer pre,.json{white-space:pre-wrap;word-break:break-word;margin:0;color:var(--secondary);font:11px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace}.safe-banner{display:flex;align-items:center;gap:12px;border-radius:9px;padding:14px;border:1px solid rgba(16,185,129,.25);background:rgba(16,185,129,.06)}.safe-banner.blocked{border-color:rgba(239,68,68,.25);background:rgba(239,68,68,.06)}.safe-icon{font-size:21px} .toolbar{display:flex;align-items:center;gap:9px;flex-wrap:wrap;margin-bottom:15px}.toolbar .input{max-width:360px}.chips{display:flex;gap:6px;flex-wrap:wrap}.chip{border:1px solid var(--border);background:transparent;color:var(--muted);border-radius:999px;padding:6px 10px;font-size:11px}.chip.active,.chip:hover{color:var(--text);background:rgba(255,255,255,.05)}.operations{display:grid;grid-template-columns:repeat(3,minmax(0,1fr));gap:10px}.operation{padding:14px;display:flex;flex-direction:column;gap:10px;min-height:145px;transition:.16s ease}.operation:hover{border-color:rgba(113,112,255,.35);transform:translateY(-1px);background:rgba(255,255,255,.035)}.operation-top{display:flex;justify-content:space-between;gap:8px}.operation-path{font:11px/1.45 ui-monospace,SFMono-Regular,Menlo,monospace;color:var(--secondary);word-break:break-all}.operation h3{font-size:12px;font-weight:590;margin:0}.operation p{font-size:11px;color:var(--muted);line-height:1.45;margin:0;flex:1}.operation-badges{display:flex;gap:5px;flex-wrap:wrap}.table-wrap{overflow:auto;border:1px solid var(--border);border-radius:9px}table{width:100%;border-collapse:collapse;min-width:760px}th,td{text-align:left;padding:10px 12px;border-bottom:1px solid var(--border2);font-size:11px}th{color:var(--dim);text-transform:uppercase;letter-spacing:.07em;font-size:9px;background:rgba(255,255,255,.02)}td{color:var(--secondary)}td.mono{font-family:ui-monospace,SFMono-Regular,Menlo,monospace}.status-code.ok{color:#a7f3d0}.status-code.err{color:#fca5a5} .login{position:fixed;inset:0;z-index:100;background:radial-gradient(circle at 50% 15%,rgba(113,112,255,.18),transparent 30%),#08090a;display:grid;place-items:center;padding:20px}.login-card{width:min(420px,100%);border:1px solid var(--border);background:#0f1011;border-radius:14px;padding:28px;box-shadow:0 28px 90px rgba(0,0,0,.5)}.login-card .logo{margin-bottom:22px}.login-card h1{font-size:25px;letter-spacing:-.5px;font-weight:510;margin:0 0 8px}.login-card p{font-size:13px;color:var(--muted);line-height:1.5;margin:0 0 20px}.login-card form{display:grid;gap:10px}.error{color:#fca5a5;font-size:11px;min-height:16px}.toast{position:fixed;right:20px;bottom:20px;z-index:120;background:#191a1b;border:1px solid var(--border);border-radius:9px;padding:11px 14px;font-size:12px;color:var(--secondary);box-shadow:0 12px 40px rgba(0,0,0,.35);transform:translateY(20px);opacity:0;pointer-events:none;transition:.2s}.toast.show{transform:none;opacity:1}.spinner{width:15px;height:15px;border:2px solid rgba(255,255,255,.15);border-top-color:var(--accent);border-radius:50%;animation:spin .7s linear infinite;display:inline-block}@keyframes spin{to{transform:rotate(360deg)}} @media(max-width:1050px){.operations{grid-template-columns:repeat(2,minmax(0,1fr))}.metrics{grid-template-columns:repeat(2,minmax(0,1fr))}.grid-2{grid-template-columns:1fr}} - @media(max-width:760px){.shell{display:block}.sidebar{position:fixed;transform:translateX(-100%);transition:.2s;width:260px}.sidebar.open{transform:none}.mobile-menu{display:inline-flex}.topbar{padding:0 15px}.content{padding:20px 14px}.hero{align-items:flex-start;flex-direction:column}.hero h2{font-size:25px}.operations{grid-template-columns:1fr}.metrics{grid-template-columns:repeat(2,1fr)}.form-row,.form-row.triple{grid-template-columns:1fr}.btn,.input{min-height:44px}.top-actions .desktop-only{display:none}} + @media(max-width:760px){.shell{display:block}.sidebar{position:fixed;transform:translateX(-100%);transition:.2s;width:260px}.sidebar.open{transform:none}.mobile-menu{display:inline-flex}.topbar{padding:0 15px}.content{padding:20px 14px}.hero{align-items:flex-start;flex-direction:column}.hero h2{font-size:25px}.operations{grid-template-columns:1fr}.metrics{grid-template-columns:repeat(2,1fr)}.form-row,.form-row.triple,.sharing-form{grid-template-columns:1fr}.btn,.input{min-height:44px}.top-actions .desktop-only{display:none}} @media(max-width:420px){.metrics{grid-template-columns:1fr}.content{padding:18px 10px}.panel{padding:14px}.metric{padding:14px}} - - -