Coverage for src/lilbee/providers/fleet/provider.py: 100%

1109 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-09-28 17:20 +0000

1"""FleetProvider: the local llama-server engine for every role. 

2 

3On first use it plans GPU placement and starts one llama-swap process per swap 

4group, each fronting that group's llama-server(s); each call routes to its role's 

5proxy by replica model id. Per-group processes let a reload restart only the 

6groups whose launches changed, so a placement or model change never unloads an 

7untouched group's model. There is no in-process fallback, so a missing role 

8surfaces a user-facing ``ProviderError``. Model management 

9(list/show/capabilities) reads the registry and GGUF headers directly and needs no 

10running server. See docs/architecture.md for swap tenancy. 

11""" 

12 

13from __future__ import annotations 

14 

15import functools 

16import logging 

17import re 

18import sys 

19import threading 

20import time 

21from contextlib import contextmanager 

22from pathlib import Path 

23from typing import TYPE_CHECKING, Any, Literal, NamedTuple, TypeVar, overload 

24 

25import httpx 

26 

27from lilbee.catalog import clean_display_name 

28from lilbee.core.config import cfg 

29from lilbee.core.health_warnings import HealthWarning 

30from lilbee.core.vectors import Vector 

31from lilbee.modelhub.registry import ModelRegistry 

32from lilbee.providers import engine_params 

33from lilbee.providers.base import ( 

34 GENERATION_RESERVE_TOKENS, 

35 ProviderError, 

36 ProviderErrorKind, 

37 prompt_token_budget, 

38) 

39from lilbee.providers.fleet import planning 

40from lilbee.providers.fleet.binary import engine_pin, resolve_llama_server 

41from lilbee.providers.fleet.client import ( 

42 ChatDeadlineError, 

43 LlamaServerClient, 

44 is_connection_failure, 

45 is_load_capacity_failure, 

46 is_rebuildable_failure, 

47 retry_on_busy, 

48) 

49from lilbee.providers.fleet.contract import ( 

50 chat_ctx_covers, 

51 contract_matches, 

52 decoded_launches, 

53 served_pairs, 

54 vision_slots_cover, 

55) 

56from lilbee.providers.fleet.groups import SwapGroup, group_for 

57from lilbee.providers.fleet.ingest_warmth import ingest_keep_warm 

58from lilbee.providers.fleet.launch import InstanceLaunch 

59from lilbee.providers.fleet.swap_config import cold_load_timeout_s 

60from lilbee.providers.fleet.swap_manager import ( 

61 SwapManager, 

62 SwapState, 

63 engine_record_exists, 

64 find_live_state, 

65 reap_stale, 

66 state_is_healthy, 

67 stop_engine, 

68) 

69from lilbee.providers.fleet.windowing import window_messages 

70from lilbee.providers.model_ref import parse_model_ref 

71from lilbee.providers.roles import MODEL_FIELD_TO_ROLE, WorkerRole, configured_model_message 

72from lilbee.providers.warm_progress import ( 

73 WarmPhase, 

74 WarmProgress, 

75 WarmProgressTracker, 

76 is_active_warm, 

77) 

78from lilbee.runtime.daemon_call import DaemonCall 

79from lilbee.runtime.engine_lock import ( 

80 ENGINE_DIR_ENV, 

81 UserLockHold, 

82 build_lock, 

83 hold_user_lock, 

84 keep_warm_requested, 

85 kernel_arbitrates_locks, 

86 live_users_exist, 

87 machine_engine_dir, 

88 private_engine_dir, 

89 request_keep_warm, 

90 withdraw_keep_warm, 

91) 

92 

93log = logging.getLogger(__name__) 

94 

95# How long a shutdown waits for an in-flight build to finish before tearing down 

96# regardless. Generous against a legitimate llama-swap spawn (a 30 s boot budget) 

97# and far short of any supervisor's patience for a process that will not exit. 

98_SHUTDOWN_BUILD_LOCK_WAIT_S = 45.0 

99 

100if TYPE_CHECKING: 

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

102 

103 from lilbee.providers.base import ( 

104 ChatMessage, 

105 ChatResult, 

106 ChatStreamItem, 

107 ChatToolResult, 

108 ClosableIterator, 

109 ) 

110 

111# User-facing name for this engine in error messages. 

112_PROVIDER_NAME = "llama-server" 

113# No-healthy-replica messages: the server stopped answering, or its process exited 

114# and the proxy restarts it on the next request. 

115_NOT_RESPONDING_MESSAGE = ( 

116 "The model server is not responding and no healthy replica is available. " 

117 "It may be restarting; try again in a moment." 

118) 

119_ENGINE_EXITED_MESSAGE = ( 

120 "The model server process exited and is being restarted. " 

121 "No other replica is available to take the request." 

122) 

123# Tokens held back from the served context for the model's own generation when the 

124# request does not cap it, plus a margin for chat-template overhead and estimate drift. 

125# Minimal input used to pre-load a role's upstream during warm-up (llama-swap 

126# starts an upstream on its first request, so warming issues one cheap call). 

127_WARM_PROMPT = "warm" 

128_WARM_MAX_TOKENS = 1 

129# Read size for paging chat shards into the page cache during warm; large enough 

130# to keep sequential reads efficient without holding much resident at once. 

131_PREWARM_CHUNK_BYTES = 8 * 1024 * 1024 

132# Shards fully paged in this boot, keyed on (path, size, mtime_ns); a fleet 

133# rebuild (e.g. a placement change) skips re-reading a hot cache. Module-level so 

134# it survives reset_services() replacing the provider instance. 

135_PREWARMED_SHARDS: set[tuple[str, int, int]] = set() 

136# Per-role client request budget: the first request covers the lazy cold load plus 

137# generation, so the weights-scaled cold-load budget plus the margin raises this floor. 

138# The embedder generates nothing, so its budget is the cold load plus the margin alone. 

139_REQUEST_TIMEOUT_FLOOR_S = 900.0 

140_REQUEST_TIMEOUT_GENERATION_MARGIN_S = 120.0 

141# Jinja chat templates flag tool support by referencing one of these names as an 

142# identifier inside a ``{% ... %}`` / ``{{ ... }}`` block (not free-text prose). 

143# The server parses tool calls natively via ``--jinja``; this probe only decides 

144# whether to offer tools to a given model at all. 

145_TOOL_TEMPLATE_PATTERN = re.compile(r"\{[%{][^}]*\b(?:tools|tool_calls|functions|function_calls)\b") 

146# Attempt cap for the busy-retry only when a page has no deadline (ocr_timeout=0, 

147# "no limit"): it backstops the retry so a persistently busy fleet can't spin 

148# forever. A page with a deadline retries until that deadline instead (see 

149# _ocr_dispatch), so the count doesn't bound the common case. 

150_VISION_BUSY_RETRIES = 18 

151# llama.cpp's repeat penalty on every vision OCR request: it damps tokens seen in 

152# the last window, which stops a page looping one line to the token cap. 

153_VISION_REPEAT_PENALTY = 1.1 

154_VISION_REPEAT_LAST_N = 64 

155# How often a waiter blocked on full replicas re-polls their health: an 

156# unhealthy replica re-admits itself by cool-down expiry, which notifies nobody. 

157_DISPATCH_HEALTH_RECHECK_S = 0.5 

158_T = TypeVar("_T") 

159 

160 

161def _prewarm_key(shard: Path) -> tuple[str, int, int]: 

162 """The prewarm identity of *shard*: same path, size, and mtime -> same pages.""" 

163 stat = shard.stat() 

164 return (str(shard), stat.st_size, stat.st_mtime_ns) 

165 

166 

167def _request_timeout_s(launch: InstanceLaunch) -> float: 

168 """Per-client request budget: the role's cold load plus margin, floored when it generates.""" 

169 budget = ( 

170 cold_load_timeout_s(launch.weights_bytes, launch.role) 

171 + _REQUEST_TIMEOUT_GENERATION_MARGIN_S 

172 ) 

173 if launch.role is WorkerRole.EMBED: 

174 return budget 

175 return max(_REQUEST_TIMEOUT_FLOOR_S, budget) 

176 

177 

178def _launches_by_group( 

179 plan: planning.FleetPlan, 

180) -> dict[SwapGroup, tuple[InstanceLaunch, ...]]: 

181 """Group a plan's launches by swap group, replica order preserved within each group. 

182 

183 Co-tenant roles land in one group, so llama-swap evicts between them rather 

184 than holding both resident. 

185 """ 

186 grouped: dict[SwapGroup, list[InstanceLaunch]] = {} 

187 for launch in plan.launches: 

188 grouped.setdefault(group_for(launch.role, plan.co_tenants), []).append(launch) 

189 return {group: tuple(group_launches) for group, group_launches in grouped.items()} 

190 

191 

192def _by_role(launches: list[InstanceLaunch]) -> dict[WorkerRole, list[InstanceLaunch]]: 

193 """Split one group's launches per role, replica order preserved.""" 

194 grouped: dict[WorkerRole, list[InstanceLaunch]] = {} 

195 for launch in launches: 

196 grouped.setdefault(launch.role, []).append(launch) 

197 return grouped 

198 

199 

200def _least_in_flight(clients: list[LlamaServerClient]) -> LlamaServerClient: 

201 """Pick the healthy client with the fewest in-flight requests. 

202 

203 Falls back to the full pool when every client is marked unhealthy, so a 

204 fully-dead pool still gets a call (which surfaces the error and lets a 

205 recovered replica mark itself healthy again). 

206 """ 

207 healthy = [client for client in clients if client.healthy] 

208 return min(healthy or clients, key=lambda c: c.in_flight) 

209 

210 

211# Serializes pick-and-reserve so concurrent routers see each other's assignment. 

212# Held only for the O(replicas) selection, never across the request itself. 

213_ROUTE_LOCK = threading.Lock() 

214 

215 

216def _reserve_least_in_flight(clients: list[LlamaServerClient]) -> LlamaServerClient: 

217 """Atomically pick the least-loaded healthy client and reserve a slot on it. 

218 

219 Selection and reservation are one critical section: without it, concurrent 

220 callers all read the same idlest replica before any of them increments its 

221 counter and route there together (a thundering herd that starves the rest of 

222 the fleet). The caller must :meth:`~LlamaServerClient.release` the slot. 

223 """ 

224 with _ROUTE_LOCK: 

225 client = _least_in_flight(clients) 

226 client.reserve() 

227 return client 

228 

229 

230def _healthy_groups_ours( 

231 states: dict[SwapGroup, SwapState], pin: str, wanted: set[tuple[WorkerRole, str]] 

232) -> bool: 

233 """Whether every healthy group in *states* is pin-equal and serves only wanted pairs. 

234 

235 True marks the incumbent as this contract's own engine that a full bind 

236 could not cover (a dead group, or config grew a role): the ladder rebuilds 

237 it in place even with live users, since those users need the rebuild too. 

238 False (a foreign pin or a model outside the contract) keeps the incumbent 

239 protected while in use. Vacuously False with no healthy group. 

240 """ 

241 if not states: 

242 return False 

243 for state in states.values(): 

244 if not contract_matches(state, (), pin): 

245 return False 

246 pairs = served_pairs(state) 

247 if pairs is None or not pairs <= wanted: 

248 return False 

249 return True 

250 

251 

252def _healthy_states(engine_dir: Path) -> dict[SwapGroup, SwapState]: 

253 """One probe pass over *engine_dir*: the recorded, answering group states. 

254 

255 The ladder's single view of a dir. Bind eligibility and the replaceability 

256 check both read this snapshot, so they cannot disagree about an engine that 

257 died between them, and one wedged proxy port is paid for once per ladder 

258 pass rather than once per decision -- all of it under the build lock, which 

259 every other lilbee start is waiting on. 

260 """ 

261 found: dict[SwapGroup, SwapState] = {} 

262 for group in SwapGroup: 

263 state = find_live_state(engine_dir, group) 

264 if state is not None and state_is_healthy(state): 

265 found[group] = state 

266 return found 

267 

268 

269def _bindable_group( 

270 state: SwapState, pin: str, wanted: set[tuple[WorkerRole, str]] 

271) -> tuple[SwapState, list[InstanceLaunch], set[tuple[WorkerRole, str]]] | None: 

272 """*state*'s launches and the wanted pairs it covers, or ``None``. 

273 

274 ``None`` for every reason an already-healthy group is not bindable by us: 

275 a foreign pin, an undecodable contract, or serving nothing we want. 

276 """ 

277 if not contract_matches(state, (), pin): 

278 # Pin mismatch or undecodable contract: not bindable by us. 

279 return None 

280 launches = decoded_launches(state) 

281 if launches is None: 

282 return None 

283 pairs = {(launch.role, launch.model) for launch in launches} & wanted 

284 return (state, launches, pairs) if pairs else None 

285 

286 

287def _log_adopted_launches( 

288 candidates: list[tuple[SwapGroup, SwapState, list[InstanceLaunch]]], 

289) -> None: 

290 """Name the engine every bound instance now runs on, and the pid that owns it.""" 

291 for _group, state, launches in candidates: 

292 for launch in launches: 

293 planning.log_engine_launch(launch, owner_pid=state.pid) 

294 

295 

296class _PrimedStream: 

297 """A stream re-fronted with its eagerly-pulled first frame. 

298 

299 close() always reaches the source stream, even before any iteration, so a 

300 caller that truncates immediately still releases the fleet's in-flight 

301 request slot (an unstarted chaining generator would silently drop it). 

302 """ 

303 

304 def __init__(self, first: ChatStreamItem, source: ClosableIterator[ChatStreamItem]) -> None: 

305 self._first: list[ChatStreamItem] = [first] 

306 self._source = source 

307 

308 def __iter__(self) -> _PrimedStream: 

309 return self 

310 

311 def __next__(self) -> ChatStreamItem: 

312 if self._first: 

313 return self._first.pop() 

314 return next(self._source) 

315 

316 def close(self) -> None: 

317 self._source.close() 

318 

319 

320def _primed_stream(items: ClosableIterator[ChatStreamItem]) -> ClosableIterator[ChatStreamItem]: 

321 """Pull the first frame of *items* now, so a dead engine raises to the caller. 

322 

323 The stream connects lazily on first iteration; without priming, a proxy 

324 that died raises only inside the consumer's loop, past any rediscovery. 

325 """ 

326 try: 

327 first = next(items) 

328 except StopIteration: 

329 return items # already exhausted; still closable 

330 return _PrimedStream(first, items) 

331 

332 

333def _call_with_failover( 

334 clients: list[LlamaServerClient], 

335 call: Callable[[LlamaServerClient], _T], 

336) -> _T: 

337 """Run *call* on the least-busy healthy client, retrying once on another replica. 

338 

339 The client is reserved at selection so concurrent ingest threads spread 

340 across replicas. A connection-level failure marks the client unhealthy and 

341 retries once on a different replica; with no other replica the failure 

342 surfaces. The reservation is released once the call resolves. 

343 """ 

344 client = _reserve_least_in_flight(clients) 

345 try: 

346 result = call(client) 

347 except Exception as exc: 

348 if not is_connection_failure(exc): 

349 raise 

350 client.mark_unhealthy() 

351 return _retry_on_other_replica(clients, client, call, exc) 

352 else: 

353 client.mark_healthy() 

354 return result 

355 finally: 

356 client.release() 

357 

358 

359def _retry_on_other_replica( 

360 clients: list[LlamaServerClient], 

361 failed: LlamaServerClient, 

362 call: Callable[[LlamaServerClient], _T], 

363 cause: Exception, 

364) -> _T: 

365 """Retry *call* once on a replica other than *failed*, marking its health.""" 

366 others = [c for c in clients if c is not failed] 

367 if not others: 

368 raise _no_healthy_replica_error(failed, cause) from cause 

369 retry = _reserve_least_in_flight(others) 

370 try: 

371 retry_result = call(retry) 

372 except Exception as retry_exc: 

373 if is_connection_failure(retry_exc): 

374 retry.mark_unhealthy() 

375 raise 

376 else: 

377 retry.mark_healthy() 

378 return retry_result 

379 finally: 

380 retry.release() 

381 

382 

383def _no_healthy_replica_error(failed: LlamaServerClient, cause: Exception) -> ProviderError: 

384 """User-facing error for a call with no healthy replica left to retry on. 

385 

386 A CONNECTION-kind ProviderError is the engine process exiting under the 

387 proxy (a transport error has no kind), so the message says so and names 

388 the log that records the exit. 

389 """ 

390 if isinstance(cause, ProviderError) and cause.kind is ProviderErrorKind.CONNECTION: 

391 message = _ENGINE_EXITED_MESSAGE 

392 if failed.engine_log is not None: 

393 message = f"{message} The engine log is {failed.engine_log}." 

394 else: 

395 message = _NOT_RESPONDING_MESSAGE 

396 return ProviderError(message, provider=_PROVIDER_NAME, kind=ProviderErrorKind.CONNECTION) 

397 

398 

399# Env vars a launch pins its devices with, one per backend (Metal has none). 

400_VISIBLE_DEVICE_ENV_VARS = ( 

401 "CUDA_VISIBLE_DEVICES", 

402 "ROCR_VISIBLE_DEVICES", 

403 "HIP_VISIBLE_DEVICES", 

404 "GGML_VK_VISIBLE_DEVICES", 

405 "ONEAPI_DEVICE_SELECTOR", 

406) 

407 

408 

409def _role_device_sets( 

410 launches: Iterable[InstanceLaunch], 

411) -> dict[WorkerRole, frozenset[str]]: 

412 """Backend-qualified device tokens each role's launches pin, by role. 

413 

414 A role's set is the union across its replicas. Roles whose launches carry 

415 no visibility env (Metal, or an unpinned backend) are absent: without 

416 pinning there is no proof of sharing, so they keep the concurrent warm. 

417 """ 

418 sets: dict[WorkerRole, set[str]] = {} 

419 for launch in launches: 

420 for var in _VISIBLE_DEVICE_ENV_VARS: 

421 value = launch.env_overrides.get(var) 

422 if value: 

423 sets.setdefault(launch.role, set()).update( 

424 f"{var}={part.strip()}" for part in value.split(",") 

425 ) 

426 return {role: frozenset(tokens) for role, tokens in sets.items()} 

427 

428 

429def _warm_chains( 

430 warm_roles: list[WorkerRole], device_sets: dict[WorkerRole, frozenset[str]] 

431) -> list[list[WorkerRole]]: 

432 """Group *warm_roles* into chains warmed sequentially; chains run in parallel. 

433 

434 Roles with overlapping device sets land in one chain, merged transitively. 

435 Within a chain chat goes last: it sizes its KV against the headroom the 

436 settled residents leave, so it must not race their loads. A role with no 

437 device set shares nothing provable and gets its own chain. 

438 """ 

439 chains: list[tuple[set[str], list[WorkerRole]]] = [] 

440 ordered = sorted(warm_roles, key=lambda r: (r is WorkerRole.CHAT, list(WorkerRole).index(r))) 

441 for role in ordered: 

442 tokens = device_sets.get(role) 

443 if not tokens: 

444 chains.append((set(), [role])) 

445 continue 

446 merged_tokens, merged_roles = set(tokens), [role] 

447 kept: list[tuple[set[str], list[WorkerRole]]] = [] 

448 for chain_tokens, chain_roles in chains: 

449 if chain_tokens & merged_tokens: 

450 merged_tokens |= chain_tokens 

451 merged_roles = chain_roles + merged_roles 

452 else: 

453 kept.append((chain_tokens, chain_roles)) 

454 kept.append((merged_tokens, merged_roles)) 

455 chains = kept 

456 return [roles for _tokens, roles in chains] 

457 

458 

459def _warm_role(role: WorkerRole, client: LlamaServerClient) -> None: 

460 """Send the cheapest request that loads *role*'s upstream behind llama-swap. 

461 

462 Vision is skipped (its load is heavy and it warms on the first OCR); chat, 

463 embed, and rerank each issue a minimal call to trigger the upstream start. 

464 """ 

465 if role is WorkerRole.CHAT: 

466 client.chat( 

467 [{"role": "user", "content": _WARM_PROMPT}], 

468 options={"max_tokens": _WARM_MAX_TOKENS}, 

469 stream=False, 

470 ) 

471 elif role is WorkerRole.EMBED: 

472 client.embed([_WARM_PROMPT]) 

473 elif role is WorkerRole.RERANK: 

474 client.rerank(_WARM_PROMPT, [_WARM_PROMPT]) 

475 

476 

477@functools.lru_cache(maxsize=32) 

478def _supports_tools_cached(path_str: str, _mtime_ns: int) -> bool: 

479 """Memoised tool-template probe keyed on the GGUF's path + mtime. 

480 

481 The mtime arg participates in the cache key only; a re-quantised file at the 

482 same path invalidates automatically because its mtime changes. 

483 """ 

484 from lilbee.providers.gguf_meta import read_gguf_metadata 

485 

486 meta = read_gguf_metadata(Path(path_str)) 

487 if not isinstance(meta, dict): 

488 return False 

489 template = meta.get("chat_template") 

490 if not isinstance(template, str): 

491 return False 

492 return _TOOL_TEMPLATE_PATTERN.search(template) is not None 

493 

494 

495class _VisionReplica(NamedTuple): 

496 """One vision server paired with its fitted ``--parallel`` slot count.""" 

497 

498 client: LlamaServerClient 

499 slots: int 

500 

501 

502class _PageBudgetExhausted(Exception): # noqa: N818 - internal control flow, not an error API 

503 """A page's document-wide OCR budget ran out before its slot came up.""" 

504 

505 

506class _VisionDispatcher: 

507 """Process-wide per-replica slot assignment for vision requests. 

508 

509 The ingest file fan-out runs many OCR requests at once; each request is 

510 assigned one specific replica and only while that replica has a free 

511 continuous-batching slot, so lilbee's own traffic can never oversubscribe a 

512 vision server into a 429 (an aggregate cap plus racy least-busy routing 

513 can). Requests past capacity wait in-process until any usable replica 

514 frees a slot; unhealthy replicas take no new work until their half-open 

515 cool-down re-admits them. 

516 """ 

517 

518 def __init__(self) -> None: 

519 self._cond = threading.Condition() 

520 self._assigned: dict[LlamaServerClient, int] = {} 

521 

522 @contextmanager 

523 def slot(self, pool: Sequence[_VisionReplica]) -> Iterator[LlamaServerClient]: 

524 """Hold one batching slot on the pool's best replica; yields that client.""" 

525 client = self._acquire(pool) 

526 try: 

527 yield client 

528 finally: 

529 self._release(client) 

530 

531 def _acquire(self, pool: Sequence[_VisionReplica]) -> LlamaServerClient: 

532 with self._cond: 

533 while True: 

534 client = self._pick(pool) 

535 if client is not None: 

536 self._assigned[client] = self._assigned.get(client, 0) + 1 

537 return client 

538 # The timed wait re-polls health: a replica can become routable 

539 # again by cool-down expiry alone, which notifies no waiter. 

540 self._cond.wait(timeout=_DISPATCH_HEALTH_RECHECK_S) 

541 

542 def _pick(self, pool: Sequence[_VisionReplica]) -> LlamaServerClient | None: 

543 """The usable replica with the most free slots, or None while all are full. 

544 

545 Falls back to the full pool when every replica is unhealthy (mirrors 

546 ``_least_in_flight``), so a dead pool surfaces the error instead of 

547 queueing forever. 

548 """ 

549 usable = [replica for replica in pool if replica.client.healthy] or list(pool) 

550 best = max(usable, key=self._free_slots) 

551 return best.client if self._free_slots(best) > 0 else None 

552 

553 def _free_slots(self, replica: _VisionReplica) -> int: 

554 return replica.slots - self._assigned.get(replica.client, 0) 

555 

556 def _release(self, client: LlamaServerClient) -> None: 

557 with self._cond: 

558 remaining = self._assigned.get(client, 0) - 1 

559 if remaining <= 0: 

560 self._assigned.pop(client, None) 

561 else: 

562 self._assigned[client] = remaining 

563 self._cond.notify_all() 

564 

565 

566_VISION_DISPATCHER = _VisionDispatcher() 

567 

568 

569def _dispatch_vision(pool: Sequence[_VisionReplica], call: Callable[[LlamaServerClient], _T]) -> _T: 

570 """Run *call* on a replica with a free batching slot, failing over once. 

571 

572 Blocks until a slot frees rather than racing requests at a full server. A 

573 connection-level failure marks the replica unhealthy and retries once on 

574 another replica's slot; with no other replica the failure surfaces. 

575 """ 

576 with _VISION_DISPATCHER.slot(pool) as client: 

577 try: 

578 result = call(client) 

579 except Exception as exc: 

580 if not is_connection_failure(exc): 

581 raise 

582 client.mark_unhealthy() 

583 failed, cause = client, exc 

584 else: 

585 client.mark_healthy() 

586 return result 

587 others = [replica for replica in pool if replica.client is not failed] 

588 if not others: 

589 raise _no_healthy_replica_error(failed, cause) from cause 

590 with _VISION_DISPATCHER.slot(others) as retry_client: 

591 try: 

592 retry_result = call(retry_client) 

593 except Exception as retry_exc: 

594 if is_connection_failure(retry_exc): 

595 retry_client.mark_unhealthy() 

596 raise 

597 retry_client.mark_healthy() 

598 return retry_result 

599 

600 

601def _vision_call( 

602 client: LlamaServerClient, messages: Sequence[Mapping[str, Any]], timeout: float | None 

603) -> str: 

604 """Run a vision chat on *client*, enforcing *timeout* like the in-process OCR. 

605 

606 The repeat penalty stops a page looping one line; ``cfg.vision_ocr_max_tokens`` 

607 caps generation if a loop still escapes it. A timeout surfaces as a 

608 ``ProviderError`` so the page-level OCR caller can fail just that page. 

609 Callers hold a dispatcher slot, so queue time isn't billed against the timeout. 

610 """ 

611 

612 options = { 

613 "max_tokens": cfg.vision_ocr_max_tokens, 

614 "repeat_penalty": _VISION_REPEAT_PENALTY, 

615 "repeat_last_n": _VISION_REPEAT_LAST_N, 

616 } 

617 if timeout and timeout > 0: 

618 return _bounded_vision_chat(client, messages, options, timeout) 

619 return client.chat(messages, options=options, stream=False) 

620 

621 

622def _bounded_vision_chat( 

623 client: LlamaServerClient, 

624 messages: Sequence[Mapping[str, Any]], 

625 options: dict[str, Any], 

626 timeout: float, 

627) -> str: 

628 """One vision chat streamed under a total *timeout*, released promptly on expiry. 

629 

630 ``chat_bounded`` streams the response in this thread and closes it (freeing the 

631 in-flight slot) once the deadline passes, so a trickling upstream can't outlive 

632 the caller. Its deadline signal is re-worded as the vision OCR timeout. 

633 """ 

634 try: 

635 return client.chat_bounded(messages, options=options, deadline_s=timeout) 

636 except ChatDeadlineError: 

637 raise ProviderError( 

638 f"Vision OCR timed out after {timeout:.0f}s.", 

639 provider=_PROVIDER_NAME, 

640 ) from None 

641 

642 

643def _ocr_dispatch( 

644 pool: Sequence[_VisionReplica], 

645 messages: Sequence[Mapping[str, Any]], 

646 deadline: float | None, 

647) -> str: 

648 """OCR *messages* on a free replica slot, retrying transient failures until *deadline*. 

649 

650 Backpressure (the dispatcher blocking until a slot frees) makes a 

651 self-inflicted 429 unreachable; a residual busy response is a still-warming 

652 server or foreign traffic, and a gateway error is a replica restarting 

653 mid-run. The retry is deadline-bound rather than 

654 attempt-bound so a page on a deep queue waits for a genuinely free slot until 

655 its own budget passes instead of dropping after a fixed count. Each attempt 

656 is bounded by the budget remaining before *deadline*; an exhausted budget 

657 raises :class:`_PageBudgetExhausted`. A ``None`` deadline (no limit) falls 

658 back to a bounded attempt count so the retry can't spin forever. 

659 """ 

660 

661 def _attempt(client: LlamaServerClient) -> str: 

662 remaining = max(0.0, deadline - time.monotonic()) if deadline is not None else None 

663 if remaining == 0.0: 

664 raise _PageBudgetExhausted 

665 return _vision_call(client, messages, remaining) 

666 

667 return retry_on_busy( 

668 lambda: _dispatch_vision(pool, _attempt), 

669 retries=_VISION_BUSY_RETRIES, 

670 deadline=deadline, 

671 ) 

672 

673 

674def _pdf_drain_budget(total_pages: int, per_page_timeout_s: float | None) -> float | None: 

675 """Total OCR wall-clock budget = pages*per_page + load grace, or None for no cap. 

676 

677 Mirrors the in-process drain budget: one document-wide deadline rather than a 

678 per-page cap, so a slow page borrows from fast ones and the vision model's cold 

679 first-inference is absorbed by the grace instead of tripping a fixed page limit. 

680 """ 

681 from lilbee.core.config import cfg 

682 

683 if not per_page_timeout_s or per_page_timeout_s <= 0: 

684 return None 

685 return total_pages * per_page_timeout_s + cfg.vision_load_budget_s 

686 

687 

688def _ocr_deadline(per_page_timeout_s: float | None) -> float | None: 

689 """Absolute monotonic deadline for one image OCR, or None when uncapped. 

690 

691 An image is a one-page document, so it gets the same budget as a PDF page: 

692 the per-page timeout plus the cold-load grace, spanning queue wait and 

693 generation together. 

694 """ 

695 budget = _pdf_drain_budget(1, per_page_timeout_s) 

696 return None if budget is None else time.monotonic() + budget 

697 

698 

699_ROLE_TO_MODEL_FIELD = {role: field for field, role in MODEL_FIELD_TO_ROLE.items()} 

700 

701 

702def _configured_model_for(role: WorkerRole) -> str: 

703 """The cfg model ref for *role*, empty when the role is unset.""" 

704 field = _ROLE_TO_MODEL_FIELD.get(role) 

705 return getattr(cfg, field) or "" if field else "" 

706 

707 

708def _unusable_engine_reason() -> str | None: 

709 """Why no server can start on this host, or None once an engine resolves. 

710 

711 Planning drops an engine-less host to serving nothing and says so only at 

712 debug, so by the time a surface has an empty pool the engine is the one cause 

713 it cannot see. Re-resolving here is also what keeps an engine installed 

714 mid-session from being reported as still missing. 

715 """ 

716 try: 

717 resolve_llama_server() 

718 except ProviderError as exc: 

719 return str(exc) 

720 return None 

721 

722 

723def _chat_needs_local_engine() -> bool: 

724 """Whether the configured chat model is one this host has to serve itself. 

725 

726 A chat ref routed to an SDK backend runs without any local engine, so a 

727 missing one is not its failure and must not be stamped on its warm. 

728 """ 

729 ref = _configured_model_for(WorkerRole.CHAT) 

730 return bool(ref) and not parse_model_ref(ref).is_remote 

731 

732 

733def _no_server_message(role: WorkerRole) -> str: 

734 """User-facing reason *role* has no server, engine state first. 

735 

736 A missing engine and a model that never placed both arrive as an empty pool, 

737 and reading the second onto the first sends the reader to a model 

738 configuration that is already correct. 

739 """ 

740 reason = _unusable_engine_reason() 

741 if reason is not None: 

742 return f"No {role.value} model server is running: {reason}" 

743 return ( 

744 f"No {role.value} model server is running. Make sure the {role.value} " 

745 "model is installed and configured, then try again." 

746 ) 

747 

748 

749class _EngineDemand(NamedTuple): 

750 """What this process needs an engine to serve: pairs plus its chat window.""" 

751 

752 pairs: set[tuple[WorkerRole, str]] 

753 # Per-slot chat tokens this process needs; 0 demands nothing. 

754 chat_ctx: int 

755 # Configured roles the plan skipped because their model is not installed. 

756 # Carried out of the demand plan so the warm tracker can name the missing 

757 # model even when the ladder never reaches _plan_and_spawn (zero installed 

758 # models fail _can_build_engine first) or binds an existing engine. 

759 skipped_not_installed: dict[WorkerRole, str] 

760 # Launches the demand plan refused for an unusable window (role -> reason); 

761 # recorded even when the ladder binds an engine or never builds one. 

762 skipped_unusable_ctx: dict[WorkerRole, str] 

763 # Vision slots this process needs at once; 0 demands nothing. 

764 vision_slots: int = 0 

765 

766 

767def _placeable_demand() -> _EngineDemand: 

768 """Configured (role, model) pairs a fresh plan would serve, and the chat window. 

769 

770 A configured role is wanted only when the planner would place it. The plan 

771 is the co-placement authority: a role that fits alone but cannot co-tenant (a 

772 unified-memory box past its budget) gets no launch, so it is dropped here too, 

773 and bind matches a running engine instead of judging it a partial cover and 

774 restarting the shared engine on every process start. The per-role check stays 

775 as the cheap gate for the reasons in its own docstring. Empty when no engine 

776 binary resolves: nothing is placeable, so the ladder serves nothing. 

777 """ 

778 from lilbee.providers.fleet.planning import ( 

779 placeable_total_vram, 

780 plan_all_launches, 

781 role_model_placeable, 

782 ) 

783 

784 try: 

785 plan = plan_all_launches() 

786 except ProviderError as exc: 

787 if exc.kind is ProviderErrorKind.NOT_FOUND: 

788 return _EngineDemand(set(), 0, {}, {}) 

789 raise 

790 placed = {launch.role for launch in plan.launches} 

791 total_vram = placeable_total_vram() 

792 pairs = { 

793 (role, model) 

794 for role in WorkerRole 

795 if role in placed 

796 and (model := _configured_model_for(role)) 

797 and role_model_placeable(role, model, total_vram) 

798 } 

799 return _EngineDemand( 

800 pairs, 

801 _demanded_chat_ctx(plan.launches, pairs), 

802 dict(plan.skipped_not_installed), 

803 dict(plan.skipped_unusable_ctx), 

804 _demanded_vision_slots(pairs), 

805 ) 

806 

807 

808def _demanded_vision_slots(pairs: set[tuple[WorkerRole, str]]) -> int: 

809 """Vision slots this process needs an engine to serve at once; 0 for none.""" 

810 if not any(role is WorkerRole.VISION for role, _ in pairs): 

811 return 0 

812 return max(1, cfg.vision_ocr_concurrency) 

813 

814 

815def _demanded_chat_ctx( 

816 launches: Iterable[InstanceLaunch], pairs: set[tuple[WorkerRole, str]] 

817) -> int: 

818 """Per-slot chat window this process needs an engine to serve; 0 for none. 

819 

820 The cfg target (a ``num_ctx`` pin, else ``chat_n_ctx_target``) capped by 

821 this process's own planned chat window: a window the plan itself cannot 

822 reach (model ceiling, hardware) is not a demand a rebuild could satisfy, 

823 so capping keeps the fit check from rebuilding the engine in a loop. 

824 

825 The cap applies only to a single-device chat plan, whose window is sized 

826 against device totals and holds regardless of what is resident. A 

827 tensor-split plan is sized against live free VRAM, which a resident 

828 incumbent deflates; capping by it would shrink the demand to whatever the 

829 incumbent left free and let the fit check pass vacuously. 

830 """ 

831 if not any(role is WorkerRole.CHAT for role, _model in pairs): 

832 return 0 

833 chat_launches = [launch for launch in launches if launch.role is WorkerRole.CHAT] 

834 planned = max((launch.ctx for launch in chat_launches), default=0) 

835 if planned <= 0: 

836 return 0 

837 # Always positive: num_ctx validates ge=1 and chat_n_ctx_target ge=512. 

838 target = cfg.num_ctx if cfg.num_ctx is not None else cfg.chat_n_ctx_target 

839 split = any(len(launch.est_vram_by_device) > 1 for launch in chat_launches) 

840 return target if split else min(target, planned) 

841 

842 

843def _can_build_engine(wanted: set[tuple[WorkerRole, str]]) -> bool: 

844 """Preconditions for a viable build, checked before stopping a warm engine. 

845 

846 A process that can serve nothing (no placeable model, an unresolvable engine 

847 binary) must not stop an engine another setup left warm and then spawn nothing. 

848 Probing the engine here resolves the binary AND enumerates devices, so a wedged 

849 GPU probe or an unusable CUDA runtime raises loud at this point -- before the 

850 caller stops a replaceable incumbent. Were the probe left to run only inside 

851 ``_plan_and_spawn`` (after the stop), that raise would kill an engine other 

852 members still hold and then skip the overflow build, leaving zero engines. This 

853 takes no memory snapshot (device enumeration reads no residency); the clean-box 

854 sizing snapshot is captured by ``_plan_and_spawn`` after the stop. 

855 """ 

856 from lilbee.providers.fleet import planning 

857 

858 if not wanted: 

859 return False 

860 try: 

861 planning.assert_engine_probeable() 

862 except ProviderError as exc: 

863 # A genuinely-missing engine binary keeps the quiet serve-nothing path; 

864 # every other probe failure must propagate (fail loud) rather than be read 

865 # as "cannot build" and silently stand down. 

866 if exc.kind is not ProviderErrorKind.NOT_FOUND: 

867 raise 

868 return False 

869 except OSError: 

870 return False 

871 return True 

872 

873 

874def _warm_ttl_seconds(*, hold_warm_for_session: bool = False) -> int: 

875 """llama-swap idle-unload timer in seconds for the spawned fleet. 

876 

877 A ttl of 0 keeps weights resident until the engine is stopped; otherwise an 

878 idle engine releases its weights after ``engine_idle_ttl_minutes`` and reloads 

879 transparently on the next prompt. The timer is held off (ttl 0) while someone 

880 who owns the engine's lifetime depends on an instant response: an interactive 

881 session (*hold_warm_for_session*) whose engine dies with it, or a bulk ingest 

882 whose unevenly loaded replicas must not idle-unload and reload cold mid-run. 

883 

884 A ``keep_engine_warm`` engine outlives every holder, so no holder may pin its 

885 weights: llama-swap's ttl is fixed at launch, and a session hold baked into a 

886 detached engine would keep the weights loaded for as long as the machine 

887 stays up. Keep-warm keeps the engine process (a small proxy, no weights) 

888 alive across launches; the weights follow the idle window like every other 

889 mode, and only the user's explicit 0 keeps them loaded. 

890 """ 

891 if not cfg.keep_engine_warm and (hold_warm_for_session or ingest_keep_warm()): 

892 return 0 

893 return cfg.engine_idle_ttl_minutes * 60 

894 

895 

896class FleetProvider: 

897 """Routes every role to the managed llama-server fleet (a fleet-of-one on one box).""" 

898 

899 def __init__(self, *, hold_warm: bool = False) -> None: 

900 # An interactive session (the TUI) owns this process for its whole 

901 # lifetime, so a fleet bound to it stays resident instead of idle-unloading 

902 # under a user who is still in the app; closing lilbee releases it. A 

903 # keep-warm fleet outlives the session, so _warm_ttl_seconds ignores the 

904 # hold for it. Set by the container that built this provider, never 

905 # mutated afterwards. 

906 self._hold_warm_for_session = hold_warm 

907 # One llama-swap per placed group, so restarting one group's servers (a 

908 # placement or per-role model change) never unloads another group's. A 

909 # co-tenant group holds chat and vision, which evict each other on load. 

910 self._swaps: dict[SwapGroup, SwapManager] = {} 

911 # The group each placed role runs in, so a role's clients and its swap 

912 # process can be reached from the role alone. 

913 self._role_group: dict[WorkerRole, SwapGroup] = {} 

914 # The launches each running group was started with, kept so a reload can 

915 # diff the fresh plan against what is running and restart only the groups 

916 # whose launches actually changed. Launch argv is port-free (ports are 

917 # injected at config render), so the comparison is stable across starts. 

918 self._launches: dict[SwapGroup, tuple[InstanceLaunch, ...]] = {} 

919 # Engine dirs this provider holds membership in (machine slot and/or 

920 # the private overflow), and the dir each running group lives in. 

921 self._engine_holds: dict[Path, UserLockHold] = {} 

922 self._group_dirs: dict[SwapGroup, Path] = {} 

923 # Latched once shutdown runs. A discarded provider (reset_services swaps 

924 # in a new one) can still have an in-flight warm-up or reload daemon 

925 # thread; without this latch that thread could start a llama-swap after 

926 # shutdown already ran, leaving a process no live provider owns. 

927 # _ensure_fleet checks it under the build lock so a post-shutdown build is 

928 # refused (the swap_manager reaper is the backstop if one slips through). 

929 self._shut_down = False 

930 # A pool of OpenAI clients per placed role (one per data-parallel replica), 

931 # all pointed at the llama-swap endpoint and routed by replica model id; 

932 # rebuilt whenever the swap process (re)starts. Requests round-robin the pool. 

933 self._clients: dict[WorkerRole, list[LlamaServerClient]] = {} 

934 # Clients retired by a reload, awaiting close. A reload's old clients may 

935 # still be held by an in-flight reader, so they are closed at the *next* 

936 # reload (by when those readers have finished) or at shutdown, never while 

937 # potentially in use. See _retire_clients. 

938 self._retiring_clients: list[LlamaServerClient] = [] 

939 # Chat batching slots and per-slot context from the chat launch, surfaced to 

940 # the concurrency gate and clients; defaults until the chat group is up. 

941 self._chat_slots = 1 

942 self._chat_ctx: int | None = None 

943 # Latest chat prefill progress reported by a streaming client, cleared 

944 # by the same stream when generation starts or the stream ends. 

945 self._chat_prefill: tuple[int, int] | None = None 

946 # Single-flight guard: the HTTP/MCP servers route concurrently, so two 

947 # first-requests must not each start a swap (double GPU allocation) or 

948 # tear one down mid-route. Reentrant: invalidate_load_cache nests calls. 

949 self._lock = threading.RLock() 

950 # Serializes the slow startup (GPU probe + GGUF parse + llama-swap spawn) 

951 # across concurrent callers, so the off-thread warm-up and an on-demand call 

952 # can't start two swaps. Held only during startup, NOT while routing. 

953 self._build_lock = threading.Lock() 

954 # Single-flight for starting vision on request: concurrent OCR pages wait 

955 # for the first page's re-plan instead of each re-planning. 

956 self._vision_request_lock = threading.Lock() 

957 # Spawn-lifecycle listeners (set by the TUI via add_spawn_listener). Stored 

958 # so warm-up can report per-role progress as it pre-loads each upstream. 

959 self._on_spawning: Callable[[WorkerRole], None] | None = None 

960 self._on_spawned: Callable[[WorkerRole], None] | None = None 

961 # Granular cold-load progress for the chat role, streamed to a launcher so 

962 # the user sees real read/engine-load progress instead of a frozen spinner. 

963 self._warm_tracker = WarmProgressTracker() 

964 # The engine's own reason a role's model failed to warm, so the launcher and 

965 # the TUI report the real cause instead of a generic "did not load". 

966 self._warm_errors: dict[WorkerRole, str] = {} 

967 # Configured roles the last plan left unplaced because their model isn't 

968 # installed (role -> ref). The warm finalizer reads it to fail a not-installed 

969 # chat with a named reason instead of clearing to a silent "not ready" retry. 

970 self._skipped_not_installed: dict[WorkerRole, str] = {} 

971 # Launches the last plan refused for an unusable window (role -> reason); 

972 # read by the warm finalizer and _require_clients. 

973 self._skipped_unusable_ctx: dict[WorkerRole, str] = {} 

974 # Single-flight guard for the off-thread warm-up: True from the moment a 

975 # warm thread is dispatched until it finishes, so a second warm_up_pool 

976 # never starts a second swap and double-allocates GPU memory. 

977 self._warming = False 

978 # Request-triggered lazy-warm single-flight: one request begins the warm 

979 # cycle, concurrent siblings share it instead of each resetting. 

980 self._lazy_warming = False 

981 # Single-flight guard for the off-thread reload: a second reload_role 

982 # while one is in flight sets the pending flag instead of dispatching, 

983 # and the in-flight thread re-runs the plan loop once per pending flag. 

984 self._reloading = False 

985 # Set when a reload arrives mid-reload: the in-flight pass may have 

986 # already snapshotted its plan, so the change must be re-applied. 

987 self._reload_pending = False 

988 # Notified when ``_reloading`` clears, so a ``reload_role(wait=True)`` caller 

989 # can block until the reload it requested (or the in-flight one that will 

990 # run its pending pass) has finished. 

991 self._reload_done = threading.Condition(self._lock) 

992 

993 def _ensure_fleet(self) -> bool: 

994 """Start one llama-swap per placed role exactly once across concurrent callers. 

995 

996 Returns whether any role group is running afterwards; ``False`` when no 

997 role is configured and installed (nothing to serve), leaving no process 

998 spawned. The startup runs under ``_build_lock`` (not the routing lock), 

999 so the off-thread warm-up and an on-demand call can't start two fleets -- 

1000 which would double-allocate GPU and parse the same GGUF twice. A second 

1001 caller blocks on the build lock and reuses the groups the first one 

1002 started. A group failing to start tears down the groups already started 

1003 in this build, so a partial fleet never leaks past the failure. 

1004 """ 

1005 with self._lock: 

1006 if self._swaps: 

1007 return True 

1008 with self._build_lock: 

1009 with self._lock: 

1010 if self._swaps: 

1011 return True 

1012 if self._shut_down: 

1013 # Provider was shut down (and likely discarded by reset_services) 

1014 # while this warm-up/reload thread was in flight; do not spawn a 

1015 # llama-swap no live provider would ever reap. 

1016 return False 

1017 

1018 return self._acquire_engine(cfg.data_root) 

1019 

1020 def _acquire_engine(self, config_root: Path) -> bool: 

1021 """The acquisition ladder: bind to a compatible engine, else build one. 

1022 

1023 Machine slot first. An incumbent is replaced in place when no live 

1024 user holds it, or when it is this contract's own engine (pin-equal, 

1025 serving only wanted models) left partially dead or partially covering: 

1026 its members are waiting for exactly that rebuild, and overflowing 

1027 around it would load duplicate weights. Only a live incompatible 

1028 engine in active use sends the build to the config root's private 

1029 overflow dir. Per-dir the step is all-or-nothing: every configured 

1030 (role, model) pair bound, or built fresh. Runs under the cross-process 

1031 build lock, so two starts never both build and stop-if-last never 

1032 races an arrival. 

1033 """ 

1034 pin = engine_pin() 

1035 demand = _placeable_demand() 

1036 # Record the demand plan's skips before walking the ladder: a bind or an 

1037 # early serve-nothing exit never reaches _plan_and_spawn, and the warm 

1038 # tracker must still be able to say "chat model X is not installed" 

1039 # rather than a retryable not-ready. 

1040 self._skipped_not_installed = dict(demand.skipped_not_installed) 

1041 self._skipped_unusable_ctx = dict(demand.skipped_unusable_ctx) 

1042 machine_dir = machine_engine_dir() 

1043 if kernel_arbitrates_locks(machine_dir): 

1044 machine = self._acquire_in_dir(machine_dir, pin, demand, is_overflow=False) 

1045 if machine is not None: 

1046 return machine 

1047 else: 

1048 # Without kernel-arbitrated locks the membership refcount cannot be 

1049 # trusted, and sharing is exactly what needs it: a probe would 

1050 # destroy a live member's lock, so the slot would look free while 

1051 # another setup is serving from it. Keep to our own dir instead. 

1052 log.warning( 

1053 "Engine dir %s is on a filesystem without working file locks; " 

1054 "using a private engine instead of the shared one. Set %s to a " 

1055 "path on a local filesystem to share one engine across lilbees.", 

1056 machine_dir, 

1057 ENGINE_DIR_ENV, 

1058 ) 

1059 # The machine slot holds a live incompatible engine in active use: overflow 

1060 # to this config root's private dir rather than evict another model setup. 

1061 private = private_engine_dir(config_root) 

1062 return self._acquire_in_dir(private, pin, demand, is_overflow=True) or False 

1063 

1064 def _acquire_in_dir( 

1065 self, engine_dir: Path, pin: str, demand: _EngineDemand, *, is_overflow: bool 

1066 ) -> bool | None: 

1067 """Bind or build one engine dir; ``None`` on the slot means overflow next. 

1068 

1069 Binds a compatible running engine. Whether an incumbent may be replaced 

1070 is decided by kernel-refcounted membership, not the proxy HTTP probe: an 

1071 engine with a live user is never reaped or stopped, so a transient probe 

1072 failure (fd exhaustion, host thrash) cannot kill a busy engine. Replace in 

1073 place only when no live user holds it or it is this contract's own engine 

1074 (pin-equal, serving only wanted models). A live incompatible engine in 

1075 active use is never evicted or stacked on: on the machine slot it returns 

1076 ``None`` (overflow), and in the overflow dir it serves nothing rather than 

1077 duplicate weights beside it. Before building, any recorded engine is cleared 

1078 -- keyed on 

1079 the state file, not the probe -- so an unprobeable incumbent is stopped 

1080 rather than double-built beside. The stop is gated on ``_can_build_engine`` 

1081 so a process that can serve nothing never destroys a warm engine it can't 

1082 replace. Held under the cross-process build lock. 

1083 """ 

1084 wanted = demand.pairs 

1085 with build_lock(engine_dir): 

1086 states = _healthy_states(engine_dir) 

1087 if wanted and self._bind_all_in_dir(engine_dir, states, pin, demand): 

1088 self._hold_membership(engine_dir) 

1089 return True 

1090 replaceable = not live_users_exist(engine_dir) or _healthy_groups_ours( 

1091 states, pin, wanted 

1092 ) 

1093 if not replaceable: 

1094 # A live engine another setup is actively using is never evicted or 

1095 # stacked on. On the machine slot that means overflow (None); in the 

1096 # overflow dir there is nowhere further to go, so serve nothing rather 

1097 # than kill the incumbent or load a second fleet's weights beside it 

1098 # (an OOM on a small-VRAM box). 

1099 return None if not is_overflow else False 

1100 if not _can_build_engine(wanted): 

1101 return False 

1102 # No live user holds this dir now (or it is ours to rebuild): reap dead 

1103 # leftovers and stop any recorded engine so planning sees true free VRAM 

1104 # and the build never lands beside an unprobeable incumbent. 

1105 reap_stale(engine_dir) 

1106 if engine_record_exists(engine_dir): 

1107 stop_engine(engine_dir) 

1108 if self._plan_and_spawn(engine_dir): 

1109 self._hold_membership(engine_dir) 

1110 return True 

1111 return False 

1112 

1113 def _bind_all_in_dir( 

1114 self, 

1115 engine_dir: Path, 

1116 states: dict[SwapGroup, SwapState], 

1117 pin: str, 

1118 demand: _EngineDemand, 

1119 ) -> bool: 

1120 """Bind every group needed to cover the demanded pairs; False leaves nothing bound. 

1121 

1122 Binding never touches groups serving models outside the demand; the dir 

1123 matches only when healthy, pin-equal groups cover every wanted pair and 

1124 the served chat window covers the demanded per-slot ctx. (Whether an 

1125 unmatched dir's engine is then replaced or overflowed around is the 

1126 ladder's call, based on live users.) 

1127 """ 

1128 wanted = demand.pairs 

1129 candidates: list[tuple[SwapGroup, SwapState, list[InstanceLaunch]]] = [] 

1130 covered: set[tuple[WorkerRole, str]] = set() 

1131 for group, state in states.items(): 

1132 found = _bindable_group(state, pin, wanted) 

1133 if found is None: 

1134 continue 

1135 bindable, launches, pairs = found 

1136 if not chat_ctx_covers(launches, demand.chat_ctx) or not vision_slots_cover( 

1137 launches, demand.vision_slots 

1138 ): 

1139 # The live chat window or vision slots are below this process's need. 

1140 return False 

1141 candidates.append((group, bindable, launches)) 

1142 covered |= pairs 

1143 if covered != wanted: 

1144 return False 

1145 bound: dict[SwapGroup, tuple[SwapManager, list[InstanceLaunch]]] = {} 

1146 for group, state, launches in candidates: 

1147 swap = SwapManager(engine_dir, group) 

1148 if not swap.bind(state): 

1149 for prior, _launches in bound.values(): 

1150 prior.shutdown() 

1151 return False 

1152 bound[group] = (swap, launches) 

1153 with self._lock: 

1154 for group, (swap, launches) in bound.items(): 

1155 self._adopt_group(group, swap, launches) 

1156 self._group_dirs[group] = engine_dir 

1157 log.info("Bound to the running engine at %s", engine_dir) 

1158 _log_adopted_launches(candidates) 

1159 return True 

1160 

1161 def _reload_dir(self) -> Path: 

1162 """The engine dir a reload rebuilds into: where our groups already live. 

1163 

1164 All this provider's groups share one dir by construction (the ladder is 

1165 all-or-nothing per dir); an empty provider rebuilds into the machine slot. 

1166 """ 

1167 with self._lock: 

1168 dirs = set(self._group_dirs.values()) 

1169 return next(iter(dirs)) if dirs else machine_engine_dir() 

1170 

1171 def _hold_membership(self, engine_dir: Path) -> None: 

1172 """Record this process as a user of *engine_dir*'s engine. 

1173 

1174 The single point every acquisition passes through, bind and build alike, 

1175 so it is also where this user's persistence opt-in is recorded against 

1176 the engine. Marking on bind (not only on build) is what makes the 

1177 setting mean what it says on a shared slot: a user who asked for a warm 

1178 engine keeps it warm even when a default-config sibling is last out. 

1179 """ 

1180 from lilbee.core.config import cfg 

1181 

1182 if engine_dir not in self._engine_holds: 

1183 self._engine_holds[engine_dir] = hold_user_lock(engine_dir) 

1184 if cfg.keep_engine_warm: 

1185 request_keep_warm(engine_dir, cfg.data_root) 

1186 

1187 def _plan_and_spawn(self, data_dir: Path) -> bool: 

1188 """Plan placement against the clean box and start one swap per group. 

1189 

1190 Caller holds the build lock. False when the engine binary is missing or 

1191 nothing is installed/configured, so the provider serves nothing. 

1192 """ 

1193 try: 

1194 # Snapshot the clean box; this plan and every later reload size 

1195 # ctx, slots, and budgets against it (a live probe under a loaded 

1196 # fleet would report our own residency as unavailable). Inside the 

1197 # try: capturing resolves the engine binary, and a binary-less 

1198 # host must serve nothing, not raise. Every other planning failure 

1199 # (a wedged GPU probe, an unusable CUDA runtime) propagates so the 

1200 # warm tracker and the caller report the real reason instead of a 

1201 # silent never-ready fleet. 

1202 planning.capture_plan_probe() 

1203 plan = planning.plan_all_launches() 

1204 except ProviderError as exc: 

1205 # Only a genuinely-missing engine binary keeps the quiet no-fleet path; 

1206 # any other planning failure (a wedged GPU probe, an unusable CUDA 

1207 # runtime) must surface to the warm tracker and on-demand callers 

1208 # rather than silently serving nothing (#540). 

1209 if exc.kind is not ProviderErrorKind.NOT_FOUND: 

1210 raise 

1211 log.debug("Engine binary unavailable; no swap started") 

1212 plan = None 

1213 # plan None (no engine binary) keeps the demand-time record from 

1214 # _acquire_engine instead of wiping it. 

1215 if plan is not None: 

1216 self._skipped_not_installed = dict(plan.skipped_not_installed) 

1217 self._skipped_unusable_ctx = dict(plan.skipped_unusable_ctx) 

1218 if plan is None or not plan.launches: 

1219 # No engine binary, or no installed/configured model: serve nothing. 

1220 return False 

1221 by_group = _launches_by_group(plan) 

1222 started: dict[SwapGroup, SwapManager] = {} 

1223 try: 

1224 for group, group_launches in by_group.items(): 

1225 swap = SwapManager(data_dir, group) 

1226 swap.start( 

1227 list(group_launches), 

1228 ttl_seconds=_warm_ttl_seconds( 

1229 hold_warm_for_session=self._hold_warm_for_session 

1230 ), 

1231 bind_lifetime=not cfg.keep_engine_warm, 

1232 ) 

1233 started[group] = swap 

1234 except BaseException: 

1235 for swap in started.values(): 

1236 swap.shutdown() 

1237 raise 

1238 with self._lock: 

1239 for group, swap in started.items(): 

1240 self._adopt_group(group, swap, list(by_group[group])) 

1241 self._group_dirs[group] = data_dir 

1242 return True 

1243 

1244 def _adopt_group( 

1245 self, group: SwapGroup, swap: SwapManager, launches: list[InstanceLaunch] 

1246 ) -> None: 

1247 """Record *group*'s freshly started swap and build a client pool per role. 

1248 

1249 Caller holds ``self._lock``. Each launch (one per replica) becomes a client 

1250 keyed by its replica model id against this group's own proxy endpoint; 

1251 the chat launch carries the slots/ctx so the capacity and served context 

1252 come from the launch, not a probe. 

1253 """ 

1254 self._swaps[group] = swap 

1255 self._launches[group] = tuple(launches) 

1256 endpoint = swap.endpoint() 

1257 for role, role_launches in _by_role(launches).items(): 

1258 # Retire the role's previous clients (a reload re-adopts over an existing 

1259 # pool): closing them now would error a reader still mid-call on an old 

1260 # client snapshot, and never closing leaks an httpx pool per replica. 

1261 old_clients = list(self._clients.get(role, [])) 

1262 self._role_group[role] = group 

1263 # token_cap truncates oversize embed/rerank inputs to the per-slot context 

1264 # (the in-process backstop); the longer timeout covers a cold upstream load. 

1265 self._clients[role] = [ 

1266 LlamaServerClient( 

1267 endpoint, 

1268 launch.model_id, 

1269 token_cap=launch.token_cap, 

1270 timeout=_request_timeout_s(launch), 

1271 rerank_mode=launch.rerank_mode, 

1272 inline_reasoning=role is WorkerRole.CHAT, 

1273 on_prefill=self._record_chat_prefill if role is WorkerRole.CHAT else None, 

1274 engine_log=swap.log_path, 

1275 # A cold embed replica 429s bulk ingest until its slots load; wait 

1276 # out the same cold-load budget llama-swap keeps it alive for so a 

1277 # burst never drops files while the server is legitimately warming. 

1278 embed_busy_deadline_s=( 

1279 cold_load_timeout_s(launch.weights_bytes, role) 

1280 if role is WorkerRole.EMBED 

1281 else None 

1282 ), 

1283 ) 

1284 for launch in role_launches 

1285 ] 

1286 if role is WorkerRole.CHAT: 

1287 self._chat_slots = role_launches[0].slots 

1288 self._chat_ctx = role_launches[0].ctx 

1289 # Every serving process passes through adoption (fresh launch, 

1290 # reload, guest bind), so the warning fires in each of them. 

1291 planning.warn_when_chat_downsized(role_launches[0]) 

1292 if role is WorkerRole.EMBED: 

1293 planning.warn_when_embed_window_below_chunk(role_launches[0]) 

1294 self._retire_clients(old_clients) 

1295 

1296 def _swap_for(self, role: WorkerRole) -> SwapManager | None: 

1297 """The swap process serving *role*, or None when the role has no server. 

1298 

1299 Caller holds ``self._lock``. 

1300 """ 

1301 group = self._role_group.get(role) 

1302 return None if group is None else self._swaps.get(group) 

1303 

1304 def _role_launches(self, role: WorkerRole) -> tuple[InstanceLaunch, ...]: 

1305 """*role*'s launch snapshot, empty when it has no server. 

1306 

1307 A co-tenant group holds more than one role's launches, so the group's 

1308 snapshot is filtered down to this role's replicas. 

1309 """ 

1310 group = self._role_group.get(role) 

1311 if group is None: 

1312 return () 

1313 return tuple(launch for launch in self._launches.get(group, ()) if launch.role is role) 

1314 

1315 def _drop_group(self, group: SwapGroup) -> SwapManager | None: 

1316 """Forget *group*'s swap/launches and every member role's clients. 

1317 

1318 Caller holds ``self._lock``. Member clients are retired (closed at a later 

1319 reload or shutdown, never while a reader could still hold one) and the chat 

1320 capacity falls back to its defaults when chat's group is dropped. 

1321 """ 

1322 swap = self._swaps.pop(group, None) 

1323 self._launches.pop(group, None) 

1324 # Prune the dir map with the group: a stale entry outliving its group makes 

1325 # _reload_dir see two dirs and pick one arbitrarily, splitting the provider. 

1326 self._group_dirs.pop(group, None) 

1327 for role in [r for r, g in self._role_group.items() if g is group]: 

1328 del self._role_group[role] 

1329 self._retire_clients(self._clients.pop(role, [])) 

1330 if role is WorkerRole.CHAT: 

1331 self._chat_slots = 1 

1332 self._chat_ctx = None 

1333 return swap 

1334 

1335 def _retire_clients(self, old_clients: list[LlamaServerClient]) -> None: 

1336 """Close the previously-retired clients, then retire *old_clients*. 

1337 

1338 Caller holds ``self._lock``. Retired clients are never handed to new 

1339 readers (they are out of ``self._clients``), so by this reload any reader 

1340 that held one from a prior reload has finished; an ``in_flight == 0`` 

1341 check confirms it before close, and any still-busy client stays retired 

1342 for the next reload. This closes idle reloaded-away pools without ever 

1343 closing one a reader could still use. Shutdown closes whatever remains. 

1344 """ 

1345 still_busy: list[LlamaServerClient] = [] 

1346 for client in self._retiring_clients: 

1347 if client.in_flight == 0: 

1348 client.close() 

1349 else: 

1350 still_busy.append(client) 

1351 self._retiring_clients = still_busy + old_clients 

1352 

1353 def _require_clients(self, role: WorkerRole) -> list[LlamaServerClient]: 

1354 """The client pool for *role*, or a user-facing error when it has no server. 

1355 

1356 A configured, placeable role gets one or more replica clients; their absence 

1357 means the role is unconfigured or did not fit memory. llama-swap loads each 

1358 upstream on its first request, so a returned client may still be cold. No 

1359 in-process fallback, so a missing pool is a hard error. 

1360 

1361 When the pool is empty but a swap was previously built and its process has 

1362 since exited (detected via ``is_live()``), a one-shot rebuild is attempted 

1363 before raising so a transient llama-swap restart recovers transparently. 

1364 """ 

1365 self._ensure_fleet() 

1366 with self._lock: 

1367 clients = self._clients.get(role) 

1368 swap = self._swap_for(role) 

1369 if not clients and swap is not None and not swap.is_live(): 

1370 self._rebuild_role(role) 

1371 with self._lock: 

1372 clients = self._clients.get(role) 

1373 if not clients: 

1374 # A refused launch's recorded reason wins over the generic line. 

1375 # BAD_REQUEST so the HTTP surfaces return this deterministic, 

1376 # actionable refusal in a 4xx body instead of a generic 500. 

1377 reason = self._skipped_unusable_ctx.get(role) 

1378 if reason is not None: 

1379 raise ProviderError( 

1380 reason, provider=_PROVIDER_NAME, kind=ProviderErrorKind.BAD_REQUEST 

1381 ) 

1382 raise ProviderError(_no_server_message(role), provider=_PROVIDER_NAME) 

1383 return list(clients) 

1384 

1385 def _with_rediscover(self, call: Callable[[], _T], *, role: WorkerRole | None = None) -> _T: 

1386 """Run *call*; on a connection-kind or load-capacity failure, retry once. 

1387 

1388 A vanished engine (its last user left on a config change, or it died) 

1389 surfaces as ProviderErrorKind.CONNECTION, or as a raw httpx transport 

1390 error when the proxy itself is gone (nothing listening to answer with 

1391 a status). Membership is still held, so dropping the swap refs and 

1392 retrying sends the call through _ensure_fleet, which rediscovers the 

1393 new proxy ports or rebuilds. One retry only; a second failure surfaces 

1394 to the caller. 

1395 

1396 A ProviderErrorKind.CAPACITY failure is the engine dying on load because 

1397 the estimate was too optimistic. Retrying it unchanged respawns the same 

1398 launch into the same death, so *role*'s auto context steps down first and 

1399 the role is rebuilt against the smaller plan. When there is no step left 

1400 to take (a user-pinned context, or already at the floor) the failure 

1401 surfaces instead: a retry that asks for the same thing is a crash loop. 

1402 """ 

1403 try: 

1404 return call() 

1405 except (ProviderError, httpx.TransportError) as err: 

1406 if is_rebuildable_failure(err) and role is not None: 

1407 return self._retry_rebuilt(call, role, err) 

1408 if not is_connection_failure(err): 

1409 raise 

1410 log.info("Engine unreachable; rediscovering before one retry") 

1411 self._drop_swap_refs() 

1412 self._release_holds() 

1413 return call() 

1414 

1415 def _retry_rebuilt(self, call: Callable[[], _T], role: WorkerRole, err: BaseException) -> _T: 

1416 """Rebuild *role* so the retry is a different launch, and run *call* again. 

1417 

1418 A held port just needs the rebuild, which picks a new one. A memory 

1419 shortfall needs the plan to come back smaller too, so the context steps 

1420 down first; when there is no step left to take, *err* is re-raised 

1421 untouched rather than rebuilding into the same death. 

1422 """ 

1423 from lilbee.providers.fleet.planning import record_ctx_downshift 

1424 

1425 if is_load_capacity_failure(err): 

1426 if not record_ctx_downshift(role): 

1427 log.warning( 

1428 "%s ran out of device memory on load and its context cannot be " 

1429 "reduced further; lower num_ctx or use a smaller model", 

1430 role.value, 

1431 ) 

1432 raise err 

1433 log.warning( 

1434 "%s ran out of device memory on load; re-planning it with a smaller " 

1435 "context where its window has room to give", 

1436 role.value, 

1437 ) 

1438 else: 

1439 log.warning("%s could not claim its port; rebuilding it on a new one", role.value) 

1440 self._rebuild_role(role) 

1441 return call() 

1442 

1443 def _rebuild_role(self, role: WorkerRole) -> None: 

1444 """Restart just *role*'s dead group (new port) from a fresh plan. 

1445 

1446 Other roles' groups keep serving; only the dead group is torn down and 

1447 respawned. Runs the same diff-driven pass as a reload, forcing *role* 

1448 into the restart set so an unchanged plan still replaces its dead swap. 

1449 """ 

1450 self._reload_pass(force=frozenset((role,))) 

1451 

1452 def role_ready(self, role: WorkerRole) -> bool: 

1453 """Whether *role*'s upstream is loaded and ready, without starting the swap. 

1454 

1455 A read-only probe for surfaces (HTTP status, SSE warming event) that want 

1456 to report cold-start state without triggering a load. False before the swap 

1457 is up or while the role's upstream is still loading. 

1458 """ 

1459 with self._lock: 

1460 swap = self._swap_for(role) 

1461 return swap is not None and swap.role_ready(role) 

1462 

1463 def max_concurrent_chats(self) -> int: 

1464 """The chat server's batching-slot capacity, so the gate admits that many. 

1465 

1466 Falls back to ``1`` before the chat group is up, so chat is serialized 

1467 until the slot count is known (the launcher warms the engine before a 

1468 client connects, so the real capacity is in effect by the first chat). 

1469 """ 

1470 with self._lock: 

1471 if WorkerRole.CHAT not in self._role_group: 

1472 return 1 

1473 return self._chat_slots 

1474 

1475 def served_chat_ctx(self) -> int | None: 

1476 """Per-slot context the chat server runs with, or None if not up.""" 

1477 with self._lock: 

1478 return self._chat_ctx if WorkerRole.CHAT in self._role_group else None 

1479 

1480 def served_chat_slots(self) -> int | None: 

1481 """Batching slots the chat server runs with, or None if not up.""" 

1482 with self._lock: 

1483 return self._chat_slots if WorkerRole.CHAT in self._role_group else None 

1484 

1485 def embed_token_cap(self) -> int | None: 

1486 """The served embed launch's token cap, or the planned one before it is up.""" 

1487 with self._lock: 

1488 launches = self._role_launches(WorkerRole.EMBED) 

1489 if launches: 

1490 return launches[0].token_cap 

1491 return planning.planned_embed_token_cap(cfg.embedding_model) 

1492 

1493 def health_warnings(self) -> list[HealthWarning]: 

1494 """Serving degradations: the embed window and any placement divergences.""" 

1495 warnings: list[HealthWarning] = [] 

1496 cap = self.embed_token_cap() 

1497 if cap is not None: 

1498 warning = engine_params.embed_window_warning(cap) 

1499 if warning is not None: 

1500 warnings.append(warning) 

1501 with self._lock: 

1502 swaps = list(self._swaps.values()) 

1503 for swap in swaps: 

1504 warnings.extend(swap.health_warnings()) 

1505 return warnings 

1506 

1507 def chat_prefill_progress(self) -> tuple[int, int] | None: 

1508 """``(processed, total)`` of a chat prefill in flight, or None when idle.""" 

1509 with self._lock: 

1510 return self._chat_prefill 

1511 

1512 def _record_chat_prefill(self, progress: tuple[int, int] | None) -> None: 

1513 """Store a stream's prefill reading. Latest writer wins across concurrent 

1514 streams; a live prefill rewrites itself on its next engine batch.""" 

1515 with self._lock: 

1516 self._chat_prefill = progress 

1517 

1518 def warm_pending(self) -> bool: 

1519 """Whether a requested warm is still running. 

1520 

1521 The tracker only stamps a phase once the chat role starts loading, which is 

1522 seconds after the swap is spawned, so ``warm_progress`` alone cannot tell a 

1523 not-yet-started warm from no warm at all. 

1524 """ 

1525 with self._lock: 

1526 return self._warming 

1527 

1528 def warm_progress(self) -> WarmProgress | None: 

1529 """Live cold-load progress for the chat role, or None before warm begins.""" 

1530 return self._warm_tracker.snapshot() 

1531 

1532 def _shutdown_swap(self, *, latch: bool = True) -> None: 

1533 """Release this process's engine use; ``latch=False`` keeps the provider reusable. 

1534 

1535 Terminal ``shutdown()`` latches ``_shut_down`` so a discarded provider's 

1536 in-flight warm/reload thread can't spawn an orphan swap, then releases 

1537 membership: the engine stops only when this was the last user and 

1538 persistence was not opted into. The cache-drop paths 

1539 (``invalidate_load_cache``, ``drop_loaded_models_async``) pass 

1540 ``latch=False``: a config change restarts the shared engine for every 

1541 user (they rediscover), and this provider rebuilds on next use. 

1542 """ 

1543 # Latched before the lock, not inside it: every _shut_down check runs after 

1544 # acquiring the build lock, so a warm or reload thread queued behind us can 

1545 # only bail early if the flag is already set when its turn comes. 

1546 if latch: 

1547 with self._lock: 

1548 self._shut_down = True 

1549 # The build lock serializes shutdown against a concurrent reload/build: 

1550 # both mutate self._swaps and the llama-swap processes, so an unserialized 

1551 # loser would overwrite the winner's state and leak a live llama-swap. 

1552 # Bounded, because a wedged engine start holds this lock and an unbounded 

1553 # wait would hang process exit outright. On timeout the teardown proceeds 

1554 # anyway: whatever the builder leaves behind is recorded in the engine 

1555 # dir's state files, so the next start's reap finds it by record, while a 

1556 # shutdown that never returns cannot be recovered from at all. 

1557 acquired = self._build_lock.acquire(timeout=_SHUTDOWN_BUILD_LOCK_WAIT_S) 

1558 if not acquired: 

1559 log.warning( 

1560 "Engine build still in progress after %.0fs; shutting down without " 

1561 "waiting for it. Leftovers are reaped from their records on the next start.", 

1562 _SHUTDOWN_BUILD_LOCK_WAIT_S, 

1563 ) 

1564 try: 

1565 # Terminal shutdown closes every client; a config-change teardown 

1566 # retires them so an in-flight reader is never severed. 

1567 self._drop_swap_refs(close_all=latch) 

1568 self._release_engines(config_changed=not latch) 

1569 finally: 

1570 if acquired: 

1571 self._build_lock.release() 

1572 

1573 def _release_engines(self, *, config_changed: bool = False) -> None: 

1574 """Drop membership in every used engine dir; stop each engine we leave last. 

1575 

1576 Runs under each dir's cross-process build lock so a departing last user 

1577 can never race an arriving binder: the arrival either sees the engine 

1578 (and its bind holds it live) or sees the slot empty and builds. 

1579 

1580 Whether the engine outlives us is the union of every user's opt-in, not 

1581 just the exiting process's config: the machine slot is shared by 

1582 installations that configure it differently, and which one leaves last 

1583 is arbitrary. 

1584 

1585 *config_changed* is the cache-drop path, where this provider's settings 

1586 or model changed. That makes the running engine stale for us, so no 

1587 persistence opt-in preserves it -- but it says nothing about the peers 

1588 still serving requests against it, so a shared engine is left running 

1589 and the next use re-runs the ladder, binding it if it happens to match 

1590 and overflowing to a private dir if it does not. 

1591 

1592 The hold map is cleared either way. Leaving a stale hold behind is not 

1593 benign: after a lazy rebuild overflows to a private dir (a foreign 

1594 process having claimed the machine slot in the gap), the next release 

1595 would iterate the stale machine hold and stop that foreign engine 

1596 mid-use, and the stale flock would keep live_users_exist true so the 

1597 foreign engine's real last user could never reap it. 

1598 """ 

1599 from lilbee.core.config import cfg 

1600 

1601 for engine_dir, hold in list(self._engine_holds.items()): 

1602 with build_lock(engine_dir, best_effort=True): 

1603 last = hold.release_and_check_last() 

1604 # A flip after binding never re-acquires, so reconcile our own mark 

1605 # here. Skipped on a config change, which stops and clears regardless. 

1606 if not config_changed: 

1607 if cfg.keep_engine_warm: 

1608 request_keep_warm(engine_dir, cfg.data_root) 

1609 else: 

1610 withdraw_keep_warm(engine_dir, cfg.data_root) 

1611 # Any remaining opt-in keeps the engine, including a peer's. 

1612 keep = not config_changed and keep_warm_requested(engine_dir) 

1613 if last and not keep: 

1614 stop_engine(engine_dir) 

1615 log.info("Engine stopped at %s (last user out)", engine_dir) 

1616 elif last: 

1617 log.info("Engine left warm at %s (last user out)", engine_dir) 

1618 elif config_changed: 

1619 log.info("Engine left running at %s (still in use by peers)", engine_dir) 

1620 self._engine_holds = {} 

1621 

1622 def _release_holds(self) -> None: 

1623 """Drop this process's engine memberships without stopping any engine. 

1624 

1625 The rediscover retry re-runs the acquisition ladder. A retained membership 

1626 would make this process count itself as a live user of the machine slot, so 

1627 the ladder would judge the slot in use and overflow to a private engine 

1628 instead of rebinding a recovered engine or rebuilding a dead one -- N 

1629 private engines and N times the VRAM after the shared engine first dies. 

1630 Nothing is stopped here: a live engine is rebound by the retry, and a dead 

1631 one is cleared by the rebuild that retry triggers. 

1632 """ 

1633 for hold in list(self._engine_holds.values()): 

1634 hold.release_and_check_last() 

1635 self._engine_holds = {} 

1636 

1637 def _drop_swap_refs(self, *, close_all: bool = False) -> None: 

1638 """Clear every group's swap/clients and the chat capacity so the next call rebuilds. 

1639 

1640 Live pools are RETIRED through the ``in_flight`` check rather than closed 

1641 outright: ``_with_rediscover`` reaches this on any connection blip, so a 

1642 chat proxy hiccup must not sever the client another thread is mid-embed or 

1643 mid-stream on (a streamed response is handed out after the retry returns, 

1644 and failures past the first frame are not retried). Idle clients close now, 

1645 busy ones stay retired for a later pass. 

1646 

1647 *close_all* is the terminal-shutdown path, where nothing will read again and 

1648 whatever remains must actually be closed. 

1649 """ 

1650 doomed: list[LlamaServerClient] = [] 

1651 with self._lock: 

1652 live = [client for pool in self._clients.values() for client in pool] 

1653 self._swaps = {} 

1654 self._launches = {} 

1655 self._role_group = {} 

1656 self._group_dirs = {} 

1657 self._clients = {} 

1658 if close_all: 

1659 doomed = live + self._retiring_clients 

1660 self._retiring_clients = [] 

1661 else: 

1662 self._retire_clients(live) 

1663 self._chat_slots = 1 

1664 self._chat_ctx = None 

1665 # A torn-down fleet's load failures describe servers that no longer 

1666 # exist; the next warm records its own. 

1667 self._warm_errors = {} 

1668 # Full teardown: the next build starts from a clean box, so it must 

1669 # re-snapshot memory rather than plan against this boot's probe. 

1670 planning.clear_plan_probe() 

1671 for client in doomed: 

1672 client.close() 

1673 

1674 def _drop_dead_swaps(self) -> None: 

1675 """Drop the refs of groups whose process is gone so the next call rebuilds them. 

1676 

1677 A no-op for groups still running (e.g. the failure was in planning), so 

1678 a live engine is never abandoned unstopped. 

1679 """ 

1680 with self._build_lock, self._lock: 

1681 for group in [g for g, swap in self._swaps.items() if not swap.running]: 

1682 self._drop_group(group) 

1683 

1684 def _require_configured_model( 

1685 self, model: str | None, configured: str, role: WorkerRole 

1686 ) -> None: 

1687 """Reject a per-call model that differs from the server's configured one. 

1688 

1689 The fleet serves the configured model for each role; switching models is 

1690 a config change that respawns the server (via ``invalidate_load_cache``), 

1691 not a per-call override. An empty/None ``model`` means "use the configured 

1692 one" and is always accepted. 

1693 """ 

1694 if model and model != configured: 

1695 raise ProviderError( 

1696 configured_model_message(role, configured, model), 

1697 provider=_PROVIDER_NAME, 

1698 kind=ProviderErrorKind.BAD_REQUEST, 

1699 ) 

1700 

1701 @overload 

1702 def chat( 

1703 self, 

1704 messages: list[ChatMessage], 

1705 *, 

1706 stream: Literal[False] = False, 

1707 options: dict[str, Any] | None = None, 

1708 model: str | None = None, 

1709 tools: list[dict[str, Any]] | None = None, 

1710 tool_choice: str | dict[str, Any] | None = None, 

1711 ) -> ChatResult: ... 

1712 

1713 @overload 

1714 def chat( 

1715 self, 

1716 messages: list[ChatMessage], 

1717 *, 

1718 stream: Literal[True], 

1719 options: dict[str, Any] | None = None, 

1720 model: str | None = None, 

1721 tools: list[dict[str, Any]] | None = None, 

1722 tool_choice: str | dict[str, Any] | None = None, 

1723 ) -> ClosableIterator[ChatStreamItem]: ... 

1724 

1725 def chat( 

1726 self, 

1727 messages: list[ChatMessage], 

1728 *, 

1729 stream: bool = False, 

1730 options: dict[str, Any] | None = None, 

1731 model: str | None = None, 

1732 tools: list[dict[str, Any]] | None = None, 

1733 tool_choice: str | dict[str, Any] | None = None, 

1734 ) -> ChatResult | ClosableIterator[ChatStreamItem]: 

1735 """Route a chat turn to the least-busy chat server. 

1736 

1737 Non-streaming returns a :class:`ChatResult` (text, tool calls, finish 

1738 reason); streaming yields :data:`ChatStreamItem` frames. ``--jinja`` on 

1739 the server parses native tool calls, so tool support needs no per-family 

1740 parser here. 

1741 """ 

1742 from lilbee.providers.engine_params import chat_options_to_kwargs 

1743 

1744 self._require_configured_model(model, str(cfg.chat_model), WorkerRole.CHAT) 

1745 with self._lazy_warm_scope(): 

1746 self._require_clients(WorkerRole.CHAT) 

1747 messages = self._fit_chat_context( 

1748 messages, tools, options, model or str(cfg.chat_model) 

1749 ) 

1750 # Translate options exactly as the in-process path did (validate via 

1751 # LLMOptions, num_predict -> max_tokens, drop num_ctx) so the server 

1752 # honors the same generation settings; a raw passthrough would drop 

1753 # num_predict and leak the load-only num_ctx. 

1754 server_options = chat_options_to_kwargs(options) or None 

1755 if stream: 

1756 # The first frame is pulled eagerly so a dead proxy fails inside 

1757 # _with_rediscover; a failure past the first frame surfaces to the 

1758 # caller as a retry error, and rediscovery covers the next call. 

1759 return self._with_rediscover( 

1760 lambda: _primed_stream( 

1761 _least_in_flight(self._require_clients(WorkerRole.CHAT)).chat_stream_items( 

1762 messages, tools=tools, tool_choice=tool_choice, options=server_options 

1763 ) 

1764 ), 

1765 role=WorkerRole.CHAT, 

1766 ) 

1767 return self._with_rediscover( 

1768 lambda: _least_in_flight(self._require_clients(WorkerRole.CHAT)).chat_result( 

1769 messages, tools=tools, tool_choice=tool_choice, options=server_options 

1770 ), 

1771 role=WorkerRole.CHAT, 

1772 ) 

1773 

1774 def chat_with_tools( 

1775 self, 

1776 messages: list[ChatMessage], 

1777 *, 

1778 tools: list[dict[str, Any]], 

1779 tool_choice: str | dict[str, Any] | None = None, 

1780 options: dict[str, Any] | None = None, 

1781 model: str | None = None, 

1782 ) -> ChatToolResult: 

1783 """Route a tool-enabled chat turn to the least-busy chat server.""" 

1784 from lilbee.providers.engine_params import chat_options_to_kwargs 

1785 

1786 self._require_configured_model(model, str(cfg.chat_model), WorkerRole.CHAT) 

1787 with self._lazy_warm_scope(): 

1788 self._require_clients(WorkerRole.CHAT) 

1789 messages = self._fit_chat_context( 

1790 messages, tools, options, model or str(cfg.chat_model) 

1791 ) 

1792 server_options = chat_options_to_kwargs(options) or None 

1793 return self._with_rediscover( 

1794 lambda: _least_in_flight(self._require_clients(WorkerRole.CHAT)).chat_tools( 

1795 messages, tools=tools, tool_choice=tool_choice, options=server_options 

1796 ), 

1797 role=WorkerRole.CHAT, 

1798 ) 

1799 

1800 def _fit_chat_context( 

1801 self, 

1802 messages: list[ChatMessage], 

1803 tools: list[dict[str, Any]] | None, 

1804 options: dict[str, Any] | None, 

1805 model: str, 

1806 ) -> list[ChatMessage]: 

1807 """Drop oldest turns so the prompt fits the served context. 

1808 

1809 A ``num_predict`` reservation larger than the default generation room is 

1810 capped to it, so an agent client that over-reserves keeps its history 

1811 instead of having it evicted; llama-server stops at the context edge 

1812 anyway. A smaller reservation is honored as-is and widens the prompt. 

1813 Raises ``ProviderError(CONTEXT_OVERFLOW)`` only when system messages, 

1814 tools, and the final turn exceed the window even with the capped 

1815 reserve (mapped to a 400 by the chat-completions route). 

1816 """ 

1817 # 0/None means the served context is unknown (no chat launch adopted yet); 

1818 # a real per-slot context is always positive, so skip windowing. 

1819 if not self._chat_ctx: 

1820 return messages 

1821 # An output reservation only ever buys the prompt MORE room, never less: 

1822 # a num_predict past the default is a ceiling on what the model may 

1823 # generate, not a claim on prompt space, and llama-server stops at the 

1824 # context edge regardless. Capping it here rather than retrying after a 

1825 # failed fit is the difference between a policy and a rescue -- an agent 

1826 # reserving most of the window leaves a budget of a few dozen tokens, in 

1827 # which the final turn still "fits" while the whole conversation is 

1828 # silently evicted. 

1829 requested = (options or {}).get("num_predict") 

1830 reserve = min(requested, GENERATION_RESERVE_TOKENS) if requested else None 

1831 budget = prompt_token_budget(self._chat_ctx, reserve) 

1832 result = window_messages(messages, tools, budget) 

1833 if not result.fits: 

1834 raise ProviderError( 

1835 f"Prompt of about {result.prompt_tokens} tokens exceeds the " 

1836 f"{budget}-token prompt budget for {model!r} " 

1837 f"({self._chat_ctx}-token window minus reserve and margin). " 

1838 "Shorten the conversation or the system prompt.", 

1839 provider=_PROVIDER_NAME, 

1840 kind=ProviderErrorKind.CONTEXT_OVERFLOW, 

1841 ) 

1842 return result.messages 

1843 

1844 def embed(self, texts: list[str]) -> list[Vector]: 

1845 return self._with_rediscover(lambda: self._embed_once(texts), role=WorkerRole.EMBED) 

1846 

1847 def _embed_once(self, texts: list[str]) -> list[Vector]: 

1848 clients = self._require_clients(WorkerRole.EMBED) 

1849 return _call_with_failover(clients, lambda client: client.embed(texts)) 

1850 

1851 def count_tokens(self, text: str) -> int: 

1852 """Exact token count of *text* under the embedding model's tokenizer. 

1853 

1854 Routes to the embed server's ``/tokenize`` so chunk sizing counts the same 

1855 tokens the embedder consumes, with the same rediscovery retry as embedding. 

1856 """ 

1857 return self._with_rediscover(lambda: self._count_once(text), role=WorkerRole.EMBED) 

1858 

1859 def _count_once(self, text: str) -> int: 

1860 clients = self._require_clients(WorkerRole.EMBED) 

1861 return _call_with_failover(clients, lambda client: client.count_tokens(text)) 

1862 

1863 def count_chat_prompt_tokens( 

1864 self, 

1865 messages: list[ChatMessage], 

1866 *, 

1867 options: dict[str, Any] | None = None, 

1868 model: str | None = None, 

1869 tools: list[dict[str, Any]] | None = None, 

1870 tool_choice: str | dict[str, Any] | None = None, 

1871 ) -> int: 

1872 """Tokens the chat server prefills for this prompt, template applied. 

1873 

1874 Takes the arguments :meth:`chat` takes and sends the same body. The 

1875 conversation is not fitted to the served window first, so the count 

1876 covers the whole prompt the caller asked about. 

1877 """ 

1878 from lilbee.providers.engine_params import chat_options_to_kwargs 

1879 

1880 self._require_configured_model(model, str(cfg.chat_model), WorkerRole.CHAT) 

1881 server_options = chat_options_to_kwargs(options) or None 

1882 with self._lazy_warm_scope(): 

1883 return self._with_rediscover( 

1884 lambda: _least_in_flight( 

1885 self._require_clients(WorkerRole.CHAT) 

1886 ).count_chat_prompt_tokens( 

1887 messages, tools=tools, tool_choice=tool_choice, options=server_options 

1888 ), 

1889 role=WorkerRole.CHAT, 

1890 ) 

1891 

1892 def vision_ocr( 

1893 self, png_bytes: bytes, model: str, prompt: str = "", *, timeout: float | None = None 

1894 ) -> str: 

1895 from lilbee.vision import build_vision_messages, resolve_ocr_prompt 

1896 

1897 self._require_configured_model(model, str(cfg.vision_model), WorkerRole.VISION) 

1898 pool = self._vision_pool() 

1899 effective = model or str(cfg.vision_model) 

1900 messages = build_vision_messages(prompt or resolve_ocr_prompt(effective), png_bytes) 

1901 try: 

1902 return _ocr_dispatch(pool, messages, _ocr_deadline(timeout)) 

1903 except _PageBudgetExhausted: 

1904 raise ProviderError( 

1905 "Vision OCR timed out waiting for a free vision slot.", 

1906 provider=_PROVIDER_NAME, 

1907 ) from None 

1908 

1909 def vision_slot_capacity(self) -> int | None: 

1910 """Total fitted ``--parallel`` slots across the running vision replicas. 

1911 

1912 ``None`` before the fleet is up (no launch snapshot yet), so the ingest 

1913 fan-out keeps its own estimate until real capacity is known. A modest 

1914 card that fit fewer slots than requested reports the smaller real number, 

1915 so the fan-out never queues more pages than the servers can serve. 

1916 """ 

1917 launches = self._role_launches(WorkerRole.VISION) 

1918 if not launches: 

1919 return None 

1920 return max(1, sum(launch.slots for launch in launches)) 

1921 

1922 def _vision_pool(self) -> list[_VisionReplica]: 

1923 """Each vision replica paired with its fitted ``--parallel`` slot count. 

1924 

1925 The fitted count can be lower than ``vision_ocr_concurrency`` when memory 

1926 forced a smaller fit; dispatching at the configured ceiling instead 

1927 over-subscribes that server. The configured ceiling applies per replica 

1928 only when no matching launch snapshot exists (a reload can momentarily 

1929 drop it between two reads). 

1930 """ 

1931 self._serve_vision_on_request() 

1932 clients = self._require_clients(WorkerRole.VISION) 

1933 launches = self._role_launches(WorkerRole.VISION) 

1934 if launches and len(launches) == len(clients): 

1935 return [ 

1936 _VisionReplica(client, max(1, launch.slots)) 

1937 for client, launch in zip(clients, launches, strict=True) 

1938 ] 

1939 fallback_slots = max(1, cfg.vision_ocr_concurrency) 

1940 return [_VisionReplica(client, fallback_slots) for client in clients] 

1941 

1942 def _serve_vision_on_request(self) -> None: 

1943 """Re-plan the fleet with vision for an OCR call that ``enable_ocr`` false left out. 

1944 

1945 Only a request that turned OCR on reaches vision OCR while the setting is 

1946 off, so the call itself is the demand. The re-plan is the diff-driven pass 

1947 a vision model change runs: it restarts the groups whose launches change, 

1948 which includes chat where chat and vision share one group. A failed 

1949 re-plan drops the grant, so the next page tries again. 

1950 """ 

1951 ref = str(cfg.vision_model) 

1952 if not ref: 

1953 return 

1954 with self._vision_request_lock: 

1955 if planning.vision_role_wanted(ref): 

1956 return 

1957 planning.grant_vision_on_request() 

1958 try: 

1959 self._dispatch_reload("fleet-vision-on-request", wait=True) 

1960 except BaseException: 

1961 planning.revoke_vision_on_request() 

1962 raise 

1963 

1964 # PDF/image OCR now runs inside xberg via the registered lilbee-vision 

1965 # backend (see data.extract.backends.vision_ocr); this provider only exposes 

1966 # single-image vision_ocr, which that backend calls. 

1967 

1968 def rerank(self, query: str, candidates: list[str]) -> list[float]: 

1969 return self._with_rediscover( 

1970 lambda: self._rerank_once(query, candidates), role=WorkerRole.RERANK 

1971 ) 

1972 

1973 def _rerank_once(self, query: str, candidates: list[str]) -> list[float]: 

1974 clients = self._require_clients(WorkerRole.RERANK) 

1975 return _call_with_failover(clients, lambda client: client.rerank(query, candidates)) 

1976 

1977 # --- model management: registry / GGUF reads, no running server needed --- 

1978 

1979 def supports_rerank(self) -> bool: 

1980 """Serve a cross-encoder (rank pooling) or an LLM reranker (yes/no logprob).""" 

1981 return True 

1982 

1983 def list_models(self) -> list[str]: 

1984 """List installed models from the registry.""" 

1985 from lilbee.app.services import get_services 

1986 

1987 registry = get_services().registry 

1988 return sorted(m.ref for m in registry.list_installed()) 

1989 

1990 def list_chat_models(self, provider: str) -> list[str]: 

1991 """The local engine has no frontier-provider catalog; always ``[]``.""" 

1992 del provider 

1993 return [] 

1994 

1995 def pull_model(self, model: str, *, on_progress: Callable[..., Any] | None = None) -> None: 

1996 """Not supported directly: ``lilbee.catalog`` handles GGUF downloads.""" 

1997 del on_progress 

1998 raise NotImplementedError( 

1999 f"The local engine cannot pull model {model!r}. " 

2000 "Download GGUF files through the catalog or 'lilbee model pull'." 

2001 ) 

2002 

2003 def show_model(self, model: str) -> dict[str, Any] | None: 

2004 """Return model metadata from GGUF headers, or ``None`` if unresolved.""" 

2005 from lilbee.providers.engine_params import resolve_model_path 

2006 from lilbee.providers.gguf_meta import read_gguf_metadata 

2007 

2008 try: 

2009 path = resolve_model_path(model) 

2010 except ProviderError: 

2011 return None 

2012 return read_gguf_metadata(path) 

2013 

2014 def get_capabilities(self, model: str) -> list[str]: 

2015 """Detect capabilities from the local GGUF files. 

2016 

2017 Cross-encoder rerank GGUFs report ``["rerank"]`` (they cannot generate); 

2018 other models report ``"completion"`` plus ``"vision"`` when an mmproj 

2019 sidecar is present. 

2020 """ 

2021 from lilbee.catalog import is_rerank_ref 

2022 from lilbee.providers.engine_params import resolve_model_path 

2023 from lilbee.providers.gguf_meta import find_mmproj_for_model 

2024 

2025 if model and is_rerank_ref(model): 

2026 return ["rerank"] 

2027 caps = ["completion"] 

2028 try: 

2029 path = resolve_model_path(model) 

2030 except ProviderError: 

2031 return caps 

2032 try: 

2033 find_mmproj_for_model(path) 

2034 caps.append("vision") 

2035 except ProviderError: 

2036 pass 

2037 return caps 

2038 

2039 def supports_tools(self, model_ref: str) -> bool: 

2040 """True iff *model_ref*'s GGUF chat template references tool tokens. 

2041 

2042 The server parses native tool calls via ``--jinja``; a template that 

2043 declares tools is the signal that the model was trained to emit them. 

2044 Cached on ``(path, mtime)`` so a tool-bearing chat doesn't re-read the 

2045 GGUF header each request; a re-quantised file at the same path 

2046 invalidates because its mtime changes. 

2047 """ 

2048 from lilbee.providers.engine_params import resolve_model_path 

2049 

2050 try: 

2051 path = resolve_model_path(model_ref) 

2052 except (ProviderError, OSError): 

2053 log.debug("supports_tools: resolve_model_path failed for %s", model_ref, exc_info=True) 

2054 return False 

2055 try: 

2056 mtime_ns = path.stat().st_mtime_ns 

2057 except OSError: 

2058 mtime_ns = 0 

2059 return _supports_tools_cached(str(path), mtime_ns) 

2060 

2061 def warm_up_pool(self) -> None: 

2062 """Pre-load every configured role off the caller's thread (idempotent). 

2063 

2064 Starting the swap and loading each role's model (seconds on a cold large 

2065 model) runs on a background thread and this returns at once: the eager-start 

2066 at TUI mount must not freeze the UI. The spawn listeners fire per role as it 

2067 loads, so the UI shows progress. A second call while warm-up is in flight 

2068 (or once the fleet is up) is a no-op. 

2069 """ 

2070 with self._lock: 

2071 if self._warming: 

2072 return 

2073 fleet_up = bool(self._swaps) 

2074 # A live swap whose model llama-swap idle-unloaded (its ttl stops only the 

2075 # llama-server child, leaving the swap handle in _swaps) reports its role 

2076 # cold. Re-warm so a prompt sent into that gap drives llama-swap's 

2077 # on-demand reload; bailing on "swaps exist" alone stranded every later 

2078 # prompt on a stale not-ready. A fully-loaded fleet still short-circuits. 

2079 # The probe runs off the lock (role_ready may hit the proxy). 

2080 if fleet_up and self._roles_ready(): 

2081 return 

2082 with self._lock: 

2083 if self._warming: 

2084 return 

2085 self._warming = True 

2086 threading.Thread( 

2087 target=self._warm_up_blocking, 

2088 name="fleet-warm-up", 

2089 daemon=True, 

2090 ).start() 

2091 

2092 def _roles_ready(self) -> bool: 

2093 """Whether every configured role's upstream is loaded (fleet fully warm).""" 

2094 with self._lock: 

2095 roles = list(self._role_group) 

2096 return bool(roles) and all(self.role_ready(role) for role in roles) 

2097 

2098 def _warm_up_blocking(self) -> None: 

2099 """Start the fleet and pre-load every role on a background thread. 

2100 

2101 Runs on a daemon thread with no caller to catch failures, so a startup 

2102 error is logged and swallowed: a role that can't load surfaces a 

2103 user-facing ProviderError on the next call, not a thread traceback. 

2104 

2105 The tracker is stamped STARTING before the fleet spawn so surfaces 

2106 show the engine coming up from the first moment (spawn plus health 

2107 check takes seconds and previously reported nothing), and stamped 

2108 ERROR with the real reason when the warm fails before the chat warm 

2109 proper begins. 

2110 """ 

2111 try: 

2112 self._warm_tracker.begin(str(cfg.chat_model)) 

2113 self._ensure_fleet() 

2114 self._preload_roles() 

2115 self._finalize_warm_if_chat_never_ran() 

2116 except Exception as exc: 

2117 if isinstance(exc, RuntimeError) and sys.is_finalizing(): 

2118 # A fast CLI exit can tear down the interpreter while this daemon 

2119 # thread is still warming; the process is leaving, so drop it quietly. 

2120 log.debug("Engine warm-up abandoned during interpreter shutdown: %s", exc) 

2121 else: 

2122 # A warm-up failure is handled (roles lazy-load on first use), so 

2123 # keep the full traceback at debug: a WARNING carrying exc_info 

2124 # reads like a crash for a condition the next real call recovers 

2125 # from. 

2126 log.warning("Engine warm-up failed: %s", exc) 

2127 log.debug("Engine warm-up failure detail.", exc_info=True) 

2128 self._fail_warm_unless_ready(str(exc)) 

2129 finally: 

2130 with self._lock: 

2131 self._warming = False 

2132 

2133 def _fail_warm_unless_ready(self, message: str) -> None: 

2134 """Stamp the warm tracker ERROR unless the chat warm already finished. 

2135 

2136 A failure in a later role's preload must not clobber a chat warm that 

2137 reached READY; every earlier failure leaves the tracker mid-phase, 

2138 where surfaces would spin forever and the prompt path could not name 

2139 the reason. 

2140 """ 

2141 snapshot = self._warm_tracker.snapshot() 

2142 if snapshot is None or snapshot.phase is not WarmPhase.READY: 

2143 self._warm_tracker.fail(message) 

2144 

2145 def _finalize_warm_if_chat_never_ran(self) -> None: 

2146 """Terminate the early STARTING stamp when no chat instance was placed. 

2147 

2148 ``_warm_chat_role`` always ends in READY or ERROR when it runs, so a 

2149 snapshot still on STARTING after a successful preload means the plan had 

2150 no chat instance. A chat model that isn't installed, one whose launch the 

2151 plan refused for an unusable window, and one with no engine to run it all 

2152 fail the warm with a user-facing reason so the prompt path renders 

2153 "failed to load" instead of spinning a "not ready" retry that can never 

2154 succeed; any other reason (a remote-routed chat has no local server to 

2155 warm) clears the stamp. 

2156 """ 

2157 snapshot = self._warm_tracker.snapshot() 

2158 if snapshot is None or snapshot.phase is not WarmPhase.STARTING: 

2159 return 

2160 missing = self._skipped_not_installed.get(WorkerRole.CHAT) 

2161 if missing is not None: 

2162 self._warm_tracker.fail(f"chat model {clean_display_name(missing)} is not installed") 

2163 return 

2164 unusable = self._skipped_unusable_ctx.get(WorkerRole.CHAT) 

2165 if unusable is not None: 

2166 self._warm_tracker.fail(unusable) 

2167 return 

2168 if _chat_needs_local_engine() and (reason := _unusable_engine_reason()) is not None: 

2169 self._warm_tracker.fail(reason) 

2170 return 

2171 self._warm_tracker.clear() 

2172 

2173 def _preload_roles(self, roles: frozenset[WorkerRole] | None = None) -> None: 

2174 """Issue a cheap request per replica so llama-swap loads each upstream now. 

2175 

2176 llama-swap starts an upstream on its first request, so warming sends a 

2177 minimal call to every replica of every role (firing the spawn listeners 

2178 around each role). A per-replica failure is logged and skipped; that replica 

2179 still loads on its first real use. The chat role routes through 

2180 :meth:`_warm_chat_role` so a launcher gets granular progress. *roles* 

2181 narrows the warm to just those roles (a reload warms only what restarted). 

2182 

2183 Roles on separate devices warm concurrently: chat is the long pole (a 

2184 large model's load dominates), so the light roles load alongside it 

2185 instead of before it. Roles whose launches pin overlapping devices warm 

2186 one at a time instead, chat last: two engines loading into the same 

2187 card at once race each other for VRAM, and the loser's first load can 

2188 OOM even though both fit once settled. 

2189 """ 

2190 with self._lock: 

2191 pools = { 

2192 role: list(clients) 

2193 for role, clients in self._clients.items() 

2194 if roles is None or role in roles 

2195 } 

2196 on_spawning, on_spawned = self._on_spawning, self._on_spawned 

2197 device_sets = _role_device_sets( 

2198 launch for launches in self._launches.values() for launch in launches 

2199 ) 

2200 

2201 if not pools: 

2202 return 

2203 listeners = (on_spawning, on_spawned) 

2204 calls = [ 

2205 DaemonCall( 

2206 functools.partial(self._warm_chain, chain, pools, listeners), 

2207 name=f"fleet-preload-{'+'.join(role.value for role in chain)}", 

2208 ) 

2209 for chain in _warm_chains(list(pools), device_sets) 

2210 ] 

2211 # Every chain finishes before the first chain's error, in chain order, is raised. 

2212 for call in calls: 

2213 call.wait() 

2214 for call in calls: 

2215 call.result() 

2216 

2217 def _warm_chain( 

2218 self, 

2219 chain: list[WorkerRole], 

2220 pools: dict[WorkerRole, list[LlamaServerClient]], 

2221 listeners: tuple[Callable[[WorkerRole], None] | None, Callable[[WorkerRole], None] | None], 

2222 ) -> None: 

2223 """Warm *chain*'s roles one at a time; every role gets its attempt. 

2224 

2225 An unexpected error warming one role (a listener blowing up) must not rob 

2226 the roles behind it of their warm, so the first error is re-raised only 

2227 after the chain finishes. 

2228 """ 

2229 on_spawning, on_spawned = listeners 

2230 first_exc: Exception | None = None 

2231 for role in chain: 

2232 try: 

2233 if on_spawning is not None: 

2234 on_spawning(role) 

2235 if role is WorkerRole.CHAT: 

2236 self._warm_chat_role(pools[role]) 

2237 else: 

2238 self._warm_role_clients(role, pools[role]) 

2239 if on_spawned is not None: 

2240 on_spawned(role) 

2241 except Exception as exc: 

2242 first_exc = first_exc or exc 

2243 if first_exc is not None: 

2244 raise first_exc 

2245 

2246 def _warm_role_clients(self, role: WorkerRole, clients: list[LlamaServerClient]) -> bool: 

2247 """Warm every replica of *role*; return whether at least one loaded. 

2248 

2249 A replica that fails to load is reported at warning level with the engine's 

2250 own message (an unsupported architecture, a corrupt file). Warm-up stays 

2251 best-effort, but the failure must not be silent: the role then serves 

2252 nothing, and a caller that never reaches it would otherwise see only an 

2253 unexplained empty answer. 

2254 """ 

2255 warmed = False 

2256 self._warm_errors.pop(role, None) 

2257 for client in clients: 

2258 try: 

2259 _warm_role(role, client) 

2260 client.mark_healthy() 

2261 warmed = True 

2262 except Exception as exc: 

2263 # A replica that cannot load is not routable. Marking it takes it 

2264 # out of the pool so calls go to a sibling on a device that works, 

2265 # instead of every request picking the dead one again. It is a 

2266 # device fault as often as a model one: an adapter that enumerates 

2267 # but cannot allocate fails here and nowhere else. The health flag 

2268 # carries its own cool-down, so a device that recovers rejoins 

2269 # without anything having to remember it was bad. 

2270 client.mark_unhealthy() 

2271 self._warm_errors[role] = str(exc) 

2272 log.warning( 

2273 "The %s model failed to load: %s", 

2274 role.value, 

2275 exc, 

2276 exc_info=log.isEnabledFor(logging.DEBUG), 

2277 ) 

2278 if warmed: 

2279 self._warm_errors.pop(role, None) 

2280 return warmed 

2281 

2282 def _warm_chat_role(self, clients: list[LlamaServerClient]) -> None: 

2283 """Warm the chat role, driving the tracker through read -> load -> ready/fail. 

2284 

2285 Readiness is decided by whether a warm request actually returned, not by 

2286 re-probing llama-swap (which can transiently report empty right after a 

2287 successful load). The terminal phase is stamped in ``finally`` so an 

2288 unexpected error mid-warm still ends the launcher's progress stream. 

2289 """ 

2290 self._warm_tracker.begin(str(cfg.chat_model)) 

2291 warmed = False 

2292 try: 

2293 self._prewarm_chat_weights() 

2294 self._warm_tracker.loading_engine() 

2295 warmed = self._warm_role_clients(WorkerRole.CHAT, clients) 

2296 finally: 

2297 if warmed: 

2298 self._warm_tracker.ready() 

2299 else: 

2300 self._warm_tracker.fail(self._chat_load_failure()) 

2301 

2302 @contextmanager 

2303 def _lazy_warm_scope(self) -> Iterator[None]: 

2304 """Report a request-triggered chat load on the warm tracker. 

2305 

2306 The request that finds the chat role cold with no warm in flight owns 

2307 the warm: it stamps STARTING on entry and READY or ERROR on exit, on 

2308 every exit path. Concurrent requests share the load and touch nothing. 

2309 """ 

2310 owner = self._begin_lazy_warm() 

2311 try: 

2312 yield 

2313 except Exception as exc: 

2314 if owner: 

2315 self._end_lazy_warm(str(exc)) 

2316 raise 

2317 if owner: 

2318 self._end_lazy_warm(None) 

2319 

2320 def _begin_lazy_warm(self) -> bool: 

2321 """Begin a warm when the chat role is cold; True when this call owns it.""" 

2322 with self._lock: 

2323 if self._lazy_warming or self._warming: 

2324 return False 

2325 if is_active_warm(self._warm_tracker.snapshot()) or self.role_ready(WorkerRole.CHAT): 

2326 return False 

2327 with self._lock: 

2328 if self._lazy_warming: 

2329 return False 

2330 self._lazy_warming = True 

2331 self._warm_tracker.begin(str(cfg.chat_model)) 

2332 return True 

2333 

2334 def _end_lazy_warm(self, error: str | None) -> None: 

2335 """Stamp the owned warm READY when the role serves, else ERROR with *error*.""" 

2336 with self._lock: 

2337 self._lazy_warming = False 

2338 if error is None or self.role_ready(WorkerRole.CHAT): 

2339 self._warm_tracker.ready() 

2340 else: 

2341 self._warm_tracker.fail(error) 

2342 

2343 def _chat_load_failure(self) -> str: 

2344 """The engine's own reason the chat model did not load, when it gave one.""" 

2345 reason = self._warm_errors.get(WorkerRole.CHAT) 

2346 if not reason: 

2347 return "The chat model did not finish loading." 

2348 return f"The chat model did not load: {reason}" 

2349 

2350 def _prewarm_chat_weights(self) -> None: 

2351 """Page the chat model's GGUF shards into the OS cache, reporting byte progress. 

2352 

2353 Reading the shards before llama-swap loads them does two things: it gives a 

2354 true read-phase percentage for the warm tracker, and it warms the page cache 

2355 so the engine's mmap faults hit memory (a large win on a network filesystem, 

2356 where random mmap faults stalled cold loads). Best-effort: any failure to 

2357 resolve or size the shards (unregistered ref, cache miss, I/O error) is 

2358 skipped, and the model still loads on the warm request. 

2359 """ 

2360 try: 

2361 shards = ModelRegistry(cfg.models_dir).shard_paths(str(cfg.chat_model)) 

2362 total = sum(shard.stat().st_size for shard in shards) 

2363 except Exception: 

2364 log.debug("Prewarm skipped; could not resolve chat shards.", exc_info=True) 

2365 return 

2366 if total <= 0: 

2367 return 

2368 keys = [_prewarm_key(shard) for shard in shards] 

2369 if all(key in _PREWARMED_SHARDS for key in keys): 

2370 # Already paged in this boot (e.g. a placement rebuild); the cache is hot. 

2371 self._warm_tracker.reading(total, total) 

2372 return 

2373 done = 0 

2374 self._warm_tracker.reading(0, total) 

2375 chunk = bytearray(_PREWARM_CHUNK_BYTES) 

2376 for index, (shard, key) in enumerate(zip(shards, keys, strict=True)): 

2377 detail = f"shard {index + 1}/{len(shards)}" if len(shards) > 1 else None 

2378 try: 

2379 with shard.open("rb", buffering=0) as handle: 

2380 while True: 

2381 read = handle.readinto(chunk) 

2382 if not read: 

2383 break 

2384 done += read 

2385 self._warm_tracker.reading(done, total, detail=detail) 

2386 _PREWARMED_SHARDS.add(key) 

2387 except OSError: 

2388 # A partial/locked shard just shortens the read bar; the engine load 

2389 # surfaces any real fault as a user-facing error on the warm request. 

2390 log.debug("Prewarm read of %s stopped early.", shard, exc_info=True) 

2391 

2392 def cancel_inference(self) -> None: 

2393 """Sever every in-flight chat stream so its blocked reader unwinds. 

2394 

2395 A cooperative worker cancel cannot reach a thread blocked in a socket 

2396 read, and the reader's own close runs only when its worker unwinds, so 

2397 the disconnect must happen here. Retired clients are swept too: a 

2398 model-swap reload retires a busy client before the cancel lands. 

2399 """ 

2400 with self._lock: 

2401 clients = [*self._clients.get(WorkerRole.CHAT, ()), *self._retiring_clients] 

2402 for client in clients: 

2403 client.abort_streams() 

2404 

2405 def reload_role(self, role: WorkerRole, *, wait: bool = False) -> None: 

2406 """Apply a model/settings change for *role* with current cfg. 

2407 

2408 The whole fleet is re-planned, but only the roles whose launches changed 

2409 restart, so the other roles' loaded models stay resident (*role* names 

2410 the change for the thread label; the diff decides what restarts). A vision 

2411 setting change drops any on-request vision, so the plan follows the setting. 

2412 """ 

2413 if role is WorkerRole.VISION: 

2414 planning.revoke_vision_on_request() 

2415 self._dispatch_reload(f"fleet-reload-{role.value}", wait=wait) 

2416 

2417 def reload_placement(self, *, wait: bool = False) -> None: 

2418 """Apply a placement change with current cfg, restarting only moved roles. 

2419 

2420 The fresh plan is diffed per role against the running fleet: a role whose 

2421 devices (and so its launch argv) did not change keeps serving through the 

2422 change -- moving the embedder never unloads a 100GB chat model. When no 

2423 fleet is up, the next use plans fresh, so this returns at once. 

2424 """ 

2425 self._dispatch_reload("fleet-reload-placement", wait=wait) 

2426 

2427 def _dispatch_reload(self, thread_name: str, *, wait: bool) -> None: 

2428 """Run the diff-driven reload once, off-thread unless *wait*. 

2429 

2430 Dispatched to a background thread because the slow restart (rewrite config + 

2431 respawn + wait-ready) must not block the settings/model-picker callback. 

2432 If no group is up yet, the next use starts the fleet with current cfg. 

2433 Single-flight: a reload while one is in flight sets the pending flag (the 

2434 in-flight pass may have already snapshotted its plan), and the in-flight 

2435 thread runs one more pass per pending flag so the change is applied, not 

2436 dropped. 

2437 

2438 ``wait=True`` runs the reload in the caller's thread and returns only once 

2439 the restart (and any reload already in flight that will run the pending 

2440 pass) has finished and the proxies are healthy again, so a caller already 

2441 off the event loop gets a real completion signal. A restarted role's model 

2442 still loads lazily (the reload kicks an off-thread warm). It propagates a 

2443 reload failure as an exception. 

2444 """ 

2445 with self._lock: 

2446 if not self._swaps: 

2447 return 

2448 if self._reloading: 

2449 self._reload_pending = True 

2450 if wait: 

2451 while self._reloading: 

2452 self._reload_done.wait() 

2453 return 

2454 self._reloading = True 

2455 self._reload_pending = False 

2456 if wait: 

2457 self._reload_blocking() 

2458 return 

2459 threading.Thread( 

2460 target=self._reload_blocking, 

2461 name=thread_name, 

2462 daemon=True, 

2463 ).start() 

2464 

2465 def _reload_blocking(self) -> None: 

2466 """Run reload passes until no further reload arrived mid-pass. 

2467 

2468 A failed pass with the pending flag set still runs the pending pass (the 

2469 fresh plan may succeed under the new cfg); only the final pass's failure 

2470 propagates, after dropping the refs to a dead swap so the next call can 

2471 rebuild. The pending check and the guard release happen under one lock 

2472 acquisition, so a reload_role landing between them cannot be acknowledged 

2473 and dropped. 

2474 """ 

2475 while True: 

2476 try: 

2477 self._reload_pass() 

2478 except BaseException: 

2479 with self._lock: 

2480 rerun = self._reload_pending 

2481 self._reload_pending = False 

2482 if not rerun: 

2483 self._reloading = False 

2484 self._reload_done.notify_all() 

2485 if rerun: 

2486 log.warning( 

2487 "Engine reload failed; retrying with the pending change.", exc_info=True 

2488 ) 

2489 continue 

2490 self._drop_dead_swaps() 

2491 raise 

2492 with self._lock: 

2493 if not self._reload_pending: 

2494 self._reloading = False 

2495 self._reload_done.notify_all() 

2496 return 

2497 self._reload_pending = False 

2498 

2499 def _rebind_or_overflow(self) -> list[WorkerRole]: 

2500 """Re-acquire a bound engine after a config change; caller holds the build lock. 

2501 

2502 A provider that rode another process's engine owns none of its groups and 

2503 cannot restart them. Restarting "in place" would spawn a duplicate fleet 

2504 into the shared slot (a bound manager's ``shutdown`` only detaches, leaving 

2505 the incumbent resident) and size it blind against VRAM the incumbent still 

2506 holds. Instead drop every binding and this process's membership, then re-run 

2507 the acquisition ladder: it rebinds to the reconfigured shared engine, builds 

2508 fresh in the machine slot if we were its last user, or overflows to a private 

2509 engine sized against a fresh probe. Returns the roles now served (to preload). 

2510 """ 

2511 from lilbee.core.config import cfg 

2512 

2513 with self._lock: 

2514 groups = list(self._swaps) 

2515 for group in groups: 

2516 with self._lock: 

2517 swap = self._drop_group(group) # also prunes _group_dirs 

2518 if swap is not None: 

2519 swap.shutdown() # bound: detaches; the shared engine keeps running 

2520 self._release_engines() # a shared engine's builder keeps it live; no stop here 

2521 if not self._acquire_engine(cfg.data_root): 

2522 return [] 

2523 with self._lock: 

2524 return list(self._role_group) 

2525 

2526 def _reload_pass(self, force: frozenset[WorkerRole] = frozenset()) -> None: 

2527 """One re-plan from current cfg, restarting only the groups that changed. 

2528 

2529 The fresh plan is diffed per swap group against the launches each running 

2530 group was started with; a group restarts only when its launches differ 

2531 (covers added and removed groups too), so an untouched group's loaded model 

2532 stays resident through a placement or per-role model change. *force* adds a 

2533 role's group to the restart set even when its plan is unchanged (dead-swap 

2534 recovery). Changed groups stop before the new ones start, so the planned 

2535 VRAM is actually free when the new servers spawn. Runs under the build lock 

2536 so a racing shutdown/build can't interleave with the restart and leak a live 

2537 llama-swap holding GPU memory. 

2538 """ 

2539 

2540 restarted: list[WorkerRole] = [] 

2541 with self._build_lock: 

2542 # The device list is structural and was captured once at boot, so a 

2543 # card that has since left keeps being planned onto. The memory 

2544 # figures beside it are deliberately not re-taken: this fleet is 

2545 # resident, and charging it against itself is what the snapshot exists 

2546 # to prevent. 

2547 planning.refresh_plan_devices() 

2548 with self._lock: 

2549 if self._shut_down: 

2550 # Terminal shutdown landed while this reload was queued; a 

2551 # rebuild here would spawn a fleet no live provider owns. 

2552 return 

2553 running = set(self._swaps) 

2554 old = dict(self._launches) 

2555 # All groups share one dir and one ownership by construction, so any 

2556 # bound manager means this provider rides another process's engine. 

2557 bound = any(swap.bound for swap in self._swaps.values()) 

2558 if bound: 

2559 # Cannot restart a shared engine's groups in place (that duplicates 

2560 # the fleet into the slot); drop the bindings and re-acquire. 

2561 restarted = self._rebind_or_overflow() 

2562 self._preload_restarted(restarted) 

2563 return 

2564 # Reap dead engines in our dirs before re-planning, as a build would. 

2565 reload_dir = self._reload_dir() 

2566 # Serialize the reload against peer acquisitions with the same 

2567 # cross-process lock a build takes. Without it, this reap_stale can kill 

2568 # a peer's swap that is spawned but not yet answering its proxy, and the 

2569 # stop-then-spawn gap lets a peer's ladder see a half-stopped slot and 

2570 # build a second fleet into the same dir, double-allocating VRAM. 

2571 with build_lock(reload_dir): 

2572 reap_stale(reload_dir) 

2573 try: 

2574 if not running: 

2575 # Nothing loaded (a resurrect after a failed pass): the box is 

2576 # clean, so refresh the plan snapshot like a first build would. 

2577 planning.capture_plan_probe() 

2578 plan = planning.plan_all_launches() 

2579 except ProviderError as exc: 

2580 # Same policy as the initial build: a genuinely-missing engine 

2581 # binary aborts the reload quietly (nothing to serve), while any 

2582 # other planning failure (a wedged GPU probe, an unusable CUDA 

2583 # runtime) propagates to fail loud. The raise lands before the 

2584 # stop phase, so a running fleet is left intact rather than half 

2585 # torn down. 

2586 if exc.kind is not ProviderErrorKind.NOT_FOUND: 

2587 raise 

2588 log.debug("Engine binary unavailable; reload left the fleet as-is") 

2589 return 

2590 # Keep the skip reasons in step with the fresh plan. 

2591 self._skipped_not_installed = dict(plan.skipped_not_installed) 

2592 self._skipped_unusable_ctx = dict(plan.skipped_unusable_ctx) 

2593 new = _launches_by_group(plan) 

2594 # A group restarts when its launches changed OR its running/planned 

2595 # presence disagrees (covers a group the new plan drops or adds). 

2596 changed = { 

2597 group 

2598 for group in running | set(new) 

2599 if (group in running) != (group in new) 

2600 or old.get(group, ()) != new.get(group, ()) 

2601 } 

2602 changed |= {group_for(role, plan.co_tenants) for role in force} 

2603 # Stop phase: free the changed groups' VRAM before their replacements 

2604 # (or another group's grown plan) spawn against it. 

2605 for group in sorted(changed, key=lambda g: g.value): 

2606 with self._lock: 

2607 swap = self._drop_group(group) 

2608 if swap is not None: 

2609 swap.shutdown() 

2610 # Start phase: spawn the changed groups present in the new plan. 

2611 for group in sorted(changed & set(new), key=lambda g: g.value): 

2612 group_launches = list(new[group]) 

2613 swap = SwapManager(reload_dir, group) 

2614 swap.start( 

2615 group_launches, 

2616 ttl_seconds=_warm_ttl_seconds( 

2617 hold_warm_for_session=self._hold_warm_for_session 

2618 ), 

2619 bind_lifetime=not cfg.keep_engine_warm, 

2620 ) 

2621 with self._lock: 

2622 self._adopt_group(group, swap, group_launches) 

2623 self._group_dirs[group] = reload_dir 

2624 restarted.extend(_by_role(group_launches)) 

2625 self._preload_restarted(restarted) 

2626 

2627 def _preload_restarted(self, restarted: list[WorkerRole]) -> None: 

2628 """Load the restarted roles' models off-thread (a no-op for none). 

2629 

2630 llama-swap spawns an upstream on its first request, so the reload returns 

2631 once the proxies answer and the UI's spawn listeners track the model loads. 

2632 """ 

2633 if not restarted: 

2634 return 

2635 threading.Thread( 

2636 target=self._preload_roles, 

2637 kwargs={"roles": frozenset(restarted)}, 

2638 name="fleet-reload-warm", 

2639 daemon=True, 

2640 ).start() 

2641 

2642 def add_spawn_listener( 

2643 self, 

2644 *, 

2645 on_spawning: Callable[[WorkerRole], None] | None = None, 

2646 on_spawned: Callable[[WorkerRole], None] | None = None, 

2647 ) -> None: 

2648 """Store spawn-lifecycle callbacks; warm-up fires them as each role loads.""" 

2649 with self._lock: 

2650 self._on_spawning = on_spawning 

2651 self._on_spawned = on_spawned 

2652 

2653 def invalidate_load_cache(self, model_path: Path | None = None) -> None: 

2654 """A model or settings change restarts the engine: drop the swap.""" 

2655 del model_path # the whole engine restarts on next use; no per-model scope. 

2656 self._shutdown_swap(latch=False) 

2657 

2658 def drop_loaded_models_async(self) -> None: 

2659 """Drop the swap off the caller's thread; next use restarts with current cfg. 

2660 

2661 ``_shutdown_swap`` stops llama-swap and waits on its process group, so a 

2662 role-agnostic load-key change (num_ctx, kv_cache_type) routes here rather 

2663 than blocking the settings callback. A no-op when no swap is up. 

2664 """ 

2665 with self._lock: 

2666 if not self._swaps: 

2667 return 

2668 threading.Thread( 

2669 target=lambda: self._shutdown_swap(latch=False), 

2670 name="fleet-drop", 

2671 daemon=True, 

2672 ).start() 

2673 

2674 def shutdown(self) -> None: 

2675 # The on-request vision grant is process-wide; the next provider starts from the setting. 

2676 planning.revoke_vision_on_request() 

2677 self._shutdown_swap()