From 4a38cd68b32bdd232f3bea0b659e77118eed551e Mon Sep 17 00:00:00 2001 From: drjones Date: Mon, 7 Sep 2026 14:47:32 -0700 Subject: [PATCH] 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 --- README.md | 17 ++++++- tenants.py | 89 ++++++++++++++++++++++++++++++++++ tests/test_tenants.py | 109 ++++++++++++++++++++++++++++++++++++++++++ vram_arbitrator.py | 106 ++++++++++++++++++++++++++++++++++++++++ 4 files changed, 319 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 212ea84..803a6c1 100644 --- a/README.md +++ b/README.md @@ -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` | | `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 — -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`. 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 ```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` diff --git a/tenants.py b/tenants.py index 1006b1f..67a6837 100644 --- a/tenants.py +++ b/tenants.py @@ -89,6 +89,13 @@ class GpuTenant: enabled: bool = True # Higher wins contention; a tenant yields to anything above it. 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) busy: BusyProbe = field(default_factory=BusyProbe) 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), enabled=d.get("enabled", True), 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 {})), busy=BusyProbe(**(d.get("busy") or {})), release=ReleaseStrategy(**(d.get("release") or {})), @@ -121,6 +130,8 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [ "name": "ollama", "kind": KIND_LLM, "priority": 60, + "needs_vram_gb": 4.0, + "idle_release_after_s": 0.0, "match": {"names": ["ollama"], "cmdline": ["llama-server", "ollama"]}, "busy": {"type": "http_count", "url": "http://localhost:11434/api/ps", "count_keys": ["models"]}, @@ -133,6 +144,8 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [ "name": "comfyui", "kind": KIND_DIFFUSION, "priority": 50, + "needs_vram_gb": 6.0, + "idle_release_after_s": 30.0, "match": {"cmdline": ["comfyui", "comfy"], "cmdline_endswith": ["main.py"]}, "busy": {"type": "http_count", "url": "http://127.0.0.1:8188/queue", "count_keys": ["queue_running", "queue_pending"], @@ -187,8 +200,20 @@ def load_tenants(force: bool = False) -> List[GpuTenant]: except Exception as 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 = [] 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: tenants.append(_tenant_from_dict(d)) except Exception as e: @@ -311,3 +336,67 @@ def describe() -> List[Dict[str, Any]]: d["release_via"] = t.release.type out.append(d) 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"), + } diff --git a/tests/test_tenants.py b/tests/test_tenants.py index aa371b0..e008f32 100644 --- a/tests/test_tenants.py +++ b/tests/test_tenants.py @@ -172,3 +172,112 @@ class TestPriority: prios = [r["priority"] for r in rows] assert prios == sorted(prios, reverse=True) 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 diff --git a/vram_arbitrator.py b/vram_arbitrator.py index 2a7118b..9fdd965 100644 --- a/vram_arbitrator.py +++ b/vram_arbitrator.py @@ -949,6 +949,9 @@ class AutoArbitrator: self._running_since: Optional[float] = None self._peak_comfy_bytes = 0 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.stats = { "yields": 0, # release confirmed @@ -1198,6 +1201,106 @@ class AutoArbitrator: return False 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: """The other direction: rescue an LLM that ComfyUI has squeezed onto the CPU. @@ -1290,6 +1393,7 @@ class AutoArbitrator: else: self.watchdog_branches["idle_check"] += 1 await self._check_ollama_starved() + await self._arbitrate() else: self.watchdog_branches["bad_status"] += 1 except Exception as e: @@ -1338,6 +1442,8 @@ class AutoArbitrator: "counters": dict(self.stats), "last_starvation_check": self.last_starvation_check, "comfy_stale_job": self.comfy_stale_job, + "last_arbitration": self.last_arbitration, + "tenant_state": self._last_tenant_state, "watchdog_branches": dict(self.watchdog_branches), "last_watchdog_error": self.last_watchdog_error, "yield_backoff": {m: round(max(t - time.time(), 0), 1)