atlas-iac/scripts/ops/cluster_settle_check.py

193 lines
9.0 KiB
Python
Raw Normal View History

#!/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()