mirror of
https://github.com/dredx/prole.git
synced 2026-09-23 11:03:59 +00:00
- make cluster-storage milestone opt-in and remove installer cluster-storage step from UI navigation\n- add cluster storage browser and GKE cluster ops helpers with CLI coverage\n- update k8s/CNPG config and install flow files for corrected cluster setup\n- add/refresh tests for storage browser, GKE ops, prod config, and service-layer navigation Co-authored-by: Junie <junie@jetbrains.com>
1263 lines
48 KiB
Python
1263 lines
48 KiB
Python
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,
|
||
)
|
||
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")
|
||
|
||
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.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 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.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_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
|
||
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("PROLE_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("PROLE_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"))
|
||
status_ok = all(
|
||
[
|
||
registry_ops.status(
|
||
namespace=env.get("REGISTRY_NAMESPACE", service_ns),
|
||
env=env,
|
||
mode=env.get("PROLE_MODE"),
|
||
),
|
||
openbao_ops.status(
|
||
namespace=service_ns,
|
||
env=env,
|
||
mode=env.get("PROLE_MODE"),
|
||
),
|
||
garage_store_ops.status(namespace=service_ns, env=env),
|
||
opentofu_ops.status(namespace=service_ns, env=env),
|
||
]
|
||
)
|
||
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 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")
|
||
|
||
# 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", ".")
|
||
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_store_ops.update(
|
||
namespace=service_ns,
|
||
env=env,
|
||
project_root=project_root,
|
||
mode=mode,
|
||
log=self.logger.info,
|
||
)
|
||
opentofu_ops.update(
|
||
namespace=service_ns,
|
||
env=env,
|
||
project_root=project_root,
|
||
mode=mode,
|
||
log=self.logger.info,
|
||
)
|
||
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"
|
||
# 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}")
|
||
cnpg_initialize(
|
||
namespace=ns,
|
||
cluster_name=cluster_name,
|
||
env=env,
|
||
project_root=project_root,
|
||
log=self.logger.info,
|
||
mode=mode,
|
||
)
|
||
|
||
# Phase 3: remaining scripts (backup, kong, optional ingress)
|
||
post_scripts = ["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"]
|
||
elif script == "init_kong.sh":
|
||
args = mode_args + ["start"]
|
||
else:
|
||
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)
|
||
|
||
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
|
||
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", "")
|
||
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", ".")
|
||
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 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)
|
||
|
||
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
|
||
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"
|
||
|
||
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"{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("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)
|