Let any tenant declare its GPU profile and event source; fix priority semantics
Two remaining pieces of the two-application coupling are gone. Overclock profiles were switched by naming 'comfy' and 'ollama' directly, so a third application could never get tuned clocks. A tenant declares overclock_profile and the arbitrator applies whichever the highest-priority *working* tenant asks for, falling back to the idle profile when nothing is running. The websocket listener parsed ComfyUI's message schema -- status, execution_start, executing, execution_success -- which tied the fast path to one application. An event source is now declarative and the messages are not parsed at all: any message means "look now", and the tenant's own busy probe decides what is true. That gives the same sub-second reaction to any application that emits anything on state change, with no knowledge of what it emits. Generalising this exposed a design error in the priority rule I had introduced. plan_release excluded candidates ranking above the demander, which broke both directions in turn. With the LLM at priority 60 and diffusion at 50, ComfyUI could never reclaim from Ollama -- the premise the whole service is built on, and preserved until now only by the ComfyUI-specific trigger that was about to be removed. Swapping the ranks then broke the reverse: a starved Ollama could no longer reclaim from an idle ComfyUI. Priority now orders rather than vetoes. Any idle reclaimable tenant is a candidate, because an idle tenant is not using its VRAM; priority decides who is asked first, and busy tenants are never interrupted whatever their rank. Diffusion outranks the LLM, whose weights reload from page cache in seconds. All three cases are pinned by tests, including that busy work is never interrupted even by a far higher-priority demander. Tests: 244 (was 242). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
14
README.md
14
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`) |
|
| `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 |
|
| `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` |
|
| `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
|
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
|
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
|
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).
|
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
|
`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
|
`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
|
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
|
## 1a. Tests
|
||||||
|
|
||||||
```bash
|
```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`
|
Hermetic: no GPU, no network, no sleeps. An autouse fixture stubs `overclock_manager._sh`
|
||||||
|
|||||||
@@ -2,7 +2,7 @@
|
|||||||
{
|
{
|
||||||
"name": "ollama",
|
"name": "ollama",
|
||||||
"kind": "llm",
|
"kind": "llm",
|
||||||
"priority": 60,
|
"priority": 50,
|
||||||
"match": {
|
"match": {
|
||||||
"names": [
|
"names": [
|
||||||
"ollama"
|
"ollama"
|
||||||
@@ -33,7 +33,7 @@
|
|||||||
{
|
{
|
||||||
"name": "comfyui",
|
"name": "comfyui",
|
||||||
"kind": "diffusion",
|
"kind": "diffusion",
|
||||||
"priority": 50,
|
"priority": 60,
|
||||||
"match": {
|
"match": {
|
||||||
"cmdline": [
|
"cmdline": [
|
||||||
"comfyui",
|
"comfyui",
|
||||||
|
|||||||
47
tenants.py
47
tenants.py
@@ -70,6 +70,22 @@ class BusyProbe:
|
|||||||
stale_after_s: float = 90.0
|
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
|
@dataclass
|
||||||
class ReleaseStrategy:
|
class ReleaseStrategy:
|
||||||
"""How to ask a tenant to give VRAM back."""
|
"""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,
|
# back. Iterating on a ComfyUI workflow should not pay a reload between every run,
|
||||||
# so this is deliberately not immediate.
|
# so this is deliberately not immediate.
|
||||||
idle_release_after_s: float = 30.0
|
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)
|
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)
|
||||||
|
events: EventSource = field(default_factory=EventSource)
|
||||||
notes: str = ""
|
notes: str = ""
|
||||||
|
|
||||||
@property
|
@property
|
||||||
@@ -116,10 +137,12 @@ def _tenant_from_dict(d: Dict[str, Any]) -> GpuTenant:
|
|||||||
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)),
|
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)),
|
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 {})),
|
||||||
|
events=EventSource(**(d.get("events") or {})),
|
||||||
notes=d.get("notes", ""),
|
notes=d.get("notes", ""),
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -129,9 +152,14 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [
|
|||||||
{
|
{
|
||||||
"name": "ollama",
|
"name": "ollama",
|
||||||
"kind": KIND_LLM,
|
"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,
|
"needs_vram_gb": 4.0,
|
||||||
"idle_release_after_s": 0.0,
|
"idle_release_after_s": 0.0,
|
||||||
|
"overclock_profile": "ollama",
|
||||||
"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"]},
|
||||||
@@ -143,9 +171,10 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [
|
|||||||
{
|
{
|
||||||
"name": "comfyui",
|
"name": "comfyui",
|
||||||
"kind": KIND_DIFFUSION,
|
"kind": KIND_DIFFUSION,
|
||||||
"priority": 50,
|
"priority": 60,
|
||||||
"needs_vram_gb": 6.0,
|
"needs_vram_gb": 6.0,
|
||||||
"idle_release_after_s": 30.0,
|
"idle_release_after_s": 30.0,
|
||||||
|
"overclock_profile": "comfy",
|
||||||
"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"],
|
||||||
@@ -153,6 +182,7 @@ DEFAULT_TENANTS: List[Dict[str, Any]] = [
|
|||||||
"release": {"type": "http_post", "url": "http://127.0.0.1:8188/free",
|
"release": {"type": "http_post", "url": "http://127.0.0.1:8188/free",
|
||||||
"body": {"unload_models": True, "free_memory": True},
|
"body": {"unload_models": True, "free_memory": True},
|
||||||
"timeout_s": 30.0},
|
"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 "
|
"notes": "Leaves dead jobs in queue_running; the queue flag is corroborated "
|
||||||
"against its own VRAM before being believed.",
|
"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"))
|
base = defaults_by_name.get(d.get("name"))
|
||||||
if base:
|
if base:
|
||||||
merged = {**base, **d}
|
merged = {**base, **d}
|
||||||
for key in ("match", "busy", "release"):
|
for key in ("match", "busy", "release", "events"):
|
||||||
if isinstance(base.get(key), dict):
|
if isinstance(base.get(key), dict):
|
||||||
merged[key] = {**base[key], **(d.get(key) or {})}
|
merged[key] = {**base[key], **(d.get(key) or {})}
|
||||||
d = merged
|
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",
|
return {"possible": True, "reason": "enough VRAM is already free",
|
||||||
"release": [], "shortfall_gb": 0.0}
|
"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 = [
|
candidates = [
|
||||||
s for s in tenants_state
|
s for s in tenants_state
|
||||||
if s["name"] != demanding
|
if s["name"] != demanding
|
||||||
and s.get("reclaimable")
|
and s.get("reclaimable")
|
||||||
and not s.get("busy")
|
and not s.get("busy")
|
||||||
and s.get("priority", 0) <= demander.get("priority", 0)
|
|
||||||
and s.get("vram_gb", 0) > 0
|
and s.get("vram_gb", 0) > 0
|
||||||
]
|
]
|
||||||
candidates.sort(key=lambda s: (s.get("priority", 0), -s.get("vram_gb", 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),
|
{"name": s["name"], "vram_gb": s.get("vram_gb", 0.0),
|
||||||
"why": ("busy" if s.get("busy") else
|
"why": ("busy" if s.get("busy") else
|
||||||
"declares no release mechanism" if not s.get("reclaimable") else
|
"declares no release mechanism" if not s.get("reclaimable") else
|
||||||
"higher priority")}
|
"enough was freed without it")}
|
||||||
for s in tenants_state
|
for s in tenants_state
|
||||||
if s["name"] != demanding and s.get("vram_gb", 0) > 0 and s["name"] not in plan
|
if s["name"] != demanding and s.get("vram_gb", 0) > 0 and s["name"] not in plan
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -224,10 +224,42 @@ class TestReleasePlanning:
|
|||||||
assert blockers["stt-relay"] == "declares no release mechanism"
|
assert blockers["stt-relay"] == "declares no release mechanism"
|
||||||
assert plan["possible"] is False
|
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})
|
state = self._state(comfyui={"priority": 99, "busy": False})
|
||||||
plan = T.plan_release("ollama", state, free_gb=1.0, needed_gb=14.9)
|
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):
|
def test_lowest_priority_is_released_first(self):
|
||||||
state = self._state() + [
|
state = self._state() + [
|
||||||
|
|||||||
@@ -932,6 +932,7 @@ class AutoArbitrator:
|
|||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.running = False
|
self.running = False
|
||||||
self.ws_task: Optional[asyncio.Task] = None
|
self.ws_task: Optional[asyncio.Task] = None
|
||||||
|
self.event_tasks: List[asyncio.Task] = []
|
||||||
self.poll_task: Optional[asyncio.Task] = None
|
self.poll_task: Optional[asyncio.Task] = None
|
||||||
self.idle_task: Optional[asyncio.Task] = None
|
self.idle_task: Optional[asyncio.Task] = None
|
||||||
self.last_yield_time = 0.0
|
self.last_yield_time = 0.0
|
||||||
@@ -958,6 +959,8 @@ class AutoArbitrator:
|
|||||||
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._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_tenant_state: Optional[Dict[str, Any]] = None
|
||||||
self.last_arbitration: 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
|
||||||
@@ -976,6 +979,12 @@ class AutoArbitrator:
|
|||||||
return
|
return
|
||||||
self.running = True
|
self.running = True
|
||||||
self.ws_task = asyncio.create_task(self._ws_listener())
|
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.poll_task = asyncio.create_task(self._poll_watchdog())
|
||||||
self.idle_task = asyncio.create_task(self._idle_purge_loop())
|
self.idle_task = asyncio.create_task(self._idle_purge_loop())
|
||||||
logger.info("AutoArbitrator background engine started (Bidirectional).")
|
logger.info("AutoArbitrator background engine started (Bidirectional).")
|
||||||
@@ -988,9 +997,10 @@ class AutoArbitrator:
|
|||||||
|
|
||||||
async def stop(self):
|
async def stop(self):
|
||||||
self.running = False
|
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:
|
if task:
|
||||||
task.cancel()
|
task.cancel()
|
||||||
|
self.event_tasks.clear()
|
||||||
await close_clients()
|
await close_clients()
|
||||||
logger.info("AutoArbitrator background engine stopped.")
|
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": True, "free_gb": round(snap["free_bytes"] / (1024**3), 2)}
|
||||||
return {"purged": False, "free_gb": round(free_gb, 2), "reason": "ComfyUI holds no VRAM"}
|
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):
|
async def _ws_listener(self):
|
||||||
client_id = "hyperswap-arbitrator"
|
client_id = "hyperswap-arbitrator"
|
||||||
ws_url = f"ws://127.0.0.1:8188/ws?clientId={client_id}"
|
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")),
|
"below_floor": bool(probe.get("below_floor")),
|
||||||
"reclaimable": t.reclaimable,
|
"reclaimable": t.reclaimable,
|
||||||
"needs_vram_gb": t.needs_vram_gb,
|
"needs_vram_gb": t.needs_vram_gb,
|
||||||
|
"overclock_profile": t.overclock_profile,
|
||||||
"idle_release_after_s": t.idle_release_after_s,
|
"idle_release_after_s": t.idle_release_after_s,
|
||||||
"reason": probe.get("reason"),
|
"reason": probe.get("reason"),
|
||||||
})
|
})
|
||||||
@@ -1251,6 +1303,22 @@ class AutoArbitrator:
|
|||||||
self.stats["tenant_releases"] = self.stats.get("tenant_releases", 0) + 1
|
self.stats["tenant_releases"] = self.stats.get("tenant_releases", 0) + 1
|
||||||
return res
|
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:
|
async def _arbitrate(self) -> None:
|
||||||
"""Generic arbitration over any number of tenants.
|
"""Generic arbitration over any number of tenants.
|
||||||
|
|
||||||
@@ -1262,6 +1330,7 @@ class AutoArbitrator:
|
|||||||
"""
|
"""
|
||||||
state = await self._tenant_state()
|
state = await self._tenant_state()
|
||||||
free_gb = self._last_tenant_state["free_gb"]
|
free_gb = self._last_tenant_state["free_gb"]
|
||||||
|
self._apply_profile_for_active(state)
|
||||||
|
|
||||||
# 1. Starvation: highest-priority demanding tenant first.
|
# 1. Starvation: highest-priority demanding tenant first.
|
||||||
for s in sorted(state, key=lambda x: -x["priority"]):
|
for s in sorted(state, key=lambda x: -x["priority"]):
|
||||||
@@ -1379,6 +1448,7 @@ class AutoArbitrator:
|
|||||||
"oc_profile": self.oc_profile,
|
"oc_profile": self.oc_profile,
|
||||||
"counters": dict(self.stats),
|
"counters": dict(self.stats),
|
||||||
"comfy_stale_job": self.comfy_stale_job,
|
"comfy_stale_job": self.comfy_stale_job,
|
||||||
|
"event_sources": dict(self.event_sources),
|
||||||
"last_arbitration": self.last_arbitration,
|
"last_arbitration": self.last_arbitration,
|
||||||
"tenant_state": self._last_tenant_state,
|
"tenant_state": self._last_tenant_state,
|
||||||
"watchdog_branches": dict(self.watchdog_branches),
|
"watchdog_branches": dict(self.watchdog_branches),
|
||||||
|
|||||||
Reference in New Issue
Block a user