diff --git a/README.md b/README.md index fcfd189..b80c223 100644 --- a/README.md +++ b/README.md @@ -83,7 +83,10 @@ transcript.](.github/images/workspace-light.png) a folder of `.gguf` files and it serves them through its own managed llama.cpp runtime -- no separate install. The `llama-server` binary is fetched once, verified against a pinned SHA-256, and cached. GPU backend selection is - `auto` (try Vulkan, fall back to CPU), `vulkan`, or `cpu`. + `auto` (try Vulkan, fall back to CPU; Vulkan is skipped when the machine has + no Vulkan loader), `vulkan`, or `cpu`. A loaded model is released after 30 + idle minutes (Settings > System; 0 keeps it loaded) or with "Unload model", + and its context window is held to what the model was trained for. - **Bring your own GGUF.** Download a model into the local folder by direct URL or Hugging Face repo, then select it from the same picker as everything else. For a repository, Settings can list its `.gguf` files (with sizes, folders diff --git a/app_factory.py b/app_factory.py index 36f4def..4fca331 100644 --- a/app_factory.py +++ b/app_factory.py @@ -128,6 +128,8 @@ def gguf_directory() -> Path: gpu_backend_setting=lambda: settings_repository.load().settings.llamacpp.gpu_backend, models_directory=gguf_directory, verify=ssl_context, + idle_unload_minutes=lambda: settings_repository.load().settings.llamacpp.idle_unload_minutes, + extra_args=lambda: settings_repository.load().settings.llamacpp.extra_args, ) gguf_model_directory = GGUFModelDirectory(gguf_directory) # ollama.Client.pull is overloaded on a Literal `stream`, one overload per diff --git a/backend/cortex_backend/api/routers/system.py b/backend/cortex_backend/api/routers/system.py index 3a0ded4..601f5c7 100644 --- a/backend/cortex_backend/api/routers/system.py +++ b/backend/cortex_backend/api/routers/system.py @@ -17,10 +17,12 @@ ) from cortex_backend.api.schemas import ( DiagnosticsResponse, + LlamaCppRuntimeStatus, ShutdownResponse, SystemResponse, ) from cortex_backend.api.security import SessionPrincipal +from cortex_backend.llamacpp.errors import LlamaCppError, RuntimeBusyError from fastapi import ( Depends, HTTPException, @@ -72,6 +74,35 @@ def system( ) + @router.post("/llamacpp/unload", response_model=LlamaCppRuntimeStatus) + def unload_llamacpp( + request: Request, + _: SessionPrincipal = Depends(require_session), + ) -> LlamaCppRuntimeStatus: + """Stop the loaded local model to free its memory; the next message loads it again. + + Safe to repeat: with nothing loaded it changes nothing and reports the + status. Refused with 409 while a response is being generated or a model + is loading, because that would cut the answer off. + """ + manager = getattr(request.app.state, "llamacpp_manager", None) + unload = getattr(manager, "unload", None) + if not callable(unload): + raise HTTPException(status_code=409, detail="The local model runtime is unavailable in this preview.") + if request.app.state.jobs.active_snapshot(kind="generation") is not None: + raise HTTPException( + status_code=409, + detail="A response is being generated. Stop it or wait for it to finish, then unload the model.", + ) + try: + unload() + except RuntimeBusyError as exc: + raise HTTPException(status_code=409, detail=exc.error) from exc + except LlamaCppError as exc: + raise HTTPException(status_code=500, detail=exc.error) from exc + return _llamacpp_status(request) + + @router.post("/system/shutdown", response_model=ShutdownResponse) def shutdown( request: Request, diff --git a/backend/cortex_backend/api/routes.py b/backend/cortex_backend/api/routes.py index 2dbc672..ea48f7d 100644 --- a/backend/cortex_backend/api/routes.py +++ b/backend/cortex_backend/api/routes.py @@ -1136,6 +1136,10 @@ def _llamacpp_status(request: Request) -> LlamaCppRuntimeStatus: last_restart_reason=live.last_restart_reason, loaded_context=live.loaded_context, last_failure_code=live.last_failure_code, + gpu_layers_offloaded=live.gpu_layers_offloaded, + gpu_layers_total=live.gpu_layers_total, + backend_note=live.backend_note, + context_note=live.context_note, ) diff --git a/backend/cortex_backend/api/schemas.py b/backend/cortex_backend/api/schemas.py index ad25237..ca38543 100644 --- a/backend/cortex_backend/api/schemas.py +++ b/backend/cortex_backend/api/schemas.py @@ -112,6 +112,19 @@ class LlamaCppRuntimeStatus(APIModel): # produced. Null when nothing failed, when the cause was not identified, # and once a server is ready. last_failure_code: LaunchFailureCode | None = None + # How many of the model's layers the running server put on the GPU, and how + # many it has, as the server reported them while loading. Null while + # nothing is ready and when the server said nothing Cortex recognises -- + # unknown, not zero. ``active_backend`` says which build launched; these say + # whether the GPU is actually in use (0 offloaded means it is not). + gpu_layers_offloaded: int | None = None + gpu_layers_total: int | None = None + # Fixed text on why the GPU build was not used when it would have been the + # default (no Vulkan loader on this machine). Null otherwise. + backend_note: str | None = None + # Fixed text saying the context window was limited to what the model was + # trained for. Null when the window is as requested. + context_note: str | None = None class SystemResponse(APIModel): diff --git a/backend/cortex_backend/core/settings.py b/backend/cortex_backend/core/settings.py index 6c87734..9c509b6 100644 --- a/backend/cortex_backend/core/settings.py +++ b/backend/cortex_backend/core/settings.py @@ -4,7 +4,9 @@ from typing import Annotated, Literal -from pydantic import BaseModel, ConfigDict, Field, StringConstraints +from pydantic import BaseModel, ConfigDict, Field, StringConstraints, field_validator + +from cortex_backend.llamacpp.extra_args import validate_extra_args ModelTag = Annotated[ @@ -57,6 +59,21 @@ class LlamaCppSettings(_SettingsModel): # "auto" tries Vulkan (broad GPU support, no extra toolkit) first and # falls back to the CPU build if Vulkan can't launch on this machine. gpu_backend: Literal["auto", "vulkan", "cpu"] = "auto" + # Minutes without a request after which the loaded model is released, so a + # 20 GB model does not keep its memory while the machine is used for + # something else. The next message loads it again. 0 keeps it loaded until + # Cortex exits or another model is chosen. + idle_unload_minutes: int = Field(default=30, ge=0, le=1440) + # Advanced runtime options appended to the launch (KV-cache types, flash + # attention, thread counts). Only an allow-list is accepted; the flags that + # define the launch -- model, context, address, key -- are refused. See + # cortex_backend.llamacpp.extra_args. + extra_args: tuple[str, ...] = () + + @field_validator("extra_args") + @classmethod + def _check_extra_args(cls, value: tuple[str, ...]) -> tuple[str, ...]: + return validate_extra_args(value) class GenerationSettings(_SettingsModel): diff --git a/backend/cortex_backend/llamacpp/binary_fetcher.py b/backend/cortex_backend/llamacpp/binary_fetcher.py index 47e4b93..e3064fb 100644 --- a/backend/cortex_backend/llamacpp/binary_fetcher.py +++ b/backend/cortex_backend/llamacpp/binary_fetcher.py @@ -11,6 +11,7 @@ import hashlib import logging import os +import re import shutil import zipfile from pathlib import Path @@ -53,6 +54,16 @@ _MAX_BINARY_DOWNLOAD_BYTES = 2 * 1024 * 1024 * 1024 # Mirrors download.py's ``MIN_FREE_SPACE_BYTES`` safety reserve. _MIN_FREE_SPACE_BYTES = 128 * 1024 * 1024 +# The only directories pruning ever considers: a release build this fetcher +# extracts (``-``, tag being llama.cpp's ``b``) and +# the ``.prune-`` name a superseded one is renamed to on its way out. In-progress +# ``.download-`` and ``.extract-`` names deliberately match neither. +_RELEASE_DIR_RE = re.compile(r"^b(\d+)-(?:cpu|vulkan)$") +_PRUNING_DIR_RE = re.compile(r"^\.prune-[0-9a-f]{32}$") +# One pruning pass removes at most this many directories; a runtime folder holds +# a handful, so reaching it means something else is going on and the rest waits +# for the next launch. +_MAX_PRUNED_PER_PASS = 8 @runtime_checkable @@ -67,6 +78,28 @@ def is_set(self) -> bool: ... +def _cancelled(cancellation_event: Cancellable | None) -> bool: + return cancellation_event is not None and cancellation_event.is_set() + + +def _directory_in_use(root: Path) -> bool: + """Whether a program is running from ``root``, or holds a file in it open. + + Opens each file for writing and closes it again without writing: a running + program's image and its loaded libraries refuse that, and so does a file + another process keeps open. Anything that cannot be listed or opened counts + as in use, because the caller's alternative is deleting it. + """ + try: + for path in root.rglob("*"): + if path.is_file(): + with path.open("r+b"): + pass + except OSError: + return True + return False + + def _raise_if_cancelled(cancellation_event: Cancellable | None) -> None: if cancellation_event is not None and cancellation_event.is_set(): raise BinaryVerificationError("Local model runtime startup was cancelled.") @@ -167,6 +200,7 @@ def __init__(self, runtime_dir: Path, *, http_client: httpx.Client | None = None self._runtime_dir = runtime_dir self._http = http_client self._verification_cache: dict[Path, tuple[_TreeIdentity, bool]] = {} + self._pruned_for: str | None = None def _verify_directory( self, @@ -248,6 +282,7 @@ def ensure_binary( if self._verify_directory( target_dir, asset, force_hash=True, cancellation_event=cancellation_event ): + self._prune_superseded_builds(release, cancellation_event) return exe_path self._runtime_dir.mkdir(parents=True, exist_ok=True) @@ -279,8 +314,78 @@ def ensure_binary( raise BinaryVerificationError( f"Downloaded llama.cpp binary for '{backend}' failed verification." ) + self._prune_superseded_builds(release, cancellation_event) return exe_path + def _prune_superseded_builds( + self, release: PinnedRelease, cancellation_event: Cancellable | None = None + ) -> None: + """Delete runtime builds older than the pinned release, once per process. + + Every pin bump used to leave a 100-200 MB build behind for good. Run only + after ``release`` itself has been verified, so there is always a working + runtime to fall back on, and never allowed to fail the launch that + triggered it: whatever cannot be removed now is tried again next time. + + Only a sibling named like one of this fetcher's own builds and with a + strictly lower build number is touched -- never the pinned release + (either backend), a newer one (a rolled-back Cortex sharing this data + folder), an in-progress download or extraction, or anything unrecognised. + A build a process is running from is left alone. Windows refuses to open + a running program's image, or a library it has loaded, for writing, so + every file is tried that way first (nothing is written); the directory + is then renamed, which is refused while any file in it is held open, and + only the renamed copy is deleted, so a half-deleted build is never left + under a name something could launch. Names are logged, never contents. + + Whatever cannot be removed now is tried again the next time Cortex starts. + """ + if self._pruned_for == release.tag: + return + self._pruned_for = release.tag + current = _RELEASE_DIR_RE.match(f"{release.tag}-cpu") + if current is None: + return + current_build = int(current.group(1)) + try: + root = self._runtime_dir.resolve() + candidates = sorted(self._runtime_dir.iterdir(), key=lambda entry: entry.name) + except OSError: + return + removed = 0 + for entry in candidates: + if removed >= _MAX_PRUNED_PER_PASS or _cancelled(cancellation_event): + break + try: + leftover = _PRUNING_DIR_RE.match(entry.name) is not None + match = _RELEASE_DIR_RE.match(entry.name) + if not leftover and (match is None or int(match.group(1)) >= current_build): + continue + # Only a real directory directly inside the runtime folder: a + # link or junction could lead anywhere, and is not ours to delete. + if entry.is_symlink() or not entry.is_dir() or entry.resolve().parent != root: + continue + if leftover: + shutil.rmtree(entry, ignore_errors=True) + else: + if _directory_in_use(entry): + logger.info("Kept an older local runtime build that is in use (%s).", entry.name) + continue + claimed = self._runtime_dir / f".prune-{uuid4().hex}" + try: + os.replace(entry, claimed) + except OSError: + logger.info("Kept an older local runtime build that is in use (%s).", entry.name) + continue + self._verification_cache.pop(entry, None) + shutil.rmtree(claimed, ignore_errors=True) + logger.info("Removed a superseded local runtime build (%s).", entry.name) + removed += 1 + except Exception as exc: + logger.warning( + "Could not remove a superseded local runtime build (%s).", type(exc).__name__ + ) + def _target_dir(self, release: PinnedRelease, backend: GpuBackend) -> Path: return self._runtime_dir / f"{release.tag}-{backend}" diff --git a/backend/cortex_backend/llamacpp/chat_client.py b/backend/cortex_backend/llamacpp/chat_client.py index cb8e843..d55f577 100644 --- a/backend/cortex_backend/llamacpp/chat_client.py +++ b/backend/cortex_backend/llamacpp/chat_client.py @@ -11,6 +11,7 @@ from __future__ import annotations +import functools import json import logging import ssl @@ -19,7 +20,7 @@ from contextlib import AbstractContextManager, nullcontext from pathlib import Path from threading import Event, Lock -from typing import Any +from typing import Any, Concatenate, ParamSpec, TypeVar import httpx @@ -36,6 +37,31 @@ # quickly is not worth waiting for before the real one. _TOKENIZE_TIMEOUT = httpx.Timeout(connect=2.0, read=10.0, write=10.0, pool=2.0) +_Params = ParamSpec("_Params") +_Result = TypeVar("_Result") + + +def _uses_the_server( + method: Callable[Concatenate[LlamaCppChatClient, _Params], _Result], +) -> Callable[Concatenate[LlamaCppChatClient, _Params], _Result]: + """Tell the provider the server is in use for the whole call. + + A generation can outlast the manager's idle period, and only this client + knows when its request begins and ends. Held from before the server is + made ready until the reply (or its failure) is over, so the server is + neither unloaded for being idle nor by a manual unload halfway through, and + the idle clock restarts when the call ends. A provider that does not track + use (a test double) is left alone. + """ + + @functools.wraps(method) + def wrapper(self: LlamaCppChatClient, *args: _Params.args, **kwargs: _Params.kwargs) -> _Result: + scope = getattr(self._provider, "request_scope", None) + with scope() if callable(scope) else nullcontext(): + return method(self, *args, **kwargs) + + return wrapper + class LlamaCppChatClient: """``ChatClient`` implementation backed by a locally-managed llama-server.""" @@ -126,6 +152,7 @@ def set_status_callback(self, callback: Callable[[str], None] | None) -> None: """ self._status_callback = callback + @_uses_the_server def chat( self, *, @@ -192,6 +219,7 @@ def chat( think=think, ) + @_uses_the_server def tokenize( self, *, diff --git a/backend/cortex_backend/llamacpp/errors.py b/backend/cortex_backend/llamacpp/errors.py index bb61301..4c3caf4 100644 --- a/backend/cortex_backend/llamacpp/errors.py +++ b/backend/cortex_backend/llamacpp/errors.py @@ -35,6 +35,14 @@ def user_message(self) -> str: return self.error +class RuntimeBusyError(LlamaCppError): + """Raised when the runtime cannot be unloaded because it is in use or loading. + + Its message is already user-facing (written by the manager, never taken + from the child), and the API reports it as a conflict. + """ + + class BinaryVerificationError(LlamaCppError): """Raised when a downloaded/cached llama-server binary fails verification.""" diff --git a/backend/cortex_backend/llamacpp/extra_args.py b/backend/cortex_backend/llamacpp/extra_args.py new file mode 100644 index 0000000..a10630b --- /dev/null +++ b/backend/cortex_backend/llamacpp/extra_args.py @@ -0,0 +1,129 @@ +"""The advanced ``llama-server`` options a user may add, and how they are checked. + +Cortex owns the launch contract: the model, the context window, the loopback +host, the ephemeral port, the API key (which travels in the environment), the +GPU-layer choice and the slot count. A user-chosen option list is appended +after that contract, so it is held to a short allow-list of tuning flags that +change how the model is run and nothing else -- the KV-cache element types, +flash attention and the CPU thread counts. Anything else is refused, and the +flags that would change what the contract guarantees are refused by name so the +error can say why. + +Only the shape of each option is checked here. The flag spellings are the +runtime's own and are passed through as written, so a value the pinned build +does not understand fails at launch like any other bad argument. + +Error messages name the position of the offending option and never repeat what +was typed: they travel back through the settings API. +""" + +from __future__ import annotations + +from collections.abc import Iterable +from dataclasses import dataclass + +MAX_EXTRA_ARGS = 12 +_MAX_ARG_CHARS = 32 +_MAX_THREADS = 1024 + +_CACHE_TYPES = frozenset( + {"f32", "f16", "bf16", "q8_0", "q4_0", "q4_1", "iq4_nl", "q5_0", "q5_1"} +) +_FLASH_ATTENTION_MODES = frozenset({"on", "off", "auto"}) + + +@dataclass(frozen=True, slots=True) +class _Option: + """One allowed flag: every spelling, and what may follow it.""" + + spellings: frozenset[str] + choices: frozenset[str] | None = None + # A whole number between 1 and this, instead of a choice. + maximum: int | None = None + # The value may be left out (older builds take ``-fa`` as a bare switch, + # newer ones ``-fa on``); a following token that is not a flag or a + # recognised value is still refused. + value_optional: bool = False + + +_OPTIONS = ( + _Option(frozenset({"-ctk", "--cache-type-k"}), choices=_CACHE_TYPES), + _Option(frozenset({"-ctv", "--cache-type-v"}), choices=_CACHE_TYPES), + _Option( + frozenset({"-fa", "--flash-attn"}), + choices=_FLASH_ATTENTION_MODES, + value_optional=True, + ), + _Option(frozenset({"-t", "--threads"}), maximum=_MAX_THREADS), + _Option(frozenset({"-tb", "--threads-batch"}), maximum=_MAX_THREADS), +) +_BY_SPELLING = {spelling: option for option in _OPTIONS for spelling in option.spellings} + +# Refused with an explanation: Cortex sets each of these itself, and letting one +# through would change the address the runtime listens on, how requests are +# authenticated, which model is loaded, or how much memory it is given. +_MANAGED_BY_CORTEX = frozenset( + { + "-m", "--model", "-mu", "--model-url", "-hf", "-hfr", "--hf-repo", "-hff", "--hf-file", + "-c", "--ctx-size", + "--host", "--port", "--path", "--api-prefix", "--reuse-port", + "--api-key", "--api-key-file", "--ssl-key-file", "--ssl-cert-file", + "-ngl", "--gpu-layers", "--n-gpu-layers", + "-np", "--parallel", + "--ui", "--no-ui", "--webui", "--no-webui", "--reasoning-format", + } +) + + +def _describe_problem(index: int, token: str, following: str | None) -> str | None: + """Why the token at ``index`` cannot stand, or ``None`` if it can.""" + place = f"option {index + 1}" + if token in _MANAGED_BY_CORTEX: + return f"{place} is set by Cortex and cannot be changed" + if token not in _BY_SPELLING: + return f"{place} is not an allowed runtime option" + option = _BY_SPELLING[token] + if following is None or following.startswith("-"): + if option.value_optional: + return None + return f"{place} needs a value" + if option.maximum is not None: + if not (following.isascii() and following.isdigit()) or not 1 <= int(following) <= option.maximum: + return f"the value after option {index + 1} must be a whole number from 1 to {option.maximum}" + return None + if option.choices is not None and following not in option.choices: + return f"the value after option {index + 1} is not one this runtime option accepts" + return None + + +def validate_extra_args(values: Iterable[str]) -> tuple[str, ...]: + """Return ``values`` as a tuple if every option is allowed, else raise ``ValueError``. + + A flag that takes a value is followed by exactly that value, each flag + appears at most once (under any of its spellings), and the list is short. + """ + tokens = tuple(values) + if len(tokens) > MAX_EXTRA_ARGS: + raise ValueError(f"at most {MAX_EXTRA_ARGS} runtime option words are allowed") + # Every word, value or flag, is one short plain word before anything is + # read into its meaning. + for position, word in enumerate(tokens, start=1): + if not isinstance(word, str) or not word or len(word) > _MAX_ARG_CHARS or not word.isascii(): + raise ValueError(f"option {position} is not a valid runtime option word") + if any(character.isspace() or not character.isprintable() for character in word): + raise ValueError(f"option {position} must be a single word") + seen: set[_Option] = set() + index = 0 + while index < len(tokens): + token = tokens[index] + following = tokens[index + 1] if index + 1 < len(tokens) else None + problem = _describe_problem(index, token, following) + if problem is not None: + raise ValueError(problem) + option = _BY_SPELLING[token] + if option in seen: + raise ValueError(f"option {index + 1} repeats a runtime option") + seen.add(option) + takes_value = following is not None and not following.startswith("-") + index += 2 if takes_value else 1 + return tokens diff --git a/backend/cortex_backend/llamacpp/launch_failure.py b/backend/cortex_backend/llamacpp/launch_failure.py index 2cda7b4..32e2ff9 100644 --- a/backend/cortex_backend/llamacpp/launch_failure.py +++ b/backend/cortex_backend/llamacpp/launch_failure.py @@ -207,6 +207,15 @@ def _rule(code: LaunchFailureCode, *patterns: str) -> tuple[LaunchFailureCode, r # made the file, so it must not be able to pick a cause. _ECHOED_MODEL_TEXT: Final = re.compile(r"\bkv\s+\d+\s*:|\bprint_info\s*:|\bgeneral\.[a-z_.]+\s*=", re.IGNORECASE) + +def echoes_model_text(line: str) -> bool: + """Whether ``line`` prints what the model file says about itself. + + Anything read from such a line is chosen by whoever made the file, so it + must not be used to report how the runtime behaved. + """ + return _ECHOED_MODEL_TEXT.search(line) is not None + # Windows exit codes (NTSTATUS) that identify the cause without any output, # which is how a process that cannot even load its libraries ends. Keys are the # unsigned 32-bit values; a signed exit code is masked to match. diff --git a/backend/cortex_backend/llamacpp/server_manager.py b/backend/cortex_backend/llamacpp/server_manager.py index 7f53c50..f6cc1bd 100644 --- a/backend/cortex_backend/llamacpp/server_manager.py +++ b/backend/cortex_backend/llamacpp/server_manager.py @@ -9,10 +9,13 @@ Lifecycle policy, stated explicitly because it is the whole point of this class: a loaded model stays resident until (a) a different model is -requested, (b) a larger context window is requested, (c) the app shuts -down, or (d) the process itself dies. Nothing here ever unloads a model -"between messages" -- if that appears to happen, one of those four causes -fired, and this class records which one (see ``last_restart_reason``). +requested, (b) a larger context window is requested, (c) different advanced +runtime options are requested, (d) the app shuts down, (e) the process itself +dies, (f) the user unloads it, or (g) it has sat unused for the configured idle +period (never while a request is in flight or a model is loading). Nothing +here ever unloads a model "between messages" for any other reason -- if that +appears to happen, one of those causes fired, and this class records which one +(see ``last_restart_reason``). This is a small, dedicated subprocess manager built directly on ``subprocess.Popen``. Running a binary Cortex itself downloaded and pinned @@ -23,6 +26,7 @@ from __future__ import annotations import ctypes +import ctypes.util from contextlib import contextmanager import json import logging @@ -37,7 +41,7 @@ from dataclasses import dataclass, field from pathlib import Path from typing import Any, Literal, Protocol -from collections.abc import Callable, Mapping +from collections.abc import Callable, Iterator, Mapping, Sequence import httpx @@ -57,13 +61,17 @@ BinaryVerificationError, CrashLoopError, LlamaCppError, + RuntimeBusyError, ServerLaunchError, ServerStartTimeoutError, ) +from .extra_args import validate_extra_args +from .gguf_metadata import read_gguf_metadata from .launch_failure import ( LaunchFailureCode, classify_child_exit, crash_loop_message, + echoes_model_text, launch_failure_message, ) @@ -120,12 +128,67 @@ # every message, short enough that a driver update or freed VRAM gets a # chance to matter within the same day rather than needing a manual reset. _KNOWN_BAD_BACKEND_TTL_SECONDS = 24.0 * 3600.0 +# How often the idle watcher looks at the clock. The setting is in whole +# minutes, so half a minute of slack is invisible, and a wake-up costs one +# comparison. +_IDLE_CHECK_INTERVAL_SECONDS = 30.0 +# A manual unload waits this long for the slow-path lock (a health +# re-verification holds it briefly) before it reports the runtime as busy; a +# model load holds it for minutes, and answering "busy" is the honest reply. +_UNLOAD_LOCK_TIMEOUT_SECONDS = 2.0 +# A trained context below this is not a plausible model, more likely a damaged +# or unusual header; it is treated as unknown rather than used to shrink the +# window to something unusable. +_MIN_TRAINED_CONTEXT = 256 +_UNLOADED_AT_REQUEST = "the model was unloaded at your request" +_NO_VULKAN_LOADER_NOTE = ( + "No Vulkan graphics loader was found on this computer, so the GPU build was " + "not downloaded and the CPU build is used. Install or update the graphics " + "driver to use the GPU." +) +_ADVANCED_OPTIONS_CHANGED = "the advanced runtime options changed" +_INVALID_ADVANCED_OPTIONS = ( + "The saved advanced runtime options are not valid. Fix or clear them in System settings." +) + + +def _idle_unload_reason(minutes: int) -> str: + return f"the model was unloaded after {minutes} minute{'' if minutes == 1 else 's'} without use" + + +def _vulkan_loader_present( + *, + platform: str = sys.platform, + environ: Mapping[str, str] = os.environ, + find_library: Callable[[str], str | None] = ctypes.util.find_library, +) -> bool: + """Whether this machine has a Vulkan loader the GPU build could load. + + The Vulkan build of llama.cpp links the loader (``vulkan-1.dll``, installed + with a graphics driver or the Vulkan runtime). Without it the roughly + 100 MB archive would be downloaded, launched, and only then found unusable. + Off Windows there is no such build to choose between, so the answer is yes + and behaviour is unchanged. + """ + if platform != "win32": + return True + system_root = environ.get("SystemRoot") or environ.get("WINDIR") + if system_root and (Path(system_root) / "System32" / "vulkan-1.dll").is_file(): + return True + try: + return find_library("vulkan-1") is not None + except OSError: + return False def _safe_restart_reason(reason: str) -> str: """Classify a restart without retaining model filenames or child text.""" if reason.startswith("the selected model changed"): return "the selected model changed" + if reason == _ADVANCED_OPTIONS_CHANGED: + return reason + if reason == _UNLOADED_AT_REQUEST or reason.startswith("the model was unloaded after "): + return reason if reason.startswith("the context window increased"): return reason if reason.startswith("the runtime process exited unexpectedly"): @@ -221,6 +284,19 @@ class LlamaCppRuntimeStatus: # said. None when nothing failed, when the cause was not identified, and # again once a server reaches ready. last_failure_code: LaunchFailureCode | None = None + # How many of the model's layers the running server put on the GPU, and how + # many it has, as the server itself reported them while loading. Both None + # when nothing is ready or when the server said nothing Cortex recognises: + # the build that launched (``active_backend``) does not say whether the GPU + # is in use, this does. 0 offloaded means the GPU build is running on the CPU. + gpu_layers_offloaded: int | None = None + gpu_layers_total: int | None = None + # Fixed text saying why the GPU build was not used when it would have been + # the default (no Vulkan loader on this machine). Never carries child text. + backend_note: str | None = None + # Fixed text saying the context window was limited to what the model was + # trained for. None when the window is as requested. + context_note: str | None = None @dataclass(frozen=True, slots=True) @@ -397,6 +473,28 @@ def _terminate_uncontained_process(process: subprocess.Popen) -> None: _LISTENING_PORT_RE = re.compile(r"\blistening on http://127\.0\.0\.1:(\d+)\b", re.IGNORECASE) +# The line llama.cpp prints once the model's layers are placed, for example +# "load_tensors: offloaded 24/33 layers to GPU". Only the two counts are kept. +# Anchored to the loader's own tag and to the end of the line, and never +# applied to a line that echoes the model file's text (see echoes_model_text). +_OFFLOADED_LAYERS_RE = re.compile( + r"\b(?:llm_)?load_tensors:\s*offloaded\s+(\d{1,5})\s*/\s*(\d{1,5})\s+layers?\s+to\s+GPU\s*$", + re.IGNORECASE, +) +_MAX_PARSED_LINE_CHARS = 512 + + +def _offloaded_layers(line: str) -> tuple[int, int] | None: + """The ``(offloaded, total)`` layer counts a loader line reports, or None.""" + if len(line) > _MAX_PARSED_LINE_CHARS or echoes_model_text(line): + return None + match = _OFFLOADED_LAYERS_RE.search(line) + if match is None: + return None + offloaded, total = int(match.group(1)), int(match.group(2)) + if total <= 0 or offloaded > total: + return None + return offloaded, total # llama-server gives every option an environment alias (LLAMA_ARG_*), and an # explicit argument only wins for the options Cortex actually passes. Anything @@ -524,6 +622,11 @@ def __init__( launcher: ProcessLauncher = default_launcher, http_client: httpx.Client | None = None, verify: ssl.SSLContext | bool = True, + idle_unload_minutes: Callable[[], int] | None = None, + extra_args: Callable[[], Sequence[str]] | None = None, + clock: Callable[[], float] = time.monotonic, + idle_check_interval_seconds: float = _IDLE_CHECK_INTERVAL_SECONDS, + vulkan_loader_probe: Callable[[], bool] = _vulkan_loader_present, ) -> None: self._runtime_dir = runtime_dir self._fetcher = fetcher @@ -548,6 +651,16 @@ def __init__( verify=verify, ) self._owns_http_client = http_client is None + # Read on every check rather than once, so a change in Settings applies + # without restarting Cortex. None means this manager never unloads on + # its own (and starts no watcher thread). + self._idle_unload_minutes = idle_unload_minutes + self._extra_args = extra_args + # Only the idle clock reads this; every deadline that guards a launch + # keeps using time.monotonic directly. A test moves it by hand. + self._clock = clock + self._idle_check_interval_seconds = idle_check_interval_seconds + self._vulkan_loader_probe = vulkan_loader_probe self._ensure_lock = threading.Lock() self._state_lock = threading.RLock() @@ -588,6 +701,21 @@ def __init__( self._preferred_backend_file = runtime_dir / "preferred_gpu_backend.json" self._api_key: str | None = None self._scrubbed_env_noted = False + # Requests that are using the running server right now (see + # request_scope), and when it was last used. Together they decide + # whether it may be unloaded for being idle. + self._active_requests = 0 + self._last_used = self._clock() + self._idle_stop = threading.Event() + self._idle_thread: threading.Thread | None = None + # The advanced options the most recent launch was given. While a server + # is ready these are what it is running with. + self._launch_extra_args: tuple[str, ...] = () + # The layer counts the running server reported (see _offloaded_layers). + self._gpu_layers: tuple[int, int] | None = None + self._backend_note: str | None = None + self._context_note: str | None = None + self._trained_context_cache: dict[Path, tuple[tuple[int, int], int | None]] = {} def close(self) -> None: """Stop the managed process and close an HTTP client owned here.""" @@ -599,6 +727,12 @@ def close(self) -> None: # starting another child after teardown has begun. with self._state_lock: self._closed = True + self._idle_stop.set() + watcher = self._idle_thread + if watcher is not None and watcher is not threading.current_thread(): + # It exits at its next wake-up; a teardown it is in the middle + # of is bounded, and stop() below waits for that lock too. + watcher.join(timeout=1.0) stop_error: Exception | None = None try: self.stop() @@ -639,6 +773,11 @@ def ensure_ready( forces a relaunch, since ``-c`` is a launch-time flag (unlike Ollama, where it's a per-request option). + A ``num_ctx`` above what the model was trained for is lowered to that + (see ``context_note`` on the status); the reuse decision and the launch + both use the lowered value, so asking for the same too-large window + again reuses the server instead of reloading it every time. + ``on_status`` is called with short, user-facing progress strings only while real work is happening (binary download, process start) -- an already-warm reused server never fires it, so no message flashes for @@ -650,7 +789,18 @@ def ensure_ready( with self._state_lock: if self._closed: raise LlamaCppError("The local model runtime manager is closed.") - verdict = self._reuse_verdict(model_path, num_ctx, token) + extra_args = self._current_extra_args() + trained_context = self._trained_context(model_path) + num_ctx, context_note = self._limit_to_trained_context(num_ctx, trained_context) + with self._state_lock: + if num_ctx is not None: + self._context_note = context_note + if extra_args != self._launch_extra_args: + # A changed option may be exactly what fixes a launch that + # kept failing; the new configuration starts with a clean slate. + self._failure_times.clear() + self._failure_key = None + verdict = self._reuse_verdict(model_path, num_ctx, token, extra_args=extra_args) if verdict.reusable: self._raise_if_stopping(token) with self._state_lock: @@ -658,7 +808,9 @@ def ensure_ready( raise LlamaCppError( "The local model runtime reported a reusable server with no address." ) - return ServerHandle(base_url=self._base_url, model_path=model_path, api_key=self._api_key) + handle = ServerHandle(base_url=self._base_url, model_path=model_path, api_key=self._api_key) + self._touch() + return handle with self._state_lock: effective_num_ctx = ( @@ -670,6 +822,10 @@ def ensure_ready( else _DEFAULT_NUM_CTX ) ) + if trained_context is not None: + # No preference was given, so a model trained for less than the + # default is still not asked for more than it knows. + effective_num_ctx = min(effective_num_ctx, trained_context) if verdict.reason is not None: self._record_restart(verdict, model_path, effective_num_ctx) @@ -680,8 +836,13 @@ def ensure_ready( raise LlamaCppError( "The previous local model runtime did not exit cleanly; restart Cortex before trying again." ) + if context_note is not None and on_status is not None: + on_status(context_note) try: - return self._start(model_path, effective_num_ctx, on_status, token) + handle = self._start(model_path, effective_num_ctx, extra_args, on_status, token) + self._touch() + self._ensure_idle_watcher() + return handle except LlamaCppError as exc: # A launch that never reaches "ready" -- the child exited # early (ServerLaunchError) or never answered its health @@ -731,6 +892,9 @@ def ready_handle(self, model_path: Path, *, num_ctx: int | None) -> ServerHandle tokens before sending it -- and is content with ``None`` when there is none. """ + # The same limit ensure_ready applies, or a window that will be lowered + # anyway would look too big for the server that is already running. + num_ctx, _ = self._limit_to_trained_context(num_ctx, self._trained_context(model_path)) with self._state_lock: if ( self._closed @@ -760,6 +924,9 @@ def status(self) -> LlamaCppRuntimeStatus: last_restart_reason = self._last_restart_reason active_backend = self._active_backend loaded_context = self._loaded_context if state == "ready" else None + gpu_layers = self._gpu_layers if state == "ready" else None + backend_note = self._backend_note + context_note = self._context_note # The expensive parts -- hashing the cached binary directory and a # settings read for the models folder -- run outside every lock, so # a status poll never stalls behind (or holds up) a model load. @@ -775,6 +942,187 @@ def status(self) -> LlamaCppRuntimeStatus: last_restart_reason=last_restart_reason, loaded_context=loaded_context, last_failure_code=last_failure_code, + gpu_layers_offloaded=gpu_layers[0] if gpu_layers is not None else None, + gpu_layers_total=gpu_layers[1] if gpu_layers is not None else None, + backend_note=backend_note, + context_note=context_note, + ) + + @contextmanager + def request_scope(self) -> Iterator[None]: + """Mark the server as in use for as long as the caller is talking to it. + + A generation can outlast any idle period, and the manager cannot see + the HTTP request the chat client makes, so the client says so. While + any scope is open the server is neither unloaded for being idle nor by + a manual unload, and the idle clock restarts when the last one closes. + """ + with self._state_lock: + self._active_requests += 1 + try: + yield + finally: + with self._state_lock: + self._active_requests -= 1 + self._last_used = self._clock() + + def unload(self) -> bool: + """Stop the loaded model now to free its memory; the next request loads it again. + + Returns True when a server was stopped and False when none was loaded + (an unload is safe to repeat). Raises :class:`RuntimeBusyError` when the + server is answering a request or a model is being loaded, and + :class:`LlamaCppError` when the process cannot be confirmed gone. + """ + if not self._ensure_lock.acquire(timeout=_UNLOAD_LOCK_TIMEOUT_SECONDS): + raise RuntimeBusyError( + "The model is being loaded or restarted. Try again when it has finished." + ) + try: + with self._state_lock: + if self._closed: + raise LlamaCppError("The local model runtime manager is closed.") + if self._active_requests > 0: + raise RuntimeBusyError( + "The model is answering a request. Stop it or wait for it to finish, then unload." + ) + if self._process is None: + return False + self._record_unload(_UNLOADED_AT_REQUEST) + if not self._terminate_and_reset(): + raise LlamaCppError( + "The local model runtime did not exit cleanly; restart Cortex before trying again." + ) + return True + finally: + self._ensure_lock.release() + + def unload_if_idle(self) -> bool: + """Unload the model if it has been unused for the configured idle period. + + Returns True when it did. Never waits behind a load or a restart (the + next check tries again), never unloads while a request is in flight, + and treats a setting it cannot read as "never". + """ + if self._idle_unload_minutes is None: + return False + try: + minutes = int(self._idle_unload_minutes()) + except Exception: + logger.debug("Could not read the idle-unload setting; not unloading.") + return False + if minutes <= 0: + return False + if not self._is_idle_for(minutes * 60.0): + return False + if not self._ensure_lock.acquire(blocking=False): + return False + try: + # Checked again now that the slow-path lock is held: a request that + # arrived in between has either bumped the clock or is waiting on + # this lock, and must not find its server gone. + if not self._is_idle_for(minutes * 60.0) or self._stop_event.is_set(): + return False + self._record_unload(_idle_unload_reason(minutes)) + return self._terminate_and_reset() + finally: + self._ensure_lock.release() + + def _is_idle_for(self, seconds: float) -> bool: + with self._state_lock: + return ( + not self._closed + and self._state == "ready" + and self._process is not None + and self._active_requests == 0 + and self._clock() - self._last_used >= seconds + ) + + def _record_unload(self, reason: str) -> None: + with self._state_lock: + self._last_restart_reason = reason + logger.info("Unloading the local model runtime (%s).", reason) + + def _touch(self) -> None: + with self._state_lock: + self._last_used = self._clock() + + def _ensure_idle_watcher(self) -> None: + """Start the thread that applies the idle period, once a model is loaded.""" + if self._idle_unload_minutes is None: + return + with self._state_lock: + if self._closed or (self._idle_thread is not None and self._idle_thread.is_alive()): + return + watcher = threading.Thread( + target=self._watch_for_idle, name="cortex-llama-idle-unload", daemon=True + ) + self._idle_thread = watcher + try: + watcher.start() + except Exception: + with self._state_lock: + self._idle_thread = None + logger.exception("Could not start the idle-unload watcher; the model stays loaded.") + + def _watch_for_idle(self) -> None: + while not self._idle_stop.wait(self._idle_check_interval_seconds): + try: + self.unload_if_idle() + except Exception: + logger.exception("The idle check for the local model runtime failed.") + + def _current_extra_args(self) -> tuple[str, ...]: + """The advanced options to launch with, re-checked against the allow-list.""" + if self._extra_args is None: + return () + try: + return validate_extra_args(self._extra_args()) + except ValueError as exc: + # Settings validate this on save, so this is a value that was + # injected or edited around them. Fail closed rather than launch + # with an option nobody approved. + raise LlamaCppError(_INVALID_ADVANCED_OPTIONS) from exc + + def _trained_context(self, model_path: Path) -> int | None: + """The context length the model file says it was trained for, if it says. + + Read from the GGUF header and remembered per file version, so a chat + message does not reopen the file. A file that cannot be read yields + None and is not retried until it changes. + """ + try: + stat = model_path.stat() + except OSError: + return None + stamp = (stat.st_size, stat.st_mtime_ns) + with self._state_lock: + cached = self._trained_context_cache.get(model_path) + if cached is not None and cached[0] == stamp: + return cached[1] + metadata = read_gguf_metadata(model_path) + length = metadata.context_length if metadata is not None else None + if length is not None and length < _MIN_TRAINED_CONTEXT: + length = None + with self._state_lock: + self._trained_context_cache[model_path] = (stamp, length) + return length + + @staticmethod + def _limit_to_trained_context( + num_ctx: int | None, trained_context: int | None + ) -> tuple[int | None, str | None]: + """``num_ctx`` lowered to the trained context, with the note that says so. + + A window past the trained length costs KV-cache memory for positions + the model never saw, and degrades its answers once a conversation + reaches them. + """ + if num_ctx is None or trained_context is None or num_ctx <= trained_context: + return num_ctx, None + return trained_context, ( + f"The context window was limited to {trained_context} tokens, the most this model " + f"was trained for ({num_ctx} were requested)." ) def stop(self) -> None: @@ -878,7 +1226,12 @@ def _ensure_guard(self, token: _CancellationToken): # -- reuse & teardown --------------------------------------------------- def _reuse_verdict( - self, model_path: Path, num_ctx: int | None, cancellation_event: _CancellationToken + self, + model_path: Path, + num_ctx: int | None, + cancellation_event: _CancellationToken, + *, + extra_args: tuple[str, ...] | None = None, ) -> _ReuseVerdict: with self._state_lock: if self._state != "ready" or self._process is None: @@ -900,6 +1253,8 @@ def _reuse_verdict( f"to {num_ctx} tokens" ), ) + if extra_args is not None and extra_args != self._launch_extra_args: + return _ReuseVerdict(reusable=False, reason=_ADVANCED_OPTIONS_CHANGED) exit_code = self._process.poll() if exit_code is not None: return _ReuseVerdict( @@ -1172,6 +1527,7 @@ def _start( self, model_path: Path, num_ctx: int, + extra_args: tuple[str, ...], on_status: StatusCallback | None, cancellation_event: _CancellationToken, ) -> ServerHandle: @@ -1181,6 +1537,11 @@ def _start( self._last_error = "The local GGUF runtime is not yet configured." raise LlamaCppError("The local GGUF runtime is not yet configured.") + with self._state_lock: + # What this launch is given, and what was decided about the GPU + # build for it: both describe the launch that is starting now. + self._launch_extra_args = extra_args + self._backend_note = None requested_backend = self._gpu_backend_setting() last_exc: Exception | None = None vulkan_launch_failed = False @@ -1197,7 +1558,7 @@ def _start( for backend in self._backend_order(requested_backend, model_path, num_ctx): try: handle = self._start_with_backend( - model_path, num_ctx, backend, on_status, cancellation_event + model_path, num_ctx, backend, extra_args, on_status, cancellation_event ) except (ServerLaunchError, BinaryVerificationError, OSError) as exc: # Not just launch failures. A backend whose archive fails its @@ -1276,7 +1637,17 @@ def _backend_order( if requested == "cpu": return ["cpu"] if requested == "vulkan": + # An explicit choice is honoured as it always was, including on a + # machine where the probe finds no loader: the probe can be wrong + # about an unusual install, and the user asked for this build. return ["vulkan"] + if not self._vulkan_loader_probe(): + # Nothing the GPU build can load is on this machine, so downloading + # and launching it would only end in a fallback to the CPU build + # after the larger download. Go straight there and say why. + with self._state_lock: + self._backend_note = _NO_VULKAN_LOADER_NOTE + return ["cpu"] if self._known_bad_backend(model_path, num_ctx) == "vulkan": return ["cpu"] return ["vulkan", "cpu"] @@ -1327,6 +1698,7 @@ def _start_with_backend( model_path: Path, num_ctx: int, backend: GpuBackend, + extra_args: tuple[str, ...], on_status: StatusCallback | None, cancellation_event: _CancellationToken, ) -> ServerHandle: @@ -1406,6 +1778,11 @@ def _start_with_backend( # spelling; ``--no-webui`` is the deprecated alias of the same # switch in the pinned build. "--no-ui", + # The user's advanced options go last. They were checked against an + # allow-list that excludes every flag above, so none of them can + # replace a launch setting, and a build that lets a later flag win + # has nothing to win against. + *extra_args, ] env, scrubbed = _child_environment(os.environ, api_key) self._note_scrubbed_environment(scrubbed) @@ -1434,6 +1811,9 @@ def _start_with_backend( self._starting_process = process stderr_tail: list[str] = [] listening_port: list[int] = [] + # The counts from the loader's offload line; the last one wins, so a + # build that prints one per attempt reports the load that was kept. + offloaded: list[tuple[int, int]] = [] listening_event = threading.Event() # When the child last wrote anything at all (see _drain_output). last_output = [time.monotonic()] @@ -1444,6 +1824,10 @@ def on_output(line: str) -> None: if match is not None and not listening_port: listening_port.append(int(match.group(1))) listening_event.set() + layers = _offloaded_layers(line) + if layers is not None and not listening_port: + # Only the load before "listening" describes this server. + offloaded.append(layers) def on_activity() -> None: last_output[0] = time.monotonic() @@ -1531,6 +1915,7 @@ def on_activity() -> None: self._last_error = None self._last_failure_code = None self._active_backend = backend + self._gpu_layers = offloaded[-1] if offloaded else None self._last_health_check = time.monotonic() self._stderr_tail = stderr_tail ready = True diff --git a/contracts/cortex-api.ts b/contracts/cortex-api.ts index 8d14008..90fd606 100644 --- a/contracts/cortex-api.ts +++ b/contracts/cortex-api.ts @@ -390,10 +390,16 @@ export interface LlamaCppRuntimeStatus { last_restart_reason?: string | null; loaded_context?: number | null; last_failure_code?: "unsupported_architecture" | "model_unreadable" | "memory" | "missing_shards" | "projector_not_a_model" | "no_gpu" | "port_unavailable" | "runtime_unusable" | "startup_timeout" | "health_check_failed" | "runtime_exited" | null; + gpu_layers_offloaded?: number | null; + gpu_layers_total?: number | null; + backend_note?: string | null; + context_note?: string | null; } export interface LlamaCppSettings { gpu_backend?: "auto" | "vulkan" | "cpu"; + idle_unload_minutes?: number; + extra_args?: Array; } export interface MemoryResponse { diff --git a/contracts/openapi.json b/contracts/openapi.json index b9e6307..176c934 100644 --- a/contracts/openapi.json +++ b/contracts/openapi.json @@ -2341,11 +2341,55 @@ ], "title": "Active Backend" }, + "backend_note": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Backend Note" + }, "binary_present": { "default": false, "title": "Binary Present", "type": "boolean" }, + "context_note": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Context Note" + }, + "gpu_layers_offloaded": { + "anyOf": [ + { + "type": "integer" + }, + { + "type": "null" + } + ], + "title": "Gpu Layers Offloaded" + }, + "gpu_layers_total": { + "anyOf": [ + { + "type": "integer" + }, + { + "type": "null" + } + ], + "title": "Gpu Layers Total" + }, "last_error": { "anyOf": [ { @@ -2444,6 +2488,14 @@ "LlamaCppSettings": { "additionalProperties": false, "properties": { + "extra_args": { + "default": [], + "items": { + "type": "string" + }, + "title": "Extra Args", + "type": "array" + }, "gpu_backend": { "default": "auto", "enum": [ @@ -2453,6 +2505,13 @@ ], "title": "Gpu Backend", "type": "string" + }, + "idle_unload_minutes": { + "default": 30, + "maximum": 1440.0, + "minimum": 0.0, + "title": "Idle Unload Minutes", + "type": "integer" } }, "title": "LlamaCppSettings", @@ -4999,6 +5058,30 @@ ] } }, + "/api/v1/llamacpp/unload": { + "post": { + "description": "Stop the loaded local model to free its memory; the next message loads it again.\n\nSafe to repeat: with nothing loaded it changes nothing and reports the\nstatus. Refused with 409 while a response is being generated or a model\nis loading, because that would cut the answer off.", + "operationId": "unload_llamacpp_api_v1_llamacpp_unload_post", + "responses": { + "200": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/LlamaCppRuntimeStatus" + } + } + }, + "description": "Successful Response" + } + }, + "security": [ + { + "CortexSession": [] + } + ], + "summary": "Unload Llamacpp" + } + }, "/api/v1/memories": { "get": { "operationId": "get_memories_api_v1_memories_get", diff --git a/frontend/src/api/client.test.ts b/frontend/src/api/client.test.ts index 73b61b3..762c00a 100644 --- a/frontend/src/api/client.test.ts +++ b/frontend/src/api/client.test.ts @@ -348,6 +348,39 @@ describe("CortexApi", () => { expect(new Headers(request.headers).get("Authorization")).toBe("Bearer session-1"); }); + it("unloads the local model with an authenticated POST and returns the new runtime status", async () => { + const fetcher = vi.fn().mockResolvedValue(new Response(JSON.stringify({ + state: "idle", + binary_present: true, + models_directory: "C:/synthetic/models", + last_restart_reason: "the model was unloaded at your request", + }), { status: 200, headers: { "Content-Type": "application/json" } })); + window.sessionStorage.setItem("cortex.session.token", "session-1"); + const api = new CortexApi("/api/v1", fetcher); + + const status = await api.unloadLlamaCpp(); + + expect(fetcher).toHaveBeenCalledWith("/api/v1/llamacpp/unload", expect.objectContaining({ method: "POST" })); + const request = fetcher.mock.calls[0]?.[1] as RequestInit; + expect(new Headers(request.headers).get("Authorization")).toBe("Bearer session-1"); + expect(status.state).toBe("idle"); + expect(status.last_restart_reason).toBe("the model was unloaded at your request"); + }); + + it("reports a refused unload with the backend's own sentence", async () => { + const fetcher = vi.fn().mockResolvedValue(new Response( + JSON.stringify({ detail: "A response is being generated. Stop it or wait for it to finish, then unload the model." }), + { status: 409, headers: { "Content-Type": "application/json" } }, + )); + window.sessionStorage.setItem("cortex.session.token", "session-1"); + const api = new CortexApi("/api/v1", fetcher); + + await expect(api.unloadLlamaCpp()).rejects.toMatchObject({ + status: 409, + detail: expect.stringContaining("being generated"), + }); + }); + it("starts a typed recipe request on the recipe route", async () => { const fetcher = vi.fn().mockResolvedValue(new Response(JSON.stringify({ job_id: "recipe-job", diff --git a/frontend/src/api/client.ts b/frontend/src/api/client.ts index 3c29553..cbd5f96 100644 --- a/frontend/src/api/client.ts +++ b/frontend/src/api/client.ts @@ -31,6 +31,7 @@ import type { JobAccepted, JobStatusResponse, HealthResponse, + LlamaCppRuntimeStatus, HandoffResponse, MemoryResponse, ModelDownloadRequest, @@ -327,6 +328,14 @@ export class CortexApi { return this.request("/system"); } + /** + * Stop the loaded local (GGUF) model to free its memory. The next message + * loads it again. Refused with 409 while a response is being generated. + */ + unloadLlamaCpp(): Promise { + return this.request("/llamacpp/unload", { method: "POST" }); + } + chats(): Promise { return this.request("/chats"); } diff --git a/frontend/src/app/App.runtimeUnload.test.tsx b/frontend/src/app/App.runtimeUnload.test.tsx new file mode 100644 index 0000000..3777188 --- /dev/null +++ b/frontend/src/app/App.runtimeUnload.test.tsx @@ -0,0 +1,132 @@ +import { render, screen, waitFor } from "@testing-library/react"; +import userEvent from "@testing-library/user-event"; +import { afterEach, beforeAll, describe, expect, it, vi } from "vitest"; +import { App } from "./App"; +import { CortexApi } from "../api/client"; +import { useModelStore } from "../stores/useModelStore"; +import { ToastProvider } from "./ToastProvider"; + +const respond = (body: unknown, status = 200) => new Response(JSON.stringify(body), { + status, + headers: { "Content-Type": "application/json" }, +}); + +const loadedStatus = { + state: "ready", + binary_present: true, + loaded_model: "gguf:demo.Q4_K_M.gguf", + models_directory: "C:\\models", + models_directory_exists: true, + active_backend: "cpu", +}; + +const unloadedStatus = { + state: "idle", + binary_present: true, + loaded_model: null, + models_directory: "C:\\models", + models_directory_exists: true, + active_backend: "cpu", + last_restart_reason: "the model was unloaded at your request", +}; + +/** A backend whose runtime status follows `runtime.current`, and whose unload answers as told. */ +function backend(unload: () => Response) { + const runtime = { current: loadedStatus as Record }; + const fetcher = vi.fn(async (input, init) => { + const url = String(input); + if (url.endsWith("/llamacpp/unload") && init?.method === "POST") { + const response = unload(); + if (response.ok) runtime.current = unloadedStatus; + return response; + } + if (url.endsWith("/system")) { + return respond({ + status: "ok", + preview: true, + session_required: true, + started_at: "2026-07-21T18:00:00Z", + llamacpp: runtime.current, + }); + } + if (url.endsWith("/chat-groups")) return respond([]); + if (url.endsWith("/chats")) return respond([]); + if (url.endsWith("/settings")) { + return respond({ settings: { models: { chat: "gguf:demo.Q4_K_M.gguf", title: null }, appearance: { theme: "dark" } } }); + } + if (url.endsWith("/memories")) return respond({ memos: [] }); + if (url.endsWith("/models")) { + return respond({ + required_models: [], + optional_models: [], + installed_models: ["gguf:demo.Q4_K_M.gguf"], + models: [{ name: "gguf:demo.Q4_K_M.gguf" }], + connection: { success: true, status: "connected", message: "Ready" }, + }); + } + return respond({ detail: "Unexpected test route." }, 404); + }); + return { fetcher, runtime }; +} + +async function openSystemSettings(user: ReturnType) { + expect(await screen.findByRole("heading", { name: "New thread" }, { timeout: 10_000 })).toBeVisible(); + await user.click(screen.getByRole("link", { name: "Settings" })); + await user.click(await screen.findByRole("button", { name: "System" }, { timeout: 10_000 })); +} + +describe("App local model unload", () => { + // The settings route is lazy; load it once so a test's own time is not spent on it. + beforeAll(async () => { + await Promise.all([ + import("../features/chat/ChatPage"), + import("../features/settings/SettingsPanel"), + ]); + }, 120_000); + + afterEach(() => { + useModelStore.getState().setLlamacppStatus(null); + window.sessionStorage.clear(); + window.history.replaceState({}, "", "/"); + }); + + it("unloads the model, shows the new state at once, and says what happened", async () => { + window.sessionStorage.setItem("cortex.session.token", "local-session"); + window.history.replaceState({}, "", "/chat/new"); + const { fetcher } = backend(() => respond(unloadedStatus)); + const user = userEvent.setup(); + render(); + await openSystemSettings(user); + const button = await screen.findByRole("button", { name: "Unload model" }); + await waitFor(() => expect(button).toBeEnabled()); + + await user.click(button); + + expect(await screen.findByText(/The local model was unloaded/)).toBeVisible(); + await waitFor(() => expect(screen.getByRole("button", { name: "Unload model" })).toBeDisabled()); + expect(screen.getByText("No local model is loaded right now.")).toBeVisible(); + const unloads = fetcher.mock.calls.filter(([input]) => String(input).endsWith("/llamacpp/unload")); + expect(unloads).toHaveLength(1); + expect(unloads[0]?.[1]).toMatchObject({ method: "POST" }); + }); + + it("reports the backend's reason when the unload is refused, and leaves the model as it was", async () => { + window.sessionStorage.setItem("cortex.session.token", "local-session"); + window.history.replaceState({}, "", "/chat/new"); + const { fetcher } = backend(() => respond( + { detail: "A response is being generated. Stop it or wait for it to finish, then unload the model." }, + 409, + )); + const user = userEvent.setup(); + render(); + await openSystemSettings(user); + const button = await screen.findByRole("button", { name: "Unload model" }); + await waitFor(() => expect(button).toBeEnabled()); + + await user.click(button); + + expect(await screen.findByText(/A response is being generated/)).toBeVisible(); + await waitFor(() => expect(screen.getByRole("button", { name: "Unload model" })).toBeEnabled()); + expect(useModelStore.getState().llamacppStatus?.state).toBe("ready"); + }); +}); diff --git a/frontend/src/app/App.tsx b/frontend/src/app/App.tsx index 515b634..c9a9ca5 100644 --- a/frontend/src/app/App.tsx +++ b/frontend/src/app/App.tsx @@ -805,6 +805,17 @@ function AuthenticatedWorkspace({ api, onSessionExpired }: { api: CortexApi; onS await finishGGUFDownload(filename); }; + const unloadLocalModel = async () => { + try { + // The answer is the new status, so the panel reflects it at once rather + // than at the next poll (which only runs while a GGUF model is selected). + setLlamacppStatus(await api.unloadLlamaCpp()); + notify("The local model was unloaded. It loads again when you send a message.", "success"); + } catch (error) { + notify(apiMessage(error, "Could not unload the local model."), "error"); + } + }; + const chooseLocalModel = async (model: string): Promise => { // Read the store rather than this render's `settings`: a GGUF download // runs for minutes before selecting what it fetched, and Settings stays @@ -881,7 +892,7 @@ function AuthenticatedWorkspace({ api, onSessionExpired }: { api: CortexApi; onS Loading workspace...}> {route.kind === "settings" - ? + ? : } diff --git a/frontend/src/features/models/ModelsPanel.runtime.test.tsx b/frontend/src/features/models/ModelsPanel.runtime.test.tsx new file mode 100644 index 0000000..8c9e306 --- /dev/null +++ b/frontend/src/features/models/ModelsPanel.runtime.test.tsx @@ -0,0 +1,283 @@ +import { render, screen, waitFor } from "@testing-library/react"; +import userEvent from "@testing-library/user-event"; +import { describe, expect, it, vi } from "vitest"; +import type { LlamaCppRuntimeStatus, ModelResponse } from "../../../../contracts/cortex-api"; +import { ModelsPanel, type RuntimeControls } from "./ModelsPanel"; + +const models: ModelResponse = { + required_models: [], + optional_models: [], + installed_models: [], + models: [], + connection: { success: true, status: "connected", message: "Connected." }, +}; + +const idleStatus: LlamaCppRuntimeStatus = { + state: "idle", + binary_present: true, + loaded_model: null, + last_error: null, + models_directory: "C:\\synthetic\\models", +}; + +const readyStatus: LlamaCppRuntimeStatus = { + ...idleStatus, + state: "ready", + loaded_model: "gguf:synthetic.Q4_K_M.gguf", + active_backend: "vulkan", +}; + +function runtimeControls(overrides: Partial = {}): RuntimeControls { + return { + idleUnloadMinutes: 30, + onIdleUnloadMinutesChange: vi.fn(), + extraArgs: [], + onExtraArgsChange: vi.fn(), + onUnload: vi.fn().mockResolvedValue(undefined), + ...overrides, + }; +} + +function renderPanel(status: LlamaCppRuntimeStatus, runtime?: RuntimeControls) { + const props = { + models, + busy: false, + progress: null, + setupUrl: "https://ollama.com/download", + onCheck: vi.fn().mockResolvedValue(undefined), + gguf: { + directory: "", + directoryDirty: false, + onDirectoryChange: vi.fn(), + onDownload: vi.fn().mockResolvedValue(undefined), + busy: false, + }, + }; + const view = render(); + return { + ...view, + rerenderWith: (nextStatus: LlamaCppRuntimeStatus, nextRuntime?: RuntimeControls) => + view.rerender(), + }; +} + +describe("ModelsPanel local runtime line", () => { + it("says how many layers are on the GPU when the runtime reported them", () => { + renderPanel({ ...readyStatus, gpu_layers_offloaded: 24, gpu_layers_total: 33 }); + + expect(screen.getByText(/Local runtime:/)).toHaveTextContent("GPU (Vulkan) · 24/33 layers on the GPU"); + }); + + it("does not call a GPU build that offloaded nothing a GPU", () => { + renderPanel({ ...readyStatus, gpu_layers_offloaded: 0, gpu_layers_total: 33 }); + + const line = screen.getByText(/Local runtime:/); + expect(line).toHaveTextContent("none of the 33 layers are on the GPU"); + expect(line).not.toHaveTextContent("GPU (Vulkan)"); + }); + + it("admits it when the runtime did not say how much is on the GPU", () => { + renderPanel(readyStatus); + + const line = screen.getByText(/Local runtime:/); + expect(line).toHaveTextContent("GPU (Vulkan)"); + expect(line).toHaveTextContent("was not reported"); + expect(line).not.toHaveTextContent("layers on the GPU"); + }); + + it("names the build only, with no layer claim, while the model is still loading", () => { + renderPanel({ ...readyStatus, state: "starting", gpu_layers_offloaded: 24, gpu_layers_total: 33 }); + + const line = screen.getByText(/Local runtime:/); + expect(line).toHaveTextContent("GPU (Vulkan)"); + expect(line).not.toHaveTextContent("layers"); + }); + + it("reports the CPU build as the CPU", () => { + renderPanel({ ...readyStatus, active_backend: "cpu", gpu_layers_offloaded: 0, gpu_layers_total: 33 }); + + expect(screen.getByText(/Local runtime:/)).toHaveTextContent(/Local runtime:\s*CPU\b/); + }); + +}); + +describe("ModelsPanel GPU skip note", () => { + it("shows why the GPU build was skipped", () => { + renderPanel({ + ...readyStatus, + active_backend: "cpu", + backend_note: "No Vulkan graphics loader was found on this computer.", + }); + + expect(screen.getByText("No Vulkan graphics loader was found on this computer.")).toBeVisible(); + }); + + it("shows no skip note when the GPU build was not skipped", () => { + renderPanel(readyStatus); + + expect(screen.queryByText(/No Vulkan graphics loader/)).not.toBeInTheDocument(); + }); + +}); + +describe("ModelsPanel context note", () => { + it("shows that the context window was limited", () => { + renderPanel({ ...readyStatus, context_note: "The context window was limited to 4096 tokens." }); + + expect(screen.getByText("The context window was limited to 4096 tokens.")).toBeVisible(); + }); + + it("shows no context note when the window is as requested", () => { + renderPanel(readyStatus); + + expect(screen.queryByText(/context window was limited/)).not.toBeInTheDocument(); + }); +}); + +describe("ModelsPanel unload control", () => { + it("unloads the loaded model when asked", async () => { + const user = userEvent.setup(); + const runtime = runtimeControls(); + renderPanel(readyStatus, runtime); + + await user.click(screen.getByRole("button", { name: "Unload model" })); + + expect(runtime.onUnload).toHaveBeenCalledTimes(1); + }); + + it("is disabled, and says so, when no model is loaded", () => { + const runtime = runtimeControls(); + renderPanel(idleStatus, runtime); + + expect(screen.getByRole("button", { name: "Unload model" })).toBeDisabled(); + expect(screen.getByText("No local model is loaded right now.")).toBeVisible(); + }); + + it("is disabled while a model is still loading", () => { + renderPanel({ ...readyStatus, state: "starting", loaded_model: null }, runtimeControls()); + + expect(screen.getByRole("button", { name: "Unload model" })).toBeDisabled(); + }); + + it("shows progress and cannot be pressed twice while the unload is in flight", async () => { + const user = userEvent.setup(); + let finish: () => void = () => {}; + const onUnload = vi.fn(() => new Promise((resolve) => { finish = resolve; })); + renderPanel(readyStatus, runtimeControls({ onUnload })); + + await user.click(screen.getByRole("button", { name: "Unload model" })); + + const busy = await screen.findByRole("button", { name: "Unloading…" }); + expect(busy).toBeDisabled(); + await user.click(busy); + expect(onUnload).toHaveBeenCalledTimes(1); + + finish(); + await waitFor(() => expect(screen.getByRole("button", { name: "Unload model" })).toBeEnabled()); + }); + + it("is left out where the build cannot unload a model", () => { + renderPanel(readyStatus, runtimeControls({ onUnload: undefined })); + + expect(screen.queryByRole("button", { name: "Unload model" })).not.toBeInTheDocument(); + expect(screen.getByLabelText(/Unload an unused model after/)).toBeVisible(); + }); + + it("shows no runtime controls when none are offered", () => { + renderPanel(readyStatus); + + expect(screen.queryByRole("button", { name: "Unload model" })).not.toBeInTheDocument(); + expect(screen.queryByLabelText(/Unload an unused model after/)).not.toBeInTheDocument(); + expect(screen.queryByLabelText("Advanced runtime options")).not.toBeInTheDocument(); + }); +}); + +describe("ModelsPanel idle period", () => { + it("shows the current period and reports a new one as it is typed", async () => { + const user = userEvent.setup(); + const runtime = runtimeControls({ idleUnloadMinutes: 30 }); + renderPanel(idleStatus, runtime); + const field = screen.getByLabelText(/Unload an unused model after/); + expect(field).toHaveValue(30); + + await user.clear(field); + await user.type(field, "45"); + + expect(runtime.onIdleUnloadMinutesChange).toHaveBeenLastCalledWith(45); + }); + + it("accepts 0 as never", async () => { + const user = userEvent.setup(); + const runtime = runtimeControls({ idleUnloadMinutes: 30 }); + renderPanel(idleStatus, runtime); + const field = screen.getByLabelText(/Unload an unused model after/); + + await user.clear(field); + await user.type(field, "0"); + + expect(runtime.onIdleUnloadMinutesChange).toHaveBeenLastCalledWith(0); + expect(screen.getByText(/0 keeps the model loaded until Cortex closes/)).toBeVisible(); + }); + + it("does not report text that is not a period from 0 to a day", async () => { + const user = userEvent.setup(); + const runtime = runtimeControls({ idleUnloadMinutes: 30 }); + renderPanel(idleStatus, runtime); + const field = screen.getByLabelText(/Unload an unused model after/); + + await user.clear(field); + await user.type(field, "1441"); + + expect(runtime.onIdleUnloadMinutesChange).not.toHaveBeenCalledWith(1441); + await user.clear(field); + await user.type(field, "-5"); + expect(runtime.onIdleUnloadMinutesChange).not.toHaveBeenCalledWith(-5); + }); + + it("follows a period that changed elsewhere", () => { + const view = renderPanel(idleStatus, runtimeControls({ idleUnloadMinutes: 30 })); + + view.rerenderWith(idleStatus, runtimeControls({ idleUnloadMinutes: 5 })); + + expect(screen.getByLabelText(/Unload an unused model after/)).toHaveValue(5); + }); +}); + +describe("ModelsPanel advanced options", () => { + it("reports the words typed, keeping the spaces while they are typed", async () => { + const user = userEvent.setup(); + const runtime = runtimeControls(); + renderPanel(idleStatus, runtime); + const field = screen.getByLabelText("Advanced runtime options"); + + await user.type(field, "-ctk q8_0 -t 8"); + + expect(field).toHaveValue("-ctk q8_0 -t 8"); + expect(runtime.onExtraArgsChange).toHaveBeenLastCalledWith(["-ctk", "q8_0", "-t", "8"]); + }); + + it("shows the options already set, and follows a change made elsewhere", () => { + const view = renderPanel(idleStatus, runtimeControls({ extraArgs: ["-fa", "on"] })); + expect(screen.getByLabelText("Advanced runtime options")).toHaveValue("-fa on"); + + view.rerenderWith(idleStatus, runtimeControls({ extraArgs: [] })); + + expect(screen.getByLabelText("Advanced runtime options")).toHaveValue(""); + }); + + it("clearing the field reports no options", async () => { + const user = userEvent.setup(); + const runtime = runtimeControls({ extraArgs: ["-t", "8"] }); + renderPanel(idleStatus, runtime); + + await user.clear(screen.getByLabelText("Advanced runtime options")); + + expect(runtime.onExtraArgsChange).toHaveBeenLastCalledWith([]); + }); + + it("says what may be set and that the launch contract stays with Cortex", () => { + renderPanel(idleStatus, runtimeControls()); + + expect(screen.getByText(/KV-cache types/)).toHaveTextContent("Cortex keeps setting the model, context window and address itself"); + }); +}); diff --git a/frontend/src/features/models/ModelsPanel.tsx b/frontend/src/features/models/ModelsPanel.tsx index 28353ad..9ade049 100644 --- a/frontend/src/features/models/ModelsPanel.tsx +++ b/frontend/src/features/models/ModelsPanel.tsx @@ -1,4 +1,5 @@ -import { ExternalLink, FolderOpen, RefreshCw, X } from "lucide-react"; +import { ExternalLink, FolderOpen, PowerOff, RefreshCw, X } from "lucide-react"; +import { useState } from "react"; import type { LlamaCppRuntimeStatus, ModelDownloadRequest, @@ -26,6 +27,18 @@ type GGUFControls = { onListFiles?: ListGGUFFiles; }; +/** How the local runtime uses memory, and the advanced options it is launched with. */ +export type RuntimeControls = { + /** Minutes without a request after which the loaded model is released; 0 keeps it loaded. */ + idleUnloadMinutes: number; + onIdleUnloadMinutesChange: (minutes: number) => void; + /** Advanced runtime options as words, e.g. `["-ctk", "q8_0", "-t", "8"]`. */ + extraArgs: readonly string[]; + onExtraArgsChange: (args: string[]) => void; + /** Absent when this build cannot unload a model; the button is then left out. */ + onUnload?: () => Promise; +}; + type Props = { models: ModelResponse; busy: boolean; @@ -36,9 +49,27 @@ type Props = { onCheck: () => Promise; llamacppStatus: LlamaCppRuntimeStatus; gguf: GGUFControls; + runtime?: RuntimeControls; }; -export function ModelsPanel({ models, busy, progress, onCancel, setupUrl, onCheck, llamacppStatus, gguf }: Props) { +/** + * What the status can honestly say about which processor is doing the work. + * The build that launched is not the same as the GPU being used: with automatic + * layer offload only some (or none) of the model may be on it, so the layer + * counts the runtime reported win over the name of the build. + */ +function describeRuntimeBackend(status: LlamaCppRuntimeStatus): string { + if (status.active_backend !== "vulkan") return "CPU"; + const { gpu_layers_offloaded: onGpu, gpu_layers_total: total } = status; + if (status.state !== "ready") return "GPU (Vulkan)"; + if (typeof onGpu !== "number" || typeof total !== "number") { + return "GPU (Vulkan) — how much of the model is on the GPU was not reported"; + } + if (onGpu === 0) return `CPU — the GPU build is running, but none of the ${total} layers are on the GPU`; + return `GPU (Vulkan) · ${onGpu}/${total} layers on the GPU`; +} + +export function ModelsPanel({ models, busy, progress, onCancel, setupUrl, onCheck, llamacppStatus, gguf, runtime }: Props) { const connection = models.connection; const missing = models.missing_models ?? []; const optionalMissing = models.optional_missing_models ?? []; @@ -113,12 +144,12 @@ export function ModelsPanel({ models, busy, progress, onCancel, setupUrl, onChec )} - + ); } -function GGUFRuntimeSection({ llamacppStatus, gguf }: { llamacppStatus: LlamaCppRuntimeStatus; gguf: GGUFControls }) { +function GGUFRuntimeSection({ llamacppStatus, gguf, runtime }: { llamacppStatus: LlamaCppRuntimeStatus; gguf: GGUFControls; runtime?: RuntimeControls }) { return (
@@ -131,7 +162,7 @@ function GGUFRuntimeSection({ llamacppStatus, gguf }: { llamacppStatus: LlamaCpp

{llamacppStatus.active_backend && (

- Local runtime: {llamacppStatus.active_backend === "vulkan" ? "GPU (Vulkan)" : "CPU"} + Local runtime: {describeRuntimeBackend(llamacppStatus)} {llamacppStatus.state === "ready" && llamacppStatus.loaded_model ? ` — currently running ${displayModelName(llamacppStatus.loaded_model)}` : llamacppStatus.state === "starting" || llamacppStatus.state === "downloading_binary" @@ -139,6 +170,9 @@ function GGUFRuntimeSection({ llamacppStatus, gguf }: { llamacppStatus: LlamaCpp : ""}

)} + {llamacppStatus.backend_note &&

{llamacppStatus.backend_note}

} + {llamacppStatus.context_note &&

{llamacppStatus.context_note}

} + {runtime && }
diff --git a/frontend/src/styles/tokens.css b/frontend/src/styles/tokens.css index 5c3791f..be8b574 100644 --- a/frontend/src/styles/tokens.css +++ b/frontend/src/styles/tokens.css @@ -949,6 +949,10 @@ input[type="range"] { width: 100%; accent-color: var(--accent); } .gguf-runtime-directory { display: flex; flex-wrap: wrap; align-items: center; gap: 8px; font-size: 0.78rem; color: var(--text-muted); } .gguf-runtime-directory input { flex: 1 1 260px; min-width: 0; } .gguf-runtime-directory-hint { display: block; margin-top: -4px; color: var(--text-faint); font-size: 0.72rem; } +.gguf-runtime-memory { display: grid; gap: 12px; } +.gguf-runtime-field { display: grid; gap: 6px; } +.gguf-runtime-unload { display: flex; flex-wrap: wrap; align-items: center; gap: 10px; } +.gguf-runtime-unload .gguf-runtime-directory-hint { margin-top: 0; } .gguf-download-form { display: grid; gap: 10px; } .gguf-download-source-toggle { display: inline-flex; gap: 6px; } .gguf-download-fields { display: grid; gap: 8px; grid-template-columns: repeat(auto-fit, minmax(180px, 1fr)); } diff --git a/tests/test_chat_client_routing.py b/tests/test_chat_client_routing.py index bd189ea..afce0e5 100644 --- a/tests/test_chat_client_routing.py +++ b/tests/test_chat_client_routing.py @@ -2191,3 +2191,150 @@ def handler(request: httpx.Request) -> httpx.Response: assert response["prompt_eval_count"] == 321 assert response["prompt_token_count"] == 321 + + +# --------------------------------------------------------------------------- +# The chat client tells the provider when the server is in use (RT-07) +# --------------------------------------------------------------------------- + + +class _ScopedProvider(_StaticProvider): + """Tracks use the way the real manager does, and records the order of events.""" + + def __init__(self, base_url: str, events: list[str]) -> None: + super().__init__(base_url) + self.events = events + + @contextmanager + def request_scope(self): + self.events.append("scope open") + try: + yield + finally: + self.events.append("scope closed") + + def ensure_ready(self, model_path: Path, *, num_ctx, on_status=None, cancellation_event=None) -> ServerHandle: + self.events.append("ensure_ready") + return super().ensure_ready( + model_path, num_ctx=num_ctx, on_status=on_status, cancellation_event=cancellation_event + ) + + def ready_handle(self, model_path: Path, *, num_ctx) -> ServerHandle: + del num_ctx + self.events.append("ready_handle") + return ServerHandle(base_url=self._base_url, model_path=model_path) + + +def _recording_http(events: list[str], response: httpx.Response) -> httpx.Client: + def handler(request: httpx.Request) -> httpx.Response: + events.append(f"http {request.url.path}") + return response + + return httpx.Client(transport=httpx.MockTransport(handler)) + + +_STREAMED_ANSWER = httpx.Response( + 200, + content=b'data: {"choices": [{"delta": {"content": "hi"}}]}\n\ndata: [DONE]\n\n', + headers={"content-type": "text/event-stream"}, +) + + +def test_the_server_is_marked_in_use_from_before_it_is_readied_until_the_reply_is_done(tmp_path: Path) -> None: + (tmp_path / "tiny.gguf").write_bytes(b"fake") + events: list[str] = [] + http_client = _recording_http(events, httpx.Response(200, json={"choices": [{"message": {"content": "ok"}}]})) + client = LlamaCppChatClient( + _ScopedProvider("http://fakellama", events), models_directory=lambda: tmp_path, http_client=http_client + ) + + client.chat(model="gguf:tiny.gguf", messages=[{"role": "user", "content": "hi"}], options={}) + + assert events == ["scope open", "ensure_ready", "http /v1/chat/completions", "scope closed"] + + +def test_a_streamed_reply_holds_the_scope_until_the_last_chunk(tmp_path: Path) -> None: + (tmp_path / "tiny.gguf").write_bytes(b"fake") + events: list[str] = [] + http_client = _recording_http(events, _STREAMED_ANSWER) + client = LlamaCppChatClient( + _ScopedProvider("http://fakellama", events), models_directory=lambda: tmp_path, http_client=http_client + ) + + client.chat( + model="gguf:tiny.gguf", + messages=[{"role": "user", "content": "hi"}], + options={}, + cancellation_event=Event(), + on_delta=lambda kind, text: events.append(f"delta {kind} {text}"), + ) + + assert events == [ + "scope open", + "ensure_ready", + "http /v1/chat/completions", + "delta content hi", + "scope closed", + ] + + +def test_the_scope_is_closed_when_the_server_fails(tmp_path: Path) -> None: + (tmp_path / "tiny.gguf").write_bytes(b"fake") + events: list[str] = [] + http_client = _recording_http(events, httpx.Response(500, json={"error": {"message": "boom"}})) + client = LlamaCppChatClient( + _ScopedProvider("http://fakellama", events), models_directory=lambda: tmp_path, http_client=http_client + ) + + with pytest.raises(LlamaCppError): + client.chat(model="gguf:tiny.gguf", messages=[{"role": "user", "content": "hi"}], options={}) + + assert events[0] == "scope open" + assert events[-1] == "scope closed" + assert events.count("scope closed") == 1 + + +def test_the_scope_is_closed_when_the_runtime_cannot_be_started(tmp_path: Path) -> None: + (tmp_path / "tiny.gguf").write_bytes(b"fake") + events: list[str] = [] + + class _FailingProvider(_ScopedProvider): + def ensure_ready(self, model_path: Path, *, num_ctx, on_status=None, cancellation_event=None): + self.events.append("ensure_ready") + raise LlamaCppError("The local model runtime could not start.") + + client = LlamaCppChatClient( + _FailingProvider("http://fakellama", events), + models_directory=lambda: tmp_path, + http_client=_recording_http(events, httpx.Response(200, json={})), + ) + + with pytest.raises(LlamaCppError): + client.chat(model="gguf:tiny.gguf", messages=[{"role": "user", "content": "hi"}], options={}) + + assert events == ["scope open", "ensure_ready", "scope closed"] + + +def test_counting_tokens_also_counts_as_use(tmp_path: Path) -> None: + (tmp_path / "tiny.gguf").write_bytes(b"fake") + events: list[str] = [] + http_client = _recording_http(events, httpx.Response(200, json={"tokens": [1, 2, 3]})) + client = LlamaCppChatClient( + _ScopedProvider("http://fakellama", events), models_directory=lambda: tmp_path, http_client=http_client + ) + + assert client.tokenize(model="gguf:tiny.gguf", text="some text", options={}) == 3 + + assert events == ["scope open", "ready_handle", "http /tokenize", "scope closed"] + + +def test_a_provider_that_does_not_track_use_is_left_alone(tmp_path: Path) -> None: + (tmp_path / "tiny.gguf").write_bytes(b"fake") + http_client = httpx.Client( + transport=httpx.MockTransport(lambda request: httpx.Response(200, json={"choices": [{"message": {"content": "ok"}}]})) + ) + client = LlamaCppChatClient(_StaticProvider("http://fakellama"), models_directory=lambda: tmp_path, http_client=http_client) + + result = client.chat(model="gguf:tiny.gguf", messages=[{"role": "user", "content": "hi"}], options={}) + + assert result["message"]["content"] == "ok" diff --git a/tests/test_llamacpp_binary_fetcher.py b/tests/test_llamacpp_binary_fetcher.py index e3a3f86..7eddb25 100644 --- a/tests/test_llamacpp_binary_fetcher.py +++ b/tests/test_llamacpp_binary_fetcher.py @@ -9,6 +9,8 @@ import hashlib import io import os +import shutil +import sys import tempfile import time import zipfile @@ -57,9 +59,9 @@ def _expected_directory_hash() -> str: return hash_directory(extract_dir) -def _release_for(archive_bytes: bytes, *, filename: str = "llama-cpu.zip") -> PinnedRelease: +def _release_for(archive_bytes: bytes, *, filename: str = "llama-cpu.zip", tag: str = "b0001") -> PinnedRelease: return PinnedRelease( - tag="b0001", + tag=tag, assets={ "cpu": AssetSpec( filename=filename, @@ -391,3 +393,204 @@ def handler(request: httpx.Request) -> httpx.Response: "cpu", cancellation_event=SimpleNamespace(is_set=lambda: stopped["value"]), ) + + +# --------------------------------------------------------------------------- +# Superseded runtime builds are removed (RT-20) +# --------------------------------------------------------------------------- + + +def _older_build(root: Path, name: str) -> Path: + """A build directory laid out like an extracted release, under ``root``.""" + directory = root / name + (directory / "nested").mkdir(parents=True) + (directory / "llama-server.exe").write_bytes(b"older stub") + (directory / "llama-server-impl.dll").write_bytes(b"older impl") + (directory / "nested" / "ggml-extra.dll").write_bytes(b"older backend") + return directory + + +def _fetch_current(root: Path, *, tag: str = "b0100") -> tuple[BinaryFetcher, PinnedRelease]: + archive_bytes = _build_archive() + release = _release_for(archive_bytes, tag=tag) + fetcher = BinaryFetcher(root, http_client=_client_returning(archive_bytes)) + fetcher.ensure_binary(release, "cpu") + return fetcher, release + + +def test_old_release_directories_are_pruned(tmp_path: Path) -> None: + older_cpu = _older_build(tmp_path, "b0050-cpu") + older_vulkan = _older_build(tmp_path, "b0050-vulkan") + oldest = _older_build(tmp_path, "b0007-cpu") + same_tag_other_backend = _older_build(tmp_path, "b0100-vulkan") + newer = _older_build(tmp_path, "b0200-cpu") + # Things that are not superseded builds and must never be touched. + in_progress_download = tmp_path / ".download-0123456789abcdef" + in_progress_download.write_bytes(b"partial archive") + in_progress_extract = tmp_path / ".extract-0123456789abcdef" + (in_progress_extract / "half").mkdir(parents=True) + marker = tmp_path / "preferred_gpu_backend.json" + marker.write_text("{}", encoding="utf-8") + unrelated = tmp_path / "notes" + unrelated.mkdir() + (unrelated / "keep.txt").write_text("keep", encoding="utf-8") + a_file_named_like_a_build = tmp_path / "b0040-cpu" + a_file_named_like_a_build.write_bytes(b"not a directory") + + fetcher, release = _fetch_current(tmp_path) + + assert not older_cpu.exists() + assert not older_vulkan.exists() + assert not oldest.exists() + # The pinned release, on both backends, and anything newer than it. + assert (tmp_path / "b0100-cpu" / "llama-server.exe").is_file() + assert same_tag_other_backend.is_dir() + assert newer.is_dir() + assert in_progress_download.is_file() + assert in_progress_extract.is_dir() + assert marker.is_file() + assert (unrelated / "keep.txt").is_file() + assert a_file_named_like_a_build.is_file() + # No half-deleted remains are left under a launchable name, and the runtime + # that was just fetched still verifies. + assert not any(entry.name.startswith(".prune-") for entry in tmp_path.iterdir()) + assert fetcher.is_cached(release, "cpu") + + +def test_pruning_also_runs_when_the_current_build_was_already_cached(tmp_path: Path) -> None: + fetcher, release = _fetch_current(tmp_path) + older = _older_build(tmp_path, "b0050-cpu") + later = BinaryFetcher(tmp_path, http_client=_client_returning(b"no download expected")) + + later.ensure_binary(release, "cpu") + + assert not older.exists() + assert fetcher.is_cached(release, "cpu") + + +def test_pruning_happens_once_per_release_and_process(tmp_path: Path) -> None: + fetcher, release = _fetch_current(tmp_path) + appeared_later = _older_build(tmp_path, "b0050-cpu") + + fetcher.ensure_binary(release, "cpu") + + assert appeared_later.is_dir() + + +@pytest.mark.skipif(sys.platform != "win32", reason="relies on Windows refusing to rename a directory that holds an open file") +def test_a_build_that_is_locked_is_kept_and_the_launch_still_succeeds(tmp_path: Path) -> None: + locked = _older_build(tmp_path, "b0050-cpu") + free = _older_build(tmp_path, "b0060-cpu") + # A file held open the way a running program holds its own files. + with (locked / "llama-server-impl.dll").open("rb"): + fetcher, release = _fetch_current(tmp_path) + + assert (locked / "llama-server.exe").is_file() + assert (locked / "nested" / "ggml-extra.dll").is_file() + assert not free.exists() + assert fetcher.is_cached(release, "cpu") + + +def test_a_build_a_program_is_running_from_is_left_whole(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + """A running image refuses a write-open; that alone is enough to keep the build.""" + in_use = _older_build(tmp_path, "b0050-cpu") + real_open = Path.open + + def open_refusing_the_running_image(self: Path, mode: str = "r", *args, **kwargs): + if self.name == "llama-server.exe" and "+" in mode: + raise PermissionError(13, "The process cannot access the file") + return real_open(self, mode, *args, **kwargs) + + monkeypatch.setattr(Path, "open", open_refusing_the_running_image) + + _fetch_current(tmp_path) + + assert (in_use / "llama-server.exe").is_file() + assert (in_use / "llama-server-impl.dll").is_file() + assert (in_use / "nested" / "ggml-extra.dll").is_file() + + +def test_a_failure_while_deleting_never_fails_the_launch_and_is_finished_next_time( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + older = _older_build(tmp_path, "b0050-cpu") + archive_bytes = _build_archive() + release = _release_for(archive_bytes, tag="b0100") # built first: it uses rmtree itself + real_rmtree = shutil.rmtree + + def refuse(*args, **kwargs): + raise PermissionError(13, "access denied") + + monkeypatch.setattr(binary_fetcher_module.shutil, "rmtree", refuse) + BinaryFetcher(tmp_path, http_client=_client_returning(archive_bytes)).ensure_binary(release, "cpu") # must not raise + + # The old build was moved out of the launchable name before the delete + # failed, so nothing can start from a half-deleted directory. + assert not older.exists() + leftovers = [entry for entry in tmp_path.iterdir() if entry.name.startswith(".prune-")] + assert len(leftovers) == 1 + + monkeypatch.setattr(binary_fetcher_module.shutil, "rmtree", real_rmtree) + BinaryFetcher(tmp_path, http_client=_client_returning(b"no download expected")).ensure_binary(release, "cpu") + + assert not any(entry.name.startswith(".prune-") for entry in tmp_path.iterdir()) + + +def test_at_most_a_handful_of_builds_are_removed_per_pass(tmp_path: Path) -> None: + for build in range(1, 13): + _older_build(tmp_path, f"b{build:04d}-cpu") + + _fetch_current(tmp_path) + + remaining = [entry.name for entry in tmp_path.iterdir() if entry.name.startswith("b00") and entry.name != "b0100-cpu"] + assert len(remaining) == 12 - binary_fetcher_module._MAX_PRUNED_PER_PASS + + +def test_a_link_that_looks_like_an_old_build_is_never_followed(tmp_path: Path) -> None: + outside = tmp_path / "outside" + (outside / "precious").mkdir(parents=True) + (outside / "precious" / "data.txt").write_text("keep", encoding="utf-8") + runtime = tmp_path / "runtime" + runtime.mkdir() + try: + (runtime / "b0050-cpu").symlink_to(outside, target_is_directory=True) + except (OSError, NotImplementedError): + pytest.skip("this account cannot create directory symlinks") + + _fetch_current(runtime) + + assert (outside / "precious" / "data.txt").read_text(encoding="utf-8") == "keep" + # And the link itself was not moved or renamed out of sight. + assert (runtime / "b0050-cpu").is_symlink() + assert not any(entry.name.startswith(".prune-") for entry in runtime.iterdir()) + + +def test_nothing_is_pruned_when_the_pinned_tag_is_not_a_build_number(tmp_path: Path) -> None: + older = _older_build(tmp_path, "b0050-cpu") + + _fetch_current(tmp_path, tag="custom-tag") + + assert older.is_dir() + + +def test_a_cancelled_pass_removes_nothing(tmp_path: Path) -> None: + fetcher, release = _fetch_current(tmp_path) + older = _older_build(tmp_path, "b0050-cpu") + fetcher._pruned_for = None + + fetcher._prune_superseded_builds(release, SimpleNamespace(is_set=lambda: True)) + + assert older.is_dir() + + +def test_pruning_logs_build_names_and_nothing_from_inside_them( + tmp_path: Path, caplog: pytest.LogCaptureFixture +) -> None: + caplog.set_level("DEBUG") + _older_build(tmp_path, "b0050-cpu") + + _fetch_current(tmp_path) + + assert "b0050-cpu" in caplog.text + for inside in ("llama-server-impl.dll", "ggml-extra.dll", "older stub"): + assert inside not in caplog.text diff --git a/tests/test_llamacpp_extra_args.py b/tests/test_llamacpp_extra_args.py new file mode 100644 index 0000000..92d0ede --- /dev/null +++ b/tests/test_llamacpp_extra_args.py @@ -0,0 +1,183 @@ +"""The advanced runtime options a user may add, and what they may not. + +Cortex owns the launch contract (model, context window, loopback address, port, +key, GPU layers, slots). What a user adds goes after it, so it is held to an +allow-list of tuning flags, and the settings model refuses anything else before +it is stored. +""" + +from __future__ import annotations + +import pytest +from pydantic import ValidationError + +from cortex_backend.core.settings import CortexSettings, LlamaCppSettings +from cortex_backend.llamacpp.extra_args import MAX_EXTRA_ARGS, validate_extra_args + + +@pytest.mark.parametrize( + "options", + [ + (), + ("-ctk", "q8_0"), + ("--cache-type-k", "f16", "--cache-type-v", "q4_0"), + ("-ctk", "bf16", "-ctv", "iq4_nl"), + ("-fa", "on"), + ("--flash-attn", "off"), + ("-fa", "auto", "-t", "8"), + # Older builds take the flag alone. + ("-fa",), + ("-fa", "-t", "8"), + ("-t", "16", "-tb", "16"), + ("--threads", "1", "--threads-batch", "1024"), + ], +) +def test_the_tuning_flags_are_accepted_as_written(options: tuple[str, ...]) -> None: + assert validate_extra_args(options) == options + + +@pytest.mark.parametrize( + "options", + [ + # Each of these would change something Cortex sets itself. + ("--host", "0.0.0.0"), + ("--port", "8080"), + ("--api-key", "not-a-real-key"), + ("--api-key-file", "keys.txt"), + ("-m", "other.gguf"), + ("--model", "other.gguf"), + ("-hf", "someone/repo"), + ("-c", "131072"), + ("--ctx-size", "131072"), + ("-ngl", "99"), + ("--n-gpu-layers", "0"), + ("-np", "8"), + ("--webui",), + ("--ui",), + ("--no-ui",), + ("--reasoning-format", "none"), + ("--ssl-key-file", "k.pem"), + # And anything else that is not on the list. + ("--unknown-flag",), + ("--mlock",), + ("--no-mmap",), + ("-ctk", "q8_0", "--host", "0.0.0.0"), + ("--host=0.0.0.0",), + ("--cache-type-k=q8_0",), + ("q8_0",), + ], +) +def test_anything_that_could_change_the_launch_contract_is_refused(options: tuple[str, ...]) -> None: + with pytest.raises(ValueError): + validate_extra_args(options) + + +@pytest.mark.parametrize( + "options", + [ + ("-ctk",), + ("-ctk", "q9_9"), + ("-ctk", "Q8_0"), + ("-ctv", "-t"), + ("-fa", "maybe"), + ("-t",), + ("-t", "0"), + ("-t", "-4"), + ("-t", "1025"), + ("-t", "4.5"), + ("-t", "four"), + ("-t", "٤"), # an Arabic-Indic digit is not an ASCII one + ("-t", "4", "4"), + ("-ctk", "q8_0", "-ctk", "f16"), + ("-ctk", "q8_0", "--cache-type-k", "f16"), + ("-t", "4", "--threads", "8"), + ("-fa", "-fa"), + ], +) +def test_a_flag_needs_exactly_its_own_value_and_appears_once(options: tuple[str, ...]) -> None: + with pytest.raises(ValueError): + validate_extra_args(options) + + +@pytest.mark.parametrize( + "options", + [ + ("",), + (" ",), + ("-t ", "4"), + ("-t", "4 "), + ("-t\n", "4"), + ("-ctk\t", "q8_0"), + ("-‑t", "4"), + ("-t", "4" * 40), + ("-" + "x" * 40,), + ], +) +def test_words_that_are_not_single_plain_words_are_refused(options: tuple[str, ...]) -> None: + with pytest.raises(ValueError): + validate_extra_args(options) + + +def test_the_list_is_bounded() -> None: + everything = ("-ctk", "q8_0", "-ctv", "q8_0", "-fa", "on", "-t", "8", "-tb", "8") + assert len(everything) <= MAX_EXTRA_ARGS + assert validate_extra_args(everything) == everything + with pytest.raises(ValueError, match="at most"): + validate_extra_args(("-fa",) * (MAX_EXTRA_ARGS + 1)) + + +def test_a_refusal_names_the_position_and_never_repeats_what_was_typed() -> None: + secretive = ("-t", "4", "--api-key", "hunter2-not-a-real-key", "--host", "10.9.8.7") + with pytest.raises(ValueError) as refused: + validate_extra_args(secretive) + + message = str(refused.value) + assert "option 3" in message + for typed in ("hunter2-not-a-real-key", "10.9.8.7", "--api-key"): + assert typed not in message + + +def test_settings_default_to_no_options() -> None: + assert CortexSettings().llamacpp.extra_args == () + + +def test_a_stored_document_from_before_this_field_still_loads_with_no_options() -> None: + settings = CortexSettings.model_validate({"llamacpp": {"gpu_backend": "cpu"}}) + + assert settings.llamacpp.gpu_backend == "cpu" + assert settings.llamacpp.extra_args == () + + +def test_settings_accept_a_list_from_json_and_keep_it_as_a_tuple() -> None: + settings = LlamaCppSettings.model_validate({"extra_args": ["-ctk", "q8_0", "-fa", "on"]}) + + assert settings.extra_args == ("-ctk", "q8_0", "-fa", "on") + + +@pytest.mark.parametrize("options", [("--host", "0.0.0.0"), ("--api-key", "x"), ("-c", "1"), ("-t",)]) +def test_settings_reject_options_the_allow_list_does_not_cover(options: tuple[str, ...]) -> None: + with pytest.raises(ValidationError): + LlamaCppSettings(extra_args=options) + + +def test_the_settings_api_refuses_a_host_override_without_echoing_it(client, headers) -> None: + current = client.get("/api/v1/settings", headers=headers).json()["settings"] + current["llamacpp"] = {**current.get("llamacpp", {}), "extra_args": ["--host", "203.0.113.9"]} + + response = client.put("/api/v1/settings", headers=headers, json={"settings": current}) + + assert response.status_code == 422 + assert "203.0.113.9" not in response.text + stored = client.get("/api/v1/settings", headers=headers).json()["settings"] + assert stored["llamacpp"].get("extra_args", []) == [] + + +def test_the_settings_api_round_trips_the_advanced_options(client, headers) -> None: + current = client.get("/api/v1/settings", headers=headers).json()["settings"] + current["llamacpp"] = {**current.get("llamacpp", {}), "extra_args": ["-ctk", "q8_0", "-t", "6"]} + + saved = client.put("/api/v1/settings", headers=headers, json={"settings": current}) + + assert saved.status_code == 200 + llamacpp = client.get("/api/v1/settings", headers=headers).json()["settings"]["llamacpp"] + assert llamacpp["extra_args"] == ["-ctk", "q8_0", "-t", "6"] diff --git a/tests/test_llamacpp_server_manager.py b/tests/test_llamacpp_server_manager.py index 13f540e..71827ec 100644 --- a/tests/test_llamacpp_server_manager.py +++ b/tests/test_llamacpp_server_manager.py @@ -23,6 +23,7 @@ BinaryVerificationError, CrashLoopError, LlamaCppError, + RuntimeBusyError, ServerLaunchError, ServerStartTimeoutError, ) @@ -253,8 +254,12 @@ def _manager( health_timeout_seconds: float = 5.0, startup_cap_seconds: float | None = None, release=_ANY_RELEASE, + **overrides, ) -> LlamaServerManager: extra = {} if startup_cap_seconds is None else {"startup_cap_seconds": startup_cap_seconds} + # The default probe reads this machine's graphics loader, which would make + # every "auto" test depend on where it runs. + extra.setdefault("vulkan_loader_probe", lambda: True) return LlamaServerManager( runtime_dir=tmp_path, fetcher=fetcher, @@ -264,7 +269,7 @@ def _manager( health_timeout_seconds=health_timeout_seconds, launcher=launcher, http_client=http_client, - **extra, + **{**extra, **overrides}, ) @@ -2733,3 +2738,921 @@ def test_the_warm_health_retry_passes_the_declared_timeout(tmp_path: Path) -> No assert "_HEALTH_RETRY_TIMEOUT_SECONDS" in source, ( "the retry path must pass the timeout it declares" ) + + +# --------------------------------------------------------------------------- +# A GPU build that cannot load is not downloaded (RT-09) +# --------------------------------------------------------------------------- + + +def _launch_with_probe(tmp_path: Path, *, gpu_backend: str, probe) -> tuple[LlamaServerManager, _FakeFetcher, _QueueLauncher]: + fetcher = _FakeFetcher() + launcher = _QueueLauncher([_FakePopen()]) + manager = _manager( + tmp_path, + fetcher=fetcher, + launcher=launcher, + http_client=_AlwaysHealthyClient(), + gpu_backend=gpu_backend, + vulkan_loader_probe=probe, + ) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + return manager, fetcher, launcher + + +def test_without_a_vulkan_loader_auto_asks_only_for_the_cpu_build(tmp_path: Path) -> None: + manager, fetcher, launcher = _launch_with_probe(tmp_path, gpu_backend="auto", probe=lambda: False) + + # The roughly 100 MB Vulkan archive was never requested, let alone launched. + assert fetcher.ensure_binary_calls == ["cpu"] + assert len(launcher.launch_args) == 1 + args = launcher.launch_args[0] + assert args[args.index("-ngl") + 1] == "0" + status = manager.status + assert status.state == "ready" + assert status.active_backend == "cpu" + assert status.backend_note is not None + assert "Vulkan" in status.backend_note + + +def test_with_a_vulkan_loader_auto_still_tries_the_gpu_build_first(tmp_path: Path) -> None: + manager, fetcher, launcher = _launch_with_probe(tmp_path, gpu_backend="auto", probe=lambda: True) + + assert fetcher.ensure_binary_calls == ["vulkan"] + args = launcher.launch_args[0] + assert args[args.index("-ngl") + 1] == "auto" + assert manager.status.active_backend == "vulkan" + assert manager.status.backend_note is None + + +def test_an_explicit_vulkan_choice_is_honoured_even_when_the_probe_finds_no_loader(tmp_path: Path) -> None: + """The probe can be wrong about an unusual install; the user asked for this build.""" + manager, fetcher, _launcher = _launch_with_probe(tmp_path, gpu_backend="vulkan", probe=lambda: False) + + assert fetcher.ensure_binary_calls == ["vulkan"] + assert manager.status.backend_note is None + + +def test_an_explicit_cpu_choice_never_consults_the_probe(tmp_path: Path) -> None: + probed: list[bool] = [] + + def probe() -> bool: + probed.append(True) + return False + + manager, fetcher, _launcher = _launch_with_probe(tmp_path, gpu_backend="cpu", probe=probe) + + assert probed == [] + assert fetcher.ensure_binary_calls == ["cpu"] + assert manager.status.backend_note is None + + +def test_the_skip_note_describes_the_latest_launch_only(tmp_path: Path) -> None: + loader_present = [False] + fetcher = _FakeFetcher() + launcher = _QueueLauncher([_FakePopen(), _FakePopen()]) + manager = _manager( + tmp_path, + fetcher=fetcher, + launcher=launcher, + http_client=_AlwaysHealthyClient(), + gpu_backend="auto", + vulkan_loader_probe=lambda: loader_present[0], + ) + manager.ensure_ready(tmp_path / "first.gguf", num_ctx=4096) + assert manager.status.backend_note is not None + + loader_present[0] = True # a driver was installed in the meantime + manager.ensure_ready(tmp_path / "second.gguf", num_ctx=4096) + + assert fetcher.ensure_binary_calls == ["cpu", "vulkan"] + assert manager.status.active_backend == "vulkan" + assert manager.status.backend_note is None + + +def test_the_loader_probe_looks_in_system32_then_on_the_search_path(tmp_path: Path) -> None: + from cortex_backend.llamacpp.server_manager import _vulkan_loader_present + + windows = tmp_path / "Windows" + (windows / "System32").mkdir(parents=True) + environ = {"SystemRoot": str(windows)} + never = lambda name: (_ for _ in ()).throw(AssertionError(f"searched for {name}")) # noqa: E731 + + assert _vulkan_loader_present(platform="win32", environ=environ, find_library=lambda name: None) is False + + (windows / "System32" / "vulkan-1.dll").write_bytes(b"loader") + assert _vulkan_loader_present(platform="win32", environ=environ, find_library=never) is True + + # Not in System32, but somewhere on the search path (a redistributable + # runtime next to the driver). + assert _vulkan_loader_present( + platform="win32", environ={}, find_library=lambda name: "C:/vulkan/vulkan-1.dll" + ) is True + + +def test_the_loader_probe_does_not_change_behaviour_off_windows() -> None: + from cortex_backend.llamacpp.server_manager import _vulkan_loader_present + + assert _vulkan_loader_present(platform="linux", environ={}, find_library=lambda name: None) is True + + +def test_a_failing_search_counts_as_no_loader() -> None: + from cortex_backend.llamacpp.server_manager import _vulkan_loader_present + + def broken(name: str) -> str: + raise OSError("search failed") + + assert _vulkan_loader_present(platform="win32", environ={}, find_library=broken) is False + + +# --------------------------------------------------------------------------- +# What the runtime says about GPU use, not just which build launched (RT-08) +# --------------------------------------------------------------------------- + +_LISTENING_BYTES = b"0.01.234.567 I srv start: listening on http://127.0.0.1:43125\n" + + +class _OutputPopen(_FakePopen): + """A child whose output is exactly ``lines``, then the line that says it is listening.""" + + def __init__(self, *lines: bytes, listening: bool = True, after: tuple[bytes, ...] = ()) -> None: + super().__init__() + self.stdout = io.BytesIO(b"".join(lines) + (_LISTENING_BYTES if listening else b"") + b"".join(after)) + + +def _ready_with_output(tmp_path: Path, popen: _FakePopen, *, gpu_backend: str = "vulkan") -> LlamaServerManager: + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=_QueueLauncher([popen]), + http_client=_AlwaysHealthyClient(), + gpu_backend=gpu_backend, + ) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + return manager + + +@pytest.mark.parametrize( + "line", + [ + b"load_tensors: offloaded 24/33 layers to GPU\n", + b"0.05.123.456 I load_tensors: offloaded 24/33 layers to GPU\n", + b"llm_load_tensors: offloaded 24/33 layers to GPU\r\n", + b"load_tensors: Offloaded 24 / 33 layers to gpu\n", + ], +) +def test_the_offload_line_before_listening_is_reported_on_the_status(tmp_path: Path, line: bytes) -> None: + manager = _ready_with_output(tmp_path, _OutputPopen(b"load_tensors: loading model tensors\n", line)) + + status = manager.status + assert status.state == "ready" + assert (status.gpu_layers_offloaded, status.gpu_layers_total) == (24, 33) + + +def test_a_gpu_build_that_offloaded_nothing_says_so(tmp_path: Path) -> None: + """The Vulkan build launching is not the GPU being used: 0 layers is the CPU.""" + manager = _ready_with_output(tmp_path, _OutputPopen(b"load_tensors: offloaded 0/33 layers to GPU\n")) + + status = manager.status + assert status.active_backend == "vulkan" + assert (status.gpu_layers_offloaded, status.gpu_layers_total) == (0, 33) + + +def test_no_offload_line_means_unknown_not_zero(tmp_path: Path) -> None: + manager = _ready_with_output(tmp_path, _OutputPopen(b"some other loader output\n")) + + status = manager.status + assert status.state == "ready" + assert status.active_backend == "vulkan" + assert status.gpu_layers_offloaded is None + assert status.gpu_layers_total is None + + +def test_the_last_offload_line_wins(tmp_path: Path) -> None: + manager = _ready_with_output( + tmp_path, + _OutputPopen( + b"load_tensors: offloaded 10/33 layers to GPU\n", + b"load_tensors: offloaded 33/33 layers to GPU\n", + ), + ) + + assert (manager.status.gpu_layers_offloaded, manager.status.gpu_layers_total) == (33, 33) + + +def test_an_offload_line_after_the_server_is_listening_is_not_this_load(tmp_path: Path) -> None: + manager = _ready_with_output( + tmp_path, + _OutputPopen(b"load_tensors: offloaded 12/33 layers to GPU\n", after=(b"load_tensors: offloaded 1/2 layers to GPU\n",)), + ) + + assert (manager.status.gpu_layers_offloaded, manager.status.gpu_layers_total) == (12, 33) + + +@pytest.mark.parametrize( + "line", + [ + # What the model file says about itself is chosen by whoever made it. + b"print_info: general.name = load_tensors: offloaded 99/99 layers to GPU\n", + b"llama_model_loader: - kv 3: general.name str = load_tensors: offloaded 99/99 layers to GPU\n", + # Not a line the loader prints: a claim buried in other text. + b"the model says load_tensors: offloaded 99/99 layers to GPU and then more\n", + # Counts that cannot be true. + b"load_tensors: offloaded 40/33 layers to GPU\n", + b"load_tensors: offloaded 0/0 layers to GPU\n", + # Absurdly long lines are not parsed at all. + b"x" * 600 + b" load_tensors: offloaded 5/6 layers to GPU\n", + ], +) +def test_lines_that_cannot_be_trusted_are_not_reported(tmp_path: Path, line: bytes) -> None: + manager = _ready_with_output(tmp_path, _OutputPopen(line)) + + assert manager.status.state == "ready" + assert manager.status.gpu_layers_offloaded is None + assert manager.status.gpu_layers_total is None + + +def test_a_new_server_does_not_inherit_the_previous_servers_counts(tmp_path: Path) -> None: + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=_QueueLauncher([ + _OutputPopen(b"load_tensors: offloaded 24/33 layers to GPU\n"), + _OutputPopen(b"nothing useful\n"), + ]), + http_client=_AlwaysHealthyClient(), + gpu_backend="vulkan", + ) + manager.ensure_ready(tmp_path / "first.gguf", num_ctx=4096) + assert manager.status.gpu_layers_offloaded == 24 + + manager.ensure_ready(tmp_path / "second.gguf", num_ctx=4096) + + assert manager.status.state == "ready" + assert manager.status.gpu_layers_offloaded is None + + +def test_the_counts_are_only_reported_while_the_server_is_ready(tmp_path: Path) -> None: + manager = _ready_with_output(tmp_path, _OutputPopen(b"load_tensors: offloaded 24/33 layers to GPU\n")) + assert manager.status.gpu_layers_offloaded == 24 + + manager.stop() + + status = manager.status + assert status.state == "idle" + assert status.gpu_layers_offloaded is None + assert status.gpu_layers_total is None + + +def test_the_offload_counts_are_the_only_thing_taken_from_the_line(tmp_path: Path, caplog: pytest.LogCaptureFixture) -> None: + caplog.set_level("DEBUG") + manager = _ready_with_output( + tmp_path, _OutputPopen(b"load_tensors: offloaded 24/33 layers to GPU\n") + ) + + assert "offloaded" not in caplog.text + assert "load_tensors" not in caplog.text + assert manager.status.gpu_layers_total == 33 + + +# --------------------------------------------------------------------------- +# Advanced runtime options and a context window the model was trained for (RT-22) +# --------------------------------------------------------------------------- + + +def _write_gguf(path: Path, *, context_length: int | None, architecture: str = "llama") -> None: + import gguf + import numpy as np + + writer = gguf.GGUFWriter(str(path), architecture) + if context_length is not None: + writer.add_context_length(context_length) + writer.add_name(path.stem) + writer.add_tensor("dummy.weight", np.zeros((2, 2), dtype=np.float32)) + writer.write_header_to_file() + writer.write_kv_data_to_file() + writer.write_tensors_to_file() + writer.close() + + +def _manager_for_a_real_file(tmp_path: Path, model_path: Path, processes: list[_FakePopen], **overrides): + client = _RecordingAttestationClient({"model_path": str(model_path), "build_info": "b10311-test"}) + launcher = _QueueLauncher(processes) + manager = _manager(tmp_path, fetcher=_FakeFetcher(), launcher=launcher, http_client=client, **overrides) + return manager, launcher + + +def _context_argument(argv: list[str]) -> str: + return argv[argv.index("-c") + 1] + + +def test_the_advanced_options_are_appended_after_the_fixed_launch_contract(tmp_path: Path) -> None: + options = ("-ctk", "q8_0", "-ctv", "q8_0", "-fa", "on", "-t", "8") + launcher = _QueueLauncher([_FakePopen()]) + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=launcher, + http_client=_AlwaysHealthyClient(), + extra_args=lambda: options, + ) + + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + argv = launcher.launch_args[0] + assert tuple(argv[-len(options):]) == options + fixed = argv[: -len(options)] + # The contract in front of them is untouched, and appears exactly once. + assert fixed[fixed.index("--host") + 1] == "127.0.0.1" + assert fixed[fixed.index("--port") + 1] == "0" + assert fixed.count("--host") == fixed.count("--port") == fixed.count("-c") == fixed.count("-m") == 1 + assert "--api-key" not in argv + + +def test_no_advanced_options_leaves_the_launch_exactly_as_it_was(tmp_path: Path) -> None: + plain = _QueueLauncher([_FakePopen()]) + empty = _QueueLauncher([_FakePopen()]) + _manager(tmp_path, fetcher=_FakeFetcher(), launcher=plain, http_client=_AlwaysHealthyClient()).ensure_ready( + tmp_path / "model.gguf", num_ctx=4096 + ) + _manager( + tmp_path, fetcher=_FakeFetcher(), launcher=empty, http_client=_AlwaysHealthyClient(), extra_args=lambda: () + ).ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + assert plain.launch_args == empty.launch_args + + +def test_changing_the_advanced_options_restarts_the_server_and_says_why(tmp_path: Path) -> None: + current: list[tuple[str, ...]] = [("-t", "4")] + first, second = _FakePopen(), _FakePopen() + launcher = _QueueLauncher([first, second]) + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=launcher, + http_client=_AlwaysHealthyClient(), + extra_args=lambda: current[0], + ) + model_path = tmp_path / "model.gguf" + manager.ensure_ready(model_path, num_ctx=4096) + manager.ensure_ready(model_path, num_ctx=4096) + assert len(launcher.launch_args) == 1 # unchanged options reuse the server + + current[0] = ("-t", "8") + manager.ensure_ready(model_path, num_ctx=4096) + + assert len(launcher.launch_args) == 2 + assert launcher.launch_args[1][-2:] == ["-t", "8"] + assert first.terminated + assert manager.status.last_restart_reason == "the advanced runtime options changed" + + +def test_clearing_the_advanced_options_relaunches_without_them(tmp_path: Path) -> None: + current: list[tuple[str, ...]] = [("-fa", "off")] + launcher = _QueueLauncher([_FakePopen(), _FakePopen()]) + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=launcher, + http_client=_AlwaysHealthyClient(), + extra_args=lambda: current[0], + ) + model_path = tmp_path / "model.gguf" + manager.ensure_ready(model_path, num_ctx=4096) + current[0] = () + + manager.ensure_ready(model_path, num_ctx=4096) + + assert len(launcher.launch_args) == 2 + assert "-fa" not in launcher.launch_args[1] + + +@pytest.mark.parametrize( + "options", + [ + ("--host", "0.0.0.0"), + ("--port", "8080"), + ("--api-key", "secret"), + ("-m", "other.gguf"), + ("-c", "999999"), + ("-ngl", "99"), + ("--unknown-flag",), + ("-t", "4", "--host", "0.0.0.0"), + ], +) +def test_options_that_would_change_the_launch_contract_never_reach_the_child( + tmp_path: Path, options: tuple[str, ...] +) -> None: + launcher = _QueueLauncher([_FakePopen()]) + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=launcher, + http_client=_AlwaysHealthyClient(), + extra_args=lambda: options, + ) + + with pytest.raises(LlamaCppError) as raised: + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + assert launcher.launch_args == [] + assert "not valid" in str(raised.value) + # The message is Cortex's own; it does not repeat what was configured. + for word in options: + assert word not in str(raised.value) + + +def test_changing_the_advanced_options_gives_a_failing_launch_a_fresh_start(tmp_path: Path) -> None: + current: list[tuple[str, ...]] = [("-t", "4")] + launcher = _QueueLauncher([_FakePopen(exit_immediately=True) for _ in range(3)] + [_FakePopen()]) + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=launcher, + http_client=_AlwaysHealthyClient(), + extra_args=lambda: current[0], + ) + model_path = tmp_path / "model.gguf" + for _ in range(3): + with pytest.raises(ServerLaunchError): + manager.ensure_ready(model_path, num_ctx=4096) + with pytest.raises(CrashLoopError): + manager.ensure_ready(model_path, num_ctx=4096) + assert len(launcher.launch_args) == 3 + + current[0] = ("-t", "2") # the option that was making it fail, fixed + manager.ensure_ready(model_path, num_ctx=4096) + + assert len(launcher.launch_args) == 4 + assert manager.status.state == "ready" + + +def test_a_window_above_the_trained_context_is_lowered_and_the_status_says_so(tmp_path: Path) -> None: + model_path = tmp_path / "small.gguf" + _write_gguf(model_path, context_length=4096) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen()]) + messages: list[str] = [] + + manager.ensure_ready(model_path, num_ctx=32768, on_status=messages.append) + + assert _context_argument(launcher.launch_args[0]) == "4096" + note = manager.status.context_note + assert note is not None + assert "4096" in note and "32768" in note + # Said once, through the progress callback, when the launch happened. + assert messages.count(note) == 1 + + +def test_asking_for_the_same_oversized_window_again_reuses_the_server(tmp_path: Path) -> None: + """The reuse decision must see the lowered value, or every message would reload the model.""" + model_path = tmp_path / "small.gguf" + _write_gguf(model_path, context_length=4096) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen(), _FakePopen()]) + messages: list[str] = [] + + manager.ensure_ready(model_path, num_ctx=32768, on_status=messages.append) + manager.ensure_ready(model_path, num_ctx=32768, on_status=messages.append) + manager.ensure_ready(model_path, num_ctx=16384, on_status=messages.append) + manager.ensure_ready(model_path, num_ctx=2048, on_status=messages.append) + + assert len(launcher.launch_args) == 1 + assert sum("limited to 4096" in message for message in messages) == 1 + assert manager.ready_handle(model_path, num_ctx=32768) is not None + + +def test_a_window_within_the_trained_context_is_used_as_asked_and_carries_no_note(tmp_path: Path) -> None: + model_path = tmp_path / "long.gguf" + _write_gguf(model_path, context_length=131072) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen()]) + + manager.ensure_ready(model_path, num_ctx=32768) + + assert _context_argument(launcher.launch_args[0]) == "32768" + assert manager.status.context_note is None + + +def test_the_note_goes_away_when_a_later_request_fits(tmp_path: Path) -> None: + model_path = tmp_path / "small.gguf" + _write_gguf(model_path, context_length=4096) + manager, _launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen()]) + manager.ensure_ready(model_path, num_ctx=32768) + assert manager.status.context_note is not None + + manager.ensure_ready(model_path, num_ctx=4096) + + assert manager.status.context_note is None + + +def test_a_model_that_does_not_say_what_it_was_trained_for_is_not_limited(tmp_path: Path) -> None: + model_path = tmp_path / "silent.gguf" + _write_gguf(model_path, context_length=None) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen()]) + + manager.ensure_ready(model_path, num_ctx=32768) + + assert _context_argument(launcher.launch_args[0]) == "32768" + assert manager.status.context_note is None + + +def test_an_implausibly_small_trained_context_is_treated_as_unknown(tmp_path: Path) -> None: + model_path = tmp_path / "odd.gguf" + _write_gguf(model_path, context_length=64) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen()]) + + manager.ensure_ready(model_path, num_ctx=8192) + + assert _context_argument(launcher.launch_args[0]) == "8192" + + +def test_a_file_that_cannot_be_read_as_gguf_is_not_limited_and_not_reread( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + import cortex_backend.llamacpp.server_manager as module + + model_path = tmp_path / "broken.gguf" + model_path.write_bytes(b"not a gguf file at all") + reads: list[Path] = [] + + def unreadable(path: Path): + reads.append(path) + return None + + monkeypatch.setattr(module, "read_gguf_metadata", unreadable) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen()]) + + manager.ensure_ready(model_path, num_ctx=8192) + manager.ensure_ready(model_path, num_ctx=8192) + + assert _context_argument(launcher.launch_args[0]) == "8192" + assert reads == [model_path] # read once for this version of the file + + +def test_a_changed_model_file_is_read_again(tmp_path: Path) -> None: + model_path = tmp_path / "grows.gguf" + _write_gguf(model_path, context_length=4096) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen(), _FakePopen()]) + manager.ensure_ready(model_path, num_ctx=32768) + assert _context_argument(launcher.launch_args[0]) == "4096" + + model_path.unlink() + # A different architecture name also makes the file a different size, so the + # change is seen even if the two writes share a timestamp tick. + _write_gguf(model_path, context_length=16384, architecture="llamax") + manager.ensure_ready(model_path, num_ctx=32768) + + # A different file (the model was replaced), so the window is re-derived and + # the larger trained context needs a bigger server than the one running. + assert _context_argument(launcher.launch_args[1]) == "16384" + + +def test_a_request_with_no_preference_is_still_held_to_the_trained_context(tmp_path: Path) -> None: + model_path = tmp_path / "tiny.gguf" + _write_gguf(model_path, context_length=2048) + manager, launcher = _manager_for_a_real_file(tmp_path, model_path, [_FakePopen()]) + + manager.ensure_ready(model_path, num_ctx=None) + + assert _context_argument(launcher.launch_args[0]) == "2048" + assert manager.status.context_note is None + + +def test_the_trained_context_limit_is_applied_before_the_crash_loop_key(tmp_path: Path) -> None: + """Three failures at an oversized request count against the value actually launched.""" + model_path = tmp_path / "small.gguf" + _write_gguf(model_path, context_length=4096) + manager, launcher = _manager_for_a_real_file( + tmp_path, model_path, [_FakePopen(exit_immediately=True) for _ in range(3)] + ) + + for requested in (8192, 32768, 65536): + with pytest.raises(ServerLaunchError): + manager.ensure_ready(model_path, num_ctx=requested) + with pytest.raises(CrashLoopError): + manager.ensure_ready(model_path, num_ctx=16384) + + assert [_context_argument(argv) for argv in launcher.launch_args] == ["4096"] * 3 + + +# --------------------------------------------------------------------------- +# Unloading the model: on request and after an idle period (RT-07) +# --------------------------------------------------------------------------- + + +class _Clock: + """A clock the test moves by hand; nothing here waits for real time.""" + + def __init__(self) -> None: + self.now = 10_000.0 + + def __call__(self) -> float: + return self.now + + def advance(self, seconds: float) -> None: + self.now += seconds + + +@pytest.fixture +def idle_managers(): + """Build managers with an idle setting, and close every one at teardown. + + The idle watcher is a ``cortex-`` thread, so a manager that is not closed + fails the session's thread check. + """ + created: list[LlamaServerManager] = [] + + def build(tmp_path: Path, *, minutes=lambda: 5, processes=None, clock=None, **overrides): + clock = clock or _Clock() + launcher = _QueueLauncher(processes if processes is not None else [_FakePopen(), _FakePopen()]) + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=launcher, + http_client=_AlwaysHealthyClient(), + idle_unload_minutes=minutes, + clock=clock, + **overrides, + ) + created.append(manager) + return manager, launcher, clock + + yield build + for manager in created: + manager.close() + + +def test_a_model_idle_for_the_whole_period_is_unloaded_and_the_reason_is_recorded(tmp_path: Path, idle_managers) -> None: + manager, launcher, clock = idle_managers(tmp_path) + process = launcher._processes[0] + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + clock.advance(5 * 60 - 1) + assert manager.unload_if_idle() is False + assert manager.status.state == "ready" + + clock.advance(1) + assert manager.unload_if_idle() is True + + status = manager.status + assert status.state == "idle" + assert status.loaded_model is None + assert process.terminated + assert status.last_restart_reason == "the model was unloaded after 5 minutes without use" + assert status.last_error is None + + +def test_the_next_request_after_an_idle_unload_starts_the_model_again(tmp_path: Path, idle_managers) -> None: + manager, launcher, clock = idle_managers(tmp_path) + model_path = tmp_path / "model.gguf" + manager.ensure_ready(model_path, num_ctx=4096) + clock.advance(10 * 60) + assert manager.unload_if_idle() is True + + handle = manager.ensure_ready(model_path, num_ctx=4096) + + assert len(launcher.launch_args) == 2 + assert handle.model_path == model_path + status = manager.status + assert status.state == "ready" + # The reason the model had to be loaded again is still on record. + assert status.last_restart_reason == "the model was unloaded after 5 minutes without use" + + +def test_an_open_request_keeps_the_model_loaded_however_long_it_takes(tmp_path: Path, idle_managers) -> None: + manager, _launcher, clock = idle_managers(tmp_path) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + with manager.request_scope(): + clock.advance(3 * 3600) # a generation far longer than the idle period + assert manager.unload_if_idle() is False + assert manager.status.state == "ready" + + # The idle period starts when the request ends, not when it began. + clock.advance(5 * 60 - 1) + assert manager.unload_if_idle() is False + clock.advance(1) + assert manager.unload_if_idle() is True + + +def test_every_use_restarts_the_idle_period(tmp_path: Path, idle_managers) -> None: + manager, launcher, clock = idle_managers(tmp_path) + model_path = tmp_path / "model.gguf" + manager.ensure_ready(model_path, num_ctx=4096) + + clock.advance(4 * 60) + manager.ensure_ready(model_path, num_ctx=4096) # a warm reuse is a use + clock.advance(4 * 60) + + assert manager.unload_if_idle() is False # eight minutes since it loaded, four since it was used + clock.advance(60) + assert manager.unload_if_idle() is True + assert len(launcher.launch_args) == 1 + + +@pytest.mark.parametrize("setting", [0, -3]) +def test_zero_turns_the_idle_unload_off(tmp_path: Path, idle_managers, setting: int) -> None: + manager, _launcher, clock = idle_managers(tmp_path, minutes=lambda: setting) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + clock.advance(30 * 24 * 3600) + + assert manager.unload_if_idle() is False + assert manager.status.state == "ready" + + +def test_a_manager_without_the_setting_never_unloads_and_starts_no_watcher(tmp_path: Path) -> None: + clock = _Clock() + manager = _manager( + tmp_path, + fetcher=_FakeFetcher(), + launcher=_QueueLauncher([_FakePopen()]), + http_client=_AlwaysHealthyClient(), + clock=clock, + ) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + clock.advance(30 * 24 * 3600) + + assert manager.unload_if_idle() is False + assert manager._idle_thread is None + assert manager.status.state == "ready" + + +def test_the_setting_is_read_at_every_check(tmp_path: Path, idle_managers) -> None: + minutes = [30] + manager, _launcher, clock = idle_managers(tmp_path, minutes=lambda: minutes[0]) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + clock.advance(6 * 60) + assert manager.unload_if_idle() is False + + minutes[0] = 5 # changed in Settings while the model stays loaded + + assert manager.unload_if_idle() is True + + +def test_an_unreadable_setting_means_never_not_unload(tmp_path: Path, idle_managers) -> None: + def broken() -> int: + raise RuntimeError("settings database unavailable") + + manager, _launcher, clock = idle_managers(tmp_path, minutes=broken) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + clock.advance(24 * 3600) + + assert manager.unload_if_idle() is False + assert manager.status.state == "ready" + + +def test_a_load_or_restart_in_flight_is_never_waited_for_or_cut_short(tmp_path: Path, idle_managers) -> None: + manager, _launcher, clock = idle_managers(tmp_path) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + clock.advance(60 * 60) + + # The slow-path lock is what a load, a restart or a health check holds. + with manager._ensure_lock: + assert manager.unload_if_idle() is False + assert manager.status.state == "ready" + + assert manager.unload_if_idle() is True + + +def test_a_request_that_arrives_while_the_model_is_being_released_gets_it_loaded_again( + tmp_path: Path, idle_managers +) -> None: + """The unload holds the same lock as a load, so a request waits for it and then loads.""" + manager, launcher, clock = idle_managers(tmp_path) + model_path = tmp_path / "model.gguf" + manager.ensure_ready(model_path, num_ctx=4096) + clock.advance(60 * 60) + lock = _ContentionSignallingLock() + manager._ensure_lock = lock # type: ignore[assignment] + original = manager._terminate_and_reset + request_result: list[object] = [] + inside_teardown = threading.Event() + + def teardown_that_waits_for_the_request() -> bool: + if threading.current_thread() is worker: + return original() + inside_teardown.set() + assert lock.contended.wait(5.0), "the request never queued behind the unload" + return original() + + def request() -> None: + assert inside_teardown.wait(5.0) + request_result.append(manager.ensure_ready(model_path, num_ctx=4096)) + + worker = threading.Thread(target=request, name="test-request") + manager._terminate_and_reset = teardown_that_waits_for_the_request # type: ignore[method-assign] + worker.start() + try: + assert manager.unload_if_idle() is True + finally: + worker.join(timeout=10.0) + + assert not worker.is_alive() + assert len(request_result) == 1 + assert len(launcher.launch_args) == 2 + assert manager.status.state == "ready" + + +def test_a_manual_unload_stops_the_model_and_is_safe_to_repeat(tmp_path: Path, idle_managers) -> None: + manager, launcher, _clock = idle_managers(tmp_path) + process = launcher._processes[0] + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + assert manager.unload() is True + + status = manager.status + assert status.state == "idle" + assert status.loaded_model is None + assert process.terminated + assert status.last_restart_reason == "the model was unloaded at your request" + assert manager.unload() is False # nothing left to unload; not an error + assert manager.status.state == "idle" + + +def test_a_manual_unload_with_nothing_loaded_changes_nothing(tmp_path: Path, idle_managers) -> None: + manager, launcher, _clock = idle_managers(tmp_path) + + assert manager.unload() is False + + assert launcher.launch_args == [] + assert manager.status.last_restart_reason is None + + +def test_a_manual_unload_is_refused_while_a_request_is_using_the_model(tmp_path: Path, idle_managers) -> None: + manager, launcher, _clock = idle_managers(tmp_path) + process = launcher._processes[0] + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + with manager.request_scope(), pytest.raises(RuntimeBusyError) as refused: + manager.unload() + + assert "answering a request" in str(refused.value) + assert not process.terminated + assert manager.status.state == "ready" + assert manager.unload() is True # once the request ends it goes through + + +def test_a_manual_unload_is_refused_while_a_model_is_loading( + tmp_path: Path, idle_managers, monkeypatch: pytest.MonkeyPatch +) -> None: + import cortex_backend.llamacpp.server_manager as module + + monkeypatch.setattr(module, "_UNLOAD_LOCK_TIMEOUT_SECONDS", 0.05) + manager, _launcher, _clock = idle_managers(tmp_path) + + with manager._ensure_lock, pytest.raises(RuntimeBusyError) as refused: + manager.unload() + + assert "being loaded" in str(refused.value) + + +def test_a_manual_unload_after_the_manager_closed_fails_closed(tmp_path: Path, idle_managers) -> None: + manager, _launcher, _clock = idle_managers(tmp_path) + manager.close() + + with pytest.raises(LlamaCppError, match="closed"): + manager.unload() + + +def test_a_process_that_will_not_exit_is_reported_and_not_called_unloaded(tmp_path: Path, idle_managers) -> None: + manager, _launcher, _clock = idle_managers(tmp_path, processes=[_UnstoppablePopen()]) + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + + with pytest.raises(LlamaCppError, match="did not exit cleanly"): + manager.unload() + + assert manager.status.state == "stopping" + + +def test_the_watcher_unloads_an_idle_model_by_itself_and_is_gone_after_close(tmp_path: Path, idle_managers) -> None: + from support import wait_until + + manager, launcher, clock = idle_managers(tmp_path, idle_check_interval_seconds=0.01) + process = launcher._processes[0] + manager.ensure_ready(tmp_path / "model.gguf", num_ctx=4096) + watcher = manager._idle_thread + assert watcher is not None and watcher.name == "cortex-llama-idle-unload" + + clock.advance(6 * 60) + + wait_until(lambda: manager.status.state == "idle", describe="the idle model to be unloaded") + assert process.terminated + manager.close() + watcher.join(timeout=5.0) + assert not watcher.is_alive() + + +def test_one_watcher_serves_every_load(tmp_path: Path, idle_managers) -> None: + manager, _launcher, _clock = idle_managers(tmp_path) + manager.ensure_ready(tmp_path / "first.gguf", num_ctx=4096) + first = manager._idle_thread + manager.ensure_ready(tmp_path / "second.gguf", num_ctx=4096) + + assert manager._idle_thread is first + assert first is not None and first.is_alive() + + +def test_the_default_check_interval_is_well_under_the_smallest_period(tmp_path: Path) -> None: + import cortex_backend.llamacpp.server_manager as module + + assert module._IDLE_CHECK_INTERVAL_SECONDS <= 60.0 diff --git a/tests/test_llamacpp_status_api.py b/tests/test_llamacpp_status_api.py index 85c0c0f..d7ace8d 100644 --- a/tests/test_llamacpp_status_api.py +++ b/tests/test_llamacpp_status_api.py @@ -73,3 +73,76 @@ def test_status_route_reports_no_failure_cause_by_default() -> None: assert _llamacpp_status(_request_with_manager(SimpleNamespace(status=live))).last_failure_code is None assert _llamacpp_status(_request_with_manager(None)).last_failure_code is None + + +def test_status_route_reports_why_the_gpu_build_was_not_used() -> None: + note = "No Vulkan graphics loader was found on this computer, so the CPU build is used." + live = LlamaCppRuntimeStatus( + state="ready", + binary_present=True, + loaded_model="gguf:model.gguf", + last_error=None, + models_directory="C:/synthetic/models", + active_backend="cpu", + backend_note=note, + ) + + status = _llamacpp_status(_request_with_manager(SimpleNamespace(status=live))) + + assert status.backend_note == note + assert status.model_dump()["backend_note"] == note + assert _llamacpp_status(_request_with_manager(None)).backend_note is None + + +def test_status_route_reports_how_many_layers_are_on_the_gpu() -> None: + live = LlamaCppRuntimeStatus( + state="ready", + binary_present=True, + loaded_model="gguf:model.gguf", + last_error=None, + models_directory="C:/synthetic/models", + active_backend="vulkan", + gpu_layers_offloaded=24, + gpu_layers_total=33, + ) + + status = _llamacpp_status(_request_with_manager(SimpleNamespace(status=live))) + + assert (status.gpu_layers_offloaded, status.gpu_layers_total) == (24, 33) + + +def test_status_route_reports_unknown_offload_as_null_not_zero() -> None: + live = LlamaCppRuntimeStatus( + state="ready", + binary_present=True, + loaded_model="gguf:model.gguf", + last_error=None, + models_directory="C:/synthetic/models", + active_backend="vulkan", + ) + + for status in ( + _llamacpp_status(_request_with_manager(SimpleNamespace(status=live))), + _llamacpp_status(_request_with_manager(None)), + ): + assert status.gpu_layers_offloaded is None + assert status.gpu_layers_total is None + + +def test_status_route_reports_that_the_context_window_was_limited() -> None: + note = "The context window was limited to 4096 tokens, the most this model was trained for (32768 were requested)." + live = LlamaCppRuntimeStatus( + state="ready", + binary_present=True, + loaded_model="gguf:model.gguf", + last_error=None, + models_directory="C:/synthetic/models", + active_backend="cpu", + loaded_context=4096, + context_note=note, + ) + + status = _llamacpp_status(_request_with_manager(SimpleNamespace(status=live))) + + assert status.context_note == note + assert _llamacpp_status(_request_with_manager(None)).context_note is None diff --git a/tests/test_llamacpp_unload_api.py b/tests/test_llamacpp_unload_api.py new file mode 100644 index 0000000..7f5c733 --- /dev/null +++ b/tests/test_llamacpp_unload_api.py @@ -0,0 +1,201 @@ +"""Unloading the local model on request: the route, its refusals, and the settings behind idle unload.""" + +from __future__ import annotations + +from pathlib import Path +from types import SimpleNamespace + +import pytest +from pydantic import ValidationError + +from cortex_backend.core.settings import CortexSettings, LlamaCppSettings +from cortex_backend.llamacpp.errors import LlamaCppError, RuntimeBusyError +from cortex_backend.llamacpp.server_manager import LlamaCppRuntimeStatus, LlamaServerManager + + +class _FakeManager: + """Stands in for the manager: records unloads, reports a status, can be told to refuse.""" + + def __init__(self, *, raises: Exception | None = None) -> None: + self.raises = raises + self.unload_calls = 0 + self.closed = False + self._loaded = True + + @property + def status(self) -> LlamaCppRuntimeStatus: + return LlamaCppRuntimeStatus( + state="ready" if self._loaded else "idle", + binary_present=True, + loaded_model="gguf:model.gguf" if self._loaded else None, + last_error=None, + models_directory="C:/synthetic/models", + active_backend="vulkan", + last_restart_reason=None if self._loaded else "the model was unloaded at your request", + ) + + def unload(self) -> bool: + self.unload_calls += 1 + if self.raises is not None: + raise self.raises + was_loaded, self._loaded = self._loaded, False + return was_loaded + + def close(self) -> None: + self.closed = True + + +@pytest.fixture +def fake_manager(app) -> _FakeManager: + manager = _FakeManager() + app.state.llamacpp_manager = manager + return manager + + +class _BusyJobs: + """The app's job registry, except that a response is being generated.""" + + def __init__(self, real) -> None: + self._real = real + + def active_snapshot(self, *, kind: str): + return object() if kind == "generation" else None + + def __getattr__(self, name: str): + return getattr(self._real, name) + + +def test_unload_stops_the_model_and_returns_the_new_status(client, headers, fake_manager) -> None: + response = client.post("/api/v1/llamacpp/unload", headers=headers) + + assert response.status_code == 200 + body = response.json() + assert body["state"] == "idle" + assert body["loaded_model"] is None + assert body["last_restart_reason"] == "the model was unloaded at your request" + assert fake_manager.unload_calls == 1 + + +def test_unload_is_safe_to_repeat(client, headers, fake_manager) -> None: + assert client.post("/api/v1/llamacpp/unload", headers=headers).status_code == 200 + again = client.post("/api/v1/llamacpp/unload", headers=headers) + + assert again.status_code == 200 + assert again.json()["state"] == "idle" + + +def test_unload_needs_a_session(client, fake_manager) -> None: + response = client.post("/api/v1/llamacpp/unload") + + assert response.status_code == 401 + assert fake_manager.unload_calls == 0 + + +def test_unload_is_refused_while_a_response_is_being_generated(app, client, headers, fake_manager) -> None: + app.state.jobs = _BusyJobs(app.state.jobs) + + response = client.post("/api/v1/llamacpp/unload", headers=headers) + + assert response.status_code == 409 + assert "being generated" in response.json()["detail"] + assert fake_manager.unload_calls == 0 + assert fake_manager.status.state == "ready" + + +def test_unload_reports_a_manager_that_is_busy_as_a_conflict(client, headers, app) -> None: + app.state.llamacpp_manager = _FakeManager( + raises=RuntimeBusyError("The model is answering a request. Stop it or wait for it to finish, then unload.") + ) + + response = client.post("/api/v1/llamacpp/unload", headers=headers) + + assert response.status_code == 409 + assert "answering a request" in response.json()["detail"] + + +def test_unload_reports_a_process_that_would_not_exit_without_claiming_success(client, headers, app) -> None: + app.state.llamacpp_manager = _FakeManager( + raises=LlamaCppError("The local model runtime did not exit cleanly; restart Cortex before trying again.") + ) + + response = client.post("/api/v1/llamacpp/unload", headers=headers) + + assert response.status_code == 500 + assert "restart Cortex" in response.json()["detail"] + + +def test_unload_without_a_runtime_says_so(client, headers, app) -> None: + assert app.state.llamacpp_manager is None + + response = client.post("/api/v1/llamacpp/unload", headers=headers) + + assert response.status_code == 409 + assert "unavailable" in response.json()["detail"] + + +def test_unload_through_the_real_manager_is_refused_while_a_request_holds_the_model( + client, headers, app, tmp_path: Path +) -> None: + manager = LlamaServerManager( + runtime_dir=tmp_path, + fetcher=SimpleNamespace(), # type: ignore[arg-type] + release=None, + gpu_backend_setting=lambda: "cpu", + models_directory=lambda: tmp_path, + ) + app.state.llamacpp_manager = manager + try: + with manager.request_scope(): + busy = client.post("/api/v1/llamacpp/unload", headers=headers) + idle = client.post("/api/v1/llamacpp/unload", headers=headers) + finally: + manager.close() + + assert busy.status_code == 409 + assert "answering a request" in busy.json()["detail"] + assert idle.status_code == 200 + assert idle.json()["state"] == "idle" + + +def test_the_unload_route_is_part_of_the_generated_contract() -> None: + from cortex_backend.api.routers import build_router + + paths = {route.path: route.methods for route in build_router().routes if getattr(route, "methods", None)} + + assert paths["/llamacpp/unload"] == {"POST"} + + +# -- the idle period that shares the runtime settings ----------------------------------------- + + +def test_the_idle_period_defaults_to_half_an_hour() -> None: + assert CortexSettings().llamacpp.idle_unload_minutes == 30 + assert CortexSettings.model_validate({"llamacpp": {"gpu_backend": "cpu"}}).llamacpp.idle_unload_minutes == 30 + + +@pytest.mark.parametrize("minutes", [-1, 1441, 10**6]) +def test_the_idle_period_is_bounded(minutes: int) -> None: + with pytest.raises(ValidationError): + LlamaCppSettings(idle_unload_minutes=minutes) + + +@pytest.mark.parametrize("minutes", [0, 1, 30, 1440]) +def test_zero_and_every_period_up_to_a_day_are_valid(minutes: int) -> None: + assert LlamaCppSettings(idle_unload_minutes=minutes).idle_unload_minutes == minutes + + +def test_the_settings_api_round_trips_the_idle_period(client, headers) -> None: + current = client.get("/api/v1/settings", headers=headers).json()["settings"] + current["llamacpp"] = {**current.get("llamacpp", {}), "idle_unload_minutes": 12} + + saved = client.put("/api/v1/settings", headers=headers, json={"settings": current}) + + assert saved.status_code == 200 + assert client.get("/api/v1/settings", headers=headers).json()["settings"]["llamacpp"]["idle_unload_minutes"] == 12 + + +def test_the_settings_api_refuses_an_idle_period_out_of_range(client, headers) -> None: + current = client.get("/api/v1/settings", headers=headers).json()["settings"] + current["llamacpp"] = {**current.get("llamacpp", {}), "idle_unload_minutes": 5000} + + assert client.put("/api/v1/settings", headers=headers, json={"settings": current}).status_code == 422