#!/usr/bin/env python3 """Read cluster health metadata and compare settling windows; never change state. Snapshots contain resource names, readiness, restart counts and warning reasons. They exclude pod environments, Secret values, log bodies and event messages. Run once after a repair, then again after 10-15 minutes with --previous. """ from __future__ import annotations import argparse from concurrent.futures import ThreadPoolExecutor from datetime import datetime, timezone import json from pathlib import Path import subprocess QUERIES = { "nodes": ["nodes"], "pods": ["pods", "-A"], "deployments": ["deployments", "-A"], "statefulsets": ["statefulsets", "-A"], "flux": ["kustomizations.kustomize.toolkit.fluxcd.io", "-A"], "helm": ["helmreleases.helm.toolkit.fluxcd.io", "-A"], "volumes": ["volumes.longhorn.io", "-n", "longhorn-system"], "warnings": ["events", "-A", "--field-selector=type=Warning"], "preemptions": ["events", "-A", "--field-selector=reason=Preempted"], } def fetch(item: tuple[str, list[str]]) -> tuple[str, list[dict]]: """Fetch one resource list; fail without echoing API error bodies.""" name, args = item result = subprocess.run( ["kubectl", "--request-timeout=20s", "get", *args, "-o", "json"], capture_output=True, text=True, timeout=30, check=False, ) if result.returncode: raise RuntimeError(f"Cannot read {name}; kubectl exit {result.returncode}") return name, json.loads(result.stdout)["items"] def identity(resource: dict) -> str: """Return the namespace/name without annotations or source content.""" meta = resource["metadata"] return "/".join(filter(None, [meta.get("namespace"), meta["name"]])) def conditions(resource: dict) -> dict: """Project condition booleans, omitting potentially content-bearing messages.""" return {c["type"]: c["status"] for c in resource.get("status", {}).get("conditions", [])} def snapshot() -> dict: """Collect a complete metadata snapshot or fail; never report partial health.""" with ThreadPoolExecutor(max_workers=4) as pool: raw = dict(pool.map(fetch, QUERIES.items())) result = {"observed_at": datetime.now(timezone.utc).isoformat()} result["nodes"] = { identity(n): {"conditions": conditions(n), "cordoned": n["spec"].get("unschedulable", False)} for n in raw["nodes"] } result["pods"] = {} for p in raw["pods"]: owners = p["metadata"].get("ownerReferences", []) # Completed Jobs and routine CronJob pod turnover are not service churn. if (any(o["kind"] == "Job" for o in owners) or p["status"]["phase"] in ("Succeeded", "Failed")): continue statuses = p["status"].get("containerStatuses", []) result["pods"][identity(p)] = { "uid": p["metadata"]["uid"], "node": p["spec"].get("nodeName"), "owner": [(o["kind"], o["name"]) for o in owners], "phase": p["status"]["phase"], "ready": conditions(p).get("Ready"), "deleting": p["metadata"].get("deletionTimestamp"), "restarts": {s["name"]: s.get("restartCount", 0) for s in statuses}, "waiting": {s["name"]: s["state"]["waiting"].get("reason") for s in statuses if "waiting" in s.get("state", {})}, } result["workloads"] = { kind + "/" + identity(w): { "desired": w["spec"].get("replicas", 1), "ready": w["status"].get("readyReplicas", 0), "generation": w["metadata"]["generation"], "observed_generation": w["status"].get("observedGeneration"), "rollout_failed": any( c.get("type") == "Progressing" and c.get("status") == "False" and c.get("reason") == "ProgressDeadlineExceeded" for c in w["status"].get("conditions", []) ), } for kind in ("deployments", "statefulsets") for w in raw[kind] } for kind in ("flux", "helm"): result[kind] = { identity(w): {"suspended": w["spec"].get("suspend", False), "ready": conditions(w).get("Ready")} for w in raw[kind] } result["volumes"] = { identity(v): {"state": v["status"].get("state"), "robustness": v["status"].get("robustness"), "requested_node": v["spec"].get("nodeID")} for v in raw["volumes"] } result["warnings"] = { e["metadata"]["uid"]: { "object": identity({"metadata": e["involvedObject"]}), "reason": e.get("reason"), "count": e.get("series", {}).get("count", e.get("count")) or 1, "last_observed": e.get("series", {}).get("lastObservedTime") or e.get("lastTimestamp") or e.get("eventTime"), } # Scheduler preemption can be a Normal event despite interrupting work. for e in [*raw["warnings"], *raw["preemptions"]] } return result def summarize(current: dict, previous: dict | None = None) -> dict: """Report current failures and deltas without hiding known offline-node noise.""" offline = {n for n, v in current["nodes"].items() if v["conditions"].get("Ready") != "True"} out = { "observed_at": current["observed_at"], "offline_nodes": sorted(offline), "pressure_nodes": {n: v["conditions"] for n, v in current["nodes"].items() if any(v["conditions"].get(k) == "True" for k in ("MemoryPressure", "DiskPressure", "PIDPressure"))}, "unready_workloads": {n: v for n, v in current["workloads"].items() if v["ready"] < v["desired"] or v.get("rollout_failed", False) or (v.get("observed_generation") or 0) < v["generation"]}, "unready_pods_on_reachable_nodes": [n for n, v in current["pods"].items() if (v["ready"] != "True" or v.get("deleting")) and v["node"] not in offline], "unready_pods_on_offline_nodes": sum(bool(v["ready"] != "True" or v.get("deleting")) and v["node"] in offline for v in current["pods"].values()), "unhealthy_attached_volumes": {n: v for n, v in current["volumes"].items() if v["state"] == "attached" and v["robustness"] != "healthy"}, "unavailable_requested_volumes": {n: v for n, v in current["volumes"].items() if v.get("requested_node") and v["state"] != "attached"}, } for kind in ("flux", "helm"): out[kind + "_not_ready"] = [n for n, v in current[kind].items() if not v["suspended"] and v["ready"] != "True"] if previous: out["window_seconds"] = round((datetime.fromisoformat(current["observed_at"]) - datetime.fromisoformat(previous["observed_at"])).total_seconds(), 1) out["new_service_pods"] = sorted(set(current["pods"]) - set(previous["pods"])) out["removed_service_pods"] = sorted(set(previous["pods"]) - set(current["pods"])) out["replaced_service_pods"] = sorted( name for name in set(current["pods"]) & set(previous["pods"]) if current["pods"][name]["uid"] != previous["pods"][name]["uid"] ) out["restart_increases"] = {} for name, pod in current["pods"].items(): old = previous["pods"].get(name, {}) if old.get("uid") == pod["uid"]: delta = sum(max(0, count - old["restarts"].get(c, 0)) for c, count in pod["restarts"].items()) if delta: out["restart_increases"][name] = delta # A newly observed event can contain an older aggregated lifetime count. # Only matching event UIDs establish a measured counter increase. out["warning_changes"] = [] for key, warning in current["warnings"].items(): old = previous["warnings"].get(key) if old and warning["count"] > old["count"]: out["warning_changes"].append(dict(warning, increase=warning["count"] - old["count"])) elif not old: out["warning_changes"].append(dict(warning, increase=None)) return out def main() -> None: """Write the requested snapshot and print its safe summary/comparison.""" parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--output", required=True, type=Path) parser.add_argument("--previous", type=Path) args = parser.parse_args() previous = json.loads(args.previous.read_text()) if args.previous else None current = snapshot() args.output.parent.mkdir(parents=True, exist_ok=True) args.output.write_text(json.dumps(current, indent=2) + "\n") print(json.dumps(summarize(current, previous), indent=2)) if __name__ == "__main__": main()