prole/knoe/core/milestones.py
chrisfu 3f34fa8b32 fix(installer): k3s --reset path hardening (kdc deploy, no-TTY 1password, context overrides)
Five fixes Junie surfaced while running the kdc-trust-reset-repeatable
Junie brief end-to-end (companion to commit 6f99f95). All hit during
the unattended `install.sh --mode k3s --reset` pipeline.

  - knoe/core/milestones.py (KerberosMilestone):
    For k3s and k3d modes, deploy the KDC pod via `init_kdc.sh start`
    before running init_kerberos.sh. init_kerberos.sh only chains into
    init_kdc.sh when PROLE_KDC_STANDALONE=1; without this hook the
    cluster came up with no KDC pod and the cross-realm trust principals
    had nowhere to land.

  - knoe/milestone.py (Milestone._get_script_env):
    Clear KUBECTL_CONTEXT in addition to KUBECONTEXT so stale entries
    from a different machine's cfg don't override the kubeconfig's
    own current-context.

  - etc/knoe_cfg.sh (_knoe_read_cfg):
    Skip KUBECTL_CONTEXT / KUBE_CONTEXT_NAME / KUBECONTEXT entries when
    reading cfg in k3s mode. Same theme: kubeconfig current-context is
    authoritative.

  - etc/init_1password.sh + knoe/core/onepassword.py:
    When running non-interactively (no TTY on stdin) and no `op`
    session exists, skip rather than hang on `op signin`. Lets the
    unattended pipeline proceed for k3s/k3d where in-cluster secrets
    are managed separately from 1Password.

Co-authored-by: Junie <junie@jetbrains.com>
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-11 13:34:10 -07:00

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