import json import os import time import urllib.error import urllib.request from datetime import datetime, timezone from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path # Shared liveness logic — the SAME module the observer and operator_ui use, so # all three agree on what "dead" means. This image only copies web.py + # index.html, so liveness.py is bind-mounted in at /app/liveness.py (see this # service's docker-compose.yml). Imported defensively: if the mount is missing # we fall back to trusting the observer-written status rather than crashing the # whole panel — the observer remains authoritative either way. try: from liveness import compute_liveness, ttls_for, node_health as _liveness_health, degrade_for_node _LIVENESS_OK = True except Exception: # pragma: no cover - only when liveness.py mount is absent _LIVENESS_OK = False STATE_DIR = Path(os.getenv("HOMELAB_STATE_ROOT", "/opt/homelab/state")) EVENTS_DIR = Path(os.getenv("HOMELAB_EVENTS_ROOT", "/opt/homelab/events")) WORLD_DIR = Path(os.getenv("HOMELAB_WORLD_ROOT", "/opt/homelab/world")) ACTIONS_DIR = Path(os.getenv("HOMELAB_ACTIONS_ROOT", "/opt/homelab/actions")) CONFIG_DIR = Path(os.getenv("HOMELAB_CONFIG_ROOT", "/opt/homelab/config")) # When set, this panel is a read-only mirror of a control-plane running # elsewhere (VPS): actions are mirrored into world/actions.json by the # runtime-materializer (same as nodes/services), and mutations (approve/ # reject) are proxied there instead of touching the local ACTIONS_DIR, which # the real executor never reads. Same variable the materializer uses to pick # its mode — keeps both components' "am I a mirror?" decision in sync. CONTROL_PLANE_URL = os.environ.get("CONTROL_PLANE_URL", "").rstrip("/") STATIC_DIR = Path(__file__).parent DEFAULT_CONFIG = { "operator_mode": "approval", "auto_mode": True, "action_thresholds": { "restart_ha": 0.8, "check_network": 0.9, }, "default_threshold": 0.9, "allowed_auto_actions": ["restart_ha"], } def read_json_file(path, default=None): if not path.exists(): return default if default is not None else [] try: return json.loads(path.read_text()) except Exception: return default if default is not None else [] def get_config(): config_path = STATE_DIR / "operator-config.json" if config_path.exists(): return read_json_file(config_path, DEFAULT_CONFIG) return DEFAULT_CONFIG def save_config(config): STATE_DIR.mkdir(parents=True, exist_ok=True) (STATE_DIR / "operator-config.json").write_text(json.dumps(config, indent=2)) def _node_liveness_map(): """Return {node_name: liveness_tier} from nodes.json (empty if unavailable).""" if not _LIVENESS_OK: return {} raw = read_json_file(WORLD_DIR / "nodes.json", default={}) if not isinstance(raw, dict): return {} return { name: compute_liveness(info.get("last_seen"), ttls=ttls_for(name, info.get("roles", []))) for name, info in raw.items() } def current_nodes(): """Nodes shaped as a list with a freshness-aware health field. nodes.json is a dict keyed by name; index.html does nodes.map() and reads node.health, so we shape to a list and compute health here. The health is the read-time safety net (worse of persisted status and last_seen freshness) so a stalled observer cannot make a dead node look nominal. """ raw = read_json_file(WORLD_DIR / "nodes.json", default={}) if isinstance(raw, list): return raw result = [] for name, info in raw.items(): health = _liveness_health(info, name=name) if _LIVENESS_OK else ( "error" if info.get("status") == "offline" else "nominal" if info.get("status") == "online" else info.get("status", "unknown") ) result.append({ "id": name, "hostname": name, "health": health, "status": info.get("status", "unknown"), "capabilities": info.get("roles", []), "connectivity": "tailscale", "incidents": 0, "last_seen": info.get("last_seen"), "disk_usage_pct": info.get("disk_usage_pct"), "mem_usage_pct": info.get("mem_usage_pct"), "cpu_usage_pct": info.get("cpu_usage_pct"), "disk_pressure": info.get("disk_pressure"), }) return result def current_services(): """Services shaped as a list, with the node-liveness cascade (variant B).""" raw = read_json_file(WORLD_DIR / "services.json", default={}) if isinstance(raw, list): return raw node_liveness = _node_liveness_map() result = [] for key, info in raw.items(): svc_status = info.get("status", "unknown") health = ("nominal" if svc_status == "healthy" else ("error" if svc_status == "unhealthy" else svc_status)) node_name = info.get("node", "") if _LIVENESS_OK and node_name in node_liveness: health = degrade_for_node(health, node_liveness[node_name]) result.append({ "id": key, "name": info.get("service", key), "node": node_name, "health": health, "actual_state": svc_status, "last_check": info.get("last_check"), "incident_id": info.get("incident_id"), }) return result def current_deployments(): return read_json_file(WORLD_DIR / "deployments.json") def current_incidents(): return read_json_file(WORLD_DIR / "incidents.json") def current_recommendations(): return read_json_file(WORLD_DIR / "recommendations.json") def current_summary(): path = WORLD_DIR / "runtime-summary.json" summary = read_json_file(path, default={}) if summary: last_update_val = summary.get("last_update") if last_update_val: try: if isinstance(last_update_val, str): last_update = datetime.fromisoformat(last_update_val.replace('Z', '+00:00')).timestamp() else: last_update = float(last_update_val) except Exception: last_update = os.path.getmtime(path) else: last_update = os.path.getmtime(path) summary["last_update"] = last_update summary["stale"] = (time.time() - last_update) > 60 return summary def current_events(): return read_json_file(WORLD_DIR / "events.json", default=[]) def current_actions(): """Actions shown to the operator. In mirror mode (CONTROL_PLANE_URL set) VPS is the single source of truth for actions, same as nodes/services: the materializer fetches VPS's /actions and writes world/actions.json, and we just read it back here. The local ACTIONS_DIR scan is only for standalone/dev use where this webui runs next to its own supervisor+executor (no CONTROL_PLANE_URL). """ if CONTROL_PLANE_URL: mirrored = read_json_file(WORLD_DIR / "actions.json", default={}) if isinstance(mirrored, dict): return mirrored actions = {} statuses = ["pending", "approved", "running", "completed", "failed", "rejected"] for status in statuses: actions[status] = [] status_dir = ACTIONS_DIR / status if status_dir.exists(): for f in status_dir.glob("*.json"): data = read_json_file(f) if data: # Injects some metadata for UI data["id"] = data.get("action_id") or f.stem data["status"] = status actions[status].append(data) return actions def _proxy_mutate(action_id, target_status): """Forward an approve/reject to the VPS control-plane API. The executor only ever polls ACTIONS_DIR/approved on VPS, so a mutation applied to this panel's local (mirrored, otherwise-unused) ACTIONS_DIR would silently vanish — the approval would look successful here but the executor would never see it. Proxying is the only way an approval made on this mirror panel actually reaches the executor. """ url = f"{CONTROL_PLANE_URL}/action/mutate" body = json.dumps({"id": action_id, "status": target_status}).encode("utf-8") req = urllib.request.Request(url, data=body, method="POST") req.add_header("Content-Type", "application/json") try: with urllib.request.urlopen(req, timeout=10) as resp: if resp.status == 200: return True, "Success" return False, f"Control plane returned HTTP {resp.status}" except urllib.error.HTTPError as e: return False, f"Control plane returned HTTP {e.code}: {e.reason}" except Exception as e: return False, f"Failed to reach control plane at {url}: {e}" def mutate_action(action_id, target_status): statuses = ["pending", "approved", "running", "completed", "failed", "rejected"] if target_status not in statuses: return False, f"Invalid target status: {target_status}" if CONTROL_PLANE_URL: return _proxy_mutate(action_id, target_status) # Find where the action is source_path = None current_status = None for status in statuses: p = ACTIONS_DIR / status / f"{action_id}.json" if p.exists(): source_path = p current_status = status break if not source_path: return False, f"Action {action_id} not found" target_dir = ACTIONS_DIR / target_status target_dir.mkdir(parents=True, exist_ok=True) target_path = target_dir / f"{action_id}.json" try: data = json.loads(source_path.read_text()) data["status"] = target_status data["updated_at"] = time.time() # Keep history of transitions history = data.get("transition_history", []) history.append({ "from": current_status, "to": target_status, "timestamp": time.time() }) data["transition_history"] = history target_path.write_text(json.dumps(data, indent=2)) if source_path != target_path: source_path.unlink() return True, "Success" except Exception as e: return False, str(e) def get_snapshot(): nodes = current_nodes() services = current_services() incidents = current_incidents() events = current_events() summary = current_summary() non_nominal = [s for s in services if s.get("health") != "nominal"] nominal_count = len(services) - len(non_nominal) return { "timestamp": datetime.now(timezone.utc).isoformat(), "summary": summary, "nodes": nodes, "non_nominal_services": non_nominal, "nominal_service_count": nominal_count, "total_service_count": len(services), "incidents": incidents, "events": events[:10], } def send_json(status, payload, handler): body = (json.dumps(payload) + "\n").encode("utf-8") handler.send_response(status) handler.send_header("Content-Type", "application/json") handler.send_header("Content-Length", str(len(body))) handler.end_headers() handler.wfile.write(body) class Handler(BaseHTTPRequestHandler): def do_GET(self): if self.path == "/config": send_json(200, get_config(), self) return if self.path == "/nodes": send_json(200, current_nodes(), self) return if self.path == "/services": send_json(200, current_services(), self) return if self.path == "/deployments": send_json(200, current_deployments(), self) return if self.path == "/incidents": send_json(200, current_incidents(), self) return if self.path == "/recommendations": send_json(200, current_recommendations(), self) return if self.path == "/summary": send_json(200, current_summary(), self) return if self.path == "/events": send_json(200, current_events(), self) return if self.path == "/actions": send_json(200, current_actions(), self) return if self.path == "/snapshot": send_json(200, get_snapshot(), self) return if self.path in ("/", "/index.html"): body = (STATIC_DIR / "index.html").read_bytes() self.send_response(200) self.send_header("Content-Type", "text/html; charset=utf-8") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) return self.send_error(404) def do_POST(self): if self.path not in ( "/config", "/action/mutate", "/mode", ): self.send_error(404) return length = int(self.headers.get("Content-Length", "0")) raw_body = self.rfile.read(length).decode("utf-8") try: payload = json.loads(raw_body) except json.JSONDecodeError: self.send_error(400, "Invalid JSON") return if self.path == "/config": config = get_config() config.update(payload) save_config(config) send_json(200, {"status": "ok"}, self) return if self.path == "/mode": mode = payload.get("mode") if not mode: self.send_error(400, "mode is required") return config = get_config() config["operator_mode"] = mode save_config(config) send_json(200, {"status": "ok"}, self) return if self.path == "/action/mutate": action_id = payload.get("id") target = payload.get("status") if not action_id or not target: self.send_error(400, "id and status are required") return success, msg = mutate_action(action_id, target) if success: send_json(200, {"status": "ok"}, self) else: self.send_error(500, msg) return def log_message(self, format, *args): return if __name__ == "__main__": # Ensure directories exist for d in [STATE_DIR, EVENTS_DIR, WORLD_DIR, ACTIONS_DIR, CONFIG_DIR]: d.mkdir(parents=True, exist_ok=True) for s in ["pending", "approved", "running", "completed", "failed", "rejected"]: (ACTIONS_DIR / s).mkdir(parents=True, exist_ok=True) port = int(os.getenv("PORT", "8080")) print(f"Operator Control Plane starting on 0.0.0.0:{port}") server = ThreadingHTTPServer(("0.0.0.0", port), Handler) server.serve_forever()