prole/knoe/core/milestones.py
chrisfu 5715add366 feat: GitLab deployment pipeline — operator fix, CI config, and namespace isolation
- etc/init_gitlab.sh: Add global.redis block (host/port/auth) to the GitLab CR so
  the chart does not fail NOTES.txt validation when redis.install: false.
  Drop the GITOPS_NAMESPACE config fallback in namespace resolution to prevent the
  Gitea namespace from bleeding into GitLab deployments; GITLAB_NAMESPACE is now the
  sole source of truth with a hard default of "gitlab".
- knoe/core/milestones.py: Fix GitOpsMilestone to route to init_gitlab.sh when
  gitops.git_provider = GitLab (was hardcoded to init_gitea.sh). Namespace resolution
  now prefers gitops.gitlab_namespace input key, then gitops.namespace, then "gitlab" —
  never picks up a stale GITLAB_NAMESPACE from the OS environment.
- conf/service/prole.cfg: Switch gitops.git_provider / GITOPS_PROVIDER to GitLab.
  Update accumulated runtime state from install runs.
- install.sh: Prefer the repo-local venv Python (PROLE_HOME/bin/python3) so that
  PyYAML and other prole_requirements.txt deps are always available.
- .gitlab-ci.yml: New CI pipeline — on every push to main, run the silent install
  (./install.sh -S -c conf/service/prole.cfg) to deploy a fresh CNPG ecosystem.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-04-01 15:26:39 -07:00

1233 lines
47 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

from __future__ import annotations
import logging
import os
import subprocess
import sys
import time
from pathlib import Path
from typing import TYPE_CHECKING
from knoe.milestone import Milestone
from knoe import config as inst_config
from knoe import prole_conf
from knoe.core.stream_exec import run_streaming_cmd
from knoe.core.env import resolve_prole_home
from knoe.core.policy import (
POLICY_CFG_KEY,
OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES,
evaluate_optional_workloads_allowed,
)
from knoe.core.ops.cloudnative_pg import (
initialize as cnpg_initialize,
deploy as cnpg_deploy,
)
from knoe.core.ops import openbao as openbao_ops
from knoe.core.ops import registry as registry_ops
from knoe.core.ops import garage_store as garage_store_ops
from knoe.core.ops import opentofu as opentofu_ops
from knoe.core.ops import monitoring as monitoring_ops
if TYPE_CHECKING:
from knoe.state import InstallerState
from knoe.milestone import ProgressCallback
def _stream_line(line: str) -> None:
try:
sys.stdout.write(line)
sys.stdout.flush()
except Exception:
pass
class DependenciesMilestone(Milestone):
def __init__(self):
super().__init__("dependencies", "Dependency Verification")
self.logger = logging.getLogger("DependenciesMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
if progress:
progress("Checking dependencies...", 0.1)
dependencies = inst_config.DEPENDENCIES
missing = []
for i, dep in enumerate(dependencies):
ok, location, version = inst_config.get_dep_info(dep)
if ok:
self.logger.info(f"[OK] {dep['name']} {version or ''}".strip())
else:
missing.append(dep)
self.logger.info(f"[MISSING] {dep['name']}")
if progress:
progress(
f"Checking {dep['name']}...", 0.1 + (i / len(dependencies)) * 0.4
)
if not missing:
self._set_status(state, "All installed")
if progress:
progress("All dependencies installed", 1.0)
return
# Check if auto-install is enabled
auto_install = self._parse_bool(
state.inputs.get("dependencies.auto_install_missing", "True")
)
if not auto_install:
self.logger.error("Dependencies missing and auto-install disabled.")
self._set_status(state, "Missing")
return
for i, dep in enumerate(missing):
install_cmd = dep.get("install_cmd")
if not install_cmd:
self.logger.error(f"No install command for {dep['name']}.")
continue
should_install = self._parse_bool(
state.inputs.get(f"dependencies.{dep['id']}.install", "True")
)
if not should_install:
self.logger.info(f"[SKIP] {dep['name']} install disabled by config.")
continue
self.logger.info(f"[INSTALL] {dep['name']} -> {install_cmd}")
if progress:
progress(f"Installing {dep['name']}...", 0.5 + (i / len(missing)) * 0.4)
rc = self._run_cmd(install_cmd)
if rc != 0:
self.logger.error(f"Install failed for {dep['name']} (code {rc})")
# Re-verify
still_missing = []
for dep in missing:
ok, _, _ = inst_config.get_dep_info(dep)
if not ok:
still_missing.append(dep["name"])
if still_missing:
self.logger.error(f"Still missing: {', '.join(still_missing)}")
self._set_status(state, "Missing")
else:
self._set_status(state, "All installed")
if progress:
progress("All dependencies installed", 1.0)
def _set_status(self, state: InstallerState, status: str):
if "Dependencies" not in state.config_data:
state.config_data["Dependencies"] = {}
state.config_data["Dependencies"]["STATUS"] = status
class NetworkScanMilestone(Milestone):
def __init__(self):
super().__init__("network_scan", "Network Configuration")
self.logger = logging.getLogger("NetworkScanMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
if not self._parse_bool(state.inputs.get("network_scan.run", "True")):
self.logger.info("Network scan disabled.")
return
existing_kdc = state.config_data.get("Network", {}).get("KDC_AUTO_DETECTED")
if existing_kdc:
self.logger.info(f"Network scan already performed. KDC: {existing_kdc}")
return
scan_binary = inst_config.get_resource_path("prole-net/prole-agent")
if not scan_binary.exists():
self.logger.error(f"Scan binary not found at {scan_binary}")
return
prole_home = resolve_prole_home(env={})
scan_dir = prole_home / "scan"
scan_dir.mkdir(parents=True, exist_ok=True)
kdc_found = None
pending_line = ""
def _handle_line(line: str):
nonlocal kdc_found
if "KDC is:" in line:
try:
ip_part = line.split("KDC is:")[1].strip()
ip = ip_part.split()[0].strip("[]():,")
if ip:
kdc_found = ip
except Exception:
pass
elif "Active Directory" in line or "88" in line:
import socket
for part in line.split():
try:
socket.inet_aton(part.strip("[]():,"))
kdc_found = part.strip("[]():,")
break
except Exception:
continue
def _handle_stdout(text: str):
nonlocal pending_line
pending_line += text
while "\n" in pending_line:
line, pending_line = pending_line.split("\n", 1)
_handle_line(line + "\n")
if progress:
progress("Scanning network...", 0.5)
rc = run_streaming_cmd(
[str(scan_binary)], cwd=str(scan_dir), on_stdout=_handle_stdout
)
if pending_line:
_handle_line(pending_line)
if kdc_found:
state.inputs["kerberos_config.kdc"] = kdc_found
state.inputs["kerberos_config.enabled"] = "True"
if "Network" not in state.config_data:
state.config_data["Network"] = {}
state.config_data["Network"]["KDC_AUTO_DETECTED"] = kdc_found
state.config_data["Network"]["KERBEROS_AUTO_ENABLED"] = "True"
self.logger.info(f"KDC detected: {kdc_found}")
if progress:
progress("Network scan complete", 1.0)
class EnvSetupMilestone(Milestone):
def __init__(self):
super().__init__("env_setup", "System Environment Setup")
self.logger = logging.getLogger("EnvSetupMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
if progress:
progress("Setting up environment...", 0.5)
env_keys = [
"PROLE_HOME",
"PROLE_CONF",
"PROLE_DATA",
"PROLE_LOGS",
"PROLE_SERVICE",
"NAMESPACE",
]
if "System Environment" not in state.config_data:
state.config_data["System Environment"] = {}
for k in env_keys:
val = state.inputs.get(f"env_setup.{k}")
if val:
state.config_data["System Environment"][k] = val
if progress:
progress("Environment setup complete", 1.0)
class SecretManagementMilestone(Milestone):
def __init__(self):
super().__init__("secret_management", "Secret Resolution")
self.logger = logging.getLogger("SecretManagementMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
if progress:
progress("Resolving secrets...", 0.5)
# Keys that might contain secrets
secret_keys = [
"init_password.db_password",
"init_password.db_password_confirm",
"kerberos_config.password",
"init_cluster.k3s_token",
]
for key in secret_keys:
val = state.inputs.get(key)
if val:
resolved = inst_config._resolve_secret_value(val)
if resolved != val:
state.inputs[key] = resolved
self.logger.info(f"Resolved secret for {key}")
# Synchronize confirm password
if state.inputs.get("init_password.db_password") and not state.inputs.get(
"init_password.db_password_confirm"
):
state.inputs["init_password.db_password_confirm"] = state.inputs[
"init_password.db_password"
]
if progress:
progress("Secret resolution complete", 1.0)
class DatabaseCreationMilestone(Milestone):
def __init__(self):
super().__init__("init_password", "Database Creation")
self.logger = logging.getLogger("DatabaseCreationMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
if progress:
progress("Preparing database configuration...", 0.5)
ns = (state.inputs.get("init_password.db_namespace", "") or "").strip()
if not ns:
ns = (
state.inputs.get("env_setup.DATABASE_NAMESPACE", "") or ""
).strip()
if not ns:
ns = (
state.inputs.get("env_setup.NAMESPACE", "") or ""
).strip() or "default"
cluster_name = (state.inputs.get("init_password.cluster_name", "") or "").strip()
if not cluster_name:
cluster_name = (state.inputs.get("env_setup.CLUSTER_NAME", "") or "").strip()
if not cluster_name:
cluster_name = "knoe-db"
user = (state.inputs.get("init_password.db_username", "") or "").strip()
p1 = state.inputs.get("init_password.db_password", "")
# In ncurses/silent mode we might not have all fields yet, but we want to ensure defaults
if not user:
import getpass
user = getpass.getuser()
state.inputs["init_password.db_username"] = user
if not p1:
import secrets
import string
alphabet = string.ascii_letters + string.digits
p1 = "".join(secrets.choice(alphabet) for i in range(24))
state.inputs["init_password.db_password"] = p1
state.inputs["init_password.db_password_confirm"] = p1
self.logger.info("Generated default database password.")
# If password is an unresolvable OpenBao or PROLE_SECRET reference
# (e.g. clean install where OpenBao hasn't been seeded yet),
# generate a fresh password so downstream scripts get a real value.
# OpenBao isn't running yet at this milestone stage, so if the ref
# can't be resolved right now we must provide a concrete password
# that will be seeded into OpenBao later by ClusterLifecycleMilestone.
if p1.startswith("${"):
resolved = inst_config._resolve_secret_value(p1)
if resolved and not resolved.startswith("${"):
p1 = resolved
else:
import secrets as _secrets
import string as _string
alphabet = _string.ascii_letters + _string.digits
p1 = "".join(_secrets.choice(alphabet) for i in range(24))
self.logger.info(
"OpenBao/secret ref unresolvable at this stage; generated new DB password for initialization."
)
state.inputs["init_password.db_password"] = p1
state.inputs["init_password.db_password_confirm"] = p1
if p1 and not p1.startswith("${") and not inst_config._is_prole_secret(p1):
if "Inputs" not in state.config_data:
state.config_data["Inputs"] = {}
# Encrypt and save to config data for persistence in prole.cfg
encrypted_pw = inst_config._encrypt_cfg_secret(p1)
state.config_data["Inputs"]["init_password.db_password"] = encrypted_pw
state.config_data["Inputs"][
"init_password.db_password_confirm"
] = encrypted_pw
state.inputs["init_password.db_namespace"] = ns
state.inputs["env_setup.DATABASE_NAMESPACE"] = ns
state.inputs["init_password.cluster_name"] = cluster_name
state.inputs["env_setup.CLUSTER_NAME"] = cluster_name
if "Database Creation" not in state.config_data:
state.config_data["Database Creation"] = {}
state.config_data["Database Creation"]["DB_USER"] = user
state.config_data["Database Creation"]["DATABASE_NAMESPACE"] = ns
state.config_data["Database Creation"]["CLUSTER_NAME"] = cluster_name
state.config_data["Database Creation"].pop("DB_NAME", None)
state.config_data["Database Creation"].pop("NAMESPACE", None)
# SSH key generation (simplified for now)
if self._parse_bool(state.inputs.get("init_password.generate_ssh_key", "True")):
key_path = Path.home() / ".ssh" / "id_prole_ed25519"
if not key_path.exists():
self._run_cmd(
[
"ssh-keygen",
"-t",
"ed25519",
"-N",
"",
"-f",
str(key_path),
"-C",
user,
]
)
if progress:
progress("Database creation preparation complete", 1.0)
class DockerBuildMilestone(Milestone):
def __init__(self):
super().__init__("init_db_build", "Docker Build")
self.logger = logging.getLogger("DockerBuildMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
if not self._parse_bool(state.inputs.get("init_db_build.run_build", "True")):
self.logger.info("Docker build disabled.")
return
# In k3s/service installs, `init_cloudnative_pg.sh` already has a robust
# pre-flight that ensures the `knoe-db` image is present/pushed (and can
# build it if missing). Building a multi-platform image here is expensive
# and redundant when initialization scripts are enabled.
from knoe.core.env import _normalize_cluster_env
env_key = (
_normalize_cluster_env(state.inputs.get("init_cluster.cluster_env", "dev"))
or "dev"
)
run_init_scripts = self._parse_bool(
state.inputs.get("init_scripts.run_scripts", "True")
)
if env_key != "dev" and run_init_scripts:
self.logger.info(
"Skipping Docker DB image build for service env; init scripts will ensure/push the image."
)
return
if progress:
progress("Building Prole database image...", 0.1)
tag = state.controller.get_knoe_db_version()
args = ["--tag", tag]
if env_key != "dev":
# Add registry prefix for non-dev envs
server = state.inputs.get(
"init_cluster.k3s_server_url", "myrddin.prole.org"
)
if "://" in server:
server = server.split("://")[1]
if ":" in server:
server = server.split(":")[0]
args = ["--tag", f"{server}:5000/knoe-db:{tag}", "--push"]
rc = state.controller.run_script("build_db.sh", args=args, on_line=_stream_line)
if rc != 0:
self.logger.error(f"Docker build failed (code {rc})")
if progress:
progress("Docker build complete", 1.0)
class ClusterLifecycleMilestone(Milestone):
def __init__(self):
super().__init__("cluster_lifecycle", "Cluster Lifecycle Management")
self.logger = logging.getLogger("ClusterLifecycleMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
from knoe.core.env import _normalize_cluster_env
cluster_env = state.inputs.get("init_cluster.cluster_env", "dev")
env_key = _normalize_cluster_env(cluster_env) or cluster_env
start_requested = self._parse_bool(
state.inputs.get("init_cluster.start_cluster", "True")
)
if not start_requested:
self.logger.info("Cluster start not requested.")
return
env = self._get_script_env(state)
if env_key == "dev":
if progress:
progress("Initializing k3d dev cluster...", 0.5)
# Ensure Docker is up (best-effort, macOS auto-start).
if not state.controller.check_docker_running():
if sys.platform == "darwin":
self.logger.info("Starting Docker...")
self._run_cmd(["open", "-a", "Docker"])
for _ in range(30):
time.sleep(2)
if state.controller.check_docker_running():
break
if not state.controller.check_docker_running():
self.logger.error("Docker is not running.")
return
cluster_name = "knoe-dev-cluster"
res = subprocess.run(
["k3d", "cluster", "list", "--no-headers"],
capture_output=True,
text=True,
)
res_stdout = res.stdout or ""
if cluster_name not in res_stdout:
from knoe.core.env import _k3d_prole_data_volume_args
prole_data = state.inputs.get(
"env_setup.PROLE_DATA",
str(inst_config.PROJECT_ROOT / "knoe-db" / "data"),
)
volume_args = _k3d_prole_data_volume_args(prole_data)
cmd = (
["k3d", "cluster", "create", cluster_name, "-a", "2"]
+ volume_args
+ ["--api-port", "0.0.0.0:6443"]
)
rc = self._run_cmd(cmd, on_stdout=_stream_line)
if rc != 0:
self.logger.error(f"k3d cluster create failed (code {rc})")
else:
if "running" not in res_stdout.lower():
self._run_cmd(
["k3d", "cluster", "start", cluster_name],
on_stdout=_stream_line,
)
try:
subprocess.run(
["kubectl", "config", "use-context", f"k3d-{cluster_name}"],
capture_output=True,
)
except Exception:
pass
# Ensure knoe-db image is available in the cluster (avoid registry pull issues).
tag = state.controller.get_knoe_db_version()
local_image = f"knoe-db:{tag}"
if (
subprocess.run(
["docker", "image", "inspect", local_image], capture_output=True
).returncode
== 0
):
db_cfg = state.config_data.get("Docker Build", {}) or {}
registry = (db_cfg.get("LOCAL_REGISTRY_INTERNAL") or "").strip()
if registry.endswith(".localhost:5000"):
registry = registry.replace(".localhost:5000", ":5000")
elif registry.endswith(".localhost"):
registry = registry.replace(".localhost", "")
if not registry:
registry = "k3d-prole-registry:5000"
remote_tag = f"{registry}/knoe-db:{tag}"
subprocess.run(
["docker", "tag", local_image, remote_tag], capture_output=True
)
self._run_cmd(
["k3d", "image", "import", remote_tag, "-c", cluster_name],
on_stdout=_stream_line,
)
elif env_key in ("service", "k3s"):
if progress:
progress("Verifying k3s connection...", 0.5)
# Additional logic for remote k3s services could go here
self.logger.info(
"Service cluster mode: verification and deployment triggered via scripts."
)
# For now, matching the behavior of silent installer's _step_init_cluster
# which mostly sets up env and calls init scripts if needed.
# Initialize OpenBao once cluster is ready. If DB password is anchored,
# start OpenBao anyway so downstream scripts can attempt to read secrets.
db_pw_raw = (state.inputs.get("init_password.db_password", "") or "").strip()
if db_pw_raw:
db_pw = inst_config._resolve_secret_value(db_pw_raw)
if db_pw and not db_pw.startswith("${"):
env_init = dict(env)
env_init["DB_PASSWORD"] = db_pw
env_init.setdefault("GRAFANA_ADMIN_PASSWORD", db_pw)
db_user = (
state.inputs.get("init_password.db_username", "") or ""
).strip()
if db_user:
env_init["KNOE_DB_USER"] = db_user
try:
openbao_ops.initialize(
namespace=env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default")),
env=env_init,
project_root=getattr(state.controller, "project_root", "."),
mode=env.get("PROLE_MODE"),
log=lambda msg: _stream_line(msg if msg.endswith("\n") else msg + "\n"),
)
except Exception as e:
self.logger.error(f"OpenBao initialization failed: {e}")
else:
self.logger.info(
"OpenBao init skipped: DB password unresolved/anchored. Ensuring OpenBao is running..."
)
try:
openbao_ops.start(
namespace=env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default")),
env=env,
project_root=getattr(state.controller, "project_root", "."),
mode=env.get("PROLE_MODE"),
log=lambda msg: _stream_line(msg if msg.endswith("\n") else msg + "\n"),
)
except Exception as e:
self.logger.error(f"OpenBao start failed: {e}")
# When common services are already healthy, run the service-layer
# migration inline so that the "Preparing Service Layer" steps are
# completed automatically — mirroring the UI shortcut behaviour.
if env_key in ("service", "k3s"):
service_ns = env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default"))
status_ok = all(
[
registry_ops.status(
namespace=env.get("REGISTRY_NAMESPACE", service_ns),
env=env,
mode=env.get("PROLE_MODE"),
),
openbao_ops.status(
namespace=service_ns,
env=env,
mode=env.get("PROLE_MODE"),
),
garage_store_ops.status(namespace=service_ns, env=env),
opentofu_ops.status(namespace=service_ns, env=env),
]
)
if status_ok:
self.logger.info(
"All common services healthy running service layer migration inline."
)
rc_migrate = state.controller.run_script(
"init_service_layer.sh",
args=["migrate"],
env=env,
on_line=_stream_line,
)
if rc_migrate != 0:
self.logger.error(
f"Service layer migration failed (code {rc_migrate})"
)
else:
self.logger.info("Service layer migration completed successfully.")
if progress:
progress("Cluster lifecycle management complete", 1.0)
class InitializationScriptsMilestone(Milestone):
def __init__(self):
super().__init__("init_scripts", "Initialization Scripts")
self.logger = logging.getLogger("InitializationScriptsMilestone")
def _regenerate_port_mapping_cfg(self, state: InstallerState, mode: str) -> None:
"""Regenerate conf/port-mapping.cfg from [Port Forwards] in prole.cfg data."""
conf_dir = inst_config.PROJECT_ROOT / "conf"
mapping_path = conf_dir / "port-mapping.cfg"
pf_section = state.config_data.get("Port Forwards", {}) or {}
prefix = (
"PORT_FORWARD_K3S_MAPPING_"
if mode == "k3s"
else "PORT_FORWARD_K3D_MAPPING_"
)
# Resolve ${NAMESPACE} from state
ns = (state.inputs.get("init_password.db_namespace", "") or "").strip()
if not ns:
ns = (
(state.config_data.get("Global", {}) or {}).get("NAMESPACE", "").strip()
)
if not ns:
ns = "default"
lines = [
"# Port mappings for Prole Tools (generated).",
"# Format: key: local=... remote=... ns=... svc=... address=...",
"",
]
for k, v in sorted(pf_section.items()):
if not k.startswith(prefix):
continue
parts = {}
for p in str(v).split(";"):
if "=" in p:
kv = p.split("=", 1)
if len(kv) == 2:
parts[kv[0].strip()] = kv[1].strip()
if "id" not in parts:
continue
m_id = parts["id"]
local = parts.get("hostPort", "")
remote = parts.get("servicePort", "")
m_ns = parts.get("namespace", "")
target = parts.get("target", "")
addr = parts.get("address", "0.0.0.0")
# Resolve ${NAMESPACE} references
if m_ns.startswith("${") and "NAMESPACE" in m_ns:
m_ns = ns
if local.startswith("${"):
continue # skip entries with unresolvable port refs
svc = target[4:] if target.startswith("svc/") else target
lines.append(
f"{m_id}: local={local} remote={remote} ns={m_ns} svc={svc} address={addr}"
)
try:
mapping_path.write_text("\n".join(lines) + "\n")
self.logger.info(f"Regenerated {mapping_path}")
except Exception as e:
self.logger.error(f"Failed to regenerate port-mapping.cfg: {e}")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
env = self._get_script_env(state)
mode = env.get("PROLE_MODE", "")
mode_args = ["--mode", mode] if mode else []
raw_min = (state.config_data.get("Global", {}) or {}).get(POLICY_CFG_KEY, "")
try:
min_nodes = int(str(raw_min).strip() or str(OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES))
except Exception:
min_nodes = OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES
opt_allowed, _count, opt_reason = evaluate_optional_workloads_allowed(
env=env, min_nodes=min_nodes
)
self.logger.info(f"[POLICY] {opt_reason}")
# Regenerate port-mapping.cfg from prole.cfg data to ensure it's clean
self._regenerate_port_mapping_cfg(state, mode or "k3d")
# Phase 1: common services (Python owners)
if progress:
progress("Running common services (python owners)...", 0.0)
service_ns = env.get("SERVICE_NAMESPACE", env.get("NAMESPACE", "default"))
registry_ns = env.get("REGISTRY_NAMESPACE", service_ns)
project_root = getattr(state.controller, "project_root", ".")
try:
registry_ops.update(
namespace=registry_ns,
env=env,
project_root=project_root,
mode=mode,
log=self.logger.info,
)
openbao_ops.update(
namespace=service_ns,
env=env,
project_root=project_root,
mode=mode,
log=self.logger.info,
)
garage_store_ops.update(
namespace=service_ns,
env=env,
project_root=project_root,
mode=mode,
log=self.logger.info,
)
opentofu_ops.update(
namespace=service_ns,
env=env,
project_root=project_root,
mode=mode,
log=self.logger.info,
)
except Exception as e:
msg = f"Common services (python owners) failed: {e}"
self.logger.error(msg)
raise Exception(msg)
# Phase 2: CNPG initialization (Python)
if progress:
progress("Initializing CloudNative-PG...", 1 / 6)
ns = env.get("DATABASE_NAMESPACE") or env.get("NAMESPACE", "default")
cluster_name = env.get("CLUSTER_NAME") or env.get("CNPG_CLUSTER_NAME") or "knoe-db"
# Ensure DB secrets exist before CNPG init (mirrors _step_init_scripts).
# env["DB_PASSWORD"] is already resolved by _script_env_for_namespace.
db_pw = (env.get("DB_PASSWORD") or state.inputs.get("init_password.db_password") or "").strip()
if db_pw and hasattr(state.controller, "ensure_db_k8s_secrets"):
try:
state.controller.ensure_db_k8s_secrets(ns, db_pw, log_fn=self.logger.info)
except Exception as _sec_err:
raise Exception(f"Failed to ensure DB secrets before CNPG init: {_sec_err}")
cnpg_initialize(
namespace=ns,
cluster_name=cluster_name,
env=env,
project_root=project_root,
log=self.logger.info,
mode=mode,
)
# Phase 3: remaining scripts (backup, kong, optional ingress)
post_scripts = ["init_cnpg_backup.sh", "init_kong.sh"]
if not opt_allowed:
self.logger.info(f"[SKIP] Monitoring disabled by policy: {opt_reason}")
state.config_data.setdefault("Monitoring", {})["STATUS"] = "Skipped"
if mode != "k3d":
post_scripts.append("init_nginx_ingress.sh")
for i, script in enumerate(post_scripts):
if progress:
progress(f"Running {script}...", (2 + i) / (2 + len(post_scripts)))
if script == "init_cnpg_backup.sh":
args: list[str] = mode_args + ["start"]
elif script == "init_kong.sh":
args = mode_args + ["start"]
else:
args = mode_args + ["initialize"]
rc = state.controller.run_script(
script, args=args, env=env, on_line=_stream_line
)
if rc != 0:
msg = f"Script {script} failed (code {rc})"
self.logger.error(msg)
raise Exception(msg)
if opt_allowed:
if progress:
progress("Running monitoring (python owner)...", (2 + len(post_scripts)) / (2 + len(post_scripts) + 1))
monitoring_ns = env.get("MONITORING_NAMESPACE", "monitoring")
try:
monitoring_ops.initialize(
namespace=monitoring_ns,
env=env,
mode=mode,
log=self.logger.info,
)
except Exception as e:
msg = f"Monitoring initialization failed: {e}"
self.logger.error(msg)
raise Exception(msg)
# Verify critical secrets
ns = env.get("NAMESPACE", "default")
self.logger.info(f"Verifying critical secrets in namespace {ns}...")
critical_secrets = ["knoe-db-user", "knoe-db-superuser", "cnpg-admin-key"]
missing_secrets = []
for secret in critical_secrets:
if self._run_cmd(f"kubectl get secret {secret} -n {ns}") != 0:
missing_secrets.append(secret)
if missing_secrets:
msg = f"CRITICAL: Missing secrets in namespace '{ns}': {', '.join(missing_secrets)}. Database initialization will fail."
self.logger.error(msg)
if progress:
progress(msg, 1.0)
raise Exception(msg)
# For k3d: start port-forward daemon and verify Grafana port-forward.
pf_ok = True
if mode == "k3d":
def _pf_line(line: str) -> None:
self.logger.info(line.rstrip("\n"))
env_pf = dict(env)
env_pf.setdefault("PORT_FORWARD_SKIP_VALIDATE", "1")
env_pf.setdefault("PORT_FORWARD_SKIP_WAIT", "1")
env_pf.setdefault("PORT_FORWARD_WAIT_TIMEOUT", "20")
env_pf.setdefault("PORT_FORWARD_WAIT_INTERVAL", "2")
rc_pf = state.controller.run_script(
"init_port_forwards.sh",
args=["--force", "--mode", mode, "start"],
env=env_pf,
on_line=_stream_line,
)
if rc_pf != 0:
pf_ok = False
self.logger.error(f"Port-forwards start failed (code {rc_pf})")
else:
for _ in range(10):
status_lines: list[str] = []
def _pf_status_line(line: str) -> None:
status_lines.append(line)
_stream_line(line)
rc_status = state.controller.run_script(
"init_port_forwards.sh",
args=["--mode", mode, "status", "grafana"],
env=env_pf,
on_line=_pf_status_line,
)
running = any(
l.lstrip().startswith("RUNNING") and "grafana" in l
for l in status_lines
)
if rc_status == 0 and running:
break
time.sleep(1)
else:
pf_ok = False
self.logger.error("Grafana port-forward not running after start.")
pf_section = state.config_data.setdefault("Port Forwards", {})
pf_section["STATUS"] = "Running" if pf_ok else "Failed"
if progress:
progress("Initialization scripts complete", 1.0)
class DeploymentMilestone(Milestone):
def __init__(self):
super().__init__("deployment", "Prole Deployment")
self.logger = logging.getLogger("DeploymentMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
if progress:
progress("Deploying Prole...", 0.3)
env = self._get_script_env(state)
mode = env.get("PROLE_MODE", "")
ns = env.get("DATABASE_NAMESPACE") or env.get("NAMESPACE", "default")
cluster_name = env.get("CLUSTER_NAME") or env.get("CNPG_CLUSTER_NAME") or "knoe-db"
project_root = getattr(state.controller, "project_root", ".")
cnpg_deploy(
namespace=ns,
cluster_name=cluster_name,
env=env,
project_root=project_root,
log=self.logger.info,
mode=mode,
)
self.logger.info("Deployment successful")
if progress:
progress("Deployment complete", 1.0)
class GitOpsMilestone(Milestone):
def __init__(self):
super().__init__("gitops", "GitOps")
self.logger = logging.getLogger("GitOpsMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
enabled = self._parse_bool(
state.inputs.get("init_cluster.gitops_enabled", "False"), default=False
)
if not enabled:
state.config_data.setdefault("GitOps", {})["STATUS"] = "Skipped"
return
# Resolve git provider: [Inputs] gitops.git_provider takes precedence,
# then [Optional Features] GITOPS_PROVIDER, then default to gitea.
provider = (
str(state.inputs.get("gitops.git_provider", "") or "").strip().lower()
or str(
(state.config_data.get("Optional Features") or {}).get(
"GITOPS_PROVIDER", ""
)
or ""
)
.strip()
.lower()
or "gitea"
)
is_gitlab = provider == "gitlab"
label = "GitLab" if is_gitlab else "Gitea"
if progress:
progress(f"Deploying {label}...", 0.4)
env = self._get_script_env(state)
if is_gitlab:
ns = (
str(state.inputs.get("gitops.gitlab_namespace", "") or "").strip()
or str(state.inputs.get("gitops.namespace", "") or "").strip()
or "gitlab"
)
env["GITLAB_NAMESPACE"] = ns
node_selector = (
str(state.inputs.get("gitops.node_selector", "") or "").strip()
)
if node_selector:
env["GITLAB_NODE_SELECTOR"] = node_selector
env["NODE_SELECTOR"] = node_selector
script = "init_gitlab.sh"
else:
ns = (
state.inputs.get("gitops.namespace", "")
or env.get("GITEA_NAMESPACE")
or "gitea"
).strip()
env["GITEA_NAMESPACE"] = ns
script = "init_gitea.sh"
mode = env.get("PROLE_MODE", "")
args: list[str] = []
if mode:
args.extend(["--mode", mode])
args.extend(["--namespace", ns])
conf_dir = prole_conf.resolve_prole_conf_dir(state.controller.project_root)
cfg_path = prole_conf.entrypoint_path(conf_dir)
if cfg_path.exists():
args.extend(["-c", str(cfg_path)])
self.logger.info(f"{script} {' '.join(args)}")
rc = state.controller.run_script(
script, args=args, env=env, on_line=_stream_line
)
if rc == 0:
gitops_sec = state.config_data.setdefault("GitOps", {})
gitops_sec["STATUS"] = "Deployed"
gitops_sec["GITOPS_NAMESPACE"] = ns
if progress:
progress(f"{label} deployed", 1.0)
else:
gitops_sec = state.config_data.setdefault("GitOps", {})
gitops_sec["STATUS"] = "Attempted"
gitops_sec["GITOPS_NAMESPACE"] = ns
msg = f"{label} deploy failed (code {rc})"
self.logger.error(msg)
raise Exception(msg)
class KerberosMilestone(Milestone):
def __init__(self):
super().__init__("kerberos", "Kerberos Initialization")
self.logger = logging.getLogger("KerberosMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
enabled = self._parse_bool(
state.inputs.get("kerberos_config.enabled", "False"), default=False
)
if not enabled:
realm = str(state.inputs.get("kerberos_config.realm", "") or "").strip()
kdc = str(state.inputs.get("kerberos_config.kdc", "") or "").strip()
if realm and kdc:
self.logger.info(
"Kerberos enabled implicitly because realm and kdc are configured."
)
enabled = True
krb_sec = state.config_data.setdefault("Kerberos Authentication", {})
if not enabled:
krb_sec["STATUS"] = "Skipped"
if progress:
progress("Kerberos disabled; skipping", 1.0)
return
env = self._get_script_env(state)
mode = env.get("PROLE_MODE", "")
args = ["--mode", mode] if mode else []
args.append("initialize")
if progress:
progress("Initializing Kerberos...", 0.4)
rc = state.controller.run_script(
"init_kerberos.sh", args=args, env=env, on_line=_stream_line
)
if rc == 0:
krb_sec["STATUS"] = "Initialized"
if progress:
progress("Kerberos initialized", 1.0)
else:
krb_sec["STATUS"] = "Attempted"
msg = f"Kerberos initialization failed (code {rc})"
self.logger.error(msg)
raise Exception(msg)
def _resolve_supabase_deploy_mode(state: InstallerState) -> str:
from knoe.core.env import _normalize_cluster_env
deploy_mode = (
os.environ.get("SUPABASE_DEPLOY_MODE") or os.environ.get("SUPABASE_MODE") or ""
).strip()
if deploy_mode:
return deploy_mode
cluster_env = _normalize_cluster_env(state.inputs.get("init_cluster.cluster_env", "dev"))
if cluster_env in ("service", "prod"):
return "k8s"
return "k3d"
class SupabaseImagePreloadMilestone(Milestone):
def __init__(self):
super().__init__("supabase_images_preload", "Supabase Image Preload")
self.logger = logging.getLogger("SupabaseImagePreloadMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
enabled = self._parse_bool(
state.inputs.get("init_cluster.supabase_enabled", "False"),
default=False,
)
supa_sec = state.config_data.setdefault("Supabase", {})
if not enabled:
supa_sec["IMAGES_STATUS"] = "Skipped"
if progress:
progress("Supabase disabled; skipping image preload", 1.0)
return
raw_min = (state.config_data.get("Global", {}) or {}).get(POLICY_CFG_KEY, "")
try:
min_nodes = int(str(raw_min).strip() or str(OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES))
except Exception:
min_nodes = OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES
opt_allowed, _count, opt_reason = evaluate_optional_workloads_allowed(
env=self._get_script_env(state), min_nodes=min_nodes
)
self.logger.info(f"[POLICY] {opt_reason}")
if not opt_allowed:
supa_sec["IMAGES_STATUS"] = "Skipped"
if progress:
progress(f"Supabase image preload skipped by policy: {opt_reason}", 1.0)
return
from knoe.core.env import _parse_bool
env = self._get_script_env(state)
script_path = state.controller.project_root / "supabase" / "deploy.sh"
if not script_path.exists():
self.logger.error(f"Supabase deploy script not found: {script_path}")
supa_sec["IMAGES_STATUS"] = "Attempted"
return
deploy_mode = _resolve_supabase_deploy_mode(state)
if deploy_mode != "k3d":
supa_sec["IMAGES_STATUS"] = "Skipped"
if progress:
progress(
f"Supabase image preload not required for mode '{deploy_mode}'",
1.0,
)
return
args = ["--mode", deploy_mode, "--prefetch-images-only"]
conf_dir = prole_conf.resolve_prole_conf_dir(state.controller.project_root)
cfg_path = prole_conf.entrypoint_path(conf_dir)
if cfg_path.exists():
args.extend(["-c", str(cfg_path)])
if _parse_bool(os.environ.get("SUPABASE_USE_DEV_COMPOSE"), False):
args.append("--with-dev-helpers")
if progress:
progress("Preloading Supabase images...", 0.4)
self.logger.info(f"supabase/deploy.sh {' '.join(args)}")
rc = subprocess.run(
["bash", str(script_path)] + args,
env=env,
).returncode
if rc == 0:
supa_sec["IMAGES_STATUS"] = "Prepared"
state.inputs["supabase.images_preloaded"] = "True"
if progress:
progress("Supabase images preloaded", 1.0)
else:
supa_sec["IMAGES_STATUS"] = "Attempted"
msg = f"Supabase image preload failed (code {rc})"
self.logger.error(msg)
raise Exception(msg)
class SupabaseMilestone(Milestone):
def __init__(self):
super().__init__("supabase", "Supabase Deployment")
self.logger = logging.getLogger("SupabaseMilestone")
def execute(
self, state: InstallerState, progress: ProgressCallback | None = None
) -> None:
enabled = self._parse_bool(
state.inputs.get("init_cluster.supabase_enabled", "False"),
default=False,
)
supa_sec = state.config_data.setdefault("Supabase", {})
if not enabled:
supa_sec["STATUS"] = "Skipped"
if progress:
progress("Supabase disabled; skipping", 1.0)
return
raw_min = (state.config_data.get("Global", {}) or {}).get(POLICY_CFG_KEY, "")
try:
min_nodes = int(str(raw_min).strip() or str(OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES))
except Exception:
min_nodes = OPTIONAL_WORKLOADS_MIN_READY_SCHEDULABLE_NODES
opt_allowed, _count, opt_reason = evaluate_optional_workloads_allowed(
env=self._get_script_env(state), min_nodes=min_nodes
)
self.logger.info(f"[POLICY] {opt_reason}")
if not opt_allowed:
supa_sec["STATUS"] = "Skipped"
if progress:
progress(f"Supabase skipped by policy: {opt_reason}", 1.0)
return
from knoe.core.env import _parse_bool
env = self._get_script_env(state)
script_path = state.controller.project_root / "supabase" / "deploy.sh"
if not script_path.exists():
self.logger.error(f"Supabase deploy script not found: {script_path}")
supa_sec["STATUS"] = "Attempted"
return
deploy_mode = _resolve_supabase_deploy_mode(state)
args = ["--mode", deploy_mode]
conf_dir = prole_conf.resolve_prole_conf_dir(state.controller.project_root)
cfg_path = prole_conf.entrypoint_path(conf_dir)
if cfg_path.exists():
args.extend(["-c", str(cfg_path)])
if _parse_bool(os.environ.get("SUPABASE_USE_DEV_COMPOSE"), False):
args.append("--with-dev-helpers")
if _parse_bool(os.environ.get("SUPABASE_FOREGROUND"), False):
args.append("--foreground")
if _parse_bool(state.inputs.get("supabase.images_preloaded", "False"), False):
args.append("--skip-prefetch")
self.logger.info(f"supabase/deploy.sh {' '.join(args)}")
rc = subprocess.run(
["bash", str(script_path)] + args,
env=env,
).returncode
if rc == 0:
supa_sec["STATUS"] = "Deployed"
if progress:
progress("Supabase deployed", 1.0)
else:
supa_sec["STATUS"] = "Attempted"
msg = f"Supabase deploy failed (code {rc})"
self.logger.error(msg)
raise Exception(msg)