ops: add read-only cluster settling comparison

This commit is contained in:
jenkins 2026-10-04 06:55:56 -05:00
parent b2a1761e77
commit 0282f1af5f
2 changed files with 357 additions and 0 deletions

View File

@ -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()

View File

@ -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()