From 17c5d8791035549d67d4a3db97c2c638fe29ed22 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:36:17 +0200 Subject: [PATCH 1/5] Add operational safety suite --- app.py | 195 +++++++++++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 191 insertions(+), 4 deletions(-) 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({ From 0449754f614a4dd4c3d767df58576b8933993cc4 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:36:18 +0200 Subject: [PATCH 2/5] Add operational safety suite --- tests/test_app.py | 111 +++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 110 insertions(+), 1 deletion(-) 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 53222b7073ac9a8501fc2a00f5bbf9c5b2439f86 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:36:19 +0200 Subject: [PATCH 3/5] Add operational safety suite --- compose.yaml | 2 ++ 1 file changed, 2 insertions(+) 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: From 5b187450c7a52fdc33bf7a409c97421bda1e7c41 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:36:19 +0200 Subject: [PATCH 4/5] Add operational safety suite --- butler.yaml | 4 ++++ 1 file changed, 4 insertions(+) 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" From 844726a54e9427cb0ea2a8dc26f405547be9b3e1 Mon Sep 17 00:00:00 2001 From: sascha Date: Sun, 16 Aug 2026 20:36:20 +0200 Subject: [PATCH 5/5] Add operational safety suite --- .gitignore | 1 + 1 file changed, 1 insertion(+) diff --git a/.gitignore b/.gitignore index d5b4207..47bca58 100644 --- a/.gitignore +++ b/.gitignore @@ -3,3 +3,4 @@ vault-sync.log __pycache__/ *.pyc +state/