refactor(web): second-pass compaction of local_models/profiles/web_git (helpers, dict literals, docstrings)

This commit is contained in:
Teknium
2026-09-02 11:19:36 -07:00
parent 15aaa359bf
commit 7c0d8f2c6f
3 changed files with 216 additions and 429 deletions

View File

@@ -1,13 +1,11 @@
"""Backend git operations for the desktop coding rail + Codex-style review pane.
"""Backend git operations for the desktop coding rail + review pane.
The desktop's git affordances (coding-rail status, worktree lanes, review pane,
branch switch) run as Electron-local git on the user's machine. On a *remote*
gateway those would operate on the wrong filesystem, so this module mirrors them
over the dashboard's authenticated REST surface — the same pattern as ``/api/fs``.
Everything shells out to the system ``git`` (and ``gh`` for ship info / PRs).
Reads degrade to ``None`` / empty on a non-repo; mutations raise so the renderer
can surface a toast. Callers pass an already path-hardened ``cwd``.
The desktop's git affordances run as Electron-local git; on a *remote* gateway
those would hit the wrong filesystem, so this module mirrors them over the
dashboard's authenticated REST surface (same pattern as ``/api/fs``). Everything
shells out to system ``git`` (and ``gh`` for PRs). Reads degrade to ``None`` /
empty on a non-repo; mutations raise so the renderer can toast. Callers pass an
already path-hardened ``cwd``.
"""
from __future__ import annotations
@@ -144,14 +142,9 @@ def _default_branch_name(cwd: str) -> str | None:
head = _git_out(cwd, ["rev-parse", "--abbrev-ref", "origin/HEAD"]).strip()
if head and head != "origin/HEAD":
return head.split("/", 1)[-1]
for ref in (
"refs/heads/main",
"refs/heads/master",
"refs/remotes/origin/main",
"refs/remotes/origin/master",
):
code, _, _ = _git(cwd, ["rev-parse", "--verify", "--quiet", ref])
if code == 0:
for ref in ("refs/heads/main", "refs/heads/master",
"refs/remotes/origin/main", "refs/remotes/origin/master"):
if _git(cwd, ["rev-parse", "--verify", "--quiet", ref])[0] == 0:
return ref.split("/")[-1]
return None
@@ -161,8 +154,7 @@ def _default_branch_name(cwd: str) -> str | None:
def _walk_entries(raw: str):
"""Yield (tag, xy, path) per changed file from ``git status --porcelain=v2 -z``,
skipping branch headers and the rename/copy origin-path records. One walker
feeds the rail, the review list, and the commit flow."""
skipping branch headers and rename/copy origin-path records."""
records = raw.split("\0")
i = 0
while i < len(records):
@@ -188,13 +180,9 @@ def _entry_staged(tag: str, xy: str) -> bool:
def _classify(tag: str, xy: str, path: str) -> dict:
y = xy[1] if len(xy) > 1 else "."
return {
"path": path,
"staged": _entry_staged(tag, xy),
"unstaged": tag == "?" or (tag in ("1", "2") and y not in (".", "?")),
"untracked": tag == "?",
"conflicted": tag == "u",
}
return {"path": path, "staged": _entry_staged(tag, xy),
"unstaged": tag == "?" or (tag in ("1", "2") and y not in (".", "?")),
"untracked": tag == "?", "conflicted": tag == "u"}
def _status_letter(tag: str, xy: str) -> str:
@@ -241,19 +229,11 @@ def repo_status(cwd: str) -> dict | None:
added += sum(_untracked_insertions(cwd, f["path"]) for f in files[:_UNTRACKED_SCAN_CAP] if f["untracked"])
return {
"branch": branch,
"defaultBranch": _default_branch_name(cwd),
"detached": detached,
"ahead": ahead,
"behind": behind,
"staged": sum(f["staged"] for f in files),
"unstaged": sum(f["unstaged"] for f in files),
"untracked": sum(f["untracked"] for f in files),
"conflicted": sum(f["conflicted"] for f in files),
"changed": len(files),
"added": added,
"removed": removed,
"files": files[:200],
"branch": branch, "defaultBranch": _default_branch_name(cwd), "detached": detached,
"ahead": ahead, "behind": behind,
"staged": sum(f["staged"] for f in files), "unstaged": sum(f["unstaged"] for f in files),
"untracked": sum(f["untracked"] for f in files), "conflicted": sum(f["conflicted"] for f in files),
"changed": len(files), "added": added, "removed": removed, "files": files[:200],
}
@@ -465,22 +445,14 @@ def _pr_query(owner: str, name: str, branches: list[str], numbers: list[int]) ->
# directly also tells us its branch — so it lands in the same by-branch map
# as everything else.
fields += [f"n{i}: pullRequest(number: {n}) {{ {_PR_NODE_FIELDS} }}" for i, n in enumerate(numbers)]
return (
f"query {{ repository(owner: {json.dumps(owner)}, name: {json.dumps(name)}) {{\n"
+ "\n".join(fields)
+ "\n} }"
)
return (f"query {{ repository(owner: {json.dumps(owner)}, name: {json.dumps(name)}) {{\n"
+ "\n".join(fields) + "\n} }")
def _pr_payload(pr: dict) -> dict:
return {
"branch": str(pr.get("headRefName")),
"draft": bool(pr.get("isDraft")),
"number": int(pr.get("number") or 0),
"state": str(pr.get("state") or "").lower(),
"title": str(pr.get("title") or ""),
"url": str(pr.get("url") or ""),
}
return {"branch": str(pr.get("headRefName")), "draft": bool(pr.get("isDraft")),
"number": int(pr.get("number") or 0), "state": str(pr.get("state") or "").lower(),
"title": str(pr.get("title") or ""), "url": str(pr.get("url") or "")}
def review_pr_list(cwd: str, branches: list[str], numbers: list[int] = None) -> dict:
@@ -594,9 +566,8 @@ def _slugify(name: str) -> str:
def _default_branch(cwd: str) -> str:
remote = _git_out(
cwd, ["symbolic-ref", "--quiet", "--short", "refs/remotes/origin/HEAD"]
).strip().replace("origin/", "", 1)
remote = _git_out(cwd, ["symbolic-ref", "--quiet", "--short", "refs/remotes/origin/HEAD"]
).strip().replace("origin/", "", 1)
if remote:
return remote
configured = _git_out(cwd, ["config", "--get", "init.defaultBranch"]).strip()
@@ -733,27 +704,11 @@ def branch_list(cwd: str) -> list[dict]:
and name.split("/", 1)[-1] not in local_set
]
return [
*(
{
"name": name,
"checkedOut": name in path_by_branch,
"isDefault": bool(trunk and name == trunk),
"isRemote": False,
"worktreePath": path_by_branch.get(name),
}
for name in locals_
),
*(
{
# No local checkout, and never the local trunk.
"name": name,
"checkedOut": False,
"isDefault": False,
"isRemote": True,
"worktreePath": None,
}
for name in remotes
),
*({"name": name, "checkedOut": name in path_by_branch, "isDefault": bool(trunk and name == trunk),
"isRemote": False, "worktreePath": path_by_branch.get(name)} for name in locals_),
# Remote rows: no local checkout, and never the local trunk.
*({"name": name, "checkedOut": False, "isDefault": False, "isRemote": True,
"worktreePath": None} for name in remotes),
]
@@ -776,9 +731,8 @@ def base_branch_list(cwd: str) -> list[dict]:
return []
# origin/HEAD when a remote exists; otherwise the local default
# (main/master/init.defaultBranch) so a no-remote repo still flags its trunk.
default = _git_out(
cwd, ["symbolic-ref", "--quiet", "--short", "refs/remotes/origin/HEAD"]
).strip() or _default_branch(cwd)
default = (_git_out(cwd, ["symbolic-ref", "--quiet", "--short", "refs/remotes/origin/HEAD"]).strip()
or _default_branch(cwd))
return [
{"name": name, "isRemote": name.startswith("origin/"),
"isDefault": bool(default and name == default)}

View File

@@ -1,14 +1,10 @@
"""Local-models dashboard routes — the desktop's window into the managed
llama.cpp runtime.
Every payload carries plain-language, pre-formatted facts the UI can show
Every payload carries plain-language, pre-formatted facts the UI shows
verbatim (what will this model do ON THIS MACHINE, how big is the download,
what is the runtime doing right now), never raw internals.
Long jobs (runtime install, model download) follow the repo's job pattern:
start-POST -> {job_id} -> GET poll with byte progress. Downloads are
byte-size checked against what the server declared (no hash verification by
design); a short download deletes the file and reports it plainly.
what is the runtime doing), never raw internals. Long jobs follow the repo's
job pattern: start-POST -> {job_id} -> GET poll with byte progress.
"""
from __future__ import annotations
@@ -65,17 +61,12 @@ def _human_gb(n: int | float) -> str:
def _job(kind: str, target: str, model_id: str | None = None) -> Dict[str, Any]:
job = {
"job_id": uuid.uuid4().hex[:12],
"kind": kind, # "runtime-install" | "model-download" | ...
"target": target,
"job_id": uuid.uuid4().hex[:12], "kind": kind, "target": target,
"model_id": model_id, # catalog id for downloads; None otherwise
"status": "running", # running | done | error
"phase": "starting", # human-readable step name
"detail": "",
"total_bytes": None,
"done_bytes": 0,
"started_at": time.time(),
"error": None,
"detail": "", "total_bytes": None, "done_bytes": 0,
"started_at": time.time(), "error": None,
}
with _JOBS_LOCK:
_JOBS[job["job_id"]] = job
@@ -98,9 +89,8 @@ def _finish(job: Dict[str, Any], detail: str) -> None:
def _spawn_job(job: Dict[str, Any], name: str, body: Callable[[], None], *,
fail_msg: str | None = None,
on_exit: Callable[[], None] | None = None) -> None:
"""Run ``body`` on a daemon thread; any exception marks the job errored
(optionally warning ``fail_msg`` with the exception). ``on_exit`` always
runs last (lock release)."""
"""Run ``body`` on a daemon thread; an exception marks the job errored
(warning ``fail_msg`` when given). ``on_exit`` always runs last."""
def _run():
try:
body()
@@ -117,9 +107,8 @@ def _spawn_job(job: Dict[str, Any], name: str, body: Callable[[], None], *,
def _refresh_runtime(skip_msg: str) -> None:
"""Bounce a running router so it rescans the models dir; the router only
scans at spawn, so a new/deleted file is invisible until then. Never
raises — the file operation already succeeded."""
"""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()
@@ -152,10 +141,9 @@ _CHUNK = 4 << 20
def _probe_range_support(url: str) -> int:
"""Total size when the server honors Range requests, else 0.
A 401/403 from the CDN means the repo is gated or the catalog names a
wrong repo — raise with a plain-language message, not a bare status."""
"""Total size when the server honors Range requests, else 0. A 401/403
means the repo is gated or the catalog names a wrong repo — raise with a
plain-language message, not a bare status."""
req = urllib.request.Request(url, headers={"Range": "bytes=0-0"})
try:
with urllib.request.urlopen(req, timeout=60) as r:
@@ -195,19 +183,15 @@ def _variant_files_on_disk(model_id: str) -> "list[Path]":
def download_file(url: str, dest: Path, job: Dict[str, Any],
*,
base_done: int = 0, keep_totals: bool = False) -> None:
"""Download url -> dest with byte progress on ``job``.
"""Download url -> dest with byte progress on ``job``; ranged-parallel when
the server supports it, single-stream otherwise. Never leaves a .part.
Ranged-parallel when the server supports it, single-stream fallback
otherwise. No integrity check against the CATALOG by design: catalog
sizes may lag an upstream re-upload, and a newer file must download
fine. Completeness is checked only against what the SERVER declared for
this transfer (range-probe total / Content-Length), so a dropped
connection still errors instead of staging a truncated file. Never
leaves a .part behind.
Multi-file variants: ``base_done`` offsets the progress so this file's
bytes accumulate onto the files before it, and ``keep_totals=True``
stops the per-file size from overwriting the variant's total.
No integrity check against the CATALOG by design (catalog sizes may lag a
re-upload); completeness is checked only against what the SERVER declared
(range-probe total / Content-Length), so a dropped connection still errors
instead of staging a truncated file. Multi-file variants: ``base_done``
offsets progress onto earlier files and ``keep_totals=True`` stops the
per-file size from overwriting the variant's total.
"""
tmp = dest.with_suffix(".part")
dest.parent.mkdir(parents=True, exist_ok=True)
@@ -241,9 +225,7 @@ def download_file(url: str, dest: Path, job: Dict[str, Any],
f.truncate(total)
job["detail"] = ""
errors: list[Exception] = []
bounds = [(i * total // _DOWNLOAD_CONNECTIONS,
(i + 1) * total // _DOWNLOAD_CONNECTIONS - 1)
for i in range(_DOWNLOAD_CONNECTIONS)]
n = _DOWNLOAD_CONNECTIONS
def fetch_range(start: int, end: int) -> None:
try:
@@ -256,9 +238,9 @@ def download_file(url: str, dest: Path, job: Dict[str, Any],
except Exception as exc: # noqa: BLE001
errors.append(exc)
threads = [threading.Thread(target=fetch_range, args=b, daemon=True,
name=f"lm-dl-{i}")
for i, b in enumerate(bounds)]
threads = [threading.Thread(target=fetch_range, daemon=True, name=f"lm-dl-{i}",
args=(i * total // n, (i + 1) * total // n - 1))
for i in range(n)]
for t in threads:
t.start()
for t in threads:
@@ -405,11 +387,9 @@ def _assign_default(job: Dict[str, Any], model_id: str) -> None:
def _loaded_models(running: Dict[str, Any]) -> "tuple[Dict[str, str], Dict[str, Any]]":
"""Which staged models are resident right now, plus how each is placed.
Placement is the granted window from the child itself and the plan's
spill facts from the preset decision — the difference between 'fast'
and 'why is my CPU busy', so it must be inspectable, not inferred."""
"""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
@@ -446,11 +426,8 @@ def _loaded_models(running: Dict[str, Any]) -> "tuple[Dict[str, str], Dict[str,
@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 come from
/api/local-models/hardware (slower, polled).
Sync def on purpose: blocking urlopen/scans run in FastAPI's threadpool
instead of stalling the event loop."""
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()
@@ -460,21 +437,18 @@ def local_models_status():
# newest installed).
tag = configured_tag if configured_tag in have else (have[0] if have else configured_tag)
# A pending engine update exists when the user runs the local engine
# (enabled + something installed) and the configured tag — pinned or the
# Hermes-release default — is newer than anything on disk. The download
# 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_installed = False
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_installed = True
runtime_backend = backend_dir.name
break
except Exception: # noqa: BLE001
@@ -493,8 +467,7 @@ def local_models_status():
running = _state_endpoint()
# Resident models from the live router; {} when down. Feeds the pane's
# Loaded pills and eject buttons.
# Resident models from the live router ({} when down): Loaded pills + eject.
loaded: Dict[str, str] = {}
placement: Dict[str, Any] = {}
if running is not None:
@@ -523,7 +496,7 @@ def local_models_status():
"tag": tag,
"configured_tag": configured_tag,
"update_available": update_available,
"runtime_installed": runtime_installed,
"runtime_installed": runtime_backend is not None,
"runtime_backend": runtime_backend,
"server_running": running is not None,
"server_base_url": (running or {}).get("base_url"),
@@ -551,22 +524,16 @@ def _loading_progress() -> Dict[str, Any]:
@router.get("/api/local-models/hardware")
def local_models_hardware():
"""The budget as plain facts. Polled by the pane and the statusbar
resource item (throttled client-side). Sync def on purpose: shells out
to nvidia-smi and probes budgets — threadpool, not loop."""
"""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 = {
"uma": budget.uma,
"vram_total_bytes": budget.total_device_bytes,
"vram_usable_bytes": budget.usable_vram_bytes,
"ram_total_bytes": ram_total,
"ram_available_bytes": ram_avail,
"vram_label": _human_gb(budget.total_device_bytes),
"gpu_name": None,
"gpu_util_percent": None,
"vram_used_bytes": None,
"uma": budget.uma, "vram_total_bytes": budget.total_device_bytes,
"vram_usable_bytes": budget.usable_vram_bytes, "ram_total_bytes": ram_total,
"ram_available_bytes": ram_avail, "vram_label": _human_gb(budget.total_device_bytes),
"gpu_name": None, "gpu_util_percent": None, "vram_used_bytes": None,
}
# GPU identity + live utilization (NVIDIA; other vendors degrade to None
# and the UI hides those readouts).
@@ -602,22 +569,18 @@ _QUANT_REASON_COMPACT = ("Compact build sized for this machine ({quant}) — "
@router.get("/api/local-models/catalog")
def local_models_catalog():
"""Every entry answers the user's three questions up front: how big is
the download, will it fit, and what context/speed shape will I get —
from the catalog's measured numbers + this machine's budget. The row
advertises the BEST build for this machine (highest quality that runs
fully on the GPU at the 64K floor; else the smallest that works, spilled
and priced). No entry is hidden; unaffordable models show WHY. Sync def
on purpose: probe_budget + catalog I/O block — threadpool, not loop."""
"""Every entry answers up front: how big is the download, will it fit,
and what context/speed shape will I get. The row advertises the BEST
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 catalog already in memory; a TTL-gated background fetch lands
# new entries for the next call (day-0 models without an app release).
# 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()
# Planning budget: price against machine capacity, not live-free VRAM —
# a loaded model must not make every row unaffordable.
# Planning budget: machine capacity, not live-free VRAM — a loaded model
# must not make every row unaffordable.
budget = hardware.probe_budget(planning=True)
# The default pick for THIS machine: quality-ranked, fit- and speed-gated.
# The reason key ships with the row so the Recommended badge's tooltip is
# the branch that actually fired, not a re-derivation that can drift.
picked = catalog.recommended_entry(budget, _eligible_entries())
@@ -630,24 +593,20 @@ def local_models_catalog():
for entry in catalog.CATALOG:
choice = catalog.select_variant(entry, budget)
# Any variant of this family on disk counts as downloaded.
downloaded_variant = next(
(v for v in entry.variants if v.model_id in staged_ids), None)
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,
"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": downloaded_variant is not None,
"downloaded_model_id": downloaded_variant.model_id if downloaded_variant else None,
"downloaded_quant": downloaded_variant.quant if downloaded_variant else None,
"mtp": entry.mtp,
"vision": entry.mmproj is not 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 in the pane until the engine
# updates; the row still renders (visible + explained beats hidden).
# 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,
}
@@ -655,9 +614,7 @@ def local_models_catalog():
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),
"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"),
@@ -675,13 +632,9 @@ def local_models_catalog():
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),
"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),
})
@@ -711,13 +664,11 @@ class RuntimeInstallBody(BaseModel):
def _runtime_progress_hook(job: Dict[str, Any]):
"""Adapter: ensure_runtime_installed's progress stream -> job fields.
Throttled to ~4 updates/s. Byte counters are CUMULATIVE across the plan:
a multi-asset engine (CUDA zip + cudart zip) reads as one growing
download, not a bar that restarts per asset; the total grows as each
asset's size becomes known. Unpack/verify keep the download's counters —
a bar bouncing back to zero after the bytes finished reads as failure."""
"""Adapter: ensure_runtime_installed's progress stream -> job fields,
throttled to ~4 updates/s. Byte counters are CUMULATIVE across the plan
(a multi-asset engine reads as one growing download, total growing as
each asset's size becomes known); unpack/verify keep the counters — a bar
bouncing back to zero after the bytes finished reads as failure."""
state = {"last": 0.0, "banked": 0, "asset": None, "asset_total": 0}
def hook(stage: str, done: int, total: int, label: str) -> None:
@@ -856,9 +807,9 @@ async def local_models_download(body: ModelDownloadBody):
@router.delete("/api/local-models/models/{model_id}")
async def local_models_delete(model_id: str):
"""Remove a staged model: every split part plus its private assets, then
bounce the router off the request thread (deleting the active file
mid-serve is exactly the stale state the refresh exists for)."""
"""Remove every split part plus private assets, then bounce the router off
the request thread (deleting the active file mid-serve is exactly the
stale state the refresh exists for)."""
files = _variant_files_on_disk(model_id)
if not files:
raise HTTPException(status_code=404, detail="model not found")
@@ -899,16 +850,12 @@ _QUICKSTART_LOCK = threading.Lock()
@router.post("/api/local-models/quickstart")
async def local_models_quickstart(body: QuickstartBody):
"""The dummy-proof path: one job that installs the runtime (if missing),
downloads this machine's build of the recommended model (if missing),
and makes it the default for new chats. Each leg is the same code the
individual routes run — this route only sequences them, so 'Configure'
and quickstart can never disagree about what gets installed.
Preflight rejects (no servable entry, engine too old) fail the POST
synchronously so the button can explain itself; everything slow runs in
the job with the usual phase/byte progress.
"""
"""One job: install the runtime (if missing), download this machine's
build of the recommended model (if missing), make it the default. Each
leg is the same code the individual routes run, so 'Configure' and
quickstart can never disagree. Preflight rejects (no servable entry,
engine too old) fail the POST synchronously so the button can explain
itself; everything slow runs in the job with phase/byte progress."""
# Resolve the target entry: explicit id, else this machine's
# recommendation, else the first catalog entry this machine can serve.
@@ -983,14 +930,9 @@ async def local_models_quickstart(body: QuickstartBody):
_spawn_job(job, "lr-quickstart", _run, fail_msg="quickstart failed: %s",
on_exit=_QUICKSTART_LOCK.release)
return {
"job_id": job["job_id"],
"model_id": entry.id,
"display_name": entry.display_name,
"needs_runtime": need_runtime,
"needs_download": need_download,
"download_bytes": download_bytes,
}
return {"job_id": job["job_id"], "model_id": entry.id, "display_name": entry.display_name,
"needs_runtime": need_runtime, "needs_download": need_download,
"download_bytes": download_bytes}
def _stop_server() -> None:
@@ -1027,10 +969,9 @@ _SERVER_ACTIONS = {"stop": _stop_server, "start": _start_server}
@router.post("/api/local-models/server")
async def local_models_server(body: ServerActionBody):
"""Turn the local engine off (stop the server, free ALL GPU memory, and
disable auto-start) or back on. The off switch is the whole-engine
counterpart of per-model eject — and unlike eject it IS durable: the
user said off, so boots stay off until they say on."""
"""Turn the local engine off (stop the server, free ALL GPU memory,
disable auto-start) or back on. Unlike per-model eject the off switch IS
durable: the user said off, so boots stay off until they say on."""
action = (body.action or "").strip().lower()
if action not in _SERVER_ACTIONS:
raise HTTPException(status_code=400, detail="action must be 'stop' or 'start'")
@@ -1050,10 +991,9 @@ class ModelEjectBody(BaseModel):
@router.post("/api/local-models/eject")
def local_models_eject(body: ModelEjectBody):
"""Free a loaded model's GPU memory now. Nothing reloads it except
demand — the next message to it (residency v2: no automatic loading
exists anywhere). Sync def on purpose: the fallback path blocks on a
urlopen with a 120s timeout — threadpool, never the event loop."""
"""Free a loaded model's GPU memory now; only demand (the next message)
reloads it — residency v2 has no automatic loading anywhere. Sync def:
the fallback path blocks on a 120s urlopen — threadpool, never the loop."""
sup = bootstrap.get_supervisor()
if sup is not None:
@@ -1082,13 +1022,11 @@ class ModelActivateBody(BaseModel):
@router.post("/api/local-models/activate")
async def local_models_activate(body: ModelActivateBody):
"""Make a downloaded model the default for new chats. Pure selection
(residency v2): a config write through the same machinery as
/api/model/set, plus making sure the server is up. NO model loading —
models load on first inference, always; an empty router costs nothing.
Kept as a job for UI continuity."""
# Split variants stage under their first part — resolve like the rest
# of the routes instead of assuming a single flat file.
"""Make a downloaded model the default for new chats: a config write via
the same machinery as /api/model/set plus making sure the server is up.
NO model loading (residency v2: models load on first inference; an empty
router costs nothing). Kept as a job for UI continuity."""
# Split variants stage under their first part — resolve like the other routes.
if body.model_id not in bootstrap.staged_model_ids():
raise HTTPException(status_code=404, detail=f"{body.model_id} is not downloaded")
@@ -1138,9 +1076,8 @@ async def local_models_job(job_id: str):
@router.get("/api/local-models/search")
async def local_models_search(q: str, limit: int = 20):
"""Full-text HF search over GGUF models — the firehose behind the
curated catalog. Per-quant fit pills come from the repo-files call once
the user opens a hit."""
"""Full-text HF search over GGUF models — the firehose behind the curated
catalog; per-quant fit pills come from the repo-files call."""
if not q.strip():
@@ -1155,9 +1092,8 @@ async def local_models_search(q: str, limit: int = 20):
@router.get("/api/local-models/search/files")
async def local_models_search_files(repo: str):
"""The servable GGUFs in one HF repo with a rough pre-download fit
verdict per quant (file size + conservative fill-ins — the GGUF header
refines it after download)."""
"""Servable GGUFs in one HF repo with a rough pre-download fit verdict per
quant (file size + conservative fill-ins; the GGUF header refines it)."""
try:
@@ -1176,10 +1112,9 @@ class BrowsedDownloadBody(BaseModel):
@router.post("/api/local-models/download-browsed")
async def local_models_download_browsed(body: BrowsedDownloadBody):
"""Download an arbitrary HF GGUF (browsed or pasted) into the managed
models dir. From the moment it lands it is a normal staged model: the
post-download bounce regenerates presets from its real header and the
fit policy owns its launch. No catalog entry — it serves 'unverified',
"""Download an arbitrary HF GGUF into the managed models dir. Once landed
it is a normal staged model (the post-download bounce regenerates presets
from its real header); with no catalog entry it serves 'unverified',
capabilities answered from the live server only."""
paths = [p for p in (body.paths or []) if p.lower().endswith(".gguf")]
@@ -1217,10 +1152,9 @@ class SideloadBody(BaseModel):
@router.post("/api/local-models/sideload")
async def local_models_sideload(body: SideloadBody):
"""Register a GGUF that already exists on this machine: link it into
the managed models dir (copy only when linking is impossible) and
bounce the router so it serves immediately. The original stays where
it is; delete-from-Hermes removes only our link."""
"""Register a GGUF already on this machine: link it into the managed
models dir (copy only when linking is impossible) and bounce the router.
The original stays put; delete-from-Hermes removes only our link."""
src = Path(body.path)
if not src.is_file() or src.suffix.lower() != ".gguf":

View File

@@ -58,10 +58,8 @@ def _warn_profile_read_error(profile: str, exc: Exception) -> None:
if key in _profile_read_warned:
return
_profile_read_warned.add(key)
_log.warning(
"profile session read failed for %r (reported only in the response "
"errors array): %s", profile, exc,
)
_log.warning("profile session read failed for %r (reported only in the response "
"errors array): %s", profile, exc)
sessions_router = APIRouter()
router = APIRouter()
@@ -203,6 +201,16 @@ def _profile_errors(log_msg: str, *args, not_found=(FileNotFoundError,),
raise HTTPException(status_code=500, detail=str(e))
def _best_effort(log_msg: str, *args, fn, default=None):
"""Run ``fn()``; on failure log ``log_msg`` with traceback and return
``default`` (the parent operation already succeeded, so never 500)."""
try:
return fn()
except Exception:
_log.exception(log_msg, *args)
return default
def _profile_targets(log_label: str, *, lightweight: bool) -> List[Tuple[str, Path]]:
"""(name, home) for every profile, falling back to ``default`` alone.
@@ -441,15 +449,17 @@ def get_profiles_sessions(
else:
targets = _profile_targets("GET /api/profiles/sessions", lightweight=True)
min_message_count = max(0, min_messages)
archived_only = archived == "only"
include_archived = archived == "include"
# Source scoping (see /api/sessions): recents pass exclude_sources=cron,
# the cron-jobs section passes source=cron — two independent lists so
# newest cron sessions can't starve the recents page.
source_filter = source or None
source_list = [s.strip() for s in (sources or "").split(",") if s.strip()]
exclude_list = [s.strip() for s in (exclude_sources or "").split(",") if s.strip()]
filters = dict(
source=source or None,
sources=[s.strip() for s in (sources or "").split(",") if s.strip()] or None,
exclude_sources=[s.strip() for s in (exclude_sources or "").split(",") if s.strip()] or None,
min_message_count=max(0, min_messages),
include_archived=archived == "include",
archived_only=archived == "only",
)
# Over-fetch per profile so the merged+sorted window is correct for the
# requested page. Capped so a huge profile can't blow up the response.
per_profile = min(max(limit + offset, limit), 500)
@@ -465,28 +475,11 @@ def get_profiles_sessions(
continue
try:
rows = db.list_sessions_rich(
source=source_filter,
sources=source_list or None,
exclude_sources=exclude_list or None,
limit=per_profile,
offset=0,
min_message_count=min_message_count,
include_archived=include_archived,
archived_only=archived_only,
order_by_last_active=order == "recent",
limit=per_profile, offset=0, order_by_last_active=order == "recent",
# Same SQL-level blob skip as /api/sessions.
compact_rows=not full,
include_pinned=True,
)
profile_total = db.session_count(
source=source_filter,
sources=source_list or None,
exclude_sources=exclude_list or None,
min_message_count=min_message_count,
include_archived=include_archived,
archived_only=archived_only,
exclude_children=True,
compact_rows=not full, include_pinned=True, **filters,
)
profile_total = db.session_count(exclude_children=True, **filters)
total += profile_total
profile_totals[name] = profile_total
merged.extend(_tag_rows(rows, name, now))
@@ -501,14 +494,8 @@ def get_profiles_sessions(
window = _pinned_window(merged, offset, limit)
if not full:
_strip_session_list_rows(window)
return {
"sessions": window,
"total": total,
"profile_totals": profile_totals,
"limit": limit,
"offset": offset,
"errors": errors,
}
return {"sessions": window, "total": total, "profile_totals": profile_totals,
"limit": limit, "offset": offset, "errors": errors}
@sessions_router.get("/api/profiles/sessions/sidebar")
@@ -556,19 +543,12 @@ def get_profiles_sessions_sidebar(
now = time.time()
def _slice(db, *, source=None, exclude=None, cap):
# include_pinned: a pinned conversation must reach the sidebar even
# when it has aged past the window, or its Pinned row renders empty.
return db.list_sessions_rich(
source=source,
exclude_sources=exclude or None,
limit=cap,
offset=0,
min_message_count=1,
include_archived=False,
archived_only=False,
order_by_last_active=True,
compact_rows=True,
# A pinned conversation must reach the sidebar even when it has
# aged past the window — otherwise its Pinned row renders empty.
include_pinned=True,
source=source, exclude_sources=exclude or None, limit=cap, offset=0,
min_message_count=1, include_archived=False, archived_only=False,
order_by_last_active=True, compact_rows=True, include_pinned=True,
)
for name, home in targets:
@@ -577,15 +557,9 @@ def get_profiles_sessions_sidebar(
db_path = Path(home) / "state.db"
if not db_path.exists():
continue
profile_cache_key = (
str(db_path),
_sidebar_db_fingerprint(db_path),
recents_cap,
tuple(recents_exclude_list),
cron_cap,
messaging_cap,
tuple(messaging_exclude_list),
)
profile_cache_key = (str(db_path), _sidebar_db_fingerprint(db_path), recents_cap,
tuple(recents_exclude_list), cron_cap, messaging_cap,
tuple(messaging_exclude_list))
slices = _sidebar_profile_cache_get(profile_cache_key)
if slices is None:
db = _open_profile_db(name, home, errors)
@@ -627,16 +601,10 @@ def get_profiles_sessions_sidebar(
return win
return {
"recents": {
"sessions": _window(recents_rows, recents_cap),
"profiles_truncated": recents_truncated,
"profiles_usage": profile_totals,
},
"recents": {"sessions": _window(recents_rows, recents_cap),
"profiles_truncated": recents_truncated, "profiles_usage": profile_totals},
"cron": {"sessions": _window(cron_rows, cron_cap)},
"messaging": {
"sessions": _window(messaging_rows, messaging_cap),
"total": len(messaging_rows),
},
"messaging": {"sessions": _window(messaging_rows, messaging_cap), "total": len(messaging_rows)},
"errors": errors,
}
@@ -739,12 +707,8 @@ def get_profiles_projects_tree(preview_limit: int = 3, session_limit: int = 2000
token = set_hermes_home_override(str(home))
try:
tree, _active_id = gateway_server._build_project_tree(
db,
preview_limit=preview_limit,
hydrate=False,
session_limit=session_limit,
include_discovered=False,
)
db, preview_limit=preview_limit, hydrate=False,
session_limit=session_limit, include_discovered=False)
_merge_profile_tree(merged, tree["projects"], name, preview_limit)
scoped_session_ids.extend(tree["scoped_session_ids"])
except Exception as exc:
@@ -756,11 +720,9 @@ def get_profiles_projects_tree(preview_limit: int = 3, session_limit: int = 2000
return {
"projects": sorted(merged.values(), key=lambda p: p.get("lastActive") or 0, reverse=True),
# Ownership is per profile, so no single project is "the active one"
# here; the desktop only reads active_id to bias its overview sort.
"active_id": None,
"scoped_session_ids": scoped_session_ids,
"errors": errors,
# Ownership is per profile, so no project is "the active one" here; the
# desktop only reads active_id to bias its overview sort.
"active_id": None, "scoped_session_ids": scoped_session_ids, "errors": errors,
}
@@ -851,13 +813,8 @@ async def create_profile_endpoint(body: ProfileCreate):
with _profile_errors("POST /api/profiles failed", not_found=(),
bad_request=(ValueError, FileExistsError, FileNotFoundError)):
path = profiles_mod.create_profile(
name=body.name,
clone_from=clone_from,
clone_all=body.clone_all,
clone_config=clone_config,
no_skills=body.no_skills,
description=body.description,
)
name=body.name, clone_from=clone_from, clone_all=body.clone_all,
clone_config=clone_config, no_skills=body.no_skills, description=body.description)
# Match the CLI's profile-create flow: fresh named profiles get the
# bundled skills. Cloning already copied the source's skills (incl.
# user-installed); no_skills wrote the opt-out marker so seeding no-ops.
@@ -874,61 +831,37 @@ async def create_profile_endpoint(body: ProfileCreate):
# dashboard page or `<profile> setup` afterward.
provider = (body.provider or "").strip()
model = (body.model or "").strip()
model_set = False
if provider and model:
try:
_write_profile_model(path, provider, model)
model_set = True
except Exception:
_log.exception("Setting model for new profile %s failed", body.name)
mcp_written = 0
if body.mcp_servers:
try:
mcp_written = _write_profile_mcp_servers(path, body.mcp_servers)
except Exception:
_log.exception("Writing MCP servers for new profile %s failed", body.name)
model_set = bool(provider and model) and _best_effort(
"Setting model for new profile %s failed", body.name,
fn=lambda: (_write_profile_model(path, provider, model), True)[1], default=False)
mcp_written = _best_effort(
"Writing MCP servers for new profile %s failed", body.name,
fn=lambda: _write_profile_mcp_servers(path, body.mcp_servers), default=0
) if body.mcp_servers else 0
# "keep" skill selection has replace semantics: disable every seeded skill
# not in the list. Skipped when empty (legacy: keep the bundle).
skills_disabled = 0
if body.keep_skills:
try:
skills_disabled = _disable_unselected_skills(path, body.keep_skills)
except Exception:
_log.exception("Applying skill selection for new profile %s failed", body.name)
skills_disabled = _best_effort(
"Applying skill selection for new profile %s failed", body.name,
fn=lambda: _disable_unselected_skills(path, body.keep_skills), default=0
) if body.keep_skills else 0
# Skills-hub installs are spawned async, scoped to the new profile via
# `-p <name>` (a fresh subprocess re-binds skills_hub.SKILLS_DIR to the
# profile's HERMES_HOME at import). PIDs go back for the UI to poll.
hub_installs: List[Dict[str, Any]] = []
for identifier in body.hub_skills:
ident = (identifier or "").strip()
if not ident:
continue
try:
proc = _spawn_hermes_action(
["-p", body.name, "skills", "install", ident, "--yes"],
_hub_action_name("install", ident),
)
hub_installs.append({"identifier": ident, "pid": proc.pid})
except Exception:
_log.exception(
"Spawning hub-skill install %s for new profile %s failed",
ident,
body.name,
)
hub_installs.append({"identifier": ident, "pid": None})
def _spawn_install(ident: str):
return _spawn_hermes_action(["-p", body.name, "skills", "install", ident, "--yes"],
_hub_action_name("install", ident)).pid
return {
"ok": True,
"name": body.name,
"path": str(path),
"model_set": model_set,
"mcp_written": mcp_written,
"skills_disabled": skills_disabled,
"hub_installs": hub_installs,
}
hub_installs: List[Dict[str, Any]] = [
{"identifier": ident, "pid": _best_effort(
"Spawning hub-skill install %s for new profile %s failed", ident, body.name,
fn=lambda: _spawn_install(ident))}
for ident in ((i or "").strip() for i in body.hub_skills) if ident
]
return {"ok": True, "name": body.name, "path": str(path), "model_set": model_set,
"mcp_written": mcp_written, "skills_disabled": skills_disabled,
"hub_installs": hub_installs}
@router.get("/api/profiles/active")
@@ -1000,27 +933,16 @@ async def open_profile_terminal_endpoint(name: str):
subprocess.Popen(["cmd.exe", "/c", "start", "", command])
elif sys.platform == "darwin":
escaped = command.replace("\\", "\\\\").replace('"', '\\"')
applescript = (
'tell application "Terminal"\n'
"activate\n"
f'do script "{escaped}"\n'
"end tell"
)
subprocess.Popen(["osascript", "-e", applescript])
subprocess.Popen(["osascript", "-e",
f'tell application "Terminal"\nactivate\ndo script "{escaped}"\nend tell'])
else:
for executable, popen_args in _linux_terminal_commands(command):
if subprocess.call(
["which", executable],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
) == 0:
if subprocess.call(["which", executable], stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL) == 0:
subprocess.Popen(popen_args)
break
else:
raise HTTPException(
status_code=400,
detail="No supported terminal emulator found",
)
raise HTTPException(status_code=400, detail="No supported terminal emulator found")
return {"ok": True, "command": command}
@@ -1042,17 +964,8 @@ async def rename_profile_endpoint(name: str, body: ProfileRename):
except ValueError:
is_default = False
if is_default:
return {
"ok": True,
"name": "default",
"display_name": body.new_name.strip(),
"path": str(path),
}
return {
"ok": True,
"name": profiles_mod.normalize_profile_name(body.new_name),
"path": str(path),
}
return {"ok": True, "name": "default", "display_name": body.new_name.strip(), "path": str(path)}
return {"ok": True, "name": profiles_mod.normalize_profile_name(body.new_name), "path": str(path)}
@router.delete("/api/profiles/{name}")
@@ -1107,14 +1020,11 @@ async def update_profile_soul(name: str, body: ProfileSoulUpdate):
# at the umask default (profiles chmods only .env to 0600) and SOUL.md
# is not a secret. (The default profile's seeder runs on every
# load_config, so its file already exists and preserve_mode applies.)
atomic_write_text(
soul_path, body.content, preserve_mode=True, create_mode=0o644
)
atomic_write_text(soul_path, body.content, preserve_mode=True, create_mode=0o644)
try:
# atomic_write_text() writes a temp file, fsyncs it and replaces the
# original — three syscalls that block for as long as the filesystem
# takes to durably commit the persona document.
# original — blocking for as long as the filesystem takes to commit.
await run_in_threadpool(_run)
except OSError as e:
_log.exception("PUT /api/profiles/%s/soul failed", name)
@@ -1134,11 +1044,7 @@ async def update_profile_description_endpoint(name: str, body: ProfileDescriptio
not_found=(), bad_request=()):
# write_profile_meta() reads profile.yaml, merges and rewrites it.
await run_in_threadpool(
profiles_mod.write_profile_meta,
profile_dir,
description=text,
description_auto=False,
)
profiles_mod.write_profile_meta, profile_dir, description=text, description_auto=False)
return {"ok": True, "description": text, "description_auto": False}
@@ -1237,29 +1143,22 @@ async def import_profile_endpoint(body: ProfileImport):
profiles_mod.import_profile, archive, name=(body.name or "").strip() or None)
imported = profile_dir.name
# Match the CLI import flow: create the wrapper alias when it's safe.
try:
def _wrapper():
if not profiles_mod.check_alias_collision(imported):
profiles_mod.create_wrapper_script(imported)
except Exception:
_log.exception("Creating wrapper for imported profile %s failed", imported)
_best_effort("Creating wrapper for imported profile %s failed", imported, fn=_wrapper)
# Surface the bundled desktop appearance overlay (if the archive carried
# one) so the desktop can apply theme/interface prefs without another
# round-trip.
desktop_overlay = None
if (profile_dir / "desktop.json").is_file():
try:
desktop_overlay = _read_desktop_overlay(profile_dir)
except Exception:
_log.exception("Reading desktop.json from imported profile %s failed", imported)
return {
"ok": True,
"name": imported,
"path": str(profile_dir),
"desktop": desktop_overlay,
}
desktop_overlay = _best_effort(
"Reading desktop.json from imported profile %s failed", imported,
fn=lambda: _read_desktop_overlay(profile_dir))
return {"ok": True, "name": imported, "path": str(profile_dir), "desktop": desktop_overlay}
@router.get("/api/profiles/{name}/desktop-overlay")