from __future__ import annotations import logging import os 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 prole_conf from knoe.core.stream_exec import run_streaming_cmd from knoe.core.env import resolve_prole_home from knoe.core.policy import ( POLICY_CFG_KEY, OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES, evaluate_optional_workloads_allowed, ) 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") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if progress: progress("Checking dependencies...", 0.1) dependencies = inst_config.DEPENDENCIES missing = [] for i, dep in enumerate(dependencies): ok, location, 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._set_status(state, "All installed") if progress: progress("All dependencies installed", 1.0) 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.logger.error("Dependencies missing and auto-install disabled.") self._set_status(state, "Missing") return 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(f"Install failed for {dep['name']} (code {rc})") # Re-verify still_missing = [] for dep in missing: ok, _, _ = inst_config.get_dep_info(dep) if not ok: still_missing.append(dep["name"]) if still_missing: self.logger.error(f"Still missing: {', '.join(still_missing)}") self._set_status(state, "Missing") else: self._set_status(state, "All installed") if progress: progress("All dependencies installed", 1.0) def _set_status(self, state: InstallerState, status: str): if "Dependencies" not in state.config_data: state.config_data["Dependencies"] = {} state.config_data["Dependencies"]["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("prole-net/prole-agent") if not scan_binary.exists(): self.logger.error(f"Scan binary not found at {scan_binary}") return prole_home = resolve_prole_home(env={}) scan_dir = prole_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 = [ "PROLE_HOME", "PROLE_CONF", "PROLE_DATA", "PROLE_LOGS", "PROLE_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.NAMESPACE", "") or "" ).strip() or "default" 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 PROLE_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_prole_secret(p1): if "Inputs" not in state.config_data: state.config_data["Inputs"] = {} # Encrypt and save to config data for persistence in prole.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.NAMESPACE"] = ns 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"]["DB_NAME"] = ns state.config_data["Database Creation"]["NAMESPACE"] = ns # 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_prole_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 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 Prole database image...", 0.1) tag = state.controller.get_knoe_db_version() args = ["--tag", tag] if env_key != "dev": # Add registry prefix for non-dev envs 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] args = ["--tag", f"{server}:5000/knoe-db:{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 == "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_prole_data_volume_args prole_data = state.inputs.get( "env_setup.PROLE_DATA", str(inst_config.PROJECT_ROOT / "knoe-db" / "data"), ) volume_args = _k3d_prole_data_volume_args(prole_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-prole-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 rc = state.controller.run_script( "init_openbao.sh", args=["initialize"], env=env_init, stdin_text=f"{db_pw}\n", on_line=_stream_line, ) if rc != 0: self.logger.error(f"OpenBao initialization failed (code {rc})") else: self.logger.info( "OpenBao init skipped: DB password unresolved/anchored. Ensuring OpenBao is running..." ) rc = state.controller.run_script( "init_openbao.sh", args=["start"], env=env, on_line=_stream_line, ) if rc != 0: self.logger.error(f"OpenBao start failed (code {rc})") # 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")) status_args = ["-n", service_ns] if self._parse_bool( state.inputs.get("init_cluster.kerberos_enabled", "False") ): status_args.append("-k") rc_status = state.controller.run_script( "status_common_services.sh", args=status_args, env=env, on_line=_stream_line, ) if rc_status == 0: 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 prole.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 Prole 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("PROLE_MODE", "") 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 prole.cfg data to ensure it's clean self._regenerate_port_mapping_cfg(state, mode or "k3d") scripts = ["init_common_services.sh", "init_cloudnative_pg.sh"] scripts.extend(["init_cnpg_backup.sh", "init_kong.sh"]) if opt_allowed: scripts.append("init_monitoring.sh") else: self.logger.info(f"[SKIP] Monitoring disabled by policy: {opt_reason}") state.config_data.setdefault("Monitoring", {})["STATUS"] = "Skipped" if mode != "k3d": scripts.append("init_nginx_ingress.sh") for i, script in enumerate(scripts): if progress: progress(f"Running {script}...", (i / len(scripts))) # Align args with legacy silent installer behavior args = ["update"] if script == "init_common_services.sh": args = mode_args + [ "-n", env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default")), ] if self._parse_bool( state.inputs.get("kerberos_config.enabled", "False") ): args.append("-k") args.append("update") elif script == "init_cloudnative_pg.sh": args = mode_args + ["initialize"] elif script == "init_cnpg_backup.sh": args = mode_args + ["start"] elif script == "init_kong.sh": args = mode_args + ["start"] elif script in ("init_monitoring.sh", "init_nginx_ingress.sh"): args = mode_args + ["initialize"] rc = state.controller.run_script( script, args=args, env=env, on_line=_stream_line ) if rc != 0: msg = f"Script {script} failed (code {rc})" self.logger.error(msg) raise Exception(msg) # Verify critical secrets ns = 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(f"kubectl get secret {secret} -n {ns}") != 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", "Prole Deployment") self.logger = logging.getLogger("DeploymentMilestone") def execute( self, state: InstallerState, progress: ProgressCallback | None = None ) -> None: if progress: progress("Deploying Prole...", 0.3) env = self._get_script_env(state) mode = env.get("PROLE_MODE", "") mode_args = ["--mode", mode] if mode else [] # Using init_cloudnative_pg.sh as the primary deployment mechanism rc = state.controller.run_script( "init_cloudnative_pg.sh", args=mode_args + ["deploy", "latest"], env=env, on_line=_stream_line, ) if rc == 0: self.logger.info("Deployment successful") else: msg = f"Deployment failed (code {rc})" self.logger.error(msg) raise Exception(msg) if progress: progress("Deployment 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 if progress: progress("Deploying Gitea...", 0.4) env = self._get_script_env(state) ns = ( state.inputs.get("gitops.namespace", "") or env.get("GITEA_NAMESPACE") or "gitea" ).strip() env["GITEA_NAMESPACE"] = ns mode = env.get("PROLE_MODE", "") args: list[str] = [] if mode: args.extend(["--mode", mode]) args.extend(["--namespace", ns]) conf_dir = prole_conf.resolve_prole_conf_dir(state.controller.project_root) cfg_path = prole_conf.entrypoint_path(conf_dir) if cfg_path.exists(): args.extend(["-c", str(cfg_path)]) self.logger.info(f"init_gitea.sh {' '.join(args)}") rc = state.controller.run_script( "init_gitea.sh", 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("Gitea deployed", 1.0) else: gitops_sec = state.config_data.setdefault("GitOps", {}) gitops_sec["STATUS"] = "Attempted" gitops_sec["GITOPS_NAMESPACE"] = ns msg = f"Gitea 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 ) 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("PROLE_MODE", "") 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 = prole_conf.resolve_prole_conf_dir(state.controller.project_root) cfg_path = prole_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 = prole_conf.resolve_prole_conf_dir(state.controller.project_root) cfg_path = prole_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)