state-hub/scripts/release_headroom_preflight.py

240 lines
9.6 KiB
Python
Raw Normal View History

#!/usr/bin/env python3
"""Read-only scheduling-headroom preflight for State Hub Helm promotion.
STATE-WP-0091-T01. Does not start Helm, does not lower requests, and does not
claim a scheduling guarantee from an aggregate remaining-millicores sum.
"""
from __future__ import annotations
import argparse
import json
import subprocess
import sys
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
CPU_FACTORS = {"n": 0.000001, "u": 0.001, "m": 1.0, "": 1000.0}
BYTE_FACTORS = {
"": 1,
"Ki": 1024,
"Mi": 1024**2,
"Gi": 1024**3,
"Ti": 1024**4,
"K": 1000,
"M": 1000**2,
"G": 1000**3,
"T": 1000**4,
}
# Historical State Hub API surge (maxSurge=1, maxUnavailable=0).
DEFAULT_API_SURGE_CPU_M = 100
DEFAULT_MCP_SURGE_CPU_M = 10
DEFAULT_MIGRATE_CPU_M = 50
DEFAULT_API_SURGE_MEMORY = 512 * 1024**2
DEFAULT_FRESHNESS_SECONDS = 15 * 60
def cpu_millicores(value: str | None) -> int:
import re
match = re.fullmatch(r"([0-9]+(?:\.[0-9]+)?)(n|u|m)?", str(value or "0"))
if not match:
raise ValueError(f"unsupported CPU quantity: {value}")
return round(float(match.group(1)) * CPU_FACTORS[match.group(2) or ""])
def bytes_value(value: str | None) -> int:
import re
match = re.fullmatch(r"([0-9]+(?:\.[0-9]+)?)(Ki|Mi|Gi|Ti|K|M|G|T)?", str(value or "0"))
if not match:
raise ValueError(f"unsupported byte quantity: {value}")
return round(float(match.group(1)) * BYTE_FACTORS[match.group(2) or ""])
def pod_effective_requests(pod: dict[str, Any]) -> tuple[int, int]:
"""Pod request is max(sum(containers), max(init containers))."""
spec = pod.get("spec") or {}
containers = spec.get("containers") or []
inits = spec.get("initContainers") or []
cpu = sum(cpu_millicores((c.get("resources") or {}).get("requests", {}).get("cpu")) for c in containers)
mem = sum(bytes_value((c.get("resources") or {}).get("requests", {}).get("memory")) for c in containers)
init_cpu = max((cpu_millicores((c.get("resources") or {}).get("requests", {}).get("cpu")) for c in inits), default=0)
init_mem = max((bytes_value((c.get("resources") or {}).get("requests", {}).get("memory")) for c in inits), default=0)
return max(cpu, init_cpu), max(mem, init_mem)
def evaluate(observation: dict[str, Any], release: dict[str, Any] | None = None) -> dict[str, Any]:
"""Return a preflight report. `ok` is whether Helm may start, not a guarantee."""
release = release or {}
api_cpu = int(release.get("api_surge_cpu_m", DEFAULT_API_SURGE_CPU_M))
mcp_cpu = int(release.get("mcp_surge_cpu_m", DEFAULT_MCP_SURGE_CPU_M))
migrate_cpu = int(release.get("migrate_cpu_m", DEFAULT_MIGRATE_CPU_M))
api_mem = int(release.get("api_surge_memory_bytes", DEFAULT_API_SURGE_MEMORY))
freshness_limit = int(release.get("freshness_seconds", DEFAULT_FRESHNESS_SECONDS))
atomic = bool(release.get("atomic", True))
reasons: list[str] = []
notes: list[str] = []
nodes = observation.get("nodes") or []
if not nodes:
reasons.append("no node capacity evidence")
remaining_cpu = 0
remaining_mem = 0
for node in nodes:
if node.get("unschedulable") or not node.get("ready", True):
reasons.append(f"node {node.get('name')} is not schedulable")
continue
remaining_cpu += int(node["allocatable_cpu_m"]) - int(node["allocated_cpu_m"])
remaining_mem += int(node["allocatable_memory_bytes"]) - int(node["allocated_memory_bytes"])
pending = observation.get("pending_unrelated") or []
pending_cpu = sum(int(p.get("cpu_m", 0)) for p in pending)
pending_mem = sum(int(p.get("memory_bytes", 0)) for p in pending)
observed_at = observation.get("observed_at")
freshness_seconds = observation.get("freshness_seconds")
if freshness_seconds is None and observed_at:
try:
stamp = datetime.fromisoformat(str(observed_at).replace("Z", "+00:00"))
freshness_seconds = int((datetime.now(timezone.utc) - stamp).total_seconds())
except ValueError:
freshness_seconds = None
reasons.append("observed_at is not parseable")
if freshness_seconds is not None and freshness_seconds > freshness_limit:
reasons.append(
f"capacity evidence is {freshness_seconds}s old (limit {freshness_limit}s)"
)
if remaining_cpu < migrate_cpu:
reasons.append(
f"migration hook needs {migrate_cpu}m CPU; remaining {remaining_cpu}m"
)
if remaining_cpu < api_cpu:
reasons.append(
f"API maxSurge=1 needs {api_cpu}m CPU; remaining {remaining_cpu}m "
f"(Insufficient CPU class; revision 61 had 65m remaining)"
)
if remaining_mem < api_mem:
reasons.append(
f"API maxSurge=1 needs {api_mem} bytes memory; remaining {remaining_mem}"
)
concurrent = api_cpu + mcp_cpu
if remaining_cpu < concurrent:
notes.append(
f"API+MCP concurrent surge would need {concurrent}m; remaining {remaining_cpu}m. "
"Helm may roll both deployments together; this is not a scheduling guarantee."
)
effective_cpu = remaining_cpu - pending_cpu
effective_mem = remaining_mem - pending_mem
if pending and effective_cpu < api_cpu:
reasons.append(
f"unrelated pending demand {pending_cpu}m leaves {effective_cpu}m "
f"against API surge {api_cpu}m"
)
if pending and effective_mem < api_mem:
reasons.append(
f"unrelated pending memory {pending_mem} leaves {effective_mem} against API surge"
)
if not atomic:
reasons.append("promotion path must keep helm --atomic; preflight will not start a non-atomic upgrade")
ok = not reasons
if ok and remaining_cpu - api_cpu < 50:
notes.append(
f"remaining {remaining_cpu}m is only {remaining_cpu - api_cpu}m above the "
f"{api_cpu}m API surge (revision 63 was 105m). Not factory capacity admission."
)
notes.append("Aggregate remaining millicores is not a kube-scheduler guarantee.")
notes.append("Preflight does not lower requests and does not start Helm.")
return {
"schema": "state-hub.release-headroom-preflight.v1",
"ok": ok,
"observed_at": observed_at,
"freshness_seconds": freshness_seconds,
"remaining_cpu_m": remaining_cpu,
"remaining_memory_bytes": remaining_mem,
"pending_unrelated_cpu_m": pending_cpu,
"api_surge_cpu_m": api_cpu,
"mcp_surge_cpu_m": mcp_cpu,
"migrate_cpu_m": migrate_cpu,
"atomic": atomic,
"reasons": reasons,
"notes": notes,
"hook_order": [
"pre-upgrade migrate job (helm.sh/hook-weight -5)",
"API RollingUpdate maxSurge=1 maxUnavailable=0",
"MCP RollingUpdate maxSurge=1 maxUnavailable=0",
],
}
def collect_live(kubectl: str = "kubectl") -> dict[str, Any]:
nodes = json.loads(subprocess.check_output([kubectl, "get", "nodes", "-o", "json"], timeout=60))
pods = json.loads(subprocess.check_output([kubectl, "get", "pods", "-A", "-o", "json"], timeout=60))
allocated: dict[str, dict[str, int]] = {}
pending_unrelated: list[dict[str, Any]] = []
now = datetime.now(timezone.utc)
for node in nodes.get("items") or []:
name = node["metadata"]["name"]
alloc = node.get("status", {}).get("allocatable") or {}
allocated[name] = {
"allocatable_cpu_m": cpu_millicores(alloc.get("cpu")),
"allocatable_memory_bytes": bytes_value(alloc.get("memory")),
"allocated_cpu_m": 0,
"allocated_memory_bytes": 0,
"ready": any(
c.get("type") == "Ready" and c.get("status") == "True"
for c in (node.get("status") or {}).get("conditions") or []
),
"unschedulable": bool((node.get("spec") or {}).get("unschedulable")),
"name": name,
}
for pod in pods.get("items") or []:
phase = (pod.get("status") or {}).get("phase")
ns = pod.get("metadata", {}).get("namespace")
name = pod.get("metadata", {}).get("name")
cpu, mem = pod_effective_requests(pod)
node_name = (pod.get("spec") or {}).get("nodeName")
if phase == "Pending" and not node_name:
if ns != "state-hub":
pending_unrelated.append({"name": f"{ns}/{name}", "cpu_m": cpu, "memory_bytes": mem})
continue
if phase in {"Succeeded", "Failed"}:
continue
if node_name in allocated:
allocated[node_name]["allocated_cpu_m"] += cpu
allocated[node_name]["allocated_memory_bytes"] += mem
return {
"observed_at": now.isoformat(),
"freshness_seconds": 0,
"nodes": list(allocated.values()),
"pending_unrelated": pending_unrelated,
"source": "kubectl",
}
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--fixture", type=Path, help="JSON observation instead of live kubectl")
parser.add_argument("--release", type=Path, help="JSON release surge description")
parser.add_argument("--kubectl", default="kubectl")
args = parser.parse_args(argv)
if args.fixture:
observation = json.loads(args.fixture.read_text())
else:
observation = collect_live(args.kubectl)
release = json.loads(args.release.read_text()) if args.release else {}
report = evaluate(observation, release)
json.dump(report, sys.stdout, indent=2)
sys.stdout.write("\n")
return 0 if report["ok"] else 1
if __name__ == "__main__":
raise SystemExit(main())