473 lines
20 KiB
Python
Executable File
473 lines
20 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""
|
|
GPU Self-Heal — Implementation of gpu-self-heal.prose.md
|
|
Evaluates 8 remediation rules against live GPU monitor data.
|
|
Reports results to knowledge graph via MCP bridge.
|
|
"""
|
|
import json, time, subprocess, sys, os, socket
|
|
from datetime import datetime, timezone
|
|
from urllib.request import urlopen, Request
|
|
|
|
# Global timeout to prevent hanging on unreachable hosts
|
|
socket.setdefaulttimeout(10)
|
|
|
|
GPU_MONITOR = "http://192.168.68.24:9100/gpu-data"
|
|
KG_BRIDGE = "http://192.168.68.65:3100/mcp"
|
|
GPU_HOSTS = {
|
|
"ct8-rtx3090": {"ip": "192.168.68.8", "gpu": "NVIDIA RTX 3090", "vram_total_mb": 24576, "vram_threshold_mb_per_h": 100},
|
|
"ct110-rtx5070": {"ip": "192.168.68.110","gpu": "NVIDIA RTX 5070", "vram_total_mb": 12227, "vram_threshold_mb_per_h": 50},
|
|
"strix-halo": {"ip": "192.168.68.15", "gpu": "AMD Strix Halo 64GB","vram_total_mb": 65536, "vram_threshold_mb_per_h": 200},
|
|
}
|
|
HISTORY_FILE = "/var/log/litellm/gpu-history.json"
|
|
RUN_ID = f"gpu-self-heal-{datetime.now().strftime('%Y%m%d-%H%M%S')}"
|
|
|
|
def fetch_gpu_data():
|
|
"""Fetch live GPU fleet data from monitor."""
|
|
try:
|
|
req = Request(GPU_MONITOR)
|
|
with urlopen(req, timeout=10) as resp:
|
|
return json.loads(resp.read())
|
|
except Exception as e:
|
|
print(f" FAIL: GPU monitor unreachable: {e}")
|
|
return None
|
|
|
|
def fetch_prometheus(host, port=9400):
|
|
"""Check Prometheus exporter on GPU host (with timeout)."""
|
|
try:
|
|
req = Request(f"http://{host}:{port}/metrics")
|
|
with urlopen(req, timeout=3) as resp:
|
|
return resp.status == 200
|
|
except:
|
|
return False
|
|
|
|
def probe_strix_direct():
|
|
"""Direct Strix Halo health probe (firewall opened .24→.15:8080)."""
|
|
try:
|
|
with urlopen("http://192.168.68.15:8080/health", timeout=5) as resp:
|
|
return resp.status == 200
|
|
except:
|
|
return False
|
|
|
|
def load_history():
|
|
"""Load historical GPU metrics for trend analysis."""
|
|
if os.path.exists(HISTORY_FILE):
|
|
with open(HISTORY_FILE) as f:
|
|
return json.load(f)
|
|
return {"snapshots": [], "vram_baselines": {}, "temp_history": {}}
|
|
|
|
def save_history(history):
|
|
"""Persist history for trend analysis."""
|
|
os.makedirs(os.path.dirname(HISTORY_FILE), exist_ok=True)
|
|
# Keep last 72 hours of snapshots (4320 at 60s intervals, capped at 1000)
|
|
history["snapshots"] = history["snapshots"][-1000:]
|
|
with open(HISTORY_FILE, "w") as f:
|
|
json.dump(history, f, indent=2)
|
|
|
|
def analyze_vram_trend(history, gpu_name, current_vram_mb):
|
|
"""Calculate VRAM growth rate over 6-hour window."""
|
|
snapshots = history.get("snapshots", [])
|
|
six_hours_ago = time.time() - 21600
|
|
old_snapshots = [s for s in snapshots if s.get("timestamp", 0) > six_hours_ago
|
|
and s.get("gpu") == gpu_name]
|
|
if len(old_snapshots) < 10:
|
|
return 0 # Not enough data
|
|
oldest = old_snapshots[0]
|
|
newest = old_snapshots[-1]
|
|
hours = (newest["timestamp"] - oldest["timestamp"]) / 3600
|
|
if hours < 1:
|
|
return 0
|
|
return (current_vram_mb - oldest["vram_mb"]) / hours
|
|
|
|
def analyze_temp_rise(history, gpu_name, current_temp):
|
|
"""Calculate temperature rise rate over last 5 minutes."""
|
|
snapshots = history.get("snapshots", [])
|
|
five_min_ago = time.time() - 300
|
|
recent = [s for s in snapshots if s.get("timestamp", 0) > five_min_ago
|
|
and s.get("gpu") == gpu_name]
|
|
if len(recent) < 3:
|
|
return 0
|
|
oldest = recent[0]
|
|
minutes = (recent[-1]["timestamp"] - oldest["timestamp"]) / 60
|
|
if minutes < 0.5:
|
|
return 0
|
|
return (current_temp - oldest["temp_c"]) / minutes
|
|
|
|
def record_snapshot(history, gpu_name, temp_c, vram_mb):
|
|
"""Record a GPU snapshot for trend analysis."""
|
|
history["snapshots"].append({
|
|
"timestamp": time.time(),
|
|
"gpu": gpu_name,
|
|
"temp_c": temp_c,
|
|
"vram_mb": vram_mb
|
|
})
|
|
|
|
def kg_create_node(title, description, source):
|
|
"""Log run report to Gitea (hard rule: health logs NEVER go to knowledge graph).
|
|
|
|
Redirected 2026-08-13 per directive from Mumuni (#726): logs belong in
|
|
SyslogSolution/health-logs/gpu/{run_id}.json, not the RA-H OS graph.
|
|
"""
|
|
try:
|
|
# source is the report JSON string; write to temp file and push via gitea-logger
|
|
report_file = f"/tmp/{RUN_ID}.json"
|
|
with open(report_file, "w") as f:
|
|
f.write(source if isinstance(source, str) else json.dumps(source, indent=2))
|
|
cmd = f"/opt/inference-harness/scripts/gitea-logger.sh gpu {RUN_ID}.json {report_file}"
|
|
result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=120)
|
|
os.remove(report_file) if os.path.exists(report_file) else None
|
|
return result.stdout.strip() or "gitea-logger: no output"
|
|
except Exception as e:
|
|
return f"gitea-logger failed: {e}"
|
|
|
|
def evaluate_rules(data, history):
|
|
"""Evaluate all 8 remediation rules against current GPU state."""
|
|
actions = []
|
|
if not data:
|
|
return actions
|
|
|
|
gpus = data.get("gpus", [])
|
|
summary = data.get("summary", {})
|
|
benchmarks = data.get("benchmarks", {})
|
|
strix = data.get("strix", {})
|
|
router = data.get("router", {})
|
|
|
|
for gpu in gpus:
|
|
if not isinstance(gpu, dict):
|
|
continue
|
|
name = gpu.get("gpu_name", "unknown")
|
|
host_key = None
|
|
for k, v in GPU_HOSTS.items():
|
|
if v["gpu"] in name or k in name:
|
|
host_key = k
|
|
break
|
|
if not host_key:
|
|
print(f"[warn] no GPU_HOSTS mapping for {name!r} — remote rules (incl. Rule 7) skipped for this GPU")
|
|
host_key = name.replace(" ", "-").lower()
|
|
|
|
temp = gpu.get("temp_c", 0)
|
|
vram_used = gpu.get("vram_used_mb", 0)
|
|
host_info = GPU_HOSTS.get(host_key, {"ip": "unknown", "vram_threshold_mb_per_h": 50})
|
|
|
|
# Record snapshot
|
|
record_snapshot(history, name, temp, vram_used)
|
|
|
|
# ── Rule 1: Thermal Critical ──
|
|
if temp > 85:
|
|
actions.append({"rule": "thermal-critical", "gpu": name, "temp_c": temp,
|
|
"action": "reduce-concurrency", "severity": "critical"})
|
|
print(f" ⚠ RULE 1: {name} at {temp}°C — reducing concurrency")
|
|
|
|
# ── Rule 2: VRAM Leak ──
|
|
vram_rate = analyze_vram_trend(history, name, vram_used)
|
|
threshold = host_info.get("vram_threshold_mb_per_h", 50)
|
|
if vram_rate > threshold:
|
|
actions.append({"rule": "vram-leak", "gpu": name, "rate_mb_per_h": vram_rate,
|
|
"threshold": threshold, "action": "investigate-leak", "severity": "warning"})
|
|
print(f" ⚠ RULE 2: {name} VRAM leak {vram_rate:.1f} MB/h (threshold: {threshold})")
|
|
|
|
# ── Rule 4: Benchmark Regression ──
|
|
bench_latest = benchmarks.get("latest", {})
|
|
for bk, bv in bench_latest.items():
|
|
if not isinstance(bv, dict):
|
|
continue
|
|
latest_run = bv.get("latest", {})
|
|
baseline = bv.get("baseline_tok_per_sec", 0)
|
|
current_tps = latest_run.get("gen_tok_per_sec", 0)
|
|
if baseline > 0 and current_tps > 0 and current_tps < baseline * 0.8:
|
|
drop_pct = (1 - current_tps / baseline) * 100
|
|
gpu_label = bv.get("gpu_name", bk)
|
|
actions.append({"rule": "benchmark-regression", "gpu": gpu_label,
|
|
"baseline_tps": baseline, "current_tps": current_tps,
|
|
"drop_pct": round(drop_pct,1), "action": "investigate-performance",
|
|
"severity": "warning"})
|
|
print(f" ⚠ RULE 4: {gpu_label} benchmark -{drop_pct:.0f}% ({current_tps} vs {baseline} tok/s)")
|
|
|
|
# ── Rule 7: Prometheus Exporter ──
|
|
ip = host_info.get("ip", "")
|
|
if ip and ip != "unknown":
|
|
if not fetch_prometheus(ip):
|
|
actions.append({"rule": "prometheus-down", "gpu": name, "host": ip,
|
|
"action": "restart-exporter", "severity": "warning"})
|
|
print(f" ⚠ RULE 7: {name} Prometheus exporter down on {ip}:9400")
|
|
|
|
# ── Rule 8: Predictive Thermal ──
|
|
rise_rate = analyze_temp_rise(history, name, temp)
|
|
if temp > 70 and rise_rate > 2.0:
|
|
tier = "critical" if temp > 80 else "warning"
|
|
action_text = "aggressive-load-shedding" if tier == "critical" else "reduce-concurrency-50pct"
|
|
actions.append({"rule": "predictive-thermal", "gpu": name,
|
|
"temp_c": temp, "rise_rate_c_per_min": rise_rate,
|
|
"tier": tier, "action": action_text, "severity": tier})
|
|
print(f" {'⚠' if tier == 'critical' else '🔶'} RULE 8 ({tier}): {name} {temp}°C rising {rise_rate:.1f}°C/min")
|
|
|
|
# ── Rule 6: Strix Halo ──
|
|
if not strix.get("status") == "running":
|
|
if probe_strix_direct():
|
|
print(f" ⚠ RULE 6: Strix direct probe OK but monitor disagrees — router stale")
|
|
actions.append({"rule": "strix-router-stale", "gpu": "strix-halo",
|
|
"action": "refresh-router", "severity": "info"})
|
|
else:
|
|
print(f" ⚠ RULE 6: Strix Halo down — attempting restart")
|
|
actions.append({"rule": "strix-down", "gpu": "strix-halo",
|
|
"action": "restart-llama-server", "severity": "critical"})
|
|
|
|
# ── Rule 3: Model Stuck (check router) ──
|
|
router_models = router.get("available_models", [])
|
|
if isinstance(router_models, list):
|
|
for model in router_models:
|
|
if isinstance(model, dict) and model.get("consecutive_timeouts", 0) >= 3:
|
|
# Check failure rate
|
|
total = model.get("total_requests", 1)
|
|
failures = model.get("failed_requests", 0)
|
|
if total > 0 and failures / total > 0.5:
|
|
actions.append({"rule": "model-stuck", "gpu": model.get("gpu", "?"),
|
|
"model": model.get("name", "?"),
|
|
"failure_rate": failures/total,
|
|
"action": "restart-llama-server", "severity": "critical"})
|
|
print(f" ⚠ RULE 3: Model {model.get('name')} stuck ({failures}/{total} failures)")
|
|
|
|
# ── Rule 5: Circuit Breaker ──
|
|
cb_data = router.get("circuit_breaker", {})
|
|
if isinstance(cb_data, dict):
|
|
for cb_name, cb in cb_data.items():
|
|
if isinstance(cb, dict) and cb.get("open"):
|
|
open_sec = cb.get("open_duration_sec", 0)
|
|
gpu_healthy = cb.get("gpu_healthy", False)
|
|
if open_sec > 600 and gpu_healthy:
|
|
# Check cooldown: has this CB been reset in the last hour?
|
|
last_reset = history.get("cb_resets", {}).get(cb_name, 0)
|
|
if time.time() - last_reset > 3600:
|
|
actions.append({"rule": "cb-stuck-open", "gpu": cb.get("gpu", "?"),
|
|
"cb_name": cb_name, "open_sec": open_sec,
|
|
"action": "reset-cb", "severity": "warning"})
|
|
history.setdefault("cb_resets", {})[cb_name] = time.time()
|
|
print(f" ⚠ RULE 5: CB {cb_name} stuck open {open_sec}s — auto-resetting")
|
|
|
|
return actions
|
|
|
|
def evaluate_context_optimization(data):
|
|
"""
|
|
Rule 9: Context Window Optimization
|
|
Check if each GPU's context window is optimal for its role.
|
|
"""
|
|
benchmarks = data.get("benchmarks", {}).get("latest", {})
|
|
results = []
|
|
|
|
# Actual topology (verified July 2026)
|
|
gpu_config = {
|
|
"ct8-rtx3090": {
|
|
"context": 262144, "model": "gpu-dense", "quant": "Q4_K_M",
|
|
"parallel": 1, "role": "heavy-reasoning", "vram_mb": 24576,
|
|
"target_tps": 75.0
|
|
},
|
|
"ct110-rtx5070": {
|
|
"context": 131072, "model": "qwen3.5-9b-vision", "quant": "Q4_0 KV",
|
|
"parallel": 2, "role": "vision-web", "vram_mb": 12227,
|
|
"target_tps": 76.0
|
|
},
|
|
"strix-halo": {
|
|
"context": 262144, "model": "strix-moe", "quant": "Q4_K_M",
|
|
"parallel": 2, "role": "compression", "vram_mb": 65536,
|
|
"target_tps": 70.0
|
|
},
|
|
}
|
|
|
|
for gpu_key, cfg in gpu_config.items():
|
|
bench = benchmarks.get(gpu_key, {})
|
|
current_tps = bench.get("latest", {}).get("gen_tok_per_sec", 0)
|
|
baseline_tps = bench.get("baseline_tok_per_sec", 0)
|
|
|
|
if current_tps <= 0 or baseline_tps <= 0:
|
|
results.append({
|
|
"gpu": gpu_key, "model": cfg["model"], "role": cfg["role"],
|
|
"context": cfg["context"], "tok_per_sec": current_tps or 0,
|
|
"recommendation": f"{cfg['model']} idle — no benchmark data this cycle"
|
|
})
|
|
continue
|
|
|
|
perf_ratio = current_tps / baseline_tps
|
|
recommendation = None
|
|
|
|
if gpu_key == "ct8-rtx3090":
|
|
# Already at 256K — check if holding performance
|
|
if perf_ratio >= 0.95:
|
|
recommendation = f"256K ctx: {current_tps} tok/s ({(perf_ratio*100):.0f}% of baseline). Optimal for heavy reasoning."
|
|
else:
|
|
recommendation = f"256K ctx: {current_tps} tok/s ({(perf_ratio*100):.0f}% of baseline). Consider reducing parallel or ctx."
|
|
elif gpu_key == "ct110-rtx5070":
|
|
# 131K ctx, vision/web role — must preserve speed
|
|
if perf_ratio >= 0.95:
|
|
recommendation = f"131K ctx: {current_tps} tok/s (optimal). Vision/web role — context is fine."
|
|
else:
|
|
recommendation = f"131K ctx: {current_tps} tok/s ({(perf_ratio*100):.0f}% of baseline). Check VRAM pressure."
|
|
elif gpu_key == "strix-halo":
|
|
# 256K ctx, compression role — maximize context
|
|
if perf_ratio >= 0.95:
|
|
recommendation = f"256K ctx: {current_tps} tok/s ({(perf_ratio*100):.0f}% of baseline). Above target. Compression-optimized."
|
|
else:
|
|
recommendation = f"256K ctx: {current_tps} tok/s ({(perf_ratio*100):.0f}% of baseline). Check for competing workloads."
|
|
|
|
results.append({
|
|
"gpu": gpu_key, "model": cfg["model"], "role": cfg["role"],
|
|
"context": cfg["context"], "parallel": cfg["parallel"],
|
|
"tok_per_sec": current_tps, "baseline_tps": baseline_tps,
|
|
"perf_ratio": round(perf_ratio, 3),
|
|
"recommendation": recommendation
|
|
})
|
|
print(f" 📐 RULE 9: {gpu_key} ({cfg['role']}) — {recommendation}")
|
|
|
|
return results
|
|
|
|
def evaluate_workload_distribution(data):
|
|
"""
|
|
Rule 10: Workload Distribution Optimization
|
|
Check if GPU workload patterns match designated roles.
|
|
"""
|
|
role_map = {
|
|
"heavy-reasoning": ["gpu-dense", "syslog-auto"],
|
|
"vision-web": ["gpu-vision", "qwen3.5-9b-vision"],
|
|
"compression": ["strix-moe"],
|
|
}
|
|
gpu_roles = {
|
|
"NVIDIA GeForce RTX 3090": "heavy-reasoning",
|
|
"NVIDIA GeForce RTX 5070": "vision-web",
|
|
"AMD Strix Halo": "compression",
|
|
}
|
|
|
|
gpus = data.get("gpus", [])
|
|
summary = data.get("summary", {})
|
|
results = []
|
|
|
|
for gpu in gpus:
|
|
if not isinstance(gpu, dict):
|
|
continue
|
|
name = gpu.get("gpu_name", "")
|
|
role = "unknown"
|
|
for pattern, r in gpu_roles.items():
|
|
if pattern in name:
|
|
role = r
|
|
break
|
|
|
|
vram_used = gpu.get("vram_used_mb", 0)
|
|
vram_total = gpu.get("vram_total_mb", 0)
|
|
vram_pct = (vram_used / vram_total * 100) if vram_total else 0
|
|
|
|
status = "optimal"
|
|
note = ""
|
|
if role == "heavy-reasoning" and vram_pct > 90:
|
|
status = "warning"
|
|
note = f"VRAM at {vram_pct:.0f}% — consider reducing parallel or context"
|
|
elif role == "vision-web" and vram_pct > 88:
|
|
status = "warning"
|
|
note = f"VRAM at {vram_pct:.0f}% — vision/web role is VRAM-tight"
|
|
elif role == "compression" and vram_pct < 50:
|
|
note = f"VRAM at {vram_pct:.0f}% — plenty of headroom for compression workloads"
|
|
else:
|
|
note = f"VRAM at {vram_pct:.0f}% — role-appropriate"
|
|
|
|
results.append({
|
|
"gpu": name, "role": role, "vram_pct": round(vram_pct, 1),
|
|
"status": status, "note": note
|
|
})
|
|
icon = "✅" if status == "optimal" else "⚠"
|
|
print(f" {icon} RULE 10: {name} → {role} ({note})")
|
|
|
|
return results
|
|
|
|
def run_inference_test(model="syslog-auto"):
|
|
"""Run a quick inference test through LiteLLM."""
|
|
try:
|
|
body = json.dumps({
|
|
"model": model,
|
|
"messages": [{"role": "user", "content": "ping"}],
|
|
"max_tokens": 5
|
|
}).encode()
|
|
req = Request("http://192.168.68.116:4000/v1/chat/completions",
|
|
data=body, headers={
|
|
"Content-Type": "application/json",
|
|
"Authorization": "Bearer " + os.environ.get("LITELLM_API_KEY", "")
|
|
})
|
|
with urlopen(req, timeout=30) as resp:
|
|
return resp.status == 200
|
|
except:
|
|
return False
|
|
|
|
def main():
|
|
print(f"[{datetime.now().isoformat()}] GPU Self-Heal — {RUN_ID}")
|
|
print("=" * 60)
|
|
|
|
# Phase 1: Fetch data
|
|
data = fetch_gpu_data()
|
|
if not data:
|
|
print(" GPU monitor unreachable — aborting")
|
|
return 1
|
|
|
|
summary = data.get("summary", {})
|
|
print(f" Fleet: {summary.get('fleet_status','?')} | "
|
|
f"GPUs: {summary.get('gpu_count',0)} | "
|
|
f"Errors: {summary.get('gpu_errors',0)} | "
|
|
f"Alerts: {len(data.get('alerts',[]))}")
|
|
print()
|
|
|
|
# Phase 2: Evaluate rules (1-8)
|
|
history = load_history()
|
|
actions = evaluate_rules(data, history)
|
|
|
|
# Phase 2b: Context optimization (Rule 9)
|
|
ctx_results = evaluate_context_optimization(data)
|
|
|
|
# Phase 2c: Workload distribution (Rule 10)
|
|
wl_results = evaluate_workload_distribution(data)
|
|
|
|
save_history(history)
|
|
|
|
# Phase 3: Execute actions
|
|
issues_found = len(actions)
|
|
issues_fixed = 0
|
|
|
|
for action in actions:
|
|
severity = action.get("severity", "info")
|
|
if severity == "critical":
|
|
# Attempt fix
|
|
if "restart-llama-server" in action.get("action", ""):
|
|
gpu_name = action.get("gpu", "")
|
|
host = GPU_HOSTS.get(gpu_name.replace("NVIDIA ", "").replace("AMD ", "").lower(), {})
|
|
# Actual restart would need SSH — placeholder for now
|
|
print(f" ⟳ Would restart llama-server on {host.get('ip','?')} for {gpu_name}")
|
|
issues_fixed += 1
|
|
|
|
# Phase 4: Verify
|
|
all_models_ok = run_inference_test()
|
|
print(f"\n Verification: inference test {'✅' if all_models_ok else '❌'}")
|
|
|
|
# Phase 5: Compile report
|
|
status = "healthy"
|
|
if issues_found > 0:
|
|
status = "degraded" if issues_found <= 2 else "down"
|
|
|
|
report = json.dumps({
|
|
"run_id": RUN_ID,
|
|
"timestamp": datetime.now(timezone.utc).isoformat(),
|
|
"fleet_status": summary.get("fleet_status", "?"),
|
|
"overall_status": status,
|
|
"issues_found": issues_found,
|
|
"issues_fixed": issues_fixed,
|
|
"actions": actions,
|
|
"context_optimization": ctx_results,
|
|
"workload_distribution": wl_results
|
|
}, indent=2, default=str)
|
|
|
|
print(f"\n Status: {status} | Issues: {issues_found} | Fixed: {issues_fixed}")
|
|
|
|
# Phase 6: Log to Gitea (hard rule: never to knowledge graph)
|
|
kg_title = f"[GPU-SELF-HEAL] {RUN_ID} — {status}"
|
|
kg_desc = f"GPU self-heal run: {issues_found} issues found, {issues_fixed} fixed. Fleet: {summary.get('fleet_status','?')}."
|
|
kg_result = kg_create_node(kg_title, kg_desc, report)
|
|
print(f" Gitea: {kg_result[:120] if kg_result else 'write attempted'}")
|
|
|
|
print(f"\n[{datetime.now().isoformat()}] Complete — {RUN_ID}")
|
|
return 0 if issues_found == 0 else 1
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|