Coverage for src/lilbee/providers/fleet/planning.py: 100%
1060 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
1"""Launch planning for the fleet: device probe, VRAM estimate, placement, argv."""
3from __future__ import annotations
5import contextvars
6import logging
7import re
8import threading
9import time
10from contextlib import contextmanager
11from dataclasses import dataclass, field, replace
12from pathlib import Path
13from typing import TYPE_CHECKING, TypeAlias
15from lilbee.core.config.enums import KV_CACHE_TYPE_BYTES, KvCacheType
16from lilbee.core.system import is_network_path
17from lilbee.providers import engine_params, model_cache
18from lilbee.providers.base import ProviderError
19from lilbee.providers.fleet import ctx as fleet_ctx
20from lilbee.providers.fleet.adapters import (
21 LLM_RERANK_CONCURRENCY,
22 ROLE_SPECS,
23 RoleServerSpec,
24 build_server_argv,
25 embed_spec,
26 rerank_spec,
27 resolve_rerank_mode,
28)
29from lilbee.providers.fleet.binary import (
30 engine_binary_identity,
31 engine_build_id,
32 llama_server_runtime_env,
33 resolve_llama_server,
34)
35from lilbee.providers.fleet.devices import (
36 VULKAN_BACKEND,
37 FleetDevice,
38 host_lacks_nvlink,
39 probe_devices,
40 visible_env,
41)
42from lilbee.providers.fleet.launch import InstanceLaunch
43from lilbee.providers.fleet.placement import (
44 InstancePlan,
45 ModelPlacementInput,
46 PeakEstimator,
47 Placement,
48 SplitCtxFitter,
49 placement_from_spec,
50 plan_placement,
51)
52from lilbee.providers.fleet.placement_spec import PlacementError, PlacementSpec
53from lilbee.providers.fleet.readback import supports_memory_readback
54from lilbee.providers.fleet.replicas import resolve_replica_count
55from lilbee.providers.fleet.vram import estimate_instance_footprint, usable_vram_fraction
56from lilbee.providers.model_cache import free_system_memory, total_system_memory
57from lilbee.providers.model_ref import parse_model_ref
58from lilbee.providers.roles import ROLE_REGISTRY, EngineBackend, RerankMode, WorkerRole
59from lilbee.runtime.progress.types import OcrBackendUsed
61log = logging.getLogger(__name__)
63if TYPE_CHECKING:
64 from collections.abc import Callable, Iterable, Iterator, Mapping, Sequence
66# Fleet-only concurrency: continuous-batching slots (--parallel) per server.
67_CHAT_SLOTS = 4
68# Stand-in in the launch log for a host whose devices could not be read.
69_NO_DEVICES = "none"
72def log_engine_launch(launch: InstanceLaunch, *, owner_pid: int | None = None) -> None:
73 """Log the binary, build, backend and devices serving *launch*.
75 *owner_pid* is the engine's owner when this process adopted it rather than
76 spawned it.
77 """
78 if owner_pid is None:
79 log.info("Launched %s serving %s on %s", launch.model_id, launch.model, _engine_id(launch))
80 return
81 log.info(
82 "Adopted %s serving %s from engine pid %d on %s",
83 launch.model_id,
84 launch.model,
85 owner_pid,
86 _engine_id(launch),
87 )
90def _engine_id(launch: InstanceLaunch) -> str:
91 """*launch*'s binary with the engine build, backend, and probed devices."""
92 devices = probed_devices()
93 names = ", ".join(f"{d.backend}{d.index}: {d.name}" for d in devices) or _NO_DEVICES
94 return (
95 f"{launch.binary} (build {engine_build_id()}, "
96 f"backend {engine_backend().value}, devices: {names})"
97 )
100def warn_when_embed_window_below_chunk(launch: InstanceLaunch) -> None:
101 """Log when the embed engine's input cap is below the configured chunk budget.
103 Runs at engine adoption so the serving process's own log names the window
104 it serves; the health route carries the same warning for clients.
105 """
106 if launch.token_cap is None:
107 return
108 warning = engine_params.embed_window_warning(launch.token_cap)
109 if warning is not None:
110 log.warning("%s", " ".join(filter(None, (warning.message, warning.remedy))))
113def warn_when_chat_downsized(launch: InstanceLaunch) -> None:
114 """Log when a chat engine's granted shape ends below the requested one.
116 Runs at engine adoption so the warning lands on the serving process's
117 own log, where the operator of that engine reads it.
118 """
119 below_ctx = launch.built_ctx_target > 0 and launch.ctx < launch.built_ctx_target
120 below_slots = launch.slots < _CHAT_SLOTS
121 if not (below_ctx or below_slots):
122 return
123 # A pre-field launch record carries built_ctx_target=0; show the granted
124 # window as the requested one rather than "x 0 context". The remedy names
125 # both target knobs because the recorded target does not say which of
126 # num_ctx (pin) or chat_n_ctx_target produced it.
127 requested_ctx = launch.built_ctx_target or launch.ctx
128 log.warning(
129 "Chat engine downsized: serving %d slot(s) x %d context "
130 "(requested: up to %d slots x %d context). Requests beyond %d at once "
131 "wait for a free slot, so parallel agents can look stalled. "
132 "A smaller model or a lower context target (num_ctx or "
133 "chat_n_ctx_target) frees room for more slots.",
134 launch.slots,
135 launch.ctx,
136 _CHAT_SLOTS,
137 requested_ctx,
138 launch.slots,
139 )
142# Slots the PLACEMENT estimate reserves KV for on a tensor-split chat: one full
143# window, the minimum any split must hold. The launch may serve more than this
144# (see _resolve_split_chat_slots) when the cards measurably have room for several
145# full windows, so this is a planning floor, not the served slot count.
146_SPLIT_CHAT_SLOTS = 1
147# Floor context the PLACEMENT estimate reserves KV for, so a large model is never
148# single-carded into a KV corner too small for real use (a 17GB model on a 24GB
149# card leaves ~no KV room -> n_ctx collapses to a few hundred tokens). Sizing the
150# placement reserve against this floor forces a tensor-split when one card cannot
151# hold weights + a usable context; the served ctx is then grown by resolve_chat_ctx
152# (single) / fit_split_ctx (split) toward the cards' real headroom, with each
153# sequence capped at the working-context target. A split may then serve several
154# such sequences, so the served total can exceed the single-window reserve; the
155# per-device headroom test in fit_split_ctx is what bounds it. A fit still at
156# its floor after the forced split is refused by plan_launches
157# (_unusable_chat_ctx_reason).
158_MIN_USABLE_CHAT_CTX = 8192
159# Offload-search upper bound meaning "no reduced offload is worth probing": the
160# search runs 0..upper, so a negative bound skips it.
161_NO_OFFLOAD_PROBE = -1
162# Embed and cross-encoder rerank serve one request at a time. Raising it was
163# tried and measured worse: on 8xA40 with an 8B Q8 embedder and one ~100-token
164# passage per request, --parallel 1 gave 133 docs/sec at 81% SM while
165# --parallel 8 gave 100 at 63%. At one slot the card is already busy, so there is
166# no stall for slot-batching to reclaim and the extra slots only add
167# continuous-batching and KV-fragmentation overhead. Batch on the request side
168# (embed_batch_sequences) instead.
169_AUX_SLOTS = 1
170# A tensor-split needs at least this many GPUs; below it the chat context objective
171# (a gguf read) is pointless because the model can only single-card or stay unplaced.
172_MIN_SPLIT_GPUS = 2
173# Pooled single-slot search roles (embed/cross-encoder rerank) whose whole input
174# batches in one pass; derived from the role registry.
175_EMBED_ROLES = tuple(role for role, info in ROLE_REGISTRY.items() if info.pooled)
176# Roles whose loaders offload every layer regardless of cfg.n_gpu_layers; only
177# chat honors cfg.n_gpu_layers.
178_ALL_LAYER_ROLES = tuple(role for role, info in ROLE_REGISTRY.items() if info.offload_all_layers)
179_FLASH_ON = "on"
180_FLASH_OFF = "off"
181_FLASH_AUTO = "auto"
182# llama-server's documented way to say "offload nothing": --device none.
183_NO_DEVICE = "none"
184# Backends pinned by the name the engine printed rather than through an env var,
185# because their variables index a different space than --list-devices reports.
186_NAME_PINNED_BACKENDS = frozenset({VULKAN_BACKEND, "SYCL"})
187# Backends whose flash-attention coverage in llama.cpp is complete enough to ask
188# for it outright. Vulkan and SYCL are behind CUDA's and have been incomplete on
189# Intel's mesa driver, so those are left to the engine's own auto, which enables
190# flash attention only where the backend really supports it.
191_TRUSTED_FLASH_BACKENDS = frozenset({"CUDA", "ROCm", "HIP", "MTL", "Metal"})
192# Roles to which flash attention applies; embed/rerank run without it.
193_FLASH_ROLES = tuple(role for role, info in ROLE_REGISTRY.items() if info.flash_attn)
196# Cap vision's own KV footprint at this fraction of usable VRAM when sizing its
197# batching slots, leaving room for the weights and any co-located role.
198_VISION_VRAM_FRACTION = 0.5
200# Cap an LLM reranker's footprint at this fraction of usable VRAM when sizing its
201# slots; its per-slot ctx is tiny, so a normal GPU fits the full fan-out and a
202# small one steps down toward 1.
203_LLM_RERANK_VRAM_FRACTION = 0.5
205# RAM kept free for the OS when placing against system memory (no discrete GPU):
206# a quarter of total RAM, capped at 4 GiB. A fixed 4 GiB floor leaves a small
207# host (7-8 GB) with no budget at all, refusing to serve even tiny models.
208_SYSTEM_MEMORY_FLOOR_DIVISOR = 4
209# A GPU driver still initializing at boot answers with no devices. Ask again
210# before letting that decide the daemon's whole run; two extra probes cost a
211# couple of seconds only on a host that has a card the engine could not see.
212_PROBE_RETRIES = 2
213_PROBE_RETRY_DELAY_S = 1.0
215# A network filesystem makes mmap dangerous (page faults served over the wire can
216# wedge the loader in uninterruptible I/O), so the chat server loads its weights
217# into a malloc'd host copy (--no-mmap) whenever that copy fits in this fraction
218# of total system RAM. Local disk keeps mmap: its lazy paging gives a faster first
219# token on a cold cache -- the common desktop first launch -- and --no-mmap's
220# buffered full read only wins on an already-hot cache (#474: 33s vs 43s for a
221# 112GB model on 3 GPUs) while pessimizing cold start. Keyed on TOTAL memory
222# (stable), not free (fluctuates), so replans do not flap the launch argv. The
223# exact ceiling is tuned on a network-volume host.
224_NO_MMAP_NETWORK_RAM_FRACTION = 0.85
226# llama.cpp split-GGUF shard naming ("%s-%05d-of-%05d.gguf"); the cold-load
227# timeout must scale with the SUM of the shards, not the first file alone.
228_SPLIT_GGUF_NAME = re.compile(r"^(?P<prefix>.+)-(?P<index>\d{5})-of-(?P<total>\d{5})\.gguf$")
231def _weights_bytes(model_path: Path) -> int:
232 """Total weights size on disk; a split GGUF sums every sibling shard."""
233 match = _SPLIT_GGUF_NAME.fullmatch(model_path.name)
234 if match is None:
235 return model_path.stat().st_size
236 return sum(
237 sibling.stat().st_size
238 for sibling in model_path.parent.iterdir()
239 if _is_sibling_shard(sibling.name, match)
240 )
243def _is_sibling_shard(name: str, match: re.Match[str]) -> bool:
244 """Whether *name* is a shard of the same split GGUF as *match*."""
245 shard = _SPLIT_GGUF_NAME.fullmatch(name)
246 return (
247 shard is not None
248 and shard["prefix"] == match["prefix"]
249 and shard["total"] == match["total"]
250 )
253def _slots_for(
254 role: WorkerRole,
255 model_path: Path,
256 ctx: int,
257 *,
258 mmproj_path: Path | None = None,
259 unified_budget: int | None = None,
260 chat_reservation: int = 0,
261 rerank_mode: RerankMode | None = None,
262 device: FleetDevice | None = None,
263) -> int:
264 """Continuous-batching slots (--parallel) for a role's server.
266 Chat batches concurrent turns; vision batches concurrent OCR pages since a
267 one-page decode underutilizes the GPU; an LLM reranker batches its per-candidate
268 chat requests; embed and cross-encoder rerank are single-slot (their batching is
269 request-side). The memory-aware roles drop toward 1 on a small or shared host
270 instead of overcommitting. ``unified_budget`` caps sizing against free system RAM
271 with no discrete GPU; ``chat_reservation`` is the search-role footprint held back
272 from chat; ``device`` is the card the role was placed on, whose memory the
273 budget comes from once placement has chosen one.
274 """
275 if role is WorkerRole.CHAT:
276 return _resolve_chat_slots(
277 model_path,
278 ctx,
279 mmproj_path=mmproj_path,
280 unified_budget=unified_budget,
281 chat_reservation=chat_reservation,
282 device=device,
283 )
284 if role is WorkerRole.VISION:
285 return _resolve_vision_slots(
286 model_path, ctx, mmproj_path=mmproj_path, unified_budget=unified_budget, device=device
287 )
288 if role is WorkerRole.RERANK and rerank_mode is RerankMode.LLM:
289 return _resolve_llm_rerank_slots(
290 model_path, ctx, unified_budget=unified_budget, device=device
291 )
292 return _AUX_SLOTS
295def _resolve_split_chat_slots(fit_fn: Callable[[int], int]) -> tuple[int, int]:
296 """Largest split-chat slot count whose sequences each keep the full window.
298 ``fit_fn(n)`` is the per-slot context that fits when serving ``n`` sequences
299 (``fit_split_ctx``, capped at the working target and verified against real
300 per-card headroom). More slots divide the KV, so a split whose cards hold
301 several full windows can serve that many agents concurrently instead of one.
302 Returns ``(slots, per_slot_ctx)``, falling to one slot when only one full
303 window fits (or the fit degenerated to the floor), which preserves the
304 max-context single-sequence behaviour on a tight card.
306 Found by bisection rather than a scan because every ``fit_fn`` call is a
307 complete binary search whose probes each shell out to gguf-parser, and the
308 whole thing runs while this process holds the cross-process build lock that
309 every other lilbee start waits on without a deadline. A descending scan paid
310 for all of ``_CHAT_SLOTS - 1`` searches in exactly the tight-card case where
311 none of them fit. Bisection is sound here because the fit is non-increasing
312 in the slot count: more sequences divide the same headroom, so once a count
313 fails no larger one can succeed.
314 """
315 full = fit_fn(1)
316 if full <= model_cache._DYNAMIC_CTX_FLOOR:
317 return 1, full
318 low, high = 1, _CHAT_SLOTS
319 while low < high:
320 mid = (low + high + 1) // 2
321 if fit_fn(mid) >= full:
322 low = mid
323 else:
324 high = mid - 1
325 return low, full
328def _resolve_chat_slots(
329 model_path: Path,
330 ctx: int,
331 *,
332 mmproj_path: Path | None = None,
333 unified_budget: int | None = None,
334 chat_reservation: int = 0,
335 device: FleetDevice | None = None,
336) -> int:
337 """Largest chat slot count (<= ``_CHAT_SLOTS``) whose footprint fits the budget
338 after reserving the search roles; steps to 1 when none fit.
340 The budget is the whole serve budget less the search roles' measured
341 footprint. ``cfg.gpu_memory_fraction`` is already the margin held back from
342 the card, and the room for co-located embed/rerank is ``chat_reservation``,
343 which is what those servers were sized at rather than a flat share. A second
344 fraction on top charged that room twice and left a fifth of the card unused
345 on a box with no search roles at all.
346 """
347 budget = _slot_budget(unified_budget, device) - chat_reservation
348 return _fit_slots(
349 _CHAT_SLOTS,
350 WorkerRole.CHAT,
351 model_path,
352 ctx,
353 mmproj_path=mmproj_path,
354 unified=unified_budget is not None,
355 budget=budget,
356 )
359def _resolve_vision_slots(
360 model_path: Path,
361 ctx: int,
362 *,
363 mmproj_path: Path | None = None,
364 unified_budget: int | None = None,
365 device: FleetDevice | None = None,
366) -> int:
367 """Largest OCR batching slot count (<= ``cfg.vision_ocr_concurrency``) that fits
368 the memory budget; 1 when the ceiling is 1 or nothing larger fits."""
369 from lilbee.core.config import cfg
371 ceiling = max(1, cfg.vision_ocr_concurrency)
372 if ceiling == 1:
373 return 1
374 return _fit_slots(
375 ceiling,
376 WorkerRole.VISION,
377 model_path,
378 ctx,
379 mmproj_path=mmproj_path,
380 unified=unified_budget is not None,
381 budget=_slot_budget(unified_budget, device, vram_fraction=_VISION_VRAM_FRACTION),
382 )
385def _resolve_llm_rerank_slots(
386 model_path: Path,
387 ctx: int,
388 *,
389 unified_budget: int | None = None,
390 device: FleetDevice | None = None,
391) -> int:
392 """Largest LLM-reranker slot count (<= ``LLM_RERANK_CONCURRENCY``) that fits the
393 memory budget; 1 when nothing larger fits. Matches the client's request fan-out."""
394 return _fit_slots(
395 LLM_RERANK_CONCURRENCY,
396 WorkerRole.RERANK,
397 model_path,
398 ctx,
399 mmproj_path=None,
400 unified=unified_budget is not None,
401 budget=_slot_budget(unified_budget, device, vram_fraction=_LLM_RERANK_VRAM_FRACTION),
402 rerank_mode=RerankMode.LLM,
403 )
406def _slot_budget(
407 unified_budget: int | None,
408 device: FleetDevice | None = None,
409 *,
410 vram_fraction: float = 1.0,
411) -> int:
412 """Memory budget for slot sizing: the usable memory on *device* (the fleet's
413 smallest when placement has not chosen one yet), capped by ``unified_budget``
414 (free system RAM) when there is no discrete GPU so the count steps down to fit
415 free memory instead of overcommitting.
417 *vram_fraction* takes less than the whole for a role that shares its card by
418 design. Chat takes all of it and subtracts what the search roles were sized
419 at, which is the same room stated as a measurement rather than a share."""
420 budget = int(plan_sizing_budget(device) * vram_fraction)
421 if unified_budget is not None:
422 budget = min(budget, unified_budget)
423 return budget
426def _fit_slots(
427 ceiling: int,
428 role: WorkerRole,
429 model_path: Path,
430 ctx: int,
431 *,
432 mmproj_path: Path | None,
433 unified: bool,
434 budget: int,
435 rerank_mode: RerankMode | None = None,
436) -> int:
437 """Largest slot count in ``1..ceiling`` whose instance footprint fits *budget*;
438 1 when none larger fit."""
439 from lilbee.providers.base import ProviderError
441 for slots in range(ceiling, 1, -1):
442 try:
443 est = estimate_instance_footprint(
444 model_path,
445 ctx=ctx,
446 slots=slots,
447 gpu_layers=_role_gpu_layers(role),
448 flash_attn=_role_flash(role, rerank_mode),
449 kv_cache_type=_role_kv_cache_type(role),
450 kv_cache_type_v=_role_kv_cache_type_v(role),
451 mmproj_path=mmproj_path,
452 expert_offload=_role_expert_offload(model_path),
453 )
454 except (ProviderError, OSError):
455 # An unsizable model runs a single slot; the load decides the rest.
456 return 1
457 if est.footprint(unified=unified) <= budget:
458 return slots
459 return 1
462def fit_chat_ctx(
463 model_path: Path,
464 meta: dict[str, str] | None,
465 *,
466 available_bytes: int,
467 ctx_ceiling: int,
468) -> engine_params.ChatFit:
469 """The offload and chat window gguf-parser says *available_bytes* backs.
471 Sizes one slot at the configured offload, as :func:`_slots_for` then grows
472 the slot count against what the window leaves. A window that cannot hold a
473 grounded prompt buys KV room by leaving layers in system memory, and the
474 largest offload whose window reaches ``min_usable_chat_ctx`` wins; the
475 alternative for those models is no service at all
476 (:func:`_unusable_chat_ctx_reason`). Raises when the estimator cannot
477 answer, which sends :func:`engine_params.resolve_chat_fit` to its
478 header-math fallback.
479 """
481 def fit_at(gpu_layers: int) -> int:
482 return fleet_ctx.fit_single_ctx(
483 model_path,
484 meta=meta,
485 slots=1,
486 available_bytes=available_bytes,
487 gpu_layers=gpu_layers,
488 flash_attn=_role_flash(WorkerRole.CHAT),
489 kv_cache_type=_role_kv_cache_type(WorkerRole.CHAT),
490 kv_cache_type_v=_role_kv_cache_type_v(WorkerRole.CHAT),
491 unified=plan_sizing_is_unified(),
492 ctx_ceiling=ctx_ceiling,
493 expert_offload=_role_expert_offload(model_path),
494 )
496 configured = _role_gpu_layers(WorkerRole.CHAT)
497 fit = engine_params.ChatFit(configured, fit_at(configured))
498 needed = engine_params.min_usable_chat_ctx()
499 if fit.ctx >= needed or not _chat_offload_is_tradable(
500 model_path,
501 meta=meta,
502 available_bytes=available_bytes,
503 ctx_ceiling=ctx_ceiling,
504 needed=needed,
505 ):
506 return fit
507 traded = _traded_chat_fit(
508 fit_at, upper=_chat_offload_probe_ceiling(meta, configured), needed=needed
509 )
510 return traded or fit
513# Architectures whose head scores every position it proposes. Their compute
514# buffer grows with the window, and the estimator prices it by the physical
515# batch instead, so the estimate falls further behind the longer the window
516# gets. The trade maximises exactly that window, which turns a fixed
517# under-estimate into an unbounded one: measured on a 1 GiB card, the trade
518# granted gemma-4-26B-A4B-eagle3 59648 tokens against a real 1413 MiB.
519_NON_TRADABLE_ARCHITECTURES = frozenset({"eagle3"})
522def _chat_offload_is_tradable(
523 model_path: Path,
524 *,
525 meta: dict[str, str] | None,
526 available_bytes: int,
527 ctx_ceiling: int,
528 needed: int,
529) -> bool:
530 """Whether a too-small chat window may buy KV room by moving layers off the card.
532 Only while the weights alone fit *available_bytes*. Past that the layers the
533 trade leaves behind are served from system memory, and a model that answers
534 from host RAM is the unusable load :func:`_unusable_chat_ctx_reason` refuses
535 on purpose. A window the user capped below *needed* is their choice, and the
536 refusal already honours it, so nothing is bought by trading for it either.
538 Never for a speculator head. The estimate the search reads is not sound over
539 the window for those, so the search optimises a number that grows wrong.
540 """
541 if (meta or {}).get("architecture") in _NON_TRADABLE_ARCHITECTURES:
542 return False
543 return ctx_ceiling >= needed and _weights_bytes(model_path) <= available_bytes
546def _chat_offload_probe_ceiling(meta: dict[str, str] | None, configured: int) -> int:
547 """Highest layer count the offload search probes, below *configured*.
549 ``_NO_OFFLOAD_PROBE`` when the architecture does not report a layer count,
550 which leaves no grid to search.
551 """
552 try:
553 layers = int((meta or {})["block_count"])
554 except (KeyError, ValueError):
555 return _NO_OFFLOAD_PROBE
556 if configured == engine_params.N_GPU_LAYERS_AUTO:
557 return layers
558 return min(configured - 1, layers)
561def _traded_chat_fit(
562 fit_at: Callable[[int], int], *, upper: int, needed: int
563) -> engine_params.ChatFit | None:
564 """Largest offload at or below *upper* whose window reaches *needed*, or ``None``.
566 Binary search: the estimate charges less VRAM the fewer layers the card
567 holds, so the fitted window is monotone in the offload.
568 """
569 lo, hi, best = 0, upper, None
570 while lo <= hi:
571 mid = (lo + hi) // 2
572 ctx = fit_at(mid)
573 if ctx >= needed:
574 best, lo = engine_params.ChatFit(mid, ctx), mid + 1
575 else:
576 hi = mid - 1
577 return best
580def _role_ctx(
581 role: WorkerRole,
582 model_path: Path,
583 meta: dict[str, str] | None,
584 device: FleetDevice | None = None,
585) -> int:
586 """Per-slot context for a role, derived as the in-process loader does.
588 Embed/rerank use the embedding model's training context; vision uses the
589 vision loader's training-context picker; chat honors ``cfg.num_ctx`` then
590 falls back to the single-GPU dynamic chat-ctx picker, sized against *device*
591 once placement has chosen one. A tensor-split chat is sized against its
592 per-device headroom instead (see :func:`fit_split_ctx`).
593 """
594 from lilbee.core.config import cfg
596 if role is WorkerRole.EMBED:
597 return engine_params.resolve_embed_ctx(meta, model_path)
598 if role is WorkerRole.RERANK:
599 if _rerank_mode_for(meta) is RerankMode.LLM:
600 return engine_params.resolve_llm_rerank_ctx(meta, model_path)
601 return engine_params.resolve_embed_ctx(meta, model_path)
602 if role is WorkerRole.VISION:
603 return engine_params.resolve_vision_ctx(model_path)
604 if cfg.num_ctx is not None:
605 return _pinned_chat_ctx(model_path, meta)
606 return engine_params.resolve_chat_ctx(
607 model_path, meta, available_bytes=plan_sizing_budget(device)
608 )
611def _pinned_chat_ctx(model_path: Path, meta: dict[str, str] | None) -> int:
612 """``cfg.num_ctx``, clamped to what the model was trained for.
614 Every unpinned resolver already clamps, and both docstrings here claimed the
615 pin did too. It did not, so a pin past the trained window was passed straight
616 to the engine, which clamps it silently and serves a different number than
617 every budget was sized for.
619 Only against a window that is actually known. A GGUF whose header cannot be
620 read falls back to a default that is a guess, and contradicting an explicit
621 pin with a guess would break the hosts where the header is the thing that is
622 broken.
623 """
624 from lilbee.core.config import cfg
626 pinned = cfg.num_ctx
627 assert pinned is not None # noqa: S101 - callers check; this documents the contract
628 ceiling = _known_chat_ceiling(model_path, meta)
629 if ceiling is None or pinned <= ceiling:
630 return pinned
631 log.warning(
632 "num_ctx is set to %d but %s was trained for %d, so %d is what will be served. "
633 "Lower num_ctx to stop planning against a window this model does not have.",
634 pinned,
635 model_path.name,
636 ceiling,
637 ceiling,
638 )
639 return ceiling
642def _known_chat_ceiling(model_path: Path, meta: dict[str, str] | None) -> int | None:
643 """The largest chat window this model is known to support, or ``None``.
645 ``None`` when the GGUF header gave no usable context length and the user set
646 no ``cfg.num_ctx_max``: there is then no measured ceiling, only a default.
647 """
648 from lilbee.core.config import cfg
649 from lilbee.providers.gguf_meta import train_ctx_from_meta
651 sentinel = -1
652 trained = train_ctx_from_meta(meta, fallback=sentinel, model_path=model_path)
653 known = [value for value in (trained, cfg.num_ctx_max) if value is not None and value > 0]
654 return min(known) if known else None
657def _rerank_mode_for(meta: dict[str, str] | None) -> RerankMode:
658 """Resolve the RERANK serving mode from cfg + the reranker GGUF arch."""
659 from lilbee.core.config import cfg
661 arch = meta.get("architecture") if meta else None
662 return resolve_rerank_mode(cfg.reranker_type, arch)
665def _role_rerank_mode(role: WorkerRole, meta: dict[str, str] | None) -> RerankMode | None:
666 """The RERANK serving mode for *role*, or ``None`` for every other role."""
667 return _rerank_mode_for(meta) if role is WorkerRole.RERANK else None
670def _server_spec(
671 role: WorkerRole, rerank_mode: RerankMode | None, meta: dict[str, str] | None
672) -> RoleServerSpec:
673 """The llama-server spec for a launch: rerank mode, decoder-aware embed pooling,
674 or the role default. EMBED forces ``--pooling last`` for decoder-only archs."""
675 if rerank_mode is not None:
676 return rerank_spec(rerank_mode)
677 if role is WorkerRole.EMBED:
678 return embed_spec(meta)
679 return ROLE_SPECS[role]
682def _pooled_batch_size(role: WorkerRole, rerank_mode: RerankMode | None, ctx: int) -> int | None:
683 """The ``--batch-size``/``--ubatch-size`` the launch raises for pooled
684 embed/cross-encoder rerank (the full context), or ``None`` for other roles."""
685 if role in _EMBED_ROLES and rerank_mode is not RerankMode.LLM:
686 return ctx
687 return None
690def _role_gpu_layers(role: WorkerRole) -> int:
691 """GPU-layer offload: chat honors ``cfg.n_gpu_layers``, others offload all layers."""
693 return engine_params.resolve_n_gpu_layers(embedding=role in _ALL_LAYER_ROLES)
696def _flash_enabled() -> bool:
697 """Flash attention is on unless ``cfg.flash_attention`` is explicitly ``False``."""
698 from lilbee.core.config import cfg
700 return cfg.flash_attention is not False
703def probed_devices() -> tuple[FleetDevice, ...]:
704 """Devices the engine enumerated, empty when they could not be read.
706 Prefers the plan snapshot so a whole planning pass answers consistently, and
707 falls back to the short-TTL read cache rather than a fresh probe.
708 """
709 probe = _current_plan_probe()
710 if probe is not None:
711 return probe.devices
712 try:
713 return tuple(_read_device_cache.get(resolve_llama_server()).devices)
714 except (ProviderError, OSError):
715 return ()
718def engine_backend() -> EngineBackend:
719 """The backend the engine selected, UNKNOWN when it could not be asked.
721 Prefers the plan snapshot so a whole planning pass answers consistently, and
722 reads it exactly as :func:`probed_devices` does, through the restating reader:
723 a snapshot taken against another engine binary re-probes here too, so the
724 answer does not depend on which reader a caller happens to reach first.
726 The placement route reports the backend of the reading its device list came
727 from instead (:func:`resolve_placement_plan`), so one payload never pairs a
728 live device list with a snapshot's backend.
729 """
730 probe = _current_plan_probe()
731 if probe is not None:
732 return probe.backend
733 try:
734 return _read_device_cache.get(resolve_llama_server()).backend
735 except (ProviderError, OSError):
736 return EngineBackend.UNKNOWN
739def _fleet_backend() -> str | None:
740 """The engine backend this host plans onto, or ``None`` when unknown."""
741 return next((device.backend for device in probed_devices()), None)
744def _flash_attention_is_trusted() -> bool:
745 """Whether to ask for flash attention outright rather than let the engine decide.
747 Unknown backends answer yes, which keeps every host that works today on the
748 argv it has now; only the backends known to lag get the engine's own auto.
749 """
750 backend = _fleet_backend()
751 return backend is None or backend in _TRUSTED_FLASH_BACKENDS
754def flash_attn_flag() -> str:
755 """``--flash-attn`` argv value for chat and vision."""
756 if not _flash_enabled():
757 return _FLASH_OFF
758 return _FLASH_ON if _flash_attention_is_trusted() else _FLASH_AUTO
761def _role_launches_with_flash(role: WorkerRole, rerank_mode: RerankMode | None = None) -> bool:
762 """Whether the launch asks the engine for flash attention on *role*.
764 The one place that answers this. The registry marks RERANK as a non-flash
765 role because a cross-encoder pools in one batch, but an LLM reranker is
766 generative and launches exactly like chat, so the mode decides there.
767 """
768 if role is WorkerRole.RERANK:
769 return rerank_mode is RerankMode.LLM
770 return role in _FLASH_ROLES
773def _role_flash(role: WorkerRole, rerank_mode: RerankMode | None = None) -> bool:
774 """Whether the estimate may assume flash attention for *role*.
776 The launch's own answer, narrowed to a definite ``on``. Under ``auto`` the
777 engine decides at load time, and assuming it would size the KV cache below
778 what the launch may need.
779 """
780 return _role_launches_with_flash(role, rerank_mode) and flash_attn_flag() == _FLASH_ON
783def _role_kv_cache_type(role: WorkerRole) -> KvCacheType:
784 """Chat honors ``cfg.kv_cache_type``; embed/rerank/vision run f16 KV."""
785 from lilbee.core.config import cfg
787 return cfg.kv_cache_type if role is WorkerRole.CHAT else KvCacheType.F16
790def _replica_count(role: WorkerRole, device_count: int) -> int:
791 """Requested data-parallel instances for *role* via the shared resolver."""
792 return resolve_replica_count(role, device_count)
795def _role_kv_cache_type_v(role: WorkerRole) -> KvCacheType:
796 """The V cache type for *role*: the configured one only when flash attention is on.
798 llama.cpp refuses a quantized V cache without flash attention ("V cache
799 quantization requires flash_attn") and the server never starts, while a
800 quantized K cache needs nothing. So V follows the setting only where flash
801 attention is certain, and is f16 under ``auto`` or ``off``. That costs memory
802 rather than a launch, and the estimate moves with it.
803 """
804 from lilbee.core.config.enums import KvCacheType
806 configured = _role_kv_cache_type(role)
807 return configured if flash_attn_flag() == _FLASH_ON else KvCacheType.F16
810def chat_cache_type_flags() -> tuple[str | None, str | None]:
811 """``(--cache-type-k, --cache-type-v)`` for chat; ``None`` leaves the f16 default."""
812 from lilbee.core.config.enums import KvCacheType
814 def flag(kind: KvCacheType) -> str | None:
815 return None if kind is KvCacheType.F16 else kind.value
817 return flag(_role_kv_cache_type(WorkerRole.CHAT)), flag(_role_kv_cache_type_v(WorkerRole.CHAT))
820def _vision_mmproj(model_ref: str) -> Path | None:
821 """Resolve a vision model's mmproj sidecar, or ``None`` if absent."""
822 from lilbee.providers.base import ProviderError
823 from lilbee.providers.gguf_meta import find_mmproj_for_model
825 try:
826 return find_mmproj_for_model(engine_params.resolve_model_path(model_ref))
827 except (ProviderError, OSError, ValueError, KeyError):
828 return None
831def _estimate_role(
832 role: WorkerRole,
833 model_ref: str,
834 *,
835 slots: int | None = None,
836 unified_budget: int | None = None,
837 chat_reservation: int = 0,
838 device_count: int = 0,
839) -> ModelPlacementInput:
840 """Estimate one role-model's footprint via gguf-parser (+ mmproj for vision).
842 ``slots`` defaults to the role's resolved batching slots (chat and vision are
843 memory-aware); ``chat_reservation`` shrinks chat to leave room for the search
844 roles; ``device_count`` resolves an auto (0) replica knob to one per GPU.
845 Charges the unified footprint with no discrete GPU, else the VRAM one.
846 """
847 from lilbee.providers.gguf_meta import read_gguf_metadata
849 path = engine_params.resolve_model_path(model_ref)
850 mmproj = _vision_mmproj(model_ref) if role is WorkerRole.VISION else None
851 meta = read_gguf_metadata(path)
852 # Size the single-instance footprint against the placement reserve (a usable
853 # KV floor for chat), so a model that fits weights-only on one card but not
854 # weights + a usable context falls through to a tensor-split instead of being
855 # single-carded into a tiny n_ctx. Non-chat roles keep their launch ctx.
856 ctx = _placement_estimate_ctx(role, path, meta)
857 rerank_mode = _role_rerank_mode(role, meta)
858 if slots is None:
859 slots = _slots_for(
860 role,
861 path,
862 ctx,
863 mmproj_path=mmproj,
864 unified_budget=unified_budget,
865 chat_reservation=chat_reservation,
866 rerank_mode=rerank_mode,
867 )
868 est = estimate_instance_footprint(
869 path,
870 ctx=ctx,
871 slots=slots,
872 gpu_layers=_role_gpu_layers(role),
873 flash_attn=_role_flash(role, rerank_mode),
874 kv_cache_type=_role_kv_cache_type(role),
875 kv_cache_type_v=_role_kv_cache_type_v(role),
876 mmproj_path=mmproj,
877 batch_size=_pooled_batch_size(role, rerank_mode, ctx),
878 expert_offload=_role_expert_offload(path),
879 )
880 fp = est.footprint(unified=unified_budget is not None)
881 if role is WorkerRole.CHAT and unified_budget is None:
882 fp = _chat_serve_budget_footprint(fp)
883 return ModelPlacementInput(
884 role=role,
885 est_vram_bytes=fp,
886 replicas=_replica_count(role, device_count),
887 est_ram_bytes=est.ram_bytes,
888 )
891def _chat_serve_budget_footprint(footprint: int) -> int:
892 """Charge a chat instance against the serve budget, not the placement headroom.
894 The planner fits instances within ``cfg.usable_vram_fraction`` of a card, but a
895 single-card chat then sizes its KV cache against the smaller
896 ``cfg.gpu_memory_fraction`` budget (``resolve_chat_ctx``). A model that fills a
897 card at 0.9 leaves no room for KV at 0.75 and collapses to a few hundred tokens,
898 so scale its placement footprint by the budget ratio: it then needs a
899 tensor-split (pooling VRAM across cards) whenever single-carding it would starve
900 its context. Small models are unaffected -- they fit the serve budget with KV
901 room to spare.
902 """
903 from lilbee.core.config import cfg
905 # Never below 1.0. The ratio only compensates while the serve budget is the
906 # smaller of the two; a gpu_memory_fraction raised past the usable fraction
907 # inverts it, and the same line that exists to charge chat more starts
908 # charging it less than the model takes.
909 return int(footprint * max(1.0, usable_vram_fraction() / cfg.gpu_memory_fraction))
912def _placement_estimate_ctx(role: WorkerRole, model_path: Path, meta: dict[str, str] | None) -> int:
913 """Per-slot context the placement estimate sizes a role against.
915 For chat this reserves KV for a usable floor (``_MIN_USABLE_CHAT_CTX``, or the
916 user's ``cfg.num_ctx`` pin), capped by the model's trained ceiling -- not the
917 single-GPU dynamic ctx (which shrinks to fit one card and then confirms a
918 single-card placement) nor the full trained ceiling (which over-reserves). A
919 model that cannot hold weights + this floor on one card is tensor-split.
920 """
921 from lilbee.core.config import cfg
923 if role is WorkerRole.CHAT:
924 if cfg.num_ctx is not None:
925 return _pinned_chat_ctx(model_path, meta)
926 return apply_ctx_downshift(
927 role,
928 min(
929 engine_params.chat_ctx_ceiling(meta, model_path),
930 max(cfg.chat_n_ctx_target, _MIN_USABLE_CHAT_CTX),
931 ),
932 )
933 return apply_ctx_downshift(role, _role_ctx(role, model_path, meta))
936def _placement_estimate_slots(role: WorkerRole, meta: dict[str, str] | None) -> int:
937 """The slot count the placement estimate reserves KV for.
939 A tensor-split chat reserves one full-context sequence here: a conservative
940 floor for the card-count decision. The launch then fills the placed cards'
941 real headroom with as many full-context slots as fit (``_resolve_split_chat_slots``),
942 never exceeding what those cards hold, so a larger launch count can't OOM.
943 """
944 from lilbee.core.config import cfg
946 if role is WorkerRole.CHAT:
947 return _SPLIT_CHAT_SLOTS
948 if role is WorkerRole.VISION:
949 return max(1, cfg.vision_ocr_concurrency)
950 if role is WorkerRole.RERANK and _rerank_mode_for(meta) is RerankMode.LLM:
951 return LLM_RERANK_CONCURRENCY
952 return _AUX_SLOTS
955def _peak_estimator(model_refs: dict[WorkerRole, str]) -> PeakEstimator:
956 """Per-device VRAM-vector estimator for the planner, bound to the configured models.
958 Estimates each role at its launch ceiling (ctx x slots) with the candidate
959 tensor-split ratio, so the planner reserves enough cards for the busiest one.
960 """
961 from lilbee.providers.gguf_meta import read_gguf_metadata
963 def estimate_peak(role: WorkerRole, ratio: tuple[int, ...]) -> tuple[int, ...]:
964 path = engine_params.resolve_model_path(model_refs[role])
965 meta = read_gguf_metadata(path)
966 mmproj = _vision_mmproj(model_refs[role]) if role is WorkerRole.VISION else None
967 slots = _placement_estimate_slots(role, meta)
968 ctx = _placement_estimate_ctx(role, path, meta)
969 rerank_mode = _role_rerank_mode(role, meta)
970 est = estimate_instance_footprint(
971 path,
972 ctx=ctx,
973 slots=slots,
974 gpu_layers=_role_gpu_layers(role),
975 flash_attn=_role_flash(role, rerank_mode),
976 kv_cache_type=_role_kv_cache_type(role),
977 kv_cache_type_v=_role_kv_cache_type_v(role),
978 mmproj_path=mmproj,
979 tensor_split=ratio,
980 batch_size=_pooled_batch_size(role, rerank_mode, ctx),
981 expert_offload=_role_expert_offload(path),
982 )
983 return est.per_device_vram
985 return estimate_peak
988def _chat_split_ctx_objective(
989 model_refs: dict[WorkerRole, str],
990) -> tuple[SplitCtxFitter | None, int]:
991 """The chat split's context fitter and target, or ``(None, 0)`` with no chat model.
993 The fitter sizes a candidate shard's served context exactly as the launch does
994 (:func:`fit_split_ctx`), so the planner widens chat onto idle cards only when a
995 tighter shard would starve KV below the target. See docs/architecture.md.
996 """
997 if WorkerRole.CHAT not in model_refs:
998 return None, 0
999 from lilbee.providers.gguf_meta import read_gguf_metadata
1001 path = engine_params.resolve_model_path(model_refs[WorkerRole.CHAT])
1002 meta = read_gguf_metadata(path)
1003 target = _placement_estimate_ctx(WorkerRole.CHAT, path, meta)
1005 def fit(ratio: tuple[int, ...], per_device_free_bytes: Sequence[int]) -> int:
1006 return fleet_ctx.fit_split_ctx(
1007 path,
1008 meta=meta,
1009 slots=_SPLIT_CHAT_SLOTS,
1010 ratio=ratio,
1011 per_device_free_bytes=per_device_free_bytes,
1012 gpu_layers=_role_gpu_layers(WorkerRole.CHAT),
1013 flash_attn=_role_flash(WorkerRole.CHAT),
1014 kv_cache_type=_role_kv_cache_type(WorkerRole.CHAT),
1015 kv_cache_type_v=_role_kv_cache_type_v(WorkerRole.CHAT),
1016 ctx_ceiling=target,
1017 expert_offload=_role_expert_offload(path),
1018 )
1020 return fit, target
1023def _search_reservation(inputs: dict[WorkerRole, ModelPlacementInput]) -> int:
1024 """Total footprint of the placed search roles (all replicas), held back ahead
1025 of chat."""
1026 return sum(
1027 inputs[role].est_vram_bytes * inputs[role].replicas
1028 for role in _EMBED_ROLES
1029 if role in inputs
1030 )
1033def _role_weights_bytes(role: WorkerRole, ref: str) -> int:
1034 """The model's weight bytes on disk (plus the mmproj for vision): a
1035 ground-truth lower bound on residency. 0 when the file cannot be resolved."""
1036 from lilbee.providers.base import ProviderError
1038 try:
1039 size = _weights_bytes(engine_params.resolve_model_path(ref))
1040 if role is WorkerRole.VISION:
1041 mmproj = _vision_mmproj(ref)
1042 if mmproj is not None:
1043 size += int(mmproj.stat().st_size)
1044 except (ProviderError, OSError):
1045 return 0
1046 return size
1049def _is_moe(meta: dict[str, str] | None) -> bool:
1050 """Whether the GGUF declares routed experts, so its experts can be offloaded."""
1051 count = (meta or {}).get("expert_count")
1052 try:
1053 return int(count) > 0 if count is not None else False
1054 except ValueError:
1055 return False
1058def expert_offload_all(meta: dict[str, str] | None) -> bool:
1059 """Whether to keep every layer's experts in system memory; MoE models only."""
1060 from lilbee.core.config import cfg
1062 return bool(cfg.cpu_moe) and _is_moe(meta)
1065def expert_offload_layers(meta: dict[str, str] | None) -> int | None:
1066 """How many layers' experts to keep in system memory, or None for no split.
1068 A non-positive ``n_cpu_moe`` offloads nothing (it would emit a no-op
1069 ``--n-cpu-moe 0``), so it reads as unset.
1070 """
1071 from lilbee.core.config import cfg
1073 if cfg.n_cpu_moe is None or cfg.n_cpu_moe < 1 or not _is_moe(meta):
1074 return None
1075 return cfg.n_cpu_moe
1078def _role_expert_offload(model_path: Path) -> tuple[str, ...]:
1079 """Expert patterns the launch will offload, for sizing the same way it runs.
1081 Reads the GGUF (cached) rather than taking metadata as an argument so every
1082 estimate site charges the same tensors the launch moves off the GPU.
1083 """
1084 from lilbee.providers.fleet.adapters import expert_offload_patterns
1085 from lilbee.providers.gguf_meta import read_gguf_metadata
1087 meta = read_gguf_metadata(model_path)
1088 return expert_offload_patterns(
1089 cpu_moe=expert_offload_all(meta), n_cpu_moe=expert_offload_layers(meta)
1090 )
1093def _expert_offload_configured() -> bool:
1094 """Whether the user asked for expert offload that would actually take effect.
1096 A non-positive ``n_cpu_moe`` offloads nothing, so it does not count.
1097 """
1098 from lilbee.core.config import cfg
1100 return bool(cfg.cpu_moe) or (cfg.n_cpu_moe is not None and cfg.n_cpu_moe >= 1)
1103def _weights_exceed_everything(size: int, *, total_vram: int, total_ram: int) -> bool:
1104 """True when a model's weights fit neither the GPUs nor system memory.
1106 File size is ground truth, not an estimate, so this bound cannot repeat the
1107 false-refusal class: no estimator error makes a 40 GiB file fit a 1 GiB box.
1108 Past both pools there is nowhere for a layer to go and no launch can win, so
1109 saying so beats a load that thrashes and then dies.
1110 """
1111 ceiling = total_vram + total_ram
1112 return ceiling > 0 and size > ceiling
1115def _weights_exceed_hardware(size: int, total_vram: int, *, is_moe: bool) -> bool:
1116 """True when this model cannot be served on this machine at all.
1118 Exceeding VRAM alone is not that. The engine chooses how many layers fit and
1119 keeps the rest in system memory, so a model larger than every card is a
1120 partial offload and lilbee's job is to launch it and say what will happen.
1121 Refusing there meant the fit never ran and the role was skipped, which left
1122 the user hand-tuning n_gpu_layers to get back what the engine does by itself.
1124 What still refuses is a model past VRAM and system memory together, where no
1125 arrangement of layers exists. A user-set n_gpu_layers or expert offload keeps
1126 standing the bound down entirely, since the user has said where the weights
1127 should go.
1128 """
1129 from lilbee.core.config import cfg
1131 if cfg.n_gpu_layers is not None:
1132 return False
1133 if is_moe and _expert_offload_configured():
1134 return False
1135 return _weights_exceed_everything(
1136 size, total_vram=total_vram, total_ram=model_cache.total_system_memory()
1137 )
1140def _vision_without_mmproj(role: WorkerRole, ref: str) -> bool:
1141 """True (with a warning) for a configured vision model whose mmproj is missing.
1143 The skip would silently disable OCR; the warning names the cause and the fix.
1144 """
1145 if role is not WorkerRole.VISION or _vision_mmproj(ref) is not None:
1146 return False
1147 log.warning(
1148 "Vision model %s has no mmproj (CLIP projector); OCR is disabled. "
1149 "Re-run 'lilbee model pull %s' to fetch the projector.",
1150 ref,
1151 ref,
1152 )
1153 return True
1156def _estimate_or_fallback(
1157 role: WorkerRole,
1158 ref: str,
1159 *,
1160 unified_budget: int | None,
1161 chat_reservation: int,
1162 device_count: int,
1163 total_vram: int,
1164 skipped_not_installed: dict[WorkerRole, str],
1165 host_committed: int = 0,
1166) -> ModelPlacementInput | None:
1167 """Size *role* for placement, degrading rather than refusing.
1169 A missing model is skipped and recorded; a sizing failure on an installed
1170 model falls back to the analytic floor; weights alone exceeding the physical
1171 VRAM refuse with a plain message (ground truth, not an estimate).
1172 """
1173 from lilbee.providers.base import ProviderError, ProviderErrorKind
1175 try:
1176 estimate = _estimate_role(
1177 role,
1178 ref,
1179 unified_budget=unified_budget,
1180 chat_reservation=chat_reservation,
1181 device_count=device_count,
1182 )
1183 except (ProviderError, OSError) as exc:
1184 if isinstance(exc, ProviderError) and exc.kind is ProviderErrorKind.NOT_FOUND:
1185 log.warning("Skipping %s server: model %r is not installed.", role.value, ref)
1186 skipped_not_installed[role] = ref
1187 return None
1188 return _sizing_failure_fallback(
1189 role,
1190 ref,
1191 exc,
1192 device_count=device_count,
1193 total_vram=total_vram,
1194 host_committed=host_committed,
1195 )
1196 return _admit_estimate(
1197 _floor_implausible_estimate(estimate, role, ref),
1198 role,
1199 ref,
1200 total_vram=total_vram,
1201 ram_bytes=estimate.est_ram_bytes,
1202 host_committed=host_committed,
1203 )
1206def _admit_estimate(
1207 estimate: ModelPlacementInput,
1208 role: WorkerRole,
1209 ref: str,
1210 *,
1211 total_vram: int,
1212 ram_bytes: int,
1213 host_committed: int = 0,
1214) -> ModelPlacementInput | None:
1215 """*estimate*, or ``None`` when this model cannot load on this machine.
1217 Two hardware bounds, one per kind of memory: the weights must fit the GPUs
1218 unless something offloads, and whatever offloading puts in system memory must
1219 fit the system.
1220 """
1221 weights = _role_weights_bytes(role, ref)
1222 if _weights_exceed_hardware(weights, total_vram, is_moe=_ref_is_moe(ref)):
1223 _warn_weights_exceed(role, ref, weights, total_vram)
1224 return None
1225 if total_vram > 0 and weights > total_vram:
1226 _warn_weights_spill(role, ref, weights, total_vram)
1227 if _host_memory_refuses(role, ref, ram_bytes, host_committed):
1228 return None
1229 return estimate
1232def _analytic_footprint_floor(
1233 weights: int, *, role: WorkerRole, meta: dict[str, str] | None, ctx: int, slots: int
1234) -> int:
1235 """The least this instance can occupy: weights, its KV cache, and overhead.
1237 Used when the estimator cannot answer. Charging weight bytes alone was a
1238 knowing under-charge: the engine allocates a KV cache sized by context and
1239 slot count, plus compute buffers, and omitting all of it lets placement fit a
1240 model that cannot fit. The comment said the load would decide, and it did, by
1241 running out of memory.
1243 A floor rather than an estimate. It is derived from the header the same way
1244 the in-process sizing path derives it, and it is deliberately the smallest
1245 defensible number, because refusing a model that would have fit is its own
1246 failure. Without a readable header the per-token fallback still applies:
1247 zero is the one answer that is certainly wrong.
1248 """
1250 kv_bytes = (
1251 model_cache.kv_bytes_per_token(
1252 meta,
1253 KV_CACHE_TYPE_BYTES[_role_kv_cache_type(role)],
1254 KV_CACHE_TYPE_BYTES[_role_kv_cache_type_v(role)],
1255 )
1256 * ctx
1257 * max(slots, 1)
1258 )
1259 overhead = int(weights * model_cache._BUFFER_OVERHEAD_FRACTION)
1260 return weights + kv_bytes + overhead
1263def _estimate_is_implausible(*, estimated: int, floor: int) -> bool:
1264 """Whether *estimated* describes a load that cannot exist.
1266 Below the bytes the card must hold there is no arrangement of memory that
1267 serves the model. A floor of zero means nothing could be computed to compare
1268 against, and a guess is not grounds to discard the only measurement there is.
1269 """
1270 return floor > 0 and 0 < estimated < floor
1273def _floor_implausible_estimate(
1274 estimate: ModelPlacementInput, role: WorkerRole, ref: str
1275) -> ModelPlacementInput:
1276 """*estimate*, or the model's weight bytes when it reports less than those.
1278 The bound is the weights alone, not the analytic footprint. That figure
1279 sizes a KV cache as though every layer ran dense attention over the whole
1280 window, which linear-attention, sliding-window and MLA models do not, so it
1281 sits far above what they hold. Weights are architecture-independent, which
1282 makes an estimate under them the one answer the planner can call impossible
1283 without modelling attention.
1285 Offload lifts the bound: the layers in system memory are weight bytes the
1286 card never holds.
1287 """
1288 if _cpu_offload_in_play():
1289 return estimate
1290 weights = _role_weights_bytes(role, ref)
1291 if not _estimate_is_implausible(estimated=estimate.est_vram_bytes, floor=weights):
1292 return estimate
1293 log.warning(
1294 "The estimator sized the %s model %s at %.1f GiB, below the %.1f GiB of "
1295 "weights the card has to hold. Charging the weights instead.",
1296 role.value,
1297 ref,
1298 estimate.est_vram_bytes / 1024**3,
1299 weights / 1024**3,
1300 )
1301 return replace(estimate, est_vram_bytes=weights)
1304def _sizing_failure_fallback(
1305 role: WorkerRole,
1306 ref: str,
1307 exc: Exception,
1308 *,
1309 device_count: int,
1310 total_vram: int,
1311 host_committed: int = 0,
1312) -> ModelPlacementInput | None:
1313 """Analytic-floor placement input for an installed model the estimator cannot
1314 size; ``None`` skips the role (the file is unresolvable, its weights alone
1315 exceed the hardware, or offloading it would exceed system memory).
1317 The host bound applies here too. Charging the whole floor to VRAM and
1318 skipping it let an unsizable model past a check every sized model faces."""
1319 weights = _role_weights_bytes(role, ref)
1320 if weights == 0:
1321 log.warning("Skipping %s server: could not size model %r (%s).", role.value, ref, exc)
1322 return None
1323 if _weights_exceed_hardware(weights, total_vram, is_moe=_ref_is_moe(ref)):
1324 _warn_weights_exceed(role, ref, weights, total_vram)
1325 return None
1326 floor = _fallback_floor_for(role, ref, weights)
1327 log.warning(
1328 "Could not size the %s model %s (%s). Charging %.1f GiB, its weights plus the "
1329 "cache and buffers it will allocate, which is a floor rather than an estimate: "
1330 "the load may still need more.",
1331 role.value,
1332 ref,
1333 exc,
1334 floor / 1024**3,
1335 )
1336 if _host_memory_refuses(role, ref, floor, host_committed):
1337 return None
1338 return ModelPlacementInput(
1339 role=role, est_vram_bytes=floor, replicas=_replica_count(role, device_count)
1340 )
1343def _fallback_floor_for(role: WorkerRole, ref: str, weights: int) -> int:
1344 """:func:`_analytic_footprint_floor` for *role*, reading what metadata it can."""
1345 from lilbee.providers.base import ProviderError
1346 from lilbee.providers.gguf_meta import read_gguf_metadata
1348 try:
1349 path = engine_params.resolve_model_path(ref)
1350 meta = read_gguf_metadata(path)
1351 except (ProviderError, OSError, ValueError):
1352 meta = None
1353 path = None
1354 ctx = _placement_estimate_ctx(role, path, meta) if path is not None else _MIN_USABLE_CHAT_CTX
1355 return _analytic_footprint_floor(
1356 weights, role=role, meta=meta, ctx=ctx, slots=_placement_estimate_slots(role, meta)
1357 )
1360def planned_embed_token_cap(ref: str) -> int | None:
1361 """The token cap an embed launch of *ref* gets, or None when *ref* cannot be read.
1363 Computed from the model's metadata the way the launch is, so the chunker can
1364 bound itself to the cap before the engine is up.
1365 """
1366 from lilbee.providers.base import ProviderError
1367 from lilbee.providers.gguf_meta import read_gguf_metadata
1369 try:
1370 path = engine_params.resolve_model_path(ref)
1371 meta = read_gguf_metadata(path)
1372 except (ProviderError, OSError, ValueError):
1373 return None
1374 return engine_params.embed_token_cap(
1375 apply_ctx_downshift(WorkerRole.EMBED, _role_ctx(WorkerRole.EMBED, path, meta))
1376 )
1379def _ref_is_moe(ref: str) -> bool:
1380 """Whether *ref*'s GGUF declares routed experts; False when it cannot be read."""
1381 from lilbee.providers.base import ProviderError
1382 from lilbee.providers.gguf_meta import read_gguf_metadata
1384 try:
1385 return _is_moe(read_gguf_metadata(engine_params.resolve_model_path(ref)))
1386 except (ProviderError, OSError):
1387 return False
1390def _cpu_offload_in_play() -> bool:
1391 """Whether this configuration puts any of a model's weights in system memory.
1393 Expert offload moves the experts, a partial ``n_gpu_layers`` moves whole
1394 layers, and zero moves the model. Without one of these the engine keeps
1395 everything on the card and the estimator's host figure describes memory
1396 nobody will allocate.
1397 """
1398 from lilbee.core.config import cfg
1400 return _expert_offload_configured() or cfg.n_gpu_layers is not None
1403def _host_bytes_must_be_resident(role: WorkerRole, ref: str) -> bool:
1404 """Whether *role*'s host bytes have to fit RAM rather than page in and out.
1406 The estimator's host figure counts mmap pages, and llama.cpp maps CPU-side
1407 weights over that mapping instead of allocating them, so with mmap they are
1408 evictable page cache: a model far larger than RAM streams from disk and
1409 serves, which is a practiced setup for a large mixture-of-experts. Only
1410 ``--no-mmap`` turns them into a buffered read that must be resident, and the
1411 single path that asks for it is a chat model on a network filesystem.
1413 Anything this cannot determine counts as mappable, because a false refusal
1414 here has no override and costs the user a model that would have run.
1415 """
1416 if role is not WorkerRole.CHAT:
1417 return False
1418 try:
1419 path = engine_params.resolve_model_path(ref)
1420 except (ProviderError, OSError, ValueError):
1421 return False
1422 if not is_network_path(path):
1423 return False
1424 return _chat_no_mmap(_role_weights_bytes(role, ref), on_network_fs=True)
1427def _host_committed(admitted: Mapping[WorkerRole, ModelPlacementInput]) -> int:
1428 """System-memory bytes the roles already admitted to this plan will hold."""
1429 return sum(inp.est_ram_bytes for inp in admitted.values())
1432def _host_memory_refuses(role: WorkerRole, ref: str, ram_bytes: int, committed: int) -> bool:
1433 """Whether *role*'s system-memory half is too big for this machine to load.
1435 Charged only when something actually offloads, and only when the bytes must
1436 be resident: refusing a mapped model that would have streamed from disk is a
1437 false refusal with no override, which is worse than a slow load.
1439 Measured against the whole plan, not this role alone. Every role was
1440 previously compared to the entire machine on its own, so two roles that each
1441 fit and together do not were both admitted.
1442 """
1443 if not _cpu_offload_in_play() or ram_bytes <= 0:
1444 return False
1445 wanted = committed + ram_bytes
1446 total = total_system_memory()
1447 if total and wanted > total and _host_bytes_must_be_resident(role, ref):
1448 log.warning(
1449 "The %s model %s cannot load: this plan puts %.1f GiB in system memory, which "
1450 "cannot be paged out here, and the machine has %.1f GiB in total. Use a smaller "
1451 "model, or offload less.",
1452 role.value,
1453 ref,
1454 wanted / 1024**3,
1455 total / 1024**3,
1456 )
1457 return True
1458 free = free_system_memory()
1459 if free and wanted > free:
1460 log.warning(
1461 "Offloading the %s model %s brings this plan to %.1f GiB in system memory and "
1462 "only %.1f GiB is free. It will still load; close other programs if it swaps "
1463 "or runs slowly.",
1464 role.value,
1465 ref,
1466 wanted / 1024**3,
1467 free / 1024**3,
1468 )
1469 return False
1472def _warn_weights_exceed(role: WorkerRole, ref: str, weights: int, total_vram: int) -> None:
1473 log.warning(
1474 "The %s model %s cannot load: its weights are %.1f GiB and this machine has "
1475 "%.1f GiB of GPU memory and %.1f GiB of system memory, so there is nowhere "
1476 "for its layers to go. Use a smaller model or a smaller quantization.",
1477 role.value,
1478 ref,
1479 weights / 1024**3,
1480 total_vram / 1024**3,
1481 model_cache.total_system_memory() / 1024**3,
1482 )
1485def _warn_weights_spill(role: WorkerRole, ref: str, weights: int, total_vram: int) -> None:
1486 """Say that a model larger than the GPUs will run partly in system memory."""
1487 log.warning(
1488 "The %s model %s is %.1f GiB and this machine has %.1f GiB of GPU memory, so "
1489 "the engine will keep the layers that fit on the GPU and the rest in system "
1490 "memory. It will run, and it will be slower than a model that fits.",
1491 role.value,
1492 ref,
1493 weights / 1024**3,
1494 total_vram / 1024**3,
1495 )
1498def placeable_total_vram() -> int:
1499 """Physical VRAM across all cards, for the weights-exceed placeability bound.
1501 Physical total is box-state-independent (a running incumbent doesn't skew
1502 it), so it is safe to read without a clean box. Reuses the plan probe when
1503 one is captured; otherwise probes best-effort and returns ``0`` on failure,
1504 which disables only the weights-exceed filter (its own ``total > 0`` guard).
1505 """
1506 probe = _current_plan_probe()
1507 if probe is not None:
1508 return sum(d.total_bytes for d in probe.devices)
1509 from lilbee.providers.base import ProviderError
1510 from lilbee.providers.fleet.gpu_env import apply_fleet_gpu_env
1512 try:
1513 apply_fleet_gpu_env()
1514 return sum(d.total_bytes for d in resolve_devices(resolve_llama_server()))
1515 except (ProviderError, OSError):
1516 return 0
1519def role_model_placeable(role: WorkerRole, ref: str, total_vram: int) -> bool:
1520 """Whether a fresh plan would actually serve *role* on *ref*.
1522 Mirrors the planner's own drop conditions (SDK-routed role, vision without a
1523 projector, model not installed, weights exceeding physical VRAM) using the
1524 same primitives, so the acquisition ladder binds and replaces against what
1525 an engine can serve rather than the raw config. Without this a
1526 configured-but-unplaceable role keeps bind from ever matching a running
1527 engine and restarts the shared engine on every process start.
1528 """
1529 if parse_model_ref(ref).is_remote or _vision_without_mmproj(role, ref):
1530 return False
1531 weights = _role_weights_bytes(role, ref) # 0 when not installed / unresolvable
1532 if weights == 0:
1533 return False
1534 return not _weights_exceed_hardware(weights, total_vram, is_moe=_ref_is_moe(ref))
1537def _server_model_inputs(
1538 roles: tuple[WorkerRole, ...] | None = None,
1539 *,
1540 unified_budget: int | None = None,
1541 device_count: int = 0,
1542 total_vram: int = 0,
1543) -> tuple[list[ModelPlacementInput], dict[WorkerRole, str], int, dict[WorkerRole, str]]:
1544 """Build placement inputs for the configured server roles.
1546 The search and vision roles are estimated first; chat is then sized against the
1547 budget minus the search footprint (the ``reservation``) so a large chat cannot
1548 starve embed/rerank on a shared-memory host. ``device_count`` resolves an auto
1549 replica knob to one per GPU. When *roles* is given, only those are considered.
1550 Skips an unconfigured optional role, a vision model with no resolvable mmproj
1551 projector, a role whose model is not installed on disk (returned as
1552 ``skipped_not_installed`` so a surface can say so), and a model whose weight
1553 bytes alone exceed ``total_vram`` (physically unloadable under all-GPU layers).
1554 A model the estimator cannot size is enrolled at its file size instead of
1555 skipped, so the load, not the estimator, decides.
1556 """
1557 from lilbee.core.config import cfg
1559 inputs: dict[WorkerRole, ModelPlacementInput] = {}
1560 model_refs: dict[WorkerRole, str] = {}
1561 skipped_not_installed: dict[WorkerRole, str] = {}
1563 def consider(role: WorkerRole, *, chat_reservation: int = 0) -> None:
1564 if roles is not None and role not in roles:
1565 return
1566 # Any role may be "" (unconfigured) -> skipped, so that role has no
1567 # server and no not-installed complaint.
1568 ref = str(getattr(cfg, ROLE_REGISTRY[role].config_field))
1569 if not ref:
1570 return # unconfigured optional role -> no server
1571 if parse_model_ref(ref).is_remote:
1572 return # SDK-routed role: no local server to plan, not a missing install
1573 if _vision_without_mmproj(role, ref):
1574 return # no projector -> vision can't run on a server
1575 estimate = _estimate_or_fallback(
1576 role,
1577 ref,
1578 unified_budget=unified_budget,
1579 chat_reservation=chat_reservation,
1580 device_count=device_count,
1581 total_vram=total_vram,
1582 skipped_not_installed=skipped_not_installed,
1583 host_committed=_host_committed(inputs),
1584 )
1585 if estimate is None:
1586 return
1587 inputs[role] = estimate
1588 model_refs[role] = ref
1590 # Estimate every non-chat role first so the search footprint is known, then size
1591 # chat against the remainder. The reservation only applies on a shared-memory
1592 # host; discrete GPUs pin each role to its own VRAM and pack independently.
1593 for role in ROLE_REGISTRY:
1594 if role is not WorkerRole.CHAT:
1595 consider(role)
1596 reservation = _search_reservation(inputs) if unified_budget is not None else 0
1597 consider(WorkerRole.CHAT, chat_reservation=reservation)
1599 ordered = [inputs[role] for role in ROLE_REGISTRY if role in inputs]
1600 return ordered, model_refs, reservation, skipped_not_installed
1603def _non_chat_reservation(
1604 instances: Sequence[InstancePlan],
1605 inputs: Sequence[ModelPlacementInput],
1606 co_tenants: frozenset[WorkerRole] = frozenset(),
1607) -> dict[int, int]:
1608 """Per-device VRAM the non-chat role servers occupy, keyed by device index.
1610 A tensor-split chat shard must size its KV against the headroom left after the
1611 embed/rerank/vision servers on the same card, not the card's raw free VRAM, or
1612 it over-commits and OOMs at launch. Chat is excluded because it sizes its own
1613 weights. Chat's own swap-group siblings are excluded too: they are evicted while
1614 chat is resident, so their VRAM is chat's to use. That only holds when chat is
1615 itself a co-tenant; a co-tenant group that does not include chat runs behind its
1616 own swap process and can be resident beside chat, so it is charged normally.
1617 Non-chat roles are single-device, so each charges its full footprint (once per
1618 replica) to its card.
1619 """
1620 chat_siblings = co_tenants if WorkerRole.CHAT in co_tenants else frozenset()
1621 charge_by_role = {inp.role: inp.est_vram_bytes for inp in inputs}
1622 reserved: dict[int, int] = {}
1623 for inst in instances:
1624 if inst.role is WorkerRole.CHAT or inst.role in chat_siblings:
1625 continue
1626 charge = charge_by_role[inst.role]
1627 for device in inst.devices:
1628 reserved[device] = reserved.get(device, 0) + charge
1629 return reserved
1632def _charge_by_device(
1633 chosen: tuple[FleetDevice, ...], ratio: tuple[int, ...], total: int
1634) -> dict[str, int]:
1635 """What each of *chosen* was charged, keyed by the name the engine prints.
1637 A single-card instance carries the whole charge. A split carries it in the
1638 proportions it launches with, which is what the planner decided and therefore
1639 what the engine's own report should be compared against.
1640 """
1641 from lilbee.providers.fleet.readback import device_label
1643 if total <= 0 or not chosen:
1644 return {}
1645 if len(chosen) == 1:
1646 return {device_label(chosen[0]): total}
1647 weights = ratio if len(ratio) == len(chosen) else (1,) * len(chosen)
1648 denominator = sum(weights) or len(chosen)
1649 return {
1650 device_label(device): total * weight // denominator
1651 for device, weight in zip(chosen, weights, strict=True)
1652 }
1655def _launch_for(
1656 plan: InstancePlan,
1657 model_ref: str,
1658 binary: Path,
1659 by_index: dict[int, FleetDevice],
1660 *,
1661 unified_budget: int | None = None,
1662 chat_reservation: int = 0,
1663 reserved_by_device: dict[int, int] | None = None,
1664 est_vram_bytes: int = 0,
1665 model_path: Path | None = None,
1666) -> InstanceLaunch:
1667 """Build the launch spec (argv + device-pinning env) for one planned instance."""
1668 from lilbee.providers.gguf_meta import read_gguf_metadata
1670 # The self-check holds a downloaded file rather than a configured reference,
1671 # so it hands the path over instead of asking for one to be resolved.
1672 model_path = model_path or engine_params.resolve_model_path(model_ref)
1673 weights_bytes = _weights_bytes(model_path)
1674 meta = read_gguf_metadata(model_path)
1675 from lilbee.core.config import cfg
1677 chosen = tuple(by_index[i] for i in plan.devices)
1678 # ctx and slots are sized against the card this role landed on, not the fleet's
1679 # smallest, which is all the pre-placement estimate had to go on. A role spread
1680 # over several cards has no one budget; the split chat, the only such role today,
1681 # sizes against per-device headroom below.
1682 placed_device = chosen[0] if len(chosen) == 1 else None
1683 is_chat = plan.role is WorkerRole.CHAT
1684 is_vision = plan.role is WorkerRole.VISION
1685 mmproj = _vision_mmproj(model_ref) if is_vision else None
1686 chat_on_network_fs = is_chat and is_network_path(model_path)
1687 if chat_on_network_fs and not _chat_no_mmap(weights_bytes, on_network_fs=True):
1688 log.warning(
1689 "Chat model %s is served from a network filesystem and is too large to load "
1690 "into host RAM; mmap over the network can stall the load in uninterruptible "
1691 "I/O. Stage it on local disk for a reliable load.",
1692 model_ref,
1693 )
1694 # A tensor-split chat serves one full-context sequence sized against the busiest
1695 # card's headroom. A cfg.num_ctx pin overrides the fit (handled by _role_ctx).
1696 multi_card_chat = is_chat and len(chosen) > 1
1697 split_chat = multi_card_chat and cfg.num_ctx is None
1698 if multi_card_chat and host_lacks_nvlink():
1699 log.warning(
1700 "Chat model %s is tensor-split across GPUs %s on a host without NVLink; "
1701 "generation is PCIe all-reduce bound and can be very slow. A model that fits "
1702 "on fewer cards will generate faster.",
1703 model_ref,
1704 list(plan.devices),
1705 )
1706 split_slots = _SPLIT_CHAT_SLOTS
1707 # A single-card chat is the one launch whose offload its own window fit can
1708 # lower; the split and pinned paths never run that fit.
1709 fit_offload = is_chat and not multi_card_chat and cfg.num_ctx is None
1710 n_gpu_layers = _role_gpu_layers(plan.role)
1711 if split_chat:
1712 reserved = reserved_by_device or {}
1713 # Headroom left after the embed/rerank servers on each shared card, not the
1714 # card's raw free VRAM, so the chat KV doesn't over-commit.
1715 per_device_free = [max(0, d.free_bytes - reserved.get(d.index, 0)) for d in chosen]
1717 def _split_fit(slots: int) -> int:
1718 return fleet_ctx.fit_split_ctx(
1719 model_path,
1720 meta=meta,
1721 slots=slots,
1722 ratio=plan.tensor_split,
1723 per_device_free_bytes=per_device_free,
1724 gpu_layers=_role_gpu_layers(WorkerRole.CHAT),
1725 flash_attn=_role_flash(WorkerRole.CHAT),
1726 kv_cache_type=_role_kv_cache_type(WorkerRole.CHAT),
1727 kv_cache_type_v=_role_kv_cache_type_v(WorkerRole.CHAT),
1728 ctx_ceiling=_placement_estimate_ctx(WorkerRole.CHAT, model_path, meta),
1729 expert_offload=_role_expert_offload(model_path),
1730 )
1732 split_slots, ctx = _resolve_split_chat_slots(_split_fit)
1733 elif fit_offload:
1734 # One call answers both halves. The window was sized against the memory
1735 # this offload leaves free, so a second call for the offload could pair a
1736 # window with an offload that never fitted it.
1737 chat_fit = engine_params.resolve_chat_fit(
1738 model_path, meta, available_bytes=plan_sizing_budget(placed_device)
1739 )
1740 n_gpu_layers = chat_fit.gpu_layers
1741 ctx = apply_ctx_downshift(plan.role, chat_fit.ctx)
1742 else:
1743 # Downshifted here and not only in the estimate: the role resolvers are
1744 # pure functions of model and config, so without this the retry after a
1745 # load OOM re-emits a byte-identical argv and dies the same way. The
1746 # split branch above already inherits it through its ctx_ceiling.
1747 ctx = apply_ctx_downshift(plan.role, _role_ctx(plan.role, model_path, meta, placed_device))
1748 rerank_mode = _role_rerank_mode(plan.role, meta)
1749 is_llm_rerank = rerank_mode is RerankMode.LLM
1750 # A multi-card chat runs as many full-context slots as its cards' KV headroom
1751 # holds (split_slots, one when a num_ctx pin skips the fit); other roles size
1752 # --parallel against the budget the same way the estimator did.
1753 slots = (
1754 split_slots
1755 if multi_card_chat
1756 else _slots_for(
1757 plan.role,
1758 model_path,
1759 ctx,
1760 mmproj_path=mmproj,
1761 unified_budget=unified_budget,
1762 chat_reservation=chat_reservation,
1763 rerank_mode=rerank_mode,
1764 device=placed_device,
1765 )
1766 )
1767 spec = _server_spec(plan.role, rerank_mode, meta)
1768 # Cross-encoder embed/rerank pools the whole input in one batch; an LLM reranker
1769 # is generative and uses the default batching plus flash attention.
1770 cross_encoder_pooled = plan.role in _EMBED_ROLES and not is_llm_rerank
1771 cache_type_k, cache_type_v = chat_cache_type_flags() if is_chat else (None, None)
1772 argv = build_server_argv(
1773 binary=binary,
1774 spec=spec,
1775 model_path=model_path,
1776 devices=plan.devices,
1777 n_gpu_layers=n_gpu_layers,
1778 slots=slots,
1779 ctx_per_slot=ctx,
1780 tensor_split=plan.tensor_split,
1781 mmproj=mmproj,
1782 flash_attn=flash_attn_flag() if _role_launches_with_flash(plan.role, rerank_mode) else None,
1783 cache_type_k=cache_type_k,
1784 cache_type_v=cache_type_v,
1785 batch_size=_pooled_batch_size(plan.role, rerank_mode, ctx),
1786 no_mmap=is_chat and _chat_no_mmap(weights_bytes, on_network_fs=chat_on_network_fs),
1787 cpu_moe=expert_offload_all(meta),
1788 n_cpu_moe=expert_offload_layers(meta),
1789 device_names=_device_names(chosen) or _cpu_pin_when_every_device_was_refused(),
1790 memory_endpoint=supports_memory_readback(binary),
1791 )
1792 return InstanceLaunch(
1793 role=plan.role,
1794 argv=argv,
1795 env_overrides={**visible_env(chosen), **llama_server_runtime_env()},
1796 model=model_ref,
1797 # token_cap drives cross-encoder/embed input truncation; the LLM rerank path
1798 # doesn't truncate (it relies on the per-slot ctx headroom), so leave it None.
1799 token_cap=engine_params.embed_token_cap(ctx) if cross_encoder_pooled else None,
1800 # Weights size scales the cold-load ready timeout (larger model = longer).
1801 weights_bytes=weights_bytes,
1802 # Slots is the chat concurrency the gate admits; ctx is what a client fits to.
1803 slots=slots,
1804 ctx=ctx,
1805 built_ctx_target=(
1806 (cfg.num_ctx if cfg.num_ctx is not None else cfg.chat_n_ctx_target) if is_chat else 0
1807 ),
1808 built_slots_target=(
1809 max(1, cfg.vision_ocr_concurrency) if plan.role is WorkerRole.VISION else 0
1810 ),
1811 replica=plan.replica,
1812 rerank_mode=rerank_mode,
1813 # What placement charged this instance, for the post-launch check against
1814 # the engine's own report of what it really allocated.
1815 est_vram_bytes=est_vram_bytes,
1816 est_vram_by_device=_charge_by_device(chosen, plan.tensor_split, est_vram_bytes),
1817 est_unreported_bytes=_unreported_bytes(plan.role, mmproj),
1818 )
1821def build_single_role_launch(role: WorkerRole, model_path: Path) -> InstanceLaunch:
1822 """The launch the fleet would build for *role* serving *model_path*, alone.
1824 One construction path. The self-check used to assemble its own beside this
1825 one and the two disagreed on slot count, on the context that follows from it,
1826 on device pinning and on the tensor split, so a green check proved nothing
1827 about the launch serving actually performs, and a red one could be a
1828 configuration serving would never have chosen.
1830 Placement is the planner's, on the devices the plan snapshot holds, so the
1831 check runs on the card the role would really land on.
1832 """
1833 from lilbee.providers.fleet.cuda_runtime import apply_cuda_runtime_env
1834 from lilbee.providers.fleet.gpu_env import apply_fleet_gpu_env
1836 apply_fleet_gpu_env()
1837 binary = resolve_llama_server()
1838 apply_cuda_runtime_env(binary)
1839 with _one_engine_per_pass():
1840 devices = _plan_devices(binary)
1841 by_index = {d.index: d for d in devices}
1842 # The whole machine, since nothing else is resident during a self-check.
1843 placed = (min(by_index),) if by_index else ()
1844 plan = InstancePlan(role=role, devices=placed)
1845 return _launch_for(
1846 plan,
1847 str(model_path),
1848 binary,
1849 by_index,
1850 unified_budget=_unified_memory_budget(devices),
1851 model_path=model_path,
1852 )
1855def resolve_devices(binary: Path) -> list[FleetDevice]:
1856 """Enumerate devices in the binary's index space, or the Vulkan VRAM probe."""
1857 return _read_devices(binary).devices
1860# The visibility variable each vendor's runtime reads, named in the warning so
1861# the reader checks the one that applies to the card they actually have.
1862_VENDOR_VISIBILITY_HINT = {
1863 "NVIDIA": "CUDA_VISIBLE_DEVICES",
1864 "AMD": "ROCR_VISIBLE_DEVICES / HIP_VISIBLE_DEVICES",
1865 "Intel": "ONEAPI_DEVICE_SELECTOR",
1866}
1869def _warn_gpu_present_but_unenumerated(binary: Path) -> None:
1870 """Say so when the host has a GPU the engine did not list.
1872 Previously asked only whether an NVIDIA card was present, so an AMD or Intel
1873 host whose engine enumerated nothing produced the identical symptom, a fleet
1874 quietly planned for CPU, and said nothing. The vendor lookup is the same one
1875 the Vulkan ICD rules use, and it works on Windows as well as Linux.
1876 """
1877 from lilbee.providers.fleet.gpu_hardware import installed_gpu_vendor_ids
1878 from lilbee.providers.fleet.gpu_select import PCIVendorID
1880 present = installed_gpu_vendor_ids()
1881 names = sorted(
1882 v.name.title() if v.name == "INTEL" else v.name for v in PCIVendorID if v in present
1883 )
1884 if not names:
1885 return
1886 hints = sorted(
1887 {_VENDOR_VISIBILITY_HINT[name] for name in names if name in _VENDOR_VISIBILITY_HINT}
1888 )
1889 log.warning(
1890 "This host has a %s GPU but the engine's device probe (%s --list-devices) "
1891 "reported none; placement is falling back to shared-memory mode with unpinned "
1892 "GPUs. Check the GPU driver, %s, and that this llama-server build supports "
1893 "that GPU.",
1894 " and ".join(names),
1895 binary,
1896 " / ".join(hints) if hints else "the vendor's visibility variable",
1897 )
1900@dataclass(frozen=True)
1901class DeviceReading:
1902 """What one ``--list-devices`` run answered, as the whole fleet reads it.
1904 The device list and the backend are independent answers. When the engine never
1905 spoke the protocol but the host's Vulkan loader supplies devices, the reading
1906 carries those devices with an ``unknown`` backend, so a client can see a
1907 non-empty GPU list beside ``engine_backend: "unknown"``.
1908 """
1910 devices: list[FleetDevice]
1911 # The backend the engine selected, UNKNOWN when it never answered. Not
1912 # derivable from ``devices``: an empty list is a CPU host and a failed probe
1913 # alike, and only the engine can say which.
1914 backend: EngineBackend
1915 # The engine listed GPUs and lilbee rejected all of them, so the plan is
1916 # CPU-shaped while the engine would still choose one of those devices.
1917 refused_all: bool = False
1920def _read_devices(binary: Path) -> DeviceReading:
1921 """Enumerate devices, the selected backend, and whether every GPU was refused.
1923 One function because all three answers come from one ``--list-devices`` run,
1924 and that run costs a subprocess against a driver that may be wedged. Asking
1925 twice would pay it twice.
1927 The binary's ``--list-devices`` is authoritative, including when it lists
1928 nothing: it prints every non-CPU device it can use, so an empty list means
1929 the engine has no usable GPU rather than that we failed to look. The Vulkan
1930 VRAM probe is consulted only when the binary produced no output at all, and
1931 it reports the same index space. A
1932 probe that times out raises instead (a wedged GPU driver); falling through
1933 to the in-process Vulkan probe there could hang this thread unkillably.
1934 """
1935 from lilbee.providers.fleet.cuda_runtime import assert_cuda_devices_usable
1936 from lilbee.providers.fleet.gpu_hardware import installed_gpu_vendor_ids
1937 from lilbee.providers.fleet.gpu_select import enumerate_gpu_vram
1938 from lilbee.providers.fleet.rocm_runtime import assert_rocm_devices_usable
1940 probe = probe_devices(binary)
1941 # An engine that answered but listed no device while a GPU is physically
1942 # present may be hitting a transient init error (the card momentarily held by
1943 # another process, e.g. an embedder served alongside -- the "ggml_cuda_init:
1944 # initialization error" symptom). Re-probe before treating the empty list as
1945 # fatal; a persistently empty list still hits the fail-loud asserts below.
1946 # The engine identity a plan snapshot carries does not cover this: the binary
1947 # is the same one across a transient init error, so only asking again tells it
1948 # apart from a host whose engine has no usable GPU.
1949 for _ in range(_DEVICE_PROBE_EMPTY_RETRIES):
1950 if probe.devices or not probe.spoke_protocol or not installed_gpu_vendor_ids():
1951 break
1952 time.sleep(_DEVICE_PROBE_EMPTY_RETRY_DELAY_S)
1953 probe = probe_devices(binary)
1954 devices = probe.devices
1955 # A GPU build that links a runtime it cannot serve the host's GPU with must
1956 # fail loud, not silently fall back to CPU (the Vulkan VRAM probe below
1957 # would mask it). Only when the engine actually answered, though: a binary
1958 # that does not support --list-devices enumerated nothing because it was
1959 # never asked, and accusing its driver of failing would be wrong and fatal.
1960 if probe.spoke_protocol:
1961 assert_cuda_devices_usable(binary, devices, probe.output)
1962 assert_rocm_devices_usable(binary, devices, probe.output)
1963 if not devices and probe.spoke_protocol:
1964 _warn_gpu_present_but_unenumerated(binary)
1965 if not devices and not probe.spoke_protocol:
1966 # Only when the binary never answered the question. An engine that ran
1967 # and listed nothing is reporting a fact, not a gap: believing the host
1968 # loader instead invents devices the engine cannot see. A CPU-only build
1969 # on a desktop with mesa is the clearest case, and the cost is not merely
1970 # a wrong device list. The fleet is planned onto GPUs, the pins are
1971 # no-ops, the shared-RAM guard is off because devices looked non-empty,
1972 # and every role then loads its full weights into system RAM while
1973 # running on the CPU anyway.
1974 #
1975 # Keyed on the exit code and the header rather than on there being no
1976 # output at all: the probe merges stderr into stdout, so a build that
1977 # predates --list-devices prints usage text and would otherwise be read
1978 # as an authoritative "no GPUs here".
1979 from lilbee.providers.fleet.gpu_select import integrated_vulkan_indices
1981 integrated = integrated_vulkan_indices()
1982 devices = [
1983 FleetDevice(
1984 VULKAN_BACKEND, idx, "", vram, free, unified=idx in integrated, from_loader=True
1985 )
1986 for idx, vram, free in (enumerate_gpu_vram() or [])
1987 ]
1988 if devices:
1989 log.warning(
1990 "The engine's device probe returned nothing, so placement is using "
1991 "the host's Vulkan loader instead and found %d device(s). If the "
1992 "engine has no Vulkan backend it will run on CPU regardless; set %s "
1993 "to override the engine location if that is wrong.",
1994 len(devices),
1995 "LILBEE_ENGINE_DIR",
1996 )
1997 return DeviceReading(devices, probe.backend, probe.refused_all)
2000_DEVICE_PROBE_TTL_S = 2.0
2001# A failed probe is cached much longer than a good one: each retry against a
2002# wedged GPU driver costs a full probe timeout, so a per-poll retry would stall
2003# every placement read for a minute at a time.
2004_DEVICE_PROBE_FAILURE_TTL_S = 60.0
2005# How long a plan read holds off re-probing an engine whose restate probe raised.
2006# The same order as the read cache's failure TTL and for the same reason (the
2007# retry ladder costs a full probe timeout), but a separate knob: the read cache
2008# bounds how long a wedged driver goes un-re-probed by view reads, this one how
2009# long a stale plan snapshot stays stale.
2010_PLAN_RESTATE_FAILURE_WAIT_S = 60.0
2011# An engine that lists no device on a GPU host may be hitting a transient GPU-init
2012# error (the card momentarily held by another process); re-probe before treating
2013# the empty list as fatal.
2014_DEVICE_PROBE_EMPTY_RETRIES = 3
2015_DEVICE_PROBE_EMPTY_RETRY_DELAY_S = 0.5
2016# Stands in for the engine binary's identity on a host that has none, so "no
2017# engine yet" and "an engine at last" are different identities.
2018_NO_ENGINE_BINARY = "no-engine"
2021class _ReadDeviceCache:
2022 """Short-TTL device-probe cache for the read/view path.
2024 Not a ``cachetools.TTLCache``: it caches the *failure* too, under its own
2025 longer TTL, and re-raises it. A memoizing cache stores return values only,
2026 so a failing probe would re-spawn the subprocess on every placement read.
2028 Inspecting placement (GET placement/gpus, preview, ``placement show``)
2029 resolves devices on every call, which spawns a ``llama-server --list-devices``
2030 subprocess; a brief TTL collapses a burst of reads onto one probe. A probe
2031 failure is cached too (with its own TTL) and re-raised to every read in the
2032 window. The launch path is never served from here -- it sizes against the
2033 clean-box plan snapshot below (captured after stale-server reaping).
2035 An entry also carries the identity of the binary that answered, so it cannot
2036 describe an engine that has since been installed or upgraded in place.
2037 """
2039 def __init__(self, ttl_s: float, failure_ttl_s: float) -> None:
2040 self._ttl_s = ttl_s
2041 self._failure_ttl_s = failure_ttl_s
2042 self._lock = threading.Lock()
2043 self._at: float | None = None
2044 self._engine: str | None = None
2045 self._reading: DeviceReading | None = None
2046 self._failure: ProviderError | None = None
2048 def _is_fresh(self, engine: str) -> bool:
2049 """Whether the entry in hand answers for *engine* and is still inside its TTL."""
2050 ttl = self._ttl_s if self._failure is None else self._failure_ttl_s
2051 return self._at is not None and self._engine == engine and time.monotonic() - self._at < ttl
2053 def get(self, binary: Path) -> DeviceReading:
2054 engine = engine_binary_identity(binary)
2055 with self._lock:
2056 fresh = self._is_fresh(engine)
2057 if fresh and self._failure is not None:
2058 raise self._failure
2059 if self._reading is None or not fresh:
2060 self._at = time.monotonic()
2061 self._engine = engine
2062 try:
2063 self._reading = _read_devices(binary)
2064 except ProviderError as exc:
2065 self._reading = None
2066 self._failure = exc
2067 raise
2068 self._failure = None
2069 return self._reading
2071 def clear(self) -> None:
2072 with self._lock:
2073 self._at = None
2074 self._engine = None
2075 self._reading = None
2076 self._failure = None
2079_read_device_cache = _ReadDeviceCache(_DEVICE_PROBE_TTL_S, _DEVICE_PROBE_FAILURE_TTL_S)
2082def clear_read_device_cache() -> None:
2083 """Drop the read-path device probe cache (e.g. after the fleet is reconfigured).
2085 Also drops what the host's Vulkan loader told us about device types, which is
2086 otherwise held for the process lifetime and would survive a driver reload or
2087 an eGPU being plugged in.
2088 """
2089 from lilbee.providers.fleet.gpu_select import (
2090 integrated_vulkan_indices,
2091 vulkan_device_types_by_name,
2092 )
2094 _read_device_cache.clear()
2095 vulkan_device_types_by_name.cache_clear()
2096 integrated_vulkan_indices.cache_clear()
2099@dataclass(frozen=True)
2100class _PlanProbe:
2101 """Clean-box memory snapshot every plan is sized against.
2103 Captured once, right after stale-server reaping and before the first build,
2104 when nothing lilbee owns is loaded. Reloads re-plan against this same
2105 snapshot instead of re-probing: a live probe under a loaded fleet reports
2106 our own residency as unavailable, which would shrink chat context and slot
2107 counts, widen splits, and (on a unified-memory host) evict roles outright.
2108 Launches stay a pure function of config + hardware + this snapshot, so the
2109 reload diff restarts only real changes. Cleared on full fleet teardown so
2110 the next boot probes the clean box afresh.
2111 """
2113 devices: tuple[FleetDevice, ...]
2114 # What one role may size its ctx and slots against, already scaled by
2115 # cfg.gpu_memory_fraction. System memory only on a host with no GPU.
2116 sizing_budget: int
2117 free_system: int
2118 # The engine binary that answered the probe. A snapshot is only an answer
2119 # about the binary it was taken against.
2120 engine: str
2121 # The engine listed GPUs and lilbee rejected all of them, so the plan is
2122 # CPU-shaped while the engine would still choose one of those devices.
2123 engine_devices_all_refused: bool = False
2124 # The backend the engine selected when this snapshot was taken.
2125 backend: EngineBackend = EngineBackend.UNKNOWN
2128@dataclass(frozen=True)
2129class _FailedProbe:
2130 """An engine identity whose restate probe raised, and when it raised."""
2132 engine: str
2133 at: float
2136# Re-reads a plan snapshot against the hardware, or answers ``None`` when the
2137# probe could not run.
2138_Restate: TypeAlias = "Callable[[_PlanProbe], _PlanProbe | None]"
2141class _PlanProbeStore:
2142 """Holds the captured plan snapshot; a single instance below (no bare global)."""
2144 def __init__(self) -> None:
2145 self._lock = threading.Lock()
2146 self._probe: _PlanProbe | None = None
2147 self._failed: _FailedProbe | None = None
2148 # Held across the probe a restate runs, so a burst of reads after an
2149 # engine swap pays one probe. Private, and never taken outside this
2150 # class, so no caller can hold it across a probe of its own. Separate
2151 # from _lock, which must never be held across a subprocess.
2152 self._restate_lock = threading.Lock()
2154 def set(self, probe: _PlanProbe) -> None:
2155 with self._lock:
2156 self._probe = probe
2157 self._failed = None
2159 def get(self) -> _PlanProbe | None:
2160 with self._lock:
2161 return self._probe
2163 def note_failed_probe(self, engine: str) -> None:
2164 """Record that the restate against *engine* could not probe."""
2165 with self._lock:
2166 self._failed = _FailedProbe(engine, time.monotonic())
2168 def probe_failed_recently(self, engine: str) -> bool:
2169 """Whether a restate against *engine* raised inside the failure wait."""
2170 with self._lock:
2171 failed = self._failed
2172 return (
2173 failed is not None
2174 and failed.engine == engine
2175 and time.monotonic() - failed.at < _PLAN_RESTATE_FAILURE_WAIT_S
2176 )
2178 def restate(self, restate: _Restate, engine: str) -> None:
2179 """Re-read the snapshot with *restate*, whatever binary it answers for now."""
2180 with self._restate_lock:
2181 probe = self.get()
2182 if probe is not None:
2183 self._store_restated(probe, restate, engine)
2185 def restate_if_stale(self, restate: _Restate, engine: str) -> _PlanProbe | None:
2186 """The snapshot, re-read with *restate* first when it answers for another binary.
2188 The staleness and the failure wait are re-checked under the lock, so a
2189 burst of stale reads pays one probe and the rest take its answer.
2190 """
2191 with self._restate_lock:
2192 probe = self.get()
2193 if probe is None or probe.engine == engine or self.probe_failed_recently(engine):
2194 return probe
2195 return self._store_restated(probe, restate, engine)
2197 def _store_restated(self, probe: _PlanProbe, restate: _Restate, engine: str) -> _PlanProbe:
2198 """Store *probe* restated, or keep it and hold off re-probing *engine* for a while.
2200 A probe that raised is not an answer about *engine*, so the snapshot keeps
2201 the identity that did answer for it and the failure is recorded instead.
2202 """
2203 fresh = restate(probe)
2204 if fresh is None:
2205 self.note_failed_probe(engine)
2206 return probe
2207 self.set(fresh)
2208 return fresh
2210 def clear(self) -> None:
2211 with self._lock:
2212 self._probe = None
2213 self._failed = None
2216_plan_probe_store = _PlanProbeStore()
2219# Where the ladder stops. Sized for chat, below which the answers are too short
2220# to be useful, so a role that still will not load here has a real problem the
2221# planner cannot size its way out of and the failure should surface. Roles whose
2222# window already sits under it (a small embedding context) are left alone rather
2223# than raised to meet it, so for them the ladder is a no-op and the failure
2224# surfaces after the one retry.
2225MIN_DOWNSHIFT_CTX = 4096
2228class _CtxDownshiftStore:
2229 """How many halvings each role's auto context has taken after a load OOM.
2231 An estimate that was too optimistic is only recoverable if the retry asks
2232 for something different. Halving the auto context does that, and keeping the
2233 count here rather than in the launch means the whole plan is re-predicted
2234 against the smaller number, including the placement it implies.
2235 """
2237 def __init__(self) -> None:
2238 self._lock = threading.Lock()
2239 self._steps: dict[WorkerRole, int] = {}
2240 # The last unshifted context each role was sized from, recorded as it is
2241 # applied. Deciding whether another halving would change anything needs
2242 # the number being halved, and this is the only place that sees it.
2243 self._base: dict[WorkerRole, int] = {}
2245 def steps(self, role: WorkerRole) -> int:
2246 with self._lock:
2247 return self._steps.get(role, 0)
2249 def note_base(self, role: WorkerRole, ctx: int) -> None:
2250 with self._lock:
2251 self._base[role] = ctx
2253 def base(self, role: WorkerRole) -> int | None:
2254 with self._lock:
2255 return self._base.get(role)
2257 def step(self, role: WorkerRole) -> int:
2258 with self._lock:
2259 taken = self._steps.get(role, 0) + 1
2260 self._steps[role] = taken
2261 return taken
2263 def clear(self, role: WorkerRole | None = None) -> None:
2264 with self._lock:
2265 if role is None:
2266 self._steps.clear()
2267 self._base.clear()
2268 return
2269 self._steps.pop(role, None)
2270 self._base.pop(role, None)
2273_ctx_downshift_store = _CtxDownshiftStore()
2276def apply_ctx_downshift(role: WorkerRole, ctx: int) -> int:
2277 """*ctx* halved once per downshift step recorded for *role*, floored.
2279 Never more than *ctx*. The floor is a stopping point, not a target: applied
2280 to a context already below it (a small embedding window, a model trained for
2281 2048 tokens) a bare floor would hand back a larger number, and the retry
2282 after a load OOM would ask for more memory than the launch that just ran out
2283 of it. Such a role simply has nothing to give back, and its failure surfaces
2284 after the one retry instead.
2286 A user's ``cfg.num_ctx`` pin is returned untouched: serving a window smaller
2287 than the one that was asked for, without being asked, is worse than failing
2288 to load and saying so.
2289 """
2290 from lilbee.core.config import cfg
2292 if role is WorkerRole.CHAT and cfg.num_ctx is not None:
2293 return ctx
2294 _ctx_downshift_store.note_base(role, ctx)
2295 return _shifted(ctx, _ctx_downshift_store.steps(role))
2298def _shifted(ctx: int, steps: int) -> int:
2299 """*ctx* halved *steps* times, never below the floor and never above *ctx*."""
2300 return min(ctx, max(MIN_DOWNSHIFT_CTX, ctx >> steps)) if steps else ctx
2303def record_ctx_downshift(role: WorkerRole) -> bool:
2304 """Take one downshift step for *role*; False when there is none left to take.
2306 False means the retry would ask for the same thing again, so the caller must
2307 surface the load failure instead of respawning an identical launch.
2308 """
2309 from lilbee.core.config import cfg
2311 if role is WorkerRole.CHAT and cfg.num_ctx is not None:
2312 return False
2313 base = _ctx_downshift_store.base(role)
2314 if base is None:
2315 # Nothing has been sized for this role yet, so there is no number to
2316 # decide against. Allow one step rather than trusting that a plan always
2317 # runs first: an unbounded grant here would let a caller that never
2318 # sizes anything loop forever.
2319 if _ctx_downshift_store.steps(role):
2320 return False
2321 _ctx_downshift_store.step(role)
2322 return True
2323 steps = _ctx_downshift_store.steps(role)
2324 if _shifted(base, steps + 1) == _shifted(base, steps):
2325 return False
2326 _ctx_downshift_store.step(role)
2327 return True
2330def clear_ctx_downshift(role: WorkerRole | None = None) -> None:
2331 """Forget *role*'s recorded downshift, or every role's, back to full size.
2333 Called when a role's engine reports ready, which is proof the reduced plan
2334 loaded: keeping the reduction after that would carry a shrunken window into
2335 a machine that has since freed memory, or into a smaller model the user
2336 switched to, and would then refuse on its first failure with a budget it
2337 had already spent.
2338 """
2339 _ctx_downshift_store.clear(role)
2342def _probe_engine_devices() -> DeviceReading:
2343 """Apply the fleet GPU/CUDA env, resolve the binary, and enumerate devices.
2345 This is the wedge point: a missing binary raises NOT_FOUND, and a CUDA build
2346 that cannot init a GPU (a broken-runtime host) raises loud from resolve_devices
2347 rather than silently degrading. Device enumeration reads no residency, so it is
2348 safe to run while an incumbent engine is still up.
2349 """
2350 from lilbee.providers.fleet.cuda_runtime import apply_cuda_runtime_env
2351 from lilbee.providers.fleet.gpu_env import apply_fleet_gpu_env
2353 apply_fleet_gpu_env()
2354 binary = resolve_llama_server()
2355 apply_cuda_runtime_env(binary)
2356 reading = _read_devices(binary)
2357 if reading.devices:
2358 return reading
2359 return _reprobe_while_a_gpu_is_installed(binary, reading)
2362def _reprobe_while_a_gpu_is_installed(binary: Path, reading: DeviceReading) -> DeviceReading:
2363 """Ask again when the host has a GPU the engine did not list.
2365 The plan snapshot is taken once, on a clean box, and is not retaken until a
2366 full teardown, so an empty first answer decides the whole run. A GPU driver
2367 that is still initializing when the daemon starts, which is ordinary under
2368 systemd or right after a container gains a device, would leave a GPU host
2369 serving on CPU until someone noticed and restarted it.
2371 Only where a card is actually installed. A host with no GPU answers empty
2372 every time and must not pay a retry for it on every start.
2373 """
2374 from lilbee.providers.fleet.gpu_hardware import installed_gpu_vendor_ids
2376 if not installed_gpu_vendor_ids():
2377 return reading
2378 for attempt in range(1, _PROBE_RETRIES + 1):
2379 log.info(
2380 "The engine listed no GPU on a host that has one; asking again in %.1fs "
2381 "(attempt %d of %d) in case the driver is still initializing.",
2382 _PROBE_RETRY_DELAY_S,
2383 attempt,
2384 _PROBE_RETRIES,
2385 )
2386 time.sleep(_PROBE_RETRY_DELAY_S)
2387 clear_read_device_cache()
2388 reading = _read_devices(binary)
2389 if reading.devices:
2390 return reading
2391 return reading
2394def assert_engine_probeable() -> None:
2395 """Raise if the engine cannot be probed; capture no snapshot.
2397 A build precondition that must run BEFORE stopping a replaceable incumbent:
2398 it surfaces a wedged GPU probe or an unusable CUDA runtime without taking the
2399 residency-dependent memory snapshot (that belongs on the clean box, after the
2400 stop, in capture_plan_probe). resolve_devices caches within its TTL, so the
2401 follow-up capture reuses this enumeration rather than re-probing the hardware.
2402 """
2403 _probe_engine_devices()
2406def _engine_identity() -> str:
2407 """Identity of the engine binary a plan snapshot has to agree with.
2409 A host with no engine binary has an identity of its own, so a snapshot taken
2410 before the engine arrived never reads as one taken after.
2411 """
2412 try:
2413 return engine_binary_identity(resolve_llama_server())
2414 except (ProviderError, OSError):
2415 return _NO_ENGINE_BINARY
2418# The binary one planning pass answers about. A pass reads the snapshot half a
2419# dozen times; resolving the identity per read let an engine that landed between
2420# two of them size the plan against one snapshot and place it against another.
2421_pass_engine: contextvars.ContextVar[str | None] = contextvars.ContextVar(
2422 "lilbee_plan_pass_engine", default=None
2423)
2426@contextmanager
2427def _one_engine_per_pass() -> Iterator[None]:
2428 """Pin the engine identity every snapshot read in this planning pass answers about.
2430 Re-entrant: a nested pass keeps the outer one's binary, so a pass that calls
2431 another still asks about one engine. A snapshot that is stale when the pass
2432 opens is restated once, at the first read.
2433 """
2434 if _pass_engine.get() is not None:
2435 yield
2436 return
2437 token = _pass_engine.set(_engine_identity())
2438 try:
2439 yield
2440 finally:
2441 _pass_engine.reset(token)
2444def _pass_engine_identity() -> str:
2445 """The engine this read answers about: the pass's binary, or the live one outside a pass."""
2446 pinned = _pass_engine.get()
2447 return pinned if pinned is not None else _engine_identity()
2450def _probe_engine_devices_and_identity() -> tuple[str, DeviceReading]:
2451 """The engine's reading, tagged with the identity read before the probe ran.
2453 The identity is read first so a binary replaced while the probe runs tags the
2454 answer with the binary that gave it, and the next read restates. The other
2455 order tags the outgoing binary's answer with the incoming binary.
2456 """
2457 engine = _engine_identity()
2458 return engine, _probe_engine_devices()
2461def capture_plan_probe() -> None:
2462 """Snapshot devices and memory for planning; call only on a clean box."""
2463 engine, reading = _probe_engine_devices_and_identity()
2464 _plan_probe_store.set(
2465 _PlanProbe(
2466 devices=tuple(reading.devices),
2467 sizing_budget=_device_sizing_budget(reading.devices),
2468 free_system=model_cache.free_system_memory(),
2469 engine=engine,
2470 engine_devices_all_refused=reading.refused_all,
2471 backend=reading.backend,
2472 )
2473 )
2476def _structural(devices: Iterable[FleetDevice]) -> tuple[FleetDevice, ...]:
2477 """*devices* with the volatile free reading zeroed, so equality is structural.
2479 This probe runs while the fleet is resident, so its free figures are deflated
2480 by the very models a reload is about to stop; comparing or keying on them
2481 would read every reload as a hardware change.
2482 """
2483 return tuple(replace(d, free_bytes=0) for d in devices)
2486def refresh_plan_devices() -> None:
2487 """Re-read which devices exist, keeping the clean-box memory figures.
2489 The snapshot is captured once and only a full teardown clears it, so an eGPU
2490 unplugged, a driver reset, or a VM hot-remove left the fleet pinning a device
2491 that is no longer there and every rebuild replanning onto it.
2493 Only the structural half is restated. The memory figures are what make a
2494 reload plan the way the boot did, and re-taking them while the fleet is
2495 resident would charge it against itself, which is the whole reason the
2496 snapshot exists. So a card that survives the refresh keeps the snapshot's
2497 free figure; only a genuinely new card contributes a fresh one.
2499 A probe that cannot run keeps the previous device list: the last known one is
2500 a better answer than none, and the loud paths for an unreachable engine live
2501 in the build, not here. A reload asks to look again, so it probes even inside
2502 the failure wait that holds the read path off.
2503 """
2504 _plan_probe_store.restate(_restated, _engine_identity())
2507def _restated(probe: _PlanProbe) -> _PlanProbe | None:
2508 """*probe* re-read against the hardware, or ``None`` when the probe could not run."""
2509 clear_read_device_cache()
2510 try:
2511 engine, reading = _probe_engine_devices_and_identity()
2512 except (ProviderError, OSError) as exc:
2513 log.debug("Device rediscovery could not run, keeping the previous list: %s", exc)
2514 return None
2515 devices = reading.devices
2516 if _structural(devices) == _structural(probe.devices):
2517 # The refusal answer is carried too: an engine that lists only devices
2518 # lilbee refuses probes as an empty list, which is structurally the same
2519 # as no GPU at all. Keeping the old binary's refusal here would leave the
2520 # CPU pin off and let ggml fall back onto the adapter just refused.
2521 return replace(
2522 probe,
2523 engine=engine,
2524 engine_devices_all_refused=reading.refused_all,
2525 backend=reading.backend,
2526 )
2527 log.info(
2528 "The set of GPUs changed since this fleet was planned (%d device(s) now, %d before); "
2529 "replanning against the ones that are here.",
2530 len(devices),
2531 len(probe.devices),
2532 )
2533 snapshot_free = {replace(d, free_bytes=0): d.free_bytes for d in probe.devices}
2534 merged = tuple(
2535 replace(d, free_bytes=snapshot_free.get(replace(d, free_bytes=0), d.free_bytes))
2536 for d in devices
2537 )
2538 return _PlanProbe(
2539 devices=merged,
2540 sizing_budget=_device_sizing_budget(merged),
2541 free_system=probe.free_system,
2542 engine=engine,
2543 engine_devices_all_refused=reading.refused_all,
2544 backend=reading.backend,
2545 )
2548def _current_plan_probe() -> _PlanProbe | None:
2549 """The plan snapshot, restated first when another engine binary is now in place.
2551 The snapshot carries the identity of the binary that answered, so an engine
2552 installed or upgraded under a running serve re-probes on the next plan read.
2553 """
2554 probe = _plan_probe_store.get()
2555 if probe is None:
2556 return None
2557 engine = _pass_engine_identity()
2558 return probe if probe.engine == engine else _restate_plan_probe(engine)
2561def _restate_plan_probe(engine: str) -> _PlanProbe | None:
2562 """Restate the snapshot for *engine*, once per burst of stale reads.
2564 An engine whose probe keeps raising is re-probed at most once per failure
2565 wait on this read path, so a broken binary costs the retry ladder
2566 occasionally rather than on every read. A reload restates through
2567 refresh_plan_devices instead, which probes every time.
2569 The wait is keyed on the binary's identity, which is a digest of its bytes,
2570 so a binary repaired in place is retried without a reload as soon as its
2571 content changes, whatever its stat fields do.
2572 """
2573 return _plan_probe_store.restate_if_stale(_restated, engine)
2576def clear_plan_probe() -> None:
2577 """Drop the plan snapshot (full fleet teardown); the next build re-captures."""
2578 _plan_probe_store.clear()
2581def _cpu_pin_when_every_device_was_refused() -> tuple[str, ...]:
2582 """``("none",)`` when the engine offered GPUs that lilbee refused, else empty.
2584 Dropping a device from lilbee's view does not stop the engine using it. With
2585 no pin at all, ggml applies its own selection, and its fallback takes the
2586 first non-CPU adapter, which is exactly the paravirtual device just refused;
2587 with every layer offloaded by default the model then runs on it while
2588 placement budgeted against system RAM. Naming no device keeps the engine on
2589 the CPU the plan was shaped for.
2590 """
2591 probe = _current_plan_probe()
2592 if probe is None or not probe.engine_devices_all_refused:
2593 return ()
2594 log.warning(
2595 "The engine listed GPU devices that lilbee will not plan onto, so it is being "
2596 "run on the CPU. Serving from one of them would be slower than the CPU or fail "
2597 "outright, and placement has been sized for system RAM."
2598 )
2599 return (_NO_DEVICE,)
2602def _plan_devices(binary: Path) -> list[FleetDevice]:
2603 """Devices the plan paths size against: the snapshot, else a live probe."""
2604 probe = _current_plan_probe()
2605 return list(probe.devices) if probe is not None else resolve_devices(binary)
2608def plan_sizing_budget(device: FleetDevice | None = None) -> int:
2609 """Usable memory for ctx/slot sizing: *device*'s own, else the snapshot, else live."""
2610 from lilbee.core.config import cfg
2612 if device is not None:
2613 return int(device.total_bytes * cfg.gpu_memory_fraction)
2614 probe = _current_plan_probe()
2615 if probe is not None:
2616 return probe.sizing_budget
2617 return _device_sizing_budget(_live_sizing_devices())
2620def plan_sizing_is_unified() -> bool:
2621 """Whether ctx sizing charges the shared-memory footprint rather than VRAM.
2623 True when no device has memory of its own -- an integrated GPU, an Apple
2624 Silicon Mac, or no GPU at all -- because there a model's would-be VRAM and
2625 its host bytes are the same memory, and charging only the VRAM half
2626 under-reports the load by everything it maps. The same test decides the pool
2627 placement charges against (:func:`_unified_memory_budget`).
2628 """
2629 probe = _current_plan_probe()
2630 devices = list(probe.devices) if probe is not None else _live_sizing_devices()
2631 return all(dev.unified for dev in devices)
2634def _device_sizing_budget(devices: Sequence[FleetDevice]) -> int:
2635 """Memory one role may size its ctx and slots against, in bytes.
2637 Read from the engine's own device report, which ran under the environment the
2638 servers will run under and states each device's memory whatever the backend.
2639 A host-memory read answers with system RAM on every host without an NVIDIA
2640 card, which gave a 24 GiB AMD card a budget the size of the machine, and on
2641 Apple Silicon it ignores that Metal will not allocate past
2642 ``recommendedMaxWorkingSetSize``, which is the figure the probe carries.
2644 The smallest device, since this is asked before placement has picked one;
2645 :func:`_launch_for` re-sizes against the card the role actually landed on.
2646 System memory only when the engine reports no device at all, where the fleet
2647 runs on the CPU and system memory is the budget.
2648 """
2649 from lilbee.core.config import cfg
2651 if devices:
2652 return int(min(d.total_bytes for d in devices) * cfg.gpu_memory_fraction)
2653 return int(model_cache.total_system_memory() * cfg.gpu_memory_fraction)
2656def _live_sizing_devices() -> list[FleetDevice]:
2657 """Devices to size against with no plan snapshot; empty when none can be read."""
2658 try:
2659 return _read_device_cache.get(resolve_llama_server()).devices
2660 except (ProviderError, OSError):
2661 return []
2664def _plan_free_system_memory() -> int:
2665 """Free system RAM for the unified-memory budget: the snapshot, else live."""
2666 probe = _current_plan_probe()
2667 return probe.free_system if probe is not None else model_cache.free_system_memory()
2670def _unreported_bytes(role: WorkerRole, mmproj: Path | None) -> int:
2671 """Estimated bytes the engine allocates without printing a buffer line.
2673 A vision projector's weights: llama.cpp allocates them in clip's own loader,
2674 which prints a size but not the "buffer size = N MiB" shape the readback
2675 reads, so the report is short by exactly this and the self-check would warn
2676 on a load that was sized correctly.
2677 """
2678 if role is not WorkerRole.VISION or mmproj is None:
2679 return 0
2680 try:
2681 return mmproj.stat().st_size
2682 except OSError:
2683 return 0
2686def _chat_no_mmap(weights_bytes: int, *, on_network_fs: bool = False) -> bool:
2687 """Whether the chat server should malloc its weights instead of mmapping them.
2689 Local disk mmaps: lazy page-fault paging gives a faster first token on a cold
2690 cache -- the common desktop first launch -- matching mmap-by-default engines.
2691 ``--no-mmap``'s buffered full read only wins on an already-hot cache and it
2692 pessimizes cold start, so it is not worth defaulting on for local disk. A
2693 network filesystem still prefers the buffered read whenever the host copy
2694 fits, because mmap page faults served over the wire can wedge the loader in
2695 uninterruptible I/O (see ``_NO_MMAP_NETWORK_RAM_FRACTION``).
2696 """
2697 if not on_network_fs:
2698 return False
2699 return weights_bytes <= model_cache.total_system_memory() * _NO_MMAP_NETWORK_RAM_FRACTION
2702def _device_names(devices: tuple[FleetDevice, ...]) -> tuple[str, ...]:
2703 """``--device`` names for *devices*, empty when the backend pins through env.
2705 Vulkan and SYCL, because neither one's environment variable speaks the space
2706 the probe enumerated. Vulkan's indexes the raw loader enumeration while the
2707 names come from the engine's filtered list, so the two disagree wherever ggml
2708 drops or merges a device. SYCL's is not an index list at all but a selector
2709 over a backend runtime, so a device the engine calls ``SYCL1`` need not be
2710 Level Zero ordinal 1: OpenCL devices interleave, discarded devices shift the
2711 numbering, and multi-tile cards appear as sub-devices.
2713 ``--device`` sidesteps both by naming devices exactly as ``--list-devices``
2714 printed them, which is where these indices were read from. CUDA and ROCm
2715 keep composing their variables, which do share the probe's space.
2716 """
2717 if not devices or devices[0].backend not in _NAME_PINNED_BACKENDS:
2718 return ()
2719 if any(d.from_loader for d in devices):
2720 # These indices are raw loader ordinals, and --device speaks the engine's
2721 # own post-filter naming, so Vulkan1 here can name Vulkan0 there or
2722 # nothing at all. Sizing against them is still worth doing; pinning by
2723 # them is not. Left unpinned, ggml applies its own device selection,
2724 # which is the filtering lilbee is trying to agree with in the first
2725 # place. The env pin is not the answer either: it takes raw ordinals but
2726 # switches off the type filter, the support check and the dedup with them.
2727 return ()
2728 return tuple(f"{d.backend}{d.index}" for d in devices)
2731def _unified_memory_budget(devices: list[FleetDevice]) -> int | None:
2732 """Shared-RAM placement budget (free RAM minus the OS floor), or ``None``.
2734 ``None`` once any device has memory of its own, since dedicated VRAM is the
2735 constraint there rather than system RAM. A host whose only devices are
2736 integrated, and a host with no devices at all, both stay inside the system
2737 budget: their GPU memory is the system's memory.
2738 """
2739 # Only a device with memory of its own lifts the system-RAM constraint. An
2740 # integrated GPU or an Apple Silicon Mac reports a slice of the same RAM the
2741 # OS is using, so treating its total as headroom over-commits the machine by
2742 # roughly the whole system footprint.
2743 if any(not device.unified for device in devices):
2744 return None
2745 return _capped_by_device_memory(
2746 max(0, _plan_free_system_memory() - _system_memory_floor()), devices
2747 )
2750def _unified_admission_budget(devices: list[FleetDevice]) -> int | None:
2751 """Shared-RAM pool a role set is *admitted* against, or ``None`` if dedicated.
2753 Total installed RAM minus the OS floor, not what happens to be free. Sizing
2754 asks a different question and keeps using free RAM: how much context can be
2755 backed right now. Admission asks whether the machine can host this fleet at
2756 all, and the plan defines the whole intended residency, so charging it
2757 against a live figure refuses a 600 MB model on a box that is merely busy at
2758 the moment, which is what happened. The GPU path already charges total
2759 capacity for exactly this reason.
2760 """
2761 if _unified_memory_budget(devices) is None:
2762 return None
2763 return _capped_by_device_memory(
2764 max(0, model_cache.total_system_memory() - _system_memory_floor()), devices
2765 )
2768def _system_memory_floor() -> int:
2769 """RAM held back for the OS when placing against system memory.
2771 ``cfg.system_memory_reserve_gb``, still capped at a quarter of installed RAM:
2772 a fixed reserve leaves a 7-8 GB host with no budget at all and refuses even
2773 tiny models, so the proportional cap holds however the reserve is set.
2774 """
2775 from lilbee.core.config import cfg
2777 total = model_cache.total_system_memory()
2778 return min(int(cfg.system_memory_reserve_gb * 1024**3), total // _SYSTEM_MEMORY_FLOOR_DIVISOR)
2781def _capped_by_device_memory(budget: int, devices: Sequence[FleetDevice]) -> int:
2782 """*budget*, never above what the devices can address between them.
2784 A shared-memory device still has a ceiling of its own: an integrated GPU
2785 addresses a fixed aperture of system RAM, and Metal will not allocate past
2786 ``recommendedMaxWorkingSetSize``. Both report that ceiling as their total, so
2787 a host budget derived from installed RAM promises memory the devices cannot
2788 reach. Unchanged where the engine reports no device, since the fleet is then
2789 running on the CPU and the host figure is the true one.
2790 """
2791 if not devices:
2792 return budget
2793 return min(budget, sum(d.total_bytes for d in devices))
2796def _device_capacity(devices: list[FleetDevice], charge_against_free: bool) -> dict[int, int]:
2797 """Per-device memory placement may charge against, keyed by device index.
2799 A card's total is what it holds, not what is going spare. A compositor, a
2800 browser, or a training job sitting on VRAM is invisible in the total, and the
2801 usable fraction placement applies covers fragmentation and driver overhead
2802 rather than other tenants, so a plan fits on paper and OOMs at load.
2804 Free bytes answer that, but only where they mean "everyone else's residency":
2805 that is the clean-box snapshot, taken after stale servers are reaped and
2806 before anything is built. Read live on a warm box they also exclude the
2807 fleet's own models, and since a plan always describes the complete intended
2808 residency, charging them there would count the fleet against itself and
2809 report a running plan as unplaceable. Those callers keep the total.
2811 Placement applies its usable fraction to whatever this returns, so a card
2812 with a tenant keeps a proportional margin rather than being packed to its
2813 last free byte, where fragmentation is worst.
2814 """
2815 packable = _packable_devices(devices)
2816 if not charge_against_free:
2817 return {d.index: d.total_bytes for d in packable}
2818 return {d.index: min(d.total_bytes, d.free_bytes) for d in packable}
2821def _packable_devices(devices: list[FleetDevice]) -> list[FleetDevice]:
2822 """The devices bin-packing may charge against.
2824 An integrated GPU's memory is the host's. Packing it beside a dedicated card
2825 promises the same RAM twice, once to its own budget and once to everything
2826 else on the machine, and its heap is often the larger number, so the packer
2827 prefers it: a 32 GiB shared heap outbids a 24 GiB card that actually has the
2828 memory. Where a dedicated device exists it is the one to serve from, and the
2829 integrated one is left to the shared-memory budget.
2831 A host with nothing but integrated devices keeps them. There is nothing else
2832 to serve from, and that path is governed by the system budget rather than by
2833 per-device packing.
2834 """
2835 dedicated = [d for d in devices if not d.unified]
2836 return dedicated or devices
2839def _resolve_placement(
2840 placement: PlacementSpec | None,
2841 inputs: list[ModelPlacementInput],
2842 model_refs: dict[WorkerRole, str],
2843 devices: list[FleetDevice],
2844 *,
2845 unified_budget: int | None,
2846 charge_against_free: bool = False,
2847) -> Placement:
2848 """Resolve a Placement from the manual spec when set, else the auto planner."""
2849 estimate_peak = _peak_estimator(model_refs)
2850 capacity = _device_capacity(devices, charge_against_free)
2851 if placement is not None:
2852 return placement_from_spec(
2853 placement,
2854 tuple(model_refs),
2855 capacity,
2856 estimate_peak=estimate_peak,
2857 )
2858 # The chat split's card count is decided against the snapshot's free VRAM (what the
2859 # launch sizes its context against) so placement and launch agree. A split needs
2860 # >=2 GPUs, so skip the chat model's gguf read entirely below that.
2861 chat_ctx_fit, chat_ctx_target = (
2862 _chat_split_ctx_objective(model_refs) if len(capacity) >= _MIN_SPLIT_GPUS else (None, 0)
2863 )
2864 return plan_placement(
2865 inputs,
2866 [(idx, budget) for idx, budget in capacity.items()],
2867 estimate_peak=estimate_peak,
2868 unified_budget=unified_budget,
2869 chat_ctx_fit=chat_ctx_fit,
2870 chat_ctx_target=chat_ctx_target,
2871 free_headroom={d.index: d.free_bytes for d in devices},
2872 )
2875def _placement_or_auto(
2876 placement: PlacementSpec | None,
2877 inputs: list[ModelPlacementInput],
2878 model_refs: dict[WorkerRole, str],
2879 devices: list[FleetDevice],
2880 *,
2881 unified_budget: int | None,
2882 charge_against_free: bool = False,
2883) -> tuple[Placement, bool]:
2884 """Resolve a saved spec, falling back to auto when it no longer fits the hardware.
2886 Returns the placement and whether the spec was the one applied. Hardware moves
2887 under a saved placement: a card is removed, a driver stops enumerating a GPU, a
2888 container starts without one. Refusing to plan there takes chat, embed and
2889 ingest down over a pin set on hardware the host no longer has, so the fleet
2890 degrades to automatic placement and logs why. An interactive apply still fails
2891 loud (:func:`lilbee.app.placement.set_placement`), where the pin is what the
2892 caller just asked for and a silent substitution would be the surprise.
2893 """
2894 if placement is None:
2895 return _resolve_placement(
2896 None,
2897 inputs,
2898 model_refs,
2899 devices,
2900 unified_budget=unified_budget,
2901 charge_against_free=charge_against_free,
2902 ), False
2903 try:
2904 return _resolve_placement(
2905 placement,
2906 inputs,
2907 model_refs,
2908 devices,
2909 unified_budget=unified_budget,
2910 charge_against_free=charge_against_free,
2911 ), True
2912 except PlacementError as exc:
2913 log.warning(
2914 "The saved GPU placement does not fit this hardware (%s); using automatic "
2915 "placement instead. Set a new placement to replace it.",
2916 exc,
2917 )
2918 return _resolve_placement(
2919 None,
2920 inputs,
2921 model_refs,
2922 devices,
2923 unified_budget=unified_budget,
2924 charge_against_free=charge_against_free,
2925 ), False
2928@dataclass(frozen=True)
2929class ResolvedPlacement:
2930 """Devices + resolved instance plans + model refs for the placement view."""
2932 devices: tuple[FleetDevice, ...]
2933 instances: tuple[InstancePlan, ...]
2934 unplaceable_roles: tuple[WorkerRole, ...]
2935 model_refs: dict[WorkerRole, str]
2936 # Roles placed anyway despite not fitting, with the shortfall in bytes. The
2937 # planner has always known this and only logged it, so a surface showed a
2938 # tight role as comfortably placed.
2939 tight_roles: dict[WorkerRole, int] = field(default_factory=dict)
2940 co_tenants: frozenset[WorkerRole] = frozenset()
2941 # False when a spec was given but did not fit the hardware, so these instances
2942 # are the auto planner's and a surface must not present them as the manual plan.
2943 spec_applied: bool = True
2944 # Roles configured but skipped because their model isn't installed (role -> ref).
2945 # Distinct from unplaceable_roles (installed but won't fit); lets a surface show
2946 # "not downloaded" instead of an empty table on a fresh install.
2947 skipped_not_installed: dict[WorkerRole, str] = field(default_factory=dict)
2948 # The backend the engine selected, which no surface may infer from ``devices``:
2949 # an empty list is a CPU host and a failed probe alike.
2950 engine_backend: EngineBackend = EngineBackend.UNKNOWN
2953def resolve_placement_plan(
2954 placement: PlacementSpec | None, *, fall_back_to_auto: bool = False
2955) -> ResolvedPlacement:
2956 """Probe devices and resolve the auto-or-manual placement, without launching.
2958 ``fall_back_to_auto`` reads *placement* as a saved setting rather than a
2959 request: one that no longer fits the hardware resolves to the auto plan with
2960 ``spec_applied`` False instead of raising.
2962 The reported backend comes from the same reading as the device list, not from
2963 :func:`engine_backend`. This path deliberately reports live hardware rather
2964 than the plan snapshot, so borrowing the snapshot's backend would pair a live
2965 device list with a backend read at fleet build.
2966 """
2967 from lilbee.providers.fleet.cuda_runtime import apply_cuda_runtime_env
2968 from lilbee.providers.fleet.gpu_env import apply_fleet_gpu_env
2970 apply_fleet_gpu_env()
2971 binary = resolve_llama_server()
2972 apply_cuda_runtime_env(binary)
2973 reading = _read_device_cache.get(binary)
2974 devices = reading.devices
2975 unified_budget = _unified_memory_budget(devices)
2976 inputs, model_refs, _, skipped_not_installed = _server_model_inputs(
2977 None, unified_budget=unified_budget, total_vram=sum(d.total_bytes for d in devices)
2978 )
2979 admission_budget = _unified_admission_budget(devices)
2980 if fall_back_to_auto:
2981 resolved, spec_applied = _placement_or_auto(
2982 placement, inputs, model_refs, devices, unified_budget=admission_budget
2983 )
2984 else:
2985 resolved = _resolve_placement(
2986 placement, inputs, model_refs, devices, unified_budget=admission_budget
2987 )
2988 spec_applied = placement is not None
2989 return ResolvedPlacement(
2990 devices=tuple(devices),
2991 instances=resolved.instances,
2992 unplaceable_roles=resolved.unplaceable_roles,
2993 model_refs=model_refs,
2994 co_tenants=resolved.co_tenants,
2995 skipped_not_installed=skipped_not_installed,
2996 spec_applied=spec_applied,
2997 tight_roles=dict(resolved.tight_roles),
2998 engine_backend=reading.backend,
2999 )
3002@dataclass(frozen=True)
3003class FleetPlan:
3004 """The servers to start, and the roles that share one swap group."""
3006 launches: tuple[InstanceLaunch, ...]
3007 co_tenants: frozenset[WorkerRole] = frozenset()
3008 # Configured roles left unplaced because their model isn't installed (role ->
3009 # ref), so the warm path can fail a not-installed chat with a named reason
3010 # instead of spinning the warm line forever.
3011 skipped_not_installed: dict[WorkerRole, str] = field(default_factory=dict)
3012 # Launches refused for a window below the minimum grounded prompt
3013 # (role -> user-facing reason with the numbers).
3014 skipped_unusable_ctx: dict[WorkerRole, str] = field(default_factory=dict)
3017def _log_placement_findings(placement: Placement, model_refs: dict[WorkerRole, str]) -> None:
3018 """Warn about placements that exceed the memory budget.
3020 Shared-memory roles that fit nowhere get no server (loading them would OOM the
3021 host). GPU roles are never refused: one whose estimate exceeds the free VRAM
3022 still loads on demand, with a warning carrying the shortfall.
3023 """
3024 for role in placement.unplaceable_roles:
3025 log.warning(
3026 "%s model %s does not fit available memory and will not be served; "
3027 "free up memory or use a smaller model.",
3028 role.value,
3029 model_refs[role],
3030 )
3031 for role, shortfall in placement.tight_roles.items():
3032 log.warning(
3033 "Memory is tight for the %s model %s: it is estimated to need %.1f GiB more "
3034 "GPU memory than is available. It will still load on demand, keeping the "
3035 "layers that fit on the GPU and the rest in system memory; if it runs "
3036 "slowly, free up GPU memory or use a smaller model.",
3037 role.value,
3038 model_refs[role],
3039 # A sub-0.05 GiB shortfall would render as "0.0 GiB more".
3040 max(shortfall / 1024**3, 0.1),
3041 )
3042 if placement.co_tenants:
3043 log.info(
3044 "%s share GPU memory and load on demand; only one is resident at a time.",
3045 ", ".join(sorted(role.value for role in placement.co_tenants)),
3046 )
3049def _unusable_chat_ctx_reason(launch: InstanceLaunch) -> str | None:
3050 """Reason to refuse a chat launch whose window cannot hold a grounded prompt.
3052 ``None`` for non-chat roles, for a window that holds the minimum grounded
3053 prompt, and for user knobs that ask for a smaller one (a ``num_ctx`` pin,
3054 a sub-minimum ``num_ctx_max`` / ``chat_n_ctx_target``).
3055 """
3056 from lilbee.core.config import cfg
3058 if launch.role is not WorkerRole.CHAT or cfg.num_ctx is not None:
3059 return None
3060 needed = engine_params.min_usable_chat_ctx()
3061 # User knobs capping the window below the minimum are honored (the num_ctx
3062 # pin bypasses above).
3063 asked = min(cfg.chat_n_ctx_target, cfg.num_ctx_max or cfg.chat_n_ctx_target)
3064 if asked < needed:
3065 return None
3066 if launch.ctx >= needed:
3067 return None
3068 return (
3069 f"The chat model {launch.model} loads, but the memory left after its weights "
3070 f"backs only a {launch.ctx}-token context, and a grounded answer needs about "
3071 f"{needed} tokens (system prompt, a retrieved source, the question, and room "
3072 "for the answer), so it will not be served. Use a smaller model or a smaller "
3073 "quant, or set num_ctx to force a larger window."
3074 )
3077def plan_launches(
3078 roles: tuple[WorkerRole, ...] | None,
3079 binary: Path,
3080 by_index: dict[int, FleetDevice],
3081 devices: list[FleetDevice],
3082) -> FleetPlan:
3083 """Plan placement for *roles* (``None`` = all configured) and build their launches."""
3084 with _one_engine_per_pass():
3085 return _planned_launches(roles, binary, by_index, devices)
3088def _planned_launches(
3089 roles: tuple[WorkerRole, ...] | None,
3090 binary: Path,
3091 by_index: dict[int, FleetDevice],
3092 devices: list[FleetDevice],
3093) -> FleetPlan:
3094 """The placement and launches for *roles*, sized and placed against one snapshot."""
3095 from lilbee.core.config import cfg
3097 unified_budget = _unified_memory_budget(devices)
3098 inputs, model_refs, reservation, skipped_not_installed = _server_model_inputs(
3099 roles,
3100 unified_budget=unified_budget,
3101 device_count=len(devices),
3102 total_vram=sum(d.total_bytes for d in devices),
3103 )
3104 spec = PlacementSpec.from_json(cfg.placement) if cfg.placement else None
3105 placement, _spec_applied = _placement_or_auto(
3106 spec,
3107 inputs,
3108 model_refs,
3109 devices,
3110 unified_budget=_unified_admission_budget(devices),
3111 # Only the clean-box snapshot's free bytes mean "what other tenants hold";
3112 # a live probe here would also be missing the fleet's own residency.
3113 charge_against_free=_current_plan_probe() is not None,
3114 )
3115 _log_placement_findings(placement, model_refs)
3116 reserved_by_device = _non_chat_reservation(placement.instances, inputs, placement.co_tenants)
3117 charged = {inp.role: inp.est_vram_bytes for inp in inputs}
3118 launches: list[InstanceLaunch] = []
3119 skipped_unusable_ctx: dict[WorkerRole, str] = {}
3120 for plan in placement.instances:
3121 launch = _launch_for(
3122 plan,
3123 model_refs[plan.role],
3124 binary,
3125 by_index,
3126 unified_budget=unified_budget,
3127 chat_reservation=reservation,
3128 reserved_by_device=reserved_by_device,
3129 est_vram_bytes=charged.get(plan.role, 0),
3130 )
3131 reason = _unusable_chat_ctx_reason(launch)
3132 if reason is not None:
3133 skipped_unusable_ctx[launch.role] = reason
3134 log.warning(reason)
3135 continue
3136 launches.append(launch)
3137 return FleetPlan(
3138 launches=tuple(launches),
3139 co_tenants=placement.co_tenants,
3140 skipped_not_installed=skipped_not_installed,
3141 skipped_unusable_ctx=skipped_unusable_ctx,
3142 )
3145class _VisionRequestGrant:
3146 """Whether a request that turned OCR on has asked for the vision role this process."""
3148 def __init__(self) -> None:
3149 self._lock = threading.Lock()
3150 self._granted = False
3152 @property
3153 def granted(self) -> bool:
3154 with self._lock:
3155 return self._granted
3157 def set(self, granted: bool) -> None:
3158 with self._lock:
3159 self._granted = granted
3162_vision_request_grant = _VisionRequestGrant()
3165def grant_vision_on_request() -> None:
3166 """Plan the vision role even with ``enable_ocr`` false, for a request that turned OCR on."""
3167 _vision_request_grant.set(True)
3170def revoke_vision_on_request() -> None:
3171 """Plan the vision role from the OCR setting alone again."""
3172 _vision_request_grant.set(False)
3175def vision_role_wanted(ref: str) -> bool:
3176 """Whether a plan serves vision on *ref*: the OCR setting uses it, or a request asked."""
3177 from lilbee.core.config import cfg
3179 chosen = OcrBackendUsed.chosen(cfg.enable_ocr, ref)
3180 return chosen is OcrBackendUsed.VISION or _vision_request_grant.granted
3183def _launched_roles() -> tuple[WorkerRole, ...]:
3184 """The roles a launch plan may serve: vision only while OCR can reach it."""
3185 from lilbee.core.config import cfg
3187 return tuple(
3188 role
3189 for role in ROLE_REGISTRY
3190 if role is not WorkerRole.VISION or vision_role_wanted(str(cfg.vision_model))
3191 )
3194def plan_all_launches() -> FleetPlan:
3195 """Apply GPU env, probe devices, and plan launches for the configured roles.
3197 Disables crash-prone Vulkan layers / dual-vendor ICDs and applies any
3198 ``cfg.gpu_devices`` pin before the probe and plan (both inherit the env).
3199 """
3200 from lilbee.providers.fleet.cuda_runtime import apply_cuda_runtime_env
3201 from lilbee.providers.fleet.gpu_env import apply_fleet_gpu_env
3203 apply_fleet_gpu_env()
3204 binary = resolve_llama_server()
3205 # Put the CUDA-runtime wheels on the process path so the device probe sees the
3206 # same runtime the servers will, before resolve_devices enumerates GPUs.
3207 apply_cuda_runtime_env()
3208 with _one_engine_per_pass():
3209 devices = _plan_devices(binary)
3210 by_index = {d.index: d for d in devices}
3211 return plan_launches(_launched_roles(), binary, by_index, devices)