refactor(local_runtime): join one-name-per-line imports and wrapped calls; guard-clause reflows in sweep_idle/_reap_orphaned_children/physics_check
This commit is contained in:
@@ -295,9 +295,7 @@ def ensure_runtime_installed(tag: str, backend: str,
|
||||
if progress is not None:
|
||||
progress("verify", 0, 0, "")
|
||||
version = verify_install(install_dir, tag)
|
||||
manifest_path.write_text(json.dumps({
|
||||
"tag": tag, "backend": plan.backend, "assets": recorded,
|
||||
"verified_version": version,
|
||||
}, indent=2), encoding="utf-8")
|
||||
manifest_path.write_text(json.dumps({"tag": tag, "backend": plan.backend, "assets": recorded,
|
||||
"verified_version": version}, indent=2), encoding="utf-8")
|
||||
logger.info("installed llama.cpp %s (%s): %s", tag, backend, version)
|
||||
return install_dir
|
||||
|
||||
@@ -106,13 +106,10 @@ def _stop_state_server(state: dict) -> None:
|
||||
|
||||
try:
|
||||
pid = int(state.get("pid"))
|
||||
except (TypeError, ValueError):
|
||||
return
|
||||
if pid <= 0:
|
||||
return
|
||||
try:
|
||||
if pid <= 0:
|
||||
return
|
||||
os.kill(pid, signal.SIGTERM)
|
||||
except (OSError, ValueError):
|
||||
except (TypeError, ValueError, OSError):
|
||||
return
|
||||
# Give it a moment to release the port and the GPU. Liveness via psutil — on Windows
|
||||
# os.kill(pid, 0) TERMINATES the process, it is not a probe.
|
||||
@@ -136,8 +133,7 @@ def refresh_local_runtime() -> bool:
|
||||
state = _state_endpoint()
|
||||
if state is None:
|
||||
return False
|
||||
logger.info("bouncing adopted llama-server (pid=%s) to rescan models",
|
||||
state.get("pid"))
|
||||
logger.info("bouncing adopted llama-server (pid=%s) to rescan models", state.get("pid"))
|
||||
_stop_state_server(state)
|
||||
else:
|
||||
shutdown_local_runtime()
|
||||
@@ -213,11 +209,7 @@ def ensure_local_runtime(config: dict, force: bool = False) -> "object | None":
|
||||
|
||||
try:
|
||||
from hermes_cli.local_runtime.binaries import (
|
||||
default_tag,
|
||||
ensure_runtime_installed,
|
||||
installed_tags,
|
||||
select_backend,
|
||||
)
|
||||
default_tag, ensure_runtime_installed, installed_tags, select_backend)
|
||||
from hermes_cli.local_runtime.supervisor import LlamaServerSupervisor
|
||||
|
||||
backend = section.get("backend", "auto")
|
||||
@@ -244,12 +236,9 @@ def ensure_local_runtime(config: dict, force: bool = False) -> "object | None":
|
||||
mdir.mkdir(parents=True, exist_ok=True)
|
||||
preset_path = _generate_presets(mdir, runtimes_root() / "presets.ini")
|
||||
|
||||
sup = LlamaServerSupervisor(
|
||||
install_dir, mdir,
|
||||
models_max=int(section.get("models_max", 4)),
|
||||
port=int(section.get("port", 0)) or None,
|
||||
preset_path=preset_path,
|
||||
)
|
||||
sup = LlamaServerSupervisor(install_dir, mdir, preset_path=preset_path,
|
||||
models_max=int(section.get("models_max", 4)),
|
||||
port=int(section.get("port", 0)) or None)
|
||||
try:
|
||||
sup.start()
|
||||
except Exception:
|
||||
@@ -259,8 +248,7 @@ def ensure_local_runtime(config: dict, force: bool = False) -> "object | None":
|
||||
sup.stop()
|
||||
raise
|
||||
_SUPERVISOR = sup
|
||||
logger.info("managed llama-server up at %s (backend=%s tag=%s)",
|
||||
sup.base_url, backend, tag)
|
||||
logger.info("managed llama-server up at %s (backend=%s tag=%s)", sup.base_url, backend, tag)
|
||||
_start_idle_sweeper(sup)
|
||||
return sup
|
||||
except Exception as exc: # noqa: BLE001 — never break session start
|
||||
@@ -294,5 +282,4 @@ def _start_idle_sweeper(sup) -> None:
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.debug("idle sweep skipped: %s", exc)
|
||||
|
||||
threading.Thread(target=_loop, daemon=True,
|
||||
name="local-runtime-idle-sweep").start()
|
||||
threading.Thread(target=_loop, daemon=True, name="local-runtime-idle-sweep").start()
|
||||
|
||||
@@ -19,17 +19,8 @@ from dataclasses import dataclass, field
|
||||
from pathlib import PurePosixPath
|
||||
|
||||
from hermes_cli.local_runtime.context_policy import (
|
||||
FLOOR,
|
||||
RUNTIME_OVERHEAD_BYTES,
|
||||
TARGET_WINDOW,
|
||||
ub_logits_bytes,
|
||||
)
|
||||
from hermes_cli.local_runtime.estimator import (
|
||||
HardwareBudget,
|
||||
LayerKind,
|
||||
ModelProfile,
|
||||
ctx_bytes,
|
||||
)
|
||||
FLOOR, RUNTIME_OVERHEAD_BYTES, TARGET_WINDOW, ub_logits_bytes)
|
||||
from hermes_cli.local_runtime.estimator import HardwareBudget, LayerKind, ModelProfile, ctx_bytes
|
||||
from hermes_cli.local_runtime.gguf import model_id_from_stem
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -122,11 +113,9 @@ class CatalogEntry:
|
||||
+ [(LayerKind.SWA, self.per_layer_f16)] * self.swa_layers
|
||||
+ [(LayerKind.RECURRENT, 0)] * self.recurrent_layers)
|
||||
return ModelProfile(
|
||||
name=variant.model_id, weights_bytes=variant.weights_bytes,
|
||||
embd_table_bytes=0, n_ctx_train=self.n_ctx_train,
|
||||
layers=layers, swa_window=self.swa_window, moe=self.moe,
|
||||
n_vocab=self.n_vocab,
|
||||
kv_scale=1.2 if self.mtp else 1.0)
|
||||
name=variant.model_id, weights_bytes=variant.weights_bytes, embd_table_bytes=0,
|
||||
n_ctx_train=self.n_ctx_train, layers=layers, swa_window=self.swa_window, moe=self.moe,
|
||||
n_vocab=self.n_vocab, kv_scale=1.2 if self.mtp else 1.0)
|
||||
|
||||
def download_files(self, variant: QuantVariant) -> tuple:
|
||||
"""Everything a download job fetches for this variant, in order."""
|
||||
@@ -191,8 +180,7 @@ _HOST_BANDWIDTH_GB_S = 80.0 # spilled weights stream over host DRAM
|
||||
PLEASANT_FLOOR_TOK_S = 20.0
|
||||
|
||||
|
||||
def predicted_decode_tok_s(entry: CatalogEntry, variant: QuantVariant,
|
||||
budget: HardwareBudget, *,
|
||||
def predicted_decode_tok_s(entry: CatalogEntry, variant: QuantVariant, budget: HardwareBudget, *,
|
||||
spilled: bool = False) -> float:
|
||||
"""Memory-bound decode prediction for ordering and floor-gating."""
|
||||
bandwidth = (_HOST_BANDWIDTH_GB_S if spilled
|
||||
@@ -251,8 +239,7 @@ _last_refresh_attempt = 0.0
|
||||
def _asset_from(d: "dict | None") -> "AssetFile | None":
|
||||
if not d:
|
||||
return None
|
||||
return AssetFile(path=d["path"], size_bytes=int(d["size_bytes"]),
|
||||
local=d.get("local"))
|
||||
return AssetFile(path=d["path"], size_bytes=int(d["size_bytes"]), local=d.get("local"))
|
||||
|
||||
|
||||
# Scalar CatalogEntry fields parsed from JSON: key -> (coerce, default); None default = required.
|
||||
@@ -274,11 +261,9 @@ def _load_catalog(doc: dict) -> "tuple[CatalogEntry, ...]":
|
||||
f"(this build reads {_SCHEMA_VERSION})")
|
||||
entries = []
|
||||
for m in doc["models"]:
|
||||
variants = tuple(
|
||||
QuantVariant(quant=v["quant"],
|
||||
files=tuple(_asset_from(f) for f in v["files"]),
|
||||
validated=bool(v.get("validated")))
|
||||
for v in m["variants"])
|
||||
variants = tuple(QuantVariant(quant=v["quant"], validated=bool(v.get("validated")),
|
||||
files=tuple(_asset_from(f) for f in v["files"]))
|
||||
for v in m["variants"])
|
||||
scalars = {k: coerce(m[k] if default is None else m.get(k, default))
|
||||
for k, (coerce, default) in _SCALAR_FIELDS.items()}
|
||||
entries.append(CatalogEntry(
|
||||
@@ -292,8 +277,7 @@ def _load_catalog(doc: dict) -> "tuple[CatalogEntry, ...]":
|
||||
def _packaged_catalog() -> "tuple[CatalogEntry, ...]":
|
||||
from importlib.resources import files
|
||||
|
||||
raw = files("hermes_cli.local_runtime").joinpath("catalog.json").read_text(
|
||||
encoding="utf-8")
|
||||
raw = files("hermes_cli.local_runtime").joinpath("catalog.json").read_text(encoding="utf-8")
|
||||
return _load_catalog(json.loads(raw))
|
||||
|
||||
|
||||
@@ -312,8 +296,7 @@ def refresh_catalog(force: bool = False) -> bool:
|
||||
return False
|
||||
_last_refresh_attempt = now
|
||||
try:
|
||||
req = urllib.request.Request(
|
||||
_CATALOG_URL, headers={"User-Agent": "hermes-local-runtime"})
|
||||
req = urllib.request.Request(_CATALOG_URL, headers={"User-Agent": "hermes-local-runtime"})
|
||||
with urllib.request.urlopen(req, timeout=10) as r:
|
||||
fetched = _load_catalog(json.load(r))
|
||||
except Exception as exc: # noqa: BLE001
|
||||
@@ -330,8 +313,7 @@ def refresh_catalog_soon() -> None:
|
||||
it already has — the refresh lands for the next one."""
|
||||
if time.monotonic() - _last_refresh_attempt < _REFRESH_TTL_S:
|
||||
return
|
||||
threading.Thread(target=refresh_catalog, daemon=True,
|
||||
name="catalog-refresh").start()
|
||||
threading.Thread(target=refresh_catalog, daemon=True, name="catalog-refresh").start()
|
||||
|
||||
|
||||
def catalog_by_id() -> dict[str, CatalogEntry]:
|
||||
|
||||
@@ -10,12 +10,7 @@ from __future__ import annotations
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
from hermes_cli.local_runtime.estimator import (
|
||||
HardwareBudget,
|
||||
ModelProfile,
|
||||
PhysicsRefusal,
|
||||
ctx_bytes,
|
||||
physics_check,
|
||||
)
|
||||
HardwareBudget, ModelProfile, PhysicsRefusal, ctx_bytes, physics_check)
|
||||
|
||||
FLOOR = 64 * 1024 # = target; one internal constant
|
||||
_LADDER_GROWTH = 1.5
|
||||
@@ -60,8 +55,7 @@ class WindowDecision:
|
||||
return self.spill_bytes > 0
|
||||
|
||||
|
||||
def initial_window(profile: ModelProfile, budget: HardwareBudget,
|
||||
*, flash_attention: bool = True,
|
||||
def initial_window(profile: ModelProfile, budget: HardwareBudget, *, flash_attention: bool = True,
|
||||
overhead_bytes: int = 0) -> WindowDecision | PhysicsRefusal:
|
||||
"""The launch decision: largest cheap rung, never below the floor.
|
||||
|
||||
@@ -103,10 +97,9 @@ def initial_window(profile: ModelProfile, budget: HardwareBudget,
|
||||
reason = f"floor held at {window // 1024}K; weights spill (deliberate price of the guarantee)"
|
||||
|
||||
kv_bytes = kv(window)
|
||||
spill = max(0, profile.weights_bytes + kv_bytes - budget.usable_vram_bytes)
|
||||
return WindowDecision(window=window, spill_bytes=spill,
|
||||
kv_on_gpu=kv_bytes <= budget.usable_vram_bytes,
|
||||
reasons=[reason])
|
||||
return WindowDecision(window=window, reasons=[reason],
|
||||
spill_bytes=max(0, profile.weights_bytes + kv_bytes - budget.usable_vram_bytes),
|
||||
kv_on_gpu=kv_bytes <= budget.usable_vram_bytes)
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -117,10 +110,8 @@ class GrowthDecision:
|
||||
|
||||
|
||||
def growth_decision(profile: ModelProfile, budget: HardwareBudget, *,
|
||||
current_window: int, session_tokens: int,
|
||||
measured_decode_tok_s: float | None,
|
||||
server_idle: bool,
|
||||
flash_attention: bool = True,
|
||||
current_window: int, session_tokens: int, measured_decode_tok_s: float | None,
|
||||
server_idle: bool, flash_attention: bool = True,
|
||||
occupancy_confirmed: bool = False) -> GrowthDecision:
|
||||
"""One growth evaluation, END-OF-TURN ONLY (recurrent state cannot rewind mid-sequence).
|
||||
|
||||
@@ -135,8 +126,7 @@ def growth_decision(profile: ModelProfile, budget: HardwareBudget, *,
|
||||
|
||||
native = profile.n_ctx_train or current_window
|
||||
if current_window >= native:
|
||||
return GrowthDecision("compress-default",
|
||||
reason="at native window; compression is the only move")
|
||||
return GrowthDecision("compress-default", reason="at native window; compression is the only move")
|
||||
|
||||
if not server_idle:
|
||||
return GrowthDecision("hold", reason="server busy; re-grant deferred to idle")
|
||||
@@ -154,8 +144,7 @@ def growth_decision(profile: ModelProfile, budget: HardwareBudget, *,
|
||||
# that no longer fits doesn't get granted.
|
||||
kv = ctx_bytes(profile, next_rung, flash_attention=flash_attention)
|
||||
if profile.weights_bytes + kv > budget.usable_vram_bytes + budget.ram_available_bytes:
|
||||
return GrowthDecision("compress-default",
|
||||
reason="next rung exceeds physics; compression instead")
|
||||
return GrowthDecision("compress-default", reason="next rung exceeds physics; compression instead")
|
||||
|
||||
return GrowthDecision("grow", next_window=next_rung,
|
||||
reason=f"rung {current_window // 1024}K -> {next_rung // 1024}K")
|
||||
@@ -172,19 +161,15 @@ def spill_overrides(profile: ModelProfile) -> list[str]:
|
||||
return [] # dense: fit's back-to-front layer cut is the only axis
|
||||
|
||||
|
||||
def launch_args(profile: ModelProfile, decision: WindowDecision, *,
|
||||
flash_attention: bool = True,
|
||||
mtp_capable: bool = False,
|
||||
mtp_draft_depth: int = 3,
|
||||
uma: bool = False,
|
||||
def launch_args(profile: ModelProfile, decision: WindowDecision, *, flash_attention: bool = True,
|
||||
mtp_capable: bool = False, mtp_draft_depth: int = 3, uma: bool = False,
|
||||
mtp_prefill: bool = False) -> list[str]:
|
||||
"""Per-model launch flags from a window decision. Explicit -c puts fit into
|
||||
spill-weights-and-hold-ctx; q8 KV cache wherever flash attention exists; -ot placement on
|
||||
spilled configs — DISCRETE cards only."""
|
||||
args = ["-c", str(decision.window)]
|
||||
if mtp_capable:
|
||||
args += ["--spec-type", "draft-mtp",
|
||||
"--spec-draft-n-max", str(mtp_draft_depth),
|
||||
args += ["--spec-type", "draft-mtp", "--spec-draft-n-max", str(mtp_draft_depth),
|
||||
"--backend-sampling", "--spec-draft-backend-sampling"]
|
||||
if mtp_prefill:
|
||||
args += ["-b", "4096", "-ub", "2048"]
|
||||
@@ -197,8 +182,7 @@ def launch_args(profile: ModelProfile, decision: WindowDecision, *,
|
||||
return args
|
||||
|
||||
|
||||
def ub_logits_bytes(n_vocab: int, *, mtp_capable: bool,
|
||||
mtp_prefill: bool = False) -> int:
|
||||
def ub_logits_bytes(n_vocab: int, *, mtp_capable: bool, mtp_prefill: bool = False) -> int:
|
||||
"""GPU logits/compute-buffer cost of the microbatch posture chosen by launch_args, priced from
|
||||
the model's own vocab and calibrated against measured server RSS (Qwen3.8 Q4, both postures,
|
||||
three windows)."""
|
||||
|
||||
@@ -42,8 +42,7 @@ def probe_port(port: int) -> DetectedServer | None:
|
||||
root = f"http://127.0.0.1:{port}"
|
||||
status, props = _get(f"{root}/props")
|
||||
if status == 401:
|
||||
return DetectedServer(base_url=f"{root}/v1", build_info="", model_path="",
|
||||
n_ctx=None, router_mode=False, auth_required=True)
|
||||
return DetectedServer(f"{root}/v1", "", "", None, router_mode=False, auth_required=True)
|
||||
if status != 200 or not isinstance(props, dict):
|
||||
return None
|
||||
build = str(props.get("build_info", ""))
|
||||
@@ -52,14 +51,10 @@ def probe_port(port: int) -> DetectedServer | None:
|
||||
dgs = props.get("default_generation_settings")
|
||||
models_status, models = _get(f"{root}/models")
|
||||
return DetectedServer(
|
||||
base_url=f"{root}/v1",
|
||||
build_info=build,
|
||||
model_path=str(props.get("model_path", "")),
|
||||
base_url=f"{root}/v1", build_info=build, model_path=str(props.get("model_path", "")),
|
||||
n_ctx=dgs.get("n_ctx") if isinstance(dgs, dict) else None,
|
||||
router_mode=(models_status == 200 and isinstance(models, dict)
|
||||
and "data" in models),
|
||||
auth_required=False,
|
||||
)
|
||||
router_mode=models_status == 200 and isinstance(models, dict) and "data" in models,
|
||||
auth_required=False)
|
||||
|
||||
|
||||
def detect_server(extra_ports: tuple[int, ...] = ()) -> DetectedServer | None:
|
||||
|
||||
@@ -95,16 +95,9 @@ def profile_from_gguf(header: GGUFHeader) -> ModelProfile:
|
||||
n_attn_seen += 1
|
||||
|
||||
return ModelProfile(
|
||||
name=header.path,
|
||||
weights_bytes=header.tensor_bytes,
|
||||
embd_table_bytes=header.embd_table_bytes,
|
||||
n_ctx_train=header.n_ctx_train,
|
||||
layers=layers,
|
||||
swa_window=header.sliding_window,
|
||||
moe=header.expert_count > 0,
|
||||
architecture=header.architecture,
|
||||
n_vocab=header.n_vocab,
|
||||
)
|
||||
name=header.path, weights_bytes=header.tensor_bytes, embd_table_bytes=header.embd_table_bytes,
|
||||
n_ctx_train=header.n_ctx_train, layers=layers, swa_window=header.sliding_window,
|
||||
moe=header.expert_count > 0, architecture=header.architecture, n_vocab=header.n_vocab)
|
||||
|
||||
|
||||
def kv_dtype_factor(flash_attention: bool) -> float:
|
||||
@@ -113,8 +106,7 @@ def kv_dtype_factor(flash_attention: bool) -> float:
|
||||
return (_Q8_BYTES_PER_ELEM / _F16_BYTES_PER_ELEM) if flash_attention else 1.0
|
||||
|
||||
|
||||
def ctx_bytes(profile: ModelProfile, window: int, *,
|
||||
flash_attention: bool = True) -> int:
|
||||
def ctx_bytes(profile: ModelProfile, window: int, *, flash_attention: bool = True) -> int:
|
||||
"""Context memory for one window: full layers linear in T, SWA layers capped at the sliding
|
||||
window, recurrent layers constant. Scaled by profile.kv_scale (MTP draft context)."""
|
||||
factor = kv_dtype_factor(flash_attention)
|
||||
@@ -145,11 +137,11 @@ def physics_check(profile: ModelProfile, budget: HardwareBudget,
|
||||
+ ctx_bytes(profile, min(floor, profile.n_ctx_train or floor),
|
||||
flash_attention=flash_attention))
|
||||
available = budget.usable_vram_bytes + budget.ram_available_bytes
|
||||
if needed > available:
|
||||
gib = 1 << 30
|
||||
return PhysicsRefusal(
|
||||
needed_bytes=needed, available_bytes=available,
|
||||
message=(f"{profile.name}: needs ~{needed / gib:.1f} GiB at the "
|
||||
f"{floor // 1024}K floor but only ~{available / gib:.1f} GiB "
|
||||
"of VRAM+RAM exist — try a smaller quant (UD-Q3/Q2)"))
|
||||
return None
|
||||
if needed <= available:
|
||||
return None
|
||||
gib = 1 << 30
|
||||
return PhysicsRefusal(
|
||||
needed_bytes=needed, available_bytes=available,
|
||||
message=(f"{profile.name}: needs ~{needed / gib:.1f} GiB at the "
|
||||
f"{floor // 1024}K floor but only ~{available / gib:.1f} GiB "
|
||||
"of VRAM+RAM exist — try a smaller quant (UD-Q3/Q2)"))
|
||||
|
||||
@@ -70,10 +70,7 @@ def maybe_grow_window(model_id: str, *, base_url: str, session_tokens: int,
|
||||
window — nothing rewinds.
|
||||
"""
|
||||
from hermes_cli.local_runtime.bootstrap import (
|
||||
get_supervisor,
|
||||
refresh_local_runtime,
|
||||
staged_models,
|
||||
)
|
||||
get_supervisor, refresh_local_runtime, staged_models)
|
||||
from hermes_cli.local_runtime.context_policy import growth_decision
|
||||
from hermes_cli.local_runtime.estimator import profile_from_gguf
|
||||
from hermes_cli.local_runtime.gguf import read_gguf_header
|
||||
@@ -83,8 +80,7 @@ def maybe_grow_window(model_id: str, *, base_url: str, session_tokens: int,
|
||||
if sup is None or not is_managed_endpoint(base_url):
|
||||
return None
|
||||
|
||||
gguf = next((p for p in staged_models()
|
||||
if p.stem.startswith(model_id) or model_id in p.stem), None)
|
||||
gguf = next((p for p in staged_models() if p.stem.startswith(model_id) or model_id in p.stem), None)
|
||||
if gguf is None:
|
||||
return None
|
||||
|
||||
|
||||
@@ -64,8 +64,7 @@ class HFFileGroup:
|
||||
fit: str = "unknown" # fits-gpu | needs-ram | too-big | unknown
|
||||
|
||||
|
||||
_QUANT_RE = re.compile(
|
||||
r"(?:IQ|Q)\d[_A-Z0-9]*|F16|BF16|F32", re.IGNORECASE)
|
||||
_QUANT_RE = re.compile(r"(?:IQ|Q)\d[_A-Z0-9]*|F16|BF16|F32", re.IGNORECASE)
|
||||
_SPLIT_RE = re.compile(r"-(\d{5})-of-(\d{5})\.gguf$", re.IGNORECASE)
|
||||
|
||||
|
||||
@@ -114,10 +113,8 @@ def repo_files(repo: str) -> list[HFFileGroup]:
|
||||
for path, size in singles]
|
||||
for stem, parts in splits.items():
|
||||
parts.sort()
|
||||
groups.append(HFFileGroup(
|
||||
label=_quant_label(stem),
|
||||
paths=tuple(p for _, p, _ in parts),
|
||||
total_bytes=sum(s for _, _, s in parts)))
|
||||
groups.append(HFFileGroup(label=_quant_label(stem), paths=tuple(p for _, p, _ in parts),
|
||||
total_bytes=sum(s for _, _, s in parts)))
|
||||
groups.sort(key=lambda g: g.total_bytes, reverse=True)
|
||||
return groups
|
||||
|
||||
@@ -134,5 +131,4 @@ def rough_fit(total_bytes: int, budget) -> str:
|
||||
|
||||
|
||||
def priced_repo_files(repo: str, budget) -> list[HFFileGroup]:
|
||||
return [replace(g, fit=rough_fit(g.total_bytes, budget))
|
||||
for g in repo_files(repo)]
|
||||
return [replace(g, fit=rough_fit(g.total_bytes, budget)) for g in repo_files(repo)]
|
||||
|
||||
@@ -9,6 +9,7 @@ support (older engines) all read as "nothing loading".
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from contextlib import suppress
|
||||
import json
|
||||
import logging
|
||||
import threading
|
||||
@@ -61,8 +62,7 @@ def _apply_event(model: str, event: str, data: dict) -> None:
|
||||
stages = [str(s) for s in (progress.get("stages") or [])]
|
||||
current = str(progress.get("current", ""))
|
||||
value = progress.get("value")
|
||||
entry = _snapshot.setdefault(model, {"stage": "", "value": 0.0,
|
||||
"percent": 0, "ts": 0.0})
|
||||
entry = _snapshot.setdefault(model, {"stage": "", "value": 0.0, "percent": 0, "ts": 0.0})
|
||||
entry["ts"] = time.monotonic()
|
||||
if current and isinstance(value, (int, float)):
|
||||
entry["stage"] = current
|
||||
@@ -87,10 +87,8 @@ def _watch() -> None:
|
||||
continue
|
||||
base, key = endpoint
|
||||
try:
|
||||
req = urllib.request.Request(
|
||||
f"{base}/models/sse",
|
||||
headers={"Authorization": f"Bearer {key}",
|
||||
"Accept": "text/event-stream"})
|
||||
req = urllib.request.Request(f"{base}/models/sse", headers={
|
||||
"Authorization": f"Bearer {key}", "Accept": "text/event-stream"})
|
||||
with urllib.request.urlopen(req, timeout=60) as r:
|
||||
buf = b""
|
||||
while True:
|
||||
@@ -103,13 +101,10 @@ def _watch() -> None:
|
||||
text = line.decode("utf-8", "replace").strip()
|
||||
if not text.startswith("data:"):
|
||||
continue
|
||||
try:
|
||||
with suppress(json.JSONDecodeError, TypeError):
|
||||
msg = json.loads(text[5:].strip())
|
||||
_apply_event(str(msg.get("model", "")),
|
||||
str(msg.get("event", "")),
|
||||
_apply_event(str(msg.get("model", "")), str(msg.get("event", "")),
|
||||
msg.get("data") or {})
|
||||
except (json.JSONDecodeError, TypeError):
|
||||
continue
|
||||
except Exception as exc: # noqa: BLE001 — watcher must never die loud
|
||||
logger.debug("load-progress SSE reconnecting: %s", exc)
|
||||
# Stream ended (router bounce, timeout, error): loading entries from the dead connection
|
||||
@@ -122,8 +117,7 @@ def _ensure_watcher() -> None:
|
||||
global _watcher
|
||||
with _lock:
|
||||
if _watcher is None or not _watcher.is_alive():
|
||||
_watcher = threading.Thread(target=_watch, daemon=True,
|
||||
name="llamacpp-load-progress")
|
||||
_watcher = threading.Thread(target=_watch, daemon=True, name="llamacpp-load-progress")
|
||||
_watcher.start()
|
||||
|
||||
|
||||
@@ -133,10 +127,8 @@ def get_loading_progress() -> dict[str, dict]:
|
||||
_ensure_watcher()
|
||||
now = time.monotonic()
|
||||
with _lock:
|
||||
return {m: {"stage": e["stage"], "value": e["value"],
|
||||
"percent": e["percent"]}
|
||||
for m, e in _snapshot.items()
|
||||
if now - e["ts"] < _STALE_ENTRY_TTL_S}
|
||||
return {m: {"stage": e["stage"], "value": e["value"], "percent": e["percent"]}
|
||||
for m, e in _snapshot.items() if now - e["ts"] < _STALE_ENTRY_TTL_S}
|
||||
|
||||
|
||||
def get_prefill_progress(model: str) -> "dict | None":
|
||||
@@ -161,11 +153,7 @@ def get_prefill_progress(model: str) -> "dict | None":
|
||||
return None
|
||||
best = 0
|
||||
for slot in slots if isinstance(slots, list) else []:
|
||||
if not slot.get("is_processing"):
|
||||
continue
|
||||
try:
|
||||
processed = int(slot.get("n_prompt_tokens_processed") or 0)
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
best = max(best, processed)
|
||||
with suppress(TypeError, ValueError):
|
||||
if slot.get("is_processing"):
|
||||
best = max(best, int(slot.get("n_prompt_tokens_processed") or 0))
|
||||
return {"processed": best} if best > 0 else None
|
||||
|
||||
@@ -9,19 +9,9 @@ from dataclasses import dataclass, replace
|
||||
from pathlib import Path
|
||||
|
||||
from hermes_cli.local_runtime.context_policy import (
|
||||
RUNTIME_OVERHEAD_BYTES,
|
||||
WindowDecision,
|
||||
initial_window,
|
||||
launch_args,
|
||||
ub_logits_bytes,
|
||||
)
|
||||
RUNTIME_OVERHEAD_BYTES, WindowDecision, initial_window, launch_args, ub_logits_bytes)
|
||||
from hermes_cli.local_runtime.estimator import (
|
||||
HardwareBudget,
|
||||
ModelProfile,
|
||||
PhysicsRefusal,
|
||||
ctx_bytes,
|
||||
profile_from_gguf,
|
||||
)
|
||||
HardwareBudget, ModelProfile, PhysicsRefusal, ctx_bytes, profile_from_gguf)
|
||||
from hermes_cli.local_runtime.gguf import model_id_from_stem, read_gguf_header
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -145,9 +135,8 @@ def _preset_for(gguf: Path, budget: HardwareBudget,
|
||||
|
||||
# The launch flags MUST match the pricing above (same entry/is_mtp/posture).
|
||||
keys = _args_to_keys(launch_args(
|
||||
profile, decision, mtp_capable=is_mtp,
|
||||
mtp_draft_depth=entry.mtp_draft_depth if entry is not None else 3,
|
||||
uma=budget.uma, mtp_prefill=mtp_prefill))
|
||||
profile, decision, mtp_capable=is_mtp, uma=budget.uma, mtp_prefill=mtp_prefill,
|
||||
mtp_draft_depth=entry.mtp_draft_depth if entry is not None else 3))
|
||||
if entry is not None and is_mtp:
|
||||
# Integrated-MTP targets sample on the backend, and so does the draft (pairing validated
|
||||
# against the vendor's published llama.cpp recipes).
|
||||
@@ -175,8 +164,7 @@ def _preset_for(gguf: Path, budget: HardwareBudget,
|
||||
spilled=decision.spilled, keys=keys)
|
||||
|
||||
|
||||
def generate_presets(models_dir: Path, budget: HardwareBudget,
|
||||
preset_path: Path,
|
||||
def generate_presets(models_dir: Path, budget: HardwareBudget, preset_path: Path,
|
||||
mtp_capable: set[str] | None = None) -> list[PresetEntry]:
|
||||
"""Walk the staged models, run the launch decision per model, and write one INI. Refused
|
||||
models get no section (the picker surfaces the refusal from the returned entries)."""
|
||||
|
||||
@@ -136,11 +136,8 @@ class LlamaServerSupervisor:
|
||||
headers = {"Authorization": f"Bearer {self.api_key}"}
|
||||
if json_type:
|
||||
headers["Content-Type"] = "application/json"
|
||||
req = urllib.request.Request(
|
||||
self._url(route),
|
||||
data=json.dumps(body).encode() if body is not None else None,
|
||||
headers=headers,
|
||||
)
|
||||
req = urllib.request.Request(self._url(route), headers=headers,
|
||||
data=json.dumps(body).encode() if body is not None else None)
|
||||
return urllib.request.urlopen(req, timeout=timeout_s)
|
||||
|
||||
def _request(self, route: str, body: dict | None = None, timeout_s: int = 30) -> dict:
|
||||
@@ -194,26 +191,21 @@ class LlamaServerSupervisor:
|
||||
self._spawn()
|
||||
self._wait_health(timeout_s)
|
||||
self._write_state()
|
||||
self._watchdog = threading.Thread(target=self._watch, daemon=True,
|
||||
name="llamacpp-supervisor")
|
||||
self._watchdog = threading.Thread(target=self._watch, daemon=True, name="llamacpp-supervisor")
|
||||
self._watchdog.start()
|
||||
|
||||
def _write_state(self) -> None:
|
||||
path = state_path()
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps({
|
||||
"base_url": self.base_url,
|
||||
"api_key": self.api_key,
|
||||
"pid": self.proc.pid if self.proc else None,
|
||||
}), encoding="utf-8")
|
||||
path.write_text(json.dumps({"base_url": self.base_url, "api_key": self.api_key,
|
||||
"pid": self.proc.pid if self.proc else None}), encoding="utf-8")
|
||||
|
||||
def _wait_health(self, timeout_s: int) -> None:
|
||||
deadline = time.monotonic() + timeout_s
|
||||
while time.monotonic() < deadline:
|
||||
if self.proc and self.proc.poll() is not None:
|
||||
raise RuntimeError(
|
||||
f"llama-server exited rc={self.proc.returncode} during startup "
|
||||
f"(log: {self.log_path})")
|
||||
raise RuntimeError(f"llama-server exited rc={self.proc.returncode} during startup "
|
||||
f"(log: {self.log_path})")
|
||||
with suppress(urllib.error.URLError, OSError, TimeoutError):
|
||||
with urllib.request.urlopen(self._url("/health"), timeout=3) as r:
|
||||
if r.status == 200:
|
||||
@@ -234,8 +226,7 @@ class LlamaServerSupervisor:
|
||||
if self._stopping:
|
||||
return
|
||||
backoff = _RESTART_BACKOFF_S[min(self._restarts, len(_RESTART_BACKOFF_S) - 1)]
|
||||
logger.warning("llama-server exited rc=%s; restart #%s in %ss",
|
||||
rc, self._restarts + 1, backoff)
|
||||
logger.warning("llama-server exited rc=%s; restart #%s in %ss", rc, self._restarts + 1, backoff)
|
||||
time.sleep(backoff)
|
||||
self._restarts += 1
|
||||
try:
|
||||
@@ -292,27 +283,22 @@ class LlamaServerSupervisor:
|
||||
exe = str(server_binary(self.install_dir))
|
||||
except Exception: # noqa: BLE001
|
||||
return
|
||||
own_pid = self.proc.pid if self.proc is not None else None
|
||||
for p in psutil.process_iter(["exe", "ppid"]):
|
||||
try:
|
||||
if p.info.get("exe") != exe:
|
||||
continue
|
||||
if self.proc is not None and p.pid == self.proc.pid:
|
||||
continue
|
||||
with suppress(psutil.NoSuchProcess, psutil.AccessDenied):
|
||||
ppid = p.info.get("ppid") or 0
|
||||
if ppid and psutil.pid_exists(ppid):
|
||||
if (p.info.get("exe") != exe or p.pid == own_pid
|
||||
or (ppid and psutil.pid_exists(ppid))):
|
||||
continue
|
||||
logger.warning("reaping orphaned llama-server child pid=%s", p.pid)
|
||||
p.kill()
|
||||
except (psutil.NoSuchProcess, psutil.AccessDenied):
|
||||
continue
|
||||
|
||||
# ── model management (router endpoints) ──────────────────
|
||||
|
||||
def models(self) -> dict:
|
||||
"""{model_id: status_value} from GET /models."""
|
||||
data = self._request("/models")
|
||||
return {m["id"]: m.get("status", {}).get("value", "unknown")
|
||||
for m in data.get("data", [])}
|
||||
for m in self._request("/models").get("data", [])}
|
||||
|
||||
def load_model(self, model_id: str, timeout_s: int = 600) -> None:
|
||||
self._request("/models/load", {"model": model_id}, timeout_s=timeout_s)
|
||||
@@ -344,15 +330,15 @@ class LlamaServerSupervisor:
|
||||
self._idle_since.pop(model_id, None)
|
||||
continue
|
||||
first_idle = self._idle_since.setdefault(model_id, now)
|
||||
if now - first_idle >= self.IDLE_UNLOAD_S:
|
||||
try:
|
||||
self.unload_model(model_id)
|
||||
self._idle_since.pop(model_id, None)
|
||||
unloaded.append(model_id)
|
||||
logger.info("idle-unloaded %s (idle %ds)", model_id,
|
||||
int(now - first_idle))
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.warning("idle unload of %s failed: %s", model_id, exc)
|
||||
if now - first_idle < self.IDLE_UNLOAD_S:
|
||||
continue
|
||||
try:
|
||||
self.unload_model(model_id)
|
||||
self._idle_since.pop(model_id, None)
|
||||
unloaded.append(model_id)
|
||||
logger.info("idle-unloaded %s (idle %ds)", model_id, int(now - first_idle))
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.warning("idle unload of %s failed: %s", model_id, exc)
|
||||
return unloaded
|
||||
|
||||
def touch_generate(self, model_id: str, timeout_s: int = 300) -> bool:
|
||||
@@ -360,10 +346,8 @@ class LlamaServerSupervisor:
|
||||
false-fail reasoning models, which spend their first tokens thinking."""
|
||||
try:
|
||||
resp = self._request("/v1/chat/completions", {
|
||||
"model": model_id,
|
||||
"messages": [{"role": "user", "content": TOUCH_PROMPT}],
|
||||
"max_tokens": 512, "temperature": 0,
|
||||
}, timeout_s=timeout_s)
|
||||
"model": model_id, "messages": [{"role": "user", "content": TOUCH_PROMPT}],
|
||||
"max_tokens": 512, "temperature": 0}, timeout_s=timeout_s)
|
||||
msg = resp["choices"][0]["message"]
|
||||
blob = (msg.get("content") or "") + " " + (msg.get("reasoning_content") or "")
|
||||
return TOUCH_EXPECT in blob.lower()
|
||||
@@ -387,10 +371,8 @@ class LlamaServerSupervisor:
|
||||
per-child and require ?model= (bare calls 400). With ``model_id`` checks that one child;
|
||||
without, every loaded child."""
|
||||
try:
|
||||
if model_id is not None:
|
||||
loaded = [model_id]
|
||||
else:
|
||||
loaded = [m for m, status in self.models().items() if status in _RESIDENT]
|
||||
loaded = ([model_id] if model_id is not None
|
||||
else [m for m, status in self.models().items() if status in _RESIDENT])
|
||||
for mid in loaded:
|
||||
slots = self._request(f"/slots?model={mid}")
|
||||
if any(s.get("is_processing") for s in slots):
|
||||
|
||||
Reference in New Issue
Block a user