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)