From 0282f1af5f85b944da24294852ba6dec1ad650f2 Mon Sep 17 00:00:00 2001 From: jenkins Date: Sun, 4 Oct 2026 06:55:56 -0500 Subject: [PATCH] ops: add read-only cluster settling comparison --- scripts/ops/cluster_settle_check.py | 187 +++++++++++++++++++++ testing/tests/test_cluster_settle_check.py | 170 +++++++++++++++++++ 2 files changed, 357 insertions(+) create mode 100755 scripts/ops/cluster_settle_check.py create mode 100644 testing/tests/test_cluster_settle_check.py diff --git a/scripts/ops/cluster_settle_check.py b/scripts/ops/cluster_settle_check.py new file mode 100755 index 00000000..8952549f --- /dev/null +++ b/scripts/ops/cluster_settle_check.py @@ -0,0 +1,187 @@ +#!/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"], +} + + +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"), + "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"), + } + for e in raw["warnings"] + } + 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" and v["node"] not in offline], + "unready_pods_on_offline_nodes": sum(v["ready"] != "True" 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() diff --git a/testing/tests/test_cluster_settle_check.py b/testing/tests/test_cluster_settle_check.py new file mode 100644 index 00000000..fc040b05 --- /dev/null +++ b/testing/tests/test_cluster_settle_check.py @@ -0,0 +1,170 @@ +"""Regression checks for interpreting quiet-window metadata without false health.""" + +from copy import deepcopy + +from scripts.ops.cluster_settle_check import summarize + + +def baseline(): + """Return a minimal healthy snapshot without source or credential data.""" + return { + "observed_at": "2026-10-04T09:00:00+00:00", + "nodes": {"worker": {"conditions": {"Ready": "True"}}}, + "workloads": {}, "pods": {}, "volumes": {}, "warnings": {}, + "flux": {}, "helm": {}, + } + + +def test_requested_detached_storage_is_not_reported_healthy(): + """A stuck requested attachment remains visible even before it is attached.""" + current = baseline() + current["volumes"]["storage/in-use"] = { + "state": "detached", "robustness": "unknown", "requested_node": "worker", + } + current["volumes"]["storage/retained"] = { + "state": "detached", "robustness": "unknown", "requested_node": "", + } + assert list(summarize(current)["unavailable_requested_volumes"]) == ["storage/in-use"] + + +def test_event_lifetime_count_is_not_a_measured_window_delta(): + """New event UIDs have unknown deltas; matching UIDs yield counter differences.""" + old = baseline() + old["warnings"]["existing"] = {"reason": "BackOff", "count": 40} + current = deepcopy(old) + current["observed_at"] = "2026-10-04T09:10:00+00:00" + current["warnings"]["existing"]["count"] = 42 + current["warnings"]["new-uid"] = {"reason": "BackOff", "count": 900} + result = summarize(current, old) + assert result["window_seconds"] == 600 + assert [w["increase"] for w in result["warning_changes"]] == [2, None] + + +def test_replaced_stateful_pod_is_distinct_from_a_container_restart(): + """The same pod name with a new UID cannot conceal controller churn.""" + old = baseline() + old["pods"]["app/stateful-0"] = { + "uid": "before", "node": "worker", "ready": "True", "restarts": {"app": 7}, + } + current = deepcopy(old) + current["pods"]["app/stateful-0"].update(uid="after", restarts={"app": 0}) + result = summarize(current, old) + assert result["replaced_service_pods"] == ["app/stateful-0"] + assert not result["restart_increases"] + + +def test_offline_pods_are_separate_from_new_reachable_node_failures(): + """Known offline-node failures stay explicit without hiding fresh failures.""" + current = baseline() + current["nodes"]["offline"] = {"conditions": {"Ready": "Unknown"}} + current["pods"] = { + "app/old": {"node": "offline", "ready": "False"}, + "app/new": {"node": "worker", "ready": "False"}, + } + result = summarize(current) + assert result["offline_nodes"] == ["offline"] + assert result["unready_pods_on_offline_nodes"] == 1 + assert result["unready_pods_on_reachable_nodes"] == ["app/new"] + + +def test_serving_old_replica_does_not_hide_failed_rollout(): + """Desired availability can coexist with a broken new release.""" + current = baseline() + current["workloads"]["deployments/app/service"] = { + "desired": 1, "ready": 1, "generation": 2, "observed_generation": 2, + "rollout_failed": True, + } + assert list(summarize(current)["unready_workloads"]) == ["deployments/app/service"] + + +def test_snapshot_omits_content_and_short_lived_jobs(monkeypatch): + """Only operational metadata leaves raw API objects; Jobs are not service churn.""" + import json + from scripts.ops import cluster_settle_check as checker + + hidden = 'SYNTHETIC_CONTENT_MUST_NOT_APPEAR' + pod = { + 'metadata': {'name': 'service', 'namespace': 'app', 'uid': 'u1'}, + 'spec': {'nodeName': 'worker', 'containers': [{'env': [{'value': hidden}]}]}, + 'status': {'phase': 'Running', 'conditions': [{'type': 'Ready', 'status': 'True'}], + 'containerStatuses': [{'name': 'app', 'restartCount': 0, 'state': {}}]}, + } + job_pod = deepcopy(pod) + job_pod['metadata'].update(name='job', ownerReferences=[{'kind': 'Job', 'name': 'timer'}]) + done_pod = deepcopy(pod) + done_pod['metadata']['name'] = 'finished' + done_pod['status']['phase'] = 'Succeeded' + data = {name: [] for name in checker.QUERIES} + data['pods'] = [pod, job_pod, done_pod] + data['warnings'] = [{ + 'metadata': {'uid': 'event'}, 'involvedObject': {'namespace': 'app', 'name': 'service'}, + 'reason': 'Unhealthy', 'count': 1, 'message': hidden, + }] + monkeypatch.setattr(checker, 'fetch', lambda item: (item[0], data[item[0]])) + current = checker.snapshot() + assert list(current['pods']) == ['app/service'] + assert hidden not in json.dumps(current) + assert current['warnings']['event']['reason'] == 'Unhealthy' + + +def test_failed_api_read_cannot_become_a_partial_health_report(monkeypatch): + """A failed required query raises without echoing an API response body.""" + from types import SimpleNamespace + import pytest + from scripts.ops import cluster_settle_check as checker + + hidden = 'SYNTHETIC_PROVIDER_RESPONSE_CONTENT' + seen = [] + + def fail_read(command, **kwargs): + """Record the safe read command and simulate a failing API process.""" + seen.append(command) + return SimpleNamespace(returncode=1, stdout=hidden, stderr=hidden) + + monkeypatch.setattr(checker.subprocess, 'run', fail_read) + with pytest.raises(RuntimeError) as failure: + checker.fetch(('nodes', ['nodes'])) + assert 'Cannot read nodes' in str(failure.value) + assert hidden not in str(failure.value) + assert seen[0][:4] == ['kubectl', '--request-timeout=20s', 'get', 'nodes'] + + +def test_restart_increase_requires_same_pod_identity(): + """Resets cannot cancel another container's observed restarts.""" + old = baseline() + old['pods']['app/service'] = { + 'uid': 'same', 'node': 'worker', 'ready': 'True', + 'restarts': {'app': 7, 'helper': 2}, + } + current = deepcopy(old) + current['pods']['app/service']['restarts'] = {'app': 0, 'helper': 5} + assert summarize(current, old)['restart_increases'] == {'app/service': 3} + + +def test_successful_api_read_projects_only_items(monkeypatch): + """The list envelope is not part of the data passed to health analysis.""" + from types import SimpleNamespace + from scripts.ops import cluster_settle_check as checker + + monkeypatch.setattr(checker.subprocess, 'run', lambda *a, **k: SimpleNamespace( + returncode=0, stdout='{"items": [], "metadata": {"resourceVersion": "1"}}', + )) + assert checker.fetch(('nodes', ['nodes'])) == ('nodes', []) + + +def test_cli_writes_snapshot_and_reports_comparison(tmp_path, monkeypatch, capsys): + """The operator's two-file workflow preserves its baseline and interval.""" + import json + from scripts.ops import cluster_settle_check as checker + + before = tmp_path / 'before.json' + before.write_text(json.dumps(baseline())) + current = baseline() + current['observed_at'] = '2026-10-04T09:15:00+00:00' + output = tmp_path / 'new' / 'after.json' + monkeypatch.setattr(checker, 'snapshot', lambda: current) + monkeypatch.setattr('sys.argv', ['check', '--output', str(output), '--previous', str(before)]) + checker.main() + assert json.loads(capsys.readouterr().out)['window_seconds'] == 900 + assert json.loads(output.read_text()) == current + assert json.loads(before.read_text()) == baseline()