Arbitrate over any number of tenants by priority

Completes the generalisation. Classification and release were already data; the
decision loop was still two hardcoded rules -- yield Ollama when ComfyUI is busy,
purge ComfyUI when Ollama is starved -- which could not express a third participant.

plan_release() works from the registry instead. A busy tenant that cannot reach its
declared needs_vram_gb is starved, and the memory comes from idle reclaimable tenants
below it in priority, lowest first, stopping once enough is freed. Tenants that cannot
be released are named as blockers rather than passed over, so an impossible plan says
which process is in the way. The plan is returned before being acted on, so the
decision is testable and is logged before anything is released. Idle release is now
per-tenant too, replacing the ComfyUI-specific purge timer.

Two bugs found by running it against the live machine rather than only in tests:

Starvation was measured against free VRAM alone, so a tenant working perfectly well on
13 GB was flagged as demanding simply because little was left over -- which is the
normal state of a busy GPU, and would have caused pointless releases from everything
else. A tenant is starved only if it cannot reach what it needs counting what it
already holds.

Fields added to the tenant schema were silently absent from the config already written
to disk, so needs_vram_gb defaulted to 0 and starvation could never trigger for the two
tenants that mattered. Shipped defaults are now merged into an existing config on load,
with explicit user values still winning.

Tests: 242 (was 231), including a three-application contention case -- the property the
hardcoded pair of rules could not express.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
drjones
2026-09-07 14:47:32 -07:00
parent 1bcfbb2335
commit 4a38cd68b3
4 changed files with 319 additions and 2 deletions

View File

@@ -201,8 +201,21 @@ A tenant is now **described as data** in `tenants.json`:
| `release` | How to ask for VRAM back — `http_post` with a body, `per_model` for Ollama's per-model unload, or `none` | | `release` | How to ask for VRAM back — `http_post` with a body, `per_model` for Ollama's per-model unload, or `none` |
| `priority` | Who wins contention | | `priority` | Who wins contention |
Two more fields drive the decision loop: **`needs_vram_gb`** (how much free memory the
application needs before it can work) and **`idle_release_after_s`** (how long it may sit
idle holding VRAM before being asked for it back — deliberately not immediate, so
iterating on a ComfyUI workflow does not reload the checkpoint between every run).
`plan_release()` then arbitrates generically: a busy tenant that cannot reach
`needs_vram_gb` *even counting what it already holds* is starved, and the memory is taken
from idle reclaimable tenants below it in priority, lowest first, stopping as soon as
enough is freed. Tenants that cannot be released are named as blockers rather than
ignored, so `possible: false` comes with the reason. The plan is returned before it is
acted on, which makes the decision testable and loggable.
Ollama, ComfyUI and the desktop compositor ship as defaults, so behaviour is unchanged — Ollama, ComfyUI and the desktop compositor ship as defaults, so behaviour is unchanged —
but nothing in the arbitration logic knows their names. Endpoints are generic: but nothing in the arbitration logic knows their names, and three applications can
contend for the card as easily as two. Endpoints are generic:
`GET /api/tenants`, `GET /api/tenants/{name}`, `POST /api/tenants/{name}/release`. `GET /api/tenants`, `GET /api/tenants/{name}`, `POST /api/tenants/{name}/release`.
A tenant with `"release": {"type": "none"}` is still worth declaring. The 842 MB speech A tenant with `"release": {"type": "none"}` is still worth declaring. The 842 MB speech
@@ -213,7 +226,7 @@ explaining that, rather than silently doing nothing.
## 1a. Tests ## 1a. Tests
```bash ```bash
/home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 231 passed in ~3.8s /home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 242 passed in ~3.8s
``` ```
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

@@ -89,6 +89,13 @@ class GpuTenant:
enabled: bool = True enabled: bool = True
# Higher wins contention; a tenant yields to anything above it. # Higher wins contention; a tenant yields to anything above it.
priority: int = 50 priority: int = 50
# How much free VRAM this application needs before it can work. Used to decide
# whether a busy tenant is actually being starved, rather than merely busy.
needs_vram_gb: float = 0.0
# How long a reclaimable tenant may sit idle holding VRAM before it is asked for it
# back. Iterating on a ComfyUI workflow should not pay a reload between every run,
# so this is deliberately not immediate.
idle_release_after_s: float = 30.0
match: ProcessMatch = field(default_factory=ProcessMatch) match: ProcessMatch = field(default_factory=ProcessMatch)
busy: BusyProbe = field(default_factory=BusyProbe) busy: BusyProbe = field(default_factory=BusyProbe)
release: ReleaseStrategy = field(default_factory=ReleaseStrategy) release: ReleaseStrategy = field(default_factory=ReleaseStrategy)
@@ -108,6 +115,8 @@ def _tenant_from_dict(d: Dict[str, Any]) -> GpuTenant:
kind=d.get("kind", KIND_OTHER), kind=d.get("kind", KIND_OTHER),
enabled=d.get("enabled", True), enabled=d.get("enabled", True),
priority=int(d.get("priority", 50)), priority=int(d.get("priority", 50)),
needs_vram_gb=float(d.get("needs_vram_gb", 0.0)),
idle_release_after_s=float(d.get("idle_release_after_s", 30.0)),
match=ProcessMatch(**(d.get("match") or {})), match=ProcessMatch(**(d.get("match") or {})),
busy=BusyProbe(**(d.get("busy") or {})), busy=BusyProbe(**(d.get("busy") or {})),
release=ReleaseStrategy(**(d.get("release") or {})), release=ReleaseStrategy(**(d.get("release") or {})),
@@ -121,6 +130,8 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [
"name": "ollama", "name": "ollama",
"kind": KIND_LLM, "kind": KIND_LLM,
"priority": 60, "priority": 60,
"needs_vram_gb": 4.0,
"idle_release_after_s": 0.0,
"match": {"names": ["ollama"], "cmdline": ["llama-server", "ollama"]}, "match": {"names": ["ollama"], "cmdline": ["llama-server", "ollama"]},
"busy": {"type": "http_count", "url": "http://localhost:11434/api/ps", "busy": {"type": "http_count", "url": "http://localhost:11434/api/ps",
"count_keys": ["models"]}, "count_keys": ["models"]},
@@ -133,6 +144,8 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [
"name": "comfyui", "name": "comfyui",
"kind": KIND_DIFFUSION, "kind": KIND_DIFFUSION,
"priority": 50, "priority": 50,
"needs_vram_gb": 6.0,
"idle_release_after_s": 30.0,
"match": {"cmdline": ["comfyui", "comfy"], "cmdline_endswith": ["main.py"]}, "match": {"cmdline": ["comfyui", "comfy"], "cmdline_endswith": ["main.py"]},
"busy": {"type": "http_count", "url": "http://127.0.0.1:8188/queue", "busy": {"type": "http_count", "url": "http://127.0.0.1:8188/queue",
"count_keys": ["queue_running", "queue_pending"], "count_keys": ["queue_running", "queue_pending"],
@@ -187,8 +200,20 @@ def load_tenants(force: bool = False) -> List[GpuTenant]:
except Exception as e: except Exception as e:
logger.warning(f"could not write {CONFIG_PATH}: {e}") logger.warning(f"could not write {CONFIG_PATH}: {e}")
# Merge in any fields a shipped default has gained since the config was written.
# Without this, adding a field silently disables the behaviour it controls for every
# existing install -- needs_vram_gb defaulted to 0, which made starvation
# undetectable for the two tenants that had been written out before it existed.
defaults_by_name = {d["name"]: d for d in DEFAULT_TENANTS}
tenants = [] tenants = []
for d in raw: for d in raw:
base = defaults_by_name.get(d.get("name"))
if base:
merged = {**base, **d}
for key in ("match", "busy", "release"):
if isinstance(base.get(key), dict):
merged[key] = {**base[key], **(d.get(key) or {})}
d = merged
try: try:
tenants.append(_tenant_from_dict(d)) tenants.append(_tenant_from_dict(d))
except Exception as e: except Exception as e:
@@ -311,3 +336,67 @@ def describe() -> List[Dict[str, Any]]:
d["release_via"] = t.release.type d["release_via"] = t.release.type
out.append(d) out.append(d)
return out return out
# ---------------------------------------------------------------- arbitration
def plan_release(demanding: str, tenants_state: List[Dict[str, Any]],
free_gb: float, needed_gb: float) -> Dict[str, Any]:
"""Decide who should give up VRAM so a starved tenant can work.
Generic over any number of applications: candidates are every *reclaimable* tenant
that is not itself busy and ranks below the demanding one, taken lowest priority
first, until enough would be freed. The two-application version of this was a pair
of hardcoded rules -- yield Ollama for ComfyUI, purge ComfyUI for Ollama -- which
could not express a third participant at all.
Returns the plan rather than performing it, so the decision is testable and can be
logged before anything is actually released.
"""
by_name = {s["name"]: s for s in tenants_state}
demander = by_name.get(demanding)
if not demander:
return {"possible": False, "reason": f"unknown tenant '{demanding}'", "release": []}
# The demander keeps what it already holds; only the remainder must be found.
shortfall = needed_gb - free_gb - demander.get("vram_gb", 0.0)
if shortfall <= 0:
return {"possible": True, "reason": "enough VRAM is already free",
"release": [], "shortfall_gb": 0.0}
candidates = [
s for s in tenants_state
if s["name"] != demanding
and s.get("reclaimable")
and not s.get("busy")
and s.get("priority", 0) <= demander.get("priority", 0)
and s.get("vram_gb", 0) > 0
]
candidates.sort(key=lambda s: (s.get("priority", 0), -s.get("vram_gb", 0)))
plan, freed = [], 0.0
for c in candidates:
if freed >= shortfall:
break
plan.append(c["name"])
freed += c.get("vram_gb", 0.0)
blockers = [
{"name": s["name"], "vram_gb": s.get("vram_gb", 0.0),
"why": ("busy" if s.get("busy") else
"declares no release mechanism" if not s.get("reclaimable") else
"higher priority")}
for s in tenants_state
if s["name"] != demanding and s.get("vram_gb", 0) > 0 and s["name"] not in plan
]
return {
"possible": freed >= shortfall,
"shortfall_gb": round(shortfall, 2),
"would_free_gb": round(freed, 2),
"release": plan,
"blockers": blockers,
"reason": (f"releasing {', '.join(plan)} frees {freed:.2f} GB of the "
f"{shortfall:.2f} GB shortfall" if plan else
"no reclaimable idle tenant holds enough VRAM"),
}

View File

@@ -172,3 +172,112 @@ class TestPriority:
prios = [r["priority"] for r in rows] prios = [r["priority"] for r in rows]
assert prios == sorted(prios, reverse=True) assert prios == sorted(prios, reverse=True)
assert all("reclaimable" in r for r in rows) assert all("reclaimable" in r for r in rows)
class TestReleasePlanning:
"""Deciding who gives up VRAM, generically over any number of applications.
The two-application version was a pair of hardcoded rules -- yield Ollama when
ComfyUI is busy, purge ComfyUI when Ollama is starved -- which could not express a
third participant at all.
"""
def _state(self, **overrides):
base = [
{"name": "desktop", "priority": 90, "vram_gb": 0.01, "busy": False,
"reclaimable": False},
{"name": "stt-relay", "priority": 70, "vram_gb": 0.8, "busy": False,
"reclaimable": False},
{"name": "ollama", "priority": 60, "vram_gb": 0.0, "busy": True,
"reclaimable": True},
{"name": "comfyui", "priority": 50, "vram_gb": 7.0, "busy": False,
"reclaimable": True},
]
for s in base:
s.update(overrides.get(s["name"], {}))
return base
def test_a_tenant_already_holding_what_it_needs_is_not_starved(self):
# A busy GPU has little free by definition. Comparing free VRAM alone flagged a
# tenant working fine on 13 GB as demanding, which would have caused pointless
# releases from everything else.
state = self._state(ollama={"vram_gb": 13.0})
plan = T.plan_release("ollama", state, free_gb=1.5, needed_gb=4.0)
assert plan["release"] == []
assert "already free" in plan["reason"]
def test_starved_tenant_reclaims_from_the_idle_one_below_it(self):
plan = T.plan_release("ollama", self._state(), free_gb=1.5, needed_gb=14.9)
assert plan["release"] == ["comfyui"]
def test_a_busy_tenant_is_never_a_victim(self):
state = self._state(comfyui={"busy": True})
plan = T.plan_release("ollama", state, free_gb=1.5, needed_gb=14.9)
assert plan["release"] == []
assert any(b["name"] == "comfyui" and b["why"] == "busy"
for b in plan["blockers"])
def test_unreclaimable_tenants_are_named_as_blockers_not_ignored(self):
# The user needs to know a third-party process is what stands in the way.
plan = T.plan_release("ollama", self._state(), free_gb=0.0, needed_gb=15.5)
blockers = {b["name"]: b["why"] for b in plan["blockers"]}
assert blockers["stt-relay"] == "declares no release mechanism"
assert plan["possible"] is False
def test_higher_priority_tenants_are_not_victims(self):
state = self._state(comfyui={"priority": 99, "busy": False})
plan = T.plan_release("ollama", state, free_gb=1.0, needed_gb=14.9)
assert "comfyui" not in plan["release"]
def test_lowest_priority_is_released_first(self):
state = self._state() + [
{"name": "batch", "priority": 10, "vram_gb": 3.0, "busy": False,
"reclaimable": True}]
plan = T.plan_release("ollama", state, free_gb=0.0, needed_gb=5.0)
assert plan["release"][0] == "batch"
def test_releases_only_as_many_tenants_as_needed(self):
state = self._state() + [
{"name": "batch", "priority": 10, "vram_gb": 9.0, "busy": False,
"reclaimable": True}]
plan = T.plan_release("ollama", state, free_gb=0.0, needed_gb=8.0)
assert plan["release"] == ["batch"] # 9 GB covers it; comfyui is left alone
def test_three_applications_can_all_participate(self):
# The property the hardcoded pair of rules could not express.
state = [
{"name": "llm", "priority": 60, "vram_gb": 0.0, "busy": True,
"reclaimable": True},
{"name": "diffusion", "priority": 50, "vram_gb": 4.0, "busy": False,
"reclaimable": True},
{"name": "trainer", "priority": 40, "vram_gb": 5.0, "busy": False,
"reclaimable": True},
]
plan = T.plan_release("llm", state, free_gb=0.0, needed_gb=9.0)
assert set(plan["release"]) == {"trainer", "diffusion"}
assert plan["possible"] is True
def test_unknown_tenant_is_rejected_cleanly(self):
plan = T.plan_release("nope", self._state(), free_gb=0.0, needed_gb=1.0)
assert plan["possible"] is False and plan["release"] == []
class TestConfigUpgrade:
def test_fields_added_later_are_merged_into_an_existing_config(self, cfg):
# A config written before needs_vram_gb existed must not silently lose the
# behaviour that field controls.
cfg.write_text(json.dumps([{
"name": "ollama",
"match": {"names": ["ollama"]},
}]))
t = T.get_tenant("ollama")
assert t.needs_vram_gb > 0
assert t.release.type == "http_post"
def test_explicit_user_values_still_win_over_defaults(self, cfg):
cfg.write_text(json.dumps([{
"name": "ollama", "priority": 5, "needs_vram_gb": 99.0,
"match": {"names": ["ollama"]},
}]))
t = T.get_tenant("ollama")
assert t.priority == 5 and t.needs_vram_gb == 99.0

View File

@@ -949,6 +949,9 @@ class AutoArbitrator:
self._running_since: Optional[float] = None self._running_since: Optional[float] = None
self._peak_comfy_bytes = 0 self._peak_comfy_bytes = 0
self.comfy_stale_job: Optional[str] = None self.comfy_stale_job: Optional[str] = None
self._idle_since: Dict[str, float] = {}
self._last_tenant_state: Optional[Dict[str, Any]] = None
self.last_arbitration: Optional[Dict[str, Any]] = None
self.last_watchdog_error: Optional[str] = None self.last_watchdog_error: Optional[str] = None
self.stats = { self.stats = {
"yields": 0, # release confirmed "yields": 0, # release confirmed
@@ -1198,6 +1201,106 @@ class AutoArbitrator:
return False return False
return True return True
async def _tenant_state(self) -> List[Dict[str, Any]]:
"""Current VRAM and busy state for every configured tenant."""
snap = get_process_vram_bytes()
stats = get_gpu_hardware_stats()
by_tenant = (stats.get("breakdown", {}) or {}).get("by_tenant_gb", {})
out = []
for t in tenants_mod.load_tenants():
if not t.enabled:
continue
bucket = _BUCKET_ALIASES.get(t.name, t.name)
vram_gb = by_tenant.get(bucket, 0.0)
probe = await tenants_mod.probe_busy(t, vram_gb=vram_gb)
out.append({
"name": t.name,
"priority": t.priority,
"vram_gb": vram_gb,
"busy": bool(probe.get("busy")),
"below_floor": bool(probe.get("below_floor")),
"reclaimable": t.reclaimable,
"needs_vram_gb": t.needs_vram_gb,
"idle_release_after_s": t.idle_release_after_s,
"reason": probe.get("reason"),
})
self._last_tenant_state = {"ts": time.time(), "free_gb":
round(snap["free_bytes"] / (1024**3), 2),
"tenants": out}
return out
async def _release_tenant(self, name: str, reason: str) -> Dict[str, Any]:
"""Release one tenant's VRAM by whatever mechanism it declares."""
t = tenants_mod.get_tenant(name)
if not t or not t.reclaimable:
return {"success": False, "reason": "not reclaimable"}
models = None
if t.release.per_model:
state = await get_ollama_live_state()
models = [m.get("name") for m in state.get("loaded_models", []) if m.get("name")]
logger.info(f"Releasing VRAM from '{name}': {reason}")
res = await tenants_mod.release_vram(t, models=models)
self.stats["tenant_releases"] = self.stats.get("tenant_releases", 0) + 1
return res
async def _arbitrate(self) -> None:
"""Generic arbitration over any number of tenants.
The two-application version was a pair of hardcoded rules -- yield Ollama when
ComfyUI is busy, purge ComfyUI when Ollama is starved -- which could not express
a third participant at all. This works from the registry instead: a busy tenant
that lacks the VRAM it declares it needs is starved, and the memory comes from
idle reclaimable tenants below it in priority, lowest first.
"""
state = await self._tenant_state()
free_gb = self._last_tenant_state["free_gb"]
# 1. Starvation: highest-priority demanding tenant first.
for s in sorted(state, key=lambda x: -x["priority"]):
if not s["busy"] or not s["needs_vram_gb"]:
continue
# Starved means it cannot reach what it needs even counting what it already
# holds. Comparing free VRAM alone flagged a tenant that was working
# perfectly well on 13 GB as demanding, purely because little was left over
# -- which is the normal state of a busy GPU, and would have caused
# pointless releases from everyone else.
if s["vram_gb"] + free_gb >= s["needs_vram_gb"]:
continue
plan = tenants_mod.plan_release(s["name"], state, free_gb, s["needs_vram_gb"])
self.last_arbitration = {"ts": time.time(), "demanding": s["name"],
"free_gb": free_gb, **plan}
if not plan["release"]:
logger.debug(f"'{s['name']}' is short of VRAM but {plan['reason']}")
return
if time.time() - self.last_reclaim_time < self.RECLAIM_COOLDOWN_S:
return
self.last_reclaim_time = time.time()
for victim in plan["release"]:
await self._release_tenant(
victim, f"{s['name']} needs {s['needs_vram_gb']} GB, {free_gb} GB free")
self.last_action = (f"Released {', '.join(plan['release'])} so "
f"'{s['name']}' could work")
return
# 2. Idle release: a tenant holding VRAM it is not using, after a grace period.
now = time.time()
for s in state:
if not s["reclaimable"] or s["vram_gb"] <= 0.25:
self._idle_since.pop(s["name"], None)
continue
if s["busy"]:
self._idle_since.pop(s["name"], None)
continue
since = self._idle_since.setdefault(s["name"], now)
grace = s["idle_release_after_s"]
if grace and (now - since) >= grace:
self._idle_since.pop(s["name"], None)
await self._release_tenant(
s["name"], f"idle {int(now - since)}s holding {s['vram_gb']} GB")
self.last_action = (f"Released idle '{s['name']}' after "
f"{int(now - since)}s")
return
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.
@@ -1290,6 +1393,7 @@ class AutoArbitrator:
else: else:
self.watchdog_branches["idle_check"] += 1 self.watchdog_branches["idle_check"] += 1
await self._check_ollama_starved() await self._check_ollama_starved()
await self._arbitrate()
else: else:
self.watchdog_branches["bad_status"] += 1 self.watchdog_branches["bad_status"] += 1
except Exception as e: except Exception as e:
@@ -1338,6 +1442,8 @@ class AutoArbitrator:
"counters": dict(self.stats), "counters": dict(self.stats),
"last_starvation_check": self.last_starvation_check, "last_starvation_check": self.last_starvation_check,
"comfy_stale_job": self.comfy_stale_job, "comfy_stale_job": self.comfy_stale_job,
"last_arbitration": self.last_arbitration,
"tenant_state": self._last_tenant_state,
"watchdog_branches": dict(self.watchdog_branches), "watchdog_branches": dict(self.watchdog_branches),
"last_watchdog_error": self.last_watchdog_error, "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)