"""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, math, hashlib from datetime import datetime, timezone 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, HTMLResponse from contextlib import asynccontextmanager from contextvars import ContextVar log = logging.getLogger("butler") VERSION = "2.3.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")) MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get("MEDIA_HANDOFF_ALLOWED_NETWORKS", "10.2.1.119/32") # --- Config loading --- _config: dict = {} def _load_config(): global _config, SERVICES, VM_CFG, TTS_CFG try: with open(CONFIG_PATH) as f: _config = yaml.safe_load(f) or {} SERVICES = _config.get("services", {}) VM_CFG = _config.get("vm", {}) TTS_CFG = _config.get("tts", {}) log.info(f"Loaded config: {len(SERVICES)} services") except FileNotFoundError: log.warning(f"No config at {CONFIG_PATH}, using defaults") SERVICES = {} VM_CFG = {} TTS_CFG = {} SERVICES: dict = {} VM_CFG: dict = {} TTS_CFG: dict = {} _load_config() # --- Audit log --- _audit_log: list[dict] = [] _audit_actor: ContextVar[str] = ContextVar("audit_actor", default="System/API") 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, actor TEXT NOT NULL DEFAULT 'Legacy/API' )""") columns = {row[1] for row in db.execute("PRAGMA table_info(audit)")} if "actor" not in columns: db.execute("ALTER TABLE audit ADD COLUMN actor TEXT NOT NULL DEFAULT 'Legacy/API'") return True except (OSError, sqlite3.Error): return False def _audit(endpoint: str, method: str, status: int, detail: str = "", dry_run: bool = False): entry = { "ts": datetime.now(timezone.utc).isoformat(), "endpoint": endpoint, "method": method, "status": status, "detail": _redact_audit_detail(detail), "dry_run": dry_run, "actor": _audit_actor.get(), } _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, actor) VALUES (?, ?, ?, ?, ?, ?, ?)", (entry["ts"], entry["endpoint"], entry["method"], entry["status"], entry["detail"], int(entry["dry_run"]), entry["actor"]), ) 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") BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "") # --- Credential cache --- _vault_cache: dict[str, str] = {} def _load_vault_cache(): """Load vault items from disk cache (written by host-side vault-sync.sh).""" global _vault_cache if not os.path.isdir(VAULT_CACHE_DIR): log.info(f"No vault cache at {VAULT_CACHE_DIR}") return new = {} for f in os.listdir(VAULT_CACHE_DIR): path = os.path.join(VAULT_CACHE_DIR, f) if os.path.isfile(path): new[f] = open(path).read().strip() _vault_cache = new log.info(f"Loaded {len(new)} vault items from cache") async def _periodic_cache_reload(): """Reload vault cache every 5 minutes (host cron writes new files).""" while True: await asyncio.sleep(300) _load_vault_cache() @asynccontextmanager async def lifespan(app: FastAPI): _load_config() _load_vault_cache() _init_audit_db() task = asyncio.create_task(_periodic_cache_reload()) yield task.cancel() app = FastAPI(title="Homelab Butler", version=VERSION, lifespan=lifespan, description="Unified API proxy + infrastructure management. AI agents: see GET / for self-onboarding.") OverviewState = Literal["healthy", "warning", "critical"] class OverviewCounts(BaseModel): critical: int warning: int healthy: int class OverviewFinding(BaseModel): severity: OverviewState code: str target: str message: str age_hours: float | None = None pct: int | None = None class OverviewModelContract(BaseModel): instruction: str severity_order: list[OverviewState] class OverviewResponse(BaseModel): schema_version: Literal[1] generated: datetime overall_state: OverviewState action_required: bool summary: OverviewCounts components: dict[str, OverviewCounts] findings: list[OverviewFinding] model_contract: OverviewModelContract details: dict | None = None # --- Credential reading (vault-first, file-fallback) --- def _read(name): """Read credential: vault cache first, then flat file.""" # Vault cache uses lowercase-hyphenated names vault_name = name.lower().replace("_", "-") if vault_name in _vault_cache: return _vault_cache[vault_name] # Try uppercase convention upper = name.upper().replace("-", "_").lower().replace("_", "-") if upper in _vault_cache: return _vault_cache[upper] # Fallback to flat file try: return open(f"{API_DIR}/{name}").read().strip() except FileNotFoundError: return None def _parse_kv(name): raw = _read(name) if not raw: return {} d = {} for line in raw.splitlines(): if ":" in line: k, v = line.split(":", 1) d[k.strip().lower()] = v.strip() return d def _parse_url_key(name): raw = _read(name) if not raw: return None, None lines = [l.strip() for l in raw.splitlines() if l.strip()] return (lines[0] if lines else None, lines[1] if len(lines) > 1 else None) # --- Service configs loaded from butler.yaml --- # --- Dockhand session --- _dockhand_cookie = None async def _dockhand_login(client): global _dockhand_cookie r = await client.post( f"{SERVICES['dockhand']['url']}/api/auth/login", json={"username": "admin", "password": _read("dockhand") or ""}, ) if r.status_code == 200: _dockhand_cookie = dict(r.cookies) return _dockhand_cookie # --- 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 _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 auth = request.headers.get("authorization", "") 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"}: 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 _clean_audit_actor(value: str) -> str: cleaned = re.sub(r"[^\w .@/\-]", "", str(value or ""))[:40].strip() return cleaned or "KI/API" @app.middleware("http") async def audit_actor_context(request: Request, call_next): if request.headers.get("authorization", "").startswith("Bearer "): actor = _clean_audit_actor(request.headers.get("x-butler-actor", "KI/API")) elif _ui_session(request): actor = "Weboberfläche" else: actor = "System/Öffentlich" token = _audit_actor.set(actor) try: return await call_next(request) finally: _audit_actor.reset(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, } def _emby_history_step(days: int) -> int: """Keep query_range below Prometheus' 11,000-points-per-series limit.""" duration_seconds = days * 86400 return max(60, math.ceil((duration_seconds / 10_500) / 60) * 60) 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 = _emby_history_step(days) url = f"{request_data['base_url'].rstrip('/')}/api/datasources/proxy/uid/{datasource_uid}/api/v1/query_range" series_by_metric: dict[str, dict] = {} chunk_seconds = 30 * 86400 try: async with httpx.AsyncClient(timeout=90) as client: chunk_start = start while chunk_start < end: chunk_end = min(chunk_start + chunk_seconds, end) response = await client.get( url, params={"query": query, "start": chunk_start, "end": chunk_end, "step": step}, headers=request_data["headers"], cookies=request_data["cookies"], ) response.raise_for_status() payload = response.json() if payload.get("status") != "success": raise HTTPException(502, "Prometheus rejected the Emby session history query") for item in payload.get("data", {}).get("result", []): metric = item.get("metric", {}) key = json.dumps(metric, sort_keys=True, separators=(",", ":")) merged = series_by_metric.setdefault(key, {"metric": metric, "values": []}) merged["values"].extend(item.get("values", [])) chunk_start = chunk_end 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") series = [] for item in series_by_metric.values(): values_by_timestamp = { float(value[0]): value for value in item["values"] if isinstance(value, (list, tuple)) and len(value) >= 2 } item["values"] = [values_by_timestamp[ts] for ts in sorted(values_by_timestamp)] series.append(item) return series, 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: return _vault_cache[vault_key] return _read(cfg.get("key_file", "")) SENSITIVE_RESPONSE_FIELDS = { "hawsertoken", "webhooksecret", "accesstoken", "refreshtoken", "password", "secret", "apikey", "api_key", "privatekey", } def _redact_response(value, extra_fields=None): """Recursively redact known secret fields in proxied JSON responses.""" sensitive = set(SENSITIVE_RESPONSE_FIELDS) sensitive.update(str(x).lower() for x in (extra_fields or [])) if isinstance(value, dict): return { key: "[REDACTED]" if str(key).lower() in sensitive else _redact_response(item, sensitive) for key, item in value.items() } if isinstance(value, list): return [_redact_response(item, sensitive) for item in value] return value def _inventory_hosts(text: str) -> list[dict]: """Parse Ansible inventory host lines and apply Pfannkuchen SSH defaults.""" hosts = [] seen = set() for raw in text.splitlines(): line = raw.strip() if not line or line.startswith(("#", "[")) or "ansible_host=" not in line: continue parts = line.split() name = parts[0] attrs = {k: v for k, v in (p.split("=", 1) for p in parts[1:] if "=" in p)} ip = attrs.get("ansible_host") if not ip or name in seen: continue user = attrs.get("ansible_user") if not user: if name.startswith("node"): user = "root" elif ip.startswith("10.7.1."): user = "chris" else: user = "sascha" hosts.append({"name": name, "ip": ip, "user": user}) seen.add(name) return hosts # --- 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") return _issue_ui_session(read_only=False) @app.get("/ui/session") async def ui_session(request: Request): 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") 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.""" svc_list = {} for name, cfg in SERVICES.items(): svc_list[name] = {"url": cfg.get("url", ""), "auth": cfg.get("auth", ""), "description": cfg.get("description", "")} return { "service": "homelab-butler", "version": VERSION, "docs": "/docs", "openapi": "/openapi.json", "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?}", "vm_status": "GET /vm/status/{vmid}", "vm_delete": "DELETE /vm/{vmid} - simple Proxmox delete (legacy)", "vm_destroy": "DELETE /vm/destroy/{vmid}?dry_run=false - complete cleanup (VM, Dockhand, Repo, Ansible, Kuma)", "inventory_add": "POST /inventory/host {name, ip, group?}", "ansible_run": "POST /ansible/run {hostname}", "tts_speak": "POST /tts/speak {text, target: speaker|telegram}", "tts_generate": "POST /tts/generate {text, voice?, language?} - return cloned WAV audio", "tts_voices": "GET /tts/voices", "tts_health": "GET /tts/health", "status": "GET /status - health of all backends", "overview": "GET /overview?details=false - deterministic homelab verdict for small models", "audit": "GET /audit - recent API calls", "sysctl_audit": "GET /system/sysctl/{host} - read-only live and persistent network tuning", }, "vault_items": len(_vault_cache), } @app.get("/health") 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", "/network/media-tunnel", "/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: return "healthy" if status_code in (401, 403): return "auth_failed" if status_code == 404: return "misconfigured" return "degraded" def _service_auth(cfg: dict) -> dict: """Build secret-bearing request data without ever returning it from an endpoint.""" auth_type = cfg.get("auth", "none") headers = {} cookies = {} base_url = cfg.get("url") if auth_type == "apikey": headers["X-Api-Key"] = _get_key(cfg) or "" elif auth_type == "apikey_urlfile": base_url, key = _parse_url_key(cfg.get("key_file", "")) headers["X-Api-Key"] = key or "" elif auth_type == "bearer": headers["Authorization"] = f"Bearer {_get_key(cfg) or ''}" elif auth_type == "n8n": headers["X-N8N-API-KEY"] = _get_key(cfg) or "" elif auth_type == "proxmox": pv = _parse_kv("proxmox") headers["Authorization"] = f"PVEAPIToken={pv.get('tokenid', '')}={pv.get('secret', '')}" return {"base_url": base_url, "headers": headers, "cookies": cookies} async def _collect_service_status() -> dict: """Run authenticated functional probes concurrently and classify their result.""" started = time.monotonic() async with httpx.AsyncClient(verify=False, timeout=5, follow_redirects=True) as client: async def probe(name: str, cfg: dict): probe_started = time.monotonic() try: request_data = _service_auth(cfg) base_url = request_data["base_url"] if not base_url: return name, { "reachable": False, "status": "misconfigured", "message": "No service URL configured" } if cfg.get("auth") == "session": request_data["cookies"] = await _dockhand_login(client) or {} health_path = cfg.get("health_path", "") target = f"{base_url.rstrip('/')}/{health_path.lstrip('/')}" if health_path else base_url expected = {int(code) for code in cfg.get("health_expected", range(200, 400))} response = await client.get( target, headers=request_data["headers"], cookies=request_data["cookies"], ) state = _classify_http_status(response.status_code, expected) messages = { "healthy": "Functional probe succeeded", "auth_failed": "Configured credentials were rejected", "misconfigured": "Configured health route was not found", "degraded": "Backend returned an unexpected HTTP status", } return name, { "reachable": True, "status": state, "http": response.status_code, "latency_ms": round((time.monotonic() - probe_started) * 1000), "message": messages[state], } except Exception as exc: return name, { "reachable": False, "status": "offline", "latency_ms": round((time.monotonic() - probe_started) * 1000), "error": type(exc).__name__, "message": "Backend could not be reached", } pairs = await asyncio.gather(*(probe(name, cfg) for name, cfg in SERVICES.items())) results = dict(pairs) results["_meta"] = {"duration_ms": round((time.monotonic() - started) * 1000)} return results @app.get("/status") async def status(_=Depends(_verify)): """Authenticated functional health check for all configured backends.""" results = await _collect_service_status() _audit("/status", "GET", 200) return results @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, actor 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") async def config_reload(_=Depends(_verify)): """Reload butler.yaml and vault cache.""" _load_config() _load_vault_cache() return {"config_services": len(SERVICES), "vault_items": len(_vault_cache)} @app.get("/info") async def info(_=Depends(_verify)): """Secret-free machine-readable context for AI agents and operators.""" return { "service": "homelab-butler", "version": VERSION, "generated": datetime.now(timezone.utc).isoformat(), "services": { name: { "url": cfg.get("url"), "auth": cfg.get("auth"), "description": cfg.get("description", ""), } for name, cfg in SERVICES.items() }, "endpoints": { "capabilities": "/capabilities", "ui": "/ui", "doctor": "/doctor/{target}", "drift": "/drift", "maintenance_preflight": "/maintenance/preflight", "status": "/status", "overview": "/overview?details=false", "audit": "/audit", "host_health": "/health/all", "backups": "/backup/status", "disk": "/disk/usage", "logs": "/logs/{host}/{container}?tail=200", "inspect": "/docker/inspect/{host}/{container}", "docs": "/docs", }, "rules": [ "Backend services are accessed through Butler or dedicated MCP servers", "VMs only; no LXC", "Docker Compose is stored in Git; no docker run", "Persistent volumes live under /app-config", "Node 7 VM SSH user is chris", ], } def _get_inventory_hosts() -> list[dict]: rc, out, _err = _ssh( AUTOMATION1, "python3 -c \"print(open('/app-config/ansible/pfannkuchen.ini').read())\"", timeout=15, ) return _inventory_hosts(out) if rc == 0 else [] def _find_inventory_host(name: str) -> dict | None: return next((host for host in _get_inventory_hosts() if host["name"] == name), None) async def _get_inventory_hosts_async() -> list[dict]: return await asyncio.to_thread(_get_inventory_hosts) async def _collect_health_all(concurrency: int = 10) -> dict: """Collect SSH and container health concurrently with bounded fan-out.""" semaphore = asyncio.Semaphore(concurrency) async def inspect_host(host: dict): async with semaphore: rc, out, err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', "echo __BUTLER_OK__; (sudo -n docker ps --format '{{.Names}}: {{.Status}}' 2>/dev/null || docker ps --format '{{.Names}}: {{.Status}}' 2>/dev/null) | head -30", 10, ) lines = out.strip().splitlines() return host["name"], { "ip": host["ip"], "user": host["user"], "reachable": rc == 0 and bool(lines) and lines[0] == "__BUTLER_OK__", "containers": lines[1:] if lines and lines[0] == "__BUTLER_OK__" else [], "error": err.strip()[:200] if rc != 0 else None, } hosts = await _get_inventory_hosts_async() pairs = await asyncio.gather(*(inspect_host(host) for host in hosts)) return dict(pairs) @app.get("/health/all") async def health_all(_=Depends(_verify)): """SSH reachability and Docker status for all inventory hosts.""" return await _collect_health_all() def _parse_backup_time(value: str | None) -> datetime | None: if not value: return None try: parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) return parsed.replace(tzinfo=timezone.utc) if parsed.tzinfo is None else parsed.astimezone(timezone.utc) except (TypeError, ValueError): return None def _backup_item(rc: int, out: str, err: str, now: datetime | None = None) -> dict: """Normalize borgmatic output into one small-model-friendly state object.""" item = {"state": "unknown", "ok": False, "last_backup": None, "age_hours": None} if rc != 0 or not out.strip(): if err: item["error"] = err.strip()[:200] return item try: data = json.loads(out) archives = data[0].get("archives", []) if isinstance(data, list) and data else [] if not archives: item["error"] = "no archives returned" return item last = archives[-1] started_at = _parse_backup_time(last.get("start")) if not started_at: item["error"] = "invalid backup timestamp" return item age_hours = max(0, ((now or datetime.now(timezone.utc)) - started_at).total_seconds() / 3600) state = "healthy" if age_hours <= 30 else "warning" if age_hours <= 48 else "critical" return { "state": state, "ok": state == "healthy", "last_backup": last.get("start"), "age_hours": round(age_hours, 1), "name": last.get("name"), } except (json.JSONDecodeError, TypeError, IndexError, KeyError): item["error"] = "invalid borgmatic JSON" return item 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, f'{host["user"]}@{host["ip"]}', "sudo -n borgmatic list --last 1 --json", 60, ) return host["name"], _backup_item(rc, out, err) 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, "exempt": 0} for item in results.values(): summary[item["state"]] += 1 return {"summary": summary, "hosts": results} @app.get("/backup/status") async def backup_status(_=Depends(_verify)): """Latest Borgmatic archive, age and severity for all VM inventory hosts.""" return await _collect_backup_status() @app.get("/backup/diagnose/{hostname}") async def backup_diagnose(hostname: str, _=Depends(_verify)): """Safe Borgmatic timer/service diagnostics without exposing backup secrets.""" host = await asyncio.to_thread(_find_inventory_host, hostname) if not host or host["name"].startswith("node"): return JSONResponse({"error": "inventory host not found"}, status_code=404) command = """echo __TIMER_ENABLED__ sudo -n systemctl is-enabled borg-backup.timer 2>&1 || true echo __TIMER_ACTIVE__ sudo -n systemctl is-active borg-backup.timer 2>&1 || true echo __TIMER_SCHEDULE__ sudo -n systemctl list-timers borg-backup.timer --all --no-pager 2>&1 || true echo __SERVICE_STATE__ sudo -n systemctl show borg-backup.service -p ActiveState -p SubState -p Result -p ExecMainStatus --no-pager 2>&1 || true echo __RECENT_LOGS__ sudo -n journalctl -u borg-backup.service --since '2 days ago' -n 120 --no-pager 2>&1 || true echo __LATEST_ARCHIVE__ timeout 50s sudo -n borgmatic list --last 1 --json 2>&1 || true """ rc, out, err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', command, 75 ) return { "hostname": host["name"], "ip": host["ip"], "user": host["user"], "rc": rc, "output": out[-30000:], "stderr": err[-2000:], } @app.post("/backup/break-lock/{hostname}") async def backup_break_lock(hostname: str, _=Depends(_verify)): """Break a stale Borg repository lock only while the backup service is inactive.""" host = await asyncio.to_thread(_find_inventory_host, hostname) if not host or host["name"].startswith("node"): return JSONResponse({"error": "inventory host not found"}, status_code=404) rc_state, state, state_err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', "sudo -n systemctl is-active borg-backup.service 2>&1 || true", 15 ) if state.strip() in ("active", "activating"): return JSONResponse( {"error": "backup service is active; refusing to break lock", "state": state.strip()}, status_code=409, ) rc, out, err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', "sudo -n borgmatic borg break-lock", 120 ) _audit(f"/backup/break-lock/{hostname}", "POST", 200 if rc == 0 else 500, f"rc={rc}") if rc != 0: return JSONResponse( {"error": "borg break-lock failed", "hostname": hostname, "stdout": out[-4000:], "stderr": err[-4000:]}, status_code=500, ) return {"status": "lock_broken", "hostname": hostname, "detail": out.strip()[-4000:]} @app.get("/backup/history/{hostname}") async def backup_history(hostname: str, _=Depends(_verify)): """Return the ten newest Borg archives for one inventory host.""" host = await asyncio.to_thread(_find_inventory_host, hostname) if not host or host["name"].startswith("node"): return JSONResponse({"error": "inventory host not found"}, status_code=404) rc, out, err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', "sudo -n borgmatic list --last 10 --json", 120 ) if rc != 0: return JSONResponse( {"error": "borg list failed", "hostname": hostname, "stdout": out[-4000:], "stderr": err[-4000:]}, status_code=500 ) try: history = json.loads(out) except (TypeError, ValueError): return JSONResponse({"error": "invalid borg JSON", "hostname": hostname}, status_code=500) return {"hostname": hostname, "history": history} @app.get("/backup/run-status/{hostname}") async def backup_run_status(hostname: str, _=Depends(_verify)): """Lightweight current backup service state and recent logs.""" host = await asyncio.to_thread(_find_inventory_host, hostname) if not host or host["name"].startswith("node"): return JSONResponse({"error": "inventory host not found"}, status_code=404) command = """echo __STATE__ sudo -n systemctl show borg-backup.service -p ActiveState -p SubState -p Result -p ExecMainStatus --no-pager 2>&1 echo __LOGS__ sudo -n journalctl -u borg-backup.service -n 30 --no-pager 2>&1 """ rc, out, err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', command, 20 ) return {"hostname": hostname, "rc": rc, "output": out[-12000:], "stderr": err[-2000:]} @app.post("/backup/run/{hostname}") async def backup_run(hostname: str, _=Depends(_verify)): """Start one inventory host Borgmatic service asynchronously.""" host = await asyncio.to_thread(_find_inventory_host, hostname) if not host or host["name"].startswith("node"): return JSONResponse({"error": "inventory host not found"}, status_code=404) command = ( "sudo -n systemctl reset-failed borg-backup.service 2>/dev/null || true; " "sudo -n systemctl start --no-block borg-backup.service; sleep 2; " "sudo -n systemctl show borg-backup.service -p ActiveState -p SubState -p Result --no-pager" ) rc, out, err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', command, 20 ) _audit(f"/backup/run/{hostname}", "POST", 202 if rc == 0 else 500, f"rc={rc}") if rc != 0: return JSONResponse( {"error": "backup start failed", "hostname": hostname, "stdout": out[-2000:], "stderr": err[-2000:]}, status_code=500, ) return JSONResponse( {"status": "started", "hostname": hostname, "detail": out.strip()}, status_code=202, ) async def _collect_disk_usage(concurrency: int = 10) -> dict: semaphore = asyncio.Semaphore(concurrency) async def inspect_disk(host: dict): async with semaphore: rc, out, _err = await asyncio.to_thread( _ssh, f'{host["user"]}@{host["ip"]}', "df -P / | tail -1", 10 ) parts = out.split() if rc != 0 or len(parts) < 6: return host["name"], None return host["name"], { "size_kib": int(parts[1]), "used_kib": int(parts[2]), "avail_kib": int(parts[3]), "pct": parts[4], "mount": parts[5], } hosts = await _get_inventory_hosts_async() pairs = await asyncio.gather(*(inspect_disk(host) for host in hosts)) return {name: item for name, item in pairs if item is not None} @app.get("/disk/usage") async def disk_usage(_=Depends(_verify)): """Root filesystem usage for all reachable inventory hosts.""" return await _collect_disk_usage() def _add_component(summary: dict, bucket: dict, state: str): normalized = state if state in ("healthy", "warning", "critical") else "warning" summary[normalized] += 1 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.""" services, hosts, backups, disks = await asyncio.gather( _collect_service_status(), _collect_health_all(), _collect_backup_status(), _collect_disk_usage(), ) summary = {"critical": 0, "warning": 0, "healthy": 0} component_summary = { name: {"critical": 0, "warning": 0, "healthy": 0} for name in ("services", "hosts", "backups", "disks") } findings = [] for name in sorted(hosts): item = hosts[name] containers = item.get("containers", []) bad_container = next( (line for line in containers if "unhealthy" in line.lower() or "restarting" in line.lower()), None, ) if not item.get("reachable"): state = "critical" findings.append({ "severity": "critical", "code": "host_unreachable", "target": name, "message": "Host is not reachable over SSH", }) elif bad_container: state = "critical" findings.append({ "severity": "critical", "code": "container_unhealthy", "target": name, "message": bad_container[:200], }) else: state = "healthy" _add_component(summary, component_summary["hosts"], state) service_states = { "healthy": "healthy", "degraded": "warning", "offline": "critical", "auth_failed": "critical", "misconfigured": "critical", } service_codes = { "degraded": "service_degraded", "offline": "service_offline", "auth_failed": "service_auth_failed", "misconfigured": "service_misconfigured", } for name in sorted(key for key in services if not key.startswith("_")): item = services[name] raw_state = item.get("status", "degraded") state = service_states.get(raw_state, "warning") _add_component(summary, component_summary["services"], state) if state != "healthy": findings.append({ "severity": state, "code": service_codes.get(raw_state, "service_degraded"), "target": name, "message": item.get("message", "Service health probe failed"), }) for name in sorted(backups.get("hosts", {})): item = backups["hosts"][name] raw_state = item.get("state", "unknown") 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({ "severity": state, "code": f"backup_{raw_state}", "target": name, "message": "Backup is missing, stale or could not be verified", "age_hours": item.get("age_hours"), }) for name in sorted(disks): pct = int(str(disks[name].get("pct", "0")).rstrip("%") or 0) state = "critical" if pct >= 90 else "warning" if pct >= 80 else "healthy" _add_component(summary, component_summary["disks"], state) if state != "healthy": findings.append({ "severity": state, "code": "disk_critical" if state == "critical" else "disk_high", "target": name, "message": f"Root filesystem usage is {pct}%", "pct": pct, }) overall_state = "critical" if summary["critical"] else "warning" if summary["warning"] else "healthy" response = { "schema_version": 1, "generated": datetime.now(timezone.utc).isoformat(), "overall_state": overall_state, "action_required": overall_state != "healthy", "summary": summary, "components": component_summary, "findings": findings, "model_contract": { "instruction": "Report overall_state, then findings in the returned order. Do not infer missing facts.", "severity_order": ["critical", "warning", "healthy"], }, } if details: response["details"] = {"services": services, "hosts": hosts, "backups": backups, "disks": disks} _audit("/overview", "GET", 200, f"state={overall_state} findings={len(findings)}") return response @app.get("/logs/{host}/{container}") async def docker_logs(host: str, container: str, tail: int = Query(200, ge=1, le=20000), _=Depends(_verify)): """Read Docker logs from an inventory host.""" if not re.fullmatch(r"[A-Za-z0-9_.-]+", host) or not re.fullmatch(r"[A-Za-z0-9_.-]+", container): raise HTTPException(400, "Invalid host or container name") target = _find_inventory_host(host) if not target: raise HTTPException(404, f"Host {host} not found") rc, out, err = _ssh( f'{target["user"]}@{target["ip"]}', f"sudo -n docker logs {container} --tail {tail} 2>&1 || docker logs {container} --tail {tail} 2>&1", timeout=30, ) if rc != 0: raise HTTPException(502, (err or out).strip()[:500]) return {"host": host, "container": container, "tail": tail, "output": out} @app.get("/docker/inspect/{host}/{container}") async def docker_inspect(host: str, container: str, _=Depends(_verify)): """Return a sanitized runtime/resource summary for a Docker container.""" if not re.fullmatch(r"[A-Za-z0-9_.-]+", host) or not re.fullmatch(r"[A-Za-z0-9_.-]+", container): raise HTTPException(400, "Invalid host or container name") target = _find_inventory_host(host) if not target: raise HTTPException(404, f"Host {host} not found") rc, out, err = _ssh( f'{target["user"]}@{target["ip"]}', f"sudo -n docker inspect {container}", timeout=20, ) if rc != 0: raise HTTPException(502, (err or out).strip()[:500]) try: raw = json.loads(out)[0] except (json.JSONDecodeError, IndexError, TypeError): raise HTTPException(502, "Invalid docker inspect response") host_cfg = raw.get("HostConfig", {}) cfg = raw.get("Config", {}) state = raw.get("State", {}) return { "name": raw.get("Name", "").lstrip("/"), "image": cfg.get("Image"), "state": { "status": state.get("Status"), "running": state.get("Running"), "started_at": state.get("StartedAt"), "exit_code": state.get("ExitCode"), "oom_killed": state.get("OOMKilled"), "restart_count": raw.get("RestartCount"), }, "runtime": host_cfg.get("Runtime"), "resources": { "memory": host_cfg.get("Memory"), "memory_reservation": host_cfg.get("MemoryReservation"), "nano_cpus": host_cfg.get("NanoCpus"), "device_requests": host_cfg.get("DeviceRequests"), }, "restart_policy": host_cfg.get("RestartPolicy"), "log_config": host_cfg.get("LogConfig"), "environment_keys": sorted(item.split("=", 1)[0] for item in cfg.get("Env", []) if "=" in item), "mounts": [ {"type": mount.get("Type"), "source": mount.get("Source"), "destination": mount.get("Destination"), "rw": mount.get("RW")} for mount in raw.get("Mounts", []) ], } @app.post("/docker/restart/{host}/{container}") async def docker_restart(host: str, container: str, _=Depends(_verify), dry_run: bool = Query(False)): """Restart a named Docker container, with optional dry-run.""" if not re.fullmatch(r"[A-Za-z0-9_.-]+", host) or not re.fullmatch(r"[A-Za-z0-9_.-]+", container): raise HTTPException(400, "Invalid host or container name") target = _find_inventory_host(host) if not target: raise HTTPException(404, f"Host {host} not found") if dry_run: return {"dry_run": True, "host": host, "container": container} rc, out, err = _ssh(f'{target["user"]}@{target["ip"]}', f"sudo -n docker restart {container}", timeout=45) _audit(f"/docker/restart/{host}/{container}", "POST", 200 if rc == 0 else 502) if rc != 0: raise HTTPException(502, (err or out).strip()[:500]) return {"success": True, "output": out.strip()} @app.post("/vault/reload") async def vault_reload(_=Depends(_verify)): _load_vault_cache() return {"reloaded": True, "items": len(_vault_cache)} # --- VPS reverse-proxy and DNS management --- VPS_SSH = "root@46.225.230.72" VPS_IPV4 = "46.225.230.72" VPS_IPV6 = "2a01:4f8:1c19:9653::1" MANAGED_DNS_ZONES = {"guck.tv"} def _validate_proxy_route(domain: str, upstream: str) -> tuple[str, str, str, str]: domain = domain.strip().lower().rstrip(".") upstream = upstream.strip().lower() 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])?)+", domain): raise ValueError("invalid domain") zone = next((item for item in MANAGED_DNS_ZONES if domain.endswith(f".{item}")), None) if not zone or domain == zone: raise ValueError("domain is outside managed DNS zones or is a zone apex") match = re.fullmatch(r"(127\.0\.0\.1|localhost):(\d{1,5})", upstream) if not match or not 1 <= int(match.group(2)) <= 65535: raise ValueError("upstream must be localhost with a valid TCP port") return domain, upstream, zone, domain[: -(len(zone) + 1)] class ProxyRouteRequest(BaseModel): domain: str upstream: str dns_token: str | None = None def _remote_python(script: str, timeout: int = 30) -> tuple[int, str, str]: encoded = base64.b64encode(script.encode()).decode() return _ssh(VPS_SSH, f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"', timeout=timeout) def _restore_caddy_backup(backup: str): if not re.fullmatch(r"/app-config/caddy/Caddyfile\.bak-\d{8}T\d{6}Z", backup): raise ValueError("invalid Caddy backup path") script = f"""from pathlib import Path Path('/app-config/caddy/Caddyfile').write_bytes(Path({backup!r}).read_bytes()) """ rc, _out, err = _remote_python(script) if rc != 0: raise RuntimeError(f"Caddy rollback write failed: {err[-300:]}") rc, _out, err = _ssh(VPS_SSH, "docker exec caddy caddy reload --config /etc/caddy/Caddyfile", timeout=30) if rc != 0: raise RuntimeError(f"Caddy rollback reload failed: {err[-300:]}") def _configure_caddy_route(domain: str, upstream: str) -> dict: timestamp = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") backup = f"/app-config/caddy/Caddyfile.bak-{timestamp}" block = f"{domain} {{\n reverse_proxy {upstream}\n}}\n\n" script = f"""from pathlib import Path import re, shutil path = Path('/app-config/caddy/Caddyfile') backup = Path({backup!r}) content = path.read_text() block = {block!r} pattern = re.compile(r'(?ms)^{re.escape(domain)}\\s*\\{{.*?^\\}}\\s*') shutil.copy2(path, backup) if pattern.search(content): content = pattern.sub(block, content, count=1) else: if content and not content.endswith('\\n'): content += '\\n' content += '\\n' + block path.write_text(content) print(backup) """ rc, out, err = _remote_python(script) if rc != 0: raise RuntimeError(f"Caddyfile update failed: {err[-300:]}") backup = out.strip() or backup rc, _out, err = _ssh(VPS_SSH, "docker exec caddy caddy validate --config /etc/caddy/Caddyfile", timeout=30) if rc != 0: _restore_caddy_backup(backup) raise RuntimeError(f"Caddy validation failed: {err[-300:]}") rc, _out, err = _ssh(VPS_SSH, "docker exec caddy caddy reload --config /etc/caddy/Caddyfile", timeout=30) if rc != 0: _restore_caddy_backup(backup) raise RuntimeError(f"Caddy reload failed: {err[-300:]}") 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()) broad = re.findall(r"(? str: token = _normalize_hetzner_dns_token(_read("HETZNER_DNS_TOKEN") or "") if token: return token 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", timeout=120, ) if rc != 0: raise RuntimeError(f"Vault cache sync failed: {err[-300:]}") _load_vault_cache() token = _read("HETZNER_DNS_TOKEN") if not token: raise RuntimeError("HETZNER_DNS_TOKEN is unavailable after vault sync") return token async def _upsert_dns_records(zone: str, name: str, token_override: str | None = None) -> dict: token = token_override or await asyncio.to_thread(_get_hetzner_dns_token) headers = {"Authorization": f"Bearer {token}", "Content-Type": "application/json"} api = "https://api.hetzner.cloud/v1" async with httpx.AsyncClient(timeout=30) as client: zones_response = await client.get(f"{api}/zones", headers=headers) zones_response.raise_for_status() zone_data = next((item for item in zones_response.json().get("zones", []) if item.get("name") == zone), None) if not zone_data: raise RuntimeError(f"DNS zone not found: {zone}") zone_id = zone_data["id"] rrsets_response = await client.get(f"{api}/zones/{zone_id}/rrsets", headers=headers) rrsets_response.raise_for_status() existing = {(item.get("name"), item.get("type")) for item in rrsets_response.json().get("rrsets", [])} for record_type, value in (("A", VPS_IPV4), ("AAAA", VPS_IPV6)): payload = {"name": name, "type": record_type, "ttl": 300, "records": [{"value": value, "comment": "Managed by Homelab Butler"}]} if (name, record_type) in existing: response = await client.put(f"{api}/zones/{zone_id}/rrsets/{name}/{record_type}", headers=headers, json=payload) else: response = await client.post(f"{api}/zones/{zone_id}/rrsets", headers=headers, json=payload) response.raise_for_status() return {"zone_id": zone_id, "records": ["A", "AAAA"]} @app.post("/vps/proxy-route") async def vps_proxy_route(req: ProxyRouteRequest, _=Depends(_verify)): try: domain, upstream, zone, name = _validate_proxy_route(req.domain, req.upstream) except ValueError as exc: raise HTTPException(400, str(exc)) from exc caddy = await asyncio.to_thread(_configure_caddy_route, domain, upstream) try: dns = await _upsert_dns_records(zone, name, req.dns_token) except Exception: await asyncio.to_thread(_restore_caddy_backup, caddy["backup"]) raise _audit("/vps/proxy-route", "POST", 200, f"{domain} -> {upstream}") return {"status": "configured", "domain": domain, "upstream": upstream, "caddy": caddy, "dns": dns} SPEEDTEST_REPO_FILES = ( ".dockerignore", "Dockerfile", "compose.yaml", "pyproject.toml", "streamscope/__init__.py", "streamscope/app.py", "streamscope/db.py", "streamscope/mtr.py", "streamscope/scoring.py", "streamscope/static/index.html", "streamscope/static/assets/app.css", "streamscope/static/assets/app.js", "streamscope/static/assets/longterm-metrics.js", ) GUCK_ADMIN_REPO_FILES = ( "guck-admin/Dockerfile", "guck-admin/compose.yaml", "guck-admin/requirements.txt", "guck-admin/src/app.py", "guck-admin/src/control.py", "guck-admin/src/sharing_watchdog.py", "guck-admin/src/templates/bandwidth.html", "guck-admin/src/templates/base.html", "guck-admin/src/templates/dashboard.html", "guck-admin/src/templates/history.html", "guck-admin/src/templates/sessions.html", "guck-admin/src/templates/settings.html", "guck-admin/src/templates/sharing.html", "guck-admin/src/templates/users.html", "guck-admin/static/admin.css", "guck-admin/static/admin.js", "guck-admin/static/icon.svg", "guck-admin/static/manifest.webmanifest", "guck-admin/static/sw.js", "guck-admin/static/world.svg", ) FORGEJO_DEPLOY_FILES = { "sascha/speedtest": frozenset(SPEEDTEST_REPO_FILES), "sascha/guck-vps": frozenset(GUCK_ADMIN_REPO_FILES), } class SpeedtestDeployRequest(BaseModel): stats_password: str session_secret: str async def _fetch_forgejo_text(repo: str, path: str) -> str: if path not in FORGEJO_DEPLOY_FILES.get(repo, frozenset()): raise ValueError("unsupported Forgejo file") cfg = SERVICES.get("forgejo", {}) base_url = cfg.get("url") token = _get_key(cfg) if not base_url or not token: raise RuntimeError("Forgejo service configuration is unavailable") url = f"{base_url}/api/v1/repos/{repo}/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_speedtest_compose(files: dict[str, str], password: str, session_secret: str) -> dict: if set(files) != set(SPEEDTEST_REPO_FILES): raise ValueError("speedtest source bundle is incomplete") compose = files["compose.yaml"] dockerfile = files["Dockerfile"] required = [ "build: .", '127.0.0.1:8080:8080', '/app-config/speedtest/data:/data', 'ADMIN_PASSWORD: "${ADMIN_PASSWORD:', 'SESSION_SECRET: "${SESSION_SECRET:', "NET_RAW", ] if any(item not in compose for item in required): raise ValueError("StreamScope compose is missing a required security or persistence setting") if "python:" not in dockerfile or "mtr-tiny" not in dockerfile or "php" in dockerfile.lower(): raise ValueError("StreamScope image must be Python-based, MTR-capable and PHP-free") secret_pattern = r"[A-Za-z0-9!@#%_+=:,.?-]{24,128}" if not re.fullmatch(secret_pattern, password): raise ValueError("stats password must be 24-128 safe characters") if not re.fullmatch(secret_pattern, session_secret): raise ValueError("session secret must be 24-128 safe characters") script = f"""from pathlib import Path import os, shutil stack = Path('/app-config/github/speedtest') backup = Path('/app-config/deployment-backups/speedtest-rollback') data = Path('/app-config/speedtest/data') if backup.exists(): shutil.rmtree(backup) if stack.exists(): backup.parent.mkdir(parents=True, exist_ok=True) shutil.copytree(stack, backup) stack.mkdir(parents=True, exist_ok=True) data.mkdir(parents=True, exist_ok=True) files = {files!r} for relative, content in files.items(): target = stack / relative target.parent.mkdir(parents=True, exist_ok=True) target.write_text(content) env = stack / '.env' env.write_text('ADMIN_PASSWORD=' + {password!r} + '\\nSESSION_SECRET=' + {session_secret!r} + '\\nSTATS_PASSWORD=' + {password!r} + '\\n') os.chmod(env, 0o600) """ rc, _out, err = _remote_python(script) if rc != 0: raise RuntimeError(f"StreamScope file deployment failed: {err[-300:]}") rollback = "rm -rf /app-config/github/speedtest && cp -a /app-config/deployment-backups/speedtest-rollback /app-config/github/speedtest && cd /app-config/github/speedtest && docker compose up -d" preflight = "cd /app-config/github/speedtest && docker compose config -q && docker compose build --pull" rc, _out, err = _ssh(VPS_SSH, preflight, timeout=600) if rc != 0: _ssh(VPS_SSH, rollback, timeout=180) raise RuntimeError(f"StreamScope build preflight failed: {err[-500:]}") deploy = "cd /app-config/github/speedtest && (docker rm -f speedtest >/dev/null 2>&1 || true) && docker compose up -d --remove-orphans" rc, out, err = _ssh(VPS_SSH, deploy, timeout=180) if rc != 0: _ssh(VPS_SSH, rollback, timeout=180) raise RuntimeError(f"StreamScope deployment failed: {(err or out)[-500:]}") health = "for i in $(seq 1 45); do curl -fsS --max-time 3 http://127.0.0.1:8080/api/health >/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"StreamScope health check failed and rollback was attempted: {err[-300:]}") return { "status": "deployed", "health": "ok", "application": "streamscope", "database": "/app-config/speedtest/data/streamscope.db", "public_port": False, "mtr": True, } @app.post("/vps/speedtest/deploy") async def vps_speedtest_deploy(req: SpeedtestDeployRequest, _=Depends(_verify)): secret_pattern = r"[A-Za-z0-9!@#%_+=:,.?-]{24,128}" if not re.fullmatch(secret_pattern, req.stats_password): raise HTTPException(400, "stats password must be 24-128 safe characters") if not re.fullmatch(secret_pattern, req.session_secret): raise HTTPException(400, "session secret must be 24-128 safe characters") contents = await asyncio.gather(*( _fetch_forgejo_text("sascha/speedtest", path) for path in SPEEDTEST_REPO_FILES )) files = dict(zip(SPEEDTEST_REPO_FILES, contents)) result = await asyncio.to_thread( _deploy_speedtest_compose, files, req.stats_password, req.session_secret ) _audit("/vps/speedtest/deploy", "POST", 200, "Git-managed StreamScope with private history and MTR") return result class GuckAdminDeployRequest(BaseModel): dry_run: bool = True def _validate_guck_admin_bundle(files: dict[str, str]) -> None: if set(files) != set(GUCK_ADMIN_REPO_FILES): raise ValueError("guck-admin source bundle is incomplete") compose = files["guck-admin/compose.yaml"] control = files["guck-admin/src/control.py"] compose_required = ( "network_mode: host", "NET_ADMIN", "/app-config/guck-admin/data:/data", "GUCK_LIMIT: /host/guck-limit.sh", ) if any(item not in compose for item in compose_required): raise ValueError("guck-admin compose is missing a required security or persistence setting") policy_required = ( "CREATE TABLE IF NOT EXISTS custom_networks", "def sync_custom_networks", "2a00:8c40:f000::/36", "45.58.235.0/24", ) if any(item not in control for item in policy_required): raise ValueError("guck-admin custom VPN policy is incomplete") def _guck_admin_remote_python(target: str, script: str, timeout: int = 60): encoded = base64.b64encode(script.encode()).decode() command = f"sudo -n python3 -c \"import base64;exec(base64.b64decode('{encoded}'))\"" return _ssh(target, command, timeout=timeout) def _deploy_guck_admin_compose(files: dict[str, str], dry_run: bool = True) -> dict: _validate_guck_admin_bundle(files) inventory = _find_inventory_host("guck-vps") if not inventory: raise RuntimeError("guck-vps is missing from Butler inventory") target = f'{inventory["user"]}@{inventory["ip"]}' if dry_run: return { "status": "validated", "dry_run": True, "host": "guck-vps", "files": len(files), "policy_networks": ["Mozilla Firefox VPN IPv6", "Fastly VPN IPv4"], } relative_files = {path.removeprefix("guck-admin/"): content for path, content in files.items()} deploy_script = f"""from pathlib import Path import os, shutil stack = Path('/app-config/guck-admin') backup = Path('/app-config/deployment-backups/guck-admin-rollback') files = {relative_files!r} if backup.exists(): shutil.rmtree(backup) backup.mkdir(parents=True, exist_ok=True) stack.mkdir(parents=True, exist_ok=True) for relative, content in files.items(): target = stack / relative old = backup / relative if target.exists(): old.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(target, old) target.parent.mkdir(parents=True, exist_ok=True) temporary = target.with_name(target.name + '.butler-new') temporary.write_text(content) os.replace(temporary, target) """ rollback_script = f"""from pathlib import Path import os, shutil stack = Path('/app-config/guck-admin') backup = Path('/app-config/deployment-backups/guck-admin-rollback') files = {tuple(relative_files)!r} for relative in files: target = stack / relative old = backup / relative if old.exists(): target.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(old, target) elif target.exists(): target.unlink() """ rc, _out, err = _guck_admin_remote_python(target, deploy_script, timeout=90) if rc != 0: raise RuntimeError(f"guck-admin file deployment failed: {err[-300:]}") def rollback(): _guck_admin_remote_python(target, rollback_script, timeout=90) _ssh(target, "cd /app-config/guck-admin && sudo -n docker compose up -d --build --remove-orphans", timeout=600) preflight = "cd /app-config/guck-admin && sudo -n docker compose config -q && sudo -n docker compose build --pull" rc, _out, err = _ssh(target, preflight, timeout=600) if rc != 0: rollback() raise RuntimeError(f"guck-admin build preflight failed: {err[-500:]}") deploy = "cd /app-config/guck-admin && sudo -n docker compose up -d --remove-orphans" rc, out, err = _ssh(target, deploy, timeout=240) if rc != 0: rollback() raise RuntimeError(f"guck-admin deployment failed: {(err or out)[-500:]}") health = "for i in $(seq 1 45); do curl -fsS --max-time 3 http://127.0.0.1:9090/health >/dev/null && exit 0; sleep 2; done; exit 1" rc, _out, err = _ssh(target, health, timeout=105) if rc != 0: rollback() raise RuntimeError(f"guck-admin health check failed and rollback was attempted: {err[-300:]}") policy = "curl -fsS -X POST --max-time 120 http://127.0.0.1:9090/actions/limiter/refresh >/dev/null && sudo -n ipset test vpn-v6 2a00:8c40:f02d:a34c::1 && sudo -n ipset test vpn-v4 45.58.235.7 && sudo -n tc class show dev ens3 | python3 -c \"import sys; s=sys.stdin.read(); raise SystemExit(0 if '1:300' in s else 1)\"" rc, out, err = _ssh(target, policy, timeout=180) if rc != 0: rollback() raise RuntimeError(f"guck-admin policy verification failed and rollback was attempted: {(err or out)[-500:]}") return { "status": "deployed", "health": "ok", "policy": "verified", "host": "guck-vps", "custom_networks": ["2a00:8c40:f000::/36", "45.58.235.0/24"], "ipv4": True, "ipv6": True, } @app.post("/vps/guck-admin/deploy") async def vps_guck_admin_deploy(req: GuckAdminDeployRequest, _=Depends(_verify)): try: contents = await asyncio.gather(*( _fetch_forgejo_text("sascha/guck-vps", path) for path in GUCK_ADMIN_REPO_FILES )) files = dict(zip(GUCK_ADMIN_REPO_FILES, contents)) result = await asyncio.to_thread(_deploy_guck_admin_compose, files, req.dry_run) except ValueError as exc: raise HTTPException(400, str(exc)) from exc except Exception as exc: raise HTTPException(502, f"guck-admin deployment failed: {str(exc)[-500:]}") from exc _audit("/vps/guck-admin/deploy", "POST", 200, f"dry_run={req.dry_run}; Git-managed custom VPN policy") return result # --- VM Lifecycle Endpoints --- import subprocess as _sp AUTOMATION1 = VM_CFG.get("automation_host", "sascha@10.5.85.5") if VM_CFG else "sascha@10.5.85.5" ISO_BUILDER = VM_CFG.get("iso_builder_path", "/app-config/ansible/iso-builder/build-iso.sh") if VM_CFG else "/app-config/ansible/iso-builder/build-iso.sh" class VMCreate(BaseModel): node: int ip: str hostname: str cores: int = 2 memory: int = 4096 disk: int = 32 def _ssh(host, cmd, timeout=600): try: r = _sp.run(["ssh","-o","ConnectTimeout=10","-o","StrictHostKeyChecking=accept-new", "-o","UserKnownHostsFile=/tmp/butler_known_hosts",host,cmd], capture_output=True, text=True, timeout=timeout) return r.returncode, r.stdout, r.stderr except _sp.TimeoutExpired: return 124, "", f"SSH command timed out after {timeout} seconds" PAPERLESS_GATEWAY = AUTOMATION1 PAPERLESS_SSH = "sascha@10.5.1.120" PAPERLESS_CONSUME_DIR = "/app-config/paperless/consume" PAPERLESS_MAX_IMPORT_BYTES = 50 * 1024 * 1024 def _ssh_bytes(host: str, cmd: str, payload: bytes, timeout: int = 120): """Send a binary payload to a fixed remote command over Butler-managed SSH.""" try: result = _sp.run( ["ssh", "-o", "ConnectTimeout=10", "-o", "StrictHostKeyChecking=accept-new", "-o", "UserKnownHostsFile=/tmp/butler_known_hosts", host, cmd], input=payload, capture_output=True, timeout=timeout, ) return result.returncode, result.stdout.decode(errors="replace"), result.stderr.decode(errors="replace") except _sp.TimeoutExpired: return 124, "", f"SSH upload timed out after {timeout} seconds" def _paperless_import_filename(filename: str, digest: str) -> str: base_name = os.path.basename(filename or "document.pdf") stem = re.sub(r"[^A-Za-z0-9._-]+", "_", base_name.rsplit(".", 1)[0]).strip("._-") if not stem: stem = "document" return f"{stem[:100]}-{digest[:12]}.pdf" def _paperless_gateway_command(command: str) -> str: """Route through automation1, whose deployment key is authorized on Homelab VMs.""" if "'" in command: raise ValueError("Paperless remote command contains an unsafe quote") return ( "ssh -o ConnectTimeout=10 -o StrictHostKeyChecking=accept-new " f"-o UserKnownHostsFile=/tmp/paperless_known_hosts {PAPERLESS_SSH} '{command}'" ) @app.post("/paperless/import") async def paperless_import(request: Request, filename: str = Query(..., min_length=1, max_length=180), _=Depends(_verify)): """Queue a PDF in Paperless without exposing Paperless or SSH to the caller.""" payload = await request.body() if not payload or not payload.startswith(b"%PDF-"): raise HTTPException(400, "Only valid PDF documents are accepted") if len(payload) > PAPERLESS_MAX_IMPORT_BYTES: raise HTTPException(413, "PDF exceeds the 50 MiB import limit") digest = hashlib.sha256(payload).hexdigest() import_name = _paperless_import_filename(filename, digest) remote_path = f"{PAPERLESS_CONSUME_DIR}/{import_name}" command = ( f"sudo -n install -d -o sascha -g sascha -m 0755 {PAPERLESS_CONSUME_DIR} && " f"tmp=$(mktemp /tmp/paperless-import.XXXXXX) && " f"cat > \"$tmp\" && sudo -n install -o sascha -g sascha -m 0644 \"$tmp\" {remote_path} && rm -f \"$tmp\"" ) rc, out, err = await asyncio.to_thread( _ssh_bytes, PAPERLESS_GATEWAY, _paperless_gateway_command(command), payload, 180 ) if rc != 0: raise HTTPException(502, (err or out).strip()[-500:] or "Paperless import transfer failed") _audit("/paperless/import", "POST", 202, f"sha256={digest}; bytes={len(payload)}") return {"status": "queued", "filename": import_name, "sha256": digest, "bytes": len(payload)} @app.get("/paperless/import/status") async def paperless_import_status(filename: str = Query(..., min_length=1, max_length=180), _=Depends(_verify)): if not re.fullmatch(r"[A-Za-z0-9._-]+\.pdf", filename): raise HTTPException(400, "Invalid import filename") remote_path = f"{PAPERLESS_CONSUME_DIR}/{filename}" command = ( f"if sudo -n test -f {remote_path}; then echo QUEUED; else echo CONSUMED; fi; " f"logs=$(sudo -n docker logs --since 15m paperless-ngx 2>&1); " f"printf \"%s\\n\" \"$logs\" | grep -F -- {filename} | tail -20 || true; " f"task=$(printf \"%s\\n\" \"$logs\" | grep -F -- {filename} | " f"grep -oE \"\\[[0-9a-f]{{8}}\\]\" | tail -1 | tr -d \"[]\"); " f"if test -n \"$task\"; then printf \"%s\\n\" \"$logs\" | grep -F -- \"[$task]\" | tail -20; fi" ) rc, out, err = await asyncio.to_thread( _ssh, PAPERLESS_GATEWAY, _paperless_gateway_command(command), 45 ) if rc != 0: raise HTTPException(502, (err or out).strip()[-500:] or "Paperless status check failed") lines = out.splitlines() state = lines[0].strip().lower() if lines else "unknown" return {"status": state, "filename": filename, "recent_log": lines[1:]} 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(["sudo", "-n", "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": 0 if fields[7] == "off" else 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} 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} SASCHA_MEDIA_VPS_HOST = "pfannkuchen" SASCHA_MEDIA_EMBY_HOST = "emby-sascha" SASCHA_MEDIA_INTERFACE = "wg-media" SASCHA_MEDIA_VPS_ADDRESS = "10.11.13.1/32" SASCHA_MEDIA_EMBY_ADDRESS = "10.11.13.3/32" SASCHA_MEDIA_PORT = 51821 SASCHA_MEDIA_MTU = 1340 SASCHA_MEDIA_CONFIRMATION = "DEPLOY_DIRECT_SASCHA_MEDIA_TUNNEL" class SaschaMediaTunnelRequest(BaseModel): dry_run: bool = True confirmation: str | None = None def _sascha_media_edge_audit_command() -> str: script = '''import json, re, subprocess +from pathlib import Path + +def run(args): + proc = subprocess.run(args, text=True, capture_output=True, timeout=30) + return {"rc": proc.returncode, "stdout": proc.stdout.strip(), "stderr": proc.stderr.strip()[-300:]} + +result = {"hostname": "pfannkuchen", "caddy": {"container_running": False, "protocols": [], "protocols_explicit": False, "upstreams": [], "site_block": []}, "network": {}} +inspect = run(["docker", "inspect", "caddy", "--format", "{{.State.Running}}"]) +result["caddy"]["container_running"] = inspect["rc"] == 0 and inspect["stdout"] == "true" +adapt = run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"]) +if adapt["rc"] == 0: + try: + config = json.loads(adapt["stdout"]) + servers = config.get("apps", {}).get("http", {}).get("servers", {}) + explicit = [] + def walk(value, matched=False): + if isinstance(value, dict): + current = matched + host = value.get("host") + if isinstance(host, list) and "tv.sascha-lutz.de" in host: + current = True + if current and isinstance(value.get("dial"), str): + result["caddy"]["upstreams"].append(value["dial"]) + for child in value.values(): walk(child, current) + elif isinstance(value, list): + for child in value: walk(child, matched) + for server in servers.values(): + protocols = server.get("protocols") + if isinstance(protocols, list): explicit.extend(protocols) + walk(server) + result["caddy"]["protocols_explicit"] = bool(explicit) + result["caddy"]["protocols"] = sorted(set(explicit)) if explicit else ["h1", "h2", "h3"] + result["caddy"]["upstreams"] = sorted(set(result["caddy"]["upstreams"])) + except Exception as exc: + result["caddy"]["adapt_error"] = str(exc)[:200] +else: + result["caddy"]["adapt_error"] = adapt["stderr"] or "caddy adapt failed" + +path = Path("/app-config/caddy/Caddyfile") +if path.exists(): + lines = path.read_text(encoding="utf-8", errors="replace").splitlines() + collecting = False; depth = 0; selected = [] + for line in lines: + if not collecting and re.match(r"^\\s*tv\\.sascha-lutz\\.de\\s*\\{", line): collecting = True + if collecting: + depth += line.count("{") - line.count("}") + if not re.search(r"(?i)(password|secret|token|private|basicauth|basic_auth|hash)", line): selected.append(line.strip()) + else: selected.append("[REDACTED SENSITIVE DIRECTIVE]") + if depth == 0: break + result["caddy"]["site_block"] = selected + +result["network"]["wg_media_active"] = subprocess.run(["systemctl", "is-active", "--quiet", "wg-quick@wg-media"]).returncode == 0 +result["network"]["wg_media_enabled"] = subprocess.run(["systemctl", "is-enabled", "--quiet", "wg-quick@wg-media"]).returncode == 0 +result["network"]["udp_51821"] = "51821" in run(["ss", "-H", "-lun"]) ["stdout"] +for name, target in (("legacy_route", "10.6.1.103"), ("direct_route", "10.11.13.3")): + result["network"][name] = run(["ip", "route", "get", target])["stdout"][:300] +print(json.dumps(result)) +'''.replace("\n+", "\n") encoded = base64.b64encode(script.encode()).decode() return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' @app.get("/media/edge/sascha") async def sascha_media_edge_audit(_=Depends(_verify)): """Return a redacted snapshot of the dedicated Hetzner edge for tv.sascha-lutz.de.""" inventory = await asyncio.to_thread(_find_inventory_host, SASCHA_MEDIA_VPS_HOST) if not inventory: raise HTTPException(404, "Hetzner media edge not found") target = f'{inventory["user"]}@{inventory["ip"]}' rc, out, err = await asyncio.to_thread(_ssh, target, _sascha_media_edge_audit_command(), 45) if rc != 0: raise HTTPException(502, (err or out).strip()[-500:] or "Sascha media edge audit failed") try: result = json.loads(out) except json.JSONDecodeError as exc: raise HTTPException(502, "Sascha media edge audit returned invalid JSON") from exc _audit("/media/edge/sascha", "GET", 200, "redacted live media-edge snapshot") return result def _media_tunnel_key_command() -> str: script = '''import json, os, subprocess +from pathlib import Path +root = Path("/app-config/wireguard-media") +private = root / "private.key" +public = root / "public.key" +root.mkdir(parents=True, exist_ok=True) +os.chmod(root, 0o700) +created = not private.exists() +if created: + key = subprocess.run(["wg", "genkey"], check=True, text=True, capture_output=True).stdout.strip() + private.write_text(key + "\\n") + os.chmod(private, 0o600) +if not public.exists() or created: + pub = subprocess.run(["wg", "pubkey"], input=private.read_text(), check=True, text=True, capture_output=True).stdout.strip() + public.write_text(pub + "\\n") + os.chmod(public, 0o644) +print(json.dumps({"public_key": public.read_text().strip(), "created": created})) +'''.replace("\n+", "\n") encoded = base64.b64encode(script.encode()).decode() return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' def _media_tunnel_install_command(role: Literal["vps", "emby"], peer_public_key: str) -> str: if not re.fullmatch(r"[A-Za-z0-9+/]{43}=", peer_public_key): raise ValueError("Invalid WireGuard public key") settings = { "vps": { "address": SASCHA_MEDIA_VPS_ADDRESS, "peer": SASCHA_MEDIA_EMBY_ADDRESS, "endpoint": None, "keepalive": None, "listen": SASCHA_MEDIA_PORT, }, "emby": { "address": SASCHA_MEDIA_EMBY_ADDRESS, "peer": SASCHA_MEDIA_VPS_ADDRESS, "endpoint": f"46.225.230.72:{SASCHA_MEDIA_PORT}", "keepalive": 15, "listen": None, }, }[role] script = f'''import json, os, shutil, subprocess, time +from pathlib import Path +root = Path("/app-config/wireguard-media") +config = root / "wg-media.conf" +etc = Path("/etc/wireguard/wg-media.conf") +backup_dir = root / "backups" +backup_dir.mkdir(parents=True, exist_ok=True) +private = (root / "private.key").read_text().strip() +listen = {settings['listen']!r} +existed = config.exists() +backup = None +if existed: + backup = backup_dir / ("wg-media.conf." + time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())) + shutil.copy2(config, backup) +lines = ["[Interface]", "Address = {settings['address']}", "MTU = {SASCHA_MEDIA_MTU}", "PrivateKey = " + private] +if listen is not None: lines.append("ListenPort = " + str(listen)) +if {role!r} == "vps": + lines.extend(["PostUp = iptables -C INPUT -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT 2>/dev/null || iptables -I INPUT 1 -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT", "PreDown = iptables -D INPUT -p udp --dport {SASCHA_MEDIA_PORT} -j ACCEPT 2>/dev/null || true"]) +lines.extend(["", "[Peer]", "PublicKey = {peer_public_key}", "AllowedIPs = {settings['peer']}"]) +if {settings['endpoint']!r}: lines.append("Endpoint = " + {settings['endpoint']!r}) +if {settings['keepalive']!r}: lines.append("PersistentKeepalive = " + str({settings['keepalive']!r})) +candidate = root / "wg-media.conf.candidate" +candidate.write_text("\\n".join(lines) + "\\n") +os.chmod(candidate, 0o600) +check = subprocess.run(["wg-quick", "strip", str(candidate)], text=True, capture_output=True) +if check.returncode != 0: + candidate.unlink(missing_ok=True) + raise RuntimeError(check.stderr.strip() or "wg-quick validation failed") +os.replace(candidate, config) +os.chmod(config, 0o600) +etc.parent.mkdir(parents=True, exist_ok=True) +if etc.is_symlink() or etc.exists(): + if etc.is_symlink() and etc.resolve() == config.resolve(): pass + elif etc.exists(): + etc_backup = backup_dir / ("etc-wg-media.conf." + time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())) + shutil.move(etc, etc_backup) + etc.symlink_to(config) +else: etc.symlink_to(config) +proc = subprocess.run(["systemctl", "enable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30) +if proc.returncode != 0: + if backup: shutil.copy2(backup, config) + else: config.unlink(missing_ok=True) + raise RuntimeError(proc.stderr.strip() or proc.stdout.strip() or "failed to start wg-media") +print(json.dumps({{"role": {role!r}, "status": "configured", "address": {settings['address']!r}, "backup": str(backup) if backup else None, "existed": existed}})) +'''.replace("\n+", "\n") encoded = base64.b64encode(script.encode()).decode() return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' def _media_tunnel_verify_command(source: str, destination: str, check_emby: bool = False) -> str: extra = '' if check_emby: extra = '''\nhttp = subprocess.run(["curl", "-sS", "--max-time", "10", "-o", "/dev/null", "-w", "%{http_code}", "http://10.11.13.3:8096/System/Ping"], text=True, capture_output=True)\nresult["emby_http"] = http.stdout.strip() if http.returncode == 0 else "000"''' script = f'''import json, subprocess, time +result = {{"ping": False, "handshake": False}} +for _ in range(10): + ping = subprocess.run(["ping", "-c", "1", "-W", "2", "-I", {source!r}, {destination!r}], text=True, capture_output=True) + hand = subprocess.run(["sudo", "-n", "wg", "show", "wg-media", "latest-handshakes"], text=True, capture_output=True) + now = int(time.time()) + stamps = [] + for line in hand.stdout.splitlines(): + try: stamps.append(int(line.split()[-1])) + except Exception: pass + result["ping"] = ping.returncode == 0 + result["handshake"] = any(stamp > 0 and now - stamp < 60 for stamp in stamps) + if result["ping"] and result["handshake"]: break + time.sleep(2) +{extra} +print(json.dumps(result)) +'''.replace("\n+", "\n") encoded = base64.b64encode(script.encode()).decode() return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' def _media_tunnel_rollback_command(remove_keys: bool) -> str: script = f'''import json, shutil, subprocess +from pathlib import Path +root = Path("/app-config/wireguard-media") +config = root / "wg-media.conf" +backups = sorted((root / "backups").glob("wg-media.conf.*")) if (root / "backups").exists() else [] +subprocess.run(["systemctl", "disable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30) +if backups: + shutil.copy2(backups[-1], config) + subprocess.run(["systemctl", "enable", "--now", "wg-quick@wg-media"], text=True, capture_output=True, timeout=30) + status = "restored" +else: + config.unlink(missing_ok=True) + Path("/etc/wireguard/wg-media.conf").unlink(missing_ok=True) + if {remove_keys!r}: + (root / "private.key").unlink(missing_ok=True); (root / "public.key").unlink(missing_ok=True) + status = "removed" +print(json.dumps({{"status": status}})) +'''.replace("\n+", "\n") encoded = base64.b64encode(script.encode()).decode() return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' def _media_host_target(name: str) -> str: inventory = _find_inventory_host(name) if not inventory: raise RuntimeError(f"Inventory host missing: {name}") return f'{inventory["user"]}@{inventory["ip"]}' def _ssh_json(target: str, command: str, timeout: int = 45) -> dict: rc, out, err = _ssh(target, command, timeout) if rc != 0: raise RuntimeError((err or out).strip()[-500:] or f"remote command failed on {target}") try: return json.loads(out) except json.JSONDecodeError as exc: raise RuntimeError(f"remote command returned invalid JSON on {target}") from exc def _deploy_sascha_media_tunnel() -> dict: vps = _media_host_target(SASCHA_MEDIA_VPS_HOST) emby = _media_host_target(SASCHA_MEDIA_EMBY_HOST) vps_key = _ssh_json(vps, _media_tunnel_key_command()) emby_key = _ssh_json(emby, _media_tunnel_key_command()) configured = [] try: vps_install = _ssh_json(vps, _media_tunnel_install_command("vps", emby_key["public_key"]), 60) configured.append((vps, bool(vps_key.get("created")))) emby_install = _ssh_json(emby, _media_tunnel_install_command("emby", vps_key["public_key"]), 60) configured.append((emby, bool(emby_key.get("created")))) vps_check = _ssh_json(vps, _media_tunnel_verify_command("10.11.13.1", "10.11.13.3", True), 35) emby_check = _ssh_json(emby, _media_tunnel_verify_command("10.11.13.3", "10.11.13.1"), 35) if not (vps_check.get("ping") and vps_check.get("handshake") and emby_check.get("ping") and emby_check.get("handshake")): raise RuntimeError("direct media tunnel verification failed") if vps_check.get("emby_http") not in {"200", "401"}: raise RuntimeError("Emby did not answer through the direct media tunnel") return { "status": "deployed", "interface": SASCHA_MEDIA_INTERFACE, "vps_address": SASCHA_MEDIA_VPS_ADDRESS, "emby_address": SASCHA_MEDIA_EMBY_ADDRESS, "listen_port": SASCHA_MEDIA_PORT, "mtu": SASCHA_MEDIA_MTU, "handshake": True, "ping_vps_to_emby": True, "ping_emby_to_vps": True, "emby_http": vps_check.get("emby_http"), "rollback_backups": [vps_install.get("backup"), emby_install.get("backup")], } except Exception: for target, created in reversed(configured): try: _ssh_json(target, _media_tunnel_rollback_command(created), 60) except Exception: pass raise @app.post("/network/media-tunnel/sascha") async def deploy_sascha_media_tunnel(req: SaschaMediaTunnelRequest, _=Depends(_verify)): """Deploy the fixed direct Hetzner-to-emby-sascha WireGuard media tunnel.""" plan = { "status": "would_deploy", "interface": SASCHA_MEDIA_INTERFACE, "vps_address": SASCHA_MEDIA_VPS_ADDRESS, "emby_address": SASCHA_MEDIA_EMBY_ADDRESS, "listen_port": SASCHA_MEDIA_PORT, "mtu": SASCHA_MEDIA_MTU, "allowed_ips": [SASCHA_MEDIA_VPS_ADDRESS, SASCHA_MEDIA_EMBY_ADDRESS], "keeps_legacy_node6_path": True, } if req.dry_run: _audit("/network/media-tunnel/sascha", "POST", 200, "dry_run=True", True) return plan if req.confirmation != SASCHA_MEDIA_CONFIRMATION: raise HTTPException(400, f"confirmation must be {SASCHA_MEDIA_CONFIRMATION}") try: result = await asyncio.to_thread(_deploy_sascha_media_tunnel) except Exception as exc: _audit("/network/media-tunnel/sascha", "POST", 502, "deployment failed; rollback attempted") raise HTTPException(502, str(exc)[-500:]) from exc _audit("/network/media-tunnel/sascha", "POST", 200, "direct media tunnel deployed") return result SYSCTL_AUDIT_KEYS = ( "net.core.default_qdisc", "net.core.rmem_default", "net.core.rmem_max", "net.core.wmem_default", "net.core.wmem_max", "net.core.netdev_max_backlog", "net.core.somaxconn", "net.ipv4.ip_forward", "net.ipv4.tcp_congestion_control", "net.ipv4.tcp_fastopen", "net.ipv4.tcp_mtu_probing", "net.ipv4.tcp_no_metrics_save", "net.ipv4.tcp_rmem", "net.ipv4.tcp_slow_start_after_idle", "net.ipv4.tcp_window_scaling", "net.ipv4.tcp_wmem", ) def _sysctl_audit_command() -> str: script = f'''import glob, json from pathlib import Path keys = {SYSCTL_AUDIT_KEYS!r} live, errors = {{}}, {{}} for key in keys: try: live[key] = Path("/proc/sys/" + key.replace(".", "/")).read_text().strip() except OSError as exc: errors[key] = str(exc)[:160] persistent = {{}} for path in ["/etc/sysctl.conf", *sorted(glob.glob("/etc/sysctl.d/*.conf"))]: try: with open(path, encoding="utf-8", errors="replace") as handle: for raw in handle: line = raw.split("#", 1)[0].strip() if "=" not in line: continue key, value = (part.strip() for part in line.split("=", 1)) if key in keys: persistent.setdefault(key, []).append({{"file": path, "value": value}}) except (FileNotFoundError, PermissionError): pass print(json.dumps({{"live": live, "persistent": persistent, "errors": errors}})) ''' encoded = base64.b64encode(script.encode()).decode() return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"' @app.get("/system/sysctl/{host}") async def system_sysctl_audit(host: str, _=Depends(_verify)): if not re.fullmatch(r"[a-z0-9][a-z0-9-]{0,62}", host): raise HTTPException(400, "Invalid host name") if host == "vps": target = VPS_SSH else: 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, _sysctl_audit_command(), 30) if rc != 0: raise HTTPException(502, (err or out).strip()[-500:] or "sysctl audit failed") try: result = json.loads(out) except json.JSONDecodeError as exc: 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, 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 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 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"), "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"), "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"), "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"), "minecraft_inventory": run("grep -in 'minecraft' /app-config/ansible/pfannkuchen.ini 2>/dev/null || true"), "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)) ''' 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 _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 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" 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): 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}) 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','')}" @app.get("/vm/list") async def vm_list(_=Depends(_verify)): auth = _pve_auth() vms = [] async with httpx.AsyncClient(verify=False, timeout=15) as c: nodes = await c.get("https://10.5.85.11:8006/api2/json/nodes", headers={"Authorization": auth}) for n in nodes.json().get("data", []): r = await c.get(f"https://10.5.85.11:8006/api2/json/nodes/{n['node']}/qemu", headers={"Authorization": auth}) for vm in r.json().get("data", []): vm["node"] = n["node"] vms.append(vm) return vms @app.post("/vm/create") async def vm_create(req: VMCreate, _=Depends(_verify), dry_run: bool = Query(False)): if dry_run: _audit("/vm/create", "POST", 200, f"dry_run: {req.hostname} {req.ip} node{req.node}", dry_run=True) return {"dry_run": True, "would_create": {"hostname": req.hostname, "ip": req.ip, "node": req.node, "cores": req.cores, "memory": req.memory, "disk": req.disk}, "steps": ["iso-builder", "wait ssh", "add inventory", "ansible setup"]} steps = [] # Step 1: Build ISO + create VM via iso-builder on automation1 cmd = f"{ISO_BUILDER} --node {req.node} --ip {req.ip} --hostname {req.hostname} --cores {req.cores} --memory {req.memory} --disk {req.disk} --password '{VM_CFG.get('default_password', 'changeme')}' --create-vm" rc, out, err = _ssh(AUTOMATION1, f"cd /app-config/ansible/iso-builder && {cmd}", timeout=300) if rc != 0: return JSONResponse({"error": "iso-builder failed", "stderr": err[-500:], "stdout": out[-500:]}, status_code=500) steps.append("iso-builder: ok") # Step 2: Wait for SSH (up to 6 min) ok = False for _ in range(36): try: rc2, out2, _ = _ssh(f"sascha@{req.ip}", "hostname", timeout=10) if rc2 == 0: ok = True steps.append(f"ssh: {out2.strip()} reachable") break except Exception: pass await asyncio.sleep(10) if not ok: return JSONResponse({"error": "SSH timeout", "steps": steps}, status_code=504) # Step 2.5: Add to Ansible inventory ini = "/app-config/ansible/pfannkuchen.ini" group = getattr(req, 'group', 'auto') inv_cmd = f"""python3 -c " lines = open('{ini}').readlines() if not any('{req.hostname} ' in l for l in lines): out = [] found = False for l in lines: out.append(l) if l.strip() == '[auto]': found = True elif found and (l.startswith('[') or l.strip() == ''): out.insert(-1, '{req.hostname} ansible_host={req.ip}\\n') found = False if found: out.append('{req.hostname} ansible_host={req.ip}\\n') open('{ini}','w').writelines(out) print('added') else: print('exists') " """ _ssh(AUTOMATION1, inv_cmd, timeout=30) _ssh(AUTOMATION1, f"mkdir -p /app-config/ansible/host_vars/{req.hostname} && printf 'ansible_host: {req.ip}\\nansible_user: sascha\\n' > /app-config/ansible/host_vars/{req.hostname}/vars.yml", timeout=30) _ssh(AUTOMATION1, f"ssh-keygen -f /home/sascha/.ssh/known_hosts -R {req.ip} 2>/dev/null; ssh -o StrictHostKeyChecking=accept-new sascha@{req.ip} hostname 2>/dev/null", timeout=30) steps.append("inventory: added") # Step 3: Ansible base setup via direct SSH (reliable fallback) rc3, _, err3 = _ssh(AUTOMATION1, f"cd /app-config/ansible && bash pfannkuchen.sh setup {req.hostname}", timeout=600) steps.append(f"ansible: {'ok' if rc3 == 0 else 'failed (rc=' + str(rc3) + ')'}") _audit("/vm/create", "POST", 200 if rc3 == 0 else 500, f"{req.hostname} {req.ip}") return {"status": "ok" if rc3 == 0 else "partial", "hostname": req.hostname, "ip": req.ip, "node": req.node, "steps": steps} @app.get("/vm/status/{vmid}") async def vm_status(vmid: int, _=Depends(_verify)): auth = _pve_auth() async with httpx.AsyncClient(verify=False, timeout=10) as c: nodes = await c.get("https://10.5.85.11:8006/api2/json/nodes", headers={"Authorization": auth}) for n in nodes.json().get("data", []): r = await c.get(f"https://10.5.85.11:8006/api2/json/nodes/{n['node']}/qemu/{vmid}/status/current", headers={"Authorization": auth}) if r.status_code == 200: return r.json().get("data", {}) return JSONResponse({"error": "VM not found"}, status_code=404) @app.delete("/vm/{vmid}") async def vm_delete(vmid: int, _=Depends(_verify)): """Simple VM delete - Proxmox only (legacy).""" auth = _pve_auth() async with httpx.AsyncClient(verify=False, timeout=30) as c: nodes = await c.get("https://10.5.85.11:8006/api2/json/nodes", headers={"Authorization": auth}) for n in nodes.json().get("data", []): r = await c.delete(f"https://10.5.85.11:8006/api2/json/nodes/{n['node']}/qemu/{vmid}", headers={"Authorization": auth}) if r.status_code == 200: return r.json() return JSONResponse({"error": "VM not found"}, status_code=404) @app.delete("/vm/destroy/{vmid}") async def vm_destroy_full(vmid: int, _=Depends(_verify), dry_run: bool = Query(False)): """ Complete VM destruction with full cleanup: 1. Stop VM (required before destroy) 2. Destroy VM (Proxmox) 3. Remove from Dockhand (by IP) 4. Delete Forgejo repo (by hostname) 5. Remove from Ansible inventory 6. Remove from Uptime Kuma monitoring Returns detailed cleanup report. """ auth = _pve_auth() results = {"vmid": vmid, "dry_run": dry_run, "steps": {}} async with httpx.AsyncClient(verify=False, timeout=30) as c: # Step 1: Find VM and get details from vm/list (more reliable than config endpoint) vm_info = None node_name = None vms = await c.get("https://10.5.85.11:8006/api2/json/nodes", headers={"Authorization": auth}) for n in vms.json().get("data", []): r = await c.get(f"https://10.5.85.11:8006/api2/json/nodes/{n['node']}/qemu", headers={"Authorization": auth}) for vm in r.json().get("data", []): if vm.get("vmid") == vmid: vm_info = vm node_name = n["node"] break if vm_info: break if not vm_info: return JSONResponse({"error": f"VM {vmid} not found"}, status_code=404) hostname = vm_info.get("name", "") or str(vmid) # Extract IP from vm list or config ip = vm_info.get("ip", "") if not ip: # Try to get from net0 config cfg = await c.get(f"https://10.5.85.11:8006/api2/json/nodes/{node_name}/qemu/{vmid}/config", headers={"Authorization": auth}) net0 = cfg.json().get("data", {}).get("net0", "") if "ip=" in net0: ip = net0.split("ip=")[-1].split(",")[0] results["vm_info"] = {"hostname": hostname, "ip": ip, "node": node_name, "status": vm_info.get("status", "unknown")} # Step 2: Stop VM (if running) if vm_info.get("status") == "running": if dry_run: results["steps"]["stop_vm"] = {"status": "dry_run", "message": f"Would stop VM {vmid}"} else: r = await c.post(f"https://10.5.85.11:8006/api2/json/nodes/{node_name}/qemu/{vmid}/status/stop", headers={"Authorization": auth}) results["steps"]["stop_vm"] = {"status": "ok" if r.status_code == 200 else "failed", "detail": r.json()} if r.status_code != 200: return JSONResponse({"error": f"Failed to stop VM: {r.json()}"}, status_code=500) # Wait for VM to stop await asyncio.sleep(5) else: results["steps"]["stop_vm"] = {"status": "skipped", "message": "VM already stopped"} # Step 3: Destroy VM if dry_run: results["steps"]["destroy_vm"] = {"status": "dry_run", "message": f"Would destroy VM {vmid}"} else: r = await c.delete(f"https://10.5.85.11:8006/api2/json/nodes/{node_name}/qemu/{vmid}", headers={"Authorization": auth}) results["steps"]["destroy_vm"] = {"status": "ok" if r.status_code == 200 else "failed", "detail": r.json()} # Step 4: Remove from Dockhand (by IP) if ip: if dry_run: results["steps"]["dockhand_remove"] = {"status": "dry_run", "message": f"Would remove Dockhand env for IP {ip}"} else: # Find environment by IP envs = await c.get("http://10.4.1.116:3000/api/environments", headers={"Authorization": f"Bearer {BUTLER_TOKEN}"}) env_id = None for env in envs.json(): if env.get("name", "").lower() == hostname.lower() or env.get("ip") == ip: env_id = env.get("id") break if env_id: r = await c.delete(f"http://10.4.1.116:3000/api/environments/{env_id}", headers={"Authorization": f"Bearer {BUTLER_TOKEN}"}) results["steps"]["dockhand_remove"] = {"status": "ok" if r.status_code in [200, 204] else "failed", "env_id": env_id} else: results["steps"]["dockhand_remove"] = {"status": "skipped", "message": "No Dockhand environment found"} else: results["steps"]["dockhand_remove"] = {"status": "skipped", "message": "No IP found"} # Step 5: Delete Forgejo repo (by hostname) if hostname: if dry_run: results["steps"]["forgejo_repo_delete"] = {"status": "dry_run", "message": f"Would delete repo sascha/{hostname}"} else: r = await c.delete(f"http://10.4.1.116:8888/forgejo/api/v1/repos/sascha/{hostname}", headers={"Authorization": f"Bearer {BUTLER_TOKEN}"}) results["steps"]["forgejo_repo_delete"] = {"status": "ok" if r.status_code in [200, 204] else "not_found", "detail": r.json() if r.status_code != 204 else "deleted"} else: results["steps"]["forgejo_repo_delete"] = {"status": "skipped", "message": "No hostname found"} # Step 6: Remove from Ansible inventory if hostname: if dry_run: results["steps"]["ansible_cleanup"] = {"status": "dry_run", "message": f"Would remove {hostname} from pfannkuchen.ini"} else: # Remove host from inventory remove_cmd = f'''python3 -c " lines = open('/app-config/ansible/pfannkuchen.ini').readlines() out = [l for l in lines if '{hostname}' not in l] open('/app-config/ansible/pfannkuchen.ini','w').writelines(out) print('removed') " ''' rc, out, err = _ssh(AUTOMATION1, remove_cmd, timeout=30) # Also remove host_vars _ssh(AUTOMATION1, f"rm -rf /app-config/ansible/host_vars/{hostname}", timeout=30) results["steps"]["ansible_cleanup"] = {"status": "ok" if rc == 0 else "failed", "detail": out.strip()} else: results["steps"]["ansible_cleanup"] = {"status": "skipped", "message": "No hostname found"} # Step 7: Remove from Uptime Kuma (if monitoring exists) if hostname: if dry_run: results["steps"]["uptime_kuma_remove"] = {"status": "dry_run", "message": f"Would remove monitor for {hostname}"} else: try: # Get all monitors from Kuma (no auth needed for local network) kuma_monitors = await c.get("http://10.200.200.1:3001/api/monitors", timeout=5) for monitor in kuma_monitors.json().get("data", []): if hostname.lower() in monitor.get("name", "").lower(): r = await c.delete(f"http://10.200.200.1:3001/api/monitors/{monitor['id']}", timeout=5) results["steps"]["uptime_kuma_remove"] = {"status": "ok" if r.status_code == 200 else "failed", "monitor_id": monitor["id"]} break else: results["steps"]["uptime_kuma_remove"] = {"status": "skipped", "message": "No Kuma monitor found"} except Exception as e: results["steps"]["uptime_kuma_remove"] = {"status": "error", "detail": str(e)} else: results["steps"]["uptime_kuma_remove"] = {"status": "skipped", "message": "No hostname found"} _audit(f"/vm/destroy/{vmid}", "DELETE", 200, f"dry_run={dry_run}") return results @app.post("/inventory/host") async def inventory_host(request: Request, _=Depends(_verify)): """Create or update an Ansible inventory host idempotently.""" body = await request.json() name, ip = body.get("name", ""), body.get("ip", "") group = body.get("group", "auto") user = body.get("user", "sascha") if not all(re.fullmatch(r"[A-Za-z0-9_.-]+", value) for value in (name, group, user)): raise HTTPException(400, "Invalid name, group, or user") try: ipaddress.ip_address(ip) except ValueError: raise HTTPException(400, "Invalid IP address") ini = "/app-config/ansible/pfannkuchen.ini" host_line = f"{name} ansible_host={ip} ansible_user={user}" script = f'''lines = open({ini!r}).readlines() name = {name!r} host_line = {host_line!r} group = {group!r} updated = False for idx, line in enumerate(lines): parts = line.split() if parts and parts[0] == name and "ansible_host=" in line: lines[idx] = host_line + "\\n" updated = True break if not updated: insert_at = None in_group = False for idx, line in enumerate(lines): if line.strip() == "[" + group + "]": in_group = True insert_at = idx + 1 continue if in_group and line.startswith("["): break if in_group: insert_at = idx + 1 if insert_at is None: lines.extend(["\\n[" + group + "]\\n", host_line + "\\n"]) else: lines.insert(insert_at, host_line + "\\n") open({ini!r}, "w").writelines(lines) print("updated" if updated else "added")''' encoded_script = base64.b64encode(script.encode()).decode() rc, out, err = _ssh( AUTOMATION1, f"python3 -c \"import base64;exec(base64.b64decode('{encoded_script}'))\"", timeout=30, ) if rc != 0: raise HTTPException(502, err.strip()[:500]) vars_script = ( f"mkdir -p /app-config/ansible/host_vars/{name} && " f"printf 'ansible_host: {ip}\\nansible_user: {user}\\n' > " f"/app-config/ansible/host_vars/{name}/vars.yml" ) rc2, _out2, err2 = _ssh(AUTOMATION1, vars_script, timeout=30) if rc2 != 0: raise HTTPException(502, err2.strip()[:500]) _audit("/inventory/host", "POST", 200, f"{name} {ip} {user}") return {"status": "ok", "name": name, "ip": ip, "group": group, "user": user, "result": out.strip()} @app.post("/ansible/run") async def ansible_run(request: Request, _=Depends(_verify)): body = await request.json() hostname = body.get("limit", body.get("hostname", "")) if not hostname: return JSONResponse({"error": "limit/hostname required"}, status_code=400) action = body.get("action", "setup") 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", "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"{prepare}{file_sync} && " f"bash pfannkuchen.sh {action} {hostname}" ) else: command = ( "cd /app-config/ansible && " "git pull --ff-only origin master && " f"bash pfannkuchen.sh {action} {hostname}" ) rc, out, err = _ssh(AUTOMATION1, command, timeout=600) _audit("/ansible/run", "POST", 200 if rc == 0 else 502, f"{action} {hostname}") if action != "setup": return { "status": "ok" if rc == 0 else "error", "action": action, "hostname": hostname, "rc": rc, "output": out[-4000:], "error": err[-1000:] if rc != 0 else "", } # After successful ansible run: sync Hawser token to Dockhand if rc == 0: try: # Get VM IP from inventory inv_path = "/app-config/ansible/pfannkuchen.ini" rc2, ip_out, _ = _ssh(AUTOMATION1, f"grep -E '^{hostname} ' {inv_path} | awk '{{print $2}}' | cut -d= -f2", timeout=10) vm_ip = ip_out.strip() if rc2 == 0 else None if vm_ip: # Read Hawser token from VM rc3, token_out, _ = _ssh(AUTOMATION1, f"ssh -o StrictHostKeyChecking=no sascha@{vm_ip} 'sudo grep ^TOKEN= /etc/hawser/config | cut -d= -f2' 2>/dev/null", timeout=15) hawser_token = token_out.strip() if rc3 == 0 and token_out.strip() else None if hawser_token: # Find environment in Dockhand by IP and update token async with httpx.AsyncClient(verify=False, timeout=10) as c: # Login to Dockhand login = await c.post(f"{SERVICES['dockhand']['url']}/api/auth/login", json={"username": "admin", "password": _read("dockhand") or ""}) if login.status_code == 200: cookie = dict(login.cookies) # Get all environments envs = await c.get(f"{SERVICES['dockhand']['url']}/api/environments", cookies=cookie) for env in envs.json(): if env.get("host") == vm_ip: # Update environment with Hawser token await c.put(f"{SERVICES['dockhand']['url']}/api/environments/{env['id']}", json={"hawserToken": hawser_token}, cookies=cookie) log.info(f"Updated Hawser token for {hostname} (env {env['id']})") break except Exception as e: log.warning(f"Hawser token sync failed for {hostname}: {e}") # NEW: SOPS + .env handling for Git-centric deployments try: log.info(f"Checking for SOPS .env setup for {hostname}") # Get SOPS age public key from automation1 rc4, age_pub, _ = _ssh(AUTOMATION1, "cat ~/.config/sops/age/keys.txt | grep '^# public key' | awk '{print $4}'", timeout=10) age_pub = age_pub.strip() if rc4 == 0 else None if age_pub: # Check if compose.yaml exists in /app-config/github/{hostname}/ rc5, compose_check, _ = _ssh(AUTOMATION1, f"test -f /app-config/github/{hostname}/compose.yaml && echo 'found' || echo 'missing'", timeout=10) if compose_check.strip() == "found": log.info(f"Found compose.yaml for {hostname}, generating .env") # Generate secrets import secrets as sec secret_key = base64.b64encode(sec.token_bytes(32)).decode() admin_pw = sec.token_urlsafe(16) db_pw = sec.token_urlsafe(16) # Store in vault cache _vault_cache[f"{hostname}_secret_key"] = secret_key _vault_cache[f"{hostname}_admin_password"] = admin_pw _vault_cache[f"{hostname}_db_password"] = db_pw # Build .env content env_lines = ["# Auto-generated by Butler", f"TZ=Europe/Berlin", "PUID=1000", "PGID=1000"] # Detect service type from hostname if "paperless" in hostname.lower(): env_lines.extend([ f"PAPERLESS_ADMIN_USER=admin", f"PAPERLESS_ADMIN_PASSWORD={admin_pw}", f"PAPERLESS_SECRET_KEY={secret_key}", f"PAPERLESS_URL=http://{vm_ip}:8000", "PAPERLESS_TIME_ZONE=Europe/Berlin", "PAPERLESS_OCR_LANGUAGE=deu", "PAPERLESS_REDIS=redis://redis:6379", "PAPERLESS_DBHOST=postgres", "PAPERLESS_DBPORT=5432", "PAPERLESS_DBNAME=paperless", "PAPERLESS_DBUSER=paperless", f"PAPERLESS_DBPASS={db_pw}", "POSTGRES_DB=paperless", "POSTGRES_USER=paperless", f"POSTGRES_PASSWORD={db_pw}", ]) env_content = "\\n".join(env_lines) + "\\n" # Write .env to automation1 env_tmp = f"/tmp/{hostname}.env" _ssh(AUTOMATION1, f"printf '%s' '{env_content}' > {env_tmp}", timeout=10) # Encrypt with SOPS enc_tmp = f"/tmp/{hostname}.env.enc" _ssh(AUTOMATION1, f"cd /tmp && SOPS_AGE_RECIPIENT={age_pub} sops --encrypted-regex 'PASSWORD|SECRET_KEY|_PASS' --encrypt {env_tmp} > {enc_tmp}", timeout=30) # Copy to VM _ssh(AUTOMATION1, f"scp -o StrictHostKeyChecking=no {env_tmp} sascha@{vm_ip}:/app-config/github/{hostname}/.env 2>/dev/null || true", timeout=15) _ssh(AUTOMATION1, f"scp -o StrictHostKeyChecking=no {enc_tmp} sascha@{vm_ip}:/app-config/github/{hostname}/.env.enc 2>/dev/null || true", timeout=15) log.info(f"SOPS .env setup completed for {hostname}") except Exception as e: log.warning(f"SOPS .env setup failed for {hostname}: {e}") return {"status": "ok" if rc == 0 else "error", "rc": rc, "output": out[-1000:]} @app.get("/ansible/status/{job_id}") async def ansible_status(job_id: int, _=Depends(_verify)): return {"info": "direct SSH mode - no async job tracking"} # --- TTS Endpoints --- class TTSRequest(BaseModel): text: str target: str = "speaker" # "speaker" or "telegram" voice: str = "deep_thought.mp3" language: str = "de" class TTSGenerateRequest(BaseModel): text: str = Field(min_length=1, max_length=2000) voice: str = Field(default="deep_thought.mp3", pattern=r"^[A-Za-z0-9_.-]+$") language: str = Field(default="de", pattern=r"^[A-Za-z]{2,8}(?:-[A-Za-z0-9]{2,8})?$") class TTSBridgeDeployRequest(BaseModel): rotate_client_token: bool = False SPEAKER_URL = TTS_CFG.get("speaker_url", "http://10.10.1.166:10800") if TTS_CFG else "http://10.10.1.166:10800" CHATTERBOX_URL = TTS_CFG.get("chatterbox_url", "http://10.2.1.104:8004/tts") if TTS_CFG else "http://10.2.1.104:8004/tts" def _chatterbox_payload(text: str, voice: str, language: str) -> dict: return { "text": text, "voice_mode": "clone", "reference_audio_filename": voice, "output_format": "wav", "language": language, "exaggeration": 0.3, "cfg_weight": 0.7, "temperature": 0.6, } @app.post("/tts/generate", response_class=Response) async def tts_generate(req: TTSGenerateRequest, _=Depends(_verify)): """Generate cloned speech and return the WAV bytes to the authenticated caller.""" async with httpx.AsyncClient(verify=False, timeout=180) as client: try: result = await client.post(CHATTERBOX_URL, json=_chatterbox_payload(req.text, req.voice, req.language)) except httpx.RequestError: _audit("/tts/generate", "POST", 502, "chatterbox request failed") raise HTTPException(status_code=502, detail="Chatterbox is unavailable") if result.status_code != 200: _audit("/tts/generate", "POST", 502, f"chatterbox_http={result.status_code}") raise HTTPException(status_code=502, detail="Chatterbox generation failed") if not result.content.startswith(b"RIFF"): _audit("/tts/generate", "POST", 502, "invalid audio response") raise HTTPException(status_code=502, detail="Chatterbox returned invalid audio") _audit("/tts/generate", "POST", 200, f"voice={req.voice} chars={len(req.text)}") return Response( content=result.content, media_type="audio/wav", headers={ "Content-Disposition": 'inline; filename="voiceclone.wav"', "Cache-Control": "no-store", "X-Content-Type-Options": "nosniff", }, ) @app.post("/tts/bridge/deploy") async def tts_bridge_deploy(req: TTSBridgeDeployRequest, _=Depends(_verify)): """Install host-local bridge secrets on automation1 without exposing them.""" client_token = _vault_cache.get("tts_bridge_client_token", "").strip() if req.rotate_client_token or not client_token: client_token = secrets.token_urlsafe(32) if not BUTLER_TOKEN: raise HTTPException(status_code=500, detail="Butler token is not configured") files = { "/app-config/tts-bridge/butler-token": BUTLER_TOKEN, "/app-config/tts-bridge/client-token": client_token, } installer = """import json, os, pathlib files = json.loads({files_json!r}) base = pathlib.Path('/app-config/tts-bridge') base.mkdir(parents=True, exist_ok=True) for filename, value in files.items(): path = pathlib.Path(filename) if path.is_dir(): path.rmdir() path.write_text(value) 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}))')}" rc, _out, err = _ssh("sascha@10.5.85.5", command, timeout=30) if rc != 0: _audit("/tts/bridge/deploy", "POST", 500, "secret installation failed") raise HTTPException(status_code=500, detail="Could not install bridge secrets") _audit("/tts/bridge/deploy", "POST", 200, f"rotated={req.rotate_client_token or not _vault_cache.get('tts_bridge_client_token')}") return { "status": "ready", "host": "automation1", "listen": "0.0.0.0:8099", "client_token": client_token, "rotated": req.rotate_client_token or not _vault_cache.get("tts_bridge_client_token"), } @app.post("/tts/speak") async def tts_speak(req: TTSRequest, _=Depends(_verify)): if req.target == "speaker": async with httpx.AsyncClient(verify=False, timeout=120) as c: r = await c.post(SPEAKER_URL, json={"text": req.text}) return {"status": "ok" if r.status_code == 200 else "error", "target": "speaker"} elif req.target == "telegram": # Generate WAV via Chatterbox, save to hermes VM as OGG for Telegram voice async with httpx.AsyncClient(verify=False, timeout=120) as c: r = await c.post(CHATTERBOX_URL, json={ "text": req.text, "voice_mode": "clone", "reference_audio_filename": req.voice, "output_format": "wav", "language": req.language, "exaggeration": 0.3, "cfg_weight": 0.7, "temperature": 0.6, }) if r.status_code != 200: return JSONResponse({"error": "chatterbox failed"}, status_code=500) # Save WAV and convert to OGG on hermes import tempfile wav_path = tempfile.mktemp(suffix=".wav") ogg_path = "/tmp/trulla_voice.ogg" with open(wav_path, "wb") as f: f.write(r.content) rc, _, _ = _ssh("sascha@10.4.1.100", f"rm -f {ogg_path}", timeout=10) # Copy WAV to hermes and convert _sp.run(["scp", "-o", "ConnectTimeout=5", wav_path, f"sascha@10.4.1.100:/tmp/trulla_voice.wav"], timeout=30) _ssh("sascha@10.4.1.100", f"ffmpeg -y -i /tmp/trulla_voice.wav -c:a libopus -b:a 64k {ogg_path} 2>/dev/null", timeout=30) os.unlink(wav_path) return {"status": "ok", "target": "telegram", "media_path": ogg_path, "hint": "Use MEDIA:/tmp/trulla_voice.ogg in response"} else: return JSONResponse({"error": f"unknown target: {req.target}"}, status_code=400) @app.get("/tts/voices") async def tts_voices(_=Depends(_verify)): async with httpx.AsyncClient(verify=False, timeout=10) as c: r = await c.get("http://10.2.1.104:8004/get_predefined_voices") return r.json() @app.get("/tts/health") async def tts_health(_=Depends(_verify)): results = {} async with httpx.AsyncClient(verify=False, timeout=5) as c: try: r = await c.get(SPEAKER_URL) results["speaker"] = r.json() except Exception as e: results["speaker"] = {"status": "offline", "error": str(e)} try: r = await c.get("http://10.2.1.104:8004/api/model-info") results["chatterbox"] = "ok" except Exception as e: results["chatterbox"] = {"status": "offline", "error": str(e)} return results class MediaHandoffPayload(BaseModel): action: Literal["start", "moved", "status", "fail"] category: str | None = Field(None, max_length=32) directory: str | None = Field(None, max_length=1000) release: str | None = Field(None, max_length=500) cleanName: str | None = Field(None, max_length=500) expectedFiles: int | None = Field(None, ge=1, le=100) jobId: str | None = Field(None, max_length=80) reason: str | None = Field(None, max_length=300) def _media_handoff_caller_allowed(request: Request) -> bool: auth = request.headers.get("authorization", "") if BUTLER_TOKEN and secrets.compare_digest(auth, f"Bearer {BUTLER_TOKEN}"): return True try: caller = ipaddress.ip_address(request.client.host if request.client else "") networks = [ ipaddress.ip_network(value.strip(), strict=False) for value in MEDIA_HANDOFF_ALLOWED_NETWORKS.split(",") if value.strip() ] except ValueError: return False return any(caller in network for network in networks) @app.post("/media/handoff") async def media_handoff(payload: MediaHandoffPayload, request: Request): """Narrow SABnzbd-to-n8n bridge; no generic unauthenticated proxy access.""" if not _media_handoff_caller_allowed(request): raise HTTPException(403, "Media handoff caller is not allowed") allowed_categories = {"serien4k", "serien", "video4k", "video"} if payload.action == "start": if payload.category not in allowed_categories: raise HTTPException(422, "Unsupported media category") expected_prefix = f"/usenet/complete/{payload.category}/" if not payload.directory or not payload.directory.startswith(expected_prefix): raise HTTPException(422, "Invalid media handoff directory") if not payload.release: raise HTTPException(422, "Release name is required") else: if not payload.jobId or not re.fullmatch(r"[a-z0-9-]{8,80}", payload.jobId): raise HTTPException(422, "Valid jobId is required") cfg = SERVICES.get("n8n") if not cfg or not cfg.get("url"): raise HTTPException(503, "n8n service is not configured") target = f"{cfg['url'].rstrip('/')}/webhook/media-handoff" body = payload.model_dump(exclude_none=True) if hasattr(payload, "model_dump") else payload.dict(exclude_none=True) try: async with httpx.AsyncClient(verify=False, timeout=15) as client: response = await client.post(target, json=body, headers={"Content-Type": "application/json"}) except httpx.HTTPError as exc: _audit("/media/handoff", "POST", 502, f"action={payload.action} error={type(exc).__name__}") raise HTTPException(502, "n8n media handoff is unavailable") from exc try: result = response.json() except Exception: result = {"ok": False, "error": "invalid_n8n_response"} _audit("/media/handoff", "POST", response.status_code, f"action={payload.action}") return JSONResponse(content=result, status_code=response.status_code) @app.api_route("/{service}/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"]) async def proxy(service: str, path: str, request: Request, _=Depends(_verify)): SKIP_SERVICES = {"vm", "inventory", "ansible", "debug", "tts", "status", "audit", "config"} if service in SKIP_SERVICES: raise HTTPException(404, f"Unknown service: {service}") cfg = SERVICES.get(service) if not cfg: raise HTTPException(404, f"Unknown service: {service}. Available: {list(SERVICES.keys())}") base_url = cfg.get("url") auth_type = cfg["auth"] headers = dict(request.headers) cookies = {} for h in ["host", "content-length", "transfer-encoding", "authorization"]: headers.pop(h, None) if auth_type == "apikey": headers["X-Api-Key"] = _get_key(cfg) or "" elif auth_type == "apikey_urlfile": url, key = _parse_url_key(cfg["key_file"]) base_url = url.rstrip("/") if url else "" headers["X-Api-Key"] = key or "" elif auth_type == "bearer": headers["Authorization"] = f"Bearer {_get_key(cfg)}" elif auth_type == "n8n": headers["X-N8N-API-KEY"] = _get_key(cfg) or "" elif auth_type == "proxmox": pv = _parse_kv("proxmox") headers["Authorization"] = f"PVEAPIToken={pv.get('tokenid', '')}={pv.get('secret', '')}" elif auth_type == "session": global _dockhand_cookie if not _dockhand_cookie: async with httpx.AsyncClient(verify=False) as c: await _dockhand_login(c) cookies = _dockhand_cookie or {} target = f"{base_url}/{path}" body = await request.body() timeout = float(cfg.get("timeout", 30)) async with httpx.AsyncClient(verify=False, timeout=timeout) as client: resp = await client.request(method=request.method, url=target, headers=headers, cookies=cookies, content=body, params=request.query_params) if auth_type == "session" and resp.status_code == 401: _dockhand_cookie = None await _dockhand_login(client) resp = await client.request(method=request.method, url=target, headers=headers, cookies=_dockhand_cookie or {}, content=body, params=request.query_params) try: data = _redact_response(resp.json(), cfg.get("redact_response_fields")) except Exception: data = resp.text _audit(f"/{service}/{path}", request.method, resp.status_code) return JSONResponse(content=data, status_code=resp.status_code)