mirror of
https://github.com/dredx/prole.git
synced 2026-09-23 11:03:59 +00:00
- Ensure k3s mode uses the k3s registry endpoint and avoid localhost/k3d image prefixes. - Make ArgoCD repo-server cmp symlink creation idempotent. - Normalize common-core provisioning to knoe-system and add repair-time dedupe of stray default-namespace installs. - Add k3s MariaDB datastore/refresh playbooks and regression tests.
1852 lines
67 KiB
Python
1852 lines
67 KiB
Python
"""Cluster lifecycle management (k3d / k3s / k8s)."""
|
|
|
|
import json
|
|
import os
|
|
import re
|
|
import subprocess
|
|
import sys
|
|
import threading
|
|
import time
|
|
import urllib.error
|
|
import urllib.request
|
|
from pathlib import Path
|
|
import tkinter as tk
|
|
from tkinter import ttk, messagebox, filedialog
|
|
from installer.screen import TerminalConsole
|
|
from installer import screen as ui
|
|
from installer.core.env import (
|
|
PROJECT_ROOT,
|
|
_default_opentofu_pipeline_url,
|
|
_deployment_mode_from_env,
|
|
_deployment_target_label,
|
|
_find_kubeconfig_file,
|
|
_k3d_prole_data_volume_args,
|
|
_kubectl_base_cmd_for_k3s as _kubectl_base_cmd_for_k3s_fn,
|
|
_looks_like_k8s_bearer_token,
|
|
_normalize_cluster_env,
|
|
_normalize_k3s_token,
|
|
_resolve_k3s_connection as _resolve_k3s_connection_fn,
|
|
_safe_str,
|
|
)
|
|
from installer.config import (
|
|
_encrypt_cfg_secret,
|
|
_merge_kubeconfig,
|
|
_write_k3s_kubeconfig,
|
|
)
|
|
from installer.core.actions import _reset_k3s_namespace
|
|
|
|
|
|
class ClusterScreenMixin:
|
|
"""Cluster lifecycle management (k3d / k3s / k8s)."""
|
|
|
|
def _render_init_cluster_page(self):
|
|
# Letterhead at top right
|
|
content_width = self.bg_canvas.winfo_width() or 975
|
|
right_margin = content_width - 48
|
|
|
|
ui.canvas_text(
|
|
self,
|
|
right_margin,
|
|
40,
|
|
"Prole",
|
|
fill="#6e6e73",
|
|
font=("SF Pro Text", 32, "bold"),
|
|
anchor="ne",
|
|
)
|
|
ui.canvas_text(
|
|
self,
|
|
right_margin,
|
|
85,
|
|
"Infrastructure Automated.",
|
|
fill="#6e6e73",
|
|
font=("SF Pro Text", 18),
|
|
anchor="ne",
|
|
)
|
|
|
|
# Tighten vertical spacing on this screen so the main status area (traffic lights +
|
|
# console output) has enough room in the fixed-size window.
|
|
self._render_title("Cluster Environment", y=110)
|
|
self._render_paragraph(
|
|
"Select a cluster environment and ensure the cluster is running and common services can be deployed.",
|
|
y=155,
|
|
)
|
|
|
|
# Cluster Selection (Radio Buttons)
|
|
x_label = 48
|
|
y = 210
|
|
self._canvas_items.append(
|
|
ui.canvas_text(
|
|
self,
|
|
x_label,
|
|
y,
|
|
"Select Environment:",
|
|
fill="black",
|
|
font=("SF Pro Text", 14, "bold"),
|
|
)
|
|
)
|
|
|
|
y += 28
|
|
cluster_options = [("dev", "Dev"), ("service", "Service"), ("prod", "Prod")]
|
|
|
|
# We need to trace cluster_env if not already traced
|
|
if not hasattr(self, "_cluster_env_trace"):
|
|
self._cluster_env_trace = self.cluster_env.trace_add(
|
|
"write", self._on_cluster_env_change
|
|
)
|
|
|
|
rb_x = x_label + 20
|
|
rb_gap = 150
|
|
for idx, (val, name) in enumerate(cluster_options):
|
|
rb = tk.Radiobutton(
|
|
self.bg_canvas,
|
|
text=name,
|
|
variable=self.cluster_env,
|
|
value=val,
|
|
bg="white",
|
|
fg="black",
|
|
activebackground="white",
|
|
selectcolor="white",
|
|
font=("SF Pro Text", 11),
|
|
)
|
|
rb_window = self.bg_canvas.create_window(
|
|
rb_x + idx * rb_gap, y, window=rb, anchor="nw"
|
|
)
|
|
self._canvas_items.append(rb_window)
|
|
self._overlay_widgets.append(rb)
|
|
|
|
selected_env_key = self._cluster_env_key()
|
|
y += 36
|
|
|
|
if selected_env_key == "dev":
|
|
# Combined Select / Add Cluster widget
|
|
self._canvas_items.append(
|
|
ui.canvas_text(
|
|
self,
|
|
x_label,
|
|
y,
|
|
"Cluster:",
|
|
fill="black",
|
|
font=("SF Pro Text", 12, "bold"),
|
|
)
|
|
)
|
|
y += 28
|
|
|
|
clusters = self._get_k3d_cluster_list()
|
|
combo = ttk.Combobox(
|
|
self.bg_canvas,
|
|
textvariable=self.selected_k3d_cluster,
|
|
values=clusters,
|
|
width=30,
|
|
)
|
|
if not self.selected_k3d_cluster.get() and clusters:
|
|
self.selected_k3d_cluster.set(clusters[0])
|
|
|
|
combo.bind(
|
|
"<<ComboboxSelected>>", lambda _evt: self._on_k3d_cluster_select()
|
|
)
|
|
|
|
combo_win = self.bg_canvas.create_window(
|
|
x_label + 20, y, window=combo, anchor="nw"
|
|
)
|
|
self._canvas_items.append(combo_win)
|
|
self._overlay_widgets.append(combo)
|
|
|
|
add_btn = tk.Button(
|
|
self.bg_canvas,
|
|
text="Add",
|
|
command=self._on_create_k3d_cluster,
|
|
bg="#F5F5DC",
|
|
fg="black",
|
|
activebackground="#E5E5D5",
|
|
activeforeground="black",
|
|
highlightbackground="#F5F5DC",
|
|
highlightthickness=0,
|
|
relief="flat",
|
|
font=("SF Pro Text", 10),
|
|
padx=10,
|
|
)
|
|
add_win = self.bg_canvas.create_window(
|
|
x_label + 290, y - 4, window=add_btn, anchor="nw"
|
|
)
|
|
self._canvas_items.append(add_win)
|
|
self._overlay_widgets.append(add_btn)
|
|
|
|
delete_btn = tk.Button(
|
|
self.bg_canvas,
|
|
text="Delete",
|
|
command=self._on_delete_k3d_cluster,
|
|
bg="#F5F5DC",
|
|
fg="#cc0000",
|
|
activebackground="#E5E5D5",
|
|
activeforeground="#cc0000",
|
|
highlightbackground="#F5F5DC",
|
|
highlightthickness=0,
|
|
relief="flat",
|
|
font=("SF Pro Text", 10),
|
|
padx=10,
|
|
)
|
|
delete_win = self.bg_canvas.create_window(
|
|
x_label + 350, y - 4, window=delete_btn, anchor="nw"
|
|
)
|
|
self._canvas_items.append(delete_win)
|
|
self._overlay_widgets.append(delete_btn)
|
|
y += 40
|
|
|
|
else:
|
|
# Service / Prod Section
|
|
# Context selection for both Service and Prod (managed via kubectx)
|
|
self._canvas_items.append(
|
|
ui.canvas_text(
|
|
self,
|
|
x_label,
|
|
y,
|
|
"Kubernetes Context:",
|
|
fill="black",
|
|
font=("SF Pro Text", 12, "bold"),
|
|
)
|
|
)
|
|
values = self._get_kubectx_list()
|
|
combo = ttk.Combobox(
|
|
self.bg_canvas,
|
|
textvariable=self.selected_kubectx,
|
|
values=values,
|
|
state="readonly",
|
|
width=40,
|
|
)
|
|
combo.bind("<<ComboboxSelected>>", self._on_kubectx_select)
|
|
if not self.selected_kubectx.get() and values:
|
|
if selected_env_key == "service" and "prole-k3s" in values:
|
|
self.selected_kubectx.set("prole-k3s")
|
|
# Auto-switch to prole-k3s if we just loaded the page and it's selected
|
|
self._switch_kubectx("prole-k3s")
|
|
else:
|
|
self.selected_kubectx.set(values[0])
|
|
# No auto-switch for others to avoid surprises, but we should probably
|
|
# if it's the only one.
|
|
combo_win = self.bg_canvas.create_window(
|
|
x_label + 190, y - 6, window=combo, anchor="nw"
|
|
)
|
|
self._canvas_items.append(combo_win)
|
|
self._overlay_widgets.append(combo)
|
|
# Keep references for UI tests / layout verification.
|
|
self._kubectx_combo = combo
|
|
self._kubectx_combo_canvas_window = combo_win
|
|
ui.canvas_text(
|
|
self,
|
|
x_label + 580,
|
|
y + 2,
|
|
"via kubectx",
|
|
fill="#6e6e73",
|
|
font=("SF Pro Text", 10),
|
|
)
|
|
y += 34
|
|
|
|
if selected_env_key == "prod":
|
|
# Staging directory input
|
|
self._canvas_items.append(
|
|
ui.canvas_text(
|
|
self,
|
|
x_label,
|
|
y,
|
|
"Local Artifact Staging Directory:",
|
|
fill="black",
|
|
font=("SF Pro Text", 12, "bold"),
|
|
)
|
|
)
|
|
entry = tk.Entry(
|
|
self.bg_canvas,
|
|
textvariable=self.prod_artifacts_path,
|
|
width=50,
|
|
bg="white",
|
|
fg="black",
|
|
insertbackground="black",
|
|
highlightbackground="#CCCCCC",
|
|
highlightthickness=1,
|
|
relief="flat",
|
|
font=("SF Pro Text", 11),
|
|
)
|
|
entry_win = self.bg_canvas.create_window(
|
|
x_label + 260, y - 6, window=entry, anchor="nw"
|
|
)
|
|
self._canvas_items.append(entry_win)
|
|
self._overlay_widgets.append(entry)
|
|
|
|
browse_btn = tk.Button(
|
|
self.bg_canvas,
|
|
text="Browse...",
|
|
command=lambda: self.prod_artifacts_path.set(
|
|
filedialog.askdirectory() or self.prod_artifacts_path.get()
|
|
),
|
|
)
|
|
browse_win = self.bg_canvas.create_window(
|
|
x_label + 770, y - 8, window=browse_btn, anchor="nw"
|
|
)
|
|
self._canvas_items.append(browse_win)
|
|
self._overlay_widgets.append(browse_btn)
|
|
y += 34
|
|
|
|
ui.canvas_text(
|
|
self,
|
|
x_label + 20,
|
|
y,
|
|
"(Used by etc/deploy_pipeline.sh --mode gcp)",
|
|
fill="#6e6e73",
|
|
font=("SF Pro Text", 10),
|
|
)
|
|
y += 34
|
|
|
|
# Common Services Status Check (Notebook)
|
|
# Show if in Service/Prod or if a cluster is selected in Dev
|
|
show_notebook = False
|
|
if selected_env_key == "dev":
|
|
if self.selected_k3d_cluster.get():
|
|
show_notebook = True
|
|
else:
|
|
show_notebook = True
|
|
|
|
if show_notebook:
|
|
ui.canvas_text(
|
|
self,
|
|
x_label,
|
|
y,
|
|
"Common Core Services Namespace:",
|
|
fill="black",
|
|
font=("SF Pro Text", 12, "bold"),
|
|
)
|
|
ns_entry = tk.Entry(
|
|
self.bg_canvas,
|
|
textvariable=self.service_namespace,
|
|
width=20,
|
|
bg="white",
|
|
fg="black",
|
|
insertbackground="black",
|
|
highlightbackground="#CCCCCC",
|
|
highlightthickness=1,
|
|
relief="flat",
|
|
font=("SF Pro Text", 11),
|
|
)
|
|
ns_win = self.bg_canvas.create_window(
|
|
x_label + 250, y - 6, window=ns_entry, anchor="nw"
|
|
)
|
|
self._canvas_items.append(ns_win)
|
|
self._overlay_widgets.append(ns_entry)
|
|
# Keep references for UI tests / layout verification.
|
|
self._service_namespace_entry = ns_entry
|
|
self._service_namespace_entry_canvas_window = ns_win
|
|
|
|
save_btn = tk.Button(
|
|
self.bg_canvas,
|
|
text="Save",
|
|
command=self._validate_and_save_cluster_config,
|
|
bg="#F5F5DC",
|
|
fg="black",
|
|
activebackground="#E5E5D5",
|
|
highlightbackground="#F5F5DC",
|
|
highlightthickness=0,
|
|
relief="flat",
|
|
font=("SF Pro Text", 10),
|
|
padx=10,
|
|
)
|
|
save_win = self.bg_canvas.create_window(
|
|
x_label + 480, y - 10, window=save_btn, anchor="nw"
|
|
)
|
|
self._canvas_items.append(save_win)
|
|
self._overlay_widgets.append(save_btn)
|
|
y += 34
|
|
|
|
# Traffic light status indicators for each service component
|
|
self._service_traffic_lights = {}
|
|
components = [
|
|
"Argo",
|
|
"CertMgr",
|
|
"Garage",
|
|
"Kong",
|
|
"OpenBao",
|
|
"OpenTofu",
|
|
]
|
|
columns = 3
|
|
col_gap = 180
|
|
row_gap = 24
|
|
start_x = x_label + 20
|
|
for idx, comp in enumerate(components):
|
|
row = idx // columns
|
|
col = idx % columns
|
|
tl_x = start_x + (col * col_gap)
|
|
tl_y = y + (row * row_gap)
|
|
# Circle indicator (default grey)
|
|
indicator = self.bg_canvas.create_oval(
|
|
tl_x,
|
|
tl_y,
|
|
tl_x + 14,
|
|
tl_y + 14,
|
|
fill="#8e8e93",
|
|
outline="#8e8e93",
|
|
)
|
|
label = ui.canvas_text(
|
|
self,
|
|
tl_x + 20,
|
|
tl_y,
|
|
comp,
|
|
fill="black",
|
|
font=("SF Pro Text", 11),
|
|
)
|
|
self._canvas_items.extend([indicator, label])
|
|
self._service_traffic_lights[comp.lower()] = indicator
|
|
rows = (len(components) + columns - 1) // columns
|
|
y += (rows * row_gap)
|
|
|
|
# Repair button should remain reachable; place it adjacent to the status area
|
|
# instead of pushing it below the console.
|
|
btn_text = "Repair"
|
|
deploy_btn = tk.Button(
|
|
self.bg_canvas,
|
|
text=btn_text,
|
|
command=self._deploy_k3s_services,
|
|
bg="#F5F5DC",
|
|
fg="black",
|
|
activebackground="#E5E5D5",
|
|
highlightbackground="#F5F5DC",
|
|
highlightthickness=0,
|
|
relief="flat",
|
|
font=("SF Pro Text", 10),
|
|
padx=10,
|
|
pady=5,
|
|
)
|
|
self._k3s_deploy_btn = deploy_btn
|
|
|
|
notebook = ttk.Notebook(self.bg_canvas)
|
|
status_frame = tk.Frame(notebook, bg="white")
|
|
deploy_frame = tk.Frame(notebook, bg="white")
|
|
status_console = ui.TerminalConsole(status_frame)
|
|
deploy_console = ui.TerminalConsole(deploy_frame)
|
|
status_console.pack(fill="both", expand=True)
|
|
deploy_console.pack(fill="both", expand=True)
|
|
notebook.add(status_frame, text="Status")
|
|
notebook.add(deploy_frame, text="Repair Output")
|
|
|
|
self._k3s_service_notebook = notebook
|
|
self._k3s_service_status_console = status_console
|
|
self._k3s_service_deploy_console = deploy_console
|
|
self._k3s_service_status_tab = status_frame
|
|
self._k3s_service_deploy_tab = deploy_frame
|
|
|
|
console_width = 900
|
|
console_height = 260
|
|
|
|
deploy_win = self.bg_canvas.create_window(
|
|
x_label + console_width - 96,
|
|
y - 34,
|
|
window=deploy_btn,
|
|
anchor="nw",
|
|
)
|
|
self._canvas_items.append(deploy_win)
|
|
self._overlay_widgets.append(deploy_btn)
|
|
self._k3s_deploy_btn_canvas_window = deploy_win
|
|
|
|
console_window = self.bg_canvas.create_window(
|
|
x_label,
|
|
y,
|
|
window=notebook,
|
|
anchor="nw",
|
|
width=console_width,
|
|
height=console_height,
|
|
)
|
|
self._canvas_items.append(console_window)
|
|
self._k3s_service_notebook_canvas_window = console_window
|
|
self._overlay_widgets.extend(
|
|
[notebook, status_frame, deploy_frame, status_console, deploy_console]
|
|
)
|
|
|
|
# Always start on the Status tab with a clean console
|
|
try:
|
|
status_console.clear()
|
|
self._k3s_service_notebook.select(self._k3s_service_status_tab)
|
|
except Exception:
|
|
pass
|
|
|
|
y += console_height + 12
|
|
self._verify_k3s_services()
|
|
|
|
def _cluster_env_key(self, env_label: str | None = None) -> str:
|
|
"""Return normalized cluster env key ('dev', 'service', 'prod')."""
|
|
value = env_label
|
|
if value is None:
|
|
try:
|
|
value = self.cluster_env.get()
|
|
except Exception:
|
|
value = ""
|
|
|
|
try:
|
|
if hasattr(value, "get") and callable(value.get):
|
|
value = value.get()
|
|
except Exception:
|
|
pass
|
|
|
|
if not value:
|
|
try:
|
|
value = (
|
|
(self.prole_cfg_data.get("Initialize Cluster", {}) or {})
|
|
.get("ENVIRONMENT", "")
|
|
.strip()
|
|
)
|
|
except Exception:
|
|
value = ""
|
|
if not value:
|
|
try:
|
|
value = (
|
|
(self.prole_cfg_data.get("Global", {}) or {})
|
|
.get("CLUSTER_ENV", "")
|
|
.strip()
|
|
)
|
|
except Exception:
|
|
value = ""
|
|
|
|
key = _normalize_cluster_env(value)
|
|
if not key:
|
|
return "dev"
|
|
return key
|
|
|
|
def _get_service_namespace(self) -> str:
|
|
"""Resolve service namespace with UI/env/config fallbacks."""
|
|
try:
|
|
ns = (self.service_namespace.get() or "").strip()
|
|
except Exception:
|
|
ns = ""
|
|
if not ns:
|
|
try:
|
|
ns = (
|
|
(self.prole_cfg_data.get("Global", {}) or {})
|
|
.get("SERVICE_NAMESPACE", "")
|
|
.strip()
|
|
)
|
|
except Exception:
|
|
ns = ""
|
|
if not ns:
|
|
ns = (os.environ.get("SERVICE_NAMESPACE") or "").strip()
|
|
|
|
default_ns = "default"
|
|
try:
|
|
if self._cluster_env_key() == "service":
|
|
default_ns = "knoe-system"
|
|
except Exception:
|
|
pass
|
|
|
|
return ns or default_ns
|
|
|
|
def _on_cluster_env_change(self, *args):
|
|
self._set_deploy_target_from_cluster_env()
|
|
# Trigger status check or refresh UI
|
|
env_key = self._cluster_env_key()
|
|
if env_key == "service":
|
|
try:
|
|
self._apply_k3s_defaults()
|
|
except Exception:
|
|
pass
|
|
# Auto-switch to prole-k3s if it exists
|
|
if hasattr(self, "selected_kubectx"):
|
|
self.selected_kubectx.set("prole-k3s")
|
|
self._switch_kubectx("prole-k3s")
|
|
elif env_key == "dev":
|
|
name = self.selected_k3d_cluster.get()
|
|
if name:
|
|
if hasattr(self, "selected_kubectx"):
|
|
self.selected_kubectx.set(f"k3d-{name}")
|
|
self._switch_kubectx(f"k3d-{name}")
|
|
self._verify_k3s_services()
|
|
self.show_page("init_cluster")
|
|
|
|
def _deployment_mode(self) -> str:
|
|
try:
|
|
env_val = self.cluster_env.get()
|
|
except Exception:
|
|
env_val = ""
|
|
return _deployment_mode_from_env(env_val)
|
|
|
|
def _deployment_pipeline_url(self, target_key: str) -> str:
|
|
env_url = (
|
|
os.environ.get("PROLE_OPENTOFU_URL") or os.environ.get("OPENTOFU_URL") or ""
|
|
).strip()
|
|
if env_url:
|
|
return env_url
|
|
section = None
|
|
if target_key == "service":
|
|
section = "Service Cluster (k3s)"
|
|
elif target_key == "prod":
|
|
section = "Prod Cluster (k8s)"
|
|
if section:
|
|
url = (
|
|
(self.prole_cfg_data.get(section, {}) or {})
|
|
.get("PIPELINE_URL", "")
|
|
.strip()
|
|
)
|
|
if url:
|
|
return url
|
|
return _default_opentofu_pipeline_url()
|
|
|
|
def _set_deploy_target_from_cluster_env(self):
|
|
if not hasattr(self, "deploy_target"):
|
|
return
|
|
if getattr(self, "_syncing_deploy_target", False):
|
|
return
|
|
self._syncing_deploy_target = True
|
|
try:
|
|
self.deploy_target.set(_deployment_target_label(self.cluster_env.get()))
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
self._syncing_deploy_target = False
|
|
|
|
def _on_deploy_target_change(self, *args):
|
|
if getattr(self, "_syncing_deploy_target", False):
|
|
return
|
|
target = ""
|
|
try:
|
|
target = self.deploy_target.get()
|
|
except Exception:
|
|
target = ""
|
|
key = _normalize_cluster_env(target)
|
|
if key == "dev":
|
|
desired = "dev"
|
|
elif key == "service":
|
|
desired = "service"
|
|
elif key == "prod":
|
|
desired = "prod"
|
|
else:
|
|
return
|
|
self._syncing_deploy_target = True
|
|
try:
|
|
if self.cluster_env.get() != desired:
|
|
self.cluster_env.set(desired)
|
|
finally:
|
|
self._syncing_deploy_target = False
|
|
|
|
|
|
def _dashboard_kong_status(self, env: str | None = None) -> dict:
|
|
"""Check Kubernetes Dashboard (Kong) pod health and return status info."""
|
|
ns = "kubernetes-dashboard"
|
|
base_cmd = self._kubectl_base_cmd(env)
|
|
|
|
def _pod_healthy(pod_status: dict) -> bool:
|
|
phase = (pod_status.get("phase") or "").strip()
|
|
if phase not in ("Running", "Succeeded"):
|
|
return False
|
|
statuses = pod_status.get("containerStatuses") or []
|
|
if not statuses:
|
|
return phase == "Succeeded"
|
|
for cs in statuses:
|
|
if cs.get("ready") is True:
|
|
continue
|
|
state = cs.get("state") or {}
|
|
waiting = state.get("waiting") or {}
|
|
reason = (waiting.get("reason") or "").lower()
|
|
if reason:
|
|
return False
|
|
return False
|
|
return True
|
|
|
|
# Ensure namespace exists
|
|
rc_ns, out_ns = self._run_cmd_capture(
|
|
base_cmd + ["get", "ns", ns, "-o", "name"], timeout=6
|
|
)
|
|
if rc_ns != 0 or ns not in (out_ns or ""):
|
|
return {
|
|
"ok": False,
|
|
"msg": "Dashboard (Kong): Missing",
|
|
"fill": "#ff9f0a",
|
|
"reset_pods": [],
|
|
}
|
|
|
|
rc, out = self._run_cmd_capture(
|
|
base_cmd + ["get", "pods", "-n", ns, "-o", "json"], timeout=8
|
|
)
|
|
if rc != 0:
|
|
return {
|
|
"ok": False,
|
|
"msg": "Dashboard (Kong): Unable to query",
|
|
"fill": "#ff9f0a",
|
|
"reset_pods": [],
|
|
}
|
|
try:
|
|
payload = json.loads(out or "{}")
|
|
except Exception:
|
|
return {
|
|
"ok": False,
|
|
"msg": "Dashboard (Kong): Unable to parse status",
|
|
"fill": "#ff9f0a",
|
|
"reset_pods": [],
|
|
}
|
|
|
|
kong_pods = []
|
|
for item in payload.get("items", []):
|
|
name = (item.get("metadata", {}) or {}).get("name", "")
|
|
name_lc = name.lower()
|
|
if "kong" in name_lc and "dashboard" in name_lc:
|
|
kong_pods.append(item)
|
|
|
|
if not kong_pods:
|
|
return {
|
|
"ok": False,
|
|
"msg": "Dashboard (Kong): Missing",
|
|
"fill": "#ff9f0a",
|
|
"reset_pods": [],
|
|
}
|
|
|
|
unhealthy = []
|
|
for item in kong_pods:
|
|
name = (item.get("metadata", {}) or {}).get("name", "")
|
|
status = item.get("status") or {}
|
|
if not _pod_healthy(status):
|
|
unhealthy.append(name)
|
|
|
|
if unhealthy:
|
|
return {
|
|
"ok": False,
|
|
"msg": "Dashboard (Kong): Unhealthy (safe to reset)",
|
|
"fill": "#ff9f0a",
|
|
"reset_pods": unhealthy,
|
|
}
|
|
|
|
return {
|
|
"ok": True,
|
|
"msg": "Dashboard (Kong): Running",
|
|
"fill": "#34c759",
|
|
"reset_pods": [],
|
|
}
|
|
|
|
def _set_cluster_env_message(
|
|
self, text: str, color: str = "#34c759", clear_after_ms: int | None = None
|
|
):
|
|
"""Set a transient status message on the Cluster Environment screen."""
|
|
|
|
def _apply():
|
|
if not hasattr(self, "_cluster_env_message_label"):
|
|
return
|
|
if not self.bg_canvas.winfo_exists():
|
|
return
|
|
try:
|
|
self.bg_canvas.itemconfig(
|
|
self._cluster_env_message_label, text=text, fill=color
|
|
)
|
|
except Exception:
|
|
pass
|
|
|
|
self.safe_after(_apply)
|
|
|
|
if clear_after_ms is not None and clear_after_ms >= 0:
|
|
|
|
def _clear():
|
|
if not hasattr(self, "_cluster_env_message_label"):
|
|
return
|
|
if not self.bg_canvas.winfo_exists():
|
|
return
|
|
try:
|
|
self.bg_canvas.itemconfig(self._cluster_env_message_label, text="")
|
|
except Exception:
|
|
pass
|
|
|
|
self.safe_after(_clear, delay=clear_after_ms)
|
|
|
|
def _authority_context_missing(self) -> bool:
|
|
"""Detect missing authority Docker context when Kerberos is enabled."""
|
|
try:
|
|
if not self.kerberos_enabled.get():
|
|
return False
|
|
except Exception:
|
|
return False
|
|
|
|
candidates = []
|
|
env_home = (os.environ.get("PROLE_HOME") or "").strip()
|
|
if env_home:
|
|
candidates.append(Path(env_home))
|
|
candidates.append(PROJECT_ROOT)
|
|
|
|
for base in candidates:
|
|
try:
|
|
if (base / "authority").is_dir():
|
|
return False
|
|
if (base / "prole" / "authority").is_dir():
|
|
return False
|
|
except Exception:
|
|
continue
|
|
return True
|
|
|
|
def _maybe_run_repair_pipeline(self, status: dict, kong_status: dict | None):
|
|
"""Run repair pipeline when cluster is ready and anomalies are detected."""
|
|
if not status.get("cluster_ok"):
|
|
return
|
|
|
|
anomalies = []
|
|
if kong_status and not kong_status.get("ok", True):
|
|
anomalies.append("dashboard")
|
|
if self._authority_context_missing():
|
|
anomalies.append("authority")
|
|
|
|
if not anomalies:
|
|
return
|
|
|
|
now = time.time()
|
|
if self._repair_inflight:
|
|
return
|
|
if now - self._repair_last_run_at < self._repair_cooldown_s:
|
|
return
|
|
|
|
self._repair_last_run_at = now
|
|
self._repair_inflight = True
|
|
|
|
def worker():
|
|
try:
|
|
env_key = status.get("env") or self._cluster_env_key()
|
|
if env_key == "service":
|
|
ns = self._get_service_namespace()
|
|
else:
|
|
ns = (self.db_namespace.get() or "").strip() or "default"
|
|
env = self._script_env_for_namespace(ns)
|
|
if env_key == "service":
|
|
env["PROLE_MODE"] = "k3s"
|
|
elif env_key == "dev":
|
|
env["PROLE_MODE"] = "k3d"
|
|
self._set_cluster_env_message(
|
|
"Repair pipeline started", "#34c759", clear_after_ms=3000
|
|
)
|
|
self.controller.run_script(
|
|
"repair_pipeline.sh",
|
|
args=["-n", ns],
|
|
env=env,
|
|
)
|
|
self._set_cluster_env_message(
|
|
"Repair pipeline complete", "#34c759", clear_after_ms=3000
|
|
)
|
|
except Exception as e:
|
|
self._set_cluster_env_message(
|
|
f"Repair pipeline failed: {e}", "#ff3b30", clear_after_ms=6000
|
|
)
|
|
finally:
|
|
self._repair_inflight = False
|
|
|
|
threading.Thread(target=worker, daemon=True).start()
|
|
|
|
def _check_opentofu_health(self) -> bool:
|
|
"""Check OpenTofu readiness in the selected namespace."""
|
|
try:
|
|
env_key = self._cluster_env_key()
|
|
except Exception:
|
|
env_key = "dev"
|
|
|
|
try:
|
|
if env_key == "service":
|
|
ns = self._get_service_namespace()
|
|
else:
|
|
ns = (self.db_namespace.get() or "").strip() or "default"
|
|
except Exception:
|
|
ns = "default"
|
|
|
|
try:
|
|
base_cmd = self._kubectl_base_cmd(env_key)
|
|
res_dep = subprocess.run(
|
|
base_cmd
|
|
+ [
|
|
"-n",
|
|
ns,
|
|
"get",
|
|
"deploy",
|
|
"opentofu",
|
|
"-o",
|
|
"jsonpath={.status.readyReplicas}",
|
|
],
|
|
capture_output=True,
|
|
text=True,
|
|
timeout=3,
|
|
)
|
|
ready_dep = res_dep.returncode == 0 and (
|
|
res_dep.stdout or "0"
|
|
).strip() not in ("", "0")
|
|
if ready_dep:
|
|
return True
|
|
|
|
res_sts = subprocess.run(
|
|
base_cmd
|
|
+ [
|
|
"-n",
|
|
ns,
|
|
"get",
|
|
"statefulset",
|
|
"opentofu",
|
|
"-o",
|
|
"jsonpath={.status.readyReplicas}",
|
|
],
|
|
capture_output=True,
|
|
text=True,
|
|
timeout=3,
|
|
)
|
|
ready_sts = res_sts.returncode == 0 and (
|
|
res_sts.stdout or "0"
|
|
).strip() not in ("", "0")
|
|
return ready_sts
|
|
except Exception:
|
|
return False
|
|
|
|
def _check_k8s_cluster(self, env_label: str) -> tuple[bool, str]:
|
|
env_key = self._cluster_env_key(env_label)
|
|
label = (
|
|
"K3s Cluster"
|
|
if env_key == "service"
|
|
else "Prod Cluster" if env_key == "prod" else "Kubernetes Cluster"
|
|
)
|
|
try:
|
|
kubectl = subprocess.run(["which", "kubectl"], capture_output=True)
|
|
if kubectl.returncode != 0:
|
|
return False, f"{label}: kubectl not found"
|
|
if env_key == "service":
|
|
kubeconfig = _find_kubeconfig_file()
|
|
if not kubeconfig:
|
|
server, token = self._k3s_connection_info()
|
|
if not server:
|
|
return False, f"{label}: missing server URL or kubeconfig"
|
|
cmd = self._kubectl_base_cmd(env_key) + ["cluster-info"]
|
|
res = subprocess.run(cmd, capture_output=True, text=True, timeout=8)
|
|
ok = res.returncode == 0
|
|
except Exception:
|
|
ok = False
|
|
msg = f"{label}: Connected" if ok else f"{label}: Not reachable"
|
|
return ok, msg
|
|
|
|
def _get_k3d_cluster_list(self) -> list[str]:
|
|
"""Get list of k3d clusters."""
|
|
try:
|
|
res = subprocess.run(
|
|
["k3d", "cluster", "list", "--no-headers"],
|
|
capture_output=True,
|
|
text=True,
|
|
timeout=5,
|
|
)
|
|
if res.returncode == 0:
|
|
# Extract first column (cluster names)
|
|
return [
|
|
line.split()[0] for line in res.stdout.splitlines() if line.strip()
|
|
]
|
|
except Exception:
|
|
pass
|
|
return []
|
|
|
|
def _on_create_k3d_cluster(self):
|
|
name = self.selected_k3d_cluster.get().strip()
|
|
if not name:
|
|
messagebox.showerror("Error", "Cluster name cannot be empty")
|
|
return
|
|
# Persist the typed name so the refreshed page keeps it selected
|
|
self.k3d_cluster_name.set(name)
|
|
|
|
# If the cluster already exists just select it (triggers notebook display)
|
|
existing = self._get_k3d_cluster_list()
|
|
if name in existing:
|
|
self.selected_k3d_cluster.set(name)
|
|
self.show_page("init_cluster")
|
|
return
|
|
|
|
# Refresh the page so the notebook/console is visible before streaming output
|
|
self.selected_k3d_cluster.set(name)
|
|
self.show_page("init_cluster")
|
|
|
|
def do_create():
|
|
console = getattr(self, "_k3s_service_deploy_console", None)
|
|
if console:
|
|
console.clear()
|
|
console.write(
|
|
f"Creating k3d cluster '{name}' @ {time.strftime('%Y-%m-%d %H:%M:%S')}\n\n"
|
|
)
|
|
try:
|
|
if self._k3s_service_notebook and self._k3s_service_deploy_tab:
|
|
self._k3s_service_notebook.select(self._k3s_service_deploy_tab)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
# Build the k3d create command
|
|
prole_data = str(self._resolve_env_dir("PROLE_DATA", "data"))
|
|
volume_args = _k3d_prole_data_volume_args(prole_data)
|
|
reg_args = []
|
|
if getattr(self, "registry_url", None):
|
|
if self.registry_url.startswith("localhost:5000"):
|
|
reg_args = ["--registry-create", "prole-registry:0.0.0.0:5000"]
|
|
else:
|
|
reg_args = ["--registry-use", self.registry_url]
|
|
cmd = (
|
|
["k3d", "cluster", "create", name, "-a", "2", "--wait"]
|
|
+ volume_args
|
|
+ reg_args
|
|
+ ["--timestamps"]
|
|
)
|
|
if console:
|
|
console.write(f"$ {' '.join(cmd)}\n\n")
|
|
|
|
proc = subprocess.Popen(
|
|
cmd,
|
|
stdout=subprocess.PIPE,
|
|
stderr=subprocess.STDOUT,
|
|
text=True,
|
|
cwd=PROJECT_ROOT,
|
|
)
|
|
for line in proc.stdout:
|
|
if console:
|
|
console.write(line)
|
|
proc.wait()
|
|
if proc.returncode != 0:
|
|
if console:
|
|
console.write(
|
|
f"\nk3d cluster create failed (exit {proc.returncode})\n"
|
|
)
|
|
return
|
|
if console:
|
|
console.write(f"\nCluster '{name}' created successfully.\n\n")
|
|
|
|
# Merge kubeconfig so kubectl can reach the new cluster
|
|
try:
|
|
subprocess.run(
|
|
[
|
|
"k3d",
|
|
"kubeconfig",
|
|
"merge",
|
|
name,
|
|
"--kubeconfig-switch-context",
|
|
],
|
|
capture_output=True,
|
|
text=True,
|
|
timeout=10,
|
|
)
|
|
if console:
|
|
console.write(f"Kubeconfig merged for cluster '{name}'.\n\n")
|
|
except Exception as kce:
|
|
if console:
|
|
console.write(f"Warning: failed to merge kubeconfig: {kce}\n")
|
|
|
|
# Refresh UI and then deploy common services check
|
|
self.safe_after(self._refresh_k3d_ui)
|
|
# Trigger common services check in the deploy console
|
|
self.safe_after(self._deploy_k3s_services)
|
|
except Exception as e:
|
|
if console:
|
|
console.write(f"\nError: {e}\n")
|
|
self.safe_after(
|
|
lambda: messagebox.showerror(
|
|
"Error", f"Failed to create k3d cluster: {e}"
|
|
)
|
|
)
|
|
|
|
threading.Thread(target=do_create, daemon=True).start()
|
|
|
|
def _refresh_k3d_ui(self):
|
|
clusters = self._get_k3d_cluster_list()
|
|
self.k3d_cluster_list.set(clusters)
|
|
name = self.k3d_cluster_name.get().strip()
|
|
if name in clusters:
|
|
self.selected_k3d_cluster.set(name)
|
|
self.show_page("init_cluster")
|
|
|
|
def _on_k3d_cluster_select(self, *args):
|
|
"""Called when a k3d cluster is selected from the dropdown."""
|
|
cluster_name = self.selected_k3d_cluster.get()
|
|
if cluster_name:
|
|
try:
|
|
# Merge kubeconfig and switch context to the selected cluster
|
|
subprocess.run(
|
|
["k3d", "kubeconfig", "merge", cluster_name, "--switch-context"],
|
|
capture_output=True,
|
|
text=True,
|
|
)
|
|
# Ensure kubectx also knows about it
|
|
self._switch_kubectx(f"k3d-{cluster_name}")
|
|
except Exception:
|
|
pass
|
|
self.show_page("init_cluster")
|
|
|
|
def _on_delete_k3d_cluster(self):
|
|
"""Delete the selected k3d cluster after user confirmation."""
|
|
name = self.selected_k3d_cluster.get().strip()
|
|
if not name:
|
|
messagebox.showerror("Error", "No cluster selected to delete.")
|
|
return
|
|
|
|
confirmed = messagebox.askyesno(
|
|
"Confirm Delete", f"Are you sure you want to delete cluster '{name}'?"
|
|
)
|
|
if not confirmed:
|
|
return
|
|
|
|
# Show the notebook/console before streaming output
|
|
self.show_page("init_cluster")
|
|
|
|
def do_delete():
|
|
console = getattr(self, "_k3s_service_deploy_console", None)
|
|
if console:
|
|
console.clear()
|
|
console.write(
|
|
f"Deleting k3d cluster '{name}' @ {time.strftime('%Y-%m-%d %H:%M:%S')}\n\n"
|
|
)
|
|
try:
|
|
if self._k3s_service_notebook and self._k3s_service_deploy_tab:
|
|
self._k3s_service_notebook.select(self._k3s_service_deploy_tab)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
cmd = ["k3d", "cluster", "delete", name]
|
|
if console:
|
|
console.write(f"$ {' '.join(cmd)}\n\n")
|
|
|
|
proc = subprocess.Popen(
|
|
cmd,
|
|
stdout=subprocess.PIPE,
|
|
stderr=subprocess.STDOUT,
|
|
text=True,
|
|
cwd=PROJECT_ROOT,
|
|
)
|
|
for line in proc.stdout:
|
|
if console:
|
|
console.write(line)
|
|
proc.wait()
|
|
if proc.returncode != 0:
|
|
if console:
|
|
console.write(
|
|
f"\nk3d cluster delete failed (exit {proc.returncode})\n"
|
|
)
|
|
return
|
|
if console:
|
|
console.write(f"\nCluster '{name}' deleted successfully.\n")
|
|
|
|
# Clear selection and refresh UI
|
|
self.selected_k3d_cluster.set("")
|
|
self.k3d_cluster_name.set("")
|
|
self.safe_after(self._refresh_k3d_ui)
|
|
except Exception as e:
|
|
if console:
|
|
console.write(f"\nError: {e}\n")
|
|
self.safe_after(
|
|
lambda: messagebox.showerror(
|
|
"Error", f"Failed to delete k3d cluster: {e}"
|
|
)
|
|
)
|
|
|
|
threading.Thread(target=do_delete, daemon=True).start()
|
|
|
|
def _on_kubectx_select(self, *args):
|
|
"""Called when a generic kubernetes context is selected from the dropdown."""
|
|
ctx = self.selected_kubectx.get()
|
|
if ctx:
|
|
self._switch_kubectx(ctx)
|
|
self._verify_k3s_services()
|
|
self.show_page("init_cluster")
|
|
|
|
def _switch_kubectx(self, context_name: str):
|
|
"""Switch the current kubernetes context using kubectx or kubectl."""
|
|
if not context_name:
|
|
return
|
|
|
|
# Determine which kubeconfig to modify.
|
|
# Use the same logic as the rest of the application.
|
|
kc = _find_kubeconfig_file()
|
|
env = os.environ.copy()
|
|
if kc:
|
|
env["KUBECONFIG"] = kc
|
|
|
|
try:
|
|
# Try kubectx first
|
|
res = subprocess.run(
|
|
["kubectx", context_name], capture_output=True, text=True, env=env
|
|
)
|
|
if res.returncode != 0:
|
|
# Fallback to kubectl config use-context
|
|
subprocess.run(
|
|
["kubectl", "config", "use-context", context_name],
|
|
capture_output=True,
|
|
text=True,
|
|
env=env
|
|
)
|
|
|
|
# Update the UI dropdown to match (but don't trigger re-select)
|
|
if hasattr(self, "selected_kubectx") and self.selected_kubectx.get() != context_name:
|
|
# Check if it's in the current values list
|
|
current_values = list(self.kubectx_list.get() or [])
|
|
if context_name in current_values:
|
|
self.selected_kubectx.set(context_name)
|
|
except Exception:
|
|
pass
|
|
|
|
def _get_kubectx_list(self) -> list[str]:
|
|
"""Get list of kubernetes contexts."""
|
|
kc = _find_kubeconfig_file()
|
|
env = os.environ.copy()
|
|
if kc:
|
|
env["KUBECONFIG"] = kc
|
|
|
|
try:
|
|
res = subprocess.run(["kubectx"], capture_output=True, text=True, env=env)
|
|
if res.returncode == 0:
|
|
return res.stdout.strip().split("\n")
|
|
|
|
# Fallback to kubectl
|
|
res = subprocess.run(
|
|
["kubectl", "config", "get-contexts", "-o", "name"],
|
|
capture_output=True,
|
|
text=True,
|
|
env=env
|
|
)
|
|
if res.returncode == 0:
|
|
return res.stdout.strip().split("\n")
|
|
except Exception:
|
|
pass
|
|
return ["default"]
|
|
|
|
def _validate_and_save_cluster_config(self) -> bool:
|
|
"""Verify the config then write the values to prole.cfg."""
|
|
env_key = self._cluster_env_key()
|
|
if env_key == "dev":
|
|
cluster_val = self.selected_k3d_cluster.get() or "knoe-dev-cluster"
|
|
if not cluster_val:
|
|
messagebox.showwarning(
|
|
"Validation", "Please select or create a k3d cluster."
|
|
)
|
|
return False
|
|
|
|
# Context switch (if context matches cluster name)
|
|
try:
|
|
subprocess.run(
|
|
["kubectl", "config", "use-context", f"k3d-{cluster_val}"],
|
|
capture_output=True,
|
|
text=True,
|
|
)
|
|
except Exception:
|
|
pass
|
|
elif env_key == "service":
|
|
if not (self.k3s_server_url.get() or "").strip():
|
|
messagebox.showwarning(
|
|
"Validation", "Please specify the K3s server URL."
|
|
)
|
|
return False
|
|
if not (self.k3s_token.get() or "").strip():
|
|
messagebox.showwarning("Validation", "Please specify the K3s token.")
|
|
return False
|
|
cluster_val = self.cluster_env.get()
|
|
elif env_key == "prod":
|
|
path = self.prod_artifacts_path.get().strip()
|
|
if not path:
|
|
messagebox.showwarning(
|
|
"Validation", "Please specify an artifact staging directory."
|
|
)
|
|
return False
|
|
p = Path(path).expanduser()
|
|
if not p.exists():
|
|
try:
|
|
p.mkdir(parents=True, exist_ok=True)
|
|
except Exception as e:
|
|
messagebox.showerror(
|
|
"Error", f"Failed to create directory {path}: {e}"
|
|
)
|
|
return False
|
|
cluster_val = self.cluster_env.get()
|
|
else:
|
|
cluster_val = self.cluster_env.get()
|
|
|
|
# Capture cluster info for config
|
|
self.prole_cfg_data["Initialize Cluster"]["ENVIRONMENT"] = cluster_val
|
|
self.prole_cfg_data["Initialize Cluster"]["K3S_SERVER_URL"] = (
|
|
self.k3s_server_url.get() or ""
|
|
).strip()
|
|
self.prole_cfg_data["Initialize Cluster"]["K3S_TOKEN"] = _encrypt_cfg_secret(
|
|
self.k3s_token.get() or ""
|
|
)
|
|
self.prole_cfg_data["Global"][
|
|
"SERVICE_NAMESPACE"
|
|
] = self._get_service_namespace()
|
|
self.prole_cfg_data["Optional Features"]["SUPABASE_ENABLED"] = str(
|
|
self.supabase_enabled.get()
|
|
)
|
|
self.prole_cfg_data["Optional Features"]["GITOPS_ENABLED"] = str(
|
|
self.gitops_enabled.get()
|
|
)
|
|
self.prole_cfg_data["Optional Features"]["KERBEROS_ENABLED"] = str(
|
|
self.kerberos_enabled.get()
|
|
)
|
|
self.prole_cfg_data["Optional Features"]["AT_REST_ENCRYPTION_ENABLED"] = str(
|
|
self.at_rest_encryption_enabled.get()
|
|
)
|
|
|
|
if env_key == "prod":
|
|
self.prole_cfg_data["Prod Cluster (k8s)"][
|
|
"STAGING_DIRECTORY"
|
|
] = self.prod_artifacts_path.get()
|
|
|
|
# Save config
|
|
self._save_prole_cfg()
|
|
return True
|
|
|
|
def _cluster_status_snapshot(self) -> dict:
|
|
env = self._cluster_env_key()
|
|
|
|
cluster_ok = False
|
|
if env == "dev":
|
|
cluster_name = self.selected_k3d_cluster.get() or "knoe-dev-cluster"
|
|
try:
|
|
res = subprocess.run(
|
|
["k3d", "cluster", "list", "--no-headers"],
|
|
capture_output=True,
|
|
text=True,
|
|
)
|
|
cluster_ok = cluster_name in (res.stdout or "")
|
|
except Exception:
|
|
cluster_ok = False
|
|
cluster_msg = (
|
|
f"K3D Cluster ({cluster_name}): Running"
|
|
if cluster_ok
|
|
else f"K3D Cluster ({cluster_name}): Not found"
|
|
)
|
|
else:
|
|
cluster_ok, cluster_msg = self._check_k8s_cluster(env)
|
|
cluster_fill = "#34c759" if cluster_ok else "#ff9f0a"
|
|
|
|
common_ok = cluster_ok
|
|
common_msg = (
|
|
"Common Services: Ready to deploy"
|
|
if common_ok
|
|
else "Common Services: Not ready"
|
|
)
|
|
common_fill = "#34c759" if common_ok else "#ff9f0a"
|
|
|
|
return {
|
|
"env": env,
|
|
"cluster_ok": cluster_ok,
|
|
"cluster_msg": cluster_msg,
|
|
"cluster_fill": cluster_fill,
|
|
"common_ok": common_ok,
|
|
"common_msg": common_msg,
|
|
"common_fill": common_fill,
|
|
}
|
|
|
|
def _cluster_ready_for_navigation(self) -> bool:
|
|
status = self._cluster_status_snapshot()
|
|
if not status["cluster_ok"]:
|
|
try:
|
|
messagebox.showerror(
|
|
"Cluster",
|
|
"Cluster is not reachable. Please verify your cluster and try again.",
|
|
)
|
|
except Exception:
|
|
pass
|
|
return False
|
|
return True
|
|
|
|
def check_cluster_status_async(self):
|
|
def worker():
|
|
notice_active = False
|
|
self._set_cluster_env_message("Collect Status", "#34c759")
|
|
status = self._cluster_status_snapshot()
|
|
self._set_cluster_env_message("Verify State", "#34c759")
|
|
kong_status = None
|
|
if status.get("cluster_ok"):
|
|
kong_status = self._dashboard_kong_status(status.get("env"))
|
|
self._set_cluster_env_message("Assess Actions", "#34c759")
|
|
|
|
# Repair step: reset Dashboard (Kong) pods if unhealthy, with cooldown
|
|
if kong_status and kong_status.get("reset_pods"):
|
|
now = time.time()
|
|
if now - self._kong_last_repair_at >= self._kong_repair_cooldown_s:
|
|
self._set_cluster_env_message(
|
|
"Execute Actions (REPAIR): reset Dashboard (Kong) pod(s)",
|
|
"#34c759",
|
|
)
|
|
pods = kong_status.get("reset_pods") or []
|
|
cmd = (
|
|
self._kubectl_base_cmd(status.get("env"))
|
|
+ ["delete", "pod", "-n", "kubernetes-dashboard"]
|
|
+ pods
|
|
)
|
|
rc, _out = self._run_cmd_capture(cmd, timeout=10)
|
|
if rc != 0:
|
|
self._set_cluster_env_message(
|
|
"Notice: Dashboard (Kong) repair failed",
|
|
"#ff3b30",
|
|
clear_after_ms=6000,
|
|
)
|
|
notice_active = True
|
|
else:
|
|
self._kong_last_repair_at = now
|
|
# Give the cluster a moment to recreate pods before re-check
|
|
time.sleep(2)
|
|
self._set_cluster_env_message("Collect Status", "#34c759")
|
|
status = self._cluster_status_snapshot()
|
|
else:
|
|
self._set_cluster_env_message(
|
|
"Execute Actions (REPAIR): skipped (recently attempted)",
|
|
"#34c759",
|
|
clear_after_ms=1500,
|
|
)
|
|
else:
|
|
self._set_cluster_env_message(
|
|
"Execute Actions (REPAIR): none needed",
|
|
"#34c759",
|
|
clear_after_ms=1500,
|
|
)
|
|
|
|
# Run the repair pipeline if anomalies are detected on a ready cluster
|
|
try:
|
|
self._maybe_run_repair_pipeline(status, kong_status)
|
|
except Exception:
|
|
pass
|
|
|
|
def update_ui():
|
|
if hasattr(self, "k3d_status_label"):
|
|
self.bg_canvas.itemconfig(
|
|
self.k3d_status_label,
|
|
text=status["cluster_msg"],
|
|
fill=status["cluster_fill"],
|
|
)
|
|
if hasattr(self, "common_services_status_label"):
|
|
self.bg_canvas.itemconfig(
|
|
self.common_services_status_label,
|
|
text=status["common_msg"],
|
|
fill=status["common_fill"],
|
|
)
|
|
if hasattr(self, "_init_cluster_button"):
|
|
try:
|
|
self._init_cluster_button.configure(text="Save")
|
|
except Exception:
|
|
pass
|
|
|
|
self.root.after(0, update_ui)
|
|
# Clear message after a short delay if no notice is active
|
|
if not notice_active:
|
|
self._set_cluster_env_message("", "#34c759", clear_after_ms=2000)
|
|
|
|
threading.Thread(target=worker, daemon=True).start()
|
|
|
|
def _services_status_snapshot(self) -> dict:
|
|
opentofu_ok = self._check_opentofu_health()
|
|
|
|
opentofu_msg = "OpenTofu: Running" if opentofu_ok else "OpenTofu: Not reachable"
|
|
|
|
opentofu_fill = "#34c759" if opentofu_ok else "#ff9f0a"
|
|
|
|
return {
|
|
"opentofu_msg": opentofu_msg,
|
|
"opentofu_fill": opentofu_fill,
|
|
}
|
|
|
|
def check_services_status_async(self):
|
|
def worker():
|
|
status = self._services_status_snapshot()
|
|
|
|
def update_ui():
|
|
if hasattr(self, "_db_opentofu_status_label"):
|
|
self.bg_canvas.itemconfig(
|
|
self._db_opentofu_status_label,
|
|
text=status["opentofu_msg"],
|
|
fill=status["opentofu_fill"],
|
|
)
|
|
|
|
self.root.after(0, update_ui)
|
|
|
|
threading.Thread(target=worker, daemon=True).start()
|
|
|
|
def ensure_cluster_ready(self):
|
|
self._action_flags["init_cluster.start_cluster"] = True
|
|
|
|
# Implementation of cluster creation/startup
|
|
def worker():
|
|
cluster_env = self._cluster_env_key()
|
|
|
|
if cluster_env == "dev":
|
|
# 1. Start Docker if not running
|
|
if not self.controller.check_docker_running():
|
|
# Attempt to start Docker on macOS
|
|
subprocess.run(["open", "-a", "Docker"], capture_output=True)
|
|
# Wait for it to start
|
|
for _ in range(30):
|
|
time.sleep(2)
|
|
if self.controller.check_docker_running():
|
|
break
|
|
|
|
if not self.controller.check_docker_running():
|
|
self.root.after(
|
|
0,
|
|
lambda: messagebox.showerror(
|
|
"Docker",
|
|
"Could not start Docker. Please start it manually.",
|
|
),
|
|
)
|
|
return
|
|
|
|
# 2. Manage or verify cluster
|
|
cluster_name = "knoe-dev-cluster"
|
|
res = subprocess.run(
|
|
["k3d", "cluster", "list", "--no-headers"],
|
|
capture_output=True,
|
|
text=True,
|
|
)
|
|
res_stdout = res.stdout or ""
|
|
if self._reset_cluster:
|
|
subprocess.run(
|
|
["k3d", "cluster", "delete", cluster_name], capture_output=True
|
|
)
|
|
res_stdout = ""
|
|
self._reset_cluster = False
|
|
if cluster_name not in res_stdout:
|
|
# Create it
|
|
# Default args based on README.md
|
|
prole_data = str(self._resolve_env_dir("PROLE_DATA", "data"))
|
|
volume_args = _k3d_prole_data_volume_args(prole_data)
|
|
cmd = [
|
|
"k3d",
|
|
"cluster",
|
|
"create",
|
|
cluster_name,
|
|
"-a",
|
|
"2",
|
|
] + volume_args
|
|
reg_args = []
|
|
try:
|
|
if self.ensure_local_registry_available():
|
|
reg_args = ["--registry-use", "k3d-prole-registry:5000"]
|
|
except Exception:
|
|
reg_args = []
|
|
cmd += reg_args + ["--api-port", "0.0.0.0:6443"]
|
|
|
|
# Run in terminal or capture output? Let's use a console window later.
|
|
# For now, run it and update status.
|
|
subprocess.run(cmd, capture_output=True)
|
|
else:
|
|
# Start it if it's stopped
|
|
subprocess.run(
|
|
["k3d", "cluster", "start", cluster_name], capture_output=True
|
|
)
|
|
|
|
self.check_cluster_status_async()
|
|
self.root.after(
|
|
0,
|
|
lambda: messagebox.showinfo(
|
|
"Cluster", f"Cluster {cluster_name} is ready."
|
|
),
|
|
)
|
|
else:
|
|
if self._reset_cluster and cluster_env in ("service", "k3s"):
|
|
ns = (self._get_service_namespace() or "").strip() or "default"
|
|
server, token = self._k3s_connection_info()
|
|
_reset_k3s_namespace(
|
|
self.controller.project_root, ns, server, token
|
|
)
|
|
self._reset_cluster = False
|
|
ok, msg = self._check_k8s_cluster(cluster_env)
|
|
if not ok:
|
|
self.root.after(0, lambda: messagebox.showerror("Cluster", msg))
|
|
return
|
|
self.check_cluster_status_async()
|
|
self.root.after(0, lambda: messagebox.showinfo("Cluster", msg))
|
|
|
|
threading.Thread(target=worker, daemon=True).start()
|
|
|
|
def _verify_k3s_services(self):
|
|
if not self._k3s_service_status_console:
|
|
self._stop_k3s_service_status_updates()
|
|
return
|
|
try:
|
|
self._k3s_service_status_console.clear()
|
|
except Exception:
|
|
pass
|
|
try:
|
|
if self._k3s_service_notebook and self._k3s_service_status_tab:
|
|
self._k3s_service_notebook.select(self._k3s_service_status_tab)
|
|
except Exception:
|
|
pass
|
|
self._start_k3s_service_status_updates()
|
|
|
|
def _start_k3s_service_status_updates(self):
|
|
if not self._k3s_status_active:
|
|
self._k3s_status_active = True
|
|
if self._k3s_status_after_id:
|
|
try:
|
|
self.root.after_cancel(self._k3s_status_after_id)
|
|
except Exception:
|
|
pass
|
|
self._k3s_status_after_id = None
|
|
self._schedule_k3s_service_status_update(0)
|
|
|
|
def _stop_k3s_service_status_updates(self):
|
|
self._k3s_status_active = False
|
|
self._k3s_status_inflight = False
|
|
if self._k3s_status_after_id:
|
|
try:
|
|
self.root.after_cancel(self._k3s_status_after_id)
|
|
except Exception:
|
|
pass
|
|
self._k3s_status_after_id = None
|
|
|
|
def _schedule_k3s_service_status_update(self, delay_ms: int):
|
|
if not self._k3s_status_active:
|
|
return
|
|
if not self.root or not self.root.winfo_exists():
|
|
return
|
|
try:
|
|
self._k3s_status_after_id = self.root.after(
|
|
delay_ms, self._run_k3s_service_status_once
|
|
)
|
|
except Exception:
|
|
self._k3s_status_after_id = None
|
|
|
|
def _run_k3s_service_status_once(self):
|
|
if self._k3s_status_inflight:
|
|
return
|
|
if not self._k3s_service_status_console:
|
|
self._k3s_status_inflight = False
|
|
self.safe_after(lambda: self._schedule_k3s_service_status_update(30000))
|
|
return
|
|
self._k3s_status_inflight = True
|
|
console = self._k3s_service_status_console
|
|
console.clear()
|
|
console.write(
|
|
f"Service status refresh @ {time.strftime('%Y-%m-%d %H:%M:%S')}\n\n"
|
|
)
|
|
|
|
def worker():
|
|
collected_lines = []
|
|
try:
|
|
ns = self._get_service_namespace()
|
|
env = self._script_env_for_namespace(ns)
|
|
env["PROLE_MODE"] = self._deployment_mode()
|
|
|
|
if not env.get("KUBECONFIG"):
|
|
console.write(
|
|
"KUBECONFIG not generated. Check cluster credentials in prole.cfg.\n"
|
|
)
|
|
return
|
|
|
|
verify_args = ["-n", ns, "verify"]
|
|
|
|
def _on_line(line):
|
|
console.write(line)
|
|
collected_lines.append(line)
|
|
|
|
rc = self.controller.run_script(
|
|
"init_common_services.sh",
|
|
args=verify_args,
|
|
env=env,
|
|
on_line=_on_line,
|
|
)
|
|
|
|
console.write(f"\nVerification Exit status: {rc}\n")
|
|
|
|
# Parse per-component traffic light status from [OK]/[FAIL] lines
|
|
component_status = self._parse_service_status_lines(collected_lines)
|
|
all_green = rc == 0
|
|
|
|
def _update_traffic_lights():
|
|
self._all_services_green = all_green
|
|
# Ensure common services success flag is synced for navigation
|
|
self._common_services_success = all_green
|
|
|
|
lights = getattr(self, "_service_traffic_lights", {})
|
|
for comp_key, indicator_id in lights.items():
|
|
if comp_key in component_status:
|
|
color = (
|
|
"#34c759" if component_status[comp_key] else "#ff3b30"
|
|
)
|
|
else:
|
|
color = "#34c759" if all_green else "#8e8e93"
|
|
try:
|
|
self.bg_canvas.itemconfig(
|
|
indicator_id, fill=color, outline=color
|
|
)
|
|
except Exception:
|
|
pass
|
|
|
|
# Update footer to enable Next button if everything is green
|
|
self.update_footer()
|
|
|
|
self.safe_after(_update_traffic_lights)
|
|
finally:
|
|
self._k3s_status_inflight = False
|
|
self.safe_after(lambda: self._schedule_k3s_service_status_update(30000))
|
|
|
|
threading.Thread(target=worker, daemon=True).start()
|
|
|
|
def _parse_service_status_lines(self, lines: list[str]) -> dict[str, bool]:
|
|
"""Parse [OK]/[FAIL] lines from status_common_services.sh output.
|
|
|
|
Returns a dict mapping component keys (argo, certmgr,
|
|
garage, kong, openbao, opentofu) to True (healthy) or False (unhealthy).
|
|
"""
|
|
# Map resource names to traffic light component keys
|
|
name_map = {
|
|
"argocd-server": "argo",
|
|
"argocd-repo-server": "argo",
|
|
"argocd-dex-server": "argo",
|
|
"argocd-applicationset-controller": "argo",
|
|
"argocd-notifications-controller": "argo",
|
|
"argocd-redis": "argo",
|
|
"argocd-application-controller": "argo",
|
|
"opentofu": "opentofu",
|
|
"garage": "garage",
|
|
"openbao": "openbao",
|
|
"prole-svc-kong": "kong",
|
|
"cert-manager": "certmgr",
|
|
"cert-manager-cainjector": "certmgr",
|
|
"cert-manager-webhook": "certmgr",
|
|
}
|
|
# Start with all components healthy (True); any FAIL sets to False
|
|
result: dict[str, bool] = {}
|
|
import re as _re
|
|
|
|
for line in lines:
|
|
m = _re.match(r"\[(OK|FAIL)\]\s+(\S+)/(\S+)", line)
|
|
if not m:
|
|
continue
|
|
status_tag = m.group(1)
|
|
resource_name = m.group(3)
|
|
comp_key = name_map.get(resource_name)
|
|
if not comp_key:
|
|
continue
|
|
if comp_key not in result:
|
|
result[comp_key] = True
|
|
if status_tag == "FAIL":
|
|
result[comp_key] = False
|
|
return result
|
|
|
|
def _deploy_k3s_services(self):
|
|
"""Deploy common services (ArgoCD/OpenTofu/Garage) to remote k3s cluster."""
|
|
|
|
def _deploy():
|
|
console = self._k3s_service_deploy_console
|
|
if console:
|
|
console.clear()
|
|
console.write(
|
|
f"Deploying common services @ {time.strftime('%Y-%m-%d %H:%M:%S')}\n\n"
|
|
)
|
|
try:
|
|
if self._k3s_service_notebook and self._k3s_service_deploy_tab:
|
|
self._k3s_service_notebook.select(self._k3s_service_deploy_tab)
|
|
except Exception:
|
|
pass
|
|
|
|
service_ns = self._get_service_namespace()
|
|
env = self._script_env_for_namespace(service_ns)
|
|
env["PROLE_MODE"] = self._deployment_mode()
|
|
log_path = self._common_services_log_path()
|
|
env["COMMON_SERVICES_INIT_LOG"] = str(log_path)
|
|
try:
|
|
log_fp = log_path.open("a", encoding="utf-8")
|
|
log_fp.write(
|
|
f"\n# Common services deploy @ {time.strftime('%Y-%m-%d %H:%M:%S')}\n"
|
|
)
|
|
log_fp.flush()
|
|
except Exception:
|
|
log_fp = None
|
|
if console:
|
|
console.write(f"Log file: {log_path}\n\n")
|
|
|
|
if not env.get("KUBECONFIG"):
|
|
if console:
|
|
console.write(
|
|
"KUBECONFIG not generated. Check cluster credentials in prole.cfg.\n"
|
|
)
|
|
if log_fp:
|
|
try:
|
|
log_fp.write(
|
|
"KUBECONFIG not generated. Check cluster credentials in prole.cfg.\n"
|
|
)
|
|
log_fp.close()
|
|
except Exception:
|
|
pass
|
|
self._verify_k3s_services()
|
|
return
|
|
|
|
try:
|
|
common_args = ["-n", env["NAMESPACE"]]
|
|
if self.kerberos_enabled.get():
|
|
common_args.append("-k")
|
|
common_args.append("update")
|
|
|
|
def _log_line(line: str):
|
|
if console:
|
|
console.write(line)
|
|
if log_fp:
|
|
try:
|
|
log_fp.write(line)
|
|
log_fp.flush()
|
|
except Exception:
|
|
pass
|
|
|
|
rc = self.controller.run_script(
|
|
"init_common_services.sh",
|
|
args=common_args,
|
|
env=env,
|
|
on_line=_log_line,
|
|
)
|
|
if console:
|
|
console.write(f"\nExit status: {rc}\n")
|
|
if log_fp:
|
|
try:
|
|
log_fp.write(f"\nExit status: {rc}\n")
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
if log_fp:
|
|
try:
|
|
log_fp.close()
|
|
except Exception:
|
|
pass
|
|
|
|
self._verify_k3s_services()
|
|
|
|
threading.Thread(target=_deploy, daemon=True).start()
|
|
|
|
def _ensure_k3s_kubeconfig_merged(self):
|
|
"""Ensure a valid kubeconfig is available and merged into ~/.kube/config.
|
|
|
|
Prefers an existing Ansible-fetched kubeconfig (with client-certificate
|
|
auth) over generating a new token-based one.
|
|
"""
|
|
# 1. Check for existing cert-based kubeconfig first
|
|
existing = _find_kubeconfig_file()
|
|
if existing:
|
|
self._managed_kubeconfig = existing
|
|
_merge_kubeconfig(existing)
|
|
return
|
|
# 2. Fall back to token-based generation
|
|
server, token = self._k3s_connection_info()
|
|
if not server or not token:
|
|
return
|
|
kubeconfig_path = str(_write_k3s_kubeconfig(server, token))
|
|
self._managed_kubeconfig = kubeconfig_path
|
|
_merge_kubeconfig(kubeconfig_path)
|
|
|
|
def _k3s_connection_info(self) -> tuple[str, str]:
|
|
# Check UI fields first; fall back to the shared 3-source resolver.
|
|
ui_server = (_safe_str(self.k3s_server_url.get()) or "").strip()
|
|
ui_token = (_safe_str(self.k3s_token.get()) or "").strip()
|
|
if ui_token:
|
|
ui_token = self._resolve_secret_value(ui_token)
|
|
ui_token = _normalize_k3s_token(ui_token)
|
|
if ui_server and not ui_server.startswith("http"):
|
|
ui_server = f"https://{ui_server}"
|
|
# If UI fields are populated, use them directly.
|
|
if ui_server and ui_token:
|
|
return ui_server, ui_token
|
|
# Otherwise delegate to the unified resolver (env → cfg → ansible).
|
|
resolved_server, resolved_token = _resolve_k3s_connection_fn(
|
|
project_root=PROJECT_ROOT,
|
|
)
|
|
return ui_server or resolved_server, ui_token or resolved_token
|
|
|
|
def _k3s_connection_raw(self) -> tuple[str, str]:
|
|
"""Return (server_url, token) without attempting secret resolution."""
|
|
server = (
|
|
_safe_str(self.k3s_server_url.get())
|
|
or os.environ.get("PROLE_K3S_SERVER")
|
|
or os.environ.get("K3S_SERVER_URL")
|
|
or ""
|
|
).strip()
|
|
token = (
|
|
_safe_str(self.k3s_token.get())
|
|
or os.environ.get("PROLE_K3S_TOKEN")
|
|
or os.environ.get("K3S_TOKEN")
|
|
or ""
|
|
).strip()
|
|
if server and not server.startswith("http"):
|
|
server = f"https://{server}"
|
|
return server, token
|