prole/knoe/core/topology.py
chrisfu 4a8d9cc90d feat: full GKE/prod deployment pipeline from UI to Artifact Registry
## GCP / Cluster Environment Screen
- Auto-populate Cloud tab from conf/prod/gcp.cfg on screen open (org_id,
  billing_account, billing_project, project_id)
- gcloud auth validity checked on screen startup; friendly modal dialog
  streams gcloud auth login output live so user never leaves the app
- Live GKE cluster browser: fetches clusters via gcloud container clusters
  list, displays with checkmark selector, auto-selects saved cluster
- Selecting a cluster runs get-credentials, sets KUBECONFIG/KUBECONTEXT,
  and syncs the region dropdown to the selected cluster's location
- Region dropdown populated live from gcloud compute regions list with
  checkmark on currently selected region; graceful fallback when offline
- New 'GCP Storage' tab with workload->StorageClass mapping (CNPG->premium-rwo,
  Redis/Monitoring->standard-rwo, Garage->garage-hdd) and Fetch from Cluster
- Provider readonly field styled correctly (no solid-black on macOS)
- Stale prole.cfg/conf/prole.cfg symlinks removed; all config I/O now
  resolves env-specific paths via prole_conf.entrypoint_path()

## GKE Autopilot Compatibility (Common Services)
- Synology iSCSI StorageClass and static PVs guarded behind PROLE_MODE!=k8s
  in init_openbao.sh (GKE Autopilot forbids hostPath/iSCSI volumes)
- In-cluster Docker registry (hostPath) skipped in k8s mode; GCP Artifact
  Registry used instead
- Kong renamed knoe-svc-kong in k8s mode; all health-check kubectl calls in
  init_common_services.sh and status_common_services.sh updated accordingly
- DNS endpoints switched from *.prole.org to *.knoe.dev in k8s mode
  (api.knoe.dev, git.knoe.dev, svc.knoe.dev); ingress uses gce class
- New GKE-clean Kong manifests under deploy/opentofu/k8s/manifests/prole/:
  no k3s node affinity, explicit Autopilot resource requests/limits

## Garage S3 Store (GKE)
- New garage-statefulset-gcp.yaml targeting garage-hdd StorageClass
  (pd-standard, avoids SSD_TOTAL_GB quota exhaustion in us-west3)
- New storageclass-gcp-hdd.yaml (pd-standard, Retain, WaitForFirstConsumer)
- GCP StorageClass manifests skipped on re-runs (Autopilot built-ins are
  immutable; skip-if-exists guard added)
- PVC deletion guard extended to cover any storageClass (not just synology)
  so stale claims are cleaned before StatefulSet recreation

## Topology (GKE Autopilot)
- DaemonSet collector skipped in prod mode (forbidden in kube-system by
  GKE Warden); Kubernetes-only node facts path used instead
- All ready GKE nodes assumed cnpg-eligible and monitoring-eligible without
  taint/synology-mount checks (skip_collector + assume_nodes_eligible flags)

## KUBECONFIG / kubectl (k8s mode)
- actions.py: new elif mode==k8s branch sets KUBECONFIG=~/.kube/config
  and injects KUBECONTEXT from prole_cfg_data into script env
- _build_kubectl_cmd falls back to Global.KUBECONTEXT when
  init_cluster.selected_kubectx is empty
- _activate_selected_gke_cluster persists KUBECONFIG/KUBECONTEXT to
  prole_cfg_data and saves prole.cfg immediately after get-credentials

## Database Build Screen (GKE)
- Registry display shows correct Artifact Registry URL
  (<region>-docker.pkg.dev/<project>/<namespace>/knoe-db) in green
- Build+push: gcloud auth configure-docker, auto-creates AR repository
  named after SERVICE_NAMESPACE (e.g. knoe-system) if missing, then
  docker tag + push; falls back to gcr.io if region unavailable
- GCP config loaded from conf/prod/gcp.cfg on every screen entry;
  keys normalised to lowercase so project_id lookup is always consistent

## Config / Namespace persistence
- prole_conf.py activate_environment: symlink creation removed; sets
  CLUSTER_ENV env-var so all subsequent calls resolve correct env directory
- knoe/ui/screens/__init__.py: startup config load uses entrypoint_path()
  instead of hardcoded conf/prole.cfg; seeds SERVICE_NAMESPACE=knoe-system
  for managed envs so Common Services never defaults to 'default'
- cfg.py _save_prole_cfg: saves to env-specific path via entrypoint_path()
- etc/prole_cfg.sh: removed all ln -snf symlink creation

Co-authored-by: Junie <junie@jetbrains.com>
2026-04-04 12:38:16 -07:00

1224 lines
43 KiB
Python

from __future__ import annotations
import json
import re
import secrets
import subprocess
import time
from dataclasses import dataclass, replace
from datetime import datetime, timezone
from pathlib import Path
from typing import Callable, Sequence
import xml.etree.ElementTree as ET
from knoe.core.storage_probe import (
DEFAULT_CACHE_TTL_SEC as STORAGE_PROBE_DEFAULT_CACHE_TTL_SEC,
DEFAULT_MAX_TARGETS_PER_NODE as STORAGE_PROBE_DEFAULT_MAX_TARGETS_PER_NODE,
DEFAULT_TIMEOUT_SEC as STORAGE_PROBE_DEFAULT_TIMEOUT_SEC,
StorageProbeResult,
probe_node_storage,
)
TOPOLOGY_SCHEMA_VERSION = "1.0"
INVENTORY_BEGIN_MARKER = "PROLE_NODE_INVENTORY_BEGIN"
INVENTORY_END_MARKER = "PROLE_NODE_INVENTORY_END"
TOPOLOGY_NAMESPACE = "kube-system"
DEFAULT_REQUIRED_MONITORING_MOUNTPOINTS: tuple[str, ...] = ("/synology/d001",)
DEFAULT_MONITORING_MIN_FREE_BYTES = 200 * 1024 * 1024 * 1024
_SYNOLOGY_DATA_MOUNT_RE = re.compile(r"^/synology/d\d{3}$")
_NETWORK_FILESYSTEM_TYPES = {
"nfs",
"nfs4",
"cifs",
"smb",
"glusterfs",
"ceph",
"cephfs",
"fuse.sshfs",
}
@dataclass(slots=True, frozen=True)
class FilesystemTopology:
device: str
mountpoint: str
fstype: str
size_bytes: int | None
used_bytes: int | None
available_bytes: int | None
transport: str
rotational: bool | None
storage_class: str
@dataclass(slots=True, frozen=True)
class CapabilityFlags:
has_synology_root: bool | None
matching_mountpoints: tuple[str, ...]
monitoring_eligible: bool | None
cnpg_eligible: bool | None
reasons: tuple[str, ...]
has_synology_root_state: str
monitoring_eligible_state: str
cnpg_eligible_state: str
@dataclass(slots=True, frozen=True)
class NodeTopology:
name: str
roles: tuple[str, ...]
labels: dict[str, str]
taints: tuple[str, ...]
ready: bool
internal_ip: str
architecture: str
operating_system: str
allocatable_cpu: str
allocatable_memory_bytes: int | None
allocatable_ephemeral_storage_bytes: int | None
cpu_model: str | None
physical_cores: int | None
logical_cores: int | None
memory_bytes: int | None
filesystems: tuple[FilesystemTopology, ...]
storage_items: tuple[StorageProbeResult, ...]
host_inventory_available: bool
host_inventory_error: str | None
capabilities: CapabilityFlags
@dataclass(slots=True, frozen=True)
class ClusterTopology:
version: str
generated_at: str
generated_epoch_ms: int
collection_status: str
required_monitoring_mountpoints: tuple[str, ...]
monitoring_min_free_bytes: int
nodes: tuple[NodeTopology, ...]
topology_root: str
messages: tuple[str, ...]
@property
def discovered_nodes(self) -> int:
return len(self.nodes)
@property
def ready_nodes(self) -> int:
return sum(1 for node in self.nodes if node.ready)
@property
def monitoring_eligible_nodes(self) -> int:
return sum(1 for node in self.nodes if node.capabilities.monitoring_eligible is True)
@property
def cnpg_eligible_nodes(self) -> int:
return sum(1 for node in self.nodes if node.capabilities.cnpg_eligible is True)
@dataclass(slots=True, frozen=True)
class TopologyDiscoveryConfig:
namespace: str = TOPOLOGY_NAMESPACE
timeout_seconds: int = 120
poll_interval_seconds: float = 2.0
required_monitoring_mountpoints: tuple[str, ...] = DEFAULT_REQUIRED_MONITORING_MOUNTPOINTS
monitoring_min_free_bytes: int = DEFAULT_MONITORING_MIN_FREE_BYTES
collector_image: str = "python:3.11-slim"
storage_probe_enabled: bool = True
storage_probe_quick: bool = True
storage_probe_force_refresh: bool = False
storage_probe_timeout_sec: int = STORAGE_PROBE_DEFAULT_TIMEOUT_SEC
storage_probe_cache_ttl_sec: int = STORAGE_PROBE_DEFAULT_CACHE_TTL_SEC
storage_probe_max_targets_per_node: int = STORAGE_PROBE_DEFAULT_MAX_TARGETS_PER_NODE
storage_probe_cache_path: Path | None = None
# When True, skip the in-cluster DaemonSet collector entirely (e.g. GKE Autopilot
# forbids creating resources in kube-system).
skip_collector: bool = False
# When True, treat all Ready nodes as monitoring- and cnpg-eligible without
# requiring host inventory or checking taints/mountpoints (e.g. GKE managed nodes).
assume_nodes_eligible: bool = False
@dataclass(slots=True, frozen=True)
class TopologyDiscoveryResult:
topology: ClusterTopology
collector_applied: bool
expected_ready_nodes: tuple[str, ...]
reported_nodes: tuple[str, ...]
missing_nodes: tuple[str, ...]
def _run_kubectl(
base_cmd: Sequence[str],
args: Sequence[str],
*,
timeout: int = 20,
env: dict[str, str] | None = None,
input_text: str | None = None,
) -> subprocess.CompletedProcess[str]:
return subprocess.run(
[*base_cmd, *args],
capture_output=True,
text=True,
timeout=timeout,
env=env,
input=input_text,
)
def _kubectl_json(
base_cmd: Sequence[str],
args: Sequence[str],
*,
timeout: int = 20,
env: dict[str, str] | None = None,
) -> dict:
res = _run_kubectl(base_cmd, [*args, "-o", "json"], timeout=timeout, env=env)
if res.returncode != 0:
stderr = (res.stderr or "").strip()
raise RuntimeError(f"kubectl command failed: {' '.join(args)}: {stderr}")
try:
return json.loads(res.stdout or "{}")
except json.JSONDecodeError as exc:
raise RuntimeError(f"Invalid kubectl JSON payload for command: {' '.join(args)}") from exc
def _parse_k8s_quantity_bytes(raw: str | None) -> int | None:
text = str(raw or "").strip()
if not text:
return None
m = re.match(r"^([0-9]+(?:\.[0-9]+)?)([A-Za-z]+)?$", text)
if not m:
return None
num = float(m.group(1))
unit = (m.group(2) or "").strip()
binary_units = {
"Ki": 1024,
"Mi": 1024**2,
"Gi": 1024**3,
"Ti": 1024**4,
"Pi": 1024**5,
"Ei": 1024**6,
}
decimal_units = {
"K": 1000,
"M": 1000**2,
"G": 1000**3,
"T": 1000**4,
"P": 1000**5,
"E": 1000**6,
}
if not unit:
return int(num)
if unit in binary_units:
return int(num * binary_units[unit])
if unit in decimal_units:
return int(num * decimal_units[unit])
return None
def _node_roles(labels: dict[str, str]) -> tuple[str, ...]:
roles: list[str] = []
for key, val in labels.items():
if key.startswith("node-role.kubernetes.io/"):
suffix = key.split("/", 1)[1].strip()
if suffix:
roles.append(suffix)
elif val:
roles.append(str(val).strip())
return tuple(sorted({r for r in roles if r}))
def _node_ready(status: dict) -> bool:
for cond in status.get("conditions") or []:
if str(cond.get("type") or "") == "Ready":
return str(cond.get("status") or "").lower() == "true"
return False
def _node_internal_ip(status: dict) -> str:
for addr in status.get("addresses") or []:
if str(addr.get("type") or "") == "InternalIP":
return str(addr.get("address") or "").strip()
return ""
def _normalize_taints(spec: dict) -> tuple[str, ...]:
out: list[str] = []
for taint in spec.get("taints") or []:
key = str(taint.get("key") or "").strip()
value = str(taint.get("value") or "").strip()
effect = str(taint.get("effect") or "").strip()
if not key:
continue
token = key
if value:
token += f"={value}"
if effect:
token += f":{effect}"
out.append(token)
return tuple(sorted(out))
def _parse_k8s_nodes(nodes_payload: dict) -> dict[str, dict]:
out: dict[str, dict] = {}
for item in nodes_payload.get("items") or []:
metadata = item.get("metadata") or {}
status = item.get("status") or {}
labels = {
str(k): str(v)
for k, v in (metadata.get("labels") or {}).items()
if str(k).strip()
}
name = str(metadata.get("name") or "").strip()
if not name:
continue
alloc = status.get("allocatable") or {}
out[name] = {
"name": name,
"roles": _node_roles(labels),
"labels": labels,
"taints": _normalize_taints(item.get("spec") or {}),
"ready": _node_ready(status),
"internal_ip": _node_internal_ip(status),
"architecture": str(
labels.get("kubernetes.io/arch")
or labels.get("beta.kubernetes.io/arch")
or ""
).strip(),
"operating_system": str(
labels.get("kubernetes.io/os")
or labels.get("beta.kubernetes.io/os")
or ""
).strip(),
"allocatable_cpu": str(alloc.get("cpu") or "").strip(),
"allocatable_memory_bytes": _parse_k8s_quantity_bytes(
str(alloc.get("memory") or "")
),
"allocatable_ephemeral_storage_bytes": _parse_k8s_quantity_bytes(
str(alloc.get("ephemeral-storage") or "")
),
}
return out
def _collector_python_script() -> str:
return r'''
import json
import os
import platform
import re
import time
from pathlib import Path
BEGIN = "PROLE_NODE_INVENTORY_BEGIN"
END = "PROLE_NODE_INVENTORY_END"
_PSEUDO_FS = {
"proc", "sysfs", "tmpfs", "cgroup", "cgroup2", "devpts", "devtmpfs", "mqueue",
"tracefs", "pstore", "securityfs", "debugfs", "ramfs", "overlay", "squashfs", "autofs"
}
_NETWORK_FS = {"nfs", "nfs4", "cifs", "smb", "glusterfs", "ceph", "cephfs", "fuse.sshfs"}
def _read_text(path: Path) -> str:
try:
return path.read_text(encoding="utf-8", errors="replace")
except Exception:
return ""
def _decode_mount_path(path: str) -> str:
def repl(match):
try:
return chr(int(match.group(1), 8))
except Exception:
return match.group(0)
return re.sub(r"\\([0-7]{3})", repl, path)
def _device_block_name(device: str) -> str:
if not device.startswith("/dev/"):
return ""
name = device.split("/")[-1]
if name.startswith("nvme") and "p" in name:
return name.split("p", 1)[0]
m = re.match(r"^(.*?)(\d+)$", name)
if m:
return m.group(1)
return name
def _read_rotational(block_name: str):
if not block_name:
return None
p = Path("/host/sys/block") / block_name / "queue/rotational"
text = _read_text(p).strip()
if text == "0":
return False
if text == "1":
return True
return None
def _read_transport(device: str, fstype: str, block_name: str) -> str:
if fstype in _NETWORK_FS:
return "network"
if "nvme" in device or block_name.startswith("nvme"):
return "nvme"
if not block_name:
return "unknown"
if block_name.startswith("sd"):
return "sata"
if block_name.startswith("vd"):
return "virtio"
if block_name.startswith("xvd"):
return "virtio"
return "unknown"
def _infer_storage_class(device: str, fstype: str, transport: str, rotational):
if fstype in _NETWORK_FS:
return "network"
if transport == "nvme" or "nvme" in device:
return "nvme"
if rotational is False:
return "ssd"
if rotational is True:
return "hdd"
return "unknown"
def _cpu_facts():
cpuinfo = _read_text(Path("/host/proc/cpuinfo"))
model = ""
logical = 0
phys_cores = set()
current_phys = ""
current_core = ""
for raw_line in cpuinfo.splitlines():
line = raw_line.strip()
if not line:
if current_phys or current_core:
phys_cores.add((current_phys, current_core))
current_phys = ""
current_core = ""
continue
if line.startswith("model name") and not model:
parts = line.split(":", 1)
if len(parts) == 2:
model = parts[1].strip()
if line.startswith("processor"):
logical += 1
if line.startswith("physical id"):
parts = line.split(":", 1)
if len(parts) == 2:
current_phys = parts[1].strip()
if line.startswith("core id"):
parts = line.split(":", 1)
if len(parts) == 2:
current_core = parts[1].strip()
if current_phys or current_core:
phys_cores.add((current_phys, current_core))
physical = len([x for x in phys_cores if x != ("", "")])
if physical <= 0 and logical > 0:
physical = logical
return model or None, physical or None, logical or None
def _memory_bytes():
meminfo = _read_text(Path("/host/proc/meminfo"))
for line in meminfo.splitlines():
if not line.startswith("MemTotal:"):
continue
parts = line.split()
if len(parts) < 2:
continue
try:
return int(parts[1]) * 1024
except Exception:
return None
return None
def _mount_rows():
mounts_text = _read_text(Path("/host/proc/1/mounts")) or _read_text(Path("/host/proc/mounts"))
rows = []
seen = set()
for line in mounts_text.splitlines():
parts = line.split()
if len(parts) < 3:
continue
device = parts[0]
mountpoint = _decode_mount_path(parts[1])
fstype = parts[2]
if fstype in _PSEUDO_FS:
continue
key = (mountpoint, device, fstype)
if key in seen:
continue
seen.add(key)
rows.append((device, mountpoint, fstype))
rows.sort(key=lambda x: x[1])
return rows
def _fs_stats(mountpoint: str):
host_mount = Path("/host") / mountpoint.lstrip("/")
try:
st = os.statvfs(str(host_mount))
size = int(st.f_frsize * st.f_blocks)
avail = int(st.f_frsize * st.f_bavail)
used = max(0, size - int(st.f_frsize * st.f_bfree))
return size, used, avail
except Exception:
return None, None, None
def _collect_filesystems():
out = []
for device, mountpoint, fstype in _mount_rows():
size, used, available = _fs_stats(mountpoint)
block_name = _device_block_name(device)
rotational = _read_rotational(block_name)
transport = _read_transport(device, fstype, block_name)
storage_class = _infer_storage_class(device, fstype, transport, rotational)
out.append({
"device": device,
"mountpoint": mountpoint,
"fstype": fstype,
"sizeBytes": size,
"usedBytes": used,
"availableBytes": available,
"transport": transport,
"rotational": rotational,
"storageClass": storage_class,
})
return out
def main():
cpu_model, physical_cores, logical_cores = _cpu_facts()
payload = {
"nodeName": os.environ.get("NODE_NAME", ""),
"reportedAt": int(time.time()),
"architecture": platform.machine(),
"cpuModel": cpu_model,
"physicalCores": physical_cores,
"logicalCores": logical_cores,
"memoryBytes": _memory_bytes(),
"filesystems": _collect_filesystems(),
}
print(BEGIN, flush=True)
print(json.dumps(payload, separators=(",", ":"), sort_keys=True), flush=True)
print(END, flush=True)
# Keep pod alive so logs can be collected exactly once without restart loops.
time.sleep(3600)
if __name__ == "__main__":
main()
'''.strip()
def _collector_manifest(name: str, *, image: str) -> dict:
labels = {
"app.kubernetes.io/name": "prole-topology-collector",
"prole.io/topology-collector": name,
}
return {
"apiVersion": "apps/v1",
"kind": "DaemonSet",
"metadata": {
"name": name,
"namespace": TOPOLOGY_NAMESPACE,
"labels": labels,
},
"spec": {
"selector": {"matchLabels": labels},
"template": {
"metadata": {"labels": labels},
"spec": {
"serviceAccountName": "default",
"terminationGracePeriodSeconds": 0,
"tolerations": [{"operator": "Exists"}],
"containers": [
{
"name": "collector",
"image": image,
"imagePullPolicy": "IfNotPresent",
"command": ["python", "-c", _collector_python_script()],
"env": [
{
"name": "NODE_NAME",
"valueFrom": {
"fieldRef": {"fieldPath": "spec.nodeName"}
},
}
],
"securityContext": {"privileged": True},
"volumeMounts": [
{
"name": "host-root",
"mountPath": "/host",
"readOnly": True,
}
],
}
],
"volumes": [{"name": "host-root", "hostPath": {"path": "/"}}],
},
},
"updateStrategy": {"type": "RollingUpdate"},
},
}
def _extract_inventory_from_log(log_text: str) -> dict | None:
text = str(log_text or "")
begin = text.find(INVENTORY_BEGIN_MARKER)
end = text.find(INVENTORY_END_MARKER)
if begin < 0 or end < 0 or end <= begin:
return None
payload = text[begin + len(INVENTORY_BEGIN_MARKER) : end].strip()
if not payload:
return None
try:
data = json.loads(payload)
except json.JSONDecodeError:
return None
return data if isinstance(data, dict) else None
def cleanup_topology_collectors(
*,
kubectl_base_cmd: Sequence[str],
namespace: str = TOPOLOGY_NAMESPACE,
env: dict[str, str] | None = None,
collector_name: str | None = None,
) -> None:
selector = (
f"prole.io/topology-collector={collector_name}"
if collector_name
else "app.kubernetes.io/name=prole-topology-collector"
)
for kind in ("daemonset", "pod"):
_run_kubectl(
kubectl_base_cmd,
[
"-n",
namespace,
"delete",
kind,
"-l",
selector,
"--ignore-not-found=true",
"--wait=false",
],
timeout=15,
env=env,
)
def _merge_node(
k8s_node: dict,
inventory: dict | None,
*,
required_monitoring_mountpoints: tuple[str, ...],
monitoring_min_free_bytes: int,
inventory_error: str | None = None,
storage_items: Sequence[StorageProbeResult] = (),
assume_eligible: bool = False,
) -> NodeTopology:
filesystems: list[FilesystemTopology] = []
if isinstance(inventory, dict):
for item in inventory.get("filesystems") or []:
if not isinstance(item, dict):
continue
filesystems.append(
FilesystemTopology(
device=str(item.get("device") or "").strip(),
mountpoint=str(item.get("mountpoint") or "").strip(),
fstype=str(item.get("fstype") or "").strip(),
size_bytes=item.get("sizeBytes")
if isinstance(item.get("sizeBytes"), int)
else None,
used_bytes=item.get("usedBytes")
if isinstance(item.get("usedBytes"), int)
else None,
available_bytes=item.get("availableBytes")
if isinstance(item.get("availableBytes"), int)
else None,
transport=str(item.get("transport") or "unknown").strip() or "unknown",
rotational=item.get("rotational")
if isinstance(item.get("rotational"), bool)
else None,
storage_class=str(item.get("storageClass") or "unknown").strip()
or "unknown",
)
)
filesystems.sort(key=lambda fs: fs.mountpoint)
host_inventory_available = bool(inventory)
mountpoints = {fs.mountpoint for fs in filesystems if fs.mountpoint}
matching_mountpoints_set = {
mp for mp in required_monitoring_mountpoints if mp in mountpoints
}
if not matching_mountpoints_set and any(
_SYNOLOGY_DATA_MOUNT_RE.fullmatch(mp)
for mp in required_monitoring_mountpoints
):
matching_mountpoints_set = {
mountpoint
for mountpoint in mountpoints
if _SYNOLOGY_DATA_MOUNT_RE.fullmatch(mountpoint)
}
matching_mountpoints = tuple(sorted(matching_mountpoints_set))
reasons: set[str] = set()
has_synology_root: bool | None
has_synology_root_state: str
monitoring_eligible: bool | None
monitoring_eligible_state: str
cnpg_eligible: bool | None
cnpg_eligible_state: str
if assume_eligible and k8s_node["ready"]:
# GKE managed nodes: skip taint/mount checks; treat all ready nodes as eligible.
has_synology_root = None
has_synology_root_state = "assumed"
monitoring_eligible = True
monitoring_eligible_state = "assumed"
cnpg_eligible = True
cnpg_eligible_state = "assumed"
host_inventory_available = True
elif not host_inventory_available:
has_synology_root = None
has_synology_root_state = "unavailable"
monitoring_eligible = None
monitoring_eligible_state = "unavailable"
cnpg_eligible = None
cnpg_eligible_state = "unavailable"
reasons.add("topology_unavailable")
else:
has_synology_root = bool(
any(fs.mountpoint.startswith("/synology/") for fs in filesystems if fs.mountpoint)
)
has_synology_root_state = "discovered"
if not k8s_node["ready"]:
monitoring_eligible = False
monitoring_eligible_state = "inferred"
reasons.add("node_not_ready")
elif not matching_mountpoints:
monitoring_eligible = False
monitoring_eligible_state = "inferred"
reasons.add("missing_required_mountpoint")
else:
max_available = max(
(
fs.available_bytes
for fs in filesystems
if fs.mountpoint in matching_mountpoints and fs.available_bytes is not None
),
default=None,
)
if max_available is None:
monitoring_eligible = None
monitoring_eligible_state = "unavailable"
reasons.add("topology_unavailable")
elif max_available < monitoring_min_free_bytes:
monitoring_eligible = False
monitoring_eligible_state = "inferred"
reasons.add("insufficient_free_space")
else:
monitoring_eligible = True
monitoring_eligible_state = "inferred"
if not k8s_node["ready"]:
cnpg_eligible = False
cnpg_eligible_state = "inferred"
reasons.add("node_not_ready")
elif has_synology_root is True:
cnpg_eligible = True
cnpg_eligible_state = "inferred"
elif has_synology_root is False:
cnpg_eligible = False
cnpg_eligible_state = "inferred"
reasons.add("missing_required_mountpoint")
else:
cnpg_eligible = None
cnpg_eligible_state = "unavailable"
reasons.add("topology_unavailable")
if inventory_error:
reasons.add("inventory_parse_error")
capabilities = CapabilityFlags(
has_synology_root=has_synology_root,
matching_mountpoints=matching_mountpoints,
monitoring_eligible=monitoring_eligible,
cnpg_eligible=cnpg_eligible,
reasons=tuple(sorted(reasons)),
has_synology_root_state=has_synology_root_state,
monitoring_eligible_state=monitoring_eligible_state,
cnpg_eligible_state=cnpg_eligible_state,
)
return NodeTopology(
name=str(k8s_node["name"]),
roles=tuple(k8s_node["roles"]),
labels=dict(sorted(k8s_node["labels"].items())),
taints=tuple(k8s_node["taints"]),
ready=bool(k8s_node["ready"]),
internal_ip=str(k8s_node["internal_ip"]),
architecture=str(k8s_node["architecture"]),
operating_system=str(k8s_node["operating_system"]),
allocatable_cpu=str(k8s_node["allocatable_cpu"]),
allocatable_memory_bytes=k8s_node["allocatable_memory_bytes"],
allocatable_ephemeral_storage_bytes=k8s_node["allocatable_ephemeral_storage_bytes"],
cpu_model=(str(inventory.get("cpuModel") or "").strip() if inventory else None) or None,
physical_cores=(inventory.get("physicalCores") if inventory else None)
if isinstance((inventory or {}).get("physicalCores"), int)
else None,
logical_cores=(inventory.get("logicalCores") if inventory else None)
if isinstance((inventory or {}).get("logicalCores"), int)
else None,
memory_bytes=(inventory.get("memoryBytes") if inventory else None)
if isinstance((inventory or {}).get("memoryBytes"), int)
else None,
filesystems=tuple(filesystems),
storage_items=tuple(storage_items),
host_inventory_available=host_inventory_available,
host_inventory_error=inventory_error,
capabilities=capabilities,
)
def _render_xml_bytes(value: int | None) -> str:
return "" if value is None else str(int(value))
def _render_xml_float(value: float | None) -> str:
return "" if value is None else f"{float(value):.3f}"
def _set_text(parent: ET.Element, tag: str, text: str) -> ET.Element:
elem = ET.SubElement(parent, tag)
elem.text = text
return elem
def _append_bool(parent: ET.Element, tag: str, value: bool | None) -> None:
_set_text(parent, tag, "unknown" if value is None else ("true" if value else "false"))
def _node_to_xml(node: NodeTopology) -> ET.Element:
root = ET.Element("nodeTopology", {"version": TOPOLOGY_SCHEMA_VERSION})
identity = ET.SubElement(root, "identity")
_set_text(identity, "name", node.name)
_set_text(identity, "internalIp", node.internal_ip)
kubernetes = ET.SubElement(root, "kubernetes")
_set_text(kubernetes, "ready", "true" if node.ready else "false")
_set_text(kubernetes, "roles", ",".join(node.roles))
_set_text(kubernetes, "architecture", node.architecture)
_set_text(kubernetes, "operatingSystem", node.operating_system)
_set_text(kubernetes, "allocatableCpu", node.allocatable_cpu)
_set_text(kubernetes, "allocatableMemoryBytes", _render_xml_bytes(node.allocatable_memory_bytes))
_set_text(
kubernetes,
"allocatableEphemeralStorageBytes",
_render_xml_bytes(node.allocatable_ephemeral_storage_bytes),
)
labels = ET.SubElement(kubernetes, "labels")
for key, val in sorted(node.labels.items()):
label = ET.SubElement(labels, "label", {"key": key})
label.text = val
taints = ET.SubElement(kubernetes, "taints")
for taint in node.taints:
_set_text(taints, "taint", taint)
hardware = ET.SubElement(root, "hardware")
_set_text(hardware, "cpuModel", node.cpu_model or "")
_set_text(hardware, "physicalCores", _render_xml_bytes(node.physical_cores))
_set_text(hardware, "logicalCores", _render_xml_bytes(node.logical_cores))
_set_text(hardware, "memoryBytes", _render_xml_bytes(node.memory_bytes))
filesystems = ET.SubElement(root, "filesystems")
for fs in node.filesystems:
elem = ET.SubElement(filesystems, "filesystem")
_set_text(elem, "device", fs.device)
_set_text(elem, "mountpoint", fs.mountpoint)
_set_text(elem, "fstype", fs.fstype)
_set_text(elem, "sizeBytes", _render_xml_bytes(fs.size_bytes))
_set_text(elem, "usedBytes", _render_xml_bytes(fs.used_bytes))
_set_text(elem, "availableBytes", _render_xml_bytes(fs.available_bytes))
_set_text(elem, "transport", fs.transport)
_set_text(
elem,
"rotational",
"unknown"
if fs.rotational is None
else ("true" if fs.rotational else "false"),
)
_set_text(elem, "storageClass", fs.storage_class)
storage = ET.SubElement(root, "storage")
for item in node.storage_items:
elem = ET.SubElement(storage, "item")
_set_text(elem, "targetId", item.target_id)
_set_text(elem, "label", item.label)
_set_text(elem, "mountpoint", item.mountpoint)
_set_text(elem, "sourceType", item.source_type)
_set_text(elem, "storageClass", item.storage_class or "")
_set_text(elem, "succeeded", "true" if item.succeeded else "false")
_set_text(elem, "skipped", "true" if item.skipped else "false")
_set_text(elem, "cached", "true" if item.cached else "false")
_set_text(elem, "approximate", "true" if item.approximate else "false")
_set_text(elem, "error", item.error or "")
_set_text(elem, "ioClass", item.io_class or "")
_set_text(elem, "timingMultiplier", _render_xml_float(item.timing_multiplier))
_set_text(elem, "observedAt", item.observed_at or "")
metrics = ET.SubElement(elem, "metrics")
_set_text(metrics, "randreadIops", _render_xml_float(item.metrics.randread_iops))
_set_text(metrics, "randwriteIops", _render_xml_float(item.metrics.randwrite_iops))
_set_text(metrics, "readBwBytes", _render_xml_bytes(item.metrics.read_bw_bytes))
_set_text(metrics, "writeBwBytes", _render_xml_bytes(item.metrics.write_bw_bytes))
_set_text(metrics, "readLatMsP95", _render_xml_float(item.metrics.read_lat_ms_p95))
_set_text(metrics, "writeLatMsP95", _render_xml_float(item.metrics.write_lat_ms_p95))
_set_text(metrics, "testRuntimeSec", _render_xml_float(item.metrics.test_runtime_sec))
notes = ET.SubElement(elem, "notes")
for note in item.notes:
_set_text(notes, "note", note)
capabilities = ET.SubElement(root, "capabilities")
_append_bool(capabilities, "hasSynologyRoot", node.capabilities.has_synology_root)
_set_text(capabilities, "hasSynologyRootState", node.capabilities.has_synology_root_state)
_set_text(
capabilities,
"matchingMountpoints",
",".join(node.capabilities.matching_mountpoints),
)
_append_bool(capabilities, "monitoringEligible", node.capabilities.monitoring_eligible)
_set_text(
capabilities,
"monitoringEligibleState",
node.capabilities.monitoring_eligible_state,
)
_append_bool(capabilities, "cnpgEligible", node.capabilities.cnpg_eligible)
_set_text(capabilities, "cnpgEligibleState", node.capabilities.cnpg_eligible_state)
reasons = ET.SubElement(root, "reasons")
for reason in node.capabilities.reasons:
_set_text(reasons, "reason", reason)
if node.host_inventory_error:
_set_text(reasons, "reason", node.host_inventory_error)
inventory = ET.SubElement(root, "hostInventory")
_set_text(inventory, "available", "true" if node.host_inventory_available else "false")
_set_text(inventory, "error", node.host_inventory_error or "")
return root
def _cluster_to_xml(cluster: ClusterTopology) -> ET.Element:
root = ET.Element("clusterTopology", {"version": TOPOLOGY_SCHEMA_VERSION})
identity = ET.SubElement(root, "identity")
_set_text(identity, "generatedAt", cluster.generated_at)
_set_text(identity, "generatedEpochMs", str(cluster.generated_epoch_ms))
_set_text(identity, "collectionStatus", cluster.collection_status)
_set_text(identity, "topologyRoot", cluster.topology_root)
summary = ET.SubElement(root, "summary")
_set_text(summary, "discoveredNodes", str(cluster.discovered_nodes))
_set_text(summary, "readyNodes", str(cluster.ready_nodes))
_set_text(summary, "monitoringEligibleNodes", str(cluster.monitoring_eligible_nodes))
_set_text(summary, "cnpgEligibleNodes", str(cluster.cnpg_eligible_nodes))
required = ET.SubElement(root, "requiredMountpoints")
for mp in cluster.required_monitoring_mountpoints:
_set_text(required, "mountpoint", mp)
_set_text(required, "monitoringMinFreeBytes", str(cluster.monitoring_min_free_bytes))
nodes = ET.SubElement(root, "nodes")
for node in cluster.nodes:
ref = ET.SubElement(nodes, "nodeRef")
_set_text(ref, "name", node.name)
_set_text(ref, "path", f"nodes/{node.name}.xml")
messages = ET.SubElement(root, "messages")
for msg in cluster.messages:
_set_text(messages, "message", msg)
return root
def _write_xml(path: Path, element: ET.Element) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
tree = ET.ElementTree(element)
ET.indent(tree, space=" ")
tree.write(path, encoding="utf-8", xml_declaration=True)
def save_topology_xml(cluster: ClusterTopology, output_root: Path) -> tuple[Path, dict[str, Path]]:
root = Path(output_root)
nodes_dir = root / "nodes"
cluster_path = root / "cluster.xml"
written_nodes: dict[str, Path] = {}
for node in cluster.nodes:
node_path = nodes_dir / f"{node.name}.xml"
_write_xml(node_path, _node_to_xml(node))
written_nodes[node.name] = node_path
_write_xml(cluster_path, _cluster_to_xml(cluster))
return cluster_path, written_nodes
def discover_cluster_topology(
*,
kubectl_base_cmd: Sequence[str],
topology_root: Path,
config: TopologyDiscoveryConfig | None = None,
progress_cb: Callable[[str], None] | None = None,
env: dict[str, str] | None = None,
) -> TopologyDiscoveryResult:
cfg = config or TopologyDiscoveryConfig()
required_mountpoints = tuple(cfg.required_monitoring_mountpoints)
monitoring_min_free_bytes = int(cfg.monitoring_min_free_bytes)
messages: list[str] = []
def _progress(message: str) -> None:
messages.append(message)
if progress_cb:
try:
progress_cb(message)
except Exception:
pass
_progress("Starting topology discovery")
_progress("Discovering Kubernetes nodes...")
try:
nodes_payload = _kubectl_json(kubectl_base_cmd, ["get", "nodes"], env=env)
k8s_nodes = _parse_k8s_nodes(nodes_payload)
except Exception as exc:
now = datetime.now(timezone.utc)
cluster = ClusterTopology(
version=TOPOLOGY_SCHEMA_VERSION,
generated_at=now.isoformat(),
generated_epoch_ms=int(now.timestamp() * 1000),
collection_status="failed",
required_monitoring_mountpoints=required_mountpoints,
monitoring_min_free_bytes=monitoring_min_free_bytes,
nodes=tuple(),
topology_root=str(topology_root),
messages=tuple(messages + [f"Topology discovery failed to query nodes: {exc}"]),
)
return TopologyDiscoveryResult(
topology=cluster,
collector_applied=False,
expected_ready_nodes=tuple(),
reported_nodes=tuple(),
missing_nodes=tuple(),
)
discovered = len(k8s_nodes)
ready_nodes = sorted(name for name, data in k8s_nodes.items() if data["ready"])
_progress(f"Discovered {discovered} Kubernetes nodes ({len(ready_nodes)} Ready)")
collector_name = f"prole-topology-{secrets.token_hex(4)}"
inventory_payloads: dict[str, dict] = {}
inventory_errors: dict[str, str] = {}
collector_applied = False
if ready_nodes and not cfg.skip_collector:
_progress("Creating temporary in-cluster topology collector")
manifest = _collector_manifest(collector_name, image=cfg.collector_image)
res_apply = _run_kubectl(
kubectl_base_cmd,
["-n", cfg.namespace, "apply", "-f", "-"],
timeout=20,
env=env,
input_text=json.dumps(manifest, separators=(",", ":"), sort_keys=True),
)
collector_applied = res_apply.returncode == 0
if not collector_applied:
stderr = (res_apply.stderr or "").strip() or "unknown_error"
_progress(f"Topology collector deployment failed: {stderr}")
else:
_progress(f"Waiting for inventory reports from Ready nodes ({len(ready_nodes)} total)")
deadline = time.monotonic() + max(5, int(cfg.timeout_seconds))
collected_for: set[str] = set()
while time.monotonic() < deadline and len(collected_for) < len(ready_nodes):
try:
pods_payload = _kubectl_json(
kubectl_base_cmd,
[
"-n",
cfg.namespace,
"get",
"pods",
"-l",
f"prole.io/topology-collector={collector_name}",
],
timeout=15,
env=env,
)
except Exception:
time.sleep(max(0.5, cfg.poll_interval_seconds))
continue
pod_by_node: dict[str, str] = {}
for item in pods_payload.get("items") or []:
spec = item.get("spec") or {}
metadata = item.get("metadata") or {}
pod_name = str(metadata.get("name") or "").strip()
node_name = str(spec.get("nodeName") or "").strip()
if pod_name and node_name:
pod_by_node[node_name] = pod_name
for node_name in ready_nodes:
if node_name in collected_for:
continue
pod_name = pod_by_node.get(node_name)
if not pod_name:
continue
res_logs = _run_kubectl(
kubectl_base_cmd,
["-n", cfg.namespace, "logs", pod_name, "--tail=-1"],
timeout=15,
env=env,
)
if res_logs.returncode != 0:
continue
payload = _extract_inventory_from_log(res_logs.stdout or "")
if not payload:
continue
inventory_payloads[node_name] = payload
collected_for.add(node_name)
_progress(f"Inventory received from node {node_name}")
if len(collected_for) >= len(ready_nodes):
break
parsed = len(collected_for)
_progress(f"Parsed topology from {parsed} of {len(ready_nodes)} Ready nodes")
time.sleep(max(0.5, cfg.poll_interval_seconds))
missing = sorted(set(ready_nodes) - set(inventory_payloads.keys()))
for node_name in missing:
inventory_errors[node_name] = "timed_out"
_progress(f"Timed out waiting for node {node_name}")
if collector_applied:
cleanup_topology_collectors(
kubectl_base_cmd=kubectl_base_cmd,
namespace=cfg.namespace,
env=env,
collector_name=collector_name,
)
merged_nodes: list[NodeTopology] = []
storage_cache_path = cfg.storage_probe_cache_path
if storage_cache_path is None:
storage_cache_path = Path(topology_root).parent / "cache" / "storage_probe_cache.json"
for node_name in sorted(k8s_nodes.keys()):
inv = inventory_payloads.get(node_name)
error = inventory_errors.get(node_name)
node = _merge_node(
k8s_nodes[node_name],
inv,
required_monitoring_mountpoints=required_mountpoints,
monitoring_min_free_bytes=monitoring_min_free_bytes,
inventory_error=error,
assume_eligible=cfg.assume_nodes_eligible,
)
if cfg.storage_probe_enabled:
storage_inventory = probe_node_storage(
node,
topology_context={
"required_monitoring_mountpoints": required_mountpoints,
"monitoring_min_free_bytes": monitoring_min_free_bytes,
},
timeout_sec=max(3, int(cfg.storage_probe_timeout_sec)),
quick=bool(cfg.storage_probe_quick),
force_refresh=bool(cfg.storage_probe_force_refresh),
cache_path=storage_cache_path,
cache_ttl_sec=max(0, int(cfg.storage_probe_cache_ttl_sec)),
max_targets=max(0, int(cfg.storage_probe_max_targets_per_node)),
)
node = replace(node, storage_items=tuple(storage_inventory.items))
_progress(
f"Storage probe for {node_name}: {len(storage_inventory.items)} target(s)"
)
merged_nodes.append(node)
if merged_nodes and all(not node.host_inventory_available for node in merged_nodes):
status = "failed"
elif any(not node.host_inventory_available for node in merged_nodes):
status = "partial"
else:
status = "complete"
if status == "partial":
missing = [node.name for node in merged_nodes if not node.host_inventory_available and node.ready]
if missing:
_progress(
"Continuing with partial topology because "
f"{len(missing)} node(s) did not report in time"
)
elif status == "failed":
_progress("Topology collection unavailable, continuing with Kubernetes-only node facts")
_progress("Writing topology XML files")
now = datetime.now(timezone.utc)
cluster = ClusterTopology(
version=TOPOLOGY_SCHEMA_VERSION,
generated_at=now.isoformat(),
generated_epoch_ms=int(now.timestamp() * 1000),
collection_status=status,
required_monitoring_mountpoints=required_mountpoints,
monitoring_min_free_bytes=monitoring_min_free_bytes,
nodes=tuple(merged_nodes),
topology_root=str(topology_root),
messages=tuple(messages),
)
save_topology_xml(cluster, topology_root)
_progress(f"Topology snapshot saved to {topology_root}")
cluster = ClusterTopology(
version=cluster.version,
generated_at=cluster.generated_at,
generated_epoch_ms=cluster.generated_epoch_ms,
collection_status=cluster.collection_status,
required_monitoring_mountpoints=cluster.required_monitoring_mountpoints,
monitoring_min_free_bytes=cluster.monitoring_min_free_bytes,
nodes=cluster.nodes,
topology_root=cluster.topology_root,
messages=tuple(messages),
)
return TopologyDiscoveryResult(
topology=cluster,
collector_applied=collector_applied,
expected_ready_nodes=tuple(ready_nodes),
reported_nodes=tuple(sorted(inventory_payloads.keys())),
missing_nodes=tuple(sorted(set(ready_nodes) - set(inventory_payloads.keys()))),
)