Stop a stale ComfyUI queue entry from disabling half the arbitration

Chasing why the reverse-direction reclaim never fired turned up something worse than
the reclaim itself.

The starvation check was never running. Instrumenting the watchdog showed busy=6,
idle_check=0: every poll took the "ComfyUI is busy" branch. ComfyUI's /queue was
reporting a WAN 2.1 i2v job in queue_running while the GPU sat at 0% and ComfyUI held
0.56 GB. The job was dead; ComfyUI had simply never cleared the row.

Believing that flag meant this service thought ComfyUI was permanently busy, so it
yielded the LLM's VRAM on every poll, never ran the idle purge, and never checked
whether the LLM had been squeezed onto the CPU. One stale row disabled half of the
arbitration, and it very likely explains the earlier burst of yields against a
cron-driven model.

A running entry is now corroborated before it is believed. The first attempt used GPU
utilisation, which does not work: utilisation is shared with Ollama and with the
third-party process on this box, so peak utilisation stayed above any sensible
threshold and a stuck entry never looked stale. ComfyUI's own VRAM is the right
signal -- a real diffusion job loads gigabytes of checkpoint, a dead one holds only
its CUDA context. After the fix the same watchdog reports busy=3, idle_check=32.

Every early return in the starvation check now records why it bailed, because with
four of them there was no way to tell which had fired. /api/health reports a stale
queue entry with its impact and how to clear it.

Also confirmed, contradicting an earlier conclusion in this branch: Ollama on this box
*does* spill to the CPU. smtek/Qwen3.8-27B:Q2_K_XL held steady at 29.2% on GPU
(size=15.59 GB, size_vram=4.56 GB) across twelve seconds of polling -- a stable
placement, not a progressive load. Both failure modes are real; which one occurs
depends on the model.

Tests: 206 (was 199).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
drjones
2026-09-07 08:48:37 -07:00
parent aeba1b47fd
commit c3d9b36035
6 changed files with 218 additions and 24 deletions

View File

@@ -21,6 +21,15 @@
### ⚡ Bidirectional VRAM Hot-Swapping & Arbitration ### ⚡ 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 * **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 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 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 ## 1a. Tests
```bash ```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` Hermetic: no GPU, no network, no sleeps. An autouse fixture stubs `overclock_manager._sh`

View File

@@ -117,6 +117,25 @@ def _check_comfy_ws() -> Dict[str, Any]:
return _check("comfyui websocket", OK, "subscribed") 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]: def _check_store() -> Dict[str, Any]:
info = telemetry_store.db_info() info = telemetry_store.db_info()
if not info.get("exists"): if not info.get("exists"):

View File

@@ -1,7 +1,6 @@
{ {
"ollama": { "ollama": {
"label": "Ollama — LLM decode (memory-bandwidth bound; measured insensitive to power and clocks)", "label": "Ollama \u2014 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).",
"power_limit_w": 320, "power_limit_w": 320,
"core_offset_mhz": 0, "core_offset_mhz": 0,
"mem_offset_mhz": 0, "mem_offset_mhz": 0,
@@ -9,11 +8,11 @@
"lock_core_max": 0, "lock_core_max": 0,
"lock_mem_mhz": 0, "lock_mem_mhz": 0,
"fan_mode": "auto", "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": { "comfy": {
"label": "ComfyUI — diffusion (compute bound; genuinely power-scaling)", "label": "ComfyUI \u2014 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.",
"power_limit_w": 370, "power_limit_w": 370,
"core_offset_mhz": 0, "core_offset_mhz": 0,
"mem_offset_mhz": 0, "mem_offset_mhz": 0,
@@ -21,18 +20,19 @@
"lock_core_max": 0, "lock_core_max": 0,
"lock_mem_mhz": 0, "lock_mem_mhz": 0,
"fan_mode": "auto", "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": { "balanced": {
"label": "Balanced — stock power and boost, automatic fans", "label": "Balanced \u2014 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": 340,
"power_limit_w": 320, "core_offset_mhz": 10,
"core_offset_mhz": 0, "mem_offset_mhz": 150,
"mem_offset_mhz": 0,
"lock_core_min": 0, "lock_core_min": 0,
"lock_core_max": 0, "lock_core_max": 0,
"lock_mem_mhz": 0, "lock_mem_mhz": 0,
"fan_mode": "auto", "fan_mode": "manual",
"fan_speed_pct": 0 "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."
} }
} }

View File

@@ -106,7 +106,7 @@ class TestAggregation:
monkeypatch.setattr(health, "_check_sudo_smi", lambda: checks[1]) monkeypatch.setattr(health, "_check_sudo_smi", lambda: checks[1])
for fn in ("_check_fan_control", "_check_profile_drift", "_check_store", for fn in ("_check_fan_control", "_check_profile_drift", "_check_store",
"_check_residency", "_check_model_dirs", "_check_comfy_ws", "_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")) monkeypatch.setattr(health, fn, lambda: health._check("x", health.OK, "d"))
async def fake_http(name, url, impact, fix): async def fake_http(name, url, impact, fix):
@@ -122,7 +122,7 @@ class TestAggregation:
monkeypatch.setattr(health, "_check_nvml", boom) monkeypatch.setattr(health, "_check_nvml", boom)
for fn in ("_check_sudo_smi", "_check_fan_control", "_check_profile_drift", for fn in ("_check_sudo_smi", "_check_fan_control", "_check_profile_drift",
"_check_store", "_check_residency", "_check_model_dirs", "_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")) monkeypatch.setattr(health, fn, lambda: health._check("x", health.OK, "d"))
async def fake_http(name, url, impact, fix): async def fake_http(name, url, impact, fix):

View File

@@ -13,6 +13,7 @@ were observed on real hardware before being encoded here:
out of memory ... unable to allocate CUDA0 buffer" out of memory ... unable to allocate CUDA0 buffer"
""" """
import asyncio import asyncio
import time
import pytest import pytest
@@ -132,3 +133,70 @@ class TestBusyBackoff:
assert "yield_deferred_busy" in arb.stats assert "yield_deferred_busy" in arb.stats
assert "yield_stalled" in arb.stats assert "yield_stalled" in arb.stats
assert "yield_timeouts" not 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")

View File

@@ -931,6 +931,14 @@ class AutoArbitrator:
self._yield_backoff_until: Dict[str, float] = {} self._yield_backoff_until: Dict[str, float] = {}
self._yield_busy_streak: Dict[str, int] = {} self._yield_busy_streak: Dict[str, int] = {}
self.last_reclaim_time = 0.0 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 = { self.stats = {
"yields": 0, # release confirmed "yields": 0, # release confirmed
"yield_deferred_busy": 0, # model mid-generation; unload queued behind it "yield_deferred_busy": 0, # model mid-generation; unload queued behind it
@@ -1122,6 +1130,62 @@ class AutoArbitrator:
backoff = min(backoff * 1.5, 15.0) backoff = min(backoff * 1.5, 15.0)
RECLAIM_COOLDOWN_S = 30.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: async def _check_ollama_starved(self) -> None:
"""The other direction: rescue an LLM that ComfyUI has squeezed onto the CPU. """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 If the LLM is spilling while ComfyUI sits idle holding VRAM, ComfyUI's cached
checkpoints are the thing to give up. 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() now = time.time()
if self.comfy_was_active or (now - self.last_reclaim_time) < self.RECLAIM_COOLDOWN_S: def bail(reason: str, **extra):
return 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() ollama = await get_ollama_live_state()
if not ollama.get("partially_offloaded"): 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() snap = get_process_vram_bytes()
if snap["comfyui_bytes"] < RECLAIM_MIN_COMFY_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 self.last_reclaim_time = now
model = ollama.get("active_model_name") model = ollama.get("active_model_name")
@@ -1185,15 +1269,25 @@ class AutoArbitrator:
resp = await client.get("/queue") resp = await client.get("/queue")
if resp.status_code == 200: if resp.status_code == 200:
q = resp.json() 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: if busy:
self.watchdog_branches["busy"] += 1
await self.trigger_comfy_priority("Watchdog saw an active queue") await self.trigger_comfy_priority("Watchdog saw an active queue")
elif self.comfy_was_active: elif self.comfy_was_active:
self.watchdog_branches["completed"] += 1
await self.trigger_comfy_completed() await self.trigger_comfy_completed()
else: else:
self.watchdog_branches["idle_check"] += 1
await self._check_ollama_starved() await self._check_ollama_starved()
except Exception: else:
pass 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) await asyncio.sleep(interval)
def suspend_oc(self, reason: str = "tuning sweep") -> None: def suspend_oc(self, reason: str = "tuning sweep") -> None:
@@ -1231,6 +1325,10 @@ class AutoArbitrator:
"idle_purge_after_s": self.COMFY_IDLE_PURGE_S, "idle_purge_after_s": self.COMFY_IDLE_PURGE_S,
"oc_profile": self.oc_profile, "oc_profile": self.oc_profile,
"counters": dict(self.stats), "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) "yield_backoff": {m: round(max(t - time.time(), 0), 1)
for m, t in self._yield_backoff_until.items() for m, t in self._yield_backoff_until.items()
if t > time.time()}, if t > time.time()},