mirror of
https://github.com/dredx/prole.git
synced 2026-09-23 10:13:58 +00:00
Minor updates aligned with min-mode and auth Round 1 integration. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
435 lines
15 KiB
Python
435 lines
15 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_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 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)
|
|
else:
|
|
# k3s / k3d: homelab local-PV setup
|
|
values_yaml = _values_yaml_k3s(grafana_password, env)
|
|
|
|
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}")
|
|
_helm(
|
|
[
|
|
"upgrade",
|
|
"--install",
|
|
release,
|
|
chart,
|
|
"--namespace",
|
|
ns,
|
|
"-f",
|
|
tmp_values.name,
|
|
"--wait",
|
|
],
|
|
env=env,
|
|
timeout=600,
|
|
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"})
|