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"]