diff --git a/README.md b/README.md index 26fd597..e533d16 100644 --- a/README.md +++ b/README.md @@ -199,13 +199,23 @@ A tenant is now **described as data** in `tenants.json`: | `match` | Which GPU processes belong to this application (name, cmdline substring, or suffix — ComfyUI is a bare `python main.py`) | | `busy` | Whether it is *genuinely* working. `http_count` sums queue lists; `vram` needs no API at all. `vram_floor_gb` catches a queue that claims work while nothing is loaded | | `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 is asked to yield **first** among idle tenants — it never protects idle memory, and never interrupts work | +| `overclock_profile` | GPU profile applied while this tenant is the active workload | +| `events` | Optional stream (e.g. a websocket) used purely as a wake-up, so reaction is sub-second rather than waiting for the next poll | 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). +**Priority orders, it does not veto.** An idle tenant is not using its VRAM, so +outranking the demander is no reason to keep it; busy tenants are never interrupted +whatever their rank. Getting this wrong broke both directions in turn — with the LLM +ranked above diffusion, ComfyUI could never preempt Ollama (the service's central +behaviour), and once the ranks were swapped, a starved Ollama could no longer reclaim +from an idle ComfyUI. Diffusion now outranks the LLM, whose weights reload from page +cache in seconds. + `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 @@ -228,7 +238,7 @@ explaining that, rather than silently doing nothing. ## 1a. Tests ```bash -/home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 242 passed in ~3.8s +/home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 244 passed in ~3.8s ``` Hermetic: no GPU, no network, no sleeps. An autouse fixture stubs `overclock_manager._sh` diff --git a/tenants.json b/tenants.json index 8d68c6b..cc6f83f 100644 --- a/tenants.json +++ b/tenants.json @@ -2,7 +2,7 @@ { "name": "ollama", "kind": "llm", - "priority": 60, + "priority": 50, "match": { "names": [ "ollama" @@ -33,7 +33,7 @@ { "name": "comfyui", "kind": "diffusion", - "priority": 50, + "priority": 60, "match": { "cmdline": [ "comfyui", diff --git a/tenants.py b/tenants.py index 67a6837..efb4226 100644 --- a/tenants.py +++ b/tenants.py @@ -70,6 +70,22 @@ class BusyProbe: stale_after_s: float = 90.0 +@dataclass +class EventSource: + """A stream that tells us *when* to look, not what to think. + + ComfyUI publishes a websocket, and the original listener parsed its message types to + decide what was happening -- which meant understanding one application's schema. Any + message is instead treated purely as a wake-up: re-run this tenant's busy probe now + rather than waiting for the next poll. That gives sub-second reaction to any + application with an event stream, with no knowledge of what it emits. + """ + type: str = "none" # none | websocket + url: Optional[str] = None + reconnect_backoff_s: float = 2.0 + max_backoff_s: float = 15.0 + + @dataclass class ReleaseStrategy: """How to ask a tenant to give VRAM back.""" @@ -96,9 +112,14 @@ class GpuTenant: # 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 + # GPU profile to apply while this tenant is the active workload. Clock and power + # tuning is workload-specific -- diffusion is compute bound, LLM decode is bandwidth + # bound -- and that was previously switched by application name in the arbitrator. + overclock_profile: Optional[str] = None match: ProcessMatch = field(default_factory=ProcessMatch) busy: BusyProbe = field(default_factory=BusyProbe) release: ReleaseStrategy = field(default_factory=ReleaseStrategy) + events: EventSource = field(default_factory=EventSource) notes: str = "" @property @@ -116,10 +137,12 @@ def _tenant_from_dict(d: Dict[str, Any]) -> GpuTenant: enabled=d.get("enabled", True), priority=int(d.get("priority", 50)), needs_vram_gb=float(d.get("needs_vram_gb", 0.0)), + overclock_profile=d.get("overclock_profile"), 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 {})), + events=EventSource(**(d.get("events") or {})), notes=d.get("notes", ""), ) @@ -129,9 +152,14 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [ { "name": "ollama", "kind": KIND_LLM, - "priority": 60, + # Lower than ComfyUI on purpose: an interactive diffusion job preempts the LLM, + # whose weights stay in the page cache and reload in seconds. Getting this the + # wrong way round silently disabled the service's central behaviour -- ComfyUI + # could never reclaim from Ollama. + "priority": 50, "needs_vram_gb": 4.0, "idle_release_after_s": 0.0, + "overclock_profile": "ollama", "match": {"names": ["ollama"], "cmdline": ["llama-server", "ollama"]}, "busy": {"type": "http_count", "url": "http://localhost:11434/api/ps", "count_keys": ["models"]}, @@ -143,9 +171,10 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [ { "name": "comfyui", "kind": KIND_DIFFUSION, - "priority": 50, + "priority": 60, "needs_vram_gb": 6.0, "idle_release_after_s": 30.0, + "overclock_profile": "comfy", "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"], @@ -153,6 +182,7 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [ "release": {"type": "http_post", "url": "http://127.0.0.1:8188/free", "body": {"unload_models": True, "free_memory": True}, "timeout_s": 30.0}, + "events": {"type": "websocket", "url": "ws://127.0.0.1:8188/ws?clientId=hyperswap"}, "notes": "Leaves dead jobs in queue_running; the queue flag is corroborated " "against its own VRAM before being believed.", }, @@ -210,7 +240,7 @@ def load_tenants(force: bool = False) -> List[GpuTenant]: base = defaults_by_name.get(d.get("name")) if base: merged = {**base, **d} - for key in ("match", "busy", "release"): + for key in ("match", "busy", "release", "events"): if isinstance(base.get(key), dict): merged[key] = {**base[key], **(d.get(key) or {})} d = merged @@ -364,12 +394,19 @@ def plan_release(demanding: str, tenants_state: List[Dict[str, Any]], return {"possible": True, "reason": "enough VRAM is already free", "release": [], "shortfall_gb": 0.0} + # Any idle reclaimable tenant is a candidate, whatever its priority. An idle tenant + # is not using its VRAM, so ranking above the demander should not protect it -- an + # earlier version filtered on priority and thereby broke both directions in turn: + # ComfyUI could not preempt Ollama, and once that was corrected a starved Ollama + # could no longer reclaim from an idle ComfyUI. + # + # Priority decides who is asked *first* (lowest gives up memory soonest) and, being + # applied only to idle tenants, never interrupts work. 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))) @@ -385,7 +422,7 @@ def plan_release(demanding: str, tenants_state: List[Dict[str, Any]], {"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")} + "enough was freed without it")} for s in tenants_state if s["name"] != demanding and s.get("vram_gb", 0) > 0 and s["name"] not in plan ] diff --git a/tests/test_tenants.py b/tests/test_tenants.py index e008f32..0c4e93f 100644 --- a/tests/test_tenants.py +++ b/tests/test_tenants.py @@ -224,10 +224,42 @@ class TestReleasePlanning: assert blockers["stt-relay"] == "declares no release mechanism" assert plan["possible"] is False - def test_higher_priority_tenants_are_not_victims(self): + def test_an_idle_tenant_yields_even_if_it_outranks_the_demander(self): + """Priority orders who is asked first; it does not protect idle memory. + + Filtering candidates by priority broke both directions in turn: with the LLM + ranked above diffusion, ComfyUI could never preempt Ollama -- the service's + central behaviour -- and once the ranks were swapped, a starved Ollama could no + longer reclaim from an idle ComfyUI. An idle tenant is not using its VRAM, so + outranking the demander is not a reason to keep it. + """ 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"] + assert plan["release"] == ["comfyui"] + + def test_priority_decides_who_is_asked_first(self): + state = [ + {"name": "demander", "priority": 50, "vram_gb": 0.0, "busy": True, + "reclaimable": True}, + {"name": "high", "priority": 90, "vram_gb": 4.0, "busy": False, + "reclaimable": True}, + {"name": "low", "priority": 10, "vram_gb": 4.0, "busy": False, + "reclaimable": True}, + ] + plan = T.plan_release("demander", state, free_gb=0.0, needed_gb=5.0) + # The lowest-priority idle tenant gives up memory first. + assert plan["release"][0] == "low" + + def test_busy_work_is_never_interrupted_whatever_the_priority(self): + state = [ + {"name": "demander", "priority": 99, "vram_gb": 0.0, "busy": True, + "reclaimable": True}, + {"name": "worker", "priority": 1, "vram_gb": 8.0, "busy": True, + "reclaimable": True}, + ] + plan = T.plan_release("demander", state, free_gb=0.0, needed_gb=8.0) + assert plan["release"] == [] + assert plan["blockers"][0]["why"] == "busy" def test_lowest_priority_is_released_first(self): state = self._state() + [ diff --git a/vram_arbitrator.py b/vram_arbitrator.py index fc5be6d..5c04f93 100644 --- a/vram_arbitrator.py +++ b/vram_arbitrator.py @@ -932,6 +932,7 @@ class AutoArbitrator: def __init__(self): self.running = False self.ws_task: Optional[asyncio.Task] = None + self.event_tasks: List[asyncio.Task] = [] self.poll_task: Optional[asyncio.Task] = None self.idle_task: Optional[asyncio.Task] = None self.last_yield_time = 0.0 @@ -958,6 +959,8 @@ class AutoArbitrator: self._peak_comfy_bytes = 0 self.comfy_stale_job: Optional[str] = None self._idle_since: Dict[str, float] = {} + self._last_event_wake = 0.0 + self.event_sources: Dict[str, str] = {} self._last_tenant_state: Optional[Dict[str, Any]] = None self.last_arbitration: Optional[Dict[str, Any]] = None self.last_watchdog_error: Optional[str] = None @@ -976,6 +979,12 @@ class AutoArbitrator: return self.running = True self.ws_task = asyncio.create_task(self._ws_listener()) + for t in tenants_mod.load_tenants(): + if t.enabled and t.events.type == "websocket" and t.events.url: + self.event_tasks.append(asyncio.create_task( + self._event_listener(t.name, t.events.url, + t.events.reconnect_backoff_s, + t.events.max_backoff_s))) self.poll_task = asyncio.create_task(self._poll_watchdog()) self.idle_task = asyncio.create_task(self._idle_purge_loop()) logger.info("AutoArbitrator background engine started (Bidirectional).") @@ -988,9 +997,10 @@ class AutoArbitrator: async def stop(self): self.running = False - for task in (self.ws_task, self.poll_task, self.idle_task): + for task in [self.ws_task, self.poll_task, self.idle_task, *self.event_tasks]: if task: task.cancel() + self.event_tasks.clear() await close_clients() logger.info("AutoArbitrator background engine stopped.") @@ -1100,6 +1110,47 @@ class AutoArbitrator: return {"purged": True, "free_gb": round(snap["free_bytes"] / (1024**3), 2)} return {"purged": False, "free_gb": round(free_gb, 2), "reason": "ComfyUI holds no VRAM"} + async def _event_listener(self, tenant_name: str, url: str, + backoff_s: float, max_backoff_s: float) -> None: + """Wake on a tenant's event stream instead of waiting for the next poll. + + Deliberately does not parse the messages. The previous listener understood + ComfyUI's schema -- status/execution_start/executing/execution_success -- which + tied the fast path to one application. Treating any message as "look now" and + letting the tenant's own busy probe decide gives the same sub-second reaction + for any application that emits anything on state change. + """ + backoff = backoff_s + while self.running: + try: + async with websockets.connect(url, ping_interval=10, ping_timeout=10) as ws: + self.event_sources[tenant_name] = "connected" + self.connected_ws = True + backoff = backoff_s + logger.info(f"Event source connected for '{tenant_name}': {url}") + while self.running: + await ws.recv() + # Coalesce bursts: a single graph emits many messages, and one + # arbitration pass per burst is enough. + now = time.time() + if now - self._last_event_wake < 0.25: + continue + self._last_event_wake = now + self.stats["event_wakeups"] = self.stats.get("event_wakeups", 0) + 1 + try: + await self._arbitrate() + except Exception as e: + logger.debug(f"arbitration from event failed: {e}") + except (websockets.exceptions.ConnectionClosed, OSError, asyncio.CancelledError): + self.event_sources[tenant_name] = "disconnected" + self.connected_ws = False + except Exception as e: + self.event_sources[tenant_name] = f"error: {str(e)[:60]}" + self.connected_ws = False + logger.debug(f"event source error for '{tenant_name}': {e}") + await asyncio.sleep(backoff) + backoff = min(backoff * 1.5, max_backoff_s) + async def _ws_listener(self): client_id = "hyperswap-arbitrator" ws_url = f"ws://127.0.0.1:8188/ws?clientId={client_id}" @@ -1229,6 +1280,7 @@ class AutoArbitrator: "below_floor": bool(probe.get("below_floor")), "reclaimable": t.reclaimable, "needs_vram_gb": t.needs_vram_gb, + "overclock_profile": t.overclock_profile, "idle_release_after_s": t.idle_release_after_s, "reason": probe.get("reason"), }) @@ -1251,6 +1303,22 @@ class AutoArbitrator: self.stats["tenant_releases"] = self.stats.get("tenant_releases", 0) + 1 return res + IDLE_PROFILE = "balanced" + + def _apply_profile_for_active(self, state: List[Dict[str, Any]]) -> None: + """Apply the GPU profile declared by whichever tenant is currently working. + + This used to be two calls naming 'comfy' and 'ollama' directly, so a third + application could never get tuned clocks. The highest-priority busy tenant wins; + with nothing working the card returns to the idle profile. + """ + busy = [s for s in state if s["busy"] and s.get("overclock_profile")] + if busy: + busy.sort(key=lambda s: -s["priority"]) + self._apply_oc_profile(busy[0]["overclock_profile"]) + else: + self._apply_oc_profile(self.IDLE_PROFILE) + async def _arbitrate(self) -> None: """Generic arbitration over any number of tenants. @@ -1262,6 +1330,7 @@ class AutoArbitrator: """ state = await self._tenant_state() free_gb = self._last_tenant_state["free_gb"] + self._apply_profile_for_active(state) # 1. Starvation: highest-priority demanding tenant first. for s in sorted(state, key=lambda x: -x["priority"]): @@ -1379,6 +1448,7 @@ class AutoArbitrator: "oc_profile": self.oc_profile, "counters": dict(self.stats), "comfy_stale_job": self.comfy_stale_job, + "event_sources": dict(self.event_sources), "last_arbitration": self.last_arbitration, "tenant_state": self._last_tenant_state, "watchdog_branches": dict(self.watchdog_branches),