prole/installer/core/actions.py

2611 lines
122 KiB
Python

"""
Installer actions and unattended workflow helpers.
"""
from __future__ import annotations
import logging
import secrets
import shlex
import socket
import string
import threading
import uuid
from installer.config import _expand_path, _collect_cfg_vars, _expand_cfg_value, _parse_bool, _is_openbao_ref, \
_is_prole_secret, _decrypt_prole_secret, _encrypt_prole_secret, _write_k3s_kubeconfig, _encrypt_cfg_secret
from installer.core.controller import ProleController
from installer.core.env import * # noqa: F401,F403
from installer.core.env import _parse_ollama_host, _bool_str, _deployment_mode_from_env, _openbao_placeholder, \
_looks_like_k8s_bearer_token, _find_kubeconfig_file, _normalize_cluster_env, _pf_upsert_mapping, \
_build_required_port_forwards, _local_registry_enabled, _collect_images_from_files, _http_ping_registry, \
_resolve_supabase_home, _push_docker_image, _detect_ansible_topology, _deployment_target_label, \
_default_opentofu_pipeline_url, _format_ollama_host, _render_prole_cfg, _k3d_prole_data_volume_args, \
_sync_opentofu_pipeline
from installer.core.milestones import (
DependenciesMilestone, NetworkScanMilestone, EnvSetupMilestone,
SecretManagementMilestone, DatabaseCreationMilestone, DockerBuildMilestone,
ClusterLifecycleMilestone, InitializationScriptsMilestone, DeploymentMilestone
)
def _configure_unbuffered_io():
os.environ.setdefault('PYTHONUNBUFFERED', '1')
for stream in (sys.stdout, sys.stderr):
try:
stream.reconfigure(line_buffering=True, write_through=True)
except Exception:
try:
stream.flush()
except Exception:
pass
class ProleSilentInstaller:
"""Console-based unattended installer driven by prole.cfg inputs."""
def __init__(self, controller: ProleController, cfg_path: str | None = None, reset_cluster: bool = False):
self.controller = controller
self.project_root = controller.project_root
self.dependencies = list(inst_config.DEPENDENCIES)
self.cfg_path = self._normalize_cfg_path(cfg_path)
self.reset_cluster = bool(reset_cluster)
self.reset_requested = bool(reset_cluster)
self.inputs = {}
self.prole_cfg_data = {
'Global': {},
'Welcome': {},
'Dependencies': {},
'Network': {},
'Port Forwards': {},
'System Environment': {},
'Kerberos Authentication': {},
'Ollama': {},
'Optional Features': {},
'Database Creation': {},
'Docker Build': {},
'Initialize Cluster': {},
'Initialization Scripts': {},
'Deployment': {},
'Dev Cluster (k3d)': {},
'Service Cluster (k3s)': {},
'Prod Cluster (k8s)': {},
'Install': {}
}
self._cfg_secret_cache = {}
self._secrets_finalized = False
self._db_built_success = False
self._scripts_success = False
self._cnpg_success = False
self._managed_kubeconfig = None
self.docker_import_dir = None
self._repair_ran = False
self._openbao_init_env = None
self._openbao_init_password = None
# ---------------- Logging helpers ----------------
def log(self, msg: str):
print(msg, flush=True)
logging.info(msg)
def err(self, msg: str):
print(msg, file=sys.stderr, flush=True)
logging.error(msg)
# ---------------- Config helpers ----------------
def _normalize_cfg_path(self, raw: str | None) -> Path:
if raw:
p = Path(_expand_path(raw))
if p.is_dir():
return p / 'prole.cfg'
return p
prole_conf = os.environ.get('PROLE_CONF')
if prole_conf:
return Path(_expand_path(prole_conf)) / 'prole.cfg'
return self.project_root / 'conf' / 'prole.cfg'
def _load_inputs_from_cfg(self) -> dict:
cfg = configparser.ConfigParser(interpolation=None)
cfg.optionxform = str
if not self.cfg_path.exists():
raise FileNotFoundError(f"prole.cfg not found at {self.cfg_path}")
cfg.read(self.cfg_path)
cfg_vars = _collect_cfg_vars(cfg)
# Load all sections into prole_cfg_data to maintain idempotency
for section in cfg.sections():
if section not in self.prole_cfg_data:
self.prole_cfg_data[section] = {}
for k, v in cfg.items(section):
self.prole_cfg_data[section][k] = _expand_cfg_value(v, cfg_vars)
if cfg.has_section('Global'):
service_ns = _expand_cfg_value(cfg.get('Global', 'SERVICE_NAMESPACE', fallback=''), cfg_vars).strip()
if service_ns:
self.prole_cfg_data['Global']['SERVICE_NAMESPACE'] = service_ns
if cfg.has_section('Port Forwards'):
self.prole_cfg_data['Port Forwards'] = {
k: _expand_cfg_value(v, cfg_vars) for k, v in cfg.items('Port Forwards')
}
for (section, key), _spec in SECRET_KEY_SPECS.items():
if cfg.has_option(section, key):
val = _expand_cfg_value(cfg.get(section, key, fallback="").strip(), cfg_vars)
if val:
self._cfg_secret_cache[(section, key)] = val
if _is_openbao_ref(val):
self._secrets_finalized = True
self._auto_input_keys = set()
auto_kdc = None
auto_enabled = None
if cfg.has_section('Network'):
sec = cfg['Network']
if 'KDC_AUTO_DETECTED' in sec:
auto_kdc = sec.get('KDC_AUTO_DETECTED', '')
if 'KERBEROS_AUTO_ENABLED' in sec:
auto_enabled = sec.get('KERBEROS_AUTO_ENABLED', '')
def mark_auto(inputs: dict):
if auto_kdc:
if inputs.get('kerberos_config.kdc', '').strip() == str(auto_kdc).strip():
self._auto_input_keys.add('kerberos_config.kdc')
if auto_enabled is not None:
input_enabled = inputs.get('kerberos_config.enabled')
input_bool = _parse_bool(input_enabled, None)
auto_bool = _parse_bool(auto_enabled, None)
if input_bool is not None and auto_bool is not None and input_bool == auto_bool:
self._auto_input_keys.add('kerberos_config.enabled')
inputs: dict[str, str] = {}
if cfg.has_section('Inputs'):
inputs.update({k: _expand_cfg_value(v, cfg_vars) for k, v in cfg.items('Inputs')})
mark_auto(inputs)
# Resolve encrypted/anchored secrets for runtime use
for key in ('init_password.db_password', 'init_password.db_password_confirm', 'kerberos_config.password', 'init_cluster.k3s_token'):
if key in inputs and inputs[key]:
inputs[key] = self._resolve_secret_value(inputs[key])
if inputs.get('init_password.db_password') and not inputs.get('init_password.db_password_confirm'):
inputs['init_password.db_password_confirm'] = inputs['init_password.db_password']
# Allow explicit runtime override of DB password in silent mode.
env_db_pw = os.environ.get("PROLE_DB_PASSWORD") or os.environ.get("DB_PASSWORD")
if env_db_pw:
raw_db_pw = inputs.get('init_password.db_password', '')
if not raw_db_pw or _is_openbao_ref(raw_db_pw) or _is_prole_secret(raw_db_pw):
inputs['init_password.db_password'] = env_db_pw
inputs['init_password.db_password_confirm'] = env_db_pw
# Fill Ollama inputs from config sections if missing
if cfg.has_section('Ollama'):
sec = cfg['Ollama']
if 'ollama_config.server_host' not in inputs and sec.get('OLLAMA_SERVER_HOST', '').strip():
inputs['ollama_config.server_host'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_HOST', ''), cfg_vars)
if 'ollama_config.server_port' not in inputs and sec.get('OLLAMA_SERVER_PORT', '').strip():
inputs['ollama_config.server_port'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_PORT', ''), cfg_vars)
if 'ollama_config.model' not in inputs and sec.get('OLLAMA_MODEL', '').strip():
inputs['ollama_config.model'] = _expand_cfg_value(sec.get('OLLAMA_MODEL', ''), cfg_vars)
if 'ollama_config.server_host' not in inputs and sec.get('OLLAMA_HOST', '').strip():
host, port = _parse_ollama_host(_expand_cfg_value(sec.get('OLLAMA_HOST', ''), cfg_vars))
if host:
inputs['ollama_config.server_host'] = host
if port and 'ollama_config.server_port' not in inputs:
inputs['ollama_config.server_port'] = port
if cfg.has_section('Global'):
sec = cfg['Global']
if 'ollama_config.server_host' not in inputs and sec.get('OLLAMA_SERVER_HOST', '').strip():
inputs['ollama_config.server_host'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_HOST', ''), cfg_vars)
if 'ollama_config.server_port' not in inputs and sec.get('OLLAMA_SERVER_PORT', '').strip():
inputs['ollama_config.server_port'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_PORT', ''), cfg_vars)
if 'ollama_config.model' not in inputs and sec.get('OLLAMA_MODEL', '').strip():
inputs['ollama_config.model'] = _expand_cfg_value(sec.get('OLLAMA_MODEL', ''), cfg_vars)
if 'ollama_config.server_host' not in inputs and sec.get('OLLAMA_HOST', '').strip():
host, port = _parse_ollama_host(_expand_cfg_value(sec.get('OLLAMA_HOST', ''), cfg_vars))
if host:
inputs['ollama_config.server_host'] = host
if port and 'ollama_config.server_port' not in inputs:
inputs['ollama_config.server_port'] = port
return inputs
# Legacy fallback mapping
legacy = {}
if cfg.has_section('System Environment'):
sec = cfg['System Environment']
for k in ('PROLE_HOME', 'PROLE_CONF', 'PROLE_DATA', 'PROLE_LOGS', 'PROLE_SERVICE'):
if k in sec:
legacy[f'env_setup.{k}'] = _expand_cfg_value(sec.get(k, ''), cfg_vars)
if cfg.has_section('Global'):
sec = cfg['Global']
if 'PROLE_HOME' in sec:
legacy['env_setup.PROLE_HOME'] = _expand_cfg_value(sec.get('PROLE_HOME', ''), cfg_vars)
if 'CLUSTER_ENV' in sec:
legacy['init_cluster.cluster_env'] = _expand_cfg_value(sec.get('CLUSTER_ENV', ''), cfg_vars)
if 'PROLE_K3S_SERVER' in sec:
legacy['init_cluster.k3s_server_url'] = _expand_cfg_value(sec.get('PROLE_K3S_SERVER', ''), cfg_vars)
if 'K3S_SERVER_URL' in sec and 'init_cluster.k3s_server_url' not in legacy:
legacy['init_cluster.k3s_server_url'] = _expand_cfg_value(sec.get('K3S_SERVER_URL', ''), cfg_vars)
if 'PROLE_K3S_TOKEN' in sec:
legacy['init_cluster.k3s_token'] = _expand_cfg_value(sec.get('PROLE_K3S_TOKEN', ''), cfg_vars)
if 'K3S_TOKEN' in sec and 'init_cluster.k3s_token' not in legacy:
legacy['init_cluster.k3s_token'] = _expand_cfg_value(sec.get('K3S_TOKEN', ''), cfg_vars)
if 'NAMESPACE' in sec:
legacy['init_password.db_namespace'] = _expand_cfg_value(sec.get('NAMESPACE', ''), cfg_vars)
legacy['env_setup.NAMESPACE'] = _expand_cfg_value(sec.get('NAMESPACE', ''), cfg_vars)
if 'PROLE_DB_USER' in sec:
legacy['init_password.db_username'] = _expand_cfg_value(sec.get('PROLE_DB_USER', ''), cfg_vars)
if 'DB_PASSWORD' in sec:
legacy['init_password.db_password'] = _expand_cfg_value(sec.get('DB_PASSWORD', ''), cfg_vars)
legacy['init_password.db_password_confirm'] = _expand_cfg_value(sec.get('DB_PASSWORD', ''), cfg_vars)
if 'DB_HOST_PORT' in sec:
legacy['init_password.db_host_port'] = _expand_cfg_value(sec.get('DB_HOST_PORT', '5432'), cfg_vars)
if 'DOCKER_IMPORT_DIR' in sec:
self.docker_import_dir = _expand_cfg_value(sec.get('DOCKER_IMPORT_DIR', ''), cfg_vars)
self.prole_cfg_data['Global']['DOCKER_IMPORT_DIR'] = self.docker_import_dir
if 'OLLAMA_SERVER_HOST' in sec:
legacy['ollama_config.server_host'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_HOST', ''), cfg_vars)
if 'OLLAMA_SERVER_PORT' in sec:
legacy['ollama_config.server_port'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_PORT', ''), cfg_vars)
if 'OLLAMA_MODEL' in sec:
legacy['ollama_config.model'] = _expand_cfg_value(sec.get('OLLAMA_MODEL', ''), cfg_vars)
if 'OLLAMA_HOST' in sec and 'ollama_config.server_host' not in legacy:
host, port = _parse_ollama_host(_expand_cfg_value(sec.get('OLLAMA_HOST', ''), cfg_vars))
if host:
legacy['ollama_config.server_host'] = host
if port:
legacy.setdefault('ollama_config.server_port', port)
if cfg.has_section('Network'):
sec = cfg['Network']
if 'KDC_AUTO_DETECTED' in sec:
legacy['kerberos_config.kdc'] = _expand_cfg_value(sec.get('KDC_AUTO_DETECTED', ''), cfg_vars)
if 'KERBEROS_AUTO_ENABLED' in sec:
legacy['kerberos_config.enabled'] = _expand_cfg_value(sec.get('KERBEROS_AUTO_ENABLED', ''), cfg_vars)
if cfg.has_section('Kerberos Authentication'):
sec = cfg['Kerberos Authentication']
for k, tgt in (
('ENABLED', 'kerberos_config.enabled'),
('REALM', 'kerberos_config.realm'),
('KDC', 'kerberos_config.kdc'),
('USER', 'kerberos_config.user'),
('PASSWORD', 'kerberos_config.password'),
):
if k in sec:
legacy[tgt] = _expand_cfg_value(sec.get(k, ''), cfg_vars)
if cfg.has_section('Ollama'):
sec = cfg['Ollama']
if 'OLLAMA_SERVER_HOST' in sec:
legacy['ollama_config.server_host'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_HOST', ''), cfg_vars)
if 'OLLAMA_SERVER_PORT' in sec:
legacy['ollama_config.server_port'] = _expand_cfg_value(sec.get('OLLAMA_SERVER_PORT', ''), cfg_vars)
if 'OLLAMA_MODEL' in sec:
legacy['ollama_config.model'] = _expand_cfg_value(sec.get('OLLAMA_MODEL', ''), cfg_vars)
if 'OLLAMA_HOST' in sec and 'ollama_config.server_host' not in legacy:
host, port = _parse_ollama_host(_expand_cfg_value(sec.get('OLLAMA_HOST', ''), cfg_vars))
if host:
legacy['ollama_config.server_host'] = host
if port:
legacy.setdefault('ollama_config.server_port', port)
if cfg.has_section('Optional Features'):
sec = cfg['Optional Features']
if 'SUPABASE_ENABLED' in sec:
legacy['init_cluster.supabase_enabled'] = sec.get('SUPABASE_ENABLED', '')
if 'KERBEROS_ENABLED' in sec:
legacy['init_cluster.kerberos_enabled'] = sec.get('KERBEROS_ENABLED', '')
legacy['kerberos_config.enabled'] = sec.get('KERBEROS_ENABLED', '')
if 'AT_REST_ENCRYPTION_ENABLED' in sec:
legacy['init_cluster.at_rest_encryption_enabled'] = sec.get('AT_REST_ENCRYPTION_ENABLED', '')
if cfg.has_section('Initialize Cluster'):
sec = cfg['Initialize Cluster']
if 'ENVIRONMENT' in sec:
legacy['init_cluster.cluster_env'] = sec.get('ENVIRONMENT', '')
if 'K3S_SERVER_URL' in sec and 'init_cluster.k3s_server_url' not in legacy:
legacy['init_cluster.k3s_server_url'] = sec.get('K3S_SERVER_URL', '')
if 'K3S_TOKEN' in sec and 'init_cluster.k3s_token' not in legacy:
legacy['init_cluster.k3s_token'] = sec.get('K3S_TOKEN', '')
if cfg.has_section('Database Creation'):
sec = cfg['Database Creation']
if 'DB_NAME' in sec:
legacy['init_password.db_namespace'] = sec.get('DB_NAME', '')
legacy['env_setup.NAMESPACE'] = sec.get('DB_NAME', '')
if 'NAMESPACE' in sec:
legacy['init_password.db_namespace'] = sec.get('NAMESPACE', '')
legacy['env_setup.NAMESPACE'] = sec.get('NAMESPACE', '')
if 'DB_USER' in sec:
legacy['init_password.db_username'] = sec.get('DB_USER', '')
mark_auto(legacy)
for key in ('init_password.db_password', 'init_password.db_password_confirm', 'kerberos_config.password', 'init_cluster.k3s_token'):
if key in legacy and legacy[key]:
legacy[key] = self._resolve_secret_value(legacy[key])
if legacy.get('init_password.db_password') and not legacy.get('init_password.db_password_confirm'):
legacy['init_password.db_password_confirm'] = legacy['init_password.db_password']
return legacy
def _default_inputs(self) -> dict:
namespace = self._initial_namespace()
owner = self._get_local_owner()
env_vals = self._env_defaults(namespace)
inputs: dict[str, str] = {}
# Dependencies
inputs['dependencies.verify_all'] = _bool_str(False)
inputs['dependencies.auto_install_missing'] = _bool_str(DEFAULT_ACTION_FLAGS.get('dependencies.auto_install_missing', True))
for dep in self.dependencies:
inputs[f"dependencies.{dep['id']}.install"] = _bool_str(True)
# Network scan
inputs['network_scan.run'] = _bool_str(DEFAULT_ACTION_FLAGS.get('network_scan.run', True))
# Environment setup
for k in ('PROLE_HOME', 'PROLE_CONF', 'PROLE_DATA', 'PROLE_LOGS', 'PROLE_SERVICE'):
inputs[f'env_setup.{k}'] = env_vals.get(k, '')
inputs['env_setup.NAMESPACE'] = namespace
# Database creation
inputs['init_password.db_namespace'] = namespace
inputs['init_password.db_username'] = owner
inputs['init_password.db_password'] = ''
inputs['init_password.db_password_confirm'] = ''
inputs['init_password.generate_ssh_key'] = _bool_str(DEFAULT_ACTION_FLAGS.get('init_password.generate_ssh_key', True))
# Build DB image
inputs['init_db_build.run_build'] = _bool_str(DEFAULT_ACTION_FLAGS.get('init_db_build.run_build', True))
# Cluster init + optional features
inputs['init_cluster.cluster_env'] = 'dev'
inputs['init_cluster.k3s_server_url'] = ''
inputs['init_cluster.k3s_token'] = ''
inputs['init_cluster.supabase_enabled'] = _bool_str(False)
inputs['init_cluster.kerberos_enabled'] = _bool_str(False)
inputs['init_cluster.at_rest_encryption_enabled'] = _bool_str(True)
inputs['init_cluster.start_cluster'] = _bool_str(DEFAULT_ACTION_FLAGS.get('init_cluster.start_cluster', True))
# Kerberos config
inputs['kerberos_config.enabled'] = _bool_str(False)
inputs['kerberos_config.realm'] = ''
inputs['kerberos_config.kdc'] = ''
inputs['kerberos_config.user'] = ''
inputs['kerberos_config.password'] = ''
inputs['kerberos_config.test_connection'] = _bool_str(DEFAULT_ACTION_FLAGS.get('kerberos_config.test_connection', False))
# Ollama config
inputs['ollama_config.server_host'] = ''
inputs['ollama_config.server_port'] = DEFAULT_OLLAMA_PORT
inputs['ollama_config.model'] = ''
# Init scripts + deploy
inputs['init_scripts.run_scripts'] = _bool_str(DEFAULT_ACTION_FLAGS.get('init_scripts.run_scripts', True))
inputs['init_cnpg_deploy.run_deploy'] = _bool_str(DEFAULT_ACTION_FLAGS.get('init_cnpg_deploy.run_deploy', True))
inputs['init_cnpg_deploy.force_rollout'] = _bool_str(DEFAULT_ACTION_FLAGS.get('init_cnpg_deploy.force_rollout', False))
# Disk selection (installer packaging)
inputs['disk_selection.disk_type'] = 'local'
inputs['disk_selection.removable_mount'] = ''
inputs['disk_selection.local_path'] = str(Path.home())
# Build tools app
inputs['build.deploy_env'] = 'Dev'
inputs['build.run_build'] = _bool_str(DEFAULT_ACTION_FLAGS.get('build.run_build', False))
return inputs
def _get_input(self, key: str, default: str | None = None) -> str:
if key in self.inputs:
return self.inputs[key]
return default if default is not None else ''
def _get_input_bool(self, key: str, default: bool = False) -> bool:
return _parse_bool(self._get_input(key, None), default=default)
def _deployment_mode(self) -> str:
return _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
def _secret_namespace(self) -> str:
ns = (self._get_input('init_password.db_namespace', '') or '').strip()
if not ns:
ns = (self._get_input('env_setup.NAMESPACE', '') or '').strip()
if not ns:
ns = (self.prole_cfg_data.get('Global', {}) or {}).get('NAMESPACE', '').strip()
return ns or 'default'
def _service_namespace(self) -> str:
ns = (self.prole_cfg_data.get('Global', {}) or {}).get('SERVICE_NAMESPACE', '').strip()
if not ns:
ns = (os.environ.get('SERVICE_NAMESPACE') or '').strip()
return ns or 'default'
def _secret_cfg_value(self, section: str, key: str, plaintext: str, leaf: str, bao_key: str) -> str:
if self._secrets_finalized:
return _openbao_placeholder(self._secret_namespace(), leaf, bao_key)
if not plaintext:
return ''
if _is_openbao_ref(plaintext) or _is_prole_secret(plaintext):
return plaintext
# Try to reuse existing encrypted value if it decrypts to the same plaintext
existing = (self.prole_cfg_data.get(section, {}) or {}).get(key)
if existing and _is_prole_secret(existing):
try:
if _decrypt_prole_secret(existing) == plaintext:
return existing
except Exception:
pass
try:
return _encrypt_prole_secret(plaintext)
except Exception:
return plaintext
def _resolve_openbao_ref(self, value: str) -> str:
if not _is_openbao_ref(value):
return value
inner = value[len(OPENBAO_PREFIX):-len(OPENBAO_SUFFIX)]
path, key = (inner.split('#', 1) + [""])[:2]
if not path or not key:
return ""
mount = "kv"
secret_path = path
if "/" in path:
maybe_mount, rest = path.split("/", 1)
if maybe_mount:
mount = maybe_mount
secret_path = rest
token = os.environ.get("OPENBAO_ROOT_TOKEN", "")
if not token:
prole_service = os.environ.get("PROLE_SERVICE")
if prole_service:
token_path = Path(prole_service) / "secrets" / "openbao-root-token"
if token_path.exists():
token = token_path.read_text().strip()
if not token:
return value
url = os.environ.get("PROLE_OPENBAO_URL")
if not url:
for p in ["8200", "18200"]:
try:
import urllib.request
with urllib.request.urlopen(f"http://127.0.0.1:{p}/v1/sys/health", timeout=0.5) as r:
if r.getcode() == 200:
url = f"http://127.0.0.1:{p}"
break
except Exception:
pass
if not url:
url = "http://127.0.0.1:8200"
url = url.rstrip("/")
try:
req = urllib.request.Request(f"{url}/v1/{mount}/data/{secret_path}")
req.add_header("X-Vault-Token", token)
with urllib.request.urlopen(req, timeout=4) as resp:
payload = json.loads(resp.read().decode("utf-8"))
return payload.get("data", {}).get("data", {}).get(key, "") or ""
except Exception:
return value
def _resolve_secret_value(self, value: str) -> str:
if _is_prole_secret(value):
return _decrypt_prole_secret(value)
if _is_openbao_ref(value):
return self._resolve_openbao_ref(value)
return value
def _sanitize_sections_for_cfg(self, sections: dict) -> dict:
sanitized = {k: dict(v) for k, v in sections.items()}
if "Kerberos Authentication" in sanitized:
val = sanitized["Kerberos Authentication"].get("PASSWORD", "")
if val or ("Kerberos Authentication", "PASSWORD") in self._cfg_secret_cache:
sanitized["Kerberos Authentication"]["PASSWORD"] = self._secret_cfg_value(
"Kerberos Authentication", "PASSWORD", val, "kerberos", "password"
)
if "Monitoring" in sanitized:
val = sanitized["Monitoring"].get("GRAFANA_ADMIN_PASSWORD", "")
if val or ("Monitoring", "GRAFANA_ADMIN_PASSWORD") in self._cfg_secret_cache:
sanitized["Monitoring"]["GRAFANA_ADMIN_PASSWORD"] = self._secret_cfg_value(
"Monitoring", "GRAFANA_ADMIN_PASSWORD", val, "monitoring", "grafana_admin_password"
)
return sanitized
def _script_env_for_namespace(self, namespace: str) -> dict:
env = os.environ.copy()
env["PROLE_HOME"] = str(self.project_root)
env["PROLE_SERVICE"] = str(self.project_root)
env["NAMESPACE"] = namespace
service_ns = self._service_namespace()
# Paths
env["PROLE_HOME"] = str(self.project_root)
env["PROLE_CONF"] = self._get_input('env_setup.PROLE_CONF', str(self.project_root / "conf"))
env["PROLE_DATA"] = self._get_input('env_setup.PROLE_DATA', str(self.project_root / "prole-db" / "data"))
env["PROLE_LOGS"] = self._get_input('env_setup.PROLE_LOGS', str(self.project_root / "logs"))
env["PROLE_SERVICE"] = self._get_input('env_setup.PROLE_SERVICE', str(self.project_root / "etc"))
mode = self._deployment_mode()
if mode:
env["PROLE_MODE"] = mode
# Kubeconfig (k3s only)
if mode == "k3s":
kc = Path(env["PROLE_SERVICE"]) / "secrets" / "k3s.kubeconfig"
if not kc.exists():
kc = self.project_root / "etc" / "secrets" / "k3s.kubeconfig"
if kc.exists():
env["KUBECONFIG"] = str(kc)
else:
env.pop("KUBECONFIG", None)
env["SERVICE_NAMESPACE"] = service_ns
db_pw = self._get_input('init_password.db_password', '').strip()
if db_pw:
db_pw = self._resolve_secret_value(db_pw)
env["DB_PASSWORD"] = db_pw
env["OPENTOFU_ADMIN_PASSWORD"] = db_pw
# Add Grafana password
grafana_pw = self.prole_cfg_data.get('Monitoring', {}).get('GRAFANA_ADMIN_PASSWORD', '').strip()
resolved_grafana = self._resolve_secret_value(grafana_pw) if grafana_pw else ""
if resolved_grafana and not (_is_openbao_ref(resolved_grafana) or _is_prole_secret(resolved_grafana)):
env["GRAFANA_ADMIN_PASSWORD"] = resolved_grafana
elif db_pw:
env["GRAFANA_ADMIN_PASSWORD"] = db_pw
# Kerberos
env["KRB5_REALM"] = self._get_input('kerberos_config.realm', '')
env["KRB5_KDC"] = self._get_input('kerberos_config.kdc', '')
env["KRB5_USER"] = self._get_input('kerberos_config.user', '')
krb_pw = self._get_input('kerberos_config.password', '')
if krb_pw:
env["KRB5_PASSWORD"] = self._resolve_secret_value(krb_pw)
else:
env["KRB5_PASSWORD"] = ""
# K3s Auth
if mode == 'k3s':
server = (self._get_input('init_cluster.k3s_server_url', '') or '').strip()
token = (self._get_input('init_cluster.k3s_token', '') or '').strip()
if token:
token = self._resolve_secret_value(token)
if server:
if not server.startswith('http'):
server = f"https://{server}"
env["PROLE_K3S_SERVER"] = server
if token:
env["PROLE_K3S_TOKEN"] = token
explicit_kubeconfig = (env.get("KUBECONFIG") or "").strip()
if explicit_kubeconfig and not os.path.exists(explicit_kubeconfig):
explicit_kubeconfig = ""
managed_set = False
if server and token and _looks_like_k8s_bearer_token(token) and not explicit_kubeconfig:
try:
if self._managed_kubeconfig and os.path.exists(self._managed_kubeconfig):
try:
os.unlink(self._managed_kubeconfig)
except Exception:
pass
self._managed_kubeconfig = str(_write_k3s_kubeconfig(server, token))
env["KUBECONFIG"] = self._managed_kubeconfig
managed_set = True
except Exception as e:
self.err(f"[WARN] Failed to generate managed KUBECONFIG: {e}")
if not managed_set:
fallback = _find_kubeconfig_file(env)
if fallback:
env["KUBECONFIG"] = fallback
# Registry settings
reg_host = self.prole_cfg_data.get('Docker Build', {}).get('LOCAL_REGISTRY')
if reg_host:
env["LOCAL_REGISTRY"] = reg_host
reg_internal = self.prole_cfg_data.get('Docker Build', {}).get('LOCAL_REGISTRY_INTERNAL')
if reg_internal:
env["LOCAL_REGISTRY_INTERNAL"] = reg_internal
return env
def _authority_context_missing(self) -> bool:
enabled = self._get_input_bool('kerberos_config.enabled', False) or \
self._get_input_bool('init_cluster.kerberos_enabled', False)
if not enabled:
return False
candidates = []
env_home = (os.environ.get('PROLE_HOME') or '').strip()
if env_home:
candidates.append(Path(env_home))
candidates.append(self.project_root)
for base in candidates:
try:
if (base / 'authority').is_dir():
return False
if (base / 'prole' / 'authority').is_dir():
return False
except Exception:
continue
return True
def _dashboard_kong_missing(self, env: dict) -> bool:
ns = "kubernetes-dashboard"
try:
res_ns = subprocess.run(
['kubectl', 'get', 'ns', ns, '-o', 'name'],
env=env, capture_output=True, text=True, timeout=6
)
if res_ns.returncode != 0 or ns not in (res_ns.stdout or ''):
return True
res = subprocess.run(
['kubectl', '-n', ns, 'get', 'pods', '--no-headers'],
env=env, capture_output=True, text=True, timeout=6
)
if res.returncode != 0:
return True
kong_lines = [l for l in (res.stdout or '').splitlines() if 'kong' in l.lower()]
if not kong_lines:
return True
for line in kong_lines:
parts = line.split()
if len(parts) < 3:
return True
if parts[2] != 'Running':
return True
return False
except Exception:
return True
def _maybe_run_repair_pipeline(self, env: dict) -> None:
if self._repair_ran:
return
anomalies = []
if self._dashboard_kong_missing(env):
anomalies.append('dashboard')
if self._authority_context_missing():
anomalies.append('authority')
if not anomalies:
return
self._repair_ran = True
service_ns = self._service_namespace()
db_ns = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
self.log(f"[INFO] Detected anomalies ({', '.join(anomalies)}); running repair pipeline...")
self._run_script(
"repair_pipeline.sh",
args=["-n", service_ns, "--db-namespace", db_ns],
env=env
)
# ---------------- Namespace helpers ----------------
def _get_local_owner(self) -> str:
try:
return getpass.getuser()
except Exception:
try:
return os.getlogin()
except Exception:
return "prole"
def _sanitize_namespace(self, name: str) -> str:
cleaned = re.sub(r'[^a-z0-9-]+', '-', (name or '').lower())
cleaned = re.sub(r'-{2,}', '-', cleaned).strip('-')
if not cleaned:
cleaned = 'prole'
if len(cleaned) > 63:
cleaned = cleaned[:63].rstrip('-')
return cleaned
def _generate_namespace_name(self) -> str:
owner = self._sanitize_namespace(self._get_local_owner())
suffix = uuid.uuid4().hex[:6]
base = f"prole-{owner}-{suffix}"
return self._sanitize_namespace(base)
def _ensure_namespace_prefix(self, name: str) -> str:
cleaned = (name or '').strip()
if not cleaned:
return NAMESPACE_PREFIX
if cleaned.startswith(NAMESPACE_PREFIX):
return cleaned
return f"{NAMESPACE_PREFIX}{cleaned}"
def _initial_namespace(self) -> str:
try:
existing = self._read_existing_env()
ns = existing.get('NAMESPACE') or existing.get('PROLE_NAMESPACE')
if ns:
return self._ensure_namespace_prefix(ns)
except Exception:
pass
ns = os.environ.get('NAMESPACE') or os.environ.get('PROLE_NAMESPACE')
if ns:
return self._ensure_namespace_prefix(ns)
return self._ensure_namespace_prefix(self._generate_namespace_name())
def _is_valid_namespace(self, name: str) -> bool:
if not name or len(name) > 63:
return False
return re.match(r'^[a-z0-9]([-a-z0-9]*[a-z0-9])?$', name) is not None
def _generate_db_password(self, length: int = 24) -> str:
alphabet = string.ascii_letters + string.digits
return ''.join(secrets.choice(alphabet) for _ in range(length))
# ---------------- Env helpers ----------------
def _env_defaults(self, namespace: str | None = None) -> dict:
default_home = Path.home() / '.prole'
resolved_home = self._resolve_env_value('PROLE_HOME', str(default_home)) or str(default_home)
home = Path(resolved_home).expanduser()
return {
'PROLE_HOME': str(home),
'PROLE_CONF': self._resolve_env_value('PROLE_CONF', str(home / 'conf')) or str(home / 'conf'),
'PROLE_DATA': self._resolve_env_value('PROLE_DATA', str(home / 'data')) or str(home / 'data'),
'PROLE_LOGS': self._resolve_env_value('PROLE_LOGS', str(home / 'logs')) or str(home / 'logs'),
'PROLE_SERVICE': self._resolve_env_value('PROLE_SERVICE', str(home / 'etc')) or str(home / 'etc'),
'NAMESPACE': namespace or '',
}
def _resolve_env_value(self, key: str, fallback: str | None = None) -> str | None:
val = os.environ.get(key)
if val:
return val
try:
env = self._read_existing_env()
val = env.get(key)
if val:
return val
except Exception:
pass
return fallback
def _resolve_env_dir(self, key: str, default_suffix: str) -> Path:
input_key = f'env_setup.{key}'
val = (self._get_input(input_key, '') or '').strip()
if val:
try:
return Path(val).expanduser()
except Exception:
pass
val = self._resolve_env_value(key)
if val:
try:
return Path(val).expanduser()
except Exception:
pass
return Path.home() / '.prole' / default_suffix
def _read_existing_env(self) -> dict:
env = {}
prole_home = os.environ.get('PROLE_HOME')
candidates = []
if prole_home:
candidates.append(Path(prole_home).expanduser() / 'env.sh')
candidates.append(Path.home() / '.prole' / 'env.sh')
for p in candidates:
try:
if p.exists():
for line in p.read_text().splitlines():
line = line.strip()
if not line or line.startswith('#'):
continue
if line.startswith('export '):
line = line[len('export '):]
if '=' in line:
k, v = line.split('=', 1)
env[k.strip()] = v.strip().strip('"')
break
except Exception:
pass
return env
def _save_env_to_file(self, values: dict):
cluster_env = self._get_input('init_cluster.cluster_env', 'dev')
env_key = _normalize_cluster_env(cluster_env)
home = Path(values['PROLE_HOME']).expanduser()
if env_key == 'dev':
home.mkdir(parents=True, exist_ok=True)
for key in ('PROLE_CONF', 'PROLE_DATA', 'PROLE_LOGS', 'PROLE_SERVICE'):
try:
Path(values[key]).expanduser().mkdir(parents=True, exist_ok=True)
except Exception:
pass
content = []
content.append('#!/usr/bin/env bash')
content.append('# Prole environment configuration')
content.append('# This file is generated by the installer. Source it in new shells, or execute as a wrapper:')
content.append('# "$PROLE_HOME/env.sh" <command> [args…]')
content.append('# shellcheck shell=bash')
for k in ('PROLE_HOME', 'PROLE_CONF', 'PROLE_DATA', 'PROLE_LOGS', 'PROLE_SERVICE'):
content.append(f'export {k}="{values[k]}"')
if values.get('NAMESPACE'):
content.append(f'export NAMESPACE="{values["NAMESPACE"]}"')
content.append('')
content.append('# Ensure PATH works for GUI-launched shells (Docker, etc.)')
content.append('_prole_add_path() { case ":${PATH}:" in *":$1:"*) ;; *) PATH="$1:${PATH:-}" ;; esac; }')
content.append('_prole_add_path "$PROLE_HOME/bin"')
content.append('_prole_add_path "/opt/homebrew/bin"')
content.append('_prole_add_path "/usr/local/bin"')
content.append('_prole_add_path "/usr/bin"')
content.append('_prole_add_path "/bin"')
content.append('_prole_add_path "/usr/sbin"')
content.append('_prole_add_path "/sbin"')
content.append('export PATH')
content.append('')
content.append('# Add custom paths below if needed (examples):')
content.append('')
content.append('# If executed with arguments (and not sourced), run them under this environment')
content.append('if [[ "${BASH_SOURCE[0]}" == "${0}" ]] && [ "$#" -gt 0 ]; then')
content.append(' exec "$@"')
content.append('fi')
new_content = '\n'.join(content) + '\n'
out = home / 'env.sh'
if env_key != 'dev':
out = Path.home() / '.prole' / 'env.sh'
out.parent.mkdir(parents=True, exist_ok=True)
should_write = True
if out.exists():
if out.read_text() == new_content:
should_write = False
self.log(f"[OK] {out} is up to date.")
if should_write:
tmp = out.parent / 'env.sh.tmp'
tmp.write_text(new_content)
tmp.replace(out)
try:
os.chmod(out, 0o755)
except Exception:
pass
self.log(f"[OK] Wrote {out}")
try:
prole_home = Path(values['PROLE_HOME']).expanduser()
prole_service = Path(values['PROLE_SERVICE']).expanduser()
init_pf_src_candidates = [
self.project_root / 'src' / 'prole' / 'etc' / 'init-port-forward.sh',
self.project_root / 'etc' / 'init-port-forward.sh',
]
init_pf_src = next((p for p in init_pf_src_candidates if p.exists()), None)
if init_pf_src is not None:
init_pf_dst = prole_home / 'init-port-forward.sh'
try:
data = init_pf_src.read_bytes()
init_pf_dst.write_bytes(data)
os.chmod(init_pf_dst, 0o755)
except Exception:
pass
etc_src_candidates = [
self.project_root / 'src' / 'prole' / 'etc',
self.project_root / 'etc',
]
etc_src = next((p for p in etc_src_candidates if p.exists()), None)
if etc_src is not None:
try:
prole_service.mkdir(parents=True, exist_ok=True)
if etc_src.resolve() != prole_service.resolve():
if prole_service.exists():
shutil.rmtree(prole_service, ignore_errors=True)
shutil.copytree(etc_src, prole_service)
except Exception:
pass
except Exception:
pass
def reload_env_from_shell(self) -> None:
home = Path(self._get_input('env_setup.PROLE_HOME', str(Path.home() / '.prole')))
env_file = home / 'env.sh'
cmd = (
f'export PROLE_HOME={shlex.quote(str(home))}; '
f'source {shlex.quote(str(env_file))}; '
'env -0'
)
try:
out = subprocess.check_output(['bash', '-lc', cmd])
except Exception as e:
raise Exception(f"Failed to reload environment: {e}")
try:
for raw in out.split(b'\x00'):
if not raw:
continue
kv = raw.decode('utf-8', errors='ignore')
if '=' not in kv:
continue
k, v = kv.split('=', 1)
if k in ('PYTHONPATH', 'PYTHONHOME'):
continue
os.environ[k] = v
except Exception:
pass
def _update_env_namespace(self, namespace: str):
try:
existing = self._read_existing_env()
defaults = self._env_defaults(namespace)
values = {**defaults, **existing}
values['NAMESPACE'] = namespace
self._save_env_to_file(values)
os.environ['NAMESPACE'] = namespace
except Exception as e:
self.err(f"[WARN] Failed to update env.sh namespace: {e}")
# ---------------- Process helpers ----------------
def _run_cmd(self, cmd, cwd=None, env=None, stdin_text=None, on_stdout=None, on_stderr=None) -> int:
stdin_handle = subprocess.PIPE if stdin_text else None
if isinstance(cmd, str):
proc = subprocess.Popen(['bash', '-lc', cmd], cwd=cwd, env=env, stdout=subprocess.PIPE, stderr=subprocess.PIPE, stdin=stdin_handle, text=True, bufsize=1)
else:
proc = subprocess.Popen(cmd, cwd=cwd, env=env, stdout=subprocess.PIPE, stderr=subprocess.PIPE, stdin=stdin_handle, text=True, bufsize=1)
if stdin_text and proc.stdin:
try:
proc.stdin.write(stdin_text)
proc.stdin.close()
except Exception:
pass
def _read_stream(stream, handler, is_err=False):
if not stream:
return
for line in iter(stream.readline, ''):
if handler:
handler(line)
else:
if is_err:
self.err(line.rstrip('\n'))
else:
self.log(line.rstrip('\n'))
err_thread = threading.Thread(target=_read_stream, args=(proc.stderr, on_stderr, True), daemon=True)
err_thread.start()
_read_stream(proc.stdout, on_stdout, False)
rc = proc.wait()
err_thread.join(timeout=2)
return rc
def _run_script(self, script_name: str, args=None, env=None, stdin_text=None, on_line=None) -> int:
def _stdout(line):
self.log(line.rstrip('\n'))
if on_line:
on_line(line)
def _stderr(line):
self.err(line.rstrip('\n'))
mode_args = []
if script_name.startswith("init_"):
mode = _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
if mode:
mode_args = ["--mode", mode]
return self.controller.run_script(
script_name,
args=mode_args + (args or []),
env=env,
stdin_text=stdin_text,
on_line=_stdout,
on_stderr_line=_stderr,
stderr_to_stdout=False,
)
def _add_port_mapping(self, mapping_str):
mode = _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
prefix = "PORT_FORWARD_K3S_MAPPING_" if mode == "k3s" else "PORT_FORWARD_K3D_MAPPING_"
pf_section = self.prole_cfg_data.get('Port Forwards', {})
if pf_section is None:
pf_section = {}
if not _pf_upsert_mapping(pf_section, prefix, mapping_str):
return
self.prole_cfg_data['Port Forwards'] = pf_section
self._write_cfg()
def _process_script_output_line(self, line):
if "GRAFANA_ADMIN_PASSWORD=" in line:
pwd = line.split("GRAFANA_ADMIN_PASSWORD=")[1].strip()
if pwd:
mon = self.prole_cfg_data.get('Monitoring', {})
if mon is None: mon = {}
mon['GRAFANA_ADMIN_PASSWORD'] = pwd
self.prole_cfg_data['Monitoring'] = mon
self._write_cfg()
if "PORT_FORWARD_MAPPING:" in line:
mapping = line.split("PORT_FORWARD_MAPPING:")[1].strip()
if mapping:
self._add_port_mapping(mapping)
def _sync_port_forward_mappings(self):
mode = _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
if not mode:
mode = 'k3d'
prefix = "PORT_FORWARD_K3S_MAPPING_" if mode == "k3s" else "PORT_FORWARD_K3D_MAPPING_"
pf_section = self.prole_cfg_data.get('Port Forwards', {})
if pf_section is None:
pf_section = {}
# Remove deprecated registry port-forward mapping if present.
for k, v in list(pf_section.items()):
if isinstance(v, str) and ("id=registry;" in v or v.strip() == "id=registry"):
del pf_section[k]
service_ns = (self._service_namespace() or '').strip() or 'default'
db_ns = (self._get_input('init_password.db_namespace', '') or '').strip()
if not db_ns:
db_ns = (self.prole_cfg_data.get('Global', {}) or {}).get('NAMESPACE', '').strip()
if not db_ns:
db_ns = (os.environ.get('NAMESPACE') or '').strip()
if not db_ns:
db_ns = 'default'
argocd_ns = (self.prole_cfg_data.get('Global', {}) or {}).get('ARGOCD_NAMESPACE', '').strip()
if not argocd_ns:
argocd_ns = (os.environ.get('ARGOCD_NAMESPACE') or '').strip()
if not argocd_ns:
argocd_ns = 'argocd'
db_host_port = (self._get_input('init_password.db_host_port', '') or '').strip()
if not db_host_port:
db_host_port = (self.prole_cfg_data.get('Global', {}) or {}).get('DB_HOST_PORT', '').strip()
if not db_host_port:
db_host_port = '5432'
supabase_enabled = self._get_input_bool('init_cluster.supabase_enabled', False)
supabase_ns = (os.environ.get('SUPABASE_NAMESPACE') or '').strip()
if not supabase_ns:
supabase_ns = (self.prole_cfg_data.get('Supabase', {}) or {}).get('NAMESPACE', '').strip()
if not supabase_ns:
supabase_ns = 'supabase'
mappings = _build_required_port_forwards(
mode=mode,
service_ns=service_ns,
argocd_ns=argocd_ns,
db_ns=db_ns,
db_host_port=db_host_port,
supabase_enabled=supabase_enabled,
supabase_namespace=supabase_ns,
)
for mapping in mappings:
_pf_upsert_mapping(pf_section, prefix, mapping)
self.prole_cfg_data['Port Forwards'] = pf_section
def _ensure_local_registry_available(self):
"""Ensure a local k3d registry is running and return (host_registry, cluster_registry)."""
if not _local_registry_enabled(self._get_input('init_cluster.cluster_env', '')):
self.log("[SKIP] Local registry disabled; skipping registry setup.")
return None
reg_name = 'prole-registry'
registry_container = f'k3d-{reg_name}'
host_registry = 'localhost:5000'
cluster_registry = f'{registry_container}.localhost:5000'
if not self.controller.check_docker_running():
self.err("[WARN] Docker not running; cannot ensure local registry.")
return None
# 1. If already responsive, just use it
if _http_ping_registry('localhost', 5000):
self.log("[INFO] Local registry already responsive on localhost:5000.")
lst = subprocess.run(['k3d', 'registry', 'list'], capture_output=True, text=True)
if registry_container not in (lst.stdout or '') and reg_name not in (lst.stdout or ''):
cluster_registry = 'localhost:5000'
self.local_registry_url = host_registry
self.local_registry_internal = cluster_registry
self.prole_cfg_data['Docker Build']['LOCAL_REGISTRY'] = host_registry
self.prole_cfg_data['Docker Build']['LOCAL_REGISTRY_INTERNAL'] = cluster_registry
return host_registry, cluster_registry
# 2. Require k3d to manage/create the local registry
k3d_exists = subprocess.run(['which', 'k3d'], capture_output=True).returncode == 0
if not k3d_exists:
self.err("[WARN] k3d not found; cannot ensure local registry.")
return None
lst = subprocess.run(['k3d', 'registry', 'list'], capture_output=True, text=True)
if registry_container not in (lst.stdout or '') and reg_name not in (lst.stdout or ''):
self.log("[INFO] Creating local k3d registry...")
subprocess.run(['k3d', 'registry', 'create', reg_name, '--port', '0.0.0.0:5000'], check=True)
running = subprocess.run(
['docker', 'ps', '--filter', f'name={registry_container}', '--format', '{{.Names}}'],
capture_output=True, text=True
)
if not (running.stdout or '').strip():
exists = subprocess.run(
['docker', 'ps', '-a', '--filter', f'name={registry_container}', '--format', '{{.Names}}'],
capture_output=True, text=True
)
if (exists.stdout or '').strip():
subprocess.run(['docker', 'start', registry_container], check=True)
self.local_registry_url = host_registry
self.local_registry_internal = cluster_registry
self.prole_cfg_data['Docker Build']['LOCAL_REGISTRY'] = host_registry
self.prole_cfg_data['Docker Build']['LOCAL_REGISTRY_INTERNAL'] = cluster_registry
return host_registry, cluster_registry
def _argocd_namespace(self) -> str:
ns = (self.prole_cfg_data.get('Global', {}) or {}).get('ARGOCD_NAMESPACE', '').strip()
if not ns:
ns = (os.environ.get('ARGOCD_NAMESPACE') or '').strip()
return ns or 'argocd'
def _registry_namespace(self) -> str:
ns = (self.prole_cfg_data.get('Global', {}) or {}).get('REGISTRY_NAMESPACE', '').strip()
if not ns:
ns = (os.environ.get('REGISTRY_NAMESPACE') or '').strip()
return ns or 'default'
def _ensure_registry_defaults(self, mode: str) -> None:
db = self.prole_cfg_data.get('Docker Build', {}) or {}
if mode == 'k3d' and _local_registry_enabled(mode):
info = self._ensure_local_registry_available()
if info:
host_registry, cluster_registry = info
db.setdefault('LOCAL_REGISTRY', host_registry)
db.setdefault('LOCAL_REGISTRY_INTERNAL', cluster_registry)
if mode == 'k3s' and _local_registry_enabled(mode):
server = (self._get_input('init_cluster.k3s_server_url', '') or '').strip()
if not server:
server = (self.prole_cfg_data.get('Global', {}) or {}).get('PROLE_K3S_SERVER', '').strip()
host = ''
if server:
host = server.replace('https://', '').replace('http://', '')
host = host.split('/')[0].split(':')[0]
if host:
cur = (db.get('LOCAL_REGISTRY') or '').strip()
if not cur or cur == 'localhost:5000':
db['LOCAL_REGISTRY'] = f"{host}:5000"
cur_internal = (db.get('LOCAL_REGISTRY_INTERNAL') or '').strip()
if not cur_internal:
registry_ns = self._registry_namespace()
db['LOCAL_REGISTRY_INTERNAL'] = f"registry.{registry_ns}.svc.cluster.local:5000"
self.prole_cfg_data['Docker Build'] = db
def _collect_dependent_images(self, include_supabase: bool, include_kerberos_proxy: bool) -> list[str]:
images = set()
for rel_dir in ('k8s/prole', 'k8s/openbao'):
base_dir = self.project_root / rel_dir
if not base_dir.exists():
continue
images.update(_collect_images_from_files(list(base_dir.glob('*.yaml'))))
if include_supabase:
supa_home = _resolve_supabase_home(self.project_root)
if supa_home:
docker_dir = supa_home / 'docker'
compose_files = [docker_dir / 'docker-compose.yml']
if os.environ.get('SUPABASE_USE_DEV_COMPOSE') == '1':
compose_files.append(docker_dir / 'dev' / 'docker-compose.dev.yml')
images.update(_collect_images_from_files(compose_files))
if not include_supabase:
images = {img for img in images if 'supabase' not in img}
if include_kerberos_proxy:
krb_img = os.environ.get('KRB5_AD_PROXY_IMAGE', 'alpine/socat')
if krb_img:
images.add(krb_img)
return sorted(images)
def _prepull_images_to_registry(self, include_supabase: bool, include_kerberos_proxy: bool) -> bool:
if not _local_registry_enabled(self._get_input('init_cluster.cluster_env', '')):
self.log("[SKIP] Local registry disabled; skipping image pre-pull.")
return False
mode = _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
if mode != 'k3d':
self.log("[SKIP] Local registry pre-pull only supported for k3d.")
return False
info = self._ensure_local_registry_available()
if not info:
self.err("[WARN] Local registry unavailable; skipping image pre-pull.")
return False
registry, _cluster_registry = info
images = self._collect_dependent_images(include_supabase, include_kerberos_proxy)
if not images:
self.log("[INFO] No dependent images found to pre-pull.")
return True
overall_ok = True
import_dir = self.docker_import_dir
for image in images:
local_tag = image
if not image.startswith(f"{registry}/"):
local_tag = f"{registry}/{image}"
# 1. Check if already in local registry
self.log(f"[INFO] Checking if {image} exists in local registry...")
check_reg = subprocess.run(['docker', 'pull', local_tag], capture_output=True, text=True)
if check_reg.returncode == 0:
self.log(f"[OK] {image} already exists in local registry as {local_tag}")
continue
# 2. Check if we already have it in local docker daemon
check_local = subprocess.run(['docker', 'image', 'inspect', image], capture_output=True, text=True)
found = (check_local.returncode == 0)
if not found and import_dir and os.path.isdir(import_dir):
# 3. Check import directory
safe_name = image.replace("/", "_").replace(":", "_")
tar_path = Path(import_dir) / f"{safe_name}.tar"
if tar_path.exists():
self.log(f"[INFO] Found {tar_path} in import directory, loading...")
load = subprocess.run(['docker', 'load', '-i', str(tar_path)], capture_output=True, text=True)
if load.returncode == 0:
found = True
else:
self.err(f"[WARN] Failed to load {tar_path}: {load.stderr}")
if not found:
# 4. Pull from Docker Hub
self.log(f"[INFO] Pulling {image} from Docker Hub...")
pull = subprocess.run(['docker', 'pull', image], capture_output=True, text=True)
if pull.returncode != 0:
overall_ok = False
self.err(pull.stdout or '')
self.err(pull.stderr or '')
self.err(f"[ERROR] docker pull failed for {image}")
continue
found = True
# If we have the image locally, tag and push to local registry
if found:
if local_tag != image:
tag = subprocess.run(['docker', 'tag', image, local_tag], capture_output=True, text=True)
if tag.returncode != 0:
overall_ok = False
self.err(tag.stdout or '')
self.err(tag.stderr or '')
self.err(f"[ERROR] docker tag failed for {image} to {local_tag}")
continue
self.log(f"[INFO] Pushing {local_tag} to local registry...")
if not _push_docker_image(local_tag, log_fn=self.log):
overall_ok = False
continue
self.log(f"[OK] Stored {image} in local registry as {local_tag}")
return overall_ok
def _export_docker_image(self, image_name: str, export_dir: Path) -> Path:
export_dir.mkdir(parents=True, exist_ok=True)
safe_name = image_name.replace("/", "_").replace(":", "_")
tar_path = export_dir / f"{safe_name}.tar"
if tar_path.exists():
self.log(f"[INFO] Reusing image tarball {tar_path}")
return tar_path
self.log(f"[INFO] Exporting {image_name} to {tar_path}")
rc = self._run_cmd(['docker', 'save', '-o', str(tar_path), image_name])
if rc != 0:
raise Exception(f"Failed to export image {image_name} (code {rc})")
return tar_path
def _import_k3s_image(self, image_name: str, force: bool = False) -> bool:
env_key = _normalize_cluster_env(self._get_input('init_cluster.cluster_env', 'dev'))
if env_key != 'service':
return True
if _local_registry_enabled(env_key) and not force:
self.log("[INFO] Local registry enabled; skipping k3s image import.")
return True
inspect = subprocess.run(['docker', 'image', 'inspect', image_name], capture_output=True)
if inspect.returncode != 0:
self.err(f"[WARN] {image_name} not found locally; cannot import to k3s.")
return False
export_dir = None
if self.docker_import_dir:
export_dir = Path(self.docker_import_dir).expanduser()
elif Path("/private/tmp").exists():
export_dir = Path("/private/tmp")
else:
export_dir = Path("/tmp")
tar_path = self._export_docker_image(image_name, export_dir)
ansible_env = os.environ.copy()
ansible_cfg = self.project_root / "ansible.cfg"
if ansible_cfg.exists():
ansible_env["ANSIBLE_CONFIG"] = str(ansible_cfg)
vault_file = os.environ.get("ANSIBLE_VAULT_PASSWORD_FILE")
if not vault_file:
candidate = self.project_root / ".vault_pass"
if candidate.exists():
vault_file = str(candidate)
host_group = "k3s_hosts"
ansible_base = ["ansible", host_group]
if vault_file:
ansible_base += ["--vault-password-file", vault_file]
dest_path = f"/tmp/{tar_path.name}"
self.log(f"[INFO] Importing {image_name} into k3s nodes (group: {host_group})")
rc = self._run_cmd(
ansible_base + ["-m", "copy", "-a", f"src={tar_path} dest={dest_path} mode=0644"],
cwd=str(self.project_root),
env=ansible_env,
)
if rc != 0:
self.err(f"[WARN] Failed to copy {tar_path} to k3s nodes (code {rc})")
return False
rc = self._run_cmd(
ansible_base + ["-m", "shell", "-a", f"ctr -n k8s.io -a /run/k3s/containerd/containerd.sock images import --no-unpack {dest_path}"],
cwd=str(self.project_root),
env=ansible_env,
)
if rc != 0:
self.err(f"[WARN] Failed to import {image_name} into k3s nodes (code {rc})")
return False
return True
# ---------------- Steps ----------------
def _apply_inputs(self):
defaults = self._default_inputs()
loaded = self._load_inputs_from_cfg()
self.inputs = {**defaults, **loaded}
try:
self._apply_ansible_defaults(set(loaded.keys()))
except Exception:
pass
mode = _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
if not mode:
mode = _deployment_mode_from_env(self._get_input('init_cluster.mode', ''))
if mode:
os.environ['PROLE_MODE'] = mode
for k in ('env_setup.PROLE_HOME', 'env_setup.PROLE_CONF', 'env_setup.PROLE_DATA', 'env_setup.PROLE_LOGS', 'env_setup.PROLE_SERVICE'):
if k in self.inputs:
self.inputs[k] = _expand_path(self.inputs[k])
def _apply_ansible_defaults(self, loaded_keys: set[str] | None = None):
info = _detect_ansible_topology(self.project_root)
if not info:
return
loaded_keys = loaded_keys or set()
auto_keys = getattr(self, '_auto_input_keys', set())
ansible_kdc = info.get('kdc_ip') or ''
try:
net = self.prole_cfg_data.get('Network', {})
if info.get('topology_json'):
net['ANSIBLE_TOPOLOGY'] = info['topology_json']
if info.get('inventory_path'):
net['ANSIBLE_INVENTORY'] = info['inventory_path']
if info.get('infrastructure_path'):
net['ANSIBLE_INFRASTRUCTURE'] = info['infrastructure_path']
if info.get('domain'):
net['ANSIBLE_DOMAIN'] = info['domain']
if info.get('realm'):
net['ANSIBLE_REALM'] = info['realm']
if info.get('ad_dc_host'):
net['AD_DC_HOST'] = info['ad_dc_host']
if info.get('ad_dc_ip'):
net['AD_DC_IP'] = info['ad_dc_ip']
if info.get('kdc_ip'):
net['KDC_ANSIBLE_DETECTED'] = info['kdc_ip']
self.prole_cfg_data['Network'] = net
except Exception:
pass
if ansible_kdc and 'kerberos_config.kdc' in auto_keys:
self.inputs['kerberos_config.kdc'] = ansible_kdc
if ansible_kdc and 'kerberos_config.enabled' in auto_keys:
self.inputs['kerberos_config.enabled'] = _bool_str(True)
if 'kerberos_config.kdc' not in loaded_keys:
if not self._get_input('kerberos_config.kdc', '').strip() and info.get('kdc_ip'):
self.inputs['kerberos_config.kdc'] = info['kdc_ip']
if 'kerberos_config.realm' not in loaded_keys:
if not self._get_input('kerberos_config.realm', '').strip() and info.get('realm'):
self.inputs['kerberos_config.realm'] = info['realm']
if 'kerberos_config.enabled' not in loaded_keys:
if info.get('kdc_ip') and not self._get_input_bool('kerberos_config.enabled', False):
self.inputs['kerberos_config.enabled'] = _bool_str(True)
if 'init_cluster.k3s_server_url' not in loaded_keys:
if not self._get_input('init_cluster.k3s_server_url', '').strip() and info.get('k3s_server_url'):
self.inputs['init_cluster.k3s_server_url'] = info['k3s_server_url']
if 'init_cluster.k3s_token' not in loaded_keys:
if not self._get_input('init_cluster.k3s_token', '').strip() and info.get('k3s_token'):
self.inputs['init_cluster.k3s_token'] = info['k3s_token']
def _write_cfg(self):
inputs = dict(self.inputs)
db_pw = self._get_input('init_password.db_password', '')
db_pw_cfg = self._secret_cfg_value('Inputs', 'init_password.db_password', db_pw, 'db', 'password')
if db_pw_cfg:
inputs['init_password.db_password'] = db_pw_cfg
inputs['init_password.db_password_confirm'] = db_pw_cfg
elif db_pw:
inputs['init_password.db_password'] = db_pw
inputs['init_password.db_password_confirm'] = self._get_input('init_password.db_password_confirm', '') or db_pw
krb_pw = self._get_input('kerberos_config.password', '')
if krb_pw:
inputs['kerberos_config.password'] = self._secret_cfg_value('Inputs', 'kerberos_config.password', krb_pw, 'kerberos', 'password')
k3s_token = self._get_input('init_cluster.k3s_token', '')
if k3s_token:
inputs['init_cluster.k3s_token'] = self._secret_cfg_value('Inputs', 'init_cluster.k3s_token', k3s_token, 'k3s', 'token')
mode = _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
target_label = _deployment_target_label(self._get_input('init_cluster.cluster_env', ''))
globals_to_save = {
'PROLE_HOME': self._get_input('env_setup.PROLE_HOME', ''),
'PROLE_DB_USER': self._get_input('init_password.db_username', ''),
'DB_PASSWORD': self._secret_cfg_value('Global', 'DB_PASSWORD', db_pw, 'db', 'password'),
'CLUSTER_ENV': self._get_input('init_cluster.cluster_env', ''),
'DEPLOYMENT_MODE': mode,
'DEPLOYMENT_TARGET': target_label,
'NAMESPACE': (self._get_input('init_password.db_namespace', '') or '').strip(),
'DB_HOST_PORT': (self._get_input('init_password.db_host_port', '5432') or '5432').strip(),
'DOCKER_IMPORT_DIR': self.docker_import_dir or '',
'PROLE_K3S_SERVER': (self._get_input('init_cluster.k3s_server_url', '') or '').strip(),
'PROLE_K3S_TOKEN': self._secret_cfg_value('Global', 'PROLE_K3S_TOKEN', self._get_input('init_cluster.k3s_token', ''), 'k3s', 'token'),
'PROLE_OPENTOFU_URL': _default_opentofu_pipeline_url(),
'SERVICE_NAMESPACE': self._service_namespace(),
}
globals_to_save.update(self.prole_cfg_data.get('Global', {}))
globals_to_save['DB_PASSWORD'] = self._secret_cfg_value('Global', 'DB_PASSWORD', db_pw, 'db', 'password')
globals_to_save['DEPLOYMENT_MODE'] = mode
globals_to_save['DEPLOYMENT_TARGET'] = target_label
globals_to_save.setdefault('ARGOCD_NAMESPACE', self._argocd_namespace())
globals_to_save.setdefault('REGISTRY_NAMESPACE', self._registry_namespace())
ollama_host = (self._get_input('ollama_config.server_host', '') or '').strip()
ollama_port = (self._get_input('ollama_config.server_port', DEFAULT_OLLAMA_PORT) or '').strip()
ollama_model = (self._get_input('ollama_config.model', '') or '').strip()
ollama_section = dict(self.prole_cfg_data.get('Ollama', {}))
if ollama_host:
if not ollama_port:
ollama_port = DEFAULT_OLLAMA_PORT
ollama_section['OLLAMA_SERVER_HOST'] = ollama_host
ollama_section['OLLAMA_SERVER_PORT'] = ollama_port
ollama_section['OLLAMA_HOST'] = _format_ollama_host(ollama_host, ollama_port)
else:
for key in ('OLLAMA_SERVER_HOST', 'OLLAMA_SERVER_PORT', 'OLLAMA_HOST'):
ollama_section.pop(key, None)
if ollama_model:
ollama_section['OLLAMA_MODEL'] = ollama_model
else:
ollama_section.pop('OLLAMA_MODEL', None)
if ollama_section:
self.prole_cfg_data['Ollama'] = ollama_section
self._sync_port_forward_mappings()
sections = {k: self.prole_cfg_data.get(k, {}) for k in [
'Welcome', 'Dependencies', 'Network', 'Port Forwards', 'System Environment',
'Monitoring', 'Kerberos Authentication', 'Ollama', 'Optional Features', 'Database Creation',
'Initialize Cluster', 'Docker Build', 'Initialization Scripts',
'Deployment', 'Install'
]}
deployment_section = dict(sections.get('Deployment', {}))
if mode:
deployment_section.setdefault('MODE', mode)
if target_label:
deployment_section.setdefault('TARGET', target_label)
sections['Deployment'] = deployment_section
sections['Dev Cluster (k3d)'] = {
**self.prole_cfg_data.get('Dev Cluster (k3d)', {}),
'MODE': 'k3d',
'CLUSTER_ENV': 'k3d-prole-dev-cluster',
'DISPLAY_NAME': 'prole-dev-cluster',
'KUBECTL_CONTEXT': self._get_input('init_cluster.cluster_env', '')
}
sections['Service Cluster (k3s)'] = {
**self.prole_cfg_data.get('Service Cluster (k3s)', {}),
'MODE': 'k3s',
'CLUSTER_ENV': 'prole-service-cluster',
'DISPLAY_NAME': 'prole-service-cluster',
'K3S_SERVER_URL': (self._get_input('init_cluster.k3s_server_url', '') or '').strip(),
'K3S_TOKEN': self._secret_cfg_value('Service Cluster (k3s)', 'K3S_TOKEN', self._get_input('init_cluster.k3s_token', ''), 'k3s', 'token'),
'PIPELINE_URL': _default_opentofu_pipeline_url()
}
sections['Prod Cluster (k8s)'] = {
**self.prole_cfg_data.get('Prod Cluster (k8s)', {}),
'MODE': 'k8s',
'CLUSTER_ENV': 'prole-prod-cluster',
'DISPLAY_NAME': 'prole-prod-cluster',
'ARTIFACTS_DIR': (self._get_input('init_cluster.prod_artifacts_path', '') or '').strip(),
'PIPELINE_URL': _default_opentofu_pipeline_url()
}
sections = self._sanitize_sections_for_cfg(sections)
# Check for existing config and timestamp to maintain idempotency
existing_text = None
existing_timestamp = None
if self.cfg_path.exists():
existing_text = self.cfg_path.read_text()
for line in existing_text.splitlines():
if line.startswith('; Generated by install.py on '):
existing_timestamp = line[len('; Generated by install.py on '):].strip()
break
cfg_text = _render_prole_cfg(inputs, globals_to_save, sections, generated_at=existing_timestamp)
should_write = True
if existing_text:
if existing_text == cfg_text:
should_write = False
if should_write:
# If content changed but we reused existing_timestamp, it might be misleading.
# But usually we want to know when it LAST changed.
# If it's different, let's generate a new one.
if existing_text and existing_text != cfg_text:
cfg_text = _render_prole_cfg(inputs, globals_to_save, sections)
self.cfg_path.parent.mkdir(parents=True, exist_ok=True)
self.cfg_path.write_text(cfg_text)
self.log(f"[CONFIG] Wrote {self.cfg_path}")
def _step_dependencies(self) -> bool:
self.log("==> Dependencies")
missing = []
for dep in self.dependencies:
ok, location, version = inst_config.get_dep_info(dep)
if ok:
self.log(f"[OK] {dep['name']} {version or ''}".strip())
continue
missing.append(dep)
self.log(f"[MISSING] {dep['name']}")
if not missing:
self.prole_cfg_data['Dependencies']['STATUS'] = 'All installed'
return True
auto_install = self._get_input_bool('dependencies.auto_install_missing', DEFAULT_ACTION_FLAGS.get('dependencies.auto_install_missing', True))
if not auto_install:
self.err("[ERROR] Dependencies missing and auto-install disabled.")
self.prole_cfg_data['Dependencies']['STATUS'] = 'Missing'
return False
for dep in missing:
install_cmd = dep.get('install_cmd')
if not install_cmd:
self.err(f"[ERROR] No install command for {dep['name']}.")
continue
should_install = self._get_input_bool(f"dependencies.{dep['id']}.install", True)
if not should_install:
self.err(f"[SKIP] {dep['name']} install disabled by config.")
continue
self.log(f"[INSTALL] {dep['name']} -> {install_cmd}")
rc = self._run_cmd(install_cmd)
if rc != 0:
self.err(f"[ERROR] Install failed for {dep['name']} (code {rc})")
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.err(f"[ERROR] Still missing: {', '.join(still_missing)}")
self.prole_cfg_data['Dependencies']['STATUS'] = 'Missing'
return False
self.prole_cfg_data['Dependencies']['STATUS'] = 'All installed'
return True
def _step_network_scan(self) -> None:
if not self._get_input_bool('network_scan.run', DEFAULT_ACTION_FLAGS.get('network_scan.run', True)):
self.log("[SKIP] Network scan disabled.")
return
existing_kdc = self.prole_cfg_data.get('Network', {}).get('KDC_AUTO_DETECTED')
if existing_kdc:
self.log(f"[OK] Network scan already performed. KDC detected: {existing_kdc}")
return
self.log("==> Network scan (prole-agent)")
ansible_kdc = ''
try:
ansible_kdc = (self.prole_cfg_data.get('Network', {}) or {}).get('KDC_ANSIBLE_DETECTED', '')
except Exception:
ansible_kdc = ''
scan_binary = get_resource_path("prole-net/prole-agent")
if not scan_binary.exists():
self.err(f"[ERROR] Scan binary not found at {scan_binary}")
return
prole_home = Path.home() / ".prole"
scan_dir = prole_home / "scan"
scan_dir.mkdir(parents=True, exist_ok=True)
kdc_found = None
def _handle_stdout(line: str):
nonlocal kdc_found
self.log(line.rstrip('\n'))
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:
for part in line.split():
try:
socket.inet_aton(part.strip('[]():,'))
kdc_found = part.strip('[]():,')
break
except Exception:
continue
rc = self._run_cmd([str(scan_binary), "-t", "10"], cwd=str(scan_dir), on_stdout=_handle_stdout)
if rc == 0:
self.log("[OK] Scan complete.")
else:
self.err(f"[ERROR] Scan failed (code {rc})")
if kdc_found and not ansible_kdc:
self.inputs['kerberos_config.kdc'] = kdc_found
self.inputs['kerberos_config.enabled'] = _bool_str(True)
self.prole_cfg_data['Network']['KDC_AUTO_DETECTED'] = kdc_found
self.prole_cfg_data['Network']['KERBEROS_AUTO_ENABLED'] = 'True'
elif kdc_found and ansible_kdc:
self.log(f"[INFO] Ansible KDC detected ({ansible_kdc}); ignoring scan-detected KDC {kdc_found}.")
def _step_env_setup(self) -> None:
self.log("==> Environment setup")
vals = {
'PROLE_HOME': self._get_input('env_setup.PROLE_HOME', str(Path.home() / '.prole')),
'PROLE_CONF': self._get_input('env_setup.PROLE_CONF', ''),
'PROLE_DATA': self._get_input('env_setup.PROLE_DATA', ''),
'PROLE_LOGS': self._get_input('env_setup.PROLE_LOGS', ''),
'PROLE_SERVICE': self._get_input('env_setup.PROLE_SERVICE', ''),
'PROLE_OPENTOFU_URL': (os.environ.get('PROLE_OPENTOFU_URL') or '').strip(),
}
defaults = self._env_defaults(self._get_input('env_setup.NAMESPACE', ''))
for k in vals:
if not vals[k]:
vals[k] = defaults.get(k, '')
if not vals.get('PROLE_OPENTOFU_URL'):
vals['PROLE_OPENTOFU_URL'] = _default_opentofu_pipeline_url()
vals['NAMESPACE'] = self._get_input('env_setup.NAMESPACE', '')
if not vals.get('PROLE_HOME'):
raise Exception("PROLE_HOME is required for env setup.")
# Persist normalized inputs
for k in ('PROLE_HOME', 'PROLE_CONF', 'PROLE_DATA', 'PROLE_LOGS', 'PROLE_SERVICE'):
self.inputs[f'env_setup.{k}'] = vals[k]
self.inputs['env_setup.NAMESPACE'] = vals.get('NAMESPACE', '')
self._save_env_to_file(vals)
self.reload_env_from_shell()
for k in ('PROLE_HOME', 'PROLE_CONF', 'PROLE_DATA', 'PROLE_LOGS', 'PROLE_SERVICE', 'PROLE_OPENTOFU_URL'):
self.prole_cfg_data['System Environment'][k] = vals[k]
def _step_init_password(self) -> None:
self.log("==> Database creation")
ns = (self._get_input('init_password.db_namespace', '') or '').strip()
if not ns:
ns = self._get_input('env_setup.NAMESPACE', '')
if not ns:
raise Exception("Database namespace cannot be empty.")
if not self._is_valid_namespace(ns):
raise Exception("Invalid database namespace.")
user = (self._get_input('init_password.db_username', '') or '').strip()
if not user:
import getpass
try:
user = getpass.getuser()
except Exception:
user = "prole"
self.log(f"[INFO] No database owner provided; using default: {user}")
p1 = self._get_input('init_password.db_password', '')
p2 = self._get_input('init_password.db_password_confirm', '') or p1
if not p1:
p1 = self._generate_db_password()
p2 = p1
self.log("[INFO] No database password provided; generated default password.")
if p1 != p2:
raise Exception("Database passwords do not match.")
# If password is an OpenBao anchor, only skip when OpenBao is already initialized.
if p1.startswith('${'):
resolved = self._resolve_secret_value(p1)
if resolved and not resolved.startswith('${'):
p1 = resolved
p2 = resolved
else:
service_dir = self._get_input('env_setup.PROLE_SERVICE', '').strip()
if not service_dir:
service_dir = self._env_defaults(ns).get('PROLE_SERVICE', '')
token_path = Path(service_dir) / 'secrets' / 'openbao-root-token'
local_token_path = Path.home() / '.prole' / 'etc' / 'secrets' / 'openbao-root-token'
has_token = bool(os.environ.get("OPENBAO_ROOT_TOKEN")) or token_path.exists() or local_token_path.exists()
if self.reset_requested or not has_token:
p1 = self._generate_db_password()
p2 = p1
self.log("[RESET] OpenBao anchor detected without token; generated new DB password for initialization.")
else:
self.log(f"[OK] OpenBao secrets already anchored for namespace {ns}. Skipping initialization.")
return
self.inputs['init_password.db_namespace'] = ns
self.inputs['env_setup.NAMESPACE'] = ns
self.inputs['init_password.db_password'] = p1
self.inputs['init_password.db_password_confirm'] = p2
self._update_env_namespace(ns)
self.prole_cfg_data['Database Creation']['DB_USER'] = user
self.prole_cfg_data['Database Creation']['DB_PASSWORD_SET'] = 'true'
self.prole_cfg_data['Database Creation']['DB_NAME'] = ns
self.prole_cfg_data['Database Creation']['NAMESPACE'] = ns
include_supabase = self._get_input_bool('init_cluster.supabase_enabled', False)
include_kerberos = self._get_input_bool('init_cluster.kerberos_enabled', False) or self._get_input_bool('kerberos_config.enabled', False)
try:
self._prepull_images_to_registry(include_supabase, include_kerberos)
except Exception as e:
self.err(f"[WARN] Image pre-pull error: {e}")
if self._get_input_bool('init_password.generate_ssh_key', DEFAULT_ACTION_FLAGS.get('init_password.generate_ssh_key', True)):
self._generate_ssh_key(user)
# Initialize OpenBao and store keys/passwords for this namespace.
# OpenBao is expected to run inside the cluster; defer initialization until cluster is ready.
env = os.environ.copy()
env["PROLE_HOME"] = str(self.project_root)
env["PROLE_SERVICE"] = str(self.project_root)
env["PROLE_DB_USER"] = user
env["DB_PASSWORD"] = p1
env["GRAFANA_ADMIN_PASSWORD"] = p1
env["NAMESPACE"] = ns
env["AT_REST_ENCRYPTION_ENABLED"] = _bool_str(self._get_input_bool('init_cluster.at_rest_encryption_enabled', False))
realm = self._get_input('kerberos_config.realm', '').strip()
kdc = self._get_input('kerberos_config.kdc', '').strip()
krb_user = self._get_input('kerberos_config.user', '').strip()
krb_pw = self._get_input('kerberos_config.password', '').strip()
if realm:
env["KRB5_REALM"] = realm
env["REALM"] = realm
env["DOMAIN"] = realm.lower()
if kdc:
env["KRB5_KDC"] = kdc
env["KRB5_ADMIN"] = kdc
if krb_user:
env["KRB5_USER"] = krb_user
if krb_pw:
env["KRB5_PASSWORD"] = krb_pw
self._openbao_init_env = env
self._openbao_init_password = p1
self.log("[INFO] Deferring OpenBao initialization until cluster is ready.")
def _generate_ssh_key(self, user: str):
self.log("==> Environment Preparation (Secure Access)")
key_path = Path.home() / ".ssh" / "id_prole_ed25519"
key_path.parent.mkdir(parents=True, exist_ok=True)
if key_path.exists():
self.log(f"[SKIP] Secure access keys already exist at {key_path}")
return
cmd = ["ssh-keygen", "-t", "ed25519", "-N", "", "-f", str(key_path), "-C", user]
rc = self._run_cmd(cmd)
if rc != 0:
self.err(f"[WARN] Preparation failed (code {rc}), trying alternative.")
cmd = ["ssh-keygen", "-t", "rsa", "-b", "4096", "-N", "", "-f", str(key_path), "-C", user]
rc = self._run_cmd(cmd)
if rc != 0:
self.err(f"[ERROR] Preparation failed (code {rc})")
def _step_db_build(self) -> None:
if not self._get_input_bool('init_db_build.run_build', DEFAULT_ACTION_FLAGS.get('init_db_build.run_build', True)):
self.log("[SKIP] DB build disabled.")
self.prole_cfg_data['Docker Build']['STATUS'] = 'Skipped'
return
tag = self.controller.get_prole_db_version()
image_name = f"prole-db:{tag}"
env_key = _normalize_cluster_env(self._get_input('init_cluster.cluster_env', 'dev'))
cluster_name = "prole-dev-cluster"
# Check if image already exists
rc_inspect = subprocess.run(['docker', 'inspect', image_name], capture_output=True).returncode
if rc_inspect == 0:
self.log(f"[OK] {image_name} already exists. Skipping build.")
if env_key == 'dev':
self.log(f"Importing image to {cluster_name}...")
self._run_cmd(['k3d', 'image', 'import', image_name, '-c', cluster_name])
elif env_key == 'service':
try:
self._import_k3s_image(image_name)
except Exception as e:
self.err(f"[WARN] k3s image import failed: {e}")
self._db_built_success = True
self.prole_cfg_data['Docker Build']['STATUS'] = 'Built'
return
self.log("==> Build prole-db image")
prole_home = Path.home() / ".prole"
build_dir = prole_home / "build" / "prole-db"
build_dir.mkdir(parents=True, exist_ok=True)
source_dir = get_resource_path("prole-db")
if source_dir.exists():
if source_dir.resolve() != build_dir.resolve():
if build_dir.exists():
shutil.rmtree(build_dir)
shutil.copytree(source_dir, build_dir)
pub_key_path = Path.home() / ".ssh" / "id_prole_ed25519.pub"
pub_key = pub_key_path.read_text().strip() if pub_key_path.exists() else ""
username = self._get_input('init_password.db_username', '')
cmd = ['docker', 'build']
cmd.extend(get_docker_build_platform_args(env_key))
cmd += [
'--build-arg', f"PROLE_USER={username}",
'--build-arg', f"PROLE_SSH_PUB_KEY={pub_key}",
'-t', image_name, '.'
]
rc = self._run_cmd(cmd, cwd=str(build_dir))
if rc == 0:
self._db_built_success = True
self.log("[OK] Build successful.")
push_ok = True
if env_key == 'dev':
cluster_name = "prole-dev-cluster"
self.log(f"Importing image to {cluster_name}...")
self._run_cmd(['k3d', 'image', 'import', image_name, '-c', cluster_name])
else:
# For non-dev, tag and push to registry
db_cfg = self.prole_cfg_data.get('Docker Build', {}) or {}
registry = (db_cfg.get('LOCAL_REGISTRY') or '').strip()
if (not registry or registry == 'localhost:5000') and env_key == 'service':
server = (self._get_input('init_cluster.k3s_server_url', '') or '').strip()
if not server:
server = (self.prole_cfg_data.get('Global', {}) or {}).get('PROLE_K3S_SERVER', '').strip()
if server:
host = server.replace('https://', '').replace('http://', '')
host = host.split('/')[0].split(':')[0]
if host:
registry = f"{host}:5000"
db_cfg['LOCAL_REGISTRY'] = registry
self.prole_cfg_data['Docker Build'] = db_cfg
if not registry:
registry = 'localhost:5000'
remote_tag = f"{registry}/{image_name}"
self.log(f"Tagging {image_name} as {remote_tag}...")
subprocess.run(['docker', 'tag', image_name, remote_tag], check=True)
self.log(f"Pushing {remote_tag} to registry...")
if not _push_docker_image(remote_tag, log_fn=self.log):
push_ok = False
self.err(f"[ERROR] Failed to push image {remote_tag}")
self._db_built_success = False
if env_key == 'service':
try:
self._import_k3s_image(image_name, force=not push_ok)
except Exception as e:
self.err(f"[WARN] k3s image import failed: {e}")
self.prole_cfg_data['Docker Build']['STATUS'] = 'Built' if self._db_built_success else 'Failed'
def _step_cluster(self) -> None:
self.log("==> Cluster setup")
cluster_env = self._get_input('init_cluster.cluster_env', 'dev')
env_key = _normalize_cluster_env(cluster_env)
self.prole_cfg_data['Initialize Cluster']['ENVIRONMENT'] = cluster_env
self.prole_cfg_data['Initialize Cluster']['K3S_SERVER_URL'] = self._get_input('init_cluster.k3s_server_url', '')
self.prole_cfg_data['Initialize Cluster']['K3S_TOKEN'] = _encrypt_cfg_secret(self._get_input('init_cluster.k3s_token', '') or '')
self.prole_cfg_data['Optional Features']['SUPABASE_ENABLED'] = _bool_str(self._get_input_bool('init_cluster.supabase_enabled', False))
self.prole_cfg_data['Optional Features']['KERBEROS_ENABLED'] = _bool_str(self._get_input_bool('init_cluster.kerberos_enabled', False))
self.prole_cfg_data['Optional Features']['AT_REST_ENCRYPTION_ENABLED'] = _bool_str(self._get_input_bool('init_cluster.at_rest_encryption_enabled', False))
mode = _deployment_mode_from_env(cluster_env)
if mode:
self._ensure_registry_defaults(mode)
if not self._get_input_bool('init_cluster.start_cluster', DEFAULT_ACTION_FLAGS.get('init_cluster.start_cluster', True)):
self.log("[SKIP] Cluster start disabled.")
return
if env_key == 'dev':
if not self.controller.check_docker_running():
if platform.system() == 'Darwin':
self.log("Starting Docker...")
self._run_cmd(['open', '-a', 'Docker'])
for _ in range(30):
time.sleep(2)
if self.controller.check_docker_running():
break
if not self.controller.check_docker_running():
raise Exception("Docker is not running.")
cluster_name = "prole-dev-cluster"
res = subprocess.run(['k3d', 'cluster', 'list', '--no-headers'], capture_output=True, text=True)
res_stdout = res.stdout or ''
if self.reset_cluster:
self.log(f"[RESET] Deleting k3d cluster {cluster_name}...")
self._run_cmd(['k3d', 'cluster', 'delete', cluster_name])
res_stdout = ''
self.reset_cluster = False
if cluster_name not in res_stdout:
prole_data = self._resolve_env_dir('PROLE_DATA', 'data')
self._create_k3d_cluster(cluster_name, prole_data)
else:
# Check if it's already running
if 'running' not in res_stdout.lower():
self._run_cmd(['k3d', 'cluster', 'start', cluster_name])
else:
self.log(f"[OK] Cluster {cluster_name} is already running.")
self.log(f"[OK] Cluster ready: {cluster_name}")
try:
os.environ.pop("KUBECONFIG", None)
subprocess.run(['kubectl', 'config', 'use-context', f'k3d-{cluster_name}'], capture_output=True)
except Exception:
pass
else:
# Set KUBECONFIG to local prole secrets if it exists
service_dir = self._get_input('env_setup.PROLE_SERVICE', '').strip()
candidates = []
if service_dir:
candidates.append(Path(service_dir) / 'secrets' / 'k3s.kubeconfig')
candidates.append(Path.home() / '.prole' / 'secrets' / 'k3s.kubeconfig')
candidates.append(self.project_root / 'etc' / 'secrets' / 'k3s.kubeconfig')
for kc in candidates:
if kc.exists():
os.environ['KUBECONFIG'] = str(kc)
self.log(f"Using local kubeconfig: {kc}")
break
kubectl = subprocess.run(['which', 'kubectl'], capture_output=True)
if kubectl.returncode != 0:
raise Exception("kubectl not found. Please install kubectl and configure access to the target cluster.")
if env_key == 'service':
if self.reset_cluster:
self.log("[RESET] Resetting k3s cluster (mode k3s)...")
self._cleanup_local_k3s_artifacts()
vault_args = []
v_file = os.environ.get('ANSIBLE_VAULT_PASSWORD_FILE')
if v_file:
vault_args = ["-v", v_file]
delete_cmd = ["./ansible.sh", "-p", "infrastructure/playbooks/k3s_delete.yml"] + vault_args
self.log(f"Running: {' '.join(delete_cmd)}")
rc = self._run_cmd(delete_cmd)
if rc != 0:
raise Exception(f"k3s delete failed with code {rc}")
# 1) --reset should run the ansible k3s/tasks/cleanup on myrddin and then reinstall.
# This is handled by k3s_reset.yml
reset_cmd = ["./ansible.sh", "-p", "infrastructure/playbooks/k3s_reset.yml", "-l", "myrddin.prole.org"] + vault_args
self.log(f"Running: {' '.join(reset_cmd)}")
rc = self._run_cmd(reset_cmd)
if rc != 0:
raise Exception(f"k3s reset failed with code {rc}")
# 2) Capture new k3s token + TLS bundle from init server and update vault
sync_token_cmd = ["./ansible.sh", "-p", "infrastructure/playbooks/k3s_sync.yml", "-l", "myrddin.prole.org"] + vault_args
self.log(f"Running: {' '.join(sync_token_cmd)}")
rc = self._run_cmd(sync_token_cmd)
if rc != 0:
raise Exception(f"k3s sync (token) failed with code {rc}")
# 3) Reinstall agent nodes with the refreshed token
reset_agents_cmd = ["./ansible.sh", "-p", "infrastructure/playbooks/k3s_reset.yml", "-l", "pi.prole.org,retropie.prole.org"] + vault_args
self.log(f"Running: {' '.join(reset_agents_cmd)}")
rc = self._run_cmd(reset_agents_cmd)
if rc != 0:
raise Exception(f"k3s reset (agents) failed with code {rc}")
self.reset_cluster = False
if not (os.environ.get("KUBECONFIG") or "").strip():
kubeconfig_path = _find_kubeconfig_file()
if kubeconfig_path:
os.environ["KUBECONFIG"] = kubeconfig_path
self.log(f"Using local kubeconfig: {kubeconfig_path}")
server = (self._get_input('init_cluster.k3s_server_url', '') or '').strip()
token = (self._get_input('init_cluster.k3s_token', '') or '').strip()
if token:
token = self._resolve_secret_value(token)
if server and token:
if not server.startswith('http'):
server = f"https://{server}"
cmd = [
'kubectl',
'--server=' + server,
'--token=' + token,
'--insecure-skip-tls-verify=true',
'cluster-info'
]
else:
kubeconfig_env = (os.environ.get('KUBECONFIG') or '').strip()
if kubeconfig_env:
cmd = ['kubectl', '--kubeconfig', kubeconfig_env, 'cluster-info']
else:
raise Exception("K3s server URL/token missing for service cluster.")
else:
cmd = ['kubectl', 'cluster-info']
res = subprocess.run(cmd, capture_output=True, text=True)
if res.returncode != 0:
raise Exception("Cluster is not reachable.")
self.log(f"[OK] Cluster ready: {env_key}")
if env_key == 'service':
try:
argocd_ns = self._argocd_namespace()
registry_ns = self._registry_namespace()
env = self._script_env_for_namespace(self._service_namespace())
env["ARGOCD_NAMESPACE"] = argocd_ns
env["REGISTRY_NAMESPACE"] = registry_ns
self.log("==> Registry/ArgoCD preflight")
self._run_script(
"init_registry.sh",
args=["-n", argocd_ns, "--registry-namespace", registry_ns, "update"],
env=env,
)
except Exception as e:
self.err(f"[WARN] Registry/ArgoCD preflight failed: {e}")
if self._openbao_init_env:
env = dict(self._openbao_init_env)
mode = _deployment_mode_from_env(self._get_input('init_cluster.cluster_env', ''))
if not mode:
if env_key == 'dev':
mode = 'k3d'
elif env_key == 'service':
mode = 'k3s'
elif env_key == 'prod':
mode = 'k8s'
env.setdefault("PROLE_MODE", mode or env_key)
self.log("==> OpenBao initialize (cluster)")
rc_bao = self._run_script("init_openbao.sh", args=["initialize"], env=env, stdin_text=f"{self._openbao_init_password}\n")
if rc_bao != 0:
raise Exception(f"OpenBao initialization failed (code {rc_bao})")
self._openbao_init_env = None
self._openbao_init_password = None
# Run repair pipeline if anomalies are detected on a ready cluster
try:
service_ns = self._service_namespace()
env = self._script_env_for_namespace(service_ns)
self._maybe_run_repair_pipeline(env)
except Exception:
pass
if env_key == 'service':
self._deploy_common_services(env_key)
else:
self._deploy_opentofu(env_key)
def _deploy_common_services(self, env_key: str) -> None:
self.log("==> Common services deploy")
service_ns = self._service_namespace()
env = self._script_env_for_namespace(service_ns)
kerberos_enabled = self._get_input_bool('kerberos_config.enabled', False) or \
self._get_input_bool('init_cluster.kerberos_enabled', False)
try:
# Check status before update
status_args = ["-n", service_ns]
if kerberos_enabled:
status_args.append("-k")
rc_status = self._run_script("status_common_services.sh", args=status_args, env=env)
if rc_status == 0:
self.log(f"[OK] Common services in {service_ns} are healthy. Skipping update.")
return
self.log(f"[INFO] Common services unhealthy or missing in {service_ns}. Attempting repair/deploy...")
svc_args = ["-n", service_ns]
if kerberos_enabled:
svc_args.append("-k")
svc_args.append("update")
rc = self._run_script("init_common_services.sh", args=svc_args, env=env)
if rc != 0:
self.err(f"[ERROR] Common services deploy failed (code {rc})")
else:
# Re-validate after fix
rc_status = self._run_script("status_common_services.sh", args=status_args, env=env)
if rc_status != 0:
self.err(f"[ERROR] Common services still unhealthy after repair attempt.")
else:
self.log(f"[OK] Common services in {service_ns} are now healthy.")
finally:
pass
def _deploy_opentofu(self, env_key: str) -> None:
self.log("==> OpenTofu deploy")
ns = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
env = self._script_env_for_namespace(ns)
try:
rc = self._run_script("init_opentofu.sh", args=["start"], env=env)
if rc != 0:
self.err(f"[ERROR] OpenTofu deploy failed (code {rc})")
finally:
pass
def _step_kerberos(self) -> None:
enabled = self._get_input_bool('kerberos_config.enabled', False)
if not enabled:
self.log("[SKIP] Kerberos disabled.")
return
realm = self._get_input('kerberos_config.realm', '')
kdc = self._get_input('kerberos_config.kdc', '')
user = self._get_input('kerberos_config.user', '')
password = self._get_input('kerberos_config.password', '')
self.prole_cfg_data['Kerberos Authentication']['ENABLED'] = _bool_str(enabled)
self.prole_cfg_data['Kerberos Authentication']['REALM'] = realm
self.prole_cfg_data['Kerberos Authentication']['KDC'] = kdc
self.prole_cfg_data['Kerberos Authentication']['SERVER'] = kdc
self.prole_cfg_data['Kerberos Authentication']['USER'] = user
self.prole_cfg_data['Kerberos Authentication']['PASSWORD'] = password
self.prole_cfg_data['Kerberos Authentication']['AD_PORT_FORWARD'] = os.environ.get('KRB5_AD_PORT_FORWARD', '1')
self.prole_cfg_data['Kerberos Authentication']['AD_TCP_PORTS'] = os.environ.get('KRB5_AD_TCP_PORTS', '88 389 445 464 636')
self.prole_cfg_data['Kerberos Authentication']['AD_UDP_PORTS'] = os.environ.get('KRB5_AD_UDP_PORTS', '88 464')
self.prole_cfg_data['Kerberos Authentication']['AD_PROXY_HOST_NETWORK'] = os.environ.get('KRB5_AD_PROXY_HOST_NETWORK', '1')
self.prole_cfg_data['Kerberos Authentication']['AD_PROXY_IMAGE'] = os.environ.get('KRB5_AD_PROXY_IMAGE', 'alpine/socat')
self.prole_cfg_data['Kerberos Authentication']['AD_PROXY_SERVICE'] = os.environ.get('KRB5_AD_SERVICE_NAME', 'prole-kerberos-ad-dc')
ns = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
env = self._script_env_for_namespace(ns)
env["KERBEROS_ENABLED"] = _bool_str(enabled)
env["ENABLED"] = env["KERBEROS_ENABLED"]
env["KRB5_REALM"] = realm
if env.get("PROLE_MODE") == "k3s":
env.setdefault("KRB5_AD_PORT_FORWARD", "0")
if not self._get_input_bool('kerberos_config.test_connection', DEFAULT_ACTION_FLAGS.get('kerberos_config.test_connection', False)):
self.log("[SKIP] Kerberos test disabled.")
return
# Check if already tested successfully
if self.prole_cfg_data.get('Kerberos Authentication', {}).get('TEST_STATUS') == 'Success':
self.log("[OK] Kerberos already tested successfully. Skipping.")
return
if not (realm and user and password and kdc):
self.err("[WARN] Kerberos test skipped: missing realm/user/password/kdc.")
return
self.log("==> Kerberos test: init_kerberos.sh test")
rc2 = self._run_script("init_kerberos.sh", args=["test"], env=env)
if rc2 != 0:
self.err(f"[ERROR] Kerberos test failed (code {rc2})")
self.prole_cfg_data.setdefault('Kerberos Authentication', {})['TEST_STATUS'] = 'Failed'
else:
self.log("[OK] Kerberos test completed successfully.")
self.prole_cfg_data.setdefault('Kerberos Authentication', {})['TEST_STATUS'] = 'Success'
def _step_init_scripts(self) -> None:
if not self._get_input_bool('init_scripts.run_scripts', DEFAULT_ACTION_FLAGS.get('init_scripts.run_scripts', True)):
self.log("[SKIP] Init scripts disabled.")
self.prole_cfg_data['Initialization Scripts']['STATUS'] = 'Skipped'
return
if self.prole_cfg_data.get('Initialization Scripts', {}).get('STATUS') == 'Finished':
self.log("[OK] Initialization scripts already finished.")
return
self.log("==> Initialization scripts")
ns = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
env = self._script_env_for_namespace(ns)
password = self._get_input('init_password.db_password', '')
kerberos_enabled = self._get_input_bool('kerberos_config.enabled', False)
svc_args = ["update"]
if kerberos_enabled:
svc_args = ["-k", "update"]
steps = [
("init_common_services.sh", svc_args, False),
("init_cloudnative_pg.sh", ["initialize"], False),
]
if kerberos_enabled:
steps.append(("init_kerberos.sh", ["initialize"], False))
steps.extend([
("init_prole-db-backup.sh", ["start"], False),
("init_monitoring.sh", ["initialize"], False),
("init_nginx_ingress.sh", ["initialize"], False),
])
overall_success = True
for script, args, needs_password in steps:
self.log(f"--> {script} {' '.join(args)}")
stdin_text = f"{password}\n" if needs_password else None
# Use custom line handler for monitoring to capture Grafana password and port mappings
rc = self._run_script(script, args=args, env=env, stdin_text=stdin_text, on_line=self._process_script_output_line)
if rc != 0:
self.err(f"[ERROR] {script} failed (code {rc})")
overall_success = False
self._scripts_success = overall_success
# Verify critical secrets
self.log(f"==> Verifying critical secrets in namespace {ns}")
critical_secrets = ["prole-db-user", "prole-db-superuser", "cnpg-admin-key"]
missing_secrets = []
for secret in critical_secrets:
rc_s = self._run_cmd(["kubectl", "get", "secret", secret, "-n", ns])
if rc_s != 0:
missing_secrets.append(secret)
if missing_secrets:
self.err(f"[CRITICAL] Missing secrets in namespace '{ns}': {', '.join(missing_secrets)}")
self.err("Database initialization will fail without these secrets.")
self._scripts_success = False
self.prole_cfg_data['Initialization Scripts']['STATUS'] = 'Completed' if self._scripts_success else 'Attempted'
def _step_cnpg_deploy(self) -> None:
if not self._get_input_bool('init_cnpg_deploy.run_deploy', DEFAULT_ACTION_FLAGS.get('init_cnpg_deploy.run_deploy', True)):
self.log("[SKIP] CnPG deploy disabled.")
self.prole_cfg_data['Deployment']['STATUS'] = 'Skipped'
return
self.log("==> Deploy CloudNative-PG")
etc_dir = self.project_root / "etc"
ns = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
env = self._script_env_for_namespace(ns)
rc = self._run_cmd(['bash', str(etc_dir / 'init_cloudnative_pg.sh'), 'deploy', 'latest'], env=env)
if rc == 0:
self._cnpg_success = True
self.prole_cfg_data['Deployment']['STATUS'] = 'Deployed'
else:
self._cnpg_success = False
self.prole_cfg_data['Deployment']['STATUS'] = 'Attempted'
self.err(f"[ERROR] CnPG deploy failed (code {rc})")
if self._get_input_bool('init_cnpg_deploy.force_rollout', DEFAULT_ACTION_FLAGS.get('init_cnpg_deploy.force_rollout', False)):
self.log("==> Force rollout")
rc2 = self._run_cmd(['bash', str(etc_dir / 'init_cloudnative_pg.sh'), 'rollout'], env=env)
if rc2 != 0:
self.err(f"[ERROR] Rollout failed (code {rc2})")
def _step_supabase(self) -> None:
if not self._get_input_bool('supabase_config.run_deploy', False):
self.log("[SKIP] Supabase deploy disabled.")
self.prole_cfg_data['Supabase'] = {'STATUS': 'Skipped'}
return
self.log("==> Deploy Supabase")
ns = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
env = self._script_env_for_namespace(ns)
script_path = self.project_root / "supabase" / "deploy.sh"
if not script_path.exists():
self.err(f"[ERROR] Supabase deploy script not found: {script_path}")
self._supabase_success = False
self.prole_cfg_data['Supabase'] = {'STATUS': 'Attempted'}
return
mode = (os.environ.get("SUPABASE_DEPLOY_MODE") or os.environ.get("SUPABASE_MODE") or "").strip()
if not mode:
# Match cluster environment to supabase deploy mode
cluster_env = _normalize_cluster_env(self._get_input('init_cluster.cluster_env', 'dev'))
if cluster_env == 'dev':
mode = "k3d"
elif cluster_env in ('service', 'prod'):
mode = "k8s"
else:
mode = "k3d"
args = ["--mode", mode]
cfg_path = None
if self.cfg_path and self.cfg_path.exists():
cfg_path = self.cfg_path
else:
candidate = self.project_root / "conf" / "prole.cfg"
if candidate.exists():
cfg_path = candidate
if cfg_path:
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")
self.log(f"--> supabase/deploy.sh {' '.join(args)}")
rc = self._run_cmd(['bash', str(script_path)] + args, env=env)
if rc == 0:
self._supabase_success = True
self.prole_cfg_data['Supabase'] = {'STATUS': 'Deployed'}
else:
self._supabase_success = False
self.prole_cfg_data['Supabase'] = {'STATUS': 'Attempted'}
self.err(f"[ERROR] Supabase deploy failed (code {rc})")
def _finalize_secrets(self) -> None:
ns = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
if not ns or ns == 'default':
ns = self._get_input('env_setup.NAMESPACE', '') or 'default'
p1 = self._get_input('init_password.db_password', '')
if p1.startswith('${'):
self.log(f"[OK] Secrets already finalized for namespace {ns}. Skipping build-a-bao.")
return
self.log("==> OpenBao: finalize secrets")
env = self._script_env_for_namespace(ns)
rc = self._run_script("build-a-bao.sh", env=env)
if rc == 0:
self._secrets_finalized = True
self.log("[OK] Secrets saved to OpenBao and prole.cfg updated.")
else:
self.err(f"[ERROR] build-a-bao.sh failed (code {rc})")
def _prepare_opentofu_pipeline(self) -> None:
env_key = _normalize_cluster_env(self._get_input('init_cluster.cluster_env', ''))
if env_key != 'service':
return
namespace = (self._get_input('init_password.db_namespace', '') or '').strip() or 'default'
k3s_server = (self._get_input('init_cluster.k3s_server_url', '') or '').strip()
k3s_token = (self._get_input('init_cluster.k3s_token', '') or '').strip()
if k3s_token:
k3s_token = self._resolve_secret_value(k3s_token)
if k3s_server and not k3s_server.startswith('http'):
k3s_server = f"https://{k3s_server}"
if not k3s_server:
self.err("[WARN] OpenTofu pipeline skipped: missing k3s server URL.")
return
if not k3s_token:
self.err("[WARN] OpenTofu pipeline missing k3s token; writing empty token.")
pipeline_dir = _sync_opentofu_pipeline(self.project_root, namespace, k3s_server, k3s_token)
self.prole_cfg_data['Deployment']['OPENTOFU_PIPELINE_DIR'] = str(pipeline_dir)
self.log(f"[OK] OpenTofu pipeline prepared at {pipeline_dir}")
def _create_k3d_cluster(self, cluster_name: str, prole_data: Path) -> None:
volume_args = _k3d_prole_data_volume_args(str(prole_data))
reg_args = []
if _local_registry_enabled('k3d'):
info = self._ensure_local_registry_available()
if info:
_host_registry, cluster_registry = info
reg_use = cluster_registry or ''
if reg_use.endswith(".localhost:5000"):
reg_use = reg_use.replace(".localhost:5000", ":5000")
elif reg_use.endswith(".localhost"):
reg_use = reg_use[:-len(".localhost")]
if reg_use:
reg_args = ['--registry-use', reg_use]
cmd = ['k3d', 'cluster', 'create', cluster_name, '-a', '2'] + volume_args + reg_args + ['--api-port', '0.0.0.0:6443']
rc = self._run_cmd(cmd)
if rc == 0:
return
if platform.system() == 'Darwin':
fallback = Path.home() / '.prole' / 'data'
try:
if prole_data.resolve() == fallback.resolve():
raise Exception("k3d cluster create failed.")
except Exception:
# If path resolution fails, still try the fallback
pass
self.err(f"[WARN] k3d create failed with {prole_data}; retrying with {fallback}...")
rc = self._run_cmd(['k3d', 'cluster', 'create', cluster_name, '-a', '2'] + _k3d_prole_data_volume_args(str(fallback)) + reg_args + ['--api-port', '0.0.0.0:6443'])
if rc == 0:
# Persist fallback for the rest of the run
self.inputs['env_setup.PROLE_DATA'] = str(fallback)
return
raise Exception(f"Failed to create k3d cluster {cluster_name} (code {rc})")
def _perform_cluster_reset(self) -> None:
if not self.reset_requested:
return
cluster_env = self._get_input('init_cluster.cluster_env', 'dev')
env_key = _normalize_cluster_env(cluster_env)
if env_key == 'dev':
self.log("==> Cluster reset (k3d)")
if not self.controller.check_docker_running():
if platform.system() == 'Darwin':
self.log("Starting Docker...")
self._run_cmd(['open', '-a', 'Docker'])
for _ in range(30):
time.sleep(2)
if self.controller.check_docker_running():
break
if not self.controller.check_docker_running():
raise Exception("Docker is not running.")
cluster_name = "prole-dev-cluster"
rc = self._run_cmd(['k3d', 'cluster', 'delete', cluster_name])
if rc != 0:
self.err(f"[WARN] Failed to delete k3d cluster {cluster_name} (code {rc}); continuing.")
prole_data = self._resolve_env_dir('PROLE_DATA', 'data')
self._create_k3d_cluster(cluster_name, prole_data)
elif env_key in ('service', 'prole-service-cluster', 'k3s'):
self.log("==> Cluster reset (k3s namespace cleanup)")
ns = self._get_input('env_setup.NAMESPACE', 'default')
server = self._get_input('init_cluster.k3s_server_url', '')
token = self._get_input('init_cluster.k3s_token', '')
_reset_k3s_namespace(self.project_root, ns, server, token)
self._cleanup_local_k3s_artifacts()
self.reset_requested = False
def _cleanup_local_k3s_artifacts(self) -> None:
paths = [
self.project_root / "etc" / "secrets" / "k3s.kubeconfig",
Path.home() / ".prole" / "secrets" / "k3s.kubeconfig",
self.project_root / "prole-k3s.kubeconfig",
]
for path in paths:
try:
if path.exists():
path.unlink()
self.log(f"[RESET] Removed local kubeconfig {path}")
except Exception as exc:
self.err(f"[WARN] Failed to remove {path}: {exc}")
def run(self) -> int:
_configure_unbuffered_io()
self.log(f"[CONFIG] Using {self.cfg_path}")
try:
self.inputs = {**self._default_inputs(), **self._load_inputs_from_cfg()}
# Synchronize controller state
self.controller.state.inputs = self.inputs
self.controller.state.config_data = self.prole_cfg_data
except Exception as e:
self.err(f"[FATAL] {e}")
return 2
self._write_cfg()
self._perform_cluster_reset()
milestones = [
DependenciesMilestone(),
NetworkScanMilestone(),
EnvSetupMilestone(),
SecretManagementMilestone(),
DatabaseCreationMilestone(),
DockerBuildMilestone(),
ClusterLifecycleMilestone(),
InitializationScriptsMilestone(),
DeploymentMilestone()
]
try:
self.controller.run_milestones(milestones, progress_callback=lambda msg, p: self.log(f"[{p*100:0.0f}%] {msg}"))
# Sync back results
self.prole_cfg_data = self.controller.state.config_data
self.prole_cfg_data['Install']['STATUS'] = 'Finished'
self._write_cfg()
self.log("[DONE] Silent install completed.")
return 0
except Exception as e:
self.err(f"[FATAL] {e}")
try:
self.prole_cfg_data['Install']['STATUS'] = 'Failed'
self._write_cfg()
except Exception:
pass
return 2
def _run_silent_install_test(project_root: Path, cfg_path: Path) -> int:
script = project_root / "tests" / "silent_install_test.sh"
if not script.exists():
print(f"[ERROR] Silent install test script not found: {script}", file=sys.stderr)
return 1
env = os.environ.copy()
env.setdefault("SILENT_INSTALL_LOG", "true")
res = subprocess.run(["bash", str(script), str(cfg_path)], env=env)
return res.returncode
def _attempt_k3s_repair(controller: ProleController, namespace: str, server: str, token: str, db_password: str) -> None:
env = os.environ.copy()
env["PROLE_HOME"] = str(controller.project_root)
env["PROLE_SERVICE"] = str(controller.project_root)
env["NAMESPACE"] = namespace
env["PROLE_MODE"] = "k3s"
if db_password:
env["DB_PASSWORD"] = db_password
env["OPENTOFU_ADMIN_PASSWORD"] = db_password
if server and token and not env.get("KUBECONFIG"):
kubeconfig_path = _write_k3s_kubeconfig(server, token)
env["KUBECONFIG"] = str(kubeconfig_path)
controller.run_script("init_openbao.sh", args=["-n", namespace, "update"], env=env)
controller.run_script("init_opentofu.sh", args=["-n", namespace, "update"], env=env)
def _reset_k3s_namespace(project_root: Path, namespace: str, server: str, token: str) -> None:
script = project_root / "scripts" / "reset-ns.sh"
if not script.exists():
print(f"[WARN] Namespace reset script not found: {script}", file=sys.stderr)
return
env = os.environ.copy()
if server and token and not env.get("KUBECONFIG"):
kubeconfig_path = _write_k3s_kubeconfig(server, token)
env["KUBECONFIG"] = str(kubeconfig_path)
subprocess.run(["bash", str(script), "-n", namespace], env=env)
def _prepare_k3s_pipeline(controller: ProleController) -> int:
project_root = controller.project_root
cfg_path = project_root / "conf" / "prole.cfg"
info = _detect_ansible_topology(project_root)
k3s_server = (info.get('k3s_server_url') or '').strip() if info else ''
k3s_token = (info.get('k3s_token') or '').strip() if info else ''
if k3s_server and not k3s_server.startswith('http'):
k3s_server = f"https://{k3s_server}"
installer = ProleSilentInstaller(controller, str(cfg_path))
try:
existing_inputs = installer._load_inputs_from_cfg()
except Exception:
existing_inputs = {}
installer.inputs = {**installer._default_inputs(), **existing_inputs}
# Override for k3s pipeline defaults
installer.inputs['init_cluster.cluster_env'] = 'prole-service-cluster'
installer.inputs['init_cluster.supabase_enabled'] = _bool_str(False)
installer.inputs['init_cluster.kerberos_enabled'] = _bool_str(False)
installer.inputs['init_cluster.at_rest_encryption_enabled'] = _bool_str(True)
installer.inputs['kerberos_config.enabled'] = _bool_str(False)
installer.inputs['kerberos_config.test_connection'] = _bool_str(False)
if k3s_server and not installer.inputs.get('init_cluster.k3s_server_url'):
installer.inputs['init_cluster.k3s_server_url'] = k3s_server
if k3s_token and not installer.inputs.get('init_cluster.k3s_token'):
installer.inputs['init_cluster.k3s_token'] = k3s_token
installer.inputs['env_setup.PROLE_HOME'] = str(project_root)
installer.inputs['env_setup.PROLE_CONF'] = str(project_root / "conf")
installer.inputs['env_setup.PROLE_DATA'] = (
installer._resolve_env_value('PROLE_DATA', str(Path.home() / '.prole' / 'data'))
or str(Path.home() / '.prole' / 'data')
)
installer.inputs['env_setup.PROLE_LOGS'] = (
installer._resolve_env_value('PROLE_LOGS', str(Path.home() / '.prole' / 'logs'))
or str(Path.home() / '.prole' / 'logs')
)
installer.inputs['env_setup.PROLE_SERVICE'] = str(project_root / "etc")
installer._write_cfg()
namespace = (installer._get_input('init_password.db_namespace', '') or '').strip()
if not namespace:
namespace = (installer._get_input('env_setup.NAMESPACE', '') or '').strip() or 'default'
db_password = installer._get_input('init_password.db_password', '').strip()
attempt = 1
total_attempts = 0
max_attempts_env = os.environ.get("PROLE_SILENT_TEST_MAX_ATTEMPTS", "").strip()
max_attempts = int(max_attempts_env) if max_attempts_env.isdigit() else 0
while True:
total_attempts += 1
rc = _run_silent_install_test(project_root, cfg_path)
if rc == 0:
break
cur_server = installer.inputs.get('init_cluster.k3s_server_url', k3s_server)
cur_token = installer.inputs.get('init_cluster.k3s_token', k3s_token)
if attempt == 1:
_attempt_k3s_repair(controller, namespace, cur_server, cur_token, db_password)
elif attempt == 2:
_reset_k3s_namespace(project_root, namespace, cur_server, cur_token)
else:
attempt = 0
attempt += 1
if max_attempts and total_attempts >= max_attempts:
return rc
cur_server = installer.inputs.get('init_cluster.k3s_server_url', k3s_server)
cur_token = installer.inputs.get('init_cluster.k3s_token', k3s_token)
_sync_opentofu_pipeline(project_root, namespace, cur_server, cur_token)
return 0