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

1"""Launch planning for the fleet: device probe, VRAM estimate, placement, argv.""" 

2 

3from __future__ import annotations 

4 

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 

14 

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 

60 

61log = logging.getLogger(__name__) 

62 

63if TYPE_CHECKING: 

64 from collections.abc import Callable, Iterable, Iterator, Mapping, Sequence 

65 

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" 

70 

71 

72def log_engine_launch(launch: InstanceLaunch, *, owner_pid: int | None = None) -> None: 

73 """Log the binary, build, backend and devices serving *launch*. 

74 

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 ) 

88 

89 

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 ) 

98 

99 

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. 

102 

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)))) 

111 

112 

113def warn_when_chat_downsized(launch: InstanceLaunch) -> None: 

114 """Log when a chat engine's granted shape ends below the requested one. 

115 

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 ) 

140 

141 

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) 

194 

195 

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 

199 

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 

204 

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 

214 

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 

225 

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$") 

229 

230 

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 ) 

241 

242 

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 ) 

251 

252 

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. 

265 

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 

293 

294 

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. 

297 

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. 

305 

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 

326 

327 

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. 

339 

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 ) 

357 

358 

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 

370 

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 ) 

383 

384 

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 ) 

404 

405 

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. 

416 

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 

424 

425 

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 

440 

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 

460 

461 

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. 

470 

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 """ 

480 

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 ) 

495 

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 

511 

512 

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"}) 

520 

521 

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. 

531 

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. 

537 

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 

544 

545 

546def _chat_offload_probe_ceiling(meta: dict[str, str] | None, configured: int) -> int: 

547 """Highest layer count the offload search probes, below *configured*. 

548 

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) 

559 

560 

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``. 

565 

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 

578 

579 

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. 

587 

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 

595 

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 ) 

609 

610 

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. 

613 

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. 

618 

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 

625 

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 

640 

641 

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``. 

644 

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 

650 

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 

655 

656 

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 

660 

661 arch = meta.get("architecture") if meta else None 

662 return resolve_rerank_mode(cfg.reranker_type, arch) 

663 

664 

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 

668 

669 

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] 

680 

681 

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 

688 

689 

690def _role_gpu_layers(role: WorkerRole) -> int: 

691 """GPU-layer offload: chat honors ``cfg.n_gpu_layers``, others offload all layers.""" 

692 

693 return engine_params.resolve_n_gpu_layers(embedding=role in _ALL_LAYER_ROLES) 

694 

695 

696def _flash_enabled() -> bool: 

697 """Flash attention is on unless ``cfg.flash_attention`` is explicitly ``False``.""" 

698 from lilbee.core.config import cfg 

699 

700 return cfg.flash_attention is not False 

701 

702 

703def probed_devices() -> tuple[FleetDevice, ...]: 

704 """Devices the engine enumerated, empty when they could not be read. 

705 

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 () 

716 

717 

718def engine_backend() -> EngineBackend: 

719 """The backend the engine selected, UNKNOWN when it could not be asked. 

720 

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. 

725 

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 

737 

738 

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) 

742 

743 

744def _flash_attention_is_trusted() -> bool: 

745 """Whether to ask for flash attention outright rather than let the engine decide. 

746 

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 

752 

753 

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 

759 

760 

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*. 

763 

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 

771 

772 

773def _role_flash(role: WorkerRole, rerank_mode: RerankMode | None = None) -> bool: 

774 """Whether the estimate may assume flash attention for *role*. 

775 

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 

781 

782 

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 

786 

787 return cfg.kv_cache_type if role is WorkerRole.CHAT else KvCacheType.F16 

788 

789 

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) 

793 

794 

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. 

797 

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 

805 

806 configured = _role_kv_cache_type(role) 

807 return configured if flash_attn_flag() == _FLASH_ON else KvCacheType.F16 

808 

809 

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 

813 

814 def flag(kind: KvCacheType) -> str | None: 

815 return None if kind is KvCacheType.F16 else kind.value 

816 

817 return flag(_role_kv_cache_type(WorkerRole.CHAT)), flag(_role_kv_cache_type_v(WorkerRole.CHAT)) 

818 

819 

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 

824 

825 try: 

826 return find_mmproj_for_model(engine_params.resolve_model_path(model_ref)) 

827 except (ProviderError, OSError, ValueError, KeyError): 

828 return None 

829 

830 

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). 

841 

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 

848 

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 ) 

889 

890 

891def _chat_serve_budget_footprint(footprint: int) -> int: 

892 """Charge a chat instance against the serve budget, not the placement headroom. 

893 

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 

904 

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)) 

910 

911 

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. 

914 

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 

922 

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)) 

934 

935 

936def _placement_estimate_slots(role: WorkerRole, meta: dict[str, str] | None) -> int: 

937 """The slot count the placement estimate reserves KV for. 

938 

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 

945 

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 

953 

954 

955def _peak_estimator(model_refs: dict[WorkerRole, str]) -> PeakEstimator: 

956 """Per-device VRAM-vector estimator for the planner, bound to the configured models. 

957 

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 

962 

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 

984 

985 return estimate_peak 

986 

987 

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. 

992 

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 

1000 

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) 

1004 

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 ) 

1019 

1020 return fit, target 

1021 

1022 

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 ) 

1031 

1032 

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 

1037 

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 

1047 

1048 

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 

1056 

1057 

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 

1061 

1062 return bool(cfg.cpu_moe) and _is_moe(meta) 

1063 

1064 

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. 

1067 

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 

1072 

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 

1076 

1077 

1078def _role_expert_offload(model_path: Path) -> tuple[str, ...]: 

1079 """Expert patterns the launch will offload, for sizing the same way it runs. 

1080 

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 

1086 

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 ) 

1091 

1092 

1093def _expert_offload_configured() -> bool: 

1094 """Whether the user asked for expert offload that would actually take effect. 

1095 

1096 A non-positive ``n_cpu_moe`` offloads nothing, so it does not count. 

1097 """ 

1098 from lilbee.core.config import cfg 

1099 

1100 return bool(cfg.cpu_moe) or (cfg.n_cpu_moe is not None and cfg.n_cpu_moe >= 1) 

1101 

1102 

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. 

1105 

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 

1113 

1114 

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. 

1117 

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. 

1123 

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 

1130 

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 ) 

1138 

1139 

1140def _vision_without_mmproj(role: WorkerRole, ref: str) -> bool: 

1141 """True (with a warning) for a configured vision model whose mmproj is missing. 

1142 

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 

1154 

1155 

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. 

1168 

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 

1174 

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 ) 

1204 

1205 

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. 

1216 

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 

1230 

1231 

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. 

1236 

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. 

1242 

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 """ 

1249 

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 

1261 

1262 

1263def _estimate_is_implausible(*, estimated: int, floor: int) -> bool: 

1264 """Whether *estimated* describes a load that cannot exist. 

1265 

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 

1271 

1272 

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. 

1277 

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. 

1284 

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) 

1302 

1303 

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). 

1316 

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 ) 

1341 

1342 

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 

1347 

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 ) 

1358 

1359 

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. 

1362 

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 

1368 

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 ) 

1377 

1378 

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 

1383 

1384 try: 

1385 return _is_moe(read_gguf_metadata(engine_params.resolve_model_path(ref))) 

1386 except (ProviderError, OSError): 

1387 return False 

1388 

1389 

1390def _cpu_offload_in_play() -> bool: 

1391 """Whether this configuration puts any of a model's weights in system memory. 

1392 

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 

1399 

1400 return _expert_offload_configured() or cfg.n_gpu_layers is not None 

1401 

1402 

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. 

1405 

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. 

1412 

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) 

1425 

1426 

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()) 

1430 

1431 

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. 

1434 

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. 

1438 

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 

1470 

1471 

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 ) 

1483 

1484 

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 ) 

1496 

1497 

1498def placeable_total_vram() -> int: 

1499 """Physical VRAM across all cards, for the weights-exceed placeability bound. 

1500 

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 

1511 

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 

1517 

1518 

1519def role_model_placeable(role: WorkerRole, ref: str, total_vram: int) -> bool: 

1520 """Whether a fresh plan would actually serve *role* on *ref*. 

1521 

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)) 

1535 

1536 

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. 

1545 

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 

1558 

1559 inputs: dict[WorkerRole, ModelPlacementInput] = {} 

1560 model_refs: dict[WorkerRole, str] = {} 

1561 skipped_not_installed: dict[WorkerRole, str] = {} 

1562 

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 

1589 

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) 

1598 

1599 ordered = [inputs[role] for role in ROLE_REGISTRY if role in inputs] 

1600 return ordered, model_refs, reservation, skipped_not_installed 

1601 

1602 

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. 

1609 

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 

1630 

1631 

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. 

1636 

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 

1642 

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 } 

1653 

1654 

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 

1669 

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 

1676 

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] 

1716 

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 ) 

1731 

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 ) 

1819 

1820 

1821def build_single_role_launch(role: WorkerRole, model_path: Path) -> InstanceLaunch: 

1822 """The launch the fleet would build for *role* serving *model_path*, alone. 

1823 

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. 

1829 

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 

1835 

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 ) 

1853 

1854 

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 

1858 

1859 

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} 

1867 

1868 

1869def _warn_gpu_present_but_unenumerated(binary: Path) -> None: 

1870 """Say so when the host has a GPU the engine did not list. 

1871 

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 

1879 

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 ) 

1898 

1899 

1900@dataclass(frozen=True) 

1901class DeviceReading: 

1902 """What one ``--list-devices`` run answered, as the whole fleet reads it. 

1903 

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 """ 

1909 

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 

1918 

1919 

1920def _read_devices(binary: Path) -> DeviceReading: 

1921 """Enumerate devices, the selected backend, and whether every GPU was refused. 

1922 

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. 

1926 

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 

1939 

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 

1980 

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) 

1998 

1999 

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" 

2019 

2020 

2021class _ReadDeviceCache: 

2022 """Short-TTL device-probe cache for the read/view path. 

2023 

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. 

2027 

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). 

2034 

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 """ 

2038 

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 

2047 

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 

2052 

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 

2070 

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 

2077 

2078 

2079_read_device_cache = _ReadDeviceCache(_DEVICE_PROBE_TTL_S, _DEVICE_PROBE_FAILURE_TTL_S) 

2080 

2081 

2082def clear_read_device_cache() -> None: 

2083 """Drop the read-path device probe cache (e.g. after the fleet is reconfigured). 

2084 

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 ) 

2093 

2094 _read_device_cache.clear() 

2095 vulkan_device_types_by_name.cache_clear() 

2096 integrated_vulkan_indices.cache_clear() 

2097 

2098 

2099@dataclass(frozen=True) 

2100class _PlanProbe: 

2101 """Clean-box memory snapshot every plan is sized against. 

2102 

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 """ 

2112 

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 

2126 

2127 

2128@dataclass(frozen=True) 

2129class _FailedProbe: 

2130 """An engine identity whose restate probe raised, and when it raised.""" 

2131 

2132 engine: str 

2133 at: float 

2134 

2135 

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]" 

2139 

2140 

2141class _PlanProbeStore: 

2142 """Holds the captured plan snapshot; a single instance below (no bare global).""" 

2143 

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() 

2153 

2154 def set(self, probe: _PlanProbe) -> None: 

2155 with self._lock: 

2156 self._probe = probe 

2157 self._failed = None 

2158 

2159 def get(self) -> _PlanProbe | None: 

2160 with self._lock: 

2161 return self._probe 

2162 

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()) 

2167 

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 ) 

2177 

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) 

2184 

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. 

2187 

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) 

2196 

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. 

2199 

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 

2209 

2210 def clear(self) -> None: 

2211 with self._lock: 

2212 self._probe = None 

2213 self._failed = None 

2214 

2215 

2216_plan_probe_store = _PlanProbeStore() 

2217 

2218 

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 

2226 

2227 

2228class _CtxDownshiftStore: 

2229 """How many halvings each role's auto context has taken after a load OOM. 

2230 

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 """ 

2236 

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] = {} 

2244 

2245 def steps(self, role: WorkerRole) -> int: 

2246 with self._lock: 

2247 return self._steps.get(role, 0) 

2248 

2249 def note_base(self, role: WorkerRole, ctx: int) -> None: 

2250 with self._lock: 

2251 self._base[role] = ctx 

2252 

2253 def base(self, role: WorkerRole) -> int | None: 

2254 with self._lock: 

2255 return self._base.get(role) 

2256 

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 

2262 

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) 

2271 

2272 

2273_ctx_downshift_store = _CtxDownshiftStore() 

2274 

2275 

2276def apply_ctx_downshift(role: WorkerRole, ctx: int) -> int: 

2277 """*ctx* halved once per downshift step recorded for *role*, floored. 

2278 

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. 

2285 

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 

2291 

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)) 

2296 

2297 

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 

2301 

2302 

2303def record_ctx_downshift(role: WorkerRole) -> bool: 

2304 """Take one downshift step for *role*; False when there is none left to take. 

2305 

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 

2310 

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 

2328 

2329 

2330def clear_ctx_downshift(role: WorkerRole | None = None) -> None: 

2331 """Forget *role*'s recorded downshift, or every role's, back to full size. 

2332 

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) 

2340 

2341 

2342def _probe_engine_devices() -> DeviceReading: 

2343 """Apply the fleet GPU/CUDA env, resolve the binary, and enumerate devices. 

2344 

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 

2352 

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) 

2360 

2361 

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. 

2364 

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. 

2370 

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 

2375 

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 

2392 

2393 

2394def assert_engine_probeable() -> None: 

2395 """Raise if the engine cannot be probed; capture no snapshot. 

2396 

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() 

2404 

2405 

2406def _engine_identity() -> str: 

2407 """Identity of the engine binary a plan snapshot has to agree with. 

2408 

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 

2416 

2417 

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) 

2424 

2425 

2426@contextmanager 

2427def _one_engine_per_pass() -> Iterator[None]: 

2428 """Pin the engine identity every snapshot read in this planning pass answers about. 

2429 

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) 

2442 

2443 

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() 

2448 

2449 

2450def _probe_engine_devices_and_identity() -> tuple[str, DeviceReading]: 

2451 """The engine's reading, tagged with the identity read before the probe ran. 

2452 

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() 

2459 

2460 

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 ) 

2474 

2475 

2476def _structural(devices: Iterable[FleetDevice]) -> tuple[FleetDevice, ...]: 

2477 """*devices* with the volatile free reading zeroed, so equality is structural. 

2478 

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) 

2484 

2485 

2486def refresh_plan_devices() -> None: 

2487 """Re-read which devices exist, keeping the clean-box memory figures. 

2488 

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. 

2492 

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. 

2498 

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()) 

2505 

2506 

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 ) 

2546 

2547 

2548def _current_plan_probe() -> _PlanProbe | None: 

2549 """The plan snapshot, restated first when another engine binary is now in place. 

2550 

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) 

2559 

2560 

2561def _restate_plan_probe(engine: str) -> _PlanProbe | None: 

2562 """Restate the snapshot for *engine*, once per burst of stale reads. 

2563 

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. 

2568 

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) 

2574 

2575 

2576def clear_plan_probe() -> None: 

2577 """Drop the plan snapshot (full fleet teardown); the next build re-captures.""" 

2578 _plan_probe_store.clear() 

2579 

2580 

2581def _cpu_pin_when_every_device_was_refused() -> tuple[str, ...]: 

2582 """``("none",)`` when the engine offered GPUs that lilbee refused, else empty. 

2583 

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,) 

2600 

2601 

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) 

2606 

2607 

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 

2611 

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()) 

2618 

2619 

2620def plan_sizing_is_unified() -> bool: 

2621 """Whether ctx sizing charges the shared-memory footprint rather than VRAM. 

2622 

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) 

2632 

2633 

2634def _device_sizing_budget(devices: Sequence[FleetDevice]) -> int: 

2635 """Memory one role may size its ctx and slots against, in bytes. 

2636 

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. 

2643 

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 

2650 

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) 

2654 

2655 

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 [] 

2662 

2663 

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() 

2668 

2669 

2670def _unreported_bytes(role: WorkerRole, mmproj: Path | None) -> int: 

2671 """Estimated bytes the engine allocates without printing a buffer line. 

2672 

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 

2684 

2685 

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. 

2688 

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 

2700 

2701 

2702def _device_names(devices: tuple[FleetDevice, ...]) -> tuple[str, ...]: 

2703 """``--device`` names for *devices*, empty when the backend pins through env. 

2704 

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. 

2712 

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) 

2729 

2730 

2731def _unified_memory_budget(devices: list[FleetDevice]) -> int | None: 

2732 """Shared-RAM placement budget (free RAM minus the OS floor), or ``None``. 

2733 

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 ) 

2748 

2749 

2750def _unified_admission_budget(devices: list[FleetDevice]) -> int | None: 

2751 """Shared-RAM pool a role set is *admitted* against, or ``None`` if dedicated. 

2752 

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 ) 

2766 

2767 

2768def _system_memory_floor() -> int: 

2769 """RAM held back for the OS when placing against system memory. 

2770 

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 

2776 

2777 total = model_cache.total_system_memory() 

2778 return min(int(cfg.system_memory_reserve_gb * 1024**3), total // _SYSTEM_MEMORY_FLOOR_DIVISOR) 

2779 

2780 

2781def _capped_by_device_memory(budget: int, devices: Sequence[FleetDevice]) -> int: 

2782 """*budget*, never above what the devices can address between them. 

2783 

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)) 

2794 

2795 

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. 

2798 

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. 

2803 

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. 

2810 

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} 

2819 

2820 

2821def _packable_devices(devices: list[FleetDevice]) -> list[FleetDevice]: 

2822 """The devices bin-packing may charge against. 

2823 

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. 

2830 

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 

2837 

2838 

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 ) 

2873 

2874 

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. 

2885 

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 

2926 

2927 

2928@dataclass(frozen=True) 

2929class ResolvedPlacement: 

2930 """Devices + resolved instance plans + model refs for the placement view.""" 

2931 

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 

2951 

2952 

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. 

2957 

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. 

2961 

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 

2969 

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 ) 

3000 

3001 

3002@dataclass(frozen=True) 

3003class FleetPlan: 

3004 """The servers to start, and the roles that share one swap group.""" 

3005 

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) 

3015 

3016 

3017def _log_placement_findings(placement: Placement, model_refs: dict[WorkerRole, str]) -> None: 

3018 """Warn about placements that exceed the memory budget. 

3019 

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 ) 

3047 

3048 

3049def _unusable_chat_ctx_reason(launch: InstanceLaunch) -> str | None: 

3050 """Reason to refuse a chat launch whose window cannot hold a grounded prompt. 

3051 

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 

3057 

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 ) 

3075 

3076 

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) 

3086 

3087 

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 

3096 

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 ) 

3143 

3144 

3145class _VisionRequestGrant: 

3146 """Whether a request that turned OCR on has asked for the vision role this process.""" 

3147 

3148 def __init__(self) -> None: 

3149 self._lock = threading.Lock() 

3150 self._granted = False 

3151 

3152 @property 

3153 def granted(self) -> bool: 

3154 with self._lock: 

3155 return self._granted 

3156 

3157 def set(self, granted: bool) -> None: 

3158 with self._lock: 

3159 self._granted = granted 

3160 

3161 

3162_vision_request_grant = _VisionRequestGrant() 

3163 

3164 

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) 

3168 

3169 

3170def revoke_vision_on_request() -> None: 

3171 """Plan the vision role from the OCR setting alone again.""" 

3172 _vision_request_grant.set(False) 

3173 

3174 

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 

3178 

3179 chosen = OcrBackendUsed.chosen(cfg.enable_ocr, ref) 

3180 return chosen is OcrBackendUsed.VISION or _vision_request_grant.granted 

3181 

3182 

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 

3186 

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 ) 

3192 

3193 

3194def plan_all_launches() -> FleetPlan: 

3195 """Apply GPU env, probe devices, and plan launches for the configured roles. 

3196 

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 

3202 

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)