From 448e1fa50cb49a1fac9a32e59e0ff99fa475d092 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 20:53:08 -0700 Subject: [PATCH] =?UTF-8?q?refactor(hermes=5Fcli/web=5Frouters):=20local?= =?UTF-8?q?=5Fmodels=20=E2=80=94=20=5Fhttp=5Ferror=20ctx,=20=5Fstep/=5Fdow?= =?UTF-8?q?nload=5Fjob/=5Fcatalog=5Frow/=5Fquickstart=5Ftarget=20helpers,?= =?UTF-8?q?=20status=20phase=20helpers=20(1174->1142=20LOC)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- hermes_cli/web_routers/local_models.py | 436 ++++++++++++------------- 1 file changed, 202 insertions(+), 234 deletions(-) diff --git a/hermes_cli/web_routers/local_models.py b/hermes_cli/web_routers/local_models.py index b07cdd3993..c7b9f4416f 100644 --- a/hermes_cli/web_routers/local_models.py +++ b/hermes_cli/web_routers/local_models.py @@ -10,6 +10,7 @@ job pattern: start-POST -> {job_id} -> GET poll with byte progress. from __future__ import annotations import asyncio +import contextlib import json import logging import os @@ -53,12 +54,26 @@ _GIB = 1 << 30 _JOBS: Dict[str, Dict[str, Any]] = {} _JOBS_LOCK = threading.Lock() _LLAMACPP_PROVIDERS = ("llamacpp", "llama.cpp", "llama-cpp") +_SPLIT_PART_RE = r"-\d{5}-of-\d{5}" def _human_gb(n: int | float) -> str: return f"{n / _GIB:.1f} GB" +def _k_label(tokens: int) -> str: + return f"{tokens // 1024}K" + + +@contextlib.contextmanager +def _http_error(status: int, prefix: str = ""): + """Map any exception to ``HTTPException(status, f"{prefix}{exc}")``.""" + try: + yield + except Exception as exc: # noqa: BLE001 + raise HTTPException(status_code=status, detail=f"{prefix}{exc}") from exc + + def _job(kind: str, target: str, model_id: str | None = None) -> Dict[str, Any]: job = { "job_id": uuid.uuid4().hex[:12], "kind": kind, "target": target, @@ -80,12 +95,16 @@ def _job_view(job: Dict[str, Any]) -> Dict[str, Any]: return out -def _finish(job: Dict[str, Any], detail: str) -> None: - job["phase"] = "done" - job["status"] = "done" +def _step(job: Dict[str, Any], phase: str, detail: str) -> None: + job["phase"] = phase job["detail"] = detail +def _finish(job: Dict[str, Any], detail: str) -> None: + _step(job, "done", detail) + job["status"] = "done" + + def _spawn_job(job: Dict[str, Any], name: str, body: Callable[[], None], *, fail_msg: str | None = None, on_exit: Callable[[], None] | None = None) -> None: @@ -110,7 +129,6 @@ def _refresh_runtime(skip_msg: str) -> None: """Bounce a running router so it rescans the models dir (it only scans at spawn). Never raises — the file operation already succeeded.""" try: - bootstrap.refresh_local_runtime() except Exception: # noqa: BLE001 logger.debug(skip_msg, exc_info=True) @@ -164,13 +182,12 @@ def _probe_range_support(url: str) -> int: def _model_id_for(gguf: Path) -> str: """Variant model id for a staged file (strips split-part suffixes).""" - return re.sub(r"-\d{5}-of-\d{5}$", "", gguf.stem) + return re.sub(_SPLIT_PART_RE + "$", "", gguf.stem) def _variant_files_on_disk(model_id: str) -> "list[Path]": """Every local file belonging to a staged model: all split parts plus its catalog-declared assets (mmproj/draft) when present.""" - files = [p for p in _models_dir().glob("*.gguf") if _model_id_for(p) == model_id] hit = catalog.find_entry_for_model(model_id) if hit is not None: @@ -198,18 +215,15 @@ def download_file(url: str, dest: Path, job: Dict[str, Any], file_done = [0] progress_lock = threading.Lock() - def bump(n: int) -> None: - with progress_lock: - file_done[0] += n - job["done_bytes"] = base_done + file_done[0] - def pump(r, f) -> None: while True: chunk = r.read(_CHUNK) if not chunk: break f.write(chunk) - bump(len(chunk)) + with progress_lock: + file_done[0] += len(chunk) + job["done_bytes"] = base_done + file_done[0] try: # The probe and the preallocation both take real seconds on a 20+ GB @@ -270,7 +284,6 @@ def download_file(url: str, dest: Path, job: Dict[str, Any], def _models_dir() -> Path: - return bootstrap.models_dir() @@ -281,7 +294,6 @@ def _hf_url(repo: str, path: str) -> str: def _download_plan(entry, variant) -> list: """Everything a variant needs: split parts + mmproj/draft assets, as (url, dest, bytes) tuples.""" - plan = [(_hf_url(entry.repo, a.path), _models_dir() / a.local_name, a.size_bytes) for a in variant.files] plan += [(_hf_url(entry.repo, a.path), bootstrap.assets_dir() / a.local_name, a.size_bytes) @@ -293,8 +305,7 @@ def _run_download_plan(job: Dict[str, Any], plan: list, label: str) -> None: """Download every missing file in ``plan``; already-present files count toward progress without a transfer.""" total = sum(p[2] for p in plan) - job["phase"] = "downloading" - job["detail"] = f"{label} — {_human_gb(total)}" + _step(job, "downloading", f"{label} — {_human_gb(total)}") done_before = 0 for url, dest, size in plan: if not dest.exists(): @@ -311,7 +322,6 @@ def _engine_too_old(min_engine: str) -> bool: if not min_engine: return False try: - tags = binaries.installed_tags() or [binaries.default_tag()] newest = max(int(t.lstrip("b")) for t in tags if t.lstrip("b").isdigit()) return newest < int(min_engine.lstrip("b")) @@ -320,7 +330,6 @@ def _engine_too_old(min_engine: str) -> bool: def _load_config() -> dict: - try: return config_mod.load_config() except Exception: # noqa: BLE001 @@ -333,7 +342,6 @@ def _runtime_section() -> dict: def _set_runtime_enabled(enabled: bool) -> dict: """Persist ``local_runtime.enabled`` and return the config written.""" - config = config_mod.load_config() config.setdefault("local_runtime", {})["enabled"] = enabled config_mod.save_config(config) @@ -341,7 +349,6 @@ def _set_runtime_enabled(enabled: bool) -> dict: def _resolve_backend(section: dict, requested: str | None = None) -> str: - backend = requested or section.get("backend", "auto") return binaries.select_backend(bootstrap._detect_gpu_vendor()) if backend == "auto" else backend @@ -349,22 +356,32 @@ def _resolve_backend(section: dict, requested: str | None = None) -> str: def _eligible_entries(): """Catalog entries this engine can activate today (engine-gated ones can't be the recommendation either).""" - return tuple(e for e in catalog.CATALOG if not _engine_too_old(e.min_engine)) +def _resolve_assets_or_400(tag: str, backend: str): + """Resolve first so an impossible combination fails the POST, not the job.""" + with _http_error(400): + return binaries.resolve_assets(tag, backend) + + +def _start_local_server(config: dict, fail_detail: str): + """Start the local server (force) and return the supervisor; raise + ``fail_detail`` when neither we nor another process ended up serving.""" + sup = bootstrap.ensure_local_runtime(config, force=True) + if sup is None and _state_endpoint() is None: + raise RuntimeError(fail_detail) + return sup + + def _ensure_server(job: Dict[str, Any], config: dict, model_id: str, *, fail_detail: str, skip_msg: str) -> None: """Start the local server if needed and self-heal a stale router: the model list is spawn-only, so a server started before ``model_id`` finished downloading can't serve it — bounce it when it doesn't know the model.""" - - job["phase"] = "starting-server" - job["detail"] = "Starting the local server" - sup = bootstrap.ensure_local_runtime(config, force=True) - if sup is None and _state_endpoint() is None: - raise RuntimeError(fail_detail) + _step(job, "starting-server", "Starting the local server") + sup = _start_local_server(config, fail_detail) if sup is not None: try: if model_id not in sup.models(): @@ -377,9 +394,7 @@ def _ensure_server(job: Dict[str, Any], config: dict, model_id: str, *, def _assign_default(job: Dict[str, Any], model_id: str) -> None: """Make ``model_id`` the main model through the same machinery as /api/model/set (late-bound so tests can stub web_deps.late).""" - - job["phase"] = "setting-default" - job["detail"] = "Making it your default" + _step(job, "setting-default", "Making it your default") web_deps.late("_apply_model_assignment_sync")("main", "llamacpp", model_id, "", "", "") @@ -390,7 +405,6 @@ def _loaded_models(running: Dict[str, Any]) -> "tuple[Dict[str, str], Dict[str, """Resident models right now, plus how each is placed (granted window from the child, spill facts from the preset decision) — the difference between 'fast' and 'why is my CPU busy', so it must be inspectable.""" - data = _router_request(running, "/models", timeout=3) # Everything resident or becoming resident: 'loading' renders as its own # state in the pane (a 20-GB load in flight is the most important thing @@ -407,7 +421,7 @@ def _loaded_models(running: Dict[str, Any]) -> "tuple[Dict[str, str], Dict[str, plan = decisions.get(model_id) if plan is not None: facts["window"] = plan.window - facts["window_label"] = f"{plan.window // 1024}K" + facts["window_label"] = _k_label(plan.window) facts["spilled"] = plan.spilled if state in ("loaded", "ready"): try: @@ -415,7 +429,7 @@ def _loaded_models(running: Dict[str, Any]) -> "tuple[Dict[str, str], Dict[str, n_ctx = props.get("default_generation_settings", {}).get("n_ctx") if n_ctx: facts["granted_window"] = int(n_ctx) - facts["granted_window_label"] = f"{int(n_ctx) // 1024}K" + facts["granted_window_label"] = _k_label(int(n_ctx)) except Exception: # noqa: BLE001 pass if facts: @@ -423,12 +437,37 @@ def _loaded_models(running: Dict[str, Any]) -> "tuple[Dict[str, str], Dict[str, return loaded, placement +def _installed_backend(tag: str) -> str | None: + """Name of the first backend dir under ``tag`` with a working server binary.""" + root = binaries.runtimes_root() / tag + if not root.exists(): + return None + for backend_dir in sorted(p for p in root.iterdir() if p.is_dir()): + try: + binaries.server_binary(backend_dir) + return backend_dir.name + except Exception: # noqa: BLE001 + continue + return None + + +def _active_llamacpp_model_id() -> str | None: + """The active main model when it is one of ours (config authority: the + same model.provider + model.default that /api/model/set writes).""" + try: + model_section = (_load_config() or {}).get("model") or {} + if str(model_section.get("provider", "")).strip().lower() in _LLAMACPP_PROVIDERS: + return str(model_section.get("default") or model_section.get("name") or "").strip() or None + except Exception: # noqa: BLE001 + pass + return None + + @router.get("/api/local-models/status") def local_models_status(): """Cheap, immediate, never blocks on probes: config state + installed runtime + staged models + supervisor state (GPU facts live in /hardware). Sync def on purpose: blocking urlopen/scans run in the threadpool.""" - section = _runtime_section() configured_tag = section.get("tag") or binaries.default_tag() have = binaries.installed_tags() @@ -440,24 +479,12 @@ def local_models_status(): # Update pending = engine in use (enabled + something installed) and the # configured tag (pinned or release default) isn't on disk. The download # is a button click, never automatic. - update_available = bool( - section.get("enabled") and have and configured_tag not in have) - - runtime_backend = None - root = binaries.runtimes_root() / tag - if root.exists(): - for backend_dir in sorted(p for p in root.iterdir() if p.is_dir()): - try: - binaries.server_binary(backend_dir) - runtime_backend = backend_dir.name - break - except Exception: # noqa: BLE001 - continue + update_available = bool(section.get("enabled") and have and configured_tag not in have) + runtime_backend = _installed_backend(tag) staged = [] mdir = _models_dir() if mdir.exists(): - for gguf in bootstrap.staged_models(): model_id = _model_id_for(gguf) # Split models: report the whole variant's bytes, not one part's. @@ -479,18 +506,6 @@ def local_models_status(): logger.warning("loaded-models read failed: %r", exc) loaded = {} - # The active main model, when it is one of ours (config authority: the - # same model.provider + model.default that /api/model/set writes). - active_model_id = None - try: - model_section = (_load_config() or {}).get("model") or {} - if str(model_section.get("provider", "")).strip().lower() in _LLAMACPP_PROVIDERS: - active_model_id = str( - model_section.get("default") or model_section.get("name") or "" - ).strip() or None - except Exception: # noqa: BLE001 - pass - return { "enabled": bool(section.get("enabled")), "tag": tag, @@ -500,7 +515,7 @@ def local_models_status(): "runtime_backend": runtime_backend, "server_running": running is not None, "server_base_url": (running or {}).get("base_url"), - "active_model_id": active_model_id, + "active_model_id": _active_llamacpp_model_id(), "loaded_models": loaded, # Live load progress per model (SSE-fed): {model_id: {stage, value, # percent}}. The chat's loading bar and the picker rows poll this. @@ -513,7 +528,6 @@ def local_models_status(): def _loading_progress() -> Dict[str, Any]: try: - return load_progress.get_loading_progress() except Exception: # noqa: BLE001 — progress is garnish, never a 500 return {} @@ -526,7 +540,6 @@ def _loading_progress() -> Dict[str, Any]: def local_models_hardware(): """The budget as plain facts, polled by the pane and statusbar. Sync def on purpose: shells out to nvidia-smi — threadpool, not loop.""" - budget = hardware.probe_budget() ram_total, ram_avail = hardware._ram_bytes() out = { @@ -538,7 +551,6 @@ def local_models_hardware(): # GPU identity + live utilization (NVIDIA; other vendors degrade to None # and the UI hides those readouts). try: - smi_exe = hardware._nvidia_smi_path() smi = subprocess.run( [smi_exe, "--query-gpu=name,utilization.gpu,memory.used", @@ -567,6 +579,70 @@ _QUANT_REASON_COMPACT = ("Compact build sized for this machine ({quant}) — " "larger than GPU memory, runs slower") +def _catalog_row(entry, budget, recommended, recommended_reason, staged_ids) -> Dict[str, Any]: + choice = catalog.select_variant(entry, budget) + # Any variant of this family on disk counts as downloaded. + dl = next((v for v in entry.variants if v.model_id in staged_ids), None) + row: Dict[str, Any] = { + "id": entry.id, "display_name": entry.display_name, "description": entry.description, + "native_context": entry.n_ctx_train, + "native_context_label": _k_label(entry.n_ctx_train), + "recommended": entry.id == recommended, + "recommended_reason": recommended_reason if entry.id == recommended else None, + "downloaded": dl is not None, + "downloaded_model_id": dl.model_id if dl else None, + "downloaded_quant": dl.quant if dl else None, + "mtp": entry.mtp, "vision": entry.mmproj is not None, + # Day-0 architectures need the llama.cpp release where their support + # landed: True gates download/activate until the engine updates, but + # the row still renders (visible + explained beats hidden). + "needs_engine": _engine_too_old(entry.min_engine), + "min_engine": entry.min_engine or None, + } + if choice is None: + smallest = min(entry.variants, key=lambda v: v.size_bytes) + smallest_total = entry.download_bytes(smallest) + row.update({ + "fits": False, "size_bytes": smallest_total, "size_label": _human_gb(smallest_total), + "fit_summary": "Needs more memory than this machine has", + "fit_detail": (f"even the most compact build ({smallest.quant}, " + f"{_human_gb(smallest_total)}) exceeds GPU + system memory"), + }) + return row + + variant = choice.variant + # Same overhead the launch decision prices (runtime buffers + vision + # projector + microbatch/MTP logits): the row must advertise the window + # the model will actually get, not a paper number. + overhead = (context_policy.RUNTIME_OVERHEAD_BYTES + + (entry.mmproj.size_bytes if entry.mmproj else 0) + + context_policy.ub_logits_bytes(entry.n_vocab, mtp_capable=entry.mtp)) + decision = context_policy.initial_window(entry.profile(variant), budget, overhead_bytes=overhead) + download_total = entry.download_bytes(variant) + row.update({ + "fits": True, "model_id": variant.model_id, "quant": variant.quant, + "quant_validated": variant.validated, "size_bytes": download_total, + "size_label": _human_gb(download_total), "variant_count": len(entry.variants), + "quant_reason": _QUANT_REASONS.get( + choice.reason_key, _QUANT_REASON_COMPACT).format(quant=variant.quant), + }) + if isinstance(decision, estimator.PhysicsRefusal): + row["fit_summary"] = row["quant_reason"] + return row + row["start_window"] = decision.window + row["start_window_label"] = _k_label(decision.window) + row["spilled"] = decision.spilled + if decision.window >= entry.n_ctx_train: + shape = f"runs at its full {row['native_context_label']} context" + else: + shape = (f"starts at {row['start_window_label']} and grows toward " + f"{row['native_context_label']} as you use it") + if decision.spilled: + shape += " (larger than your GPU memory — runs slower)" + row["fit_summary"] = shape + return row + + @router.get("/api/local-models/catalog") def local_models_catalog(): """Every entry answers up front: how big is the download, will it fit, @@ -574,7 +650,6 @@ def local_models_catalog(): build for this machine (highest quality fully on GPU at the 64K floor; else the smallest that works, spilled and priced). No entry is hidden; unaffordable models show WHY. Sync def: blocking I/O -> threadpool.""" - # Serve the in-memory catalog; a TTL-gated background fetch lands new # entries for the next call (day-0 models without an app release). catalog.refresh_catalog_soon() @@ -589,71 +664,8 @@ def local_models_catalog(): # Completeness-checked staging (split parts all present) — same answer the # picker and router see, so a mid-download model never reads as downloaded. staged_ids = set(bootstrap.staged_model_ids()) - entries = [] - for entry in catalog.CATALOG: - choice = catalog.select_variant(entry, budget) - # Any variant of this family on disk counts as downloaded. - dl = next((v for v in entry.variants if v.model_id in staged_ids), None) - row: Dict[str, Any] = { - "id": entry.id, "display_name": entry.display_name, "description": entry.description, - "native_context": entry.n_ctx_train, - "native_context_label": f"{entry.n_ctx_train // 1024}K", - "recommended": entry.id == recommended, - "recommended_reason": recommended_reason if entry.id == recommended else None, - "downloaded": dl is not None, - "downloaded_model_id": dl.model_id if dl else None, - "downloaded_quant": dl.quant if dl else None, - "mtp": entry.mtp, "vision": entry.mmproj is not None, - # Day-0 architectures need the llama.cpp release where their support - # landed: True gates download/activate until the engine updates, but - # the row still renders (visible + explained beats hidden). - "needs_engine": _engine_too_old(entry.min_engine), - "min_engine": entry.min_engine or None, - } - if choice is None: - smallest = min(entry.variants, key=lambda v: v.size_bytes) - smallest_total = entry.download_bytes(smallest) - row.update({ - "fits": False, "size_bytes": smallest_total, "size_label": _human_gb(smallest_total), - "fit_summary": "Needs more memory than this machine has", - "fit_detail": (f"even the most compact build ({smallest.quant}, " - f"{_human_gb(smallest_total)}) exceeds GPU + system memory"), - }) - entries.append(row) - continue - - variant = choice.variant - # Same overhead the launch decision prices (runtime buffers + vision - # projector + microbatch/MTP logits): the row must advertise the window - # the model will actually get, not a paper number. - overhead = (context_policy.RUNTIME_OVERHEAD_BYTES - + (entry.mmproj.size_bytes if entry.mmproj else 0) - + context_policy.ub_logits_bytes(entry.n_vocab, mtp_capable=entry.mtp)) - decision = context_policy.initial_window(entry.profile(variant), budget, overhead_bytes=overhead) - download_total = entry.download_bytes(variant) - row.update({ - "fits": True, "model_id": variant.model_id, "quant": variant.quant, - "quant_validated": variant.validated, "size_bytes": download_total, - "size_label": _human_gb(download_total), "variant_count": len(entry.variants), - "quant_reason": _QUANT_REASONS.get( - choice.reason_key, _QUANT_REASON_COMPACT).format(quant=variant.quant), - }) - if not isinstance(decision, estimator.PhysicsRefusal): - row["start_window"] = decision.window - row["start_window_label"] = f"{decision.window // 1024}K" - row["spilled"] = decision.spilled - if decision.window >= entry.n_ctx_train: - shape = f"runs at its full {row['native_context_label']} context" - else: - shape = (f"starts at {row['start_window_label']} and grows toward " - f"{row['native_context_label']} as you use it") - if decision.spilled: - shape += " (larger than your GPU memory — runs slower)" - row["fit_summary"] = shape - else: - row["fit_summary"] = row["quant_reason"] - entries.append(row) - return {"models": entries} + return {"models": [_catalog_row(e, budget, recommended, recommended_reason, staged_ids) + for e in catalog.CATALOG]} # ── runtime install (job) ──────────────────────────────────── @@ -693,35 +705,25 @@ def _runtime_progress_hook(job: Dict[str, Any]): job["done_bytes"] = plan_done job["total_bytes"] = plan_total or None elif stage == "extract": - job["phase"] = "unpacking-runtime" pct = f" — {min(100, round(done / total * 100))}%" if total else "" - job["detail"] = f"Unpacking the engine{suffix}{pct}" + _step(job, "unpacking-runtime", f"Unpacking the engine{suffix}{pct}") else: # verify - job["phase"] = "verifying-runtime" - job["detail"] = f"Verifying the engine{suffix}" + _step(job, "verifying-runtime", f"Verifying the engine{suffix}") return hook @router.post("/api/local-models/runtime/install") async def local_models_runtime_install(body: RuntimeInstallBody): - section = _runtime_section() tag = section.get("tag") or binaries.default_tag() backend = _resolve_backend(section, body.backend) - # Resolve first so an impossible combination fails the POST, not the job. - try: - plan = binaries.resolve_assets(tag, backend) - except Exception as exc: # noqa: BLE001 - raise HTTPException(status_code=400, detail=str(exc)) - + plan = _resolve_assets_or_400(tag, backend) job = _job("runtime-install", f"llama.cpp {tag} ({backend})") def _run(): - previous = binaries.installed_tags() - job["phase"] = "downloading" - job["detail"] = f"Fetching {len(plan.assets)} package(s) for {backend}" + _step(job, "downloading", f"Fetching {len(plan.assets)} package(s) for {backend}") binaries.ensure_runtime_installed(tag, backend, progress=_runtime_progress_hook(job)) # Engine update path: a server already running on an older tag moves @@ -729,10 +731,8 @@ async def local_models_runtime_install(body: RuntimeInstallBody): # server) skip this; Use/boot handles their start. restarted = False try: - if bootstrap.get_supervisor() is not None and previous and tag not in previous: - job["phase"] = "restarting" - job["detail"] = "Switching the running server to the new build" + _step(job, "restarting", "Switching the running server to the new build") bootstrap.shutdown_local_runtime() bootstrap.ensure_local_runtime(_load_config(), force=True) restarted = True @@ -761,11 +761,22 @@ class ModelDownloadBody(BaseModel): model_id: str +def _download_job(job: Dict[str, Any], body: Callable[[], None], name: str, label: str, + fail_msg: str | None = None) -> None: + """Spawn a download job: ``body`` fetches, then the job finishes as + "