mirror of
https://github.com/dredx/prole.git
synced 2026-09-23 12:03:59 +00:00
kube-prometheus-stack CRD re-application stalls on every installer run even when the release is in deployed state with all pods Running. For k3d this is purely idempotency noise — skip the upgrade and return immediately. Broken state (failed status, unbound PVCs) still triggers a full purge+reinstall. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
531 lines
19 KiB
Python
531 lines
19 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
import os
|
|
import tempfile
|
|
import time
|
|
|
|
from ._services_common import _LogFn, _detect_mode, _helm, _kubectl, _log, _namespace, _to_bool
|
|
|
|
|
|
def _monitoring_namespace(namespace: str | None, env: dict | None) -> str:
|
|
if env:
|
|
explicit = str(env.get("MONITORING_NAMESPACE") or "").strip()
|
|
if explicit:
|
|
return explicit
|
|
return _namespace(namespace, env, default="monitoring")
|
|
|
|
|
|
def _release_name(env: dict | None) -> str:
|
|
return str((env or {}).get("MONITORING_RELEASE") or "prometheus")
|
|
|
|
|
|
def _is_app_cluster_context(env: dict | None) -> bool:
|
|
app_ctx = str((env or {}).get("APP_CLUSTER_KUBECONTEXT") or "").strip()
|
|
current_ctx = str((env or {}).get("KUBECONTEXT") or "").strip()
|
|
if not app_ctx:
|
|
return True
|
|
if not current_ctx:
|
|
return False
|
|
return current_ctx == app_ctx
|
|
|
|
|
|
def _helm_status_value(release: str, namespace: str, env: dict | None) -> str:
|
|
res = _helm(["status", release, "-n", namespace, "-o", "json"], env=env, timeout=60)
|
|
if res.returncode != 0:
|
|
return ""
|
|
try:
|
|
data = json.loads(res.stdout or "{}")
|
|
except json.JSONDecodeError:
|
|
return ""
|
|
return str(((data.get("info") or {}).get("status") or "")).strip().lower()
|
|
|
|
|
|
def _helm_history_json(release: str, namespace: str, env: dict | None) -> list[dict]:
|
|
res = _helm(["history", release, "-n", namespace, "-o", "json"], env=env, timeout=60)
|
|
if res.returncode != 0:
|
|
return []
|
|
try:
|
|
data = json.loads(res.stdout or "[]")
|
|
except json.JSONDecodeError:
|
|
return []
|
|
if isinstance(data, list):
|
|
return [entry for entry in data if isinstance(entry, dict)]
|
|
return []
|
|
|
|
|
|
def _helm_history_text(release: str, namespace: str, env: dict | None) -> str:
|
|
res = _helm(["history", release, "-n", namespace], env=env, timeout=60)
|
|
return (res.stdout or "").strip()
|
|
|
|
|
|
def _revision_number(entry: dict) -> int:
|
|
try:
|
|
return int(str(entry.get("revision") or "").strip())
|
|
except (TypeError, ValueError):
|
|
return -1
|
|
|
|
|
|
def _latest_history_entry(entries: list[dict]) -> dict | None:
|
|
if not entries:
|
|
return None
|
|
return max(entries, key=_revision_number)
|
|
|
|
|
|
def _latest_pending_status(entries: list[dict]) -> tuple[str, int]:
|
|
latest = _latest_history_entry(entries) or {}
|
|
latest_status = str(latest.get("status") or "").strip().lower()
|
|
latest_revision = _revision_number(latest)
|
|
return latest_status, latest_revision
|
|
|
|
|
|
def _last_deployed_revision(entries: list[dict]) -> int:
|
|
deployed_revs = [_revision_number(entry) for entry in entries if str(entry.get("status") or "").strip().lower() == "deployed"]
|
|
deployed_revs = [rev for rev in deployed_revs if rev > 0]
|
|
if not deployed_revs:
|
|
return -1
|
|
return max(deployed_revs)
|
|
|
|
|
|
def _recover_stale_prometheus_helm_lock(
|
|
*,
|
|
release: str,
|
|
namespace: str,
|
|
env: dict | None,
|
|
log: _LogFn | None,
|
|
) -> bool:
|
|
pending_statuses = {"pending-install", "pending-upgrade", "pending-rollback"}
|
|
status = _helm_status_value(release, namespace, env)
|
|
history_entries = _helm_history_json(release, namespace, env)
|
|
latest_status, latest_revision = _latest_pending_status(history_entries)
|
|
|
|
if latest_status not in pending_statuses:
|
|
return False
|
|
|
|
if status and status != latest_status:
|
|
_log(
|
|
log,
|
|
f"[MONITORING] Helm release {release} latest history={latest_status} but current status={status}; treating as progressing, skipping stale-lock recovery.",
|
|
)
|
|
return False
|
|
|
|
poll_count = 2
|
|
poll_sleep_seconds = 30
|
|
_log(
|
|
log,
|
|
f"[MONITORING] Detected pending Helm state for {release}: status={latest_status} rev={latest_revision}. Verifying stability over {poll_count} poll(s) / {poll_count * poll_sleep_seconds}s.",
|
|
)
|
|
|
|
for poll in range(1, poll_count + 1):
|
|
time.sleep(poll_sleep_seconds)
|
|
polled_entries = _helm_history_json(release, namespace, env)
|
|
polled_status, polled_revision = _latest_pending_status(polled_entries)
|
|
_log(
|
|
log,
|
|
f"[MONITORING] Helm stale-lock poll {poll}/{poll_count}: latest status={polled_status or '<none>'} rev={polled_revision}.",
|
|
)
|
|
if polled_status not in pending_statuses:
|
|
_log(log, f"[MONITORING] Helm release {release} left pending state during verification; not stale.")
|
|
return False
|
|
if polled_status != latest_status or polled_revision != latest_revision:
|
|
_log(
|
|
log,
|
|
f"[MONITORING] Helm release {release} pending state changed during verification; treating as actively progressing.",
|
|
)
|
|
return False
|
|
|
|
last_deployed_revision = _last_deployed_revision(history_entries)
|
|
if last_deployed_revision <= 0:
|
|
history_text = _helm_history_text(release, namespace, env)
|
|
raise RuntimeError(
|
|
"REPAIR_BLOCKED: stale Helm lock detected for "
|
|
f"release '{release}' in namespace '{namespace}' (latest status={latest_status}, revision={latest_revision}), "
|
|
"but no prior deployed revision exists.\n"
|
|
f"helm history output:\n{history_text or '<empty>'}"
|
|
)
|
|
|
|
_log(
|
|
log,
|
|
f"[MONITORING] Recovering stale Helm lock for {release}: rolling back to deployed revision {last_deployed_revision}.",
|
|
)
|
|
_helm(
|
|
["rollback", release, str(last_deployed_revision), "-n", namespace],
|
|
env=env,
|
|
timeout=600,
|
|
check=True,
|
|
)
|
|
post_rollback_status = _helm_status_value(release, namespace, env)
|
|
if post_rollback_status != "deployed":
|
|
raise RuntimeError(
|
|
"REPAIR_BLOCKED: Helm rollback did not restore deployed state for "
|
|
f"release '{release}' in namespace '{namespace}' (status={post_rollback_status or '<unknown>'})."
|
|
)
|
|
_log(log, f"[MONITORING] Helm stale-lock recovery succeeded for {release}; status is deployed.")
|
|
return True
|
|
|
|
|
|
def initialize(
|
|
*,
|
|
namespace: str | None = None,
|
|
env: dict | None = None,
|
|
log: _LogFn | None = None,
|
|
mode: str | None = None,
|
|
) -> None:
|
|
update(namespace=namespace, env=env, log=log, mode=mode)
|
|
|
|
|
|
def start(
|
|
*,
|
|
namespace: str | None = None,
|
|
env: dict | None = None,
|
|
log: _LogFn | None = None,
|
|
mode: str | None = None,
|
|
) -> None:
|
|
update(namespace=namespace, env=env, log=log, mode=mode)
|
|
|
|
|
|
def _values_yaml_k8s(grafana_password: str) -> str:
|
|
"""Helm values for GKE / cloud-managed Kubernetes (no local PV affinity)."""
|
|
return (
|
|
"prometheus:\n"
|
|
" prometheusSpec:\n"
|
|
" storageSpec:\n"
|
|
" volumeClaimTemplate:\n"
|
|
" spec:\n"
|
|
" storageClassName: standard\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" resources:\n"
|
|
" requests:\n"
|
|
" storage: 10Gi\n"
|
|
"alertmanager:\n"
|
|
" alertmanagerSpec:\n"
|
|
" storage:\n"
|
|
" volumeClaimTemplate:\n"
|
|
" spec:\n"
|
|
" storageClassName: standard\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" resources:\n"
|
|
" requests:\n"
|
|
" storage: 5Gi\n"
|
|
"grafana:\n"
|
|
" adminUser: admin\n"
|
|
f" adminPassword: {grafana_password}\n"
|
|
" persistence:\n"
|
|
" type: sts\n"
|
|
" enabled: true\n"
|
|
" storageClassName: standard\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" size: 5Gi\n"
|
|
)
|
|
|
|
|
|
def _values_yaml_k3d(grafana_password: str, env: dict | None) -> str:
|
|
"""Helm values for k3d local dev cluster (local-path storage, no node affinity)."""
|
|
sc = str((env or {}).get("MONITORING_STORAGE_CLASS") or "local-path")
|
|
return (
|
|
"grafana:\n"
|
|
" adminUser: admin\n"
|
|
f" adminPassword: {grafana_password}\n"
|
|
" persistence:\n"
|
|
" type: sts\n"
|
|
" enabled: true\n"
|
|
f" storageClassName: {sc}\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" size: 10Gi\n"
|
|
"prometheus:\n"
|
|
" prometheusSpec:\n"
|
|
" storageSpec:\n"
|
|
" volumeClaimTemplate:\n"
|
|
" spec:\n"
|
|
f" storageClassName: {sc}\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" resources:\n"
|
|
" requests:\n"
|
|
" storage: 10Gi\n"
|
|
"alertmanager:\n"
|
|
" alertmanagerSpec:\n"
|
|
" storage:\n"
|
|
" volumeClaimTemplate:\n"
|
|
" spec:\n"
|
|
f" storageClassName: {sc}\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" resources:\n"
|
|
" requests:\n"
|
|
" storage: 2Gi\n"
|
|
)
|
|
|
|
|
|
def _values_yaml_k3s(grafana_password: str, env: dict | None) -> str:
|
|
"""Helm values for k3s homelab (local iSCSI PVs pinned to merlin.knoe.org)."""
|
|
data_dir = str((env or {}).get("PROLE_MONITORING_DATA_DIR") or "/synology/d004").rstrip("/")
|
|
volume_id = os.path.basename(data_dir) # e.g. "d004"
|
|
sc_prom = f"merlin-local-iscsi-{volume_id}-prometheus"
|
|
sc_alert = f"merlin-local-iscsi-{volume_id}-alertmanager"
|
|
sc_grafana = f"merlin-local-iscsi-{volume_id}-grafana"
|
|
|
|
monitoring_node = str((env or {}).get("MONITORING_PRIMARY_NODE") or "merlin.prole.org")
|
|
excluded_nodes = str((env or {}).get("MONITORING_NODE_EXPORTER_EXCLUDE_NODES") or "pi.knoe.org")
|
|
excluded_list = [n.strip() for n in excluded_nodes.split(",") if n.strip()]
|
|
|
|
node_affinity_yaml = (
|
|
" affinity:\n"
|
|
" nodeAffinity:\n"
|
|
" requiredDuringSchedulingIgnoredDuringExecution:\n"
|
|
" nodeSelectorTerms:\n"
|
|
" - matchExpressions:\n"
|
|
" - key: kubernetes.io/hostname\n"
|
|
" operator: In\n"
|
|
" values:\n"
|
|
f" - {monitoring_node}\n"
|
|
)
|
|
|
|
values = "prometheus-node-exporter:\n"
|
|
if excluded_list:
|
|
values += (
|
|
" affinity:\n"
|
|
" nodeAffinity:\n"
|
|
" requiredDuringSchedulingIgnoredDuringExecution:\n"
|
|
" nodeSelectorTerms:\n"
|
|
" - matchExpressions:\n"
|
|
" - key: kubernetes.io/hostname\n"
|
|
" operator: NotIn\n"
|
|
" values:\n"
|
|
)
|
|
for node in excluded_list:
|
|
values += f" - {node}\n"
|
|
|
|
values += (
|
|
"prometheus:\n"
|
|
" prometheusSpec:\n"
|
|
+ node_affinity_yaml
|
|
+ " storageSpec:\n"
|
|
" volumeClaimTemplate:\n"
|
|
" spec:\n"
|
|
f" storageClassName: {sc_prom}\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" resources:\n"
|
|
" requests:\n"
|
|
" storage: 30Gi\n"
|
|
)
|
|
values += (
|
|
"alertmanager:\n"
|
|
" alertmanagerSpec:\n"
|
|
+ node_affinity_yaml
|
|
+ " storage:\n"
|
|
" volumeClaimTemplate:\n"
|
|
" spec:\n"
|
|
f" storageClassName: {sc_alert}\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" resources:\n"
|
|
" requests:\n"
|
|
" storage: 5Gi\n"
|
|
)
|
|
values += (
|
|
"grafana:\n"
|
|
" adminUser: admin\n"
|
|
f" adminPassword: {grafana_password}\n"
|
|
" persistence:\n"
|
|
" type: sts\n"
|
|
" enabled: true\n"
|
|
f" storageClassName: {sc_grafana}\n"
|
|
" accessModes: [ReadWriteOnce]\n"
|
|
" size: 10Gi\n"
|
|
)
|
|
return values
|
|
|
|
|
|
def _purge_broken_monitoring(ns: str, release: str, env: dict | None, log: _LogFn | None) -> None:
|
|
"""Uninstall a broken monitoring release and delete only its unbound (Pending) PVCs.
|
|
|
|
Bound PVCs are left intact so that any pod that came up after a prior timeout
|
|
can still access its data on the next install, and to avoid wiping working state.
|
|
"""
|
|
_log(log, f"[MONITORING] Purging broken {release} release in {ns}")
|
|
_helm(["uninstall", release, "-n", ns], env=env, timeout=300)
|
|
pvc_res = _kubectl(["-n", ns, "get", "pvc", "--no-headers"], env=env, timeout=30)
|
|
for line in (pvc_res.stdout or "").splitlines():
|
|
parts = line.split()
|
|
if len(parts) >= 2 and parts[1] == "Pending":
|
|
pvc_name = parts[0].strip()
|
|
_log(log, f"[MONITORING] Deleting unbound PVC {pvc_name}")
|
|
_kubectl(["-n", ns, "delete", f"pvc/{pvc_name}", "--ignore-not-found=true"], env=env, timeout=60)
|
|
|
|
|
|
def update(
|
|
*,
|
|
namespace: str | None = None,
|
|
env: dict | None = None,
|
|
log: _LogFn | None = None,
|
|
mode: str | None = None,
|
|
) -> None:
|
|
effective_mode = _detect_mode(mode, env)
|
|
ns = _monitoring_namespace(namespace, env)
|
|
release = _release_name(env)
|
|
chart = str((env or {}).get("MONITORING_CHART") or "prometheus-community/kube-prometheus-stack")
|
|
grafana_password = str((env or {}).get("GRAFANA_ADMIN_PASSWORD") or "knoe")
|
|
|
|
_log(log, "[MONITORING] Ensuring helm repos")
|
|
_helm(["repo", "add", "prometheus-community", "https://prometheus-community.github.io/helm-charts"], env=env)
|
|
_helm(["repo", "update"], env=env)
|
|
|
|
_kubectl(["create", "namespace", ns], env=env, timeout=60)
|
|
|
|
if effective_mode == "k8s":
|
|
values_yaml = _values_yaml_k8s(grafana_password)
|
|
elif effective_mode == "k3d":
|
|
values_yaml = _values_yaml_k3d(grafana_password, env)
|
|
else:
|
|
# k3s: homelab local-PV setup pinned to merlin.knoe.org
|
|
values_yaml = _values_yaml_k3s(grafana_password, env)
|
|
|
|
# Detect broken state: Helm release is "failed" AND there are unbound PVCs.
|
|
# Gating on "failed" (not just Pending pods) avoids purging a still-converging
|
|
# install or a release that timed out but whose pods eventually came up.
|
|
helm_status = _helm_status_value(release, ns, env)
|
|
if helm_status == "failed":
|
|
pvcs_res = _kubectl(["-n", ns, "get", "pvc", "--no-headers"], env=env, timeout=30)
|
|
has_unbound_pvcs = any(
|
|
len(ln.split()) >= 2 and ln.split()[1] == "Pending"
|
|
for ln in (pvcs_res.stdout or "").splitlines()
|
|
if ln.strip()
|
|
)
|
|
if has_unbound_pvcs:
|
|
_log(log, f"[MONITORING] Helm release {release} is failed with unbound PVCs in {ns}; purging.")
|
|
_purge_broken_monitoring(ns, release, env, log)
|
|
elif helm_status == "deployed" and effective_mode == "k3d":
|
|
# k3d idempotency: CRD re-application on every run stalls for large charts.
|
|
# If all monitoring pods are Running, skip the upgrade.
|
|
pods_res = _kubectl(["-n", ns, "get", "pods", "--no-headers"], env=env, timeout=30)
|
|
pod_lines = [ln for ln in (pods_res.stdout or "").splitlines() if ln.strip()]
|
|
not_running = [
|
|
ln for ln in pod_lines
|
|
if len(ln.split()) >= 3 and ln.split()[2] not in ("Running", "Completed")
|
|
]
|
|
if pod_lines and not not_running:
|
|
_log(log, f"[MONITORING] Release {release} already deployed and all pods Running — skipping upgrade")
|
|
return
|
|
|
|
tmp_values = tempfile.NamedTemporaryFile(
|
|
mode="w", suffix=".yaml", prefix="monitoring-values-", delete=False
|
|
)
|
|
try:
|
|
tmp_values.write(values_yaml)
|
|
tmp_values.flush()
|
|
tmp_values.close()
|
|
|
|
recovered_stale_lock = False
|
|
if (
|
|
effective_mode == "k8s"
|
|
and _is_app_cluster_context(env)
|
|
and release == "prometheus"
|
|
and ns == "monitoring"
|
|
):
|
|
recovered_stale_lock = _recover_stale_prometheus_helm_lock(
|
|
release=release,
|
|
namespace=ns,
|
|
env=env,
|
|
log=log,
|
|
)
|
|
if recovered_stale_lock:
|
|
_log(log, f"[MONITORING] Retrying helm upgrade --install for {release} after stale-lock recovery.")
|
|
|
|
_log(log, f"[MONITORING] Deploying {release} in namespace {ns}")
|
|
if effective_mode == "k3d":
|
|
# k3d: the Prometheus operator converges StatefulSets asynchronously
|
|
# after Helm exits; --wait would block on all components (prometheus,
|
|
# alertmanager, node-exporter) and time out. Let Helm apply resources
|
|
# and return immediately; the status check in status_common_services.sh
|
|
# verifies readiness independently.
|
|
_helm(
|
|
[
|
|
"upgrade", "--install", release, chart,
|
|
"--namespace", ns,
|
|
"-f", tmp_values.name,
|
|
"--timeout", "5m",
|
|
],
|
|
env=env,
|
|
timeout=360,
|
|
check=True,
|
|
)
|
|
else:
|
|
_helm(
|
|
[
|
|
"upgrade", "--install", release, chart,
|
|
"--namespace", ns,
|
|
"-f", tmp_values.name,
|
|
"--wait",
|
|
"--timeout", "10m",
|
|
],
|
|
env=env,
|
|
timeout=660,
|
|
check=True,
|
|
)
|
|
|
|
# Cross-cluster: install Prometheus Operator CRDs on the DB cluster so that
|
|
# CNPG (running on knoe-cnpg-0) can create PodMonitor resources without
|
|
# the reconciler looping on "PodMonitor CRD not present".
|
|
db_ctx = str((env or {}).get("DB_CLUSTER_KUBECONTEXT") or "").strip()
|
|
if db_ctx and effective_mode == "k8s":
|
|
_log(log, f"[MONITORING] Installing Prometheus Operator CRDs on DB cluster ({db_ctx})")
|
|
# helm show crds is client-side — clear KUBECONTEXT so no --kube-context flag
|
|
crd_env = {**(env or {}), "KUBECONTEXT": ""}
|
|
crd_res = _helm(["show", "crds", chart], env=crd_env, timeout=60)
|
|
if crd_res.returncode == 0 and crd_res.stdout.strip():
|
|
db_env = {**(env or {}), "KUBECONTEXT": db_ctx}
|
|
_kubectl(
|
|
["apply", "--server-side", "-f", "-"],
|
|
env=db_env,
|
|
input_text=crd_res.stdout,
|
|
timeout=120,
|
|
)
|
|
_log(log, "[MONITORING] Prometheus Operator CRDs installed on DB cluster.")
|
|
else:
|
|
_log(log, "[MONITORING] WARNING: could not fetch CRDs from chart — skipping DB cluster CRD install.")
|
|
finally:
|
|
os.unlink(tmp_values.name)
|
|
|
|
|
|
def restart(
|
|
*,
|
|
namespace: str | None = None,
|
|
env: dict | None = None,
|
|
log: _LogFn | None = None,
|
|
) -> None:
|
|
ns = _monitoring_namespace(namespace, env)
|
|
_log(log, f"[MONITORING] Restarting grafana deployment in namespace {ns}")
|
|
_kubectl(["-n", ns, "rollout", "restart", "deployment/prometheus-grafana"], env=env, timeout=180)
|
|
|
|
|
|
def stop(
|
|
*,
|
|
namespace: str | None = None,
|
|
env: dict | None = None,
|
|
log: _LogFn | None = None,
|
|
) -> None:
|
|
ns = _monitoring_namespace(namespace, env)
|
|
release = _release_name(env)
|
|
_log(log, f"[MONITORING] Uninstalling release {release} in namespace {ns}")
|
|
_helm(["uninstall", release, "--namespace", ns], env=env, timeout=180)
|
|
|
|
|
|
def status(
|
|
*,
|
|
namespace: str | None = None,
|
|
env: dict | None = None,
|
|
) -> bool:
|
|
ns = _monitoring_namespace(namespace, env)
|
|
dep = _kubectl(
|
|
[
|
|
"-n",
|
|
ns,
|
|
"get",
|
|
"deployment",
|
|
"prometheus-grafana",
|
|
"-o",
|
|
"jsonpath={.status.readyReplicas}",
|
|
],
|
|
env=env,
|
|
timeout=20,
|
|
)
|
|
return _to_bool(dep.returncode == 0 and (dep.stdout or "0").strip() not in {"", "0"})
|