diff --git a/README.md b/README.md index 1b23626..8adc719 100644 --- a/README.md +++ b/README.md @@ -107,8 +107,8 @@ only if the card actually needs it. * **Server-Sent Events (SSE)**: A single background sampler produces one 1 Hz snapshot and fans it out to every subscriber via `GET /api/stream`. Previously each connected client independently re-ran the whole snapshot — NVML, `/proc/meminfo`, an HTTP round-trip each to Ollama and ComfyUI, and a recursive walk of the ComfyUI models tree with a `stat()` per checkpoint — once per second, so opening the dashboard in three tabs tripled the load on the thing it was measuring. Slow clients drop stale frames instead of stalling the sampler. ### 🤖 Model Context Protocol (MCP 2.0) Server -* **12 Native Agentic Tools**: Allows AI agents (Antigravity CLI, Claude Desktop, Cursor) to manage GPU resources, trigger model hot-swaps, tune fan curves, and inspect telemetry. -* **3 Live MCP Resources**: Exposes live metrics, model catalogs, and switch logs as streamable resources (`gpu://metrics/live`, `gpu://models/catalog`, `gpu://history/switches`). +* **23 Native Agentic Tools**: Allows AI agents (Antigravity CLI, Claude Desktop, Cursor) to manage GPU resources, trigger model hot-swaps, measure page-cache residency, read persisted performance analytics, drive the thermal governor, and run overclock sweeps. +* **6 Live MCP Resources**: Live metrics, model catalog, switch log, measured cache residency, per-profile analytics, and the overclock profiles with the evidence behind each setting. * **Dual Transport Support**: Run via standard input/output (`--stdio`) or network Server-Sent Events (`--sse --port 8001`). ### ⏱️ Automated Latency & Throughput Benchmark Engine @@ -135,19 +135,32 @@ flowchart TD REST["REST API & OpenAPI Docs"] MCP["Model Context Protocol (MCP 2.0)"] SSE["1Hz Real-Time SSE Stream"] - Arbitrator["VRAM Arbitrator (15ms Soft-Yield)"] + Arbitrator["VRAM Arbitrator (confirmed yield)"] Overclock["Overclock & Fan Manager"] Warmer["Page Cache Pre-Warmer"] end - HostRAM <== "PCIe 4.0 x16 Bus (~31.5 GB/s Hot-Swap)" ==> GPU + HostRAM <== "PCIe 4.0 x16 Bus (measured 2.6 GB/s warm model load)" ==> GPU Orchestrator --> GPU Orchestrator --> HostRAM ``` ### The Physics of Sub-Second Switching * **Host RAM as Staging**: Active LLMs and diffusion checkpoints remain resident in the 64GB Linux Page Cache. -* **PCIe 4.0 x16 Hot-Swapping**: Transferring weights across PCIe 4.0 x16 achieves **~31.5 GB/s** bandwidth, reducing model loads from 30+ seconds (disk) to **under 1.5 seconds**. +* **Warm vs cold model loads, measured.** The same 12.87 GB model, loaded through Ollama on this box: + + | Page-cache residency | Load time | Effective rate | + | :--- | :--- | :--- | + | 3.1% (dropped with `FADV_DONTNEED`) | 34.3 s | 0.38 GB/s | + | 100% (force-warmed) | 4.9 s | 2.63 GB/s | + + A **6.9× speedup**, and the reason the page cache matters. Note the effective rate is + well below the PCIe 4.0 x16 bus rate and below the 6.4 GB/s the page cache itself + reads at: Ollama's `load_duration` also covers host-to-device transfer and model + initialisation, not just the file read. Classification thresholds are calibrated + against these measured numbers rather than the theoretical bus bandwidth — an earlier + 5 GB/s cache-hit bar sat above what a fully warm load can even achieve, so every warm + load was misreported as a partial hit. * **Soft-Yielding**: Dropping Ollama's VRAM allocation via `keep_alive: 0` preserves the weights in host RAM. Measured on this box: the HTTP request returns in **~63 ms**, and the driver finishes releasing 14.9 GB **~77 ms after that**. HyperSwap waits for the second number before handing VRAM to ComfyUI — the earlier "~15 ms" figure timed the request, not the release. --- @@ -220,19 +233,33 @@ HyperSwap includes a native **MCP 2.0 server** (`mcp_server.py`) exposing orches | **`set_gpu_fan_speed`** | `mode` (str), `percent` (optional int) | Sets fan speed mode (`auto`\|`manual`) and target PWM % (30–100%). | | **`get_host_memory_status`** | *None* | 64GB host RAM breakdown, active page cache size, and cache ratio. | | **`switch_ollama_model`** | `model_name` (str), `keep_alive` (str) | Hot-swaps active LLM in VRAM, measures latency (ms) and tokens/sec. | -| **`soft_yield_ollama_vram`** | `model_name` (optional str) | Yields Ollama VRAM to 0 MB in ~15ms while keeping model weights in RAM cache. | +| **`soft_yield_ollama_vram`** | `model_name` (optional str) | Yields Ollama VRAM to 0 MB and waits for NVML to confirm the driver actually released it. Returns the request/confirm split. | | **`purge_comfyui_vram`** | *None* | Purges loaded diffusion models from ComfyUI pipeline VRAM. | -| **`prewarm_all_models_to_ram`** | *None* | Faults all local LLM and diffusion checkpoints into Linux OS page cache. | +| **`prewarm_all_models_to_ram`** | *None* | Warms the highest-value models into page cache within a byte budget, skipping what is already resident. | | **`prewarm_single_model`** | `model_name` (optional str), `filepath` (optional str) | Pre-warms a single GGUF or Safetensors file into RAM. | | **`list_available_models`** | *None* | Lists all installed Ollama models and discovered ComfyUI Safetensors on disk. | | **`get_switch_history`** | `limit` (int, default 20) | Retrieves recent switch events, millisecond latencies, and RAM hit status. | | **`run_model_switch_benchmark`**| `iterations` (int, default 2) | Automated round-trip latency benchmark between installed models. | +| **`get_page_cache_residency`** | `include_files` (bool) | Measured page-cache residency per model file, with the measurement method used for each. | +| **`get_warm_plan`** | `budget_gb` (optional float) | Previews what warming would read and skip, ranked by recency/frequency. Does not warm. | +| **`request_vram_for_ollama`** | `needed_gb` (float) | Purges ComfyUI's checkpoints immediately if VRAM headroom is short, bypassing the idle timer. | +| **`get_profile_performance`** | `days` (float, default 7) | Measured tok/s and thermals per overclock profile, from persisted history. | +| **`get_thermal_governor_status`** | *None* | Current derate level, the reason for it, and escalation history. | +| **`set_thermal_governor`** | `enabled` (optional bool), `reset` (bool) | Enable/disable the governor, or clear an active derate. | +| **`get_overclock_status`** | *None* | Active profile, all profiles with their evidence, and which levers this driver honours. | +| **`apply_overclock_profile`** | `profile` (str) | Apply `ollama` \| `comfy` \| `balanced`. | +| **`restore_stock_gpu_state`** | *None* | Drop clock locks and offsets, restore default power limit, return fans to automatic. | +| **`run_overclock_sweep`** | `knob`, `profile`, `workload`, `start`, `stop`, `repeats`, `apply_best` | Sweep a knob against a real workload and report the fastest stable value. Verifies the knob moves the hardware first. Takes minutes. | +| **`get_autotune_status`** | *None* | Sweep progress, the last result table, and all recorded autotune steps. | ### MCP Resources List * `gpu://metrics/live`: Real-time snapshot of GPU sensors and RAM page cache. * `gpu://models/catalog`: Catalog of all discovered GGUF and Safetensors models. * `gpu://history/switches`: Event log of recent model transitions and swap speeds. +* `gpu://cache/residency`: Measured page-cache residency across every model on disk. +* `gpu://analytics/profiles`: Measured throughput and thermals per overclock profile. +* `gpu://overclock/profiles`: Overclock profiles including the measurement behind each setting. --- diff --git a/mcp_server.py b/mcp_server.py index 7851760..47e037a 100644 --- a/mcp_server.py +++ b/mcp_server.py @@ -7,16 +7,19 @@ import logging from typing import Dict, List, Any, Optional from mcp.server import MCPServer -import ram_optimizer -import vram_arbitrator +import autotune import overclock_manager +import ram_optimizer +import telemetry_store +import thermal_governor +import vram_arbitrator logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s") logger = logging.getLogger("gpu_swapper_mcp") mcp = MCPServer( name="gpu-program-swapper", - version="1.0.0", + version="2.0.0", description="Orchestrates high-speed GPU VRAM hot-swaps between Ollama LLMs and ComfyUI with 64GB RAM cache telemetry." ) @@ -133,6 +136,105 @@ def set_gpu_fan_speed(mode: str = "auto", percent: Optional[int] = None) -> str: res = overclock_manager.set_fan_auto() return json.dumps(res, indent=2) +@mcp.tool() +def get_page_cache_residency(include_files: bool = True) -> str: + """Measure how much of each model on disk is genuinely resident in the Linux page cache. + + Uses cachestat(2) where the kernel permits it and a read-rate probe where it does not + (Ollama blobs are owned by another user). Reports which method was used per file, and + marks anything it cannot measure rather than guessing.""" + report = ram_optimizer.get_cache_report(include_files=include_files) + report["capability"] = ram_optimizer.residency_capability() + return json.dumps(report, indent=2, default=str) + + +@mcp.tool() +def get_warm_plan(budget_gb: Optional[float] = None) -> str: + """Preview which models pre-warming would load into RAM, in what order, and what it + would skip — ranked by recency/frequency and capped by a byte budget. Does not warm.""" + return json.dumps(ram_optimizer.build_warm_plan(budget_gb), indent=2, default=str) + + +@mcp.tool() +async def request_vram_for_ollama(needed_gb: float = 0.0) -> str: + """Free VRAM for an LLM right now: purges ComfyUI's cached checkpoints immediately if + there is not enough headroom, instead of waiting for the normal idle timer.""" + res = await vram_arbitrator.arbitrator.request_vram_for_ollama(needed_gb) + return json.dumps(res, indent=2, default=str) + + +@mcp.tool() +def get_profile_performance(days: float = 7.0) -> str: + """Compare measured decode throughput and thermals per overclock profile, from + persisted history. Answers whether a given profile is actually delivering more tok/s.""" + return json.dumps({ + "window_days": days, + "profiles": telemetry_store.profile_comparison(days), + "swaps": telemetry_store.swap_stats(days), + }, indent=2, default=str) + + +@mcp.tool() +def get_thermal_governor_status() -> str: + """Current thermal derate level, why it was applied, and the escalation history.""" + return json.dumps(thermal_governor.governor.get_status(), indent=2, default=str) + + +@mcp.tool() +def set_thermal_governor(enabled: Optional[bool] = None, reset: bool = False) -> str: + """Enable or disable the thermal governor, or clear an active derate and reapply the + full profile.""" + if enabled is not None: + thermal_governor.governor.set_enabled(enabled) + if reset: + thermal_governor.governor.reset() + return json.dumps(thermal_governor.governor.get_status(), indent=2, default=str) + + +@mcp.tool() +def get_overclock_status() -> str: + """Active overclock profile, all profiles with the evidence behind their settings, and + which hardware levers this driver actually honours (clock offsets are ignored on some).""" + return json.dumps(overclock_manager.get_status(), indent=2, default=str) + + +@mcp.tool() +def apply_overclock_profile(profile: str) -> str: + """Apply an overclock profile by name: ollama | comfy | balanced.""" + return json.dumps(overclock_manager.apply_profile(profile), indent=2, default=str) + + +@mcp.tool() +def restore_stock_gpu_state() -> str: + """Drop all clock locks and offsets, restore the default power limit, and return the + fans to automatic control.""" + return json.dumps(overclock_manager.restore_safe("MCP request"), indent=2, default=str) + + +@mcp.tool() +async def run_overclock_sweep(knob: str = "power_limit_w", profile: str = "ollama", + workload: str = "auto", start: Optional[int] = None, + stop: Optional[int] = None, repeats: int = 1, + apply_best: bool = False) -> str: + """Sweep one GPU knob against a real workload and report the fastest stable value. + + knob: power_limit_w | lock_mem_mhz | lock_core_max | mem_offset_mhz | core_offset_mhz + workload: 'ollama' (decode tok/s), 'comfy' (SDXL it/s), or 'auto' to match the profile. + + Verifies the knob actually moves the hardware before sweeping, refuses to run while + ComfyUI is busy, and always restores the original profile. Takes minutes.""" + res = await autotune.sweep(knob=knob, profile=profile, workload=workload, + start=start, stop=stop, repeats=repeats, + apply_best=apply_best) + return json.dumps(res, indent=2, default=str) + + +@mcp.tool() +def get_autotune_status() -> str: + """Sweep progress, the last sweep's full result table, and every recorded autotune step.""" + return json.dumps(autotune.get_status(), indent=2, default=str) + + # ========================================== # MCP RESOURCES # ========================================== @@ -154,6 +256,21 @@ def get_switch_history_resource() -> str: """Recent model switch events and latencies.""" return json.dumps(vram_arbitrator.get_switch_history(), indent=2) +@mcp.resource("gpu://cache/residency") +def get_cache_residency_resource() -> str: + """Measured page-cache residency across every model on disk.""" + return json.dumps(ram_optimizer.get_cache_report(include_files=True), indent=2, default=str) + +@mcp.resource("gpu://analytics/profiles") +def get_profile_analytics_resource() -> str: + """Measured throughput and thermals per overclock profile, from persisted history.""" + return json.dumps(telemetry_store.profile_comparison(7.0), indent=2, default=str) + +@mcp.resource("gpu://overclock/profiles") +def get_overclock_profiles_resource() -> str: + """Overclock profiles, including the measurement recorded behind each setting.""" + return json.dumps(overclock_manager.get_status(), indent=2, default=str) + if __name__ == "__main__": import argparse @@ -163,6 +280,11 @@ if __name__ == "__main__": parser.add_argument("--port", type=int, default=8001, help="Port for SSE transport") args = parser.parse_args() + # Only when run as a standalone server. server.py imports this module for the + # benchmark tool, and starting the store at import scope would spin up a writer as a + # side effect of that import. + telemetry_store.start() + if args.sse: mcp.run(transport="sse", host="0.0.0.0", port=args.port) else: diff --git a/ram_optimizer.py b/ram_optimizer.py index bf0bff2..0a452ef 100644 --- a/ram_optimizer.py +++ b/ram_optimizer.py @@ -143,7 +143,7 @@ PROBE_WINDOW_BYTES = 2 * 1024 * 1024 PROBE_CACHED_GBPS = 1.5 -def _throughput_probe(fd: int, size: int) -> Dict[str, Any]: +def _throughput_probe(fd: int, size: int, windows_override: Optional[int] = None) -> Dict[str, Any]: """Infer residency by timing reads of small windows spread across the file. Used only where cachestat is not permitted (Ollama's blobs are owned by uid `ollama`). @@ -158,7 +158,7 @@ def _throughput_probe(fd: int, size: int) -> Dict[str, Any]: pollution the probe itself created, and leaving them behind would slowly warm the cache with data nobody asked for. """ - windows = min(PROBE_WINDOWS, max(int(size // PROBE_WINDOW_BYTES), 1)) + windows = min(windows_override or PROBE_WINDOWS, max(int(size // PROBE_WINDOW_BYTES), 1)) if windows <= 0: return {"resident_pct": 0.0, "windows": 0} @@ -194,7 +194,8 @@ def _throughput_probe(fd: int, size: int) -> Dict[str, Any]: } -def page_residency(filepath: str, allow_probe: bool = True) -> Dict[str, Any]: +def page_residency(filepath: str, allow_probe: bool = True, + probe_windows: Optional[int] = None) -> Dict[str, Any]: """Measure what fraction of a file is resident in the Linux page cache.""" try: size = os.path.getsize(filepath) @@ -216,7 +217,7 @@ def page_residency(filepath: str, allow_probe: bool = True) -> Dict[str, Any]: method, measurable = "cachestat", True extra = {"dirty_pages": cs.nr_dirty, "evicted_pages": cs.nr_evicted} elif allow_probe: - probe = _throughput_probe(fd, size) + probe = _throughput_probe(fd, size, probe_windows) pct = probe["resident_pct"] method, measurable = "probe", True extra = {"probe_windows": probe["windows"], "probe_median_gbps": probe.get("median_gbps")} @@ -236,6 +237,11 @@ def page_residency(filepath: str, allow_probe: bool = True) -> Dict[str, Any]: "method": method, "measurable": measurable, "warm": pct >= WARM_SKIP_THRESHOLD_PCT, + # Only an exact measurement is trustworthy enough to skip work on. A probe of a + # dozen 2 MB windows can clear 90% on a file that is mostly cold -- observed + # here as a 12.87 GB "already resident" blob that then loaded at 2.44 GB/s. + "warm_confident": (method == "cachestat" and pct >= WARM_SKIP_THRESHOLD_PCT) + or (method == "probe" and pct >= 100.0), **extra, } except Exception as e: @@ -468,16 +474,19 @@ def get_cache_report(include_files: bool = True, force_refresh: bool = False) -> # ---------------------------------------------------------------- warming def warm_file_to_ram(filepath: str, chunk_size: int = 16 * 1024 * 1024, - skip_if_warm: bool = True) -> Dict[str, Any]: - """Pre-fault a file into the Linux page cache, skipping it if already resident.""" + skip_if_warm: bool = True, force: bool = False) -> Dict[str, Any]: + """Pre-fault a file into the Linux page cache, skipping it only if confidently resident.""" if not os.path.exists(filepath): return {"success": False, "error": f"File not found: {filepath}", "duration_ms": 0} - before = page_residency(filepath) - if skip_if_warm and before.get("warm"): + # Probe densely here: this decision skips real work, so it is worth 32 samples + # rather than 12. + before = page_residency(filepath, probe_windows=32) + if skip_if_warm and not force and before.get("warm_confident"): return { "success": True, "filepath": filepath, "skipped": True, "reason": "already resident", "resident_pct": before.get("resident_pct"), + "method": before.get("method"), "size_mb": round(before.get("size_bytes", 0) / (1024**2), 2), "duration_ms": 0.0, "bytes_read": 0, } @@ -499,7 +508,7 @@ def warm_file_to_ram(filepath: str, chunk_size: int = 16 * 1024 * 1024, bytes_read += n duration = time.perf_counter() - t0 - after = page_residency(filepath) + after = page_residency(filepath, probe_windows=32) return { "success": True, "filepath": filepath, @@ -551,11 +560,11 @@ async def warm_ollama_model(model_name: str, keep_alive: str = "5m") -> Dict[str "duration_ms": round((time.perf_counter() - t0) * 1000, 2)} -def warm_ollama_blob(model_name: str) -> Dict[str, Any]: +def warm_ollama_blob(model_name: str, force: bool = False) -> Dict[str, Any]: """Warm a specific Ollama model's GGUF into page cache without touching VRAM.""" for f in find_ollama_model_files(): if f["model"] == model_name: - res = warm_file_to_ram(f["full_path"]) + res = warm_file_to_ram(f["full_path"], force=force) res["model"] = model_name return res return {"success": False, "error": f"no blob found for model '{model_name}'"} @@ -603,13 +612,13 @@ def build_warm_plan(budget_gb: Optional[float] = None) -> Dict[str, Any]: if c["full_path"] in seen_paths: continue seen_paths.add(c["full_path"]) - res = page_residency(c["full_path"]) + res = page_residency(c["full_path"], probe_windows=32) entry = { "name": c["name"], "kind": c["kind"], "full_path": c["full_path"], "size_gb": c.get("size_gb", 0), "score": round(c["score"], 4), "resident_pct": res.get("resident_pct", 0.0), } - if res.get("warm"): + if res.get("warm_confident"): entry["action"] = "already-warm" skipped.append(entry) continue diff --git a/server.py b/server.py index 07eab1e..e89d7d1 100644 --- a/server.py +++ b/server.py @@ -227,6 +227,7 @@ class WarmRequest(BaseModel): model_name: Optional[str] = Field(None, description="Ollama model name to warm", example="gemma4:26b") filepath: Optional[str] = Field(None, description="Absolute file path of Safetensors/GGUF to warm into RAM") blob_only: bool = Field(False, description="Warm the model's weights into page cache without loading VRAM") + force: bool = Field(False, description="Warm even if residency sampling thinks it is already resident") class WarmAllRequest(BaseModel): budget_gb: Optional[float] = Field(None, description="Byte budget for warming; defaults to 70% of MemAvailable", example=24.0) @@ -362,11 +363,11 @@ async def api_warm_plan(budget_gb: Optional[float] = Query(None, description="Ov async def api_warm_model(req: WarmRequest): """Pre-warm a specific Ollama model or file path into the Linux page cache.""" if req.model_name and req.blob_only: - return ram_optimizer.warm_ollama_blob(req.model_name) + return ram_optimizer.warm_ollama_blob(req.model_name, force=req.force) if req.model_name: return await ram_optimizer.warm_ollama_model(req.model_name, keep_alive="1m") if req.filepath: - return ram_optimizer.warm_file_to_ram(req.filepath) + return ram_optimizer.warm_file_to_ram(req.filepath, force=req.force) raise HTTPException(status_code=400, detail="model_name or filepath required") @app.get("/api/cache/report", summary="Measured Page-Cache Residency", tags=["Memory Optimization"]) diff --git a/vram_arbitrator.py b/vram_arbitrator.py index 74752d6..9610080 100644 --- a/vram_arbitrator.py +++ b/vram_arbitrator.py @@ -29,10 +29,20 @@ COMFY_API_BASE = "http://127.0.0.1:8188" # Circular buffer for transition events (the durable log lives in telemetry_store) SWITCH_HISTORY = deque(maxlen=50) -# Bandwidth thresholds used to classify how a model actually got into VRAM. -# PCIe 4.0 x16 tops out near 31.5 GB/s; this NVMe sustains well under 2 GB/s. -RAM_HIT_GBPS = 5.0 -PARTIAL_HIT_GBPS = 1.5 +# Bandwidth thresholds for classifying how a model reached VRAM, calibrated by measuring +# the same 12.87 GB model loaded cold and warm on this box (2026-08-28): +# +# 3.1% resident -> 34.3 s -> 0.38 GB/s +# 100% resident -> 4.9 s -> 2.63 GB/s +# +# The first cut at these numbers assumed a page-cache-fed load would approach the bus +# rate and set the cache-hit bar at 5 GB/s. It does not: Ollama's load_duration covers +# host-to-device transfer and model initialisation as well as the file read, so a fully +# resident model still reports ~2.6 GB/s while the page cache itself reads at 6.4 GB/s. +# A 5 GB/s bar could therefore never be met, and every warm load was being reported as +# a partial hit. Thresholds now sit either side of the measured 6.9x separation. +RAM_HIT_GBPS = 2.0 +PARTIAL_HIT_GBPS = 0.8 # How long Ollama's VRAM may take to actually drain before we stop waiting. # Ollama will not unload a model while a generation is in flight, so a short ceiling @@ -557,11 +567,14 @@ def classify_load(size_bytes: int, load_duration_ms: float) -> Dict[str, Any]: return {"cache_status": "Already in VRAM", "load_gbps": None, "is_ram_hit": True} if not size_bytes: # No size on record — fall back to the old heuristic, but say so. + # Without a size we cannot compute bandwidth at all; this is a guess and is + # labelled as one. 8s roughly splits the measured warm (4.9s) and cold (34.3s) + # loads for a mid-size model, but it is meaningless for very small or large ones. return { - "cache_status": "RAM Cache Hit ⚡" if load_duration_ms < 2500 else "Cold Disk Load 💾", + "cache_status": "RAM Cache Hit ⚡" if load_duration_ms < 8000 else "Cold Disk Load 💾", "load_gbps": None, - "is_ram_hit": load_duration_ms < 2500, - "detail": "size unknown, fell back to duration heuristic", + "is_ram_hit": load_duration_ms < 8000, + "detail": "size unknown, fell back to a duration guess", } gbps = (size_bytes / (1024**3)) / (load_duration_ms / 1000.0) if gbps >= RAM_HIT_GBPS: