Add Butler operational safety suite #41

Merged
sascha merged 5 commits from feature/doctor-drift-preflight-audit into main 2026-08-16 20:36:23 +02:00
Showing only changes of commit 17c5d87910 - Show all commits

195
app.py
View file

@ -1,7 +1,7 @@
"""Homelab Butler v2.1 – Unified API proxy for Pfannkuchen homelab.
Reads service config from butler.yaml, credentials from Vaultwarden cache with flat-file fallback."""
import os, json, asyncio, logging, time, base64, re, subprocess, ipaddress, secrets
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({