prole/knoe/ui/screens/cluster_nodes.py
chrisfu 3312c39b1f feat: GCP/GKE CNPG hardening, Artifact Registry traffic light, and knoe-system namespace fixes
UI screens
- database.py: fix mode detection to use env_key priority (prod→k8s, service→k3s) so stale DEPLOYMENT_MODE never overrides the user's chosen environment
- database.py: Registry status reads ARTIFACT_REGISTRY_AVAILABLE persisted by cluster screen; uses SERVICE_NAMESPACE for Artifact Registry repo name
- cluster.py: add Artifact Registry traffic light (amber→green/red) to prod section; _check_artifact_registry_async persists ARTIFACT_REGISTRY_AVAILABLE into Global cfg
- cluster.py: re-trigger Artifact Registry check after GKE cluster selection so the light re-evaluates once region is available from KUBECONTEXT
- cluster_nodes.py: fix TclError on Python 3.14 — pady=(2,0) tuple → pady=2 scalar
- __init__.py: seed knoe-system namespace when saved value is "default", not only when empty
- services.py: replace hardcoded "Prole DB" log string with dynamic cnpg_cluster name

Core ops
- cloudnative_pg.py: replace one-shot Barman plugin retry with 6-attempt loop; first cert-manager/x509 failure triggers rollout restart + 30 s CA propagation wait; subsequent failures back off up to 60 s per attempt
- cloudnative_pg.py: TLS CA CN now uses cluster_name instead of hardcoded "Prole CNPG CA"
- registry.py, garage_store.py: refactored into per-mode modules (k3d/k3s/k8s registry and garage store, shared _garage_common)

Deploy / config
- deploy/gcp/gke/knoe-db.yaml: GKE-specific CNPG cluster manifest (rw/ro/r on separate nodes with premium-rwo storage)
- etc/init_common_services.sh, modes/k8s/knoe-db/.version: updated for current deploy
- kong-deployment.yaml: updated manifest

Tests
- test_cluster_nodes_render_smoke.py: add pack/grid, winfo_children, winfo_reqheight, update_idletasks, grid_slaves to dummy widgets; monkeypatch tk.Label so CNPG placement render completes without a real Tkinter root

Co-authored-by: Junie <junie@jetbrains.com>
2026-04-04 19:36:08 -07:00

930 lines
36 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.cnpg_placement import (
load_cnpg_placement_plan,
plan_cnpg_placement,
save_cnpg_placement_plan,
)
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,
"storage_probe_skip": False,
"storage_probe_quick": True,
"storage_probe_force_refresh": False,
"cnpg_placement_plan": None,
"cnpg_rebalance_requested": False,
}
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_placement_plan_path(self) -> Path:
conf_root = self._resolve_env_dir("PROLE_CONF", "conf")
return conf_root / "cnpg-placement" / "ecosystem-0-knoe-db.json"
def _cluster_nodes_compute_placement(
self, eligible_nodes: list[str], rebalance: bool = False
) -> dict:
plan_path = self._cluster_nodes_placement_plan_path()
prior = load_cnpg_placement_plan(plan_path)
plan = plan_cnpg_placement(
cluster_name="knoe-db",
desired_instances=3,
candidate_nodes=eligible_nodes,
prior_plan=prior,
rebalance=rebalance,
)
try:
save_cnpg_placement_plan(plan_path, plan)
except Exception:
pass
return plan
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()
is_gke = env_key == "prod"
cfg = TopologyDiscoveryConfig(
required_monitoring_mountpoints=(
() if is_gke else self._cluster_nodes_required_mountpoints()
),
monitoring_min_free_bytes=self._cluster_nodes_monitoring_min_free_bytes(),
storage_probe_enabled=(
False if is_gke else not bool(state.get("storage_probe_skip"))
),
storage_probe_quick=bool(state.get("storage_probe_quick", True)),
storage_probe_force_refresh=bool(force or state.get("storage_probe_force_refresh")),
skip_collector=is_gke,
assume_nodes_eligible=is_gke,
)
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:
eligible_names = sorted(node.name for node in nodes)
glob["CNPG_ELIGIBLE_NODES"] = ",".join(eligible_names)
glob["CNPG_STAGE1_NODE"] = eligible_names[0]
state = self._cluster_nodes_state()
rebalance = state.pop("cnpg_rebalance_requested", False)
plan = self._cluster_nodes_compute_placement(eligible_names, rebalance=rebalance)
state["cnpg_placement_plan"] = plan
glob["CNPG_PLACEMENT_PLAN_ID"] = str(plan.get("plan_id") or "")
glob["CNPG_PLACEMENT_PLAN_HASH"] = str(plan.get("plan_hash") or "")
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 node.storage_items:
parts: list[str] = []
for item in node.storage_items[:2]:
rr = item.metrics.randread_iops
rw = item.metrics.randwrite_iops
rr_txt = "?" if rr is None else f"{rr:.0f}"
rw_txt = "?" if rw is None else f"{rw:.0f}"
klass = item.io_class or "unknown"
parts.append(f"{item.label} [{klass}] R:{rr_txt}/W:{rw_txt}")
return ", ".join(parts)
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:
req_w = table.winfo_reqwidth()
canvas_w = canvas.winfo_width()
canvas.itemconfigure(inner, width=max(canvas_w, req_w))
canvas.configure(scrollregion=canvas.bbox("all"))
except Exception:
pass
table.bind("<Configure>", _cfg)
canvas.bind("<Configure>", _cfg)
headers = [
"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)
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:
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=0):
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 (3, 6, 7, 8, 9) 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
for item in node.storage_items:
rr = item.metrics.randread_iops
rw = item.metrics.randwrite_iops
rr_txt = "?" if rr is None else f"{rr:.0f}"
rw_txt = "?" if rw is None else f"{rw:.0f}"
class_txt = item.io_class or "unknown"
lat_parts = [v for v in [item.metrics.read_lat_ms_p95, item.metrics.write_lat_ms_p95] if v is not None]
lat_txt = f" · p95 {max(lat_parts):.1f} ms" if lat_parts else ""
cache_txt = " [cached]" if item.cached else ""
txt = (
f" ↳ {item.label} ({item.mountpoint}) "
f"R: {rr_txt} IOPS / W: {rw_txt} IOPS "
f"{class_txt}{lat_txt}{cache_txt}"
)
child = tk.Label(
table,
text=txt,
bg="#fcfcff",
fg="#3a3a3c",
font=("SF Pro Text", 9),
padx=6,
pady=3,
anchor="w",
justify="left",
wraplength=860,
)
child.grid(row=row, column=0, columnspan=len(headers), sticky="nsew")
self._overlay_widgets.append(child)
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"
storage_rows = []
for item in node.storage_items:
storage_rows.append(
f"- {item.label} mount={item.mountpoint} source={item.source_type} "
f"class={item.io_class or 'unknown'} cached={item.cached} approximate={item.approximate} "
f"readIops={item.metrics.randread_iops} writeIops={item.metrics.randwrite_iops} "
f"readP95Ms={item.metrics.read_lat_ms_p95} writeP95Ms={item.metrics.write_lat_ms_p95} "
f"runtimeSec={item.metrics.test_runtime_sec}"
)
storage_text = "\n".join(storage_rows) if storage_rows else "- none"
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}\n"
"Storage probe items:\n"
f"{storage_text}"
)
def _cluster_nodes_render_cnpg_placement(
self, topology: ClusterTopology | None, x: int, y: int
) -> int:
DESIRED = 3
CARD_WIDTH = 905
eligible = [n for n in (topology.nodes if topology else []) if n.capabilities.cnpg_eligible is True]
plan: dict | None = self._cluster_nodes_state().get("cnpg_placement_plan")
outer = tk.Frame(
self.bg_canvas,
bg="#f6f8fb",
highlightthickness=1,
highlightbackground="#d1d1d6",
)
self._overlay_widgets.append(outer)
win = self.bg_canvas.create_window(x, y, anchor="nw", width=CARD_WIDTH, window=outer)
self._canvas_items.append(win)
# Header row
hdr_frame = tk.Frame(outer, bg="#f0f4ff")
hdr_frame.pack(fill="x", padx=0, pady=0)
self._overlay_widgets.append(hdr_frame)
tk.Label(
hdr_frame,
text="CNPG Node Placement",
bg="#f0f4ff",
fg="#1d1d1f",
font=("SF Pro Text", 11, "bold"),
padx=10,
pady=6,
anchor="w",
).pack(side="left")
self._overlay_widgets.append(hdr_frame.winfo_children()[-1])
meta_txt = "3 instances · round-robin · premium-rwo · required node separation"
tk.Label(
hdr_frame,
text=meta_txt,
bg="#f0f4ff",
fg="#6e6e73",
font=("SF Pro Text", 10),
padx=6,
pady=6,
anchor="e",
).pack(side="right")
self._overlay_widgets.append(hdr_frame.winfo_children()[-1])
body = tk.Frame(outer, bg="#f6f8fb")
body.pack(fill="x", padx=10, pady=(4, 0))
self._overlay_widgets.append(body)
if len(eligible) < DESIRED:
warn_txt = (
f"Warning: only {len(eligible)} CNPG-eligible node(s) found; "
f"{DESIRED} required for separate-node placement."
)
tk.Label(
body,
text=warn_txt,
bg="#fff3cd",
fg="#856404",
font=("SF Pro Text", 10),
padx=8,
pady=6,
anchor="w",
justify="left",
).pack(fill="x", pady=(2, 6))
self._overlay_widgets.append(body.winfo_children()[-1])
elif plan:
assignments: dict = plan.get("assignments") or {}
col_headers = ["Ordinal", "Node", "Data storage", "WAL storage"]
table = tk.Frame(body, bg="#f6f8fb")
table.pack(fill="x", pady=(2, 4))
self._overlay_widgets.append(table)
for c, hdr in enumerate(col_headers):
tk.Label(
table,
text=hdr,
bg="#e8eaf6",
fg="#1d1d1f",
font=("SF Pro Text", 10, "bold"),
padx=6,
pady=4,
anchor="w",
width=22 if c == 1 else 14,
).grid(row=0, column=c, sticky="nsew", padx=(0, 1))
self._overlay_widgets.append(table.grid_slaves(row=0, column=c)[0])
for ordinal in range(DESIRED):
node_name = assignments.get(str(ordinal)) or assignments.get(ordinal) or "—"
row_vals = [str(ordinal), node_name, "premium-rwo 100 Gi", "premium-rwo 20 Gi"]
row_bg = "white" if ordinal % 2 == 0 else "#f9fbff"
for c, val in enumerate(row_vals):
tk.Label(
table,
text=val,
bg=row_bg,
fg="#1d1d1f",
font=("SF Pro Text", 10),
padx=6,
pady=3,
anchor="w",
width=22 if c == 1 else 14,
).grid(row=ordinal + 1, column=c, sticky="nsew", padx=(0, 1))
self._overlay_widgets.append(table.grid_slaves(row=ordinal + 1, column=c)[0])
meta = plan.get("metadata") or {}
reason = str(meta.get("reason") or "")
badge = "reused" if meta.get("reused") else reason.replace("_", " ")
plan_id = str(plan.get("plan_id") or "")
plan_line = f"{plan_id} [{badge}]" if plan_id else ""
if plan_line:
tk.Label(
body,
text=f"Plan: {plan_line}",
bg="#f6f8fb",
fg="#6e6e73",
font=("SF Pro Text", 9),
padx=4,
pady=2,
anchor="w",
).pack(anchor="w")
self._overlay_widgets.append(body.winfo_children()[-1])
note = "rw and ro roles float — CNPG assigns dynamically, not pinned by ordinal"
tk.Label(
body,
text=note,
bg="#f6f8fb",
fg="#6e6e73",
font=("SF Pro Text", 9, "italic"),
padx=4,
pady=2,
anchor="w",
).pack(anchor="w")
self._overlay_widgets.append(body.winfo_children()[-1])
btn_frame = tk.Frame(outer, bg="#f6f8fb")
btn_frame.pack(anchor="w", padx=10, pady=(4, 8))
self._overlay_widgets.append(btn_frame)
def _rebalance():
self._cluster_nodes_state()["cnpg_rebalance_requested"] = True
self._cluster_nodes_start_discovery(force=True)
rebalance_btn = tk.Button(
btn_frame,
text="Rebalance",
command=_rebalance,
bg="#f5f5dc",
fg="black",
activebackground="#e5e5d5",
highlightthickness=0,
relief="flat",
font=("SF Pro Text", 10),
padx=10,
pady=4,
)
rebalance_btn.pack(side="left")
self._overlay_widgets.append(rebalance_btn)
outer.update_idletasks()
card_h = outer.winfo_reqheight()
return y + card_h + 8
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):
self._clear_canvas_page()
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) if messages else "Starting topology discovery"
msg_scroll = tk.Scrollbar(msg_frame, orient="vertical")
msg_scroll.pack(side="right", fill="y")
msg_out = tk.Text(
msg_frame,
bg="#f6f8fb",
fg="#1d1d1f",
font=("SF Pro Text", 10),
padx=10,
pady=8,
relief="flat",
highlightthickness=0,
bd=0,
wrap="word",
yscrollcommand=msg_scroll.set,
)
msg_out.pack(side="left", fill="both", expand=True)
msg_scroll.configure(command=msg_out.yview)
msg_out.insert("1.0", msg_text)
msg_out.see("end")
msg_out.configure(state="disabled")
self._overlay_widgets.extend([msg_scroll, msg_out])
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)
y = self._cluster_nodes_render_cnpg_placement(topology, 48, y + 12)
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,
)
self._overlay_widgets.append(refresh_btn)
self._canvas_items.append(
self.bg_canvas.create_window(48, y + 8, window=refresh_btn, anchor="nw", width=200)
)
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