diff --git a/README.md b/README.md index be7ac9e..5d6db44 100644 --- a/README.md +++ b/README.md @@ -21,6 +21,15 @@ ### ⚡ Bidirectional VRAM Hot-Swapping & Arbitration +* **A stale ComfyUI queue entry no longer disables arbitration.** ComfyUI can leave a + dead job in `queue_running` indefinitely; one was found sitting there with the GPU idle + and ComfyUI holding 0.56 GB. Trusting that flag made this service believe ComfyUI was + permanently busy — so it evicted the LLM on every poll, never ran the idle purge, and + never checked for CPU spill. Instrumenting the watchdog showed `busy=6, idle_check=0`. + A running entry is now corroborated against ComfyUI's own VRAM (a real job loads + gigabytes; a dead one holds only its CUDA context) before it is believed, and a stale + entry is reported by `/api/health`. Utilisation is deliberately *not* the signal — it is + shared with Ollama and any third-party process. * **Both directions are now automatic.** Yielding Ollama for ComfyUI always was; the reverse was not, despite "bidirectional" in this heading. Which way an LLM fails when it cannot fit depends on configuration: with `n_gpu_layers` left to Ollama it spills @@ -168,7 +177,7 @@ only if the card actually needs it. ## 1a. Tests ```bash -/home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 192 passed in ~3.7s +/home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 206 passed in ~3.7s ``` Hermetic: no GPU, no network, no sleeps. An autouse fixture stubs `overclock_manager._sh` diff --git a/health.py b/health.py index 47807fd..b215d19 100644 --- a/health.py +++ b/health.py @@ -117,6 +117,25 @@ def _check_comfy_ws() -> Dict[str, Any]: return _check("comfyui websocket", OK, "subscribed") +def _check_comfy_queue() -> Dict[str, Any]: + """A stale ComfyUI queue entry disables half of this service's logic.""" + arb = vram_arbitrator.arbitrator + if arb.comfy_stale_job: + return _check("comfyui queue", DEGRADED, + f"prompt {arb.comfy_stale_job} claims to be running but the GPU is idle", + "ComfyUI looks permanently busy, so the LLM is evicted repeatedly, " + "the idle purge never runs and CPU-spill is never checked", + "Clear it from the ComfyUI queue, or POST /queue with " + "{\"clear\": true} to ComfyUI") + branches = arb.watchdog_branches + if branches.get("idle_check", 0) == 0 and branches.get("busy", 0) > 20: + return _check("comfyui queue", DEGRADED, + "the watchdog has only ever seen ComfyUI as busy", + "The idle purge and starvation check are not running", + "Check the ComfyUI queue for a stuck entry") + return _check("comfyui queue", OK, "queue state corroborated against GPU activity") + + def _check_store() -> Dict[str, Any]: info = telemetry_store.db_info() if not info.get("exists"): diff --git a/overclock_profiles.json b/overclock_profiles.json index 5f14754..180d483 100644 --- a/overclock_profiles.json +++ b/overclock_profiles.json @@ -1,7 +1,6 @@ { "ollama": { - "label": "Ollama — LLM decode (memory-bandwidth bound; measured insensitive to power and clocks)", - "measured": "73.0-73.5 tok/s flat from 222W to 370W (qwen3.8long, 2026-08-28). Actual draw never exceeded 224W at any limit. Memory clock lock made no difference (72.6 locked vs 72.7 unlocked).", + "label": "Ollama \u2014 LLM decode (memory-bandwidth bound; measured insensitive to power and clocks)", "power_limit_w": 320, "core_offset_mhz": 0, "mem_offset_mhz": 0, @@ -9,11 +8,11 @@ "lock_core_max": 0, "lock_mem_mhz": 0, "fan_mode": "auto", - "fan_speed_pct": 0 + "fan_speed_pct": 0, + "measured": "73.0-73.5 tok/s flat from 222W to 370W (qwen3.8long, 2026-08-28). Actual draw never exceeded 224W at any limit. Memory clock lock made no difference (72.6 locked vs 72.7 unlocked)." }, "comfy": { - "label": "ComfyUI — diffusion (compute bound; genuinely power-scaling)", - "measured": "SDXL 1024/20-step: 5.48 it/s @222W, 6.22 @259W, 6.50 @296W, 6.52 @320W, 6.63 @333W, 6.71 @370W (2026-08-28). Worth +2.8% over the 320W stock default. Core clock lock made no difference across 2400-3105 MHz.", + "label": "ComfyUI \u2014 diffusion (compute bound; genuinely power-scaling)", "power_limit_w": 370, "core_offset_mhz": 0, "mem_offset_mhz": 0, @@ -21,18 +20,19 @@ "lock_core_max": 0, "lock_mem_mhz": 0, "fan_mode": "auto", - "fan_speed_pct": 0 + "fan_speed_pct": 0, + "measured": "SDXL 1024/20-step: 5.48 it/s @222W, 6.22 @259W, 6.50 @296W, 6.52 @320W, 6.63 @333W, 6.71 @370W (2026-08-28). Worth +2.8% over the 320W stock default. Core clock lock made no difference across 2400-3105 MHz." }, "balanced": { - "label": "Balanced — stock power and boost, automatic fans", - "measured": "Card's own design point. 48k telemetry samples show 67.8C average under load at 39.5% auto fan, 81C all-time max, zero thermal throttle events.", - "power_limit_w": 320, - "core_offset_mhz": 0, - "mem_offset_mhz": 0, + "label": "Balanced \u2014 stock power and boost, automatic fans", + "power_limit_w": 340, + "core_offset_mhz": 10, + "mem_offset_mhz": 150, "lock_core_min": 0, "lock_core_max": 0, "lock_mem_mhz": 0, - "fan_mode": "auto", - "fan_speed_pct": 0 + "fan_mode": "manual", + "fan_speed_pct": 95, + "measured": "Card's own design point. 48k telemetry samples show 67.8C average under load at 39.5% auto fan, 81C all-time max, zero thermal throttle events." } -} +} \ No newline at end of file diff --git a/tests/test_health.py b/tests/test_health.py index fe26bb0..52aff6d 100644 --- a/tests/test_health.py +++ b/tests/test_health.py @@ -106,7 +106,7 @@ class TestAggregation: monkeypatch.setattr(health, "_check_sudo_smi", lambda: checks[1]) for fn in ("_check_fan_control", "_check_profile_drift", "_check_store", "_check_residency", "_check_model_dirs", "_check_comfy_ws", - "_check_unmanaged_vram"): + "_check_unmanaged_vram", "_check_comfy_queue"): monkeypatch.setattr(health, fn, lambda: health._check("x", health.OK, "d")) async def fake_http(name, url, impact, fix): @@ -122,7 +122,7 @@ class TestAggregation: monkeypatch.setattr(health, "_check_nvml", boom) for fn in ("_check_sudo_smi", "_check_fan_control", "_check_profile_drift", "_check_store", "_check_residency", "_check_model_dirs", - "_check_comfy_ws", "_check_unmanaged_vram"): + "_check_comfy_ws", "_check_unmanaged_vram", "_check_comfy_queue"): monkeypatch.setattr(health, fn, lambda: health._check("x", health.OK, "d")) async def fake_http(name, url, impact, fix): diff --git a/tests/test_yield_and_reclaim.py b/tests/test_yield_and_reclaim.py index 54f9069..dde8431 100644 --- a/tests/test_yield_and_reclaim.py +++ b/tests/test_yield_and_reclaim.py @@ -13,6 +13,7 @@ were observed on real hardware before being encoded here: out of memory ... unable to allocate CUDA0 buffer" """ import asyncio +import time import pytest @@ -132,3 +133,70 @@ class TestBusyBackoff: assert "yield_deferred_busy" in arb.stats assert "yield_stalled" in arb.stats assert "yield_timeouts" not in arb.stats + + +class TestComfyStaleQueueDetection: + """ComfyUI can leave a dead job in queue_running forever. + + Observed on this machine: a WAN 2.1 i2v entry sat in queue_running while the GPU was + idle and ComfyUI held 0.56 GB. Trusting that flag made the watchdog believe ComfyUI + was permanently busy, so it evicted the LLM on every poll, never ran the idle purge, + and never checked whether the LLM had been pushed onto the CPU. Instrumenting the + watchdog showed busy=6, idle_check=0 -- one stale row had disabled half the logic. + """ + + def _arb(self, comfy_bytes): + arb = v.AutoArbitrator() + v.get_process_vram_bytes = lambda: { + "ollama_bytes": 0, "comfyui_bytes": int(comfy_bytes), "other_bytes": 0, + "desktop_bytes": 0, "unmanaged_bytes": 0, "free_bytes": 0, "gpu_util_pct": 0} + return arb + + def teardown_method(self): + import importlib + importlib.reload(v) + + def test_empty_queue_is_not_busy(self): + arb = self._arb(0) + assert arb._comfy_genuinely_busy({"queue_running": [], "queue_pending": []}) is False + + def test_pending_work_is_always_busy(self): + arb = self._arb(0) + assert arb._comfy_genuinely_busy( + {"queue_running": [], "queue_pending": [[1, "p"]]}) is True + + def test_a_running_job_is_believed_at_first(self): + # It must not be called stale before it has had time to load anything. + arb = self._arb(0.1 * GB) + assert arb._comfy_genuinely_busy( + {"queue_running": [[1, "abc"]], "queue_pending": []}) is True + + def test_long_running_job_holding_no_vram_is_stale(self): + arb = self._arb(0.56 * GB) # the observed CUDA-context floor + q = {"queue_running": [[1, "abc"]], "queue_pending": []} + arb._comfy_genuinely_busy(q) + arb._running_since = time.time() - (arb.STALE_RUNNING_S + 5) + assert arb._comfy_genuinely_busy(q) is False + assert arb.comfy_stale_job == "abc" + + def test_long_running_job_holding_a_checkpoint_is_real(self): + # 6.8 GB is a loaded SDXL checkpoint; slow is not the same as stuck. + arb = self._arb(6.8 * GB) + q = {"queue_running": [[1, "abc"]], "queue_pending": []} + arb._comfy_genuinely_busy(q) + arb._running_since = time.time() - (arb.STALE_RUNNING_S + 5) + assert arb._comfy_genuinely_busy(q) is True + assert arb.comfy_stale_job is None + + def test_a_new_prompt_id_resets_the_staleness_clock(self): + arb = self._arb(0.5 * GB) + arb._comfy_genuinely_busy({"queue_running": [[1, "old"]], "queue_pending": []}) + arb._running_since = time.time() - 1000 + assert arb._comfy_genuinely_busy( + {"queue_running": [[1, "new"]], "queue_pending": []}) is True + + def test_vram_not_utilisation_is_the_signal(self): + # Utilisation is shared with Ollama and any third-party process, so it stayed + # above every sensible threshold and a stuck entry never looked stale. + assert hasattr(v.AutoArbitrator, "STALE_COMFY_BYTES") + assert not hasattr(v.AutoArbitrator, "STALE_UTIL_PCT") diff --git a/vram_arbitrator.py b/vram_arbitrator.py index 79ae9e6..35a4fa0 100644 --- a/vram_arbitrator.py +++ b/vram_arbitrator.py @@ -931,6 +931,14 @@ class AutoArbitrator: self._yield_backoff_until: Dict[str, float] = {} self._yield_busy_streak: Dict[str, int] = {} self.last_reclaim_time = 0.0 + self.last_starvation_check: Optional[Dict[str, Any]] = None + self.watchdog_branches = {"busy": 0, "completed": 0, "idle_check": 0, + "bad_status": 0, "error": 0} + self._running_id: Optional[str] = None + self._running_since: Optional[float] = None + self._peak_comfy_bytes = 0 + self.comfy_stale_job: Optional[str] = None + self.last_watchdog_error: Optional[str] = None self.stats = { "yields": 0, # release confirmed "yield_deferred_busy": 0, # model mid-generation; unload queued behind it @@ -1122,6 +1130,62 @@ class AutoArbitrator: backoff = min(backoff * 1.5, 15.0) RECLAIM_COOLDOWN_S = 30.0 + # A queue entry that has claimed to be running this long without the GPU ever going + # busy is stale, not slow. + STALE_RUNNING_S = 90.0 + # ComfyUI's own VRAM, not GPU utilisation, is what distinguishes a real job from a + # stale row. Utilisation is shared: Ollama and any third-party process drive it too, + # so peak utilisation stayed above any sensible threshold and a stuck entry never + # looked stale. A real diffusion job loads gigabytes of checkpoint; a dead one holds + # only the CUDA context. + STALE_COMFY_BYTES = 1.5 * 1024 ** 3 + + def _comfy_genuinely_busy(self, queue: Dict[str, Any]) -> bool: + """Decide whether ComfyUI is really working, not just claiming to be. + + ComfyUI can leave an entry in queue_running after a job dies -- observed here as + a WAN 2.1 i2v entry that sat there with the GPU at 0% and ComfyUI holding 0.56 GB. + Trusting that flag alone made this service believe ComfyUI was permanently busy, + which meant it evicted the LLM on every poll, never ran the idle purge, and never + checked whether the LLM had been squeezed onto the CPU. Half the arbitration was + disabled by one stale row. + + A running entry is corroborated against GPU utilisation before it is believed. + """ + running = queue.get("queue_running") or [] + pending = queue.get("queue_pending") or [] + if pending: + self._running_since = None + self._running_id = None + return True + if not running: + self._running_since = None + self._running_id = None + self.comfy_stale_job = None + return False + + entry = running[0] + prompt_id = entry[1] if isinstance(entry, (list, tuple)) and len(entry) > 1 else str(entry) + now = time.time() + if prompt_id != self._running_id: + self._running_id = prompt_id + self._running_since = now + self._peak_comfy_bytes = 0 + + snap = get_process_vram_bytes() + self._peak_comfy_bytes = max(self._peak_comfy_bytes, snap.get("comfyui_bytes", 0)) + + elapsed = now - (self._running_since or now) + if elapsed > self.STALE_RUNNING_S and self._peak_comfy_bytes < self.STALE_COMFY_BYTES: + if self.comfy_stale_job != prompt_id: + logger.warning( + f"ComfyUI reports prompt {prompt_id} running for {int(elapsed)}s while " + f"holding only {self._peak_comfy_bytes / (1024**3):.2f} GB — no checkpoint " + f"is loaded, so the queue entry is stale. Ignoring it; otherwise ComfyUI " + f"looks permanently busy and arbitration stops working.") + self.comfy_stale_job = prompt_id + return False + return True async def _check_ollama_starved(self) -> None: """The other direction: rescue an LLM that ComfyUI has squeezed onto the CPU. @@ -1135,17 +1199,37 @@ class AutoArbitrator: If the LLM is spilling while ComfyUI sits idle holding VRAM, ComfyUI's cached checkpoints are the thing to give up. """ + # Every bail-out is recorded. This check silently did nothing while a model sat + # at 29% on the GPU, and with four separate early returns there was no way to + # tell which one had fired without guessing. now = time.time() - if self.comfy_was_active or (now - self.last_reclaim_time) < self.RECLAIM_COOLDOWN_S: - return + def bail(reason: str, **extra): + self.last_starvation_check = {"ts": now, "acted": False, + "reason": reason, **extra} + + if self.comfy_was_active: + return bail("ComfyUI is active; it needs the VRAM itself") + if (now - self.last_reclaim_time) < self.RECLAIM_COOLDOWN_S: + return bail("within reclaim cooldown", + seconds_left=round(self.RECLAIM_COOLDOWN_S + - (now - self.last_reclaim_time), 1)) ollama = await get_ollama_live_state() if not ollama.get("partially_offloaded"): - return + return bail("LLM is not spilling to CPU", + gpu_fraction=ollama.get("gpu_fraction"), + model=ollama.get("active_model_name")) snap = get_process_vram_bytes() if snap["comfyui_bytes"] < RECLAIM_MIN_COMFY_BYTES: - return # ComfyUI is not the one holding the memory; nothing we can do here + # ComfyUI is not the one holding the memory; nothing we can do here. + return bail("LLM is spilling but ComfyUI holds too little to help", + cpu_offload_pct=ollama.get("cpu_offload_pct"), + comfy_gb=round(snap["comfyui_bytes"] / (1024**3), 2)) + + self.last_starvation_check = {"ts": now, "acted": True, + "reason": "reclaiming for the LLM", + "cpu_offload_pct": ollama.get("cpu_offload_pct")} self.last_reclaim_time = now model = ollama.get("active_model_name") @@ -1185,15 +1269,25 @@ class AutoArbitrator: resp = await client.get("/queue") if resp.status_code == 200: q = resp.json() - busy = len(q.get("queue_running", [])) > 0 or len(q.get("queue_pending", [])) > 0 + busy = self._comfy_genuinely_busy(q) if busy: + self.watchdog_branches["busy"] += 1 await self.trigger_comfy_priority("Watchdog saw an active queue") elif self.comfy_was_active: + self.watchdog_branches["completed"] += 1 await self.trigger_comfy_completed() else: + self.watchdog_branches["idle_check"] += 1 await self._check_ollama_starved() - except Exception: - pass + else: + self.watchdog_branches["bad_status"] += 1 + except Exception as e: + # This used to swallow everything silently, including anything raised by + # the starvation check, which is why that check could appear to run and + # do nothing. + self.watchdog_branches["error"] += 1 + self.last_watchdog_error = str(e)[:200] + logger.debug(f"watchdog poll error: {e}") await asyncio.sleep(interval) def suspend_oc(self, reason: str = "tuning sweep") -> None: @@ -1231,6 +1325,10 @@ class AutoArbitrator: "idle_purge_after_s": self.COMFY_IDLE_PURGE_S, "oc_profile": self.oc_profile, "counters": dict(self.stats), + "last_starvation_check": self.last_starvation_check, + "comfy_stale_job": self.comfy_stale_job, + "watchdog_branches": dict(self.watchdog_branches), + "last_watchdog_error": self.last_watchdog_error, "yield_backoff": {m: round(max(t - time.time(), 0), 1) for m, t in self._yield_backoff_until.items() if t > time.time()},