Files

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())