prole/installer/ui/screens/cluster.py

1406 lines
60 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,
_looks_like_k8s_bearer_token,
_normalize_cluster_env,
_normalize_k3s_token,
)
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')
self._render_title('Cluster Environment', y=150)
self._render_paragraph('Select a cluster environment and ensure the cluster is running and common services can be deployed.', y=200)
# Cluster Selection (Radio Buttons)
x_label = 48
y = 280
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 += 44
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
if selected_env_key == 'service':
# K3s connection settings
try:
self._apply_k3s_defaults()
except Exception:
pass
self._canvas_items.append(ui.canvas_text(self, x_label, y, 'K3s Connection:', fill='black', font=('SF Pro Text', 12, 'bold')))
y += 28
ui.canvas_text(self, x_label + 20, y, 'Server URL:', fill='black', font=('SF Pro Text', 11))
server_entry = tk.Entry(self.bg_canvas, textvariable=self.k3s_server_url, width=60)
server_win = self.bg_canvas.create_window(x_label + 140, y - 8, window=server_entry, anchor='nw')
self._canvas_items.append(server_win)
self._overlay_widgets.append(server_entry)
y += 28
ui.canvas_text(self, x_label + 20, y, 'Token:', fill='black', font=('SF Pro Text', 11))
token_entry = tk.Entry(self.bg_canvas, textvariable=self.k3s_token, show='*', width=60)
token_win = self.bg_canvas.create_window(x_label + 200, y - 8, window=token_entry, anchor='nw')
self._canvas_items.append(token_win)
self._overlay_widgets.append(token_entry)
y += 34
# When server + token are available, generate and merge kubeconfig
# so the prole-k3s context appears in the dropdown.
if selected_env_key == 'service':
try:
self._ensure_k3s_kubeconfig_merged()
except Exception:
pass
# Context selection for both Service and Prod
self._canvas_items.append(ui.canvas_text(self, x_label, y, 'Kubernetes Context:', fill='black', font=('SF Pro Text', 12, 'bold')))
y += 28
values = self._get_kubectx_list()
combo = ttk.Combobox(self.bg_canvas, textvariable=self.selected_kubectx, values=values, state='readonly', width=40)
if not self.selected_kubectx.get() and values:
if selected_env_key == 'service' and 'prole-k3s' in values:
self.selected_kubectx.set('prole-k3s')
else:
self.selected_kubectx.set(values[0])
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)
y += 40
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')))
y += 30
entry = tk.Entry(self.bg_canvas, textvariable=self.prod_artifacts_path, width=60)
entry_win = self.bg_canvas.create_window(x_label + 20, y, 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 + 440, y - 5, window=browse_btn, anchor='nw')
self._canvas_items.append(browse_win)
self._overlay_widgets.append(browse_btn)
y += 30
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:
service_label_text = f"Common Services (namespace: {self._get_service_namespace()})"
service_label = ui.canvas_text(self, x_label, y, service_label_text, fill='black', font=('SF Pro Text', 12, 'bold'))
self._canvas_items.append(service_label)
y += 28
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='Deploy 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 = 210
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._overlay_widgets.extend([notebook, status_frame, deploy_frame, status_console, deploy_console])
restored = self._restore_common_services_log(deploy_console)
if restored:
try:
self._k3s_service_notebook.select(self._k3s_service_deploy_tab)
except Exception:
pass
y += console_height + 12
deploy_btn = tk.Button(
self.bg_canvas,
text='Deploy Missing Services',
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,
)
deploy_win = self.bg_canvas.create_window(x_label + 20, y, window=deploy_btn, anchor='nw')
self._canvas_items.append(deploy_win)
self._overlay_widgets.append(deploy_btn)
y += 40
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()
return ns or 'default'
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
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 _openbao_url(self) -> str:
url = (os.environ.get("PROLE_OPENBAO_URL") or "http://127.0.0.1:18200").strip()
return url.rstrip('/')
def _check_openbao_health(self) -> bool:
"""Check OpenBao health in the currently selected namespace/pod.
Preference order:
1) If in dev k3d or service k3s cluster, check k8s resources in selected namespace (deployment/statefulset ready)
2) Fallback to HTTP health endpoint if PROLE_OPENBAO_URL (or default) responds
"""
try:
env_key = self._cluster_env_key()
except Exception:
env_key = 'dev'
# 1) Check k8s readiness in selected namespace when using dev/service clusters
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)
# deployment/openbao readiness
res_dep = subprocess.run(base_cmd + ['-n', ns, 'get', 'deploy', 'openbao', '-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'))
# statefulset/openbao readiness
res_sts = subprocess.run(base_cmd + ['-n', ns, 'get', 'statefulset', 'openbao', '-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'))
if ready_dep or ready_sts:
return True
except Exception:
pass
# 2) Fallback to HTTP health
url = self._openbao_url()
if not url:
return False
health_url = f"{url}/v1/sys/health"
try:
with urllib.request.urlopen(health_url, timeout=2):
return True
except urllib.error.HTTPError:
# OpenBao responds with non-200 for sealed/standby; still reachable
return True
except Exception:
return 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 or not token:
return False, f"{label}: missing server URL/token 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)
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 _get_kubectx_list(self) -> list[str]:
"""Get list of kubernetes contexts."""
try:
res = subprocess.run(["kubectx"], capture_output=True, text=True)
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)
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 'prole-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']['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 "prole-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:
openbao_ok = self._check_openbao_health()
opentofu_ok = self._check_opentofu_health()
openbao_msg = 'OpenBao: Running' if openbao_ok else 'OpenBao: Not reachable'
opentofu_msg = 'OpenTofu: Running' if opentofu_ok else 'OpenTofu: Not reachable'
openbao_fill = '#34c759' if openbao_ok else '#ff9f0a'
opentofu_fill = '#34c759' if opentofu_ok else '#ff9f0a'
return {
'openbao_msg': openbao_msg,
'openbao_fill': openbao_fill,
'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_openbao_status_label'):
self.bg_canvas.itemconfig(self._db_openbao_status_label, text=status['openbao_msg'], fill=status['openbao_fill'])
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 = 'prole-dev-cluster'
res = subprocess.run(['k3d', 'cluster', 'list', '--no-headers'], capture_output=True, text=True)
res_stdout = res.stdout or ''
if self._reset_cluster:
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
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():
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
status_args = ["-n", ns]
if self.kerberos_enabled.get():
status_args.append("-k")
rc = self.controller.run_script(
"status_common_services.sh",
args=status_args,
env=env,
on_line=lambda l: console.write(l)
)
console.write(f"\nExit status: {rc}\n")
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 _deploy_k3s_services(self):
"""Deploy common services (ArgoCD/OpenBao/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):
"""Generate a kubeconfig from current K3s credentials and merge it into ~/.kube/config."""
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]:
server = (self.k3s_server_url.get() or os.environ.get('PROLE_K3S_SERVER') or os.environ.get('K3S_SERVER_URL') or '').strip()
token = (self.k3s_token.get() or os.environ.get('PROLE_K3S_TOKEN') or os.environ.get('K3S_TOKEN') or '').strip()
if token:
token = self._resolve_secret_value(token)
token = _normalize_k3s_token(token)
if server and not server.startswith('http'):
server = f"https://{server}"
return server, token
def _k3s_connection_raw(self) -> tuple[str, str]:
"""Return (server_url, token) without attempting secret resolution."""
server = (self.k3s_server_url.get() or os.environ.get('PROLE_K3S_SERVER') or os.environ.get('K3S_SERVER_URL') or '').strip()
token = (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
def _script_env_for_namespace(self, namespace: str) -> dict:
env = os.environ.copy()
env["PROLE_HOME"] = str(PROJECT_ROOT)
env["PROLE_SERVICE"] = str(PROJECT_ROOT)
env["PROLE_DB_USER"] = self.db_username.get().strip()
if (self.db_password.get() or '').strip():
env["DB_PASSWORD"] = self.db_password.get().strip()
env["OPENTOFU_ADMIN_PASSWORD"] = self.db_password.get().strip()
env["NAMESPACE"] = namespace
env["SERVICE_NAMESPACE"] = self._get_service_namespace()
mode = self._deployment_mode()
if mode:
env["PROLE_MODE"] = mode
env["DEPLOYMENT_MODE"] = mode
env["DEPLOYMENT_TARGET"] = _deployment_target_label(self.cluster_env.get())
if self.kerberos_realm.get().strip():
env["KRB5_REALM"] = self.kerberos_realm.get().strip()
env["REALM"] = self.kerberos_realm.get().strip()
env["DOMAIN"] = self.kerberos_realm.get().strip().lower()
if self.kerberos_kdc.get().strip():
env["KRB5_KDC"] = self.kerberos_kdc.get().strip()
env["KRB5_ADMIN"] = self.kerberos_kdc.get().strip()
if self.kerberos_user.get().strip():
env["KRB5_USER"] = self.kerberos_user.get().strip()
if self.kerberos_password.get().strip():
env["KRB5_PASSWORD"] = self.kerberos_password.get().strip()
if self._cluster_env_key() == 'dev':
# For k3d dev clusters, resolve KUBECONFIG idempotently.
# Even when the cluster is already running we must ensure
# KUBECONFIG points at a valid file so downstream scripts
# (status_common_services.sh, init_common_services.sh, …) work.
if not (env.get("KUBECONFIG") or "").strip():
cluster_name = (self.selected_k3d_cluster.get() or '').strip()
if not cluster_name:
cluster_name = 'prole-dev-cluster'
# 1. Merge k3d kubeconfig into ~/.kube/config and switch context
try:
subprocess.run(
['k3d', 'kubeconfig', 'merge', cluster_name, '--kubeconfig-switch-context'],
capture_output=True, text=True, timeout=10
)
except Exception:
pass
default_kube = str(Path.home() / '.kube' / 'config')
if Path(default_kube).exists():
env["KUBECONFIG"] = default_kube
else:
# 2. Fallback: write a standalone kubeconfig via k3d
try:
import tempfile
tmp = tempfile.NamedTemporaryFile(
prefix='k3d-kubeconfig-', suffix='.yaml', delete=False
)
tmp.close()
res = subprocess.run(
['k3d', 'kubeconfig', 'get', cluster_name],
capture_output=True, text=True, timeout=10
)
if res.returncode == 0 and res.stdout.strip():
Path(tmp.name).write_text(res.stdout)
env["KUBECONFIG"] = tmp.name
else:
os.unlink(tmp.name)
except Exception:
pass
elif self._cluster_env_key() == 'service':
raw_server, raw_token = self._k3s_connection_raw()
if raw_server:
env["PROLE_K3S_SERVER"] = raw_server
if raw_token:
env["PROLE_K3S_TOKEN"] = raw_token
k3s_server, k3s_token = self._k3s_connection_info()
explicit_kubeconfig = (env.get("KUBECONFIG") or "").strip()
if explicit_kubeconfig and not os.path.exists(explicit_kubeconfig):
explicit_kubeconfig = ""
managed_set = False
if k3s_server and k3s_token and not explicit_kubeconfig:
# Always generate a fresh kubeconfig to ensure no environment leakage
try:
if self._managed_kubeconfig and os.path.exists(self._managed_kubeconfig):
try:
os.unlink(self._managed_kubeconfig)
except Exception:
pass
self._managed_kubeconfig = str(_write_k3s_kubeconfig(k3s_server, k3s_token))
env["KUBECONFIG"] = self._managed_kubeconfig
managed_set = True
_merge_kubeconfig(self._managed_kubeconfig)
except Exception as e:
print(f"[WARN] Failed to generate managed KUBECONFIG: {e}")
if not managed_set:
fallback = _find_kubeconfig_file(env)
if fallback:
env["KUBECONFIG"] = fallback
# Registry settings
reg_host = self.prole_cfg_data.get('Docker Build', {}).get('LOCAL_REGISTRY')
if reg_host:
env["LOCAL_REGISTRY"] = reg_host
reg_internal = self.prole_cfg_data.get('Docker Build', {}).get('LOCAL_REGISTRY_INTERNAL')
if reg_internal:
env["LOCAL_REGISTRY_INTERNAL"] = reg_internal
return env
def _k3s_kubectl_base_cmd(self) -> list[str]:
if self._managed_kubeconfig and os.path.exists(self._managed_kubeconfig):
return ["kubectl", "--kubeconfig", self._managed_kubeconfig]
server, token = self._k3s_connection_info()
if server and token:
# Token is already resolved by _k3s_connection_info
return [
"kubectl",
"--server=" + server,
"--token=" + token,
"--insecure-skip-tls-verify=true"
]
kubeconfig = _find_kubeconfig_file()
if kubeconfig:
return ["kubectl", "--kubeconfig", kubeconfig]
return ["kubectl"]
def _kubectl_base_cmd(self, env: str | None = None) -> list[str]:
env_key = self._cluster_env_key(env)
if env_key == 'service':
return self._k3s_kubectl_base_cmd()
kubeconfig_env = (os.environ.get("KUBECONFIG") or "").strip()
if kubeconfig_env:
return ["kubectl", "--kubeconfig", kubeconfig_env]
if env_key == 'prod':
kubeconfig = _find_kubeconfig_file()
if kubeconfig:
return ["kubectl", "--kubeconfig", kubeconfig]
return ["kubectl"]