homelab-butler/app.py

4304 lines
195 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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.4.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,10.5.85.12/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)}
# --- Media file integrity verification ---
MEDIA_VERIFY_PROBES = {
"arrapps": ("bazarrUHD", "sascha"),
"arr-chris": (None, "chris"),
"arr-chris-live": (None, "chris"),
}
class MediaVerifyRequest(BaseModel):
host: str = Field(..., max_length=64)
path: str = Field(..., max_length=1000)
expected_minutes: float | None = Field(None, gt=0, le=1440)
tolerance_pct: float = Field(20.0, gt=0, le=100)
@app.post("/media/verify")
async def media_verify(payload: MediaVerifyRequest, _=Depends(_verify)):
"""Probe a media file's real duration via ffprobe to catch truncated imports.
Runs `ffprobe` inside an existing container (bazarrUHD on arrapps, which already
mounts /data and ships ffmpeg) so no new service is required. Compares the
measured duration against an expected runtime (minutes) supplied by the caller
(e.g. Sonarr/Radarr's runtime field) within tolerance_pct.
"""
if not re.fullmatch(r"[A-Za-z0-9_.-]+", payload.host):
raise HTTPException(400, "Invalid host name")
if payload.host not in MEDIA_VERIFY_PROBES:
raise HTTPException(400, f"No ffprobe container configured for host {payload.host}")
if not re.fullmatch(r"/data/[^\x00]+\.(mkv|mp4|avi|m4v|ts)", payload.path):
raise HTTPException(400, "Path must be an absolute /data media file")
if ".." in payload.path or "'" in payload.path:
raise HTTPException(400, "Path contains unsafe characters")
container, _default_user = MEDIA_VERIFY_PROBES[payload.host]
if not container:
raise HTTPException(400, f"Host {payload.host} has no configured ffprobe container yet")
target = _find_inventory_host(payload.host)
if not target:
raise HTTPException(404, f"Host {payload.host} not found in inventory")
probe_cmd = (
f"docker exec {container} ffprobe -v error "
f"-show_entries format=duration -of default=noprint_wrappers=1:nokey=1 '{payload.path}'"
)
rc, out, err = await asyncio.to_thread(
_ssh, f'{target["user"]}@{target["ip"]}', f"sudo -n {probe_cmd} || {probe_cmd}", 30
)
if rc != 0:
_audit("/media/verify", "POST", 502, f"host={payload.host} path={payload.path}")
raise HTTPException(502, (err or out).strip()[:500] or "ffprobe failed")
raw_duration = out.strip()
try:
duration_seconds = float(raw_duration)
except ValueError:
_audit("/media/verify", "POST", 502, f"host={payload.host} unparsable duration")
raise HTTPException(502, "ffprobe returned no parsable duration; file is likely corrupt")
duration_minutes = duration_seconds / 60
result = {
"host": payload.host,
"path": payload.path,
"duration_seconds": round(duration_seconds, 1),
"duration_minutes": round(duration_minutes, 2),
"expected_minutes": payload.expected_minutes,
"verified": True,
"suspect": False,
}
if payload.expected_minutes:
deviation_pct = abs(duration_minutes - payload.expected_minutes) / payload.expected_minutes * 100
result["deviation_pct"] = round(deviation_pct, 1)
result["suspect"] = deviation_pct > payload.tolerance_pct
_audit(
"/media/verify", "POST", 200,
f"host={payload.host} dur={duration_minutes:.1f}m suspect={result['suspect']}",
)
return result
# --- 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
class DnsRrsetUpsertRequest(BaseModel):
record_type: Literal["A", "AAAA", "CNAME"] = "A"
value: str
ttl: int = Field(default=300, ge=60, le=86400)
comment: str = Field(default="Managed by Homelab Butler", max_length=200)
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"(?<![A-Za-z0-9._-])([A-Za-z0-9_-]{32,})(?![A-Za-z0-9._-])", raw)
candidates.extend(value for value in broad if value.isascii())
for line in raw.splitlines():
if "token" not in line.lower() and "api key" not in line.lower() and "api-key" not in line.lower():
continue
value = re.split(r"[:=]", line, maxsplit=1)[-1].strip().strip("`'\"")
if value.isascii() and re.fullmatch(r"[A-Za-z0-9._-]{24,}", value):
candidates.append(value)
return candidates[0] if len(set(candidates)) == 1 else None
def _get_hetzner_dns_token() -> str:
token = _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
class SaschaMediaEdgeOptimizeRequest(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": [], "emby_snippet": []}, "network": {}}
+inspect = run(["docker", "inspect", "caddy", "--format", "{{.State.Running}}"])
+result["caddy"]["container_running"] = inspect["rc"] == 0 and inspect["stdout"] == "true"
+adapt = run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"])
+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
+ collecting = False; depth = 0; selected = []
+ for line in lines:
+ if not collecting and re.match(r"^\\s*\\(emby_config\\)\\s*\\{", line): collecting = True
+ if collecting:
+ depth += line.count("{") - line.count("}")
+ if not re.search(r"(?i)(password|secret|token|private|basicauth|basic_auth|hash)", line): selected.append(line.strip())
+ else: selected.append("[REDACTED SENSITIVE DIRECTIVE]")
+ if depth == 0: break
+ result["caddy"]["emby_snippet"] = selected
+
+result["network"]["wg_media_active"] = subprocess.run(["systemctl", "is-active", "--quiet", "wg-quick@wg-media"]).returncode == 0
+result["network"]["wg_media_enabled"] = subprocess.run(["systemctl", "is-enabled", "--quiet", "wg-quick@wg-media"]).returncode == 0
+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 _sascha_media_edge_optimize_command() -> str:
script = '''import hashlib, json, os, re, shutil, subprocess, time
+from pathlib import Path
+path = Path("/app-config/caddy/Caddyfile")
+text = path.read_text(encoding="utf-8")
+old_pattern = re.compile(r"(?m)^(\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+)10\\.6\\.1\\.103:8096(\\s*)$")
+new_pattern = re.compile(r"(?m)^\\s*import\\s+emby_config\\s+tv\\.sascha-lutz\\.de\\s+10\\.11\\.13\\.3:8096\\s*$")
+old_count = len(old_pattern.findall(text)); new_count = len(new_pattern.findall(text))
+if old_count == 1:
+ text = old_pattern.sub(r"\\g<1>10.11.13.3:8096\\g<2>", text)
+elif not (old_count == 0 and new_count == 1):
+ raise RuntimeError("expected exactly one tv.sascha-lutz.de upstream")
+lines = text.splitlines(keepends=True)
+first = next((i for i, line in enumerate(lines) if line.strip() and not line.lstrip().startswith("#")), None)
+if first is None or lines[first].strip() != "{":
+ lines[0:0] = ["{\\n", " servers {\\n", " protocols h1 h2\\n", " }\\n", "}\\n", "\\n"]
+else:
+ depth = 0; global_end = None; servers_start = None; servers_end = None
+ for i in range(first, len(lines)):
+ stripped = lines[i].strip(); before = depth
+ if before == 1 and re.match(r"^servers(?:\\s+\\S+)?\\s*\\{$", stripped): servers_start = i
+ depth += lines[i].count("{") - lines[i].count("}")
+ if servers_start is not None and i > servers_start and depth == 1 and servers_end is None: servers_end = i
+ if i > first and depth == 0: global_end = i; break
+ if global_end is None: raise RuntimeError("unbalanced Caddy global options block")
+ if servers_start is None:
+ lines[global_end:global_end] = [" servers {\\n", " protocols h1 h2\\n", " }\\n"]
+ else:
+ if servers_end is None: raise RuntimeError("unbalanced Caddy servers block")
+ protocol_lines = [i for i in range(servers_start + 1, servers_end) if re.match(r"^\\s*protocols\\s+", lines[i])]
+ if len(protocol_lines) > 1: raise RuntimeError("multiple Caddy protocol directives")
+ if protocol_lines: lines[protocol_lines[0]] = re.sub(r"protocols\\s+.*", "protocols h1 h2", lines[protocol_lines[0]])
+ else: lines.insert(servers_start + 1, " protocols h1 h2\\n")
+text = "".join(lines)
+if re.search(r"(?im)^\\s*header(?:_down)?\\s+Alt-Svc", text): raise RuntimeError("manual Alt-Svc directive requires review")
+stamp = time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())
+backup = path.with_name("Caddyfile.pre-sascha-media-" + stamp)
+candidate = path.with_name("Caddyfile.sascha-media-candidate")
+shutil.copy2(path, backup); candidate.write_text(text, encoding="utf-8"); os.chmod(candidate, path.stat().st_mode)
+def run(args, timeout=30):
+ proc = subprocess.run(args, text=True, capture_output=True, timeout=timeout)
+ if proc.returncode != 0: raise RuntimeError((proc.stderr or proc.stdout).strip()[-500:] or "command failed")
+ return proc.stdout.strip()
+try:
+ run(["docker", "cp", str(candidate), "caddy:/tmp/Caddyfile.sascha-media-candidate"])
+ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/tmp/Caddyfile.sascha-media-candidate"])
+ with path.open("w", encoding="utf-8") as handle:
+ handle.write(text); handle.flush(); os.fsync(handle.fileno())
+ host_hash = hashlib.sha256(path.read_bytes()).hexdigest()
+ container_hash = run(["docker", "exec", "caddy", "sha256sum", "/etc/caddy/Caddyfile"]).split()[0]
+ if host_hash != container_hash: raise RuntimeError("host/container Caddyfile hash mismatch")
+ run(["docker", "exec", "caddy", "caddy", "validate", "--config", "/etc/caddy/Caddyfile"])
+ run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"])
+ adapted = json.loads(run(["docker", "exec", "caddy", "caddy", "adapt", "--config", "/etc/caddy/Caddyfile"]))
+ protocols = []
+ for server in adapted.get("apps", {}).get("http", {}).get("servers", {}).values(): protocols.extend(server.get("protocols", []))
+ if sorted(set(protocols)) != ["h1", "h2"]: raise RuntimeError("Caddy did not load h1+h2-only protocols")
+ if not run(["curl", "-fsS", "--max-time", "15", "http://10.11.13.3:8096/System/Ping"]): raise RuntimeError("Emby direct-path ping was empty")
+except Exception:
+ with path.open("w", encoding="utf-8") as handle:
+ handle.write(backup.read_text(encoding="utf-8")); handle.flush(); os.fsync(handle.fileno())
+ subprocess.run(["docker", "exec", "caddy", "caddy", "reload", "--config", "/etc/caddy/Caddyfile"], text=True, capture_output=True, timeout=30)
+ raise
+finally:
+ candidate.unlink(missing_ok=True)
+ subprocess.run(["docker", "exec", "caddy", "rm", "-f", "/tmp/Caddyfile.sascha-media-candidate"], text=True, capture_output=True)
+print(json.dumps({"status": "optimized", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "backup": str(backup), "sha256": host_hash}))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'sudo -n python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _optimize_sascha_media_edge() -> dict:
return _ssh_json(_media_host_target(SASCHA_MEDIA_VPS_HOST), _sascha_media_edge_optimize_command(), 90)
@app.post("/media/edge/sascha/optimize")
async def optimize_sascha_media_edge(req: SaschaMediaEdgeOptimizeRequest, _=Depends(_verify)):
"""Switch only tv.sascha-lutz.de to wg-media and disable HTTP/3 on its dedicated Hetzner edge."""
plan = {"status": "would_optimize", "hostname": "tv.sascha-lutz.de", "upstream": "10.11.13.3:8096", "protocols": ["h1", "h2"], "zero_downtime_reload": True}
if req.dry_run:
_audit("/media/edge/sascha/optimize", "POST", 200, "dry_run=True", True)
return plan
if req.confirmation != "OPTIMIZE_TV_SASCHA_LUTZ_DE":
raise HTTPException(400, "confirmation must be OPTIMIZE_TV_SASCHA_LUTZ_DE")
try:
result = await asyncio.to_thread(_optimize_sascha_media_edge)
except Exception as exc:
_audit("/media/edge/sascha/optimize", "POST", 502, "optimization failed; Caddy rollback attempted")
raise HTTPException(502, str(exc)[-500:]) from exc
_audit("/media/edge/sascha/optimize", "POST", 200, "direct upstream and h1+h2 enabled")
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 / "wgmtest.conf"
+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_tunnel_key_cleanup_command() -> str:
script = '''import json
+from pathlib import Path
+root = Path("/app-config/wireguard-media")
+if not (root / "wg-media.conf").exists():
+ (root / "private.key").unlink(missing_ok=True)
+ (root / "public.key").unlink(missing_ok=True)
+ status = "new_keys_removed"
+else:
+ status = "kept_for_existing_config"
+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
installed_targets = {target for target, _created in configured}
for target, key in ((vps, vps_key), (emby, emby_key)):
if target not in installed_targets and key.get("created"):
try: _ssh_json(target, _media_tunnel_key_cleanup_command(), 30)
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
def _media_benchmark_server_command() -> str:
server = '''import socket
+payload = b"\\0" * (1024 * 1024)
+with socket.socket() as listener:
+ listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
+ listener.bind(("0.0.0.0", 5209)); listener.listen(4); listener.settimeout(80)
+ for _ in range(2):
+ conn, _addr = listener.accept()
+ with conn:
+ conn.settimeout(30)
+ for _ in range(256): conn.sendall(payload)
+'''.replace("\n+", "\n")
encoded = base64.b64encode(server.encode()).decode()
command = f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
return (
"sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true; "
"sudo -n systemd-run --unit=butler-media-benchmark --collect --property=RuntimeMaxSec=90 "
f"/bin/sh -c {json.dumps(command)}"
)
def _media_benchmark_client_command() -> str:
script = '''import json, socket, time
+results = {}
+for name, host in (("legacy_node6", "10.6.1.103"), ("direct_wg_media", "10.11.13.3")):
+ total = 0; started = time.monotonic()
+ with socket.create_connection((host, 5209), timeout=10) as conn:
+ conn.settimeout(40)
+ while True:
+ chunk = conn.recv(1024 * 1024)
+ if not chunk: break
+ total += len(chunk)
+ elapsed = time.monotonic() - started
+ results[name] = {"bytes": total, "seconds": round(elapsed, 3), "mbit_s": round(total * 8 / elapsed / 1000000, 1)}
+print(json.dumps(results))
+'''.replace("\n+", "\n")
encoded = base64.b64encode(script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
def _benchmark_sascha_media_paths() -> dict:
vps = _media_host_target(SASCHA_MEDIA_VPS_HOST)
emby = _media_host_target(SASCHA_MEDIA_EMBY_HOST)
rc, out, err = _ssh(emby, _media_benchmark_server_command(), 30)
if rc != 0:
raise RuntimeError((err or out).strip()[-500:] or "failed to start benchmark server")
time.sleep(2)
try:
result = _ssh_json(vps, _media_benchmark_client_command(), 90)
finally:
_ssh(emby, "sudo -n systemctl stop butler-media-benchmark.service >/dev/null 2>&1 || true", 20)
if any(item.get("bytes") != 256 * 1024 * 1024 for item in result.values()):
raise RuntimeError("benchmark transferred an unexpected byte count")
return result
@app.post("/network/media-tunnel/sascha/benchmark")
async def benchmark_sascha_media_paths(_=Depends(_verify)):
"""Compare the legacy node6 route with the direct WireGuard media path using fixed transient TCP streams."""
try:
result = await asyncio.to_thread(_benchmark_sascha_media_paths)
except Exception as exc:
_audit("/network/media-tunnel/sascha/benchmark", "POST", 502, "benchmark failed")
raise HTTPException(502, str(exc)[-500:]) from exc
_audit("/network/media-tunnel/sascha/benchmark", "POST", 200, "fixed 256 MiB TCP comparison")
return {"status": "completed", "direction": "emby-sascha_to_hetzner", "results": result}
SYSCTL_AUDIT_KEYS = (
"net.core.default_qdisc",
"net.core.rmem_default",
"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.put("/dns/rrset/{zone_name}/{record_name}")
async def dns_rrset_upsert(zone_name: str, record_name: str, req: DnsRrsetUpsertRequest, _=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)]
record_type = req.record_type.upper()
value = req.value.strip().rstrip("." if record_type == "CNAME" else "")
try:
if record_type == "A" and ipaddress.ip_address(value).version != 4:
raise ValueError
if record_type == "AAAA" and ipaddress.ip_address(value).version != 6:
raise ValueError
if record_type == "CNAME":
_validate_managed_hostname(value)
except ValueError as exc:
raise HTTPException(400, f"Invalid {record_type} record value") from exc
zone, rrsets, headers = await _hetzner_zone_and_rrsets(zone_name)
existing = next((item for item in rrsets if item.get("name") == record_name and item.get("type") == record_type), None)
payload = {
"name": record_name,
"type": record_type,
"ttl": req.ttl,
"records": [{"value": value, "comment": req.comment}],
}
async with httpx.AsyncClient(timeout=30) as client:
if existing:
response = await client.put(
f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets/{record_name}/{record_type}",
headers=headers,
json=payload,
)
else:
response = await client.post(
f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets",
headers=headers,
json=payload,
)
if response.status_code not in {200, 201}:
raise HTTPException(502, f"Hetzner RRSet upsert failed: HTTP {response.status_code}")
verify = await client.get(f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets", headers=headers)
current = [item for item in verify.json().get("rrsets", []) if item.get("name") == record_name and item.get("type") == record_type] if verify.status_code == 200 else []
if not current or not any(record.get("value") == value for record in current[0].get("records", [])):
raise HTTPException(502, "RRSet read-back does not contain requested value")
_audit(f"/dns/rrset/{zone_name}/{record_name}", "PUT", 200, f"type={record_type} value={value}")
return {"zone": zone_name, "zone_id": zone.get("id"), "name": record_name, "created": existing is None, "rrset": current[0]}
@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
def _sab_history_command() -> str:
container_script = r'''import json, re, subprocess, urllib.parse, urllib.request
text = open('/config/sabnzbd.ini', encoding='utf-8', errors='replace').read()
match = re.search(r'^api_key\s*=\s*(\S+)', text, re.M)
port_match = re.search(r'^port\s*=\s*(\d+)', text, re.M)
api_key = match.group(1) if match else ''
port = port_match.group(1) if port_match else '7777'
params = urllib.parse.urlencode({'mode': 'history', 'limit': 100, 'output': 'json', 'apikey': api_key})
with urllib.request.urlopen('http://127.0.0.1:' + port + '/api?' + params, timeout=20) as response:
data = json.load(response)
slots = data.get('history', {}).get('slots', [])
allowed = ('nzo_id', 'name', 'category', 'status', 'script', 'script_line', 'fail_message', 'completed', 'storage', 'path')
print(json.dumps([{key: item.get(key) for key in allowed} for item in slots]))
'''
container_encoded = base64.b64encode(container_script.encode()).decode()
host_script = f'''import subprocess, sys
command = ["sudo", "-n", "docker", "exec", "sabnzbd", "python3", "-c", "import base64;exec(base64.b64decode('{container_encoded}'))"]
proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
if proc.returncode != 0:
command = command[2:]
proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
if proc.returncode != 0:
print((proc.stderr or proc.stdout)[-500:], file=sys.stderr)
raise SystemExit(proc.returncode)
print(proc.stdout)
'''
encoded = base64.b64encode(host_script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
@app.get("/media/handoff/sab-history")
async def media_handoff_sab_history(_=Depends(_verify)):
"""Return sanitized SAB history without exposing the SAB API key."""
inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
if not inventory:
raise HTTPException(404, "Host sabnzbd not found")
target = f'{inventory["user"]}@{inventory["ip"]}'
rc, out, err = await asyncio.to_thread(_ssh, target, _sab_history_command(), 40)
if rc != 0:
raise HTTPException(502, (err or out).strip()[-500:] or "SAB history failed")
try:
return {"history": json.loads(out)}
except json.JSONDecodeError as exc:
raise HTTPException(502, "invalid SAB history response") from exc
def _media_handoff_diagnostics_command() -> str:
script = r'''import hashlib, json, os, subprocess, urllib.error, urllib.request
def run(args):
proc = subprocess.run(args, capture_output=True, text=True, timeout=20)
return proc.returncode, proc.stdout.strip(), proc.stderr.strip()
def docker_exec(command):
for prefix in (["sudo", "-n", "docker"], ["docker"]):
rc, out, err = run(prefix + ["exec", "sabnzbd", "sh", "-lc", command])
if rc == 0 or "not found" not in (err + out).lower():
return rc, out, err
return rc, out, err
inspect_rc, inspect_out, _ = run(["sudo", "-n", "docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
if inspect_rc != 0:
inspect_rc, inspect_out, _ = run(["docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
py_rc, py_out, _ = docker_exec("command -v python3")
curl_rc, curl_out, _ = docker_exec("command -v curl")
stat_rc, stat_out, _ = docker_exec("test -f /usenet/scripts/movetdarr.sh && stat -c '%a %s' /usenet/scripts/movetdarr.sh && sha256sum /usenet/scripts/movetdarr.sh")
log_rc, log_out, _ = docker_exec("tail -n 200 /usenet/scripts/postprocess.log 2>/dev/null || true")
source_rc, source_out, _ = docker_exec("find /usenet/complete -mindepth 2 -maxdepth 2 -type d ! -name '_UNPACK_*' -print 2>/dev/null | sort | tail -100")
target_rc, target_out, _ = docker_exec("find /tdarr/complete -mindepth 2 -maxdepth 2 -type d -print 2>/dev/null | sort | tail -100")
probe_code = 0
probe_body = ""
probe = urllib.request.Request(
"http://10.5.85.2:8888/media/handoff",
method="POST",
headers={"Content-Type": "application/json"},
data=b'{"action":"status","jobId":"diagnostic-probe"}',
)
try:
with urllib.request.urlopen(probe, timeout=10) as response:
probe_code = response.status
probe_body = response.read(500).decode(errors="replace")
except urllib.error.HTTPError as exc:
probe_code = exc.code
probe_body = exc.read(500).decode(errors="replace")
except Exception as exc:
probe_body = type(exc).__name__
script_info = {"exists": stat_rc == 0, "executable": False, "sha256": None, "mode": None, "size": None}
if stat_rc == 0:
lines = stat_out.splitlines()
if lines:
parts = lines[0].split()
if len(parts) >= 2:
script_info.update(mode=parts[0], size=int(parts[1]), executable=any(ch in parts[0][-3:] for ch in "1357"))
if len(lines) > 1:
script_info["sha256"] = lines[1].split()[0]
print(json.dumps({
"container_running": inspect_rc == 0 and inspect_out == "true",
"tools": {"python3": py_rc == 0 and bool(py_out), "curl": curl_rc == 0 and bool(curl_out)},
"script": script_info,
"caller_probe": {"http_status": probe_code, "body": probe_body},
"source_directories": source_out.splitlines() if source_rc == 0 else [],
"target_directories": target_out.splitlines() if target_rc == 0 else [],
"recent_log": log_out.splitlines()[-200:],
}))
'''
encoded = base64.b64encode(script.encode()).decode()
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
@app.get("/media/handoff/diagnostics")
async def media_handoff_diagnostics(_=Depends(_verify)):
"""Read-only diagnostics for the fixed SABnzbd handoff script and caller path."""
inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
if not inventory:
raise HTTPException(404, "Host sabnzbd not found")
target = f'{inventory["user"]}@{inventory["ip"]}'
rc, out, err = await asyncio.to_thread(_ssh, target, _media_handoff_diagnostics_command(), 30)
if rc != 0:
raise HTTPException(502, (err or out).strip()[-500:] or "media handoff diagnostics failed")
try:
return json.loads(out)
except json.JSONDecodeError as exc:
raise HTTPException(502, "invalid media handoff diagnostic response") from exc
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):
caller = request.client.host if request.client else "unknown"
_audit("/media/handoff", "POST", 403, f"caller={caller} action={payload.action}")
raise HTTPException(403, "Media handoff caller is not allowed")
# 16.09.2026: serienen/videoen ergaenzt. Die englischen Arr-Instanzen
# sonarrEN (Port 8991, Root /data/FHD/serienen) und radarrEN (Port 7880,
# Root /data/FHD/videoen) nutzen eigene SAB-Kategorien. Ohne sie brach der
# Handoff mit 422 ab und Releases blieben in /usenet/complete liegen.
allowed_categories = {"serien4k", "serien", "serienen", "video4k", "video", "videoen"}
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)