from __future__ import annotations import logging import os import shlex import shutil import subprocess import sys import time from pathlib import Path from typing import TYPE_CHECKING from knoe.milestone import Milestone from knoe import config as inst_config from knoe import knoe_conf from knoe.core.stream_exec import run_streaming_cmd from knoe.core.env import resolve_knoe_home from knoe.core.policy import ( POLICY_CFG_KEY, OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES, evaluate_optional_workloads_allowed, ) from knoe.core.ops.cloudnative_pg import ( initialize as cnpg_initialize, deploy as cnpg_deploy, ) from knoe.core.ops import openbao as openbao_ops from knoe.core.ops import registry as registry_ops from knoe.core.ops import garage_store as garage_store_ops from knoe.core.ops import opentofu as opentofu_ops from knoe.core.ops import monitoring as monitoring_ops if TYPE_CHECKING: from knoe.state import InstallerState from knoe.milestone import ProgressCallback def _stream_line(line: str) -> None: try: sys.stdout.write(line) sys.stdout.flush() except Exception: pass class DependenciesMilestone(Milestone): def __init__(self): super().__init__("dependencies", "Dependency Verification") self.logger = logging.getLogger("DependenciesMilestone") self._last_gke_context_failure_reason = "" def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: platform_info = inst_config.detect_dependency_platform() verify_all = self._parse_bool(state.inputs.get("dependencies.verify_all", "False")) self.logger.info( "Dependency platform detected: %s", inst_config.get_dependency_platform_summary(), ) if progress: progress("Checking dependencies...", 0.1) if verify_all: dependencies = inst_config.get_platform_dependencies() else: dependencies = inst_config.get_required_dependencies(state.inputs) required_ids = inst_config.get_required_dependency_ids(state.inputs) backend_ok, backend_reason = inst_config._validate_dependency_backend(dependencies) if not backend_ok: self._fatal_dependency_failure( state, f"Dependency backend unsupported: {backend_reason}", ) if not dependencies: reason = platform_info.get("reason") or "No dependencies configured for detected platform" self._fatal_dependency_failure(state, f"Dependency resolution failed: {reason}") missing: list[dict] = [] for i, dep in enumerate(dependencies): ok, _, version = inst_config.get_dep_info(dep) if ok: self.logger.info(f"[OK] {dep['name']} {version or ''}".strip()) else: missing.append(dep) self.logger.info(f"[MISSING] {dep['name']}") if progress: progress( f"Checking {dep['name']}...", 0.1 + (i / len(dependencies)) * 0.4 ) if not missing: self._finalize_dependency_prerequisites(state, progress) return # Check if auto-install is enabled auto_install = self._parse_bool( state.inputs.get("dependencies.auto_install_missing", "True") ) if not auto_install: self._fatal_missing_dependencies( state, missing, required_ids, "Dependencies missing and auto-install disabled", ) for i, dep in enumerate(missing): install_cmd = dep.get("install_cmd") if not install_cmd: self.logger.error(f"No install command for {dep['name']}.") continue should_install = self._parse_bool( state.inputs.get(f"dependencies.{dep['id']}.install", "True") ) if not should_install: self.logger.info(f"[SKIP] {dep['name']} install disabled by config.") continue self.logger.info(f"[INSTALL] {dep['name']} -> {install_cmd}") if progress: progress(f"Installing {dep['name']}...", 0.5 + (i / len(missing)) * 0.4) rc = self._run_cmd(install_cmd) if rc != 0: self.logger.error( "Install failed for %s (code %s) [backend=%s, cmd=%s]", dep["name"], rc, platform_info.get("package_manager") or "unknown", install_cmd, ) # Re-verify still_missing: list[dict] = [] for dep in missing: ok, _, _ = inst_config.get_dep_info(dep) if not ok: still_missing.append(dep) if still_missing: self._fatal_missing_dependencies( state, still_missing, required_ids, "Dependencies unresolved after install attempts", ) else: self._finalize_dependency_prerequisites(state, progress) def _finalize_dependency_prerequisites( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if not self._ensure_gcloud_auth_for_k8s(state): self._fatal_dependency_failure( state, "GKE bootstrap failed: gcloud authentication is required before cluster operations.", ) if not self._ensure_gke_auth_plugin_for_k8s(state): self._fatal_dependency_failure( state, "GKE bootstrap failed: gke-gcloud-auth-plugin is required before acquiring cluster contexts.", ) if not self._ensure_gke_contexts_for_k8s(state): context_reason = str(getattr(self, "_last_gke_context_failure_reason", "") or "").strip() message = "GKE bootstrap failed: required kube contexts are missing or could not be acquired." if context_reason: message = f"GKE bootstrap failed: {context_reason}" self._fatal_dependency_failure( state, message, ) self._set_status(state, "All installed") if progress: progress("All dependencies installed", 1.0) def _fatal_dependency_failure(self, state: InstallerState, message: str) -> None: self.logger.error(message) self._set_status(state, "Missing") raise RuntimeError(message) def _fatal_missing_dependencies( self, state: InstallerState, missing: list[dict], required_ids: set[str], reason: str, ) -> None: required_missing = [ str(dep.get("name") or dep.get("id")) for dep in missing if str(dep.get("id")) in required_ids ] optional_missing = [ str(dep.get("name") or dep.get("id")) for dep in missing if str(dep.get("id")) not in required_ids ] if optional_missing: self.logger.warning("Optional dependencies still missing: %s", ", ".join(optional_missing)) if required_missing: self._fatal_dependency_failure( state, f"{reason}. Required dependencies still missing: {', '.join(required_missing)}", ) def _is_gke_mode(self, state: InstallerState) -> bool: return "gcloud" in inst_config.get_required_dependency_ids(state.inputs) def _ensure_gke_auth_plugin_for_k8s(self, state: InstallerState) -> bool: if not self._is_gke_mode(state): return True return bool(shutil.which("gke-gcloud-auth-plugin")) def _gke_context_targets(self, state: InstallerState) -> list[dict]: def _first_non_empty(*keys: str) -> str: for key in keys: value = str(state.inputs.get(key, "")).strip() if value: return value return "" project = _first_non_empty("init_cluster.project_id") targets: list[dict] = [] for kind in ("app", "db"): context = _first_non_empty( f"init_cluster.{kind}_cluster_kubecontext", f"env_setup.{kind.upper()}_CLUSTER_KUBECONTEXT", ) cluster = _first_non_empty( f"init_cluster.{kind}_cluster_name", f"env_setup.{kind.upper()}_CLUSTER_NAME", ) region = _first_non_empty(f"init_cluster.{kind}_cluster_region") zone = _first_non_empty(f"init_cluster.{kind}_cluster_zone") if not context: continue targets.append( { "kind": kind, "context": context, "cluster": cluster, "region": region, "zone": zone, "project": project, } ) deduped: list[dict] = [] seen: set[str] = set() for target in targets: context = str(target.get("context") or "") if context and context not in seen: deduped.append(target) seen.add(context) return deduped def _resolve_gke_kubeconfig_target(self) -> tuple[str | None, str | None]: preferred = (Path.home() / ".kube" / "config").expanduser() raw_kubeconfig = str(os.environ.get("KUBECONFIG", "")).strip() candidate_raw = "" if raw_kubeconfig: first_entry = raw_kubeconfig.split(os.pathsep)[0].strip() if first_entry: candidate_raw = first_entry if candidate_raw: candidate = Path(candidate_raw).expanduser() if str(candidate) == "/etc/rancher/k3s/k3s.yaml": self.logger.warning( "Ignoring KUBECONFIG='%s' for GKE bootstrap (local k3s config); using '%s' instead.", candidate, preferred, ) elif self._is_kubeconfig_target_writable(candidate): return str(candidate), None else: self.logger.warning( "Ignoring non-writable KUBECONFIG='%s' for GKE bootstrap; using '%s' instead.", candidate, preferred, ) try: preferred.parent.mkdir(parents=True, exist_ok=True) except Exception as exc: return None, ( f"kubeconfig target '{preferred}' is not writable " f"(failed to create parent directory: {exc})" ) if not self._is_kubeconfig_target_writable(preferred): return None, ( f"kubeconfig target '{preferred}' is not writable. " "Set KUBECONFIG to a user-writable file and retry." ) return str(preferred), None def _is_kubeconfig_target_writable(self, target: Path) -> bool: try: if target.exists(): return target.is_file() and os.access(target, os.W_OK) return target.parent.exists() and target.parent.is_dir() and os.access(target.parent, os.W_OK) except Exception: return False def _get_kube_contexts(self, kubeconfig_target: str | None = None) -> set[str] | None: kube_env = inst_config._augment_env_for_dependency_backend(os.environ.copy()) if kubeconfig_target: kube_env["KUBECONFIG"] = kubeconfig_target try: res = subprocess.run( ["kubectl", "config", "get-contexts", "-o", "name"], capture_output=True, text=True, env=kube_env, timeout=20, ) except Exception as exc: self.logger.error("Failed to query kube contexts: %s", exc) return None if res.returncode != 0: self.logger.error( "Failed to list kube contexts (code %s): %s", res.returncode, (res.stderr or "").strip(), ) return None return { line.strip() for line in (res.stdout or "").splitlines() if line.strip() } def _ensure_gke_contexts_for_k8s(self, state: InstallerState) -> bool: self._last_gke_context_failure_reason = "" if not self._is_gke_mode(state): return True targets = self._gke_context_targets(state) if not targets: return True kubeconfig_target, kubeconfig_error = self._resolve_gke_kubeconfig_target() if kubeconfig_error: self._last_gke_context_failure_reason = kubeconfig_error self.logger.error( "Unable to select writable kubeconfig target for GKE bootstrap: %s", kubeconfig_error, ) return False self.logger.info("Using kubeconfig target for GKE bootstrap: %s", kubeconfig_target) kube_contexts = self._get_kube_contexts(kubeconfig_target) if kube_contexts is None: self._last_gke_context_failure_reason = "failed to list kube contexts for the selected kubeconfig target" return False missing_targets = [ target for target in targets if str(target.get("context")) not in kube_contexts ] if not missing_targets: return True gcloud_bin = shutil.which("gcloud") if not gcloud_bin: self._last_gke_context_failure_reason = "gcloud is missing while required GKE contexts are absent" self.logger.error( "Missing gcloud; cannot acquire required GKE contexts: %s", ", ".join(str(t.get("context") or "") for t in missing_targets), ) return False for target in missing_targets: context = str(target.get("context") or "") cluster = str(target.get("cluster") or "") if not cluster: self._last_gke_context_failure_reason = ( f"cannot acquire missing context '{context}': cluster name is not configured" ) self.logger.error( "Cannot acquire missing context '%s': cluster name is not configured.", context, ) return False cmd = [gcloud_bin, "container", "clusters", "get-credentials", cluster] region = str(target.get("region") or "") zone = str(target.get("zone") or "") project = str(target.get("project") or "") if region: cmd += ["--region", region] elif zone: cmd += ["--zone", zone] if project: cmd += ["--project", project] self.logger.info( "[ACTION] Acquiring GKE context '%s' via: %s [kubeconfig=%s]", context, " ".join(shlex.quote(part) for part in cmd), kubeconfig_target, ) rc = self._run_cmd(cmd, env={"KUBECONFIG": str(kubeconfig_target)}) if rc != 0: self._last_gke_context_failure_reason = ( f"failed to acquire GKE context '{context}' for cluster '{cluster}'" ) self.logger.error( "Failed to acquire GKE context '%s' for cluster '%s' (code %s)", context, cluster, rc, ) return False refreshed_contexts = self._get_kube_contexts(kubeconfig_target) if refreshed_contexts is None: self._last_gke_context_failure_reason = "failed to refresh kube contexts after GKE credential acquisition" return False still_missing = [ str(target.get("context") or "") for target in targets if str(target.get("context") or "") not in refreshed_contexts ] if still_missing: self._last_gke_context_failure_reason = ( "required kube contexts are still missing after credential acquisition: " + ", ".join(still_missing) ) self.logger.error( "GKE contexts still missing after credential acquisition: %s", ", ".join(still_missing), ) return False return True def _ensure_gcloud_auth_for_k8s(self, state: InstallerState) -> bool: if not self._is_gke_mode(state): return True gcloud_env = inst_config._augment_env_for_dependency_backend(os.environ.copy()) cmd_base = ["gcloud"] token_cmd = cmd_base + ["auth", "print-access-token", "--quiet"] try: token_res = subprocess.run( token_cmd, capture_output=True, text=True, env=gcloud_env, timeout=20, ) if token_res.returncode == 0 and (token_res.stdout or "").strip(): return True except Exception: pass interactive = bool(getattr(sys.stdin, "isatty", lambda: False)()) if interactive: login_cmd = cmd_base + ["auth", "login", "--no-launch-browser"] self.logger.info( "[ACTION] No active gcloud session found. Starting login flow: %s", " ".join(shlex.quote(part) for part in login_cmd), ) try: login_rc = self._run_cmd(login_cmd) if login_rc == 0: token_res = subprocess.run( token_cmd, capture_output=True, text=True, env=gcloud_env, timeout=20, ) if token_res.returncode == 0 and (token_res.stdout or "").strip(): return True except Exception: pass self.logger.error( "gcloud is installed but no active auth session is available for GKE. " "Run `gcloud auth login --no-launch-browser` in a terminal, complete the URL/code flow on a browser-capable machine, then retry deploy." ) return False def _set_status(self, state: InstallerState, status: str): sections = ["Dependencies", inst_config.get_dependency_milestone_title()] for section in sections: if section not in state.config_data: state.config_data[section] = {} state.config_data[section]["STATUS"] = status class NetworkScanMilestone(Milestone): def __init__(self): super().__init__("network_scan", "Network Configuration") self.logger = logging.getLogger("NetworkScanMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if not self._parse_bool(state.inputs.get("network_scan.run", "True")): self.logger.info("Network scan disabled.") return existing_kdc = state.config_data.get("Network", {}).get("KDC_AUTO_DETECTED") if existing_kdc: self.logger.info(f"Network scan already performed. KDC: {existing_kdc}") return scan_binary = inst_config.get_resource_path("scan/network-agent") if not scan_binary.exists(): self.logger.error(f"Scan binary not found at {scan_binary}") return knoe_home = resolve_knoe_home(env={}) scan_dir = knoe_home / "scan" scan_dir.mkdir(parents=True, exist_ok=True) kdc_found = None pending_line = "" def _handle_line(line: str): nonlocal kdc_found if "KDC is:" in line: try: ip_part = line.split("KDC is:")[1].strip() ip = ip_part.split()[0].strip("[]():,") if ip: kdc_found = ip except Exception: pass elif "Active Directory" in line or "88" in line: import socket for part in line.split(): try: socket.inet_aton(part.strip("[]():,")) kdc_found = part.strip("[]():,") break except Exception: continue def _handle_stdout(text: str): nonlocal pending_line pending_line += text while "\n" in pending_line: line, pending_line = pending_line.split("\n", 1) _handle_line(line + "\n") if progress: progress("Scanning network...", 0.5) rc = run_streaming_cmd( [str(scan_binary)], cwd=str(scan_dir), on_stdout=_handle_stdout ) if pending_line: _handle_line(pending_line) if kdc_found: state.inputs["kerberos_config.kdc"] = kdc_found state.inputs["kerberos_config.enabled"] = "True" if "Network" not in state.config_data: state.config_data["Network"] = {} state.config_data["Network"]["KDC_AUTO_DETECTED"] = kdc_found state.config_data["Network"]["KERBEROS_AUTO_ENABLED"] = "True" self.logger.info(f"KDC detected: {kdc_found}") if progress: progress("Network scan complete", 1.0) class EnvSetupMilestone(Milestone): def __init__(self): super().__init__("env_setup", "System Environment Setup") self.logger = logging.getLogger("EnvSetupMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if progress: progress("Setting up environment...", 0.5) env_keys = [ "KNOE_HOME", "KNOE_CONF", "PROLE_DATA", "PROLE_LOGS", "KNOE_SERVICE", "NAMESPACE", ] if "System Environment" not in state.config_data: state.config_data["System Environment"] = {} for k in env_keys: val = state.inputs.get(f"env_setup.{k}") if val: state.config_data["System Environment"][k] = val if progress: progress("Environment setup complete", 1.0) class SecretManagementMilestone(Milestone): def __init__(self): super().__init__("secret_management", "Secret Resolution") self.logger = logging.getLogger("SecretManagementMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if progress: progress("Resolving secrets...", 0.5) # Keys that might contain secrets secret_keys = [ "init_password.db_password", "init_password.db_password_confirm", "kerberos_config.password", "init_cluster.k3s_token", ] for key in secret_keys: val = state.inputs.get(key) if val: resolved = inst_config._resolve_secret_value(val) if resolved != val: state.inputs[key] = resolved self.logger.info(f"Resolved secret for {key}") # Synchronize confirm password if state.inputs.get("init_password.db_password") and not state.inputs.get( "init_password.db_password_confirm" ): state.inputs["init_password.db_password_confirm"] = state.inputs[ "init_password.db_password" ] if progress: progress("Secret resolution complete", 1.0) class DatabaseCreationMilestone(Milestone): def __init__(self): super().__init__("init_password", "Database Creation") self.logger = logging.getLogger("DatabaseCreationMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if progress: progress("Preparing database configuration...", 0.5) ns = (state.inputs.get("init_password.db_namespace", "") or "").strip() if not ns: ns = ( state.inputs.get("env_setup.DATABASE_NAMESPACE", "") or "" ).strip() if not ns: ns = ( state.inputs.get("env_setup.NAMESPACE", "") or "" ).strip() or "default" cluster_name = (state.inputs.get("init_password.cluster_name", "") or "").strip() if not cluster_name: cluster_name = (state.inputs.get("env_setup.CLUSTER_NAME", "") or "").strip() if not cluster_name: cluster_name = "knoe-db" app_cluster_name = ( state.inputs.get("init_password.app_cluster_name", "") or state.inputs.get("env_setup.APP_CLUSTER_NAME", "") or "knoe-dev-0" ).strip() or "knoe-dev-0" db_cluster_name = ( state.inputs.get("init_password.db_cluster_name", "") or state.inputs.get("env_setup.DB_CLUSTER_NAME", "") or "knoe-cnpg-0" ).strip() or "knoe-cnpg-0" user = (state.inputs.get("init_password.db_username", "") or "").strip() p1 = state.inputs.get("init_password.db_password", "") # In ncurses/silent mode we might not have all fields yet, but we want to ensure defaults if not user: import getpass user = getpass.getuser() state.inputs["init_password.db_username"] = user if not p1: import secrets import string alphabet = string.ascii_letters + string.digits p1 = "".join(secrets.choice(alphabet) for i in range(24)) state.inputs["init_password.db_password"] = p1 state.inputs["init_password.db_password_confirm"] = p1 self.logger.info("Generated default database password.") # If password is an unresolvable OpenBao or KNOE_SECRET reference # (e.g. clean install where OpenBao hasn't been seeded yet), # generate a fresh password so downstream scripts get a real value. # OpenBao isn't running yet at this milestone stage, so if the ref # can't be resolved right now we must provide a concrete password # that will be seeded into OpenBao later by ClusterLifecycleMilestone. if p1.startswith("${"): resolved = inst_config._resolve_secret_value(p1) if resolved and not resolved.startswith("${"): p1 = resolved else: import secrets as _secrets import string as _string alphabet = _string.ascii_letters + _string.digits p1 = "".join(_secrets.choice(alphabet) for i in range(24)) self.logger.info( "OpenBao/secret ref unresolvable at this stage; generated new DB password for initialization." ) state.inputs["init_password.db_password"] = p1 state.inputs["init_password.db_password_confirm"] = p1 if p1 and not p1.startswith("${") and not inst_config._is_knoe_secret(p1): if "Inputs" not in state.config_data: state.config_data["Inputs"] = {} # Encrypt and save to config data for persistence in knoe.cfg encrypted_pw = inst_config._encrypt_cfg_secret(p1) state.config_data["Inputs"]["init_password.db_password"] = encrypted_pw state.config_data["Inputs"][ "init_password.db_password_confirm" ] = encrypted_pw state.inputs["init_password.db_namespace"] = ns state.inputs["env_setup.DATABASE_NAMESPACE"] = ns state.inputs["init_password.cluster_name"] = cluster_name state.inputs["env_setup.CLUSTER_NAME"] = cluster_name state.inputs["init_password.app_cluster_name"] = app_cluster_name state.inputs["env_setup.APP_CLUSTER_NAME"] = app_cluster_name state.inputs["init_password.db_cluster_name"] = db_cluster_name state.inputs["env_setup.DB_CLUSTER_NAME"] = db_cluster_name self.logger.info(f"Application Cluster: {app_cluster_name} (Autopilot)") self.logger.info(f"Database Cluster: {db_cluster_name} (Standard)") self.logger.info( f"CloudNativePG targets dedicated DB cluster '{db_cluster_name}' with CNPG cluster name '{cluster_name}'." ) if "Database Creation" not in state.config_data: state.config_data["Database Creation"] = {} state.config_data["Database Creation"]["DB_USER"] = user state.config_data["Database Creation"]["DATABASE_NAMESPACE"] = ns state.config_data["Database Creation"]["CLUSTER_NAME"] = cluster_name state.config_data["Database Creation"]["APP_CLUSTER_NAME"] = app_cluster_name state.config_data["Database Creation"]["DB_CLUSTER_NAME"] = db_cluster_name state.config_data["Database Creation"].pop("DB_NAME", None) state.config_data["Database Creation"].pop("NAMESPACE", None) # SSH key generation (simplified for now) if self._parse_bool(state.inputs.get("init_password.generate_ssh_key", "True")): key_path = Path.home() / ".ssh" / "id_knoe_ed25519" if not key_path.exists(): self._run_cmd( [ "ssh-keygen", "-t", "ed25519", "-N", "", "-f", str(key_path), "-C", user, ] ) if progress: progress("Database creation preparation complete", 1.0) class DockerBuildMilestone(Milestone): def __init__(self): super().__init__("init_db_build", "Docker Build") self.logger = logging.getLogger("DockerBuildMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if not self._parse_bool(state.inputs.get("init_db_build.run_build", "True")): self.logger.info("Docker build disabled.") return # In k3s/service installs, `init_cloudnative_pg.sh` already has a robust # pre-flight that ensures the `knoe-db` image is present/pushed (and can # build it if missing). Building a multi-platform image here is expensive # and redundant when initialization scripts are enabled. from knoe.core.env import _normalize_cluster_env env_key = ( _normalize_cluster_env(state.inputs.get("init_cluster.cluster_env", "dev")) or "dev" ) run_init_scripts = self._parse_bool( state.inputs.get("init_scripts.run_scripts", "True") ) if env_key != "dev" and env_key != "min" and run_init_scripts: self.logger.info( "Skipping Docker DB image build for service env; init scripts will ensure/push the image." ) return if progress: progress("Building knoe-db database image...", 0.1) tag = state.controller.get_knoe_db_version() args = ["--tag", tag] if env_key != "dev" and env_key != "min": # Add registry prefix for non-dev envs. # Prioritize KNOE_IMAGE_REGISTRY from [Global] if set. registry = (state.config_data.get("Global", {}) or {}).get("KNOE_IMAGE_REGISTRY", "").strip() if not registry or registry == "localhost:5000": server = state.inputs.get("init_cluster.k3s_server_url", "myrddin.prole.org") if "://" in server: server = server.split("://")[1] if ":" in server: server = server.split(":")[0] registry = f"{server}:5000" # Preflight: check registry reachability from current host from knoe.core.env import _http_ping_registry, _registry_image_ref_exists reg_host = registry reg_port = 80 if ":" in registry: parts = registry.split(":") reg_host = parts[0] try: reg_port = int(parts[1]) except: pass self.logger.info(f"Preflight: checking registry reachability at {reg_host}:{reg_port}...") if not _http_ping_registry(reg_host, reg_port): self.logger.warning(f"Registry {registry} is unreachable from this node; proceeding anyway.") remote_tag = f"{registry}/knoe-db:{tag}" # Check if tag already exists in registry if _registry_image_ref_exists(remote_tag): self.logger.info(f"[SKIP] {remote_tag} already exists in registry; skipping build.") if progress: progress(f"Using existing {remote_tag}", 1.0) return args = ["--tag", remote_tag, "--push"] rc = state.controller.run_script("build_db.sh", args=args, on_line=_stream_line) if rc != 0: self.logger.error(f"Docker build failed (code {rc})") if progress: progress("Docker build complete", 1.0) class ClusterLifecycleMilestone(Milestone): def __init__(self): super().__init__("cluster_lifecycle", "Cluster Lifecycle Management") self.logger = logging.getLogger("ClusterLifecycleMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: from knoe.core.env import _normalize_cluster_env cluster_env = state.inputs.get("init_cluster.cluster_env", "dev") env_key = _normalize_cluster_env(cluster_env) or cluster_env start_requested = self._parse_bool( state.inputs.get("init_cluster.start_cluster", "True") ) if not start_requested: self.logger.info("Cluster start not requested.") return env = self._get_script_env(state) if env_key == "min": if progress: progress("Initializing minimal containerd environment...", 0.5) rc = state.controller.run_script("init_min.sh", args=["initialize"], env=env, on_line=_stream_line) if rc != 0: raise RuntimeError(f"Minimal environment initialization failed (code {rc})") if progress: progress("Minimal environment initialized", 1.0) return if env_key == "dev": if progress: progress("Initializing k3d dev cluster...", 0.5) # Ensure Docker is up (best-effort, macOS auto-start). if not state.controller.check_docker_running(): if sys.platform == "darwin": self.logger.info("Starting Docker...") self._run_cmd(["open", "-a", "Docker"]) for _ in range(30): time.sleep(2) if state.controller.check_docker_running(): break if not state.controller.check_docker_running(): self.logger.error("Docker is not running.") return cluster_name = "knoe-dev-cluster" res = subprocess.run( ["k3d", "cluster", "list", "--no-headers"], capture_output=True, text=True, ) res_stdout = res.stdout or "" if cluster_name not in res_stdout: from knoe.core.env import _k3d_knoe_data_volume_args knoe_data = state.inputs.get( "env_setup.PROLE_DATA", str(inst_config.PROJECT_ROOT / "knoe-db" / "data"), ) volume_args = _k3d_knoe_data_volume_args(knoe_data) cmd = ( ["k3d", "cluster", "create", cluster_name, "-a", "2"] + volume_args + ["--api-port", "0.0.0.0:6443"] ) rc = self._run_cmd(cmd, on_stdout=_stream_line) if rc != 0: self.logger.error(f"k3d cluster create failed (code {rc})") else: if "running" not in res_stdout.lower(): self._run_cmd( ["k3d", "cluster", "start", cluster_name], on_stdout=_stream_line, ) try: subprocess.run( ["kubectl", "config", "use-context", f"k3d-{cluster_name}"], capture_output=True, ) except Exception: pass # Ensure knoe-db image is available in the cluster (avoid registry pull issues). tag = state.controller.get_knoe_db_version() local_image = f"knoe-db:{tag}" if ( subprocess.run( ["docker", "image", "inspect", local_image], capture_output=True ).returncode == 0 ): db_cfg = state.config_data.get("Docker Build", {}) or {} registry = (db_cfg.get("LOCAL_REGISTRY_INTERNAL") or "").strip() if registry.endswith(".localhost:5000"): registry = registry.replace(".localhost:5000", ":5000") elif registry.endswith(".localhost"): registry = registry.replace(".localhost", "") if not registry: registry = "k3d-knoe-registry:5000" remote_tag = f"{registry}/knoe-db:{tag}" subprocess.run( ["docker", "tag", local_image, remote_tag], capture_output=True ) self._run_cmd( ["k3d", "image", "import", remote_tag, "-c", cluster_name], on_stdout=_stream_line, ) elif env_key in ("service", "k3s"): if progress: progress("Verifying k3s connection...", 0.5) # Additional logic for remote k3s services could go here self.logger.info( "Service cluster mode: verification and deployment triggered via scripts." ) # For now, matching the behavior of silent installer's _step_init_cluster # which mostly sets up env and calls init scripts if needed. # Initialize OpenBao once cluster is ready. If DB password is anchored, # start OpenBao anyway so downstream scripts can attempt to read secrets. db_pw_raw = (state.inputs.get("init_password.db_password", "") or "").strip() if db_pw_raw: db_pw = inst_config._resolve_secret_value(db_pw_raw) if db_pw and not db_pw.startswith("${"): env_init = dict(env) env_init["DB_PASSWORD"] = db_pw env_init.setdefault("GRAFANA_ADMIN_PASSWORD", db_pw) db_user = ( state.inputs.get("init_password.db_username", "") or "" ).strip() if db_user: env_init["KNOE_DB_USER"] = db_user try: openbao_ops.initialize( namespace=env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default")), env=env_init, project_root=getattr(state.controller, "project_root", "."), mode=env.get("KNOE_MODE"), log=lambda msg: _stream_line(msg if msg.endswith("\n") else msg + "\n"), ) except Exception as e: self.logger.error(f"OpenBao initialization failed: {e}") else: self.logger.info( "OpenBao init skipped: DB password unresolved/anchored. Ensuring OpenBao is running..." ) try: openbao_ops.start( namespace=env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default")), env=env, project_root=getattr(state.controller, "project_root", "."), mode=env.get("KNOE_MODE"), log=lambda msg: _stream_line(msg if msg.endswith("\n") else msg + "\n"), ) except Exception as e: self.logger.error(f"OpenBao start failed: {e}") # When common services are already healthy, run the service-layer # migration inline so that the "Preparing Service Layer" steps are # completed automatically — mirroring the UI shortcut behaviour. if env_key in ("service", "k3s"): service_ns = env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default")) mode = str(env.get("KNOE_MODE", "")).strip().lower() status_checks = [ registry_ops.status( namespace=env.get("REGISTRY_NAMESPACE", service_ns), env=env, mode=env.get("KNOE_MODE"), ), openbao_ops.status( namespace=service_ns, env=env, mode=env.get("KNOE_MODE"), ), garage_store_ops.status(namespace=service_ns, env=env), ] if mode != "k8s": status_checks.append(opentofu_ops.status(namespace=service_ns, env=env)) else: self.logger.info("k8s mode: skipping OpenTofu health check.") status_ok = all(status_checks) if status_ok: self.logger.info( "All common services healthy – running service layer migration inline." ) rc_migrate = state.controller.run_script( "init_service_layer.sh", args=["migrate"], env=env, on_line=_stream_line, ) if rc_migrate != 0: self.logger.error( f"Service layer migration failed (code {rc_migrate})" ) else: self.logger.info("Service layer migration completed successfully.") if progress: progress("Cluster lifecycle management complete", 1.0) class InitializationScriptsMilestone(Milestone): def __init__(self): super().__init__("init_scripts", "Initialization Scripts") self.logger = logging.getLogger("InitializationScriptsMilestone") def _regenerate_port_mapping_cfg(self, state: InstallerState, mode: str) -> None: """Regenerate conf/port-mapping.cfg from [Port Forwards] in knoe.cfg data.""" conf_dir = inst_config.PROJECT_ROOT / "conf" mapping_path = conf_dir / "port-mapping.cfg" pf_section = state.config_data.get("Port Forwards", {}) or {} prefix = ( "PORT_FORWARD_K3S_MAPPING_" if mode == "k3s" else "PORT_FORWARD_K3D_MAPPING_" ) # Resolve ${NAMESPACE} from state ns = (state.inputs.get("init_password.db_namespace", "") or "").strip() if not ns: ns = ( (state.config_data.get("Global", {}) or {}).get("NAMESPACE", "").strip() ) if not ns: ns = "default" lines = [ "# Port mappings for Knoe Tools (generated).", "# Format: key: local=... remote=... ns=... svc=... address=...", "", ] for k, v in sorted(pf_section.items()): if not k.startswith(prefix): continue parts = {} for p in str(v).split(";"): if "=" in p: kv = p.split("=", 1) if len(kv) == 2: parts[kv[0].strip()] = kv[1].strip() if "id" not in parts: continue m_id = parts["id"] local = parts.get("hostPort", "") remote = parts.get("servicePort", "") m_ns = parts.get("namespace", "") target = parts.get("target", "") addr = parts.get("address", "0.0.0.0") # Resolve ${NAMESPACE} references if m_ns.startswith("${") and "NAMESPACE" in m_ns: m_ns = ns if local.startswith("${"): continue # skip entries with unresolvable port refs svc = target[4:] if target.startswith("svc/") else target lines.append( f"{m_id}: local={local} remote={remote} ns={m_ns} svc={svc} address={addr}" ) try: mapping_path.write_text("\n".join(lines) + "\n") self.logger.info(f"Regenerated {mapping_path}") except Exception as e: self.logger.error(f"Failed to regenerate port-mapping.cfg: {e}") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: env = self._get_script_env(state) mode = env.get("KNOE_MODE", "") manage_opentofu = str(mode).strip().lower() != "k8s" mode_args = ["--mode", mode] if mode else [] raw_min = (state.config_data.get("Global", {}) or {}).get(POLICY_CFG_KEY, "") try: min_nodes = int(str(raw_min).strip() or str(OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES)) except Exception: min_nodes = OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES opt_allowed, _count, opt_reason = evaluate_optional_workloads_allowed( env=env, min_nodes=min_nodes ) self.logger.info(f"[POLICY] {opt_reason}") # Regenerate port-mapping.cfg from knoe.cfg data to ensure it's clean self._regenerate_port_mapping_cfg(state, mode or "k3d") # Phase 1: common services (Python owners) if progress: progress("Running common services (python owners)...", 0.0) service_ns = env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default")) registry_ns = env.get("REGISTRY_NAMESPACE", service_ns) project_root = getattr(state.controller, "project_root", ".") mode = env.get("KNOE_MODE", "") # In k8s mode the default KUBECONTEXT is the app cluster; # if a DB cluster context is provided, some services (like Garage) # may need to run there to be close to the data. app_ctx = env.get("APP_CLUSTER_KUBECONTEXT", "").strip() db_ctx = env.get("DB_CLUSTER_KUBECONTEXT", "").strip() is_split_gke = mode == "k8s" and db_ctx and app_ctx and db_ctx != app_ctx try: registry_ops.update( namespace=registry_ns, env=env, project_root=project_root, mode=mode, log=self.logger.info, ) openbao_ops.update( namespace=service_ns, env=env, project_root=project_root, mode=mode, log=self.logger.info, ) # Garage on DB cluster for split-GKE garage_env = dict(env) if is_split_gke: self.logger.info( f"[GARAGE] Split-cluster GKE detected. Using DB cluster context {db_ctx} " f"as authoritative for Garage (namespace: {service_ns})." ) garage_env["KUBECONTEXT"] = db_ctx else: self.logger.info(f"[GARAGE] Deploying Garage to default context (namespace: {service_ns}).") garage_store_ops.update( namespace=service_ns, env=garage_env, project_root=project_root, mode=mode, log=self.logger.info, ) if is_split_gke: self.logger.info(f"[GARAGE] Fetching ILB IP for service 'garage-s3-ilb' from DB cluster {db_ctx}...") # Retry for up to 5 minutes, as GKE ILB address assignment can be slow. garage_ip = None for i in range(30): garage_ip = garage_store_ops.get_endpoint_ip( namespace=service_ns, env=garage_env, log=self.logger.info ) if garage_ip: break self.logger.info(f"[GARAGE] Waiting for Garage ILB IP (attempt {i+1}/30)...") time.sleep(10) if garage_ip: endpoint = f"http://{garage_ip}:3900" self.logger.info(f"[GARAGE] Resolved authoritative DB-cluster Garage private endpoint: {endpoint}") # Persist to state config data for downstream GitLab use state.config_data.setdefault("Global", {})["GARAGE_PRIVATE_S3_ENDPOINT"] = endpoint # Also update current env so immediate next steps see it env["GARAGE_PRIVATE_S3_ENDPOINT"] = endpoint env["GARAGE_S3_ENDPOINT"] = endpoint # Persist to disk (knoe.cfg / gke.cfg) so reruns see it immediately if hasattr(state.controller, "write_config"): self.logger.info("[GARAGE] Persisting resolved endpoint to config file...") state.controller.write_config() else: msg = ( f"[GARAGE] [REPAIR_BLOCKED] Failed to resolve Garage ILB IP on DB cluster {db_ctx}. " "Ensure the 'garage-s3-ilb' Service exists with 'networking.gke.io/load-balancer-type: Internal' " "annotation and has been assigned an IP." ) self.logger.error(msg) raise Exception(msg) if manage_opentofu: opentofu_ops.update( namespace=service_ns, env=env, project_root=project_root, mode=mode, log=self.logger.info, ) else: self.logger.info("k8s mode: skipping OpenTofu deploy.") except Exception as e: msg = f"Common services (python owners) failed: {e}" self.logger.error(msg) raise Exception(msg) # Phase 2: CNPG initialization (Python) if progress: progress("Initializing CloudNative-PG...", 1 / 6) ns = env.get("DATABASE_NAMESPACE") or env.get("NAMESPACE", "default") cluster_name = env.get("CLUSTER_NAME") or env.get("CNPG_CLUSTER_NAME") or "knoe-db" # Build a DB-cluster-scoped env for CNPG operations. # In k8s mode the default KUBECONTEXT is the app cluster; CNPG lives on the # dedicated DB cluster, so override KUBECONTEXT to DB_CLUSTER_KUBECONTEXT. cnpg_env = dict(env) _db_ctx = env.get("DB_CLUSTER_KUBECONTEXT", "").strip() if _db_ctx: cnpg_env["KUBECONTEXT"] = _db_ctx self.logger.info(f"CNPG operations will use DB cluster context: {_db_ctx}") # Ensure DB secrets exist before CNPG init (mirrors _step_init_scripts). # env["DB_PASSWORD"] is already resolved by _script_env_for_namespace. db_pw = (env.get("DB_PASSWORD") or state.inputs.get("init_password.db_password") or "").strip() if db_pw and hasattr(state.controller, "ensure_db_k8s_secrets"): try: state.controller.ensure_db_k8s_secrets(ns, db_pw, log_fn=self.logger.info) except Exception as _sec_err: raise Exception(f"Failed to ensure DB secrets before CNPG init: {_sec_err}") # Multi-cluster: iterate over [CNPG Clusters] registry when declared. # Each cluster is initialized to its fixed declared namespace. # Falls back to single-cluster behavior when no registry is configured. cluster_registry = [] try: _registry_fn = getattr(state.controller, "_cnpg_cluster_registry", None) if _registry_fn is not None: cluster_registry = _registry_fn() except Exception as _reg_err: self.logger.warning(f"Could not read CNPG cluster registry: {_reg_err}") if cluster_registry: for _cluster in cluster_registry: _ns = _cluster.get("namespace") or ns _name = _cluster.get("name") or cluster_name _cluster_env = dict(cnpg_env) if _cluster.get("image"): _cluster_env["CNPG_IMAGE_NAME"] = _cluster["image"] if db_pw and hasattr(state.controller, "ensure_db_k8s_secrets"): try: state.controller.ensure_db_k8s_secrets(_ns, db_pw, log_fn=self.logger.info) except Exception as _sec_err: self.logger.warning(f"DB secrets for {_ns}: {_sec_err}") cnpg_initialize( namespace=_ns, cluster_name=_name, env=_cluster_env, project_root=project_root, log=self.logger.info, mode=mode, ) else: cnpg_initialize( namespace=ns, cluster_name=cluster_name, env=cnpg_env, project_root=project_root, log=self.logger.info, mode=mode, ) # Phase 3: remaining scripts (backup, kong, optional ingress, redis) post_scripts = ["init_redis.sh", "init_cnpg_backup.sh", "init_kong.sh"] if not opt_allowed: self.logger.info(f"[SKIP] Monitoring disabled by policy: {opt_reason}") state.config_data.setdefault("Monitoring", {})["STATUS"] = "Skipped" if mode != "k3d": post_scripts.append("init_nginx_ingress.sh") for i, script in enumerate(post_scripts): if progress: progress(f"Running {script}...", (2 + i) / (2 + len(post_scripts))) if script == "init_cnpg_backup.sh": args: list[str] = mode_args + ["start"] # CNPG backup runs on the DB cluster, not the app cluster script_env = cnpg_env elif script == "init_kong.sh" or script == "init_redis.sh": args = mode_args + ["start"] script_env = env else: args = mode_args + ["initialize"] script_env = env rc = state.controller.run_script( script, args=args, env=script_env, on_line=_stream_line ) if rc != 0: msg = f"Script {script} failed (code {rc})" self.logger.error(msg) raise Exception(msg) if opt_allowed: if progress: progress("Running monitoring (python owner)...", (2 + len(post_scripts)) / (2 + len(post_scripts) + 1)) monitoring_ns = env.get("MONITORING_NAMESPACE", "monitoring") try: monitoring_ops.initialize( namespace=monitoring_ns, env=env, mode=mode, log=self.logger.info, ) except Exception as e: msg = f"Monitoring initialization failed: {e}" self.logger.error(msg) raise Exception(msg) # Verify critical secrets in the identity cluster's namespace. if cluster_registry: _identity = next( (c for c in cluster_registry if c.get("identity", "").lower() == "true"), cluster_registry[0], ) ns = _identity.get("namespace") or ns else: ns = env.get("DATABASE_NAMESPACE") or env.get("NAMESPACE", "default") self.logger.info(f"Verifying critical secrets in namespace {ns}...") critical_secrets = ["knoe-db-user", "knoe-db-superuser", "cnpg-admin-key"] missing_secrets = [] for secret in critical_secrets: if ( self._run_cmd( ["kubectl", "get", "secret", secret, "-n", ns], env=cnpg_env, ) != 0 ): missing_secrets.append(secret) if missing_secrets: msg = f"CRITICAL: Missing secrets in namespace '{ns}': {', '.join(missing_secrets)}. Database initialization will fail." self.logger.error(msg) if progress: progress(msg, 1.0) raise Exception(msg) # For k3d: start port-forward daemon and verify Grafana port-forward. pf_ok = True if mode == "k3d": def _pf_line(line: str) -> None: self.logger.info(line.rstrip("\n")) env_pf = dict(env) env_pf.setdefault("PORT_FORWARD_SKIP_VALIDATE", "1") env_pf.setdefault("PORT_FORWARD_SKIP_WAIT", "1") env_pf.setdefault("PORT_FORWARD_WAIT_TIMEOUT", "20") env_pf.setdefault("PORT_FORWARD_WAIT_INTERVAL", "2") rc_pf = state.controller.run_script( "init_port_forwards.sh", args=["--force", "--mode", mode, "start"], env=env_pf, on_line=_stream_line, ) if rc_pf != 0: pf_ok = False self.logger.error(f"Port-forwards start failed (code {rc_pf})") else: for _ in range(10): status_lines: list[str] = [] def _pf_status_line(line: str) -> None: status_lines.append(line) _stream_line(line) rc_status = state.controller.run_script( "init_port_forwards.sh", args=["--mode", mode, "status", "grafana"], env=env_pf, on_line=_pf_status_line, ) running = any( l.lstrip().startswith("RUNNING") and "grafana" in l for l in status_lines ) if rc_status == 0 and running: break time.sleep(1) else: pf_ok = False self.logger.error("Grafana port-forward not running after start.") pf_section = state.config_data.setdefault("Port Forwards", {}) pf_section["STATUS"] = "Running" if pf_ok else "Failed" if progress: progress("Initialization scripts complete", 1.0) class DeploymentMilestone(Milestone): def __init__(self): super().__init__("deployment", "Knoe.dev Deployment") self.logger = logging.getLogger("DeploymentMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if progress: progress("Deploying Knoe.dev ...", 0.3) env = self._get_script_env(state) mode = env.get("KNOE_MODE", "") ns = env.get("DATABASE_NAMESPACE") or env.get("NAMESPACE", "default") cluster_name = env.get("CLUSTER_NAME") or env.get("CNPG_CLUSTER_NAME") or "knoe-db" project_root = getattr(state.controller, "project_root", ".") # Ensure DB context is set for CNPG operations in k8s mode _db_ctx = env.get("DB_CLUSTER_KUBECONTEXT", "").strip() if _db_ctx: env["KUBECONTEXT"] = _db_ctx self.logger.info(f"CNPG deployment will use DB cluster context: {_db_ctx}") # Multi-cluster: iterate over [CNPG Clusters] registry when declared. cluster_registry = [] try: _registry_fn = getattr(state.controller, "_cnpg_cluster_registry", None) if _registry_fn is not None: cluster_registry = _registry_fn() except Exception as _reg_err: self.logger.warning(f"Could not read CNPG cluster registry: {_reg_err}") if cluster_registry: for _cluster in cluster_registry: _ns = _cluster.get("namespace") or ns _name = _cluster.get("name") or cluster_name _cluster_env = dict(env) if _cluster.get("image"): _cluster_env["CNPG_IMAGE_NAME"] = _cluster["image"] cnpg_deploy( namespace=_ns, cluster_name=_name, env=_cluster_env, project_root=project_root, log=self.logger.info, mode=mode, ) else: cnpg_deploy( namespace=ns, cluster_name=cluster_name, env=env, project_root=project_root, log=self.logger.info, mode=mode, ) self.logger.info("Deployment successful") if progress: progress("Deployment complete", 1.0) class SecurityMilestone(Milestone): """Enforce uniform master password across all services via update.sh. Runs update.sh (reads vault_knoe_db_master_password from the Ansible vault) and applies the password to k8s secrets, PostgreSQL, and Grafana. Only active in k8s mode — silently skipped for k3s/k3d. """ def __init__(self): super().__init__("security", "Security Hardening") self.logger = logging.getLogger("SecurityMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: env = self._get_script_env(state) mode = env.get("KNOE_MODE", "") if mode != "k8s": self.logger.info( f"[SECURITY] Mode={mode!r}: skipping live secret rotation (non-k8s)." ) if progress: progress("Security hardening skipped (non-k8s mode)", 1.0) return project_root = getattr(state.controller, "project_root", ".") update_sh = str(Path(project_root) / "update.sh") if not Path(update_sh).exists(): self.logger.warning( f"[SECURITY] update.sh not found at {update_sh} — skipping." ) if progress: progress("update.sh not found — skipping security hardening", 1.0) return if progress: progress("Rotating master password via update.sh...", 0.2) self.logger.info("[SECURITY] Running update.sh to enforce vault master password.") script_env = { **env, "PATH": os.environ.get("PATH", ""), } rc = subprocess.call([update_sh], env=script_env) if rc != 0: raise Exception( f"update.sh failed with exit code {rc}. " "Run './update.sh --prompt' to recreate the vault file if it is corrupt." ) self.logger.info("[SECURITY] Master password rotation complete.") state.config_data.setdefault("Security", {})["STATUS"] = "Secured" if progress: progress("Security hardening complete", 1.0) class GitOpsMilestone(Milestone): def __init__(self): super().__init__("gitops", "GitOps") self.logger = logging.getLogger("GitOpsMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: enabled = self._parse_bool( state.inputs.get("init_cluster.gitops_enabled", "False"), default=False ) if not enabled: state.config_data.setdefault("GitOps", {})["STATUS"] = "Skipped" return # Resolve git provider: [Inputs] gitops.git_provider takes precedence, # then [Optional Features] GITOPS_PROVIDER, then default to gitea. provider = ( str(state.inputs.get("gitops.git_provider", "") or "").strip().lower() or str( (state.config_data.get("Optional Features") or {}).get( "GITOPS_PROVIDER", "" ) or "" ) .strip() .lower() or "gitea" ) is_gitlab = provider == "gitlab" label = "GitLab" if is_gitlab else "Gitea" if progress: progress(f"Deploying {label}...", 0.4) env = self._get_script_env(state) mode = env.get("KNOE_MODE", "") def _first_non_empty(*values: object) -> str: for value in values: candidate = str(value or "").strip() if candidate: return candidate return "" def _resolve_secret_value(raw: str, *fallback_env_keys: str) -> str: candidate = str(raw or "").strip() if not candidate: return "" def _fallback_from_env() -> str: for key in fallback_env_keys: env_candidate = _first_non_empty(env.get(key, ""), os.environ.get(key, "")) if env_candidate: return env_candidate return "" try: resolved = state.controller._resolve_secret_value(candidate) resolved_str = str(resolved or "").strip() if resolved_str and not resolved_str.startswith("secretref://"): return resolved_str if candidate.startswith("secretref://"): fallback = _fallback_from_env() if fallback: return fallback return resolved_str or candidate except Exception: if candidate.startswith("secretref://"): fallback = _fallback_from_env() if fallback: return fallback return candidate if is_gitlab: # Always prefer the explicit gitlab_namespace key; never fall back to # gitops.namespace which may be 'gitea' from a prior Gitea run. ns = ( str(state.inputs.get("gitops.gitlab_namespace", "") or "").strip() or str( (state.config_data.get("Global") or {}).get( "GITLAB_NAMESPACE", "" ) or "" ).strip() or "gitlab" ) env["GITLAB_NAMESPACE"] = ns global_cfg = state.config_data.setdefault("Global", {}) oidc_client_id_ref = _first_non_empty( state.inputs.get("auth.clientId", ""), state.inputs.get("oidc.clientId", ""), global_cfg.get("GITLAB_OIDC_CLIENT_ID", ""), global_cfg.get("OIDC_CLIENT_ID", ""), global_cfg.get("GOOGLE_OIDC_CLIENT_ID", ""), global_cfg.get("OIDC_CLIENT_ID_REF", ""), global_cfg.get("GOOGLE_OIDC_CLIENT_ID_REF", ""), os.environ.get("GITLAB_OIDC_CLIENT_ID", ""), os.environ.get("OIDC_CLIENT_ID", ""), os.environ.get("GOOGLE_OIDC_CLIENT_ID", ""), ) oidc_client_secret_ref = _first_non_empty( state.inputs.get("auth.clientSecret", ""), state.inputs.get("oidc.clientSecret", ""), global_cfg.get("GITLAB_OIDC_CLIENT_SECRET", ""), global_cfg.get("OIDC_CLIENT_SECRET", ""), global_cfg.get("GOOGLE_OIDC_CLIENT_SECRET", ""), global_cfg.get("OIDC_CLIENT_SECRET_REF", ""), global_cfg.get("GOOGLE_OIDC_CLIENT_SECRET_REF", ""), os.environ.get("GITLAB_OIDC_CLIENT_SECRET", ""), os.environ.get("OIDC_CLIENT_SECRET", ""), os.environ.get("GOOGLE_OIDC_CLIENT_SECRET", ""), ) oidc_client_id = _resolve_secret_value( oidc_client_id_ref, "GITLAB_OIDC_CLIENT_ID", "OIDC_CLIENT_ID", "GOOGLE_OIDC_CLIENT_ID", "GOOGLE_CLIENT_ID", ) oidc_client_secret = _resolve_secret_value( oidc_client_secret_ref, "GITLAB_OIDC_CLIENT_SECRET", "OIDC_CLIENT_SECRET", "GOOGLE_OIDC_CLIENT_SECRET", "GOOGLE_CLIENT_SECRET", ) if oidc_client_id: env["GITLAB_OIDC_CLIENT_ID"] = oidc_client_id if oidc_client_secret: env["GITLAB_OIDC_CLIENT_SECRET"] = oidc_client_secret oidc_provider = _first_non_empty( env.get("GITLAB_OIDC_PROVIDER_NAME", ""), global_cfg.get("GITLAB_OIDC_PROVIDER_NAME", ""), os.environ.get("GITLAB_OIDC_PROVIDER_NAME", ""), "openid_connect", ) env["GITLAB_OIDC_PROVIDER_NAME"] = oidc_provider gitlab_host = _first_non_empty( env.get("GITLAB_DOMAIN", ""), global_cfg.get("GITLAB_DOMAIN", ""), global_cfg.get("GITLAB_HOST", ""), state.inputs.get("gitops.gitlab_host", ""), state.inputs.get("routing.gitHost", ""), os.environ.get("GITLAB_DOMAIN", ""), os.environ.get("GITLAB_HOST", ""), "git.knoe.dev" if mode == "k8s" else "", ) oidc_issuer = _first_non_empty( env.get("GITLAB_OIDC_ISSUER", ""), state.inputs.get("auth.issuer", ""), global_cfg.get("GITLAB_OIDC_ISSUER", ""), global_cfg.get("OIDC_ISSUER", ""), os.environ.get("GITLAB_OIDC_ISSUER", ""), os.environ.get("OIDC_ISSUER", ""), "https://api.knoe.dev/auth" if mode == "k8s" else "", ) if oidc_issuer: env["GITLAB_OIDC_ISSUER"] = oidc_issuer default_redirect_uri = ( "https://git.knoe.dev/users/auth/openid_connect/callback" if mode == "k8s" else ( f"https://{gitlab_host}/users/auth/openid_connect/callback" if gitlab_host else "" ) ) oidc_redirect_uri = _first_non_empty( env.get("GITLAB_OIDC_REDIRECT_URI", ""), global_cfg.get("GITLAB_OIDC_REDIRECT_URI", ""), os.environ.get("GITLAB_OIDC_REDIRECT_URI", ""), default_redirect_uri, ) if oidc_redirect_uri: env["GITLAB_OIDC_REDIRECT_URI"] = oidc_redirect_uri # Knoe User provisioning integration surface. # Keep JIT login creation active while exposing stable metadata for a # future proactive user create/sync workflow. global_cfg["KNOE_USER_GITLAB_AUTH_PROVIDER"] = oidc_provider global_cfg["KNOE_USER_GITLAB_JIT_AUTO_CREATE_USERS"] = "true" global_cfg["KNOE_USER_GITLAB_PROVISIONING_READY"] = "true" if gitlab_host: global_cfg["KNOE_USER_GITLAB_API_BASE"] = f"https://{gitlab_host}/api/v4" if oidc_issuer: global_cfg["KNOE_USER_GITLAB_OIDC_ISSUER"] = oidc_issuer if oidc_redirect_uri: global_cfg["KNOE_USER_GITLAB_OIDC_REDIRECT_URI"] = oidc_redirect_uri node_selector = ( str(state.inputs.get("gitops.node_selector", "") or "").strip() ) if node_selector: env["GITLAB_NODE_SELECTOR"] = node_selector env["NODE_SELECTOR"] = node_selector script = "init_gitlab.sh" else: ns = ( state.inputs.get("gitops.namespace", "") or env.get("GITEA_NAMESPACE") or "gitea" ).strip() env["GITEA_NAMESPACE"] = ns script = "init_gitea.sh" args: list[str] = [] if mode: args.extend(["--mode", mode]) args.extend(["--namespace", ns]) conf_dir = knoe_conf.resolve_knoe_conf_dir(state.controller.project_root) cfg_path = knoe_conf.entrypoint_path(conf_dir) if cfg_path.exists(): args.extend(["-c", str(cfg_path)]) self.logger.info(f"{script} {' '.join(args)}") rc = state.controller.run_script( script, args=args, env=env, on_line=_stream_line ) if rc == 0: gitops_sec = state.config_data.setdefault("GitOps", {}) gitops_sec["STATUS"] = "Deployed" gitops_sec["GITOPS_NAMESPACE"] = ns if progress: progress(f"{label} deployed", 1.0) else: gitops_sec = state.config_data.setdefault("GitOps", {}) gitops_sec["STATUS"] = "Attempted" gitops_sec["GITOPS_NAMESPACE"] = ns msg = f"{label} deploy failed (code {rc})" self.logger.error(msg) raise Exception(msg) class KerberosMilestone(Milestone): def __init__(self): super().__init__("kerberos", "Kerberos Initialization") self.logger = logging.getLogger("KerberosMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: enabled = self._parse_bool( state.inputs.get("kerberos_config.enabled", "False"), default=False ) if not enabled: realm = str(state.inputs.get("kerberos_config.realm", "") or "").strip() kdc = str(state.inputs.get("kerberos_config.kdc", "") or "").strip() if realm and kdc: self.logger.info( "Kerberos enabled implicitly because realm and kdc are configured." ) enabled = True krb_sec = state.config_data.setdefault("Kerberos Authentication", {}) if not enabled: krb_sec["STATUS"] = "Skipped" if progress: progress("Kerberos disabled; skipping", 1.0) return env = self._get_script_env(state) mode = env.get("KNOE_MODE", "") if mode == "min": # For minimal mode, we use init_min.sh to bootstrap or just proceed pass args = ["--mode", mode] if mode else [] args.append("initialize") if progress: progress("Initializing Kerberos...", 0.4) rc = state.controller.run_script( "init_kerberos.sh", args=args, env=env, on_line=_stream_line ) if rc == 0: krb_sec["STATUS"] = "Initialized" if progress: progress("Kerberos initialized", 1.0) else: krb_sec["STATUS"] = "Attempted" msg = f"Kerberos initialization failed (code {rc})" self.logger.error(msg) raise Exception(msg) def _resolve_supabase_deploy_mode(state: InstallerState) -> str: from knoe.core.env import _normalize_cluster_env deploy_mode = ( os.environ.get("SUPABASE_DEPLOY_MODE") or os.environ.get("SUPABASE_MODE") or "" ).strip() if deploy_mode: return deploy_mode cluster_env = _normalize_cluster_env(state.inputs.get("init_cluster.cluster_env", "dev")) if cluster_env in ("service", "prod"): return "k8s" return "k3d" class SupabaseImagePreloadMilestone(Milestone): def __init__(self): super().__init__("supabase_images_preload", "Supabase Image Preload") self.logger = logging.getLogger("SupabaseImagePreloadMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: enabled = self._parse_bool( state.inputs.get("init_cluster.supabase_enabled", "False"), default=False, ) supa_sec = state.config_data.setdefault("Supabase", {}) if not enabled: supa_sec["IMAGES_STATUS"] = "Skipped" if progress: progress("Supabase disabled; skipping image preload", 1.0) return raw_min = (state.config_data.get("Global", {}) or {}).get(POLICY_CFG_KEY, "") try: min_nodes = int(str(raw_min).strip() or str(OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES)) except Exception: min_nodes = OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES opt_allowed, _count, opt_reason = evaluate_optional_workloads_allowed( env=self._get_script_env(state), min_nodes=min_nodes ) self.logger.info(f"[POLICY] {opt_reason}") if not opt_allowed: supa_sec["IMAGES_STATUS"] = "Skipped" if progress: progress(f"Supabase image preload skipped by policy: {opt_reason}", 1.0) return from knoe.core.env import _parse_bool env = self._get_script_env(state) script_path = state.controller.project_root / "supabase" / "deploy.sh" if not script_path.exists(): self.logger.error(f"Supabase deploy script not found: {script_path}") supa_sec["IMAGES_STATUS"] = "Attempted" return deploy_mode = _resolve_supabase_deploy_mode(state) if deploy_mode != "k3d": supa_sec["IMAGES_STATUS"] = "Skipped" if progress: progress( f"Supabase image preload not required for mode '{deploy_mode}'", 1.0, ) return args = ["--mode", deploy_mode, "--prefetch-images-only"] conf_dir = knoe_conf.resolve_knoe_conf_dir(state.controller.project_root) cfg_path = knoe_conf.entrypoint_path(conf_dir) if cfg_path.exists(): args.extend(["-c", str(cfg_path)]) if _parse_bool(os.environ.get("SUPABASE_USE_DEV_COMPOSE"), False): args.append("--with-dev-helpers") if progress: progress("Preloading Supabase images...", 0.4) self.logger.info(f"supabase/deploy.sh {' '.join(args)}") rc = subprocess.run( ["bash", str(script_path)] + args, env=env, ).returncode if rc == 0: supa_sec["IMAGES_STATUS"] = "Prepared" state.inputs["supabase.images_preloaded"] = "True" if progress: progress("Supabase images preloaded", 1.0) else: supa_sec["IMAGES_STATUS"] = "Attempted" msg = f"Supabase image preload failed (code {rc})" self.logger.error(msg) raise Exception(msg) class SupabaseMilestone(Milestone): def __init__(self): super().__init__("supabase", "Supabase Deployment") self.logger = logging.getLogger("SupabaseMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: enabled = self._parse_bool( state.inputs.get("init_cluster.supabase_enabled", "False"), default=False, ) supa_sec = state.config_data.setdefault("Supabase", {}) if not enabled: supa_sec["STATUS"] = "Skipped" if progress: progress("Supabase disabled; skipping", 1.0) return raw_min = (state.config_data.get("Global", {}) or {}).get(POLICY_CFG_KEY, "") try: min_nodes = int(str(raw_min).strip() or str(OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES)) except Exception: min_nodes = OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES opt_allowed, _count, opt_reason = evaluate_optional_workloads_allowed( env=self._get_script_env(state), min_nodes=min_nodes ) self.logger.info(f"[POLICY] {opt_reason}") if not opt_allowed: supa_sec["STATUS"] = "Skipped" if progress: progress(f"Supabase skipped by policy: {opt_reason}", 1.0) return from knoe.core.env import _parse_bool env = self._get_script_env(state) script_path = state.controller.project_root / "supabase" / "deploy.sh" if not script_path.exists(): self.logger.error(f"Supabase deploy script not found: {script_path}") supa_sec["STATUS"] = "Attempted" return deploy_mode = _resolve_supabase_deploy_mode(state) args = ["--mode", deploy_mode] conf_dir = knoe_conf.resolve_knoe_conf_dir(state.controller.project_root) cfg_path = knoe_conf.entrypoint_path(conf_dir) if cfg_path.exists(): args.extend(["-c", str(cfg_path)]) if _parse_bool(os.environ.get("SUPABASE_USE_DEV_COMPOSE"), False): args.append("--with-dev-helpers") if _parse_bool(os.environ.get("SUPABASE_FOREGROUND"), False): args.append("--foreground") if _parse_bool(state.inputs.get("supabase.images_preloaded", "False"), False): args.append("--skip-prefetch") self.logger.info(f"supabase/deploy.sh {' '.join(args)}") rc = subprocess.run( ["bash", str(script_path)] + args, env=env, ).returncode if rc == 0: supa_sec["STATUS"] = "Deployed" if progress: progress("Supabase deployed", 1.0) else: supa_sec["STATUS"] = "Attempted" msg = f"Supabase deploy failed (code {rc})" self.logger.error(msg) raise Exception(msg)