From 48fbde040cea8747639db3a433db723a44d82805 Mon Sep 17 00:00:00 2001 From: drjones Date: Mon, 7 Sep 2026 17:41:55 -0700 Subject: [PATCH] Test the job queue, and surface it in the dashboard and MCP The queue shipped with no tests despite having produced four bugs during development, which is the wrong order. 21 tests now cover the parts whose failure modes are not obvious from reading the code: ordering by priority then FIFO, that depth is genuinely unbounded and survives a restart, that running work is never cancelled, that a job abandoned mid-run is requeued rather than left RUNNING forever, and that an LLM job's VRAM requirement comes from its own model rather than a tenant-wide figure -- the mistake that dispatched a 14.9 GB model into 8 GB of free memory and killed llama-server three times. The dashboard had no view of the queue at all, so a stuck queue was indistinguishable from an empty one. The Job Queue panel shows what is running and what it released to get there, pending jobs in execution order with their wait time and a cancel control, and -- when the scheduler is blocked -- how long it has been waiting and why. MCP gains queue_job, get_job_queue and cancel_job, so an agent can line work up rather than firing a request and hoping the GPU is free. Tests: 271. Co-Authored-By: Claude Opus 5 --- README.md | 7 +- mcp_server.py | 28 +++++++ static/app.js | 71 +++++++++++++++++ static/index.html | 22 ++++++ tests/test_jobs.py | 184 +++++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 310 insertions(+), 2 deletions(-) create mode 100644 tests/test_jobs.py diff --git a/README.md b/README.md index d9376e8..6a7e6d2 100644 --- a/README.md +++ b/README.md @@ -194,7 +194,10 @@ scarce resource this service exists to hand between applications, and overlappin would just recreate the contention it resolves. `GET /api/jobs` · `GET /api/jobs/{id}` · `DELETE /api/jobs/{id}` (pending only — running -work is never killed) · `DELETE /api/jobs` to clear the queue. +work is never killed) · `DELETE /api/jobs` to clear the queue. Agents get the same through +MCP (`queue_job`, `get_job_queue`, `cancel_job`), and the dashboard's **Job Queue** panel +shows what is running, what it released to get there, and why the scheduler is waiting if +it is. **A job that cannot run yet waits; a job that can never run fails with the reason.** Dispatching into insufficient VRAM does not fail gracefully — it kills `llama-server` @@ -266,7 +269,7 @@ explaining that, rather than silently doing nothing. ## 1a. Tests ```bash -/home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 250 passed in ~3.8s +/home/drjones/comfy-mcp-venv/bin/python -m pytest tests/ -q # 271 passed in ~4.0s ``` Hermetic: no GPU, no network, no sleeps. An autouse fixture stubs `overclock_manager._sh` diff --git a/mcp_server.py b/mcp_server.py index f5650a3..b1d685c 100644 --- a/mcp_server.py +++ b/mcp_server.py @@ -10,6 +10,7 @@ from mcp.server import MCPServer import autotune import engines import health +import jobs as jobs_mod import overclock_manager import ram_optimizer import telemetry_store @@ -146,6 +147,33 @@ async def get_engine_config() -> str: return json.dumps(await engines.get_engine_config(), indent=2, default=str) +@mcp.tool() +def queue_job(tenant: str, payload: Dict[str, Any], priority: Optional[int] = None, + label: Optional[str] = None) -> str: + """Queue work for a GPU application without waiting for it. + + tenant: 'ollama' (payload: model, prompt, options) or 'comfyui' (payload: {"prompt": + }). The queue is on disk, so there is no depth limit; jobs run one at a + time, highest priority first, with VRAM arbitrated before each starts.""" + return json.dumps(jobs_mod.submit(tenant, payload, priority, label), indent=2, + default=str) + + +@mcp.tool() +def get_job_queue(state: Optional[str] = None, limit: int = 50) -> str: + """Queued and recent jobs, plus what the scheduler is doing and why it may be + waiting. Pending jobs are listed in the order they will run.""" + return json.dumps({"jobs": jobs_mod.listing(state, limit), + "scheduler": jobs_mod.scheduler.get_status()}, + indent=2, default=str) + + +@mcp.tool() +def cancel_job(job_id: str) -> str: + """Cancel a job that has not started. Running work is never killed.""" + return json.dumps(jobs_mod.cancel(job_id), indent=2, default=str) + + @mcp.tool() async def check_system_health() -> str: """Check every dependency HyperSwap needs (NVML, sudo nvidia-smi, fan control via the diff --git a/static/app.js b/static/app.js index c008843..3929659 100644 --- a/static/app.js +++ b/static/app.js @@ -1189,3 +1189,74 @@ function renderTenants(arb, gpu) { `${a.demanding} short by ${a.shortfall_gb} GB — ${a.reason}` + (blockers ? `
blocked by: ${blockers}
` : ''); } + + +// ---------------------------------------------------------------- job queue + +async function fetchQueue() { + const body = document.getElementById('queue-body'); + if (!body) return; + try { + const d = await (await fetch('/api/jobs?limit=40')).json(); + const s = d.scheduler || {}; + document.getElementById('queue-summary').textContent = + `${s.queue_depth ?? 0} queued · ${s.completed ?? 0} done · ${s.failed ?? 0} failed`; + + // What the scheduler is doing right now, including why it is waiting. A job that + // cannot get VRAM used to sit silent, which made a stuck queue indistinguishable + // from an empty one. + const cur = document.getElementById('queue-current'); + if (s.current) { + cur.innerHTML = `running ` + + `${s.current.tenant}` + + ` · ${s.current.label || s.current.id}` + + (s.current.room && s.current.room.released && s.current.room.released.length + ? ` · released ${s.current.room.released.join(', ')}` : ''); + } else if (s.blocked) { + cur.innerHTML = `waiting ${s.blocked.waited_s}s ` + + `${s.blocked.tenant}` + + ` · ${s.blocked.reason || ''}`; + } else { + cur.innerHTML = 'scheduler idle'; + } + + const colour = { + pending: 'text-slate-400', running: 'text-emerald-400', done: 'text-cyan-500', + failed: 'text-rose-400', cancelled: 'text-slate-600', + }; + body.innerHTML = (d.jobs || []).map(job => { + const when = job.state === 'pending' + ? `waiting ${job.waiting_s}s` + : (job.duration_s != null ? `${job.duration_s}s` : ''); + return `
+ + ${job.state} + p${job.priority} + ${job.tenant} + ${job.label ? ' · ' + job.label : ''} + + ${when}${ + job.state === 'pending' + ? ` ` + : ''} +
${job.error ? `
${job.error}
` : ''}`; + }).join('') || '
No jobs yet.
'; + } catch (e) { + body.innerHTML = `
${e}
`; + } +} + +async function cancelJob(id) { + await fetch(`/api/jobs/${id}`, { method: 'DELETE' }); + fetchQueue(); +} + +async function clearQueue() { + await fetch('/api/jobs', { method: 'DELETE' }); + fetchQueue(); +} + +document.addEventListener('DOMContentLoaded', () => { + fetchQueue(); + setInterval(fetchQueue, 2000); +}); diff --git a/static/index.html b/static/index.html index 4da909c..0403204 100644 --- a/static/index.html +++ b/static/index.html @@ -681,6 +681,28 @@ + + +
+
+
+
+ +
+
+

Job Queue

+

Work lined up across every application, highest priority first

+
+
+
+ — + +
+
+
+
+
+
diff --git a/tests/test_jobs.py b/tests/test_jobs.py new file mode 100644 index 0000000..b320a47 --- /dev/null +++ b/tests/test_jobs.py @@ -0,0 +1,184 @@ +"""Tests for the cross-tenant job queue. + +Every case below corresponds to something that actually went wrong while building this, +because the failure modes are not obvious from the code: + + * dispatching without room does not fail gracefully -- a CUDA OOM kills llama-server, + and three queued jobs were destroyed in a row; + * the VRAM requirement is a property of the job's model, not of the tenant, and a flat + 4 GB let a 14.9 GB model be dispatched into 8 GB of free memory; + * waiting forever is as wrong as failing immediately, when the job can never fit; + * an exception during dispatch left the row RUNNING while the scheduler moved on. +""" +import asyncio +import json +import time + +import pytest + +import jobs as J +import telemetry_store + + +@pytest.fixture +def queue(tmp_path, monkeypatch): + monkeypatch.setattr(telemetry_store, "DB_PATH", str(tmp_path / "q.db")) + J.init() + return tmp_path + + +class TestQueueBasics: + def test_submit_and_read_back(self, queue): + r = J.submit("ollama", {"model": "m"}, label="first") + assert r["success"] + job = J.get(r["id"]) + assert job["state"] == J.PENDING and job["label"] == "first" + assert job["payload"] == {"model": "m"} + + def test_unknown_tenant_is_rejected(self, queue): + assert J.submit("nope", {})["success"] is False + + def test_priority_defaults_to_the_tenants_own(self, queue): + import tenants as T + r = J.submit("comfyui", {}) + assert r["priority"] == T.get_tenant("comfyui").priority + + def test_pending_jobs_are_listed_in_execution_order(self, queue): + J.submit("ollama", {}, priority=10, label="low") + J.submit("ollama", {}, priority=90, label="high") + J.submit("ollama", {}, priority=50, label="mid") + order = [j["label"] for j in J.listing(J.PENDING)] + assert order == ["high", "mid", "low"] + + def test_equal_priority_is_first_in_first_out(self, queue): + a = J.submit("ollama", {}, priority=50, label="a")["id"] + time.sleep(0.01) + J.submit("ollama", {}, priority=50, label="b") + assert J._next_job()["id"] == a + + def test_the_queue_has_no_depth_limit(self, queue): + # "Unbounded" is the point; it lives on disk, not in memory. + for i in range(500): + J.submit("ollama", {}, label=f"j{i}") + assert J.stats()["queue_depth"] == 500 + + def test_queue_survives_a_restart(self, queue): + J.submit("ollama", {}, label="persisted") + J._cache = None # nothing in-process is holding it + assert [j["label"] for j in J.listing(J.PENDING)] == ["persisted"] + + +class TestCancellation: + def test_pending_jobs_can_be_cancelled(self, queue): + jid = J.submit("ollama", {})["id"] + assert J.cancel(jid)["success"] is True + assert J.get(jid)["state"] == J.CANCELLED + + def test_running_work_is_never_cancelled(self, queue): + # This service frees VRAM by asking, never by killing work in flight. + jid = J.submit("ollama", {})["id"] + J._mark(jid, J.RUNNING, started_at=time.time()) + assert J.cancel(jid)["success"] is False + assert J.get(jid)["state"] == J.RUNNING + + def test_clearing_the_queue_leaves_running_work_alone(self, queue): + running = J.submit("ollama", {})["id"] + J._mark(running, J.RUNNING, started_at=time.time()) + J.submit("ollama", {}) + J.submit("ollama", {}) + assert J.clear_pending()["cancelled"] == 2 + assert J.get(running)["state"] == J.RUNNING + + +class TestOrphanRecovery: + def test_jobs_left_running_by_a_dead_process_are_requeued(self, queue): + """RUNNING means "this process is working on it". + + An exception during dispatch left the row RUNNING while the scheduler moved on, + so the job never finished and never retried. + """ + jid = J.submit("ollama", {})["id"] + J._mark(jid, J.RUNNING, started_at=time.time()) + assert J.requeue_orphans() == 1 + job = J.get(jid) + assert job["state"] == J.PENDING and job["started_at"] is None + + def test_finished_jobs_are_untouched_by_recovery(self, queue): + done = J.submit("ollama", {})["id"] + J._mark(done, J.DONE, finished_at=time.time()) + assert J.requeue_orphans() == 0 + assert J.get(done)["state"] == J.DONE + + +class TestPerJobVramRequirement: + """A tenant-wide figure cannot be right for an LLM.""" + + def test_llm_requirement_comes_from_the_model_being_loaded(self, queue, monkeypatch): + import vram_arbitrator + monkeypatch.setattr(vram_arbitrator, "_model_size_bytes", + lambda m: int(12.87 * 1024 ** 3)) + s = J.Scheduler() + # 12.87 GB on disk occupies ~14.9 GB once context and KV cache are allocated. + assert 14.5 < s._job_vram_requirement("ollama", {"model": "big"}) < 15.5 + + def test_a_small_model_needs_correspondingly_less(self, queue, monkeypatch): + import vram_arbitrator + monkeypatch.setattr(vram_arbitrator, "_model_size_bytes", + lambda m: int(1.96 * 1024 ** 3)) + s = J.Scheduler() + assert s._job_vram_requirement("ollama", {"model": "small"}) < 3.0 + + def test_falls_back_to_the_tenant_figure_when_the_model_is_unknown(self, queue, monkeypatch): + import vram_arbitrator + monkeypatch.setattr(vram_arbitrator, "_model_size_bytes", lambda m: 0) + s = J.Scheduler() + import tenants as T + assert (s._job_vram_requirement("ollama", {"model": "?"}) + == T.get_tenant("ollama").needs_vram_gb) + + def test_non_llm_tenants_use_their_declared_figure(self, queue): + import tenants as T + s = J.Scheduler() + assert (s._job_vram_requirement("comfyui", {}) + == T.get_tenant("comfyui").needs_vram_gb) + + +class TestSchedulerStatus: + def test_stats_report_depth_and_the_oldest_wait(self, queue): + J.submit("ollama", {}) + J.submit("comfyui", {}) + st = J.stats() + assert st["queue_depth"] == 2 + assert st["pending_by_tenant"] == {"ollama": 1, "comfyui": 1} + assert st["oldest_pending_s"] is not None + + def test_status_includes_queue_stats(self, queue): + J.submit("ollama", {}) + s = J.Scheduler().get_status() + assert s["queue_depth"] == 1 and s["running"] is False + + def test_a_blocked_job_reports_why_it_is_waiting(self, queue): + # Silence here is what left three jobs pending indefinitely with no explanation. + s = J.Scheduler() + s.blocked = {"id": "x", "tenant": "ollama", "reason": "needs 14.93 GB", + "waited_s": 42.0, "since": time.time()} + assert "14.93" in s.get_status()["blocked"]["reason"] + + +class TestDispatchers: + def test_every_tenant_kind_that_can_run_work_has_a_dispatcher(self): + import tenants as T + assert T.KIND_LLM in J.DISPATCHERS + assert T.KIND_DIFFUSION in J.DISPATCHERS + + def test_ollama_dispatch_reports_failure_rather_than_raising(self, queue, monkeypatch): + class _R: + status_code = 500 + text = '{"error":"llama-server process has terminated: cudaMalloc failed"}' + class _C: + async def __aenter__(self): return self + async def __aexit__(self, *a): return False + async def post(self, url, json=None): return _R() + monkeypatch.setattr(J.httpx, "AsyncClient", lambda **k: _C()) + res = asyncio.run(J._dispatch_ollama({"model": "m", "prompt": "hi"})) + assert res["ok"] is False and "cudaMalloc" in res["error"]