mirror of
https://github.com/dredx/prole.git
synced 2026-09-23 11:03:59 +00:00
1224 lines
43 KiB
Python
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": "knoe-topology-collector",
|
|
"knoe.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"knoe.io/topology-collector={collector_name}"
|
|
if collector_name
|
|
else "app.kubernetes.io/name=knoe-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"knoe-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"knoe.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()))),
|
|
)
|