prole/knoe/ui/screens/cluster_nodes.py
chrisfu cf51335ced Fix kube context handling and monitoring eligibility
- detect local k3s node kubeconfig and skip kubectx/use-context mutation when already targeting local API\n- add configurable KUBE_CONTEXT_NAME resolution with compatibility fallbacks and switch only when required\n- update init scripts to use ensure_kube_context helper naming\n- broaden monitoring eligibility to discovered /synology/d### mounts so /synology/d004 qualifies\n- add focused kube-context and topology tests covering local/remote and read-only kubeconfig cases

Co-authored-by: Junie <junie@jetbrains.com>
2026-03-23 13:50:29 -07:00

676 lines
26 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

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

"""Cluster Nodes screen: topology discovery + node capability view."""
from __future__ import annotations
import json
import threading
import tkinter as tk
from pathlib import Path
from tkinter import messagebox
from knoe import screen as ui
from knoe.core.topology import (
ClusterTopology,
DEFAULT_MONITORING_MIN_FREE_BYTES,
DEFAULT_REQUIRED_MONITORING_MOUNTPOINTS,
NodeTopology,
TopologyDiscoveryConfig,
discover_cluster_topology,
)
class ClusterNodesScreenMixin:
CLUSTER_NODE_SERVICES: list[tuple[str, str]] = [
("argocd", "ArgoCD"),
("gitlab", "GitLab"),
("gitea", "Gitea"),
("cloudnativepg", "CloudNative-PG"),
("pihole", "Pi-hole"),
]
ADVANCED_ASSIGNMENT_GRID_KEY = "ADVANCED_ASSIGNMENT_GRID"
def _cluster_nodes_state(self) -> dict:
state = getattr(self, "_cluster_nodes_topology_state", None)
if isinstance(state, dict):
return state
state = {
"phase": "pending",
"messages": [],
"topology": None,
"topology_path": "",
"selected": set(),
"busy": False,
"expanded": None,
}
self._cluster_nodes_topology_state = state
return state
def _cluster_nodes_add_message(self, message: str, rerender: bool = True) -> None:
state = self._cluster_nodes_state()
messages = state.setdefault("messages", [])
if not messages or messages[-1] != message:
messages.append(message)
if rerender:
try:
self.root.after(0, self._render_cluster_nodes_page)
except Exception:
pass
def _cluster_nodes_required_mountpoints(self) -> tuple[str, ...]:
return DEFAULT_REQUIRED_MONITORING_MOUNTPOINTS
def _cluster_nodes_monitoring_min_free_bytes(self) -> int:
return DEFAULT_MONITORING_MIN_FREE_BYTES
def _cluster_nodes_topology_root(self) -> Path:
return self._resolve_env_dir("PROLE_DATA", "data") / "topology"
def _cluster_nodes_start_discovery(self, force: bool = False) -> None:
state = self._cluster_nodes_state()
if state.get("busy"):
return
if state.get("topology") is not None and not force:
return
state["busy"] = True
state["phase"] = "discovering"
if force:
state["messages"] = []
state["topology"] = None
def worker() -> None:
try:
self._cluster_nodes_add_message("Starting topology discovery", rerender=False)
env_key = self._cluster_env_key()
base_cmd = self._kubectl_base_cmd(env_key)
output_root = self._cluster_nodes_topology_root()
cfg = TopologyDiscoveryConfig(
required_monitoring_mountpoints=self._cluster_nodes_required_mountpoints(),
monitoring_min_free_bytes=self._cluster_nodes_monitoring_min_free_bytes(),
)
result = discover_cluster_topology(
kubectl_base_cmd=base_cmd,
topology_root=output_root,
config=cfg,
progress_cb=lambda msg: self._cluster_nodes_add_message(msg, rerender=False),
)
topology = result.topology
state["topology"] = topology
state["phase"] = topology.collection_status
state["topology_path"] = str(output_root)
selected = state.setdefault("selected", set())
if not selected:
for node in topology.nodes:
if node.capabilities.monitoring_eligible is True:
selected.add(node.name)
self._cluster_nodes_apply_topology_config(topology)
except Exception as exc:
state["phase"] = "failed"
self._cluster_nodes_add_message(
f"Topology discovery failed; continuing with Kubernetes-only data: {exc}",
rerender=False,
)
finally:
state["busy"] = False
try:
self.root.after(0, self._render_cluster_nodes_page)
except Exception:
pass
threading.Thread(target=worker, daemon=True).start()
def _cluster_nodes_apply_topology_config(self, topology: ClusterTopology) -> None:
nodes = [n for n in topology.nodes if n.capabilities.cnpg_eligible is True]
synology_roots = sorted(
{
fs.mountpoint
for node in topology.nodes
for fs in node.filesystems
if fs.mountpoint.startswith("/synology/")
}
)
cluster_sec = self.prole_cfg_data.setdefault("Cluster Nodes", {})
cluster_sec["TOPOLOGY_STATUS"] = topology.collection_status
cluster_sec["TOPOLOGY_PATH"] = topology.topology_root
glob = self.prole_cfg_data.setdefault("Global", {})
if nodes:
glob["CNPG_ELIGIBLE_NODES"] = ",".join(sorted(node.name for node in nodes))
glob["CNPG_STAGE1_NODE"] = sorted(node.name for node in nodes)[0]
if synology_roots:
glob["SYNOLOGY_ROOTS"] = ",".join(synology_roots)
try:
self._save_prole_cfg()
except Exception:
pass
@staticmethod
def _fmt_bytes(value: int | None) -> str:
if value is None:
return "Unknown"
units = ["B", "KiB", "MiB", "GiB", "TiB", "PiB"]
val = float(value)
idx = 0
while val >= 1024.0 and idx < len(units) - 1:
val /= 1024.0
idx += 1
return f"{val:.1f} {units[idx]}"
@staticmethod
def _fmt_bool(value: bool | None) -> str:
if value is None:
return "Unknown"
return "Yes" if value else "No"
def _cluster_nodes_storage_summary(self, node: NodeTopology) -> str:
if not node.filesystems:
return "Topology unavailable"
rows: list[str] = []
for fs in sorted(node.filesystems, key=lambda x: x.mountpoint):
if fs.mountpoint == "/":
rows.append(f"1× {fs.storage_class.upper()} root")
continue
if fs.mountpoint.startswith("/synology/"):
free_txt = self._fmt_bytes(fs.available_bytes)
rows.append(f"{fs.fstype.upper()} {fs.mountpoint} {free_txt} free")
return ", ".join(rows[:2]) if rows else "No Synology mounts"
def _cluster_nodes_status_color(self, phase: str) -> str:
return {
"pending": "#6e6e73",
"discovering": "#0a84ff",
"complete": "#34c759",
"partial": "#ff9f0a",
"failed": "#ff3b30",
}.get(phase, "#6e6e73")
def _cluster_nodes_toggle_selection(self, node_name: str, selected: bool) -> None:
state = self._cluster_nodes_state()
chosen: set[str] = state.setdefault("selected", set())
if selected:
chosen.add(node_name)
else:
chosen.discard(node_name)
def _cluster_nodes_render_summary_cards(self, topology: ClusterTopology | None, x: int, y: int) -> int:
cards = [
("Discovered", str(topology.discovered_nodes if topology else 0)),
("Ready", str(topology.ready_nodes if topology else 0)),
("Monitoring eligible", str(topology.monitoring_eligible_nodes if topology else 0)),
("CNPG eligible", str(topology.cnpg_eligible_nodes if topology else 0)),
]
width = 210
for idx, (label, value) in enumerate(cards):
px = x + idx * (width + 10)
frame = tk.Frame(self.bg_canvas, bg="#f6f8fb", bd=0, highlightthickness=1, highlightbackground="#d1d1d6")
self._overlay_widgets.append(frame)
win = self.bg_canvas.create_window(px, y, anchor="nw", width=width, height=76, window=frame)
self._canvas_items.append(win)
lbl = tk.Label(frame, text=label, bg="#f6f8fb", fg="#6e6e73", font=("SF Pro Text", 10))
lbl.pack(anchor="w", padx=10, pady=(8, 0))
val = tk.Label(frame, text=value, bg="#f6f8fb", fg="#1d1d1f", font=("SF Pro Text", 20, "bold"))
val.pack(anchor="w", padx=10, pady=(2, 8))
self._overlay_widgets.extend([lbl, val])
return y + 88
def _cluster_nodes_render_table(self, topology: ClusterTopology | None, x: int, y: int) -> int:
frame = tk.Frame(self.bg_canvas, bg="white", highlightthickness=1, highlightbackground="#d1d1d6")
self._overlay_widgets.append(frame)
win = self.bg_canvas.create_window(x, y, anchor="nw", width=905, height=420, window=frame)
self._canvas_items.append(win)
canvas = tk.Canvas(frame, bg="white", highlightthickness=0, bd=0)
xbar = tk.Scrollbar(frame, orient="horizontal", command=canvas.xview)
ybar = tk.Scrollbar(frame, orient="vertical", command=canvas.yview)
canvas.configure(xscrollcommand=xbar.set, yscrollcommand=ybar.set)
xbar.pack(side="bottom", fill="x")
ybar.pack(side="right", fill="y")
canvas.pack(side="left", fill="both", expand=True)
table = tk.Frame(canvas, bg="white")
inner = canvas.create_window(0, 0, window=table, anchor="nw")
def _cfg(_evt=None):
try:
canvas.configure(scrollregion=canvas.bbox("all"))
canvas.itemconfigure(inner, height=canvas.winfo_height())
except Exception:
pass
table.bind("<Configure>", _cfg)
canvas.bind("<Configure>", _cfg)
headers = [
"Select",
"Node",
"Role",
"Arch",
"CPU",
"RAM",
"Ready",
"Storage summary",
"Synology mounts",
"Monitoring eligible",
"CNPG eligible",
"Details",
]
for c, header in enumerate(headers):
table.grid_columnconfigure(c, minsize=120 if c > 0 else 70)
cell = tk.Label(
table,
text=header,
bg="#f5f6fa",
fg="#1d1d1f",
font=("SF Pro Text", 10, "bold"),
padx=6,
pady=6,
anchor="w",
)
cell.grid(row=0, column=c, sticky="nsew")
self._overlay_widgets.append(cell)
nodes = list(topology.nodes) if topology else []
state = self._cluster_nodes_state()
selected: set[str] = state.setdefault("selected", set())
expanded = state.get("expanded")
row = 1
for node in nodes:
select_var = tk.BooleanVar(master=self.root, value=node.name in selected)
cb = tk.Checkbutton(
table,
variable=select_var,
bg="white",
activebackground="white",
command=lambda n=node.name, v=select_var: self._cluster_nodes_toggle_selection(n, bool(v.get())),
)
cb.grid(row=row, column=0, sticky="nsew")
self._overlay_widgets.append(cb)
role = ",".join(node.roles) or "worker"
cpu_text = " · ".join(
p
for p in [
node.cpu_model or "Unknown CPU",
(
f"{node.physical_cores or '?'}c/{node.logical_cores or '?'}t"
if node.physical_cores or node.logical_cores
else ""
),
]
if p
)
synology = ", ".join(node.capabilities.matching_mountpoints) or "—"
values = [
node.name,
role,
node.architecture or "Unknown",
cpu_text,
self._fmt_bytes(node.memory_bytes),
"Ready" if node.ready else "NotReady",
self._cluster_nodes_storage_summary(node),
synology,
f"{self._fmt_bool(node.capabilities.monitoring_eligible)} ({node.capabilities.monitoring_eligible_state})",
f"{self._fmt_bool(node.capabilities.cnpg_eligible)} ({node.capabilities.cnpg_eligible_state})",
]
for idx, txt in enumerate(values, start=1):
cell = tk.Label(
table,
text=txt,
bg="white",
fg="#1d1d1f",
font=("SF Pro Text", 10),
padx=6,
pady=4,
anchor="w",
justify="left",
wraplength=240 if idx in (4, 7, 8, 9, 10) else 120,
)
cell.grid(row=row, column=idx, sticky="nsew")
self._overlay_widgets.append(cell)
toggle_txt = "Hide" if expanded == node.name else "Show"
btn = tk.Button(
table,
text=toggle_txt,
command=lambda n=node.name: self._cluster_nodes_toggle_details(n),
bg="#f5f5dc",
fg="black",
relief="flat",
font=("SF Pro Text", 10),
)
btn.grid(row=row, column=len(headers) - 1, sticky="nsew", padx=4)
self._overlay_widgets.append(btn)
row += 1
if expanded == node.name:
detail = self._cluster_nodes_detail_text(node)
detail_lbl = tk.Label(
table,
text=detail,
bg="#f9fbff",
fg="#1d1d1f",
justify="left",
anchor="w",
padx=8,
pady=8,
wraplength=860,
font=("SF Pro Text", 10),
)
detail_lbl.grid(row=row, column=0, columnspan=len(headers), sticky="nsew")
self._overlay_widgets.append(detail_lbl)
row += 1
if not nodes:
label = tk.Label(
table,
text="Topology data is not available yet.",
bg="white",
fg="#6e6e73",
font=("SF Pro Text", 11),
)
label.grid(row=1, column=0, columnspan=len(headers), sticky="nsew", padx=8, pady=10)
self._overlay_widgets.append(label)
return y + 432
def _cluster_nodes_toggle_details(self, node_name: str) -> None:
state = self._cluster_nodes_state()
state["expanded"] = None if state.get("expanded") == node_name else node_name
self._render_cluster_nodes_page()
def _cluster_nodes_detail_text(self, node: NodeTopology) -> str:
labels = ", ".join(f"{k}={v}" for k, v in sorted(node.labels.items())) or "None"
taints = ", ".join(node.taints) or "None"
reasons = ", ".join(node.capabilities.reasons) or "None"
fs_rows = []
for fs in node.filesystems:
fs_rows.append(
f"- {fs.mountpoint} ({fs.fstype}) device={fs.device} sizeBytes={fs.size_bytes} "
f"usedBytes={fs.used_bytes} availableBytes={fs.available_bytes} "
f"transport={fs.transport} rotational={fs.rotational} storageClass={fs.storage_class}"
)
filesystems = "\n".join(fs_rows) if fs_rows else "- topology_unavailable"
return (
f"Labels: {labels}\n"
f"Taints: {taints}\n"
f"Allocatable: cpu={node.allocatable_cpu} memoryBytes={node.allocatable_memory_bytes} "
f"ephemeralStorageBytes={node.allocatable_ephemeral_storage_bytes}\n"
f"Host hardware: cpuModel={node.cpu_model or 'Unknown'} physicalCores={node.physical_cores} "
f"logicalCores={node.logical_cores} memoryBytes={node.memory_bytes}\n"
f"Capabilities reasons: {reasons}\n"
"Filesystems:\n"
f"{filesystems}"
)
def _cluster_nodes_save_selected(self) -> None:
state = self._cluster_nodes_state()
selected = sorted(state.get("selected") or [])
if "Cluster Nodes" not in self.prole_cfg_data:
self.prole_cfg_data["Cluster Nodes"] = {}
self.prole_cfg_data["Cluster Nodes"]["SELECTED_NODES"] = ",".join(selected)
try:
self._save_prole_cfg()
messagebox.showinfo("Cluster Nodes", "Saved selected cluster nodes.")
except Exception:
pass
def _render_cluster_nodes_page(self):
if not getattr(self, "_should_show_cluster_nodes_screen", lambda: True)():
self._render_title("Cluster Nodes", y=150)
self._render_paragraph(
"This screen is only shown for multi-node, non-k3d clusters.",
y=200,
wrap=860,
)
ui.canvas_text(
self,
48,
260,
"Continue to the next step using the sidebar or Next.",
fill="#6e6e73",
font=("SF Pro Text", 11),
)
return
self._render_title("Cluster Nodes", y=150)
self._render_paragraph(
"Discovering node topology from inside the cluster. The table below combines Kubernetes facts "
"with host inventory and shows monitoring/CNPG eligibility with reasons.",
y=200,
wrap=900,
)
state = self._cluster_nodes_state()
if not state.get("busy") and state.get("topology") is None:
self._cluster_nodes_start_discovery(force=False)
phase = str(state.get("phase") or "pending")
color = self._cluster_nodes_status_color(phase)
messages: list[str] = list(state.get("messages") or [])
latest = messages[-1] if messages else "Topology discovery is pending."
ui.canvas_text(
self,
48,
248,
f"Status: {phase.upper()} — {latest}",
fill=color,
font=("SF Pro Text", 11, "bold"),
)
if state.get("topology_path"):
ui.canvas_text(
self,
48,
272,
f"Topology snapshot path: {state['topology_path']}",
fill="#1d1d1f",
font=("SF Pro Text", 10),
)
msg_frame = tk.Frame(self.bg_canvas, bg="#f6f8fb", highlightthickness=1, highlightbackground="#d1d1d6")
self._overlay_widgets.append(msg_frame)
msg_win = self.bg_canvas.create_window(48, 292, window=msg_frame, anchor="nw", width=905, height=86)
self._canvas_items.append(msg_win)
msg_text = "\n".join(messages[-4:]) if messages else "Starting topology discovery"
msg_lbl = tk.Label(
msg_frame,
text=msg_text,
bg="#f6f8fb",
fg="#1d1d1f",
justify="left",
anchor="w",
font=("SF Pro Text", 10),
padx=10,
pady=8,
)
msg_lbl.pack(fill="both", expand=True)
self._overlay_widgets.append(msg_lbl)
topology = state.get("topology") if isinstance(state.get("topology"), ClusterTopology) else None
y = self._cluster_nodes_render_summary_cards(topology, 48, 388)
y = self._cluster_nodes_render_table(topology, 48, y)
refresh_btn = tk.Button(
self.bg_canvas,
text="Refresh topology",
command=lambda: self._cluster_nodes_start_discovery(force=True),
bg="#f5f5dc",
fg="black",
activebackground="#e5e5d5",
highlightthickness=0,
relief="flat",
font=("SF Pro Text", 11),
padx=12,
pady=6,
)
save_btn = tk.Button(
self.bg_canvas,
text="Save selected nodes",
command=self._cluster_nodes_save_selected,
bg="#f5f5dc",
fg="black",
activebackground="#e5e5d5",
highlightthickness=0,
relief="flat",
font=("SF Pro Text", 11),
padx=12,
pady=6,
)
self._overlay_widgets.extend([refresh_btn, save_btn])
self._canvas_items.append(
self.bg_canvas.create_window(48, y + 8, window=refresh_btn, anchor="nw", width=200)
)
self._canvas_items.append(
self.bg_canvas.create_window(260, y + 8, window=save_btn, anchor="nw", width=220)
)
def _cluster_nodes_advanced_enabled(self) -> bool:
sec = (getattr(self, "prole_cfg_data", None) or {}).get("Cluster Nodes", {})
raw = (sec or {}).get(self.ADVANCED_ASSIGNMENT_GRID_KEY, False)
if isinstance(raw, bool):
return raw
return str(raw).strip().lower() in {"1", "true", "yes", "on"}
def _cluster_nodes_set_advanced_enabled(self, enabled: bool) -> None:
if "Cluster Nodes" not in self.prole_cfg_data:
self.prole_cfg_data["Cluster Nodes"] = {}
self.prole_cfg_data["Cluster Nodes"][self.ADVANCED_ASSIGNMENT_GRID_KEY] = "true" if enabled else "false"
try:
self._save_prole_cfg()
except Exception:
pass
def _cluster_nodes_monitoring_host(self, hosts: list[str]) -> str:
for host in hosts:
if "d004" in host.lower():
return host
return hosts[-1] if hosts else ""
def _cluster_nodes_default_policy(self, hosts: list[str]) -> dict:
services = {}
for sid, _title in self.CLUSTER_NODE_SERVICES:
services[sid] = {"primary_host": "", "enabled_hosts": list(hosts)}
return {
"defaults": {
"cnpg_allocation": "round-robin",
"monitoring_storage_root": "/synology/d004",
"monitoring_host": self._cluster_nodes_monitoring_host(hosts),
},
"hosts": {h: {"protected": False, "reserved_hostports": False} for h in hosts},
"services": services,
}
def _cluster_nodes_hosts(self) -> list[str]:
topology = self._cluster_nodes_state().get("topology")
if isinstance(topology, ClusterTopology) and topology.nodes:
return [node.name for node in topology.nodes]
topo = getattr(self, "ansible_topology", None) or {}
groups = topo.get("groups") or {}
domain = (topo.get("domain") or "").strip()
hosts: list[str] = []
k3s_hosts = groups.get("k3s_hosts") or []
if isinstance(k3s_hosts, list) and k3s_hosts:
hosts = list(k3s_hosts)
else:
ip_hosts = topo.get("hosts") or {}
if isinstance(ip_hosts, dict):
hosts = sorted(ip_hosts.keys())
out: list[str] = []
seen = set()
for host in hosts:
host = str(host or "").strip()
if not host:
continue
display = f"{host}.{domain}" if domain and "." not in host else host
if display not in seen:
seen.add(display)
out.append(display)
return out
def _cluster_nodes_load_policy(self, hosts: list[str]) -> dict:
sec = (getattr(self, "prole_cfg_data", None) or {}).get("Cluster Nodes", {})
raw = (sec or {}).get("POLICY_JSON", "")
if raw:
try:
pol = json.loads(raw)
if isinstance(pol, dict):
return self._cluster_nodes_normalize_policy(pol, hosts)
except Exception:
pass
return self._cluster_nodes_default_policy(hosts)
def _cluster_nodes_normalize_policy(self, pol: dict, hosts: list[str]) -> dict:
out = {"defaults": {}, "hosts": {}, "services": {}}
host_set = set(hosts)
defaults_in = pol.get("defaults") if isinstance(pol.get("defaults"), dict) else {}
out["defaults"] = {
"cnpg_allocation": str(defaults_in.get("cnpg_allocation") or "round-robin"),
"monitoring_storage_root": str(defaults_in.get("monitoring_storage_root") or "/synology/d004"),
"monitoring_host": str(defaults_in.get("monitoring_host") or self._cluster_nodes_monitoring_host(hosts)),
}
for host in hosts:
hpol = ((pol.get("hosts") or {}).get(host) or {}) if isinstance(pol.get("hosts"), dict) else {}
out["hosts"][host] = {
"protected": bool(hpol.get("protected")),
"reserved_hostports": bool(hpol.get("reserved_hostports")),
}
services_in = pol.get("services") if isinstance(pol.get("services"), dict) else {}
for sid, _title in self.CLUSTER_NODE_SERVICES:
svc = (services_in.get(sid) or {}) if isinstance(services_in, dict) else {}
enabled = [host for host in (svc.get("enabled_hosts") or []) if host in host_set]
primary = (svc.get("primary_host") or "").strip()
if primary and primary not in host_set:
primary = ""
if primary and primary not in enabled:
enabled = sorted(set(enabled + [primary]))
out["services"][sid] = {"primary_host": primary, "enabled_hosts": enabled}
return out
def _validate_and_save_cluster_nodes_policy(self, policy: dict) -> bool:
for sid, _title in self.CLUSTER_NODE_SERVICES:
svc = policy.get("services", {}).get(sid, {})
primary = (svc.get("primary_host") or "").strip()
enabled = set(svc.get("enabled_hosts") or [])
if primary and primary not in enabled:
try:
messagebox.showerror(
"Cluster Nodes",
f"Service '{sid}': Primary Host must also be Enabled Here.",
)
except Exception:
pass
return False
try:
raw = json.dumps(policy, indent=2, sort_keys=True)
except Exception as exc:
try:
messagebox.showerror("Cluster Nodes", f"Could not serialize policy: {exc}")
except Exception:
pass
return False
if "Cluster Nodes" not in self.prole_cfg_data:
self.prole_cfg_data["Cluster Nodes"] = {}
self.prole_cfg_data["Cluster Nodes"]["POLICY_JSON"] = raw
try:
self._save_prole_cfg()
except Exception:
pass
return True