prole/knoe/core/ops/monitoring.py
chrisfu 8dcec846cd fix(monitoring): stop purge cycle on k3d retry
Three issues caused the purge-and-reinstall loop:

1. _purge_broken_monitoring deleted ALL PVCs including Bound ones; grafana's
   working PVC was wiped on every retry. Now only delete PVCs in Pending state.

2. Broken-state detection keyed on Pending pods + unbound PVCs, which is true
   during any still-converging install (including ones that timed out but
   whose pods eventually came up). Gate on helm status=='failed' + unbound PVCs.

3. k3d install used --wait, which blocks on all kube-prometheus-stack components
   (prometheus, alertmanager, node-exporter). They converge async after the
   Prometheus operator starts; --wait always timed out. Drop --wait for k3d;
   the status_common_services.sh check verifies readiness independently.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-03 12:55:11 -07:00

519 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)
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"})