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

563 statements  

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

1"""Supervise the single llama-swap process that fronts every fleet role. 

2 

3llama-swap owns each role's llama-server lifecycle; this manages the one proxy 

4process and exposes its endpoint and readiness. See docs/architecture.md. 

5""" 

6 

7from __future__ import annotations 

8 

9import contextlib 

10import itertools 

11import json 

12import logging 

13import os 

14import signal 

15import socket 

16import subprocess 

17import sys 

18import threading 

19import time 

20from collections.abc import Iterable, Iterator 

21from dataclasses import dataclass 

22from functools import lru_cache 

23from pathlib import Path 

24from typing import TYPE_CHECKING, BinaryIO 

25 

26import httpx 

27import psutil 

28 

29from lilbee.core.health_warnings import HealthWarning 

30from lilbee.providers.base import ProviderError, ProviderErrorKind 

31from lilbee.providers.fleet.binary import engine_pin, resolve_llama_swap 

32from lilbee.providers.fleet.child_guard import release_death_pipe, spawn_bound_child 

33from lilbee.providers.fleet.groups import SwapGroup 

34from lilbee.providers.fleet.launch import role_model_prefix 

35from lilbee.providers.fleet.planning import clear_ctx_downshift, log_engine_launch 

36from lilbee.providers.fleet.readback import ( 

37 MEMORY_FLAG, 

38 check_launch, 

39 check_memory_report, 

40 report_missing_log, 

41) 

42from lilbee.providers.fleet.swap_config import PORT_FLAG, build_swap_config 

43from lilbee.runtime.engine_lock import clear_keep_warm 

44 

45if TYPE_CHECKING: 

46 from lilbee.providers.fleet.launch import InstanceLaunch 

47 from lilbee.providers.roles import WorkerRole 

48 

49log = logging.getLogger(__name__) 

50 

51_HOST = "127.0.0.1" 

52# One llama-swap per swap group: the group name lands in the config filename so 

53# each group's processes are identified (and stopped) by their own config path, 

54# and a placement change can restart one group without touching the others. 

55# The writer pid segment is uniqueness, not ownership: the build lock ensures 

56# one builder per engine dir, and reaping cleans dead writers' leftovers. 

57_CONFIG_FILENAME_TEMPLATE = "llama-swap-{group}.{pid}.json" 

58_CONFIG_FILE_GLOB = "llama-swap-*.json" 

59# llama-swap's own stdout/stderr (its HTTP access log) is captured to a file in a 

60# ``logs/`` dir inside the engine dir, which is the machine slot rather than any 

61# one lilbee's data root, so the log sits beside the engine it belongs to instead 

62# of beside server.log. Capturing it at all, rather than inheriting the parent's 

63# fd, is because a TUI or CLI parent owns the terminal and an inherited fd would bleed 

64# llama-swap's request log onto the screen and corrupt the render. Per-model 

65# upstream logs are unaffected (those go to llama-swap's /logs API). 

66_LOGS_SUBDIR = "logs" 

67_LOG_FILENAME_TEMPLATE = "llama-swap-{group}.log" 

68# Each writer's state file records its swap's pid/pgid so a later start can 

69# stop a dead or unhealthy engine. Health, not ownership, decides sparing. 

70_STATE_FILENAME_PREFIX = "llama-swap.state." 

71_STATE_FILENAME_SUFFIX = ".json" 

72# Also matches the legacy single shared state file ("llama-swap.state.json"). 

73_STATE_FILE_GLOB = f"{_STATE_FILENAME_PREFIX}*" 

74_STATE_KEY_PID = "pid" 

75_STATE_KEY_PGID = "pgid" 

76_STATE_KEY_CREATED_AT = "created_at" 

77_STATE_KEY_NAME = "name" 

78_STATE_KEY_MEMBER_PORTS = "member_ports" 

79_STATE_KEY_PROXY_PORT = "proxy_port" 

80_STATE_KEY_LAUNCHES = "launches" 

81_STATE_KEY_ENGINE_PIN = "engine_pin" 

82# Atomic state writes: the dot prefix keeps half-written tmp files out of the 

83# reap scan's glob. 

84_STATE_TMP_PREFIX = "." 

85_STATE_TMP_SUFFIX = ".tmp" 

86# Pid reuse guard: a live process at a recorded pid whose create time differs 

87# from the recorded one by more than this is a different process. 

88_CREATE_TIME_TOLERANCE_S = 1.0 

89_LLAMA_SWAP_PROCESS_NAME = "llama-swap" 

90_LLAMA_SERVER_PROCESS_NAME = "llama-server" 

91_CONFIG_FLAG = "-config" 

92_LISTEN_FLAG = "-listen" 

93_HEALTH_PATH = "/health" 

94_RUNNING_PATH = "/running" 

95_HTTP_TIMEOUT_S = 10.0 

96# llama-swap's own proxy answers within a second; upstream model loads have their 

97# own (longer) budget inside llama-swap, so this only covers the proxy coming up. 

98_BOOT_TIMEOUT_S = 30.0 

99_BOOT_POLL_S = 0.25 

100# Cap on the captured llama-swap output a boot-failure error carries. 

101_BOOT_LOG_TAIL_CHARS = 2000 

102# Per-group SIGTERM grace before SIGKILL on the manager shutdown/reload path. A 

103# hard kill is safe (llama-server holds no persistent state). Note this is NOT 

104# the constant the serve handoff waits on: that path goes through stop_engine -> 

105# _stop_stale_swap and spends _ORPHAN_STOP_TIMEOUT_S plus the kill/reap waits, so 

106# SERVER_LOCK_TIMEOUT budgets only a teardown whose SIGTERMs are honored. 

107_STOP_TIMEOUT_S = 2.5 

108# Grace for a llama-server that outlived llama-swap before it is force-killed. 

109_ORPHAN_STOP_TIMEOUT_S = 5.0 

110# Grace for a SIGKILLed process to exit (and release its VRAM) before the next 

111# free-memory probe runs. 

112_KILL_WAIT_TIMEOUT_S = 5.0 

113_PROBE_TIMEOUT_S = 5.0 

114# Liveness probes talk to a loopback proxy, so they get their own short budget 

115# rather than the module's 10 s general HTTP one. The ladder runs this probe for 

116# every group while holding the cross-process build lock, so one wedged port 

117# (SYN-accepted but unresponsive) would otherwise stall every other lilbee start 

118# for tens of seconds. A local proxy that cannot answer /running this fast is 

119# not usable for inference either. 

120_LIVENESS_TIMEOUT = httpx.Timeout(connect=0.5, read=2.0, write=2.0, pool=2.0) 

121 

122 

123@lru_cache(maxsize=1) 

124def _probe_client() -> httpx.Client: 

125 """One shared client for the localhost engine probes. 

126 

127 ``httpx.get`` builds a fresh ``Client`` per call, and every ``Client`` 

128 construction creates an SSL context, which loads the system CA bundle. These 

129 probes are plain HTTP to 127.0.0.1, so none of that TLS setup is ever used -- 

130 and the readiness probe runs on the task bar's timer (up to 10 Hz), which made 

131 ``ssl.create_default_context`` 23% of TUI CPU in a py-spy profile. One client 

132 builds that at most once and keeps the connection alive between polls. 

133 ``trust_env`` is off so a proxy env var cannot redirect a loopback probe. 

134 """ 

135 return httpx.Client(trust_env=False) 

136 

137 

138_PROVIDER = "llama-server" 

139# /running JSON shape: {"running": [{"model": <id>, "state": "ready", ...}, ...]}. 

140_KEY_RUNNING = "running" 

141_KEY_MODEL = "model" 

142_KEY_STATE = "state" 

143_STATE_READY = "ready" 

144 

145 

146def _platform_const(module: object, name: str, default: int) -> int: 

147 """A platform-conditional stdlib constant (absent on some OSes -> default).""" 

148 return getattr(module, name, default) 

149 

150 

151_CREATE_NEW_PROCESS_GROUP: int = _platform_const(subprocess, "CREATE_NEW_PROCESS_GROUP", 0) 

152_SIGKILL: int = _platform_const(signal, "SIGKILL", signal.SIGTERM) 

153 

154 

155def _atomic_write(path: Path, text: str) -> None: 

156 """Write *text* to *path* via a temp file in the same dir, then rename over it. 

157 

158 A plain write truncates the destination first, so a process dying mid-write 

159 (OOM kill, SIGKILL, disk full) leaves an empty or half-written file behind. 

160 For the llama-swap config that means the next spawn hands the engine a file 

161 it cannot start from; for the state file it means a sibling's reap scan 

162 reads a torn record. 

163 

164 The temp name carries the destination's name, and both config and state 

165 filenames embed the writing process's pid, so a crash leftover can be told 

166 from a live writer's file in flight -- see ``_clean_stale_tmp_files``. 

167 """ 

168 tmp_path = path.with_name(f"{_STATE_TMP_PREFIX}{path.name}{_STATE_TMP_SUFFIX}") 

169 tmp_path.write_text(text, encoding="utf-8") 

170 os.replace(tmp_path, path) 

171 

172 

173def _state_filename(owner_pid: int, group: str) -> str: 

174 """The per-owner, per-group state filename for the lilbee process *owner_pid*.""" 

175 return f"{_STATE_FILENAME_PREFIX}{group}.{owner_pid}{_STATE_FILENAME_SUFFIX}" 

176 

177 

178@dataclass(frozen=True) 

179class SwapState: 

180 """A running llama-swap's recorded identity and serving contract. 

181 

182 Read back from the engine dir's state file, so it describes engines this 

183 process did not start. The currency the bind/build ladder is written in: 

184 swap_manager records it, provider reads it to decide what a slot is 

185 serving, and contract matches it against what this process wants. 

186 """ 

187 

188 pid: int 

189 pgid: int | None 

190 created_at: float | None = None 

191 member_ports: tuple[int, ...] = () 

192 proxy_port: int | None = None 

193 launches: tuple[dict, ...] = () 

194 engine_pin: str | None = None 

195 

196 

197class SwapManager: 

198 """Owns one llama-swap process fronting one role group's servers. 

199 

200 The provider runs one manager per role, so restarting a group (a placement 

201 or model change) never touches another group's loaded servers. 

202 """ 

203 

204 def __init__(self, data_dir: Path, group: SwapGroup) -> None: 

205 self._data_dir = data_dir 

206 self._group = group 

207 self._config_path = data_dir / _config_filename(os.getpid(), group.value) 

208 self._log_path = data_dir / _LOGS_SUBDIR / _LOG_FILENAME_TEMPLATE.format(group=group.value) 

209 # Instances whose engine report has already been compared to the estimate, 

210 # so the check runs once per start rather than on every readiness poll. 

211 self._estimate_checked: set[str] = set() 

212 self._launch_by_model: dict[str, InstanceLaunch] = {} 

213 # Placement divergences by model id, surfaced on health. Guarded because 

214 # role_ready probes off the provider lock while health reads anytime. 

215 self._placement_warnings: dict[str, HealthWarning] = {} 

216 self._warnings_lock = threading.Lock() 

217 self._state_path = data_dir / _state_filename(os.getpid(), group.value) 

218 self._proc: subprocess.Popen[bytes] | None = None 

219 self._log_file: BinaryIO | None = None 

220 # Where this boot's output starts in the append-mode log file. 

221 self._log_offset = 0 

222 self._port: int | None = None 

223 self._member_ports: list[int] = [] 

224 # Which member port serves which model id, for the direct GET /memory 

225 # readback. Only launches this manager spawned are in here: a bound 

226 # manager's state record carries the ports but not the mapping, and it 

227 # has no launches to check either. 

228 self._member_port_by_model: dict[str, int] = {} 

229 # The serving contract (per-role model/ctx/slots) persisted in every 

230 # state write, so a guest lilbee can bind to this live fleet. 

231 self._launches_payload: list[dict] = [] 

232 # True when this manager uses an engine another process built: it then 

233 # never writes state, never reaps, and never signals engine processes. 

234 self._bound = False 

235 

236 def start( 

237 self, launches: list[InstanceLaunch], *, ttl_seconds: int = 0, bind_lifetime: bool = True 

238 ) -> None: 

239 """Write the config and spawn llama-swap, waiting for its proxy to answer. 

240 

241 The proxy and every member get a freshly allocated free port, which is 

242 why llama-swap's own startPort is not used: that assigns a fixed 

243 sequential range at config load, so it would collide with a previous 

244 instance's server still shutting down (the new llama-server then fails 

245 its bind and llama-swap reports it only as "exited prematurely"). 

246 

247 This narrows that collision rather than removing a race. The ports are 

248 picked by binding and closing ephemeral sockets, while llama-swap starts 

249 each upstream lazily on its first request, so a member port can sit 

250 unbound for as long as it takes that request to arrive and anything else 

251 on the box may take it in between. Nothing in llama-swap offers a 

252 spawn-time probe to close that window; warming the roles up front 

253 shortens it for the roles that are warmed. 

254 

255 ``bind_lifetime`` binds the engine to this process so a crash cannot orphan 

256 it; it is False for a keep-warm fleet that is meant to outlive lilbee. 

257 """ 

258 # Idempotent safety net; the provider reaps before planning so the GPU 

259 # probe already saw the real free memory. 

260 self.reap_stale() 

261 # Singleton guard: one llama-swap per data_dir for this lilbee. Reap any 

262 # llama-swap we already started against this config (a leaked duplicate 

263 # from a prior race/reload) before spawning, so they cannot accumulate 

264 # and double-book a GPU. 

265 _stop_own_fleet(self._config_path, tuple(self._member_ports)) 

266 ports = _pick_free_ports(1 + len(launches)) 

267 member_ports = dict(zip([launch.model_id for launch in launches], ports[1:], strict=True)) 

268 self._member_ports = sorted(member_ports.values()) 

269 self._member_port_by_model = dict(member_ports) 

270 self._launches_payload = [launch.to_state() for launch in launches] 

271 self._launch_by_model = {launch.model_id: launch for launch in launches} 

272 self._estimate_checked.clear() 

273 with self._warnings_lock: 

274 self._placement_warnings.clear() 

275 self._config_path.parent.mkdir(parents=True, exist_ok=True) 

276 self._log_path.parent.mkdir(parents=True, exist_ok=True) 

277 _atomic_write( 

278 self._config_path, 

279 build_swap_config( 

280 launches, 

281 member_ports, 

282 swap=self._group.swaps, 

283 ttl_seconds=ttl_seconds, 

284 engine_log_dir=self._log_path.parent, 

285 ), 

286 ) 

287 self._port = ports[0] 

288 # Capture llama-swap's stdout/stderr to a file so its access log never 

289 # reaches an inherited terminal (a TUI/CLI parent) and garbles the screen. 

290 self._close_log() 

291 self._log_path.parent.mkdir(parents=True, exist_ok=True) 

292 self._log_file = self._log_path.open("ab") 

293 self._log_offset = self._log_path.stat().st_size 

294 self._proc = spawn_bound_child( 

295 [ 

296 str(resolve_llama_swap()), 

297 _CONFIG_FLAG, 

298 str(self._config_path), 

299 _LISTEN_FLAG, 

300 f"{_HOST}:{self._port}", 

301 ], 

302 bind_lifetime=bind_lifetime, 

303 stdout=self._log_file, 

304 stderr=subprocess.STDOUT, 

305 start_new_session=True, 

306 creationflags=_CREATE_NEW_PROCESS_GROUP, 

307 ) 

308 self._write_state() 

309 self._await_health() 

310 for launch in launches: 

311 log_engine_launch(launch) 

312 

313 def reap_stale(self) -> None: 

314 """Kill every dead or unhealthy recorded engine; see :func:`reap_stale`.""" 

315 reap_stale(self._data_dir) 

316 

317 def _process_identity(self) -> tuple[int, int | None, float | None] | None: 

318 """(pid, pgid, create time) of the swap this manager runs, or None.""" 

319 if self._proc is not None: 

320 pid = self._proc.pid 

321 pgid: int | None = None 

322 if sys.platform != "win32": 

323 with contextlib.suppress(ProcessLookupError): 

324 pgid = os.getpgid(pid) 

325 created_at: float | None = None 

326 with contextlib.suppress(psutil.NoSuchProcess, psutil.AccessDenied): 

327 created_at = psutil.Process(pid).create_time() 

328 return pid, pgid, created_at 

329 return None 

330 

331 def _write_state(self) -> None: 

332 """Record the swap's pid/pgid/create time, member ports, and our identity. 

333 

334 The write is atomic (tmp file then ``os.replace``) so a sibling's reap 

335 scan can never read a torn file and mistake this live record for junk. 

336 """ 

337 identity = self._process_identity() 

338 if identity is None: 

339 return 

340 swap_pid, pgid, created_at = identity 

341 state = { 

342 _STATE_KEY_PID: swap_pid, 

343 _STATE_KEY_PGID: pgid, 

344 _STATE_KEY_CREATED_AT: created_at, 

345 _STATE_KEY_NAME: _LLAMA_SWAP_PROCESS_NAME, 

346 _STATE_KEY_MEMBER_PORTS: self._member_ports, 

347 _STATE_KEY_PROXY_PORT: self._port, 

348 _STATE_KEY_LAUNCHES: self._launches_payload, 

349 _STATE_KEY_ENGINE_PIN: engine_pin(), 

350 } 

351 _atomic_write(self._state_path, json.dumps(state)) 

352 

353 @property 

354 def log_path(self) -> Path: 

355 """The llama-swap log for this group, which records each member's exit.""" 

356 return self._log_path 

357 

358 def endpoint(self) -> str: 

359 """Base URL of the llama-swap OpenAI-compatible proxy.""" 

360 if self._port is None: 

361 raise ProviderError( 

362 "The local model engine is not running.", 

363 provider=_PROVIDER, 

364 kind=ProviderErrorKind.SERVER, 

365 ) 

366 return f"http://{_HOST}:{self._port}" 

367 

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

369 """Whether at least one of *role*'s replica servers is loaded and ready.""" 

370 prefix = role_model_prefix(role) 

371 ready = self._ready_models() 

372 self._check_estimates(ready) 

373 return any(model.startswith(prefix) for model in ready) 

374 

375 def _check_estimates(self, ready: set[str]) -> None: 

376 """Compare each newly-ready engine's own report against what it was planned for. 

377 

378 The plan is otherwise open-loop, and a wrong estimate only ever surfaces 

379 as a failed request much later. Once per instance per start: readiness is 

380 polled, and the answer does not change once the engine has loaded. 

381 """ 

382 for model_id in ready - self._estimate_checked: 

383 launch = self._launch_by_model.get(model_id) 

384 self._estimate_checked.add(model_id) 

385 if launch is None: 

386 continue 

387 # Ready means this role's context loaded, so any reduction taken to 

388 # get here has done its job and must not follow the role into the 

389 # next plan, a freed machine, or a model the user switched to. 

390 clear_ctx_downshift(launch.role) 

391 # A launch carrying --memory serves its own report on GET /memory 

392 # and was given no trace log to read (swap_config), so the log-side 

393 # checks would misfire on it by construction. 

394 if MEMORY_FLAG in launch.argv: 

395 self._check_memory_estimate(model_id, launch) 

396 continue 

397 # The engine is ready, so a missing log is not "too early" any more. 

398 if report_missing_log(self._log_path.parent, model_id, launch.role): 

399 continue 

400 if launch.est_vram_bytes <= 0: 

401 continue 

402 warning = check_launch( 

403 self._log_path.parent, 

404 model_id, 

405 launch.role, 

406 launch.model, 

407 launch.est_vram_bytes, 

408 launch.est_vram_by_device, 

409 launch.est_unreported_bytes, 

410 ) 

411 if warning is not None: 

412 with self._warnings_lock: 

413 self._placement_warnings[model_id] = warning 

414 

415 def _check_memory_estimate(self, model_id: str, launch: InstanceLaunch) -> None: 

416 """Run the estimate check off the engine's own ``GET /memory``. 

417 

418 Straight to the member's port rather than through llama-swap: the model 

419 is ready, so the server owns its port, and the proxy adds only a routing 

420 layer that has nothing to route on for a bare GET. An endpoint that does 

421 not answer after the engine took the flag is reported, not swallowed -- 

422 it is this mode's analog of a ready engine that wrote no log. 

423 """ 

424 port = self._member_port_by_model.get(model_id) 

425 if port is None: # bound to another process's engine; nothing was planned here 

426 return 

427 try: 

428 resp = _probe_client().get(f"http://{_HOST}:{port}/memory", timeout=_PROBE_TIMEOUT_S) 

429 payload = resp.json() if resp.status_code == httpx.codes.OK else None 

430 except (httpx.HTTPError, ValueError): 

431 payload = None 

432 if payload is None: 

433 log.warning( 

434 "The %s engine was launched with %s but its /memory endpoint did not " 

435 "answer, so its memory use could not be checked against the estimate. " 

436 "Placement estimates for this model are unverified.", 

437 launch.role.value, 

438 MEMORY_FLAG, 

439 ) 

440 return 

441 if launch.est_vram_bytes <= 0: 

442 return 

443 warning = check_memory_report( 

444 launch.role, 

445 launch.model, 

446 launch.est_vram_bytes, 

447 launch.est_vram_by_device, 

448 payload, 

449 ) 

450 if warning is not None: 

451 with self._warnings_lock: 

452 self._placement_warnings[model_id] = warning 

453 

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

455 """Placement divergences recorded from ready engines in this group.""" 

456 with self._warnings_lock: 

457 return list(self._placement_warnings.values()) 

458 

459 def is_live(self) -> bool: 

460 """Whether the swap process is up and its proxy answers ``/running``.""" 

461 if self._proc is None or self._proc.poll() is not None: 

462 return False 

463 if self._port is None: 

464 return False 

465 return self._proxy_answers() 

466 

467 @property 

468 def running(self) -> bool: 

469 """Whether this manager currently has a spawned llama-swap process.""" 

470 return self._proc is not None 

471 

472 @property 

473 def bound(self) -> bool: 

474 """Whether this manager rides an engine built by another process.""" 

475 return self._bound 

476 

477 def bind(self, state: SwapState) -> bool: 

478 """Use a running engine's proxy without taking any ownership of it. 

479 

480 The engine's own state record stays untouched: the binder writes 

481 nothing, and shutdown() merely drops the binding. 

482 """ 

483 if state.proxy_port is None: 

484 return False 

485 self._port = state.proxy_port 

486 self._member_ports = list(state.member_ports) 

487 if not self._proxy_answers(): 

488 self._port = None 

489 self._member_ports = [] 

490 return False 

491 self._launches_payload = [dict(launch) for launch in state.launches] 

492 self._bound = True 

493 return True 

494 

495 def _proxy_answers(self) -> bool: 

496 """Whether the bound proxy port serves llama-swap's running endpoint. 

497 

498 Shares state_is_healthy's identity check via _running_endpoint_answers, so 

499 bind and reap agree on what "answering" means by construction rather than by 

500 two hand-kept-identical copies. 

501 """ 

502 return _running_endpoint_answers(self.endpoint()) 

503 

504 def shutdown(self) -> None: 

505 """Stop every llama-swap this lilbee owns at our config and reap servers. 

506 

507 Authoritative teardown keyed on config-path identity, not the single 

508 tracked ``Popen``: a warm-up/reset race or a reload can leave several 

509 llama-swap processes this lilbee spawned, any of them reparented to init 

510 (still holding the engine binary open) -- trusting one handle would leak 

511 them. Every llama-swap running against our config is reaped. Unlinks only 

512 this owner's state file; another instance's record stays. 

513 """ 

514 if self._bound: 

515 # Not ours to stop: drop the binding and leave the engine serving. 

516 self._bound = False 

517 self._port = None 

518 self._member_ports = [] 

519 self._launches_payload = [] 

520 return 

521 _stop_own_fleet(self._config_path, tuple(self._member_ports)) 

522 # Nothing is coming back to bind these, so the picker can offer them again. 

523 release_reserved_ports([*self._member_ports, *([self._port] if self._port else [])]) 

524 self._state_path.unlink(missing_ok=True) 

525 if self._proc is not None: 

526 # Free this engine's death pipe so its watcher exits now, not at our death. 

527 release_death_pipe(self._proc.pid) 

528 self._proc = None 

529 self._port = None 

530 self._close_log() 

531 

532 def _close_log(self) -> None: 

533 """Close the captured llama-swap log handle, if one is open.""" 

534 if self._log_file is not None: 

535 with contextlib.suppress(OSError): 

536 self._log_file.close() 

537 self._log_file = None 

538 

539 def _await_health(self) -> None: 

540 """Poll the proxy's /health until it answers, or fail with a clear error.""" 

541 url = f"{self.endpoint()}{_HEALTH_PATH}" 

542 deadline = time.monotonic() + _BOOT_TIMEOUT_S 

543 while time.monotonic() < deadline: 

544 if self._proc is not None and self._proc.poll() is not None: 

545 self._fail("The local model engine exited before it was ready.") 

546 with contextlib.suppress(httpx.HTTPError): 

547 if _probe_client().get(url, timeout=_PROBE_TIMEOUT_S).status_code == httpx.codes.OK: 

548 return 

549 time.sleep(_BOOT_POLL_S) 

550 self._fail("The local model engine did not start in time.") 

551 

552 def _ready_models(self) -> set[str]: 

553 """Model ids whose upstream is loaded and ready, per llama-swap's /running. 

554 

555 A read-only probe: a concurrent shutdown can clear ``_port`` between the 

556 caller's check and ``endpoint()``, raising ProviderError, so that is 

557 suppressed too and the probe reports "nothing ready" rather than throwing. 

558 """ 

559 with contextlib.suppress(httpx.HTTPError, ValueError, KeyError, TypeError, ProviderError): 

560 payload = ( 

561 _probe_client() 

562 .get(f"{self.endpoint()}{_RUNNING_PATH}", timeout=_PROBE_TIMEOUT_S) 

563 .json() 

564 ) 

565 return { 

566 entry[_KEY_MODEL] 

567 for entry in payload[_KEY_RUNNING] 

568 if entry.get(_KEY_STATE) == _STATE_READY 

569 } 

570 return set() 

571 

572 def _boot_log_tail(self) -> str: 

573 """The current boot's captured llama-swap output, capped for an error message.""" 

574 try: 

575 with self._log_path.open("rb") as handle: 

576 handle.seek(self._log_offset) 

577 data = handle.read() 

578 except OSError: 

579 return "" 

580 return data.decode(errors="replace").strip()[-_BOOT_LOG_TAIL_CHARS:] 

581 

582 def _fail(self, message: str) -> None: 

583 """Tear down and raise a user-facing engine-start error carrying the boot log.""" 

584 self.shutdown() 

585 tail = self._boot_log_tail() 

586 if tail: 

587 message = f"{message} Engine log ({self._log_path}):\n{tail}" 

588 raise ProviderError(message, provider=_PROVIDER, kind=ProviderErrorKind.SERVER) 

589 

590 

591# Linux publishes the range here; every other platform is asked via sysctl. 

592_PROC_PORT_RANGE = Path("/proc/sys/net/ipv4/ip_local_port_range") 

593 

594 

595def _port_range_from(path: Path) -> tuple[int, int] | None: 

596 """The two integers in *path*, or ``None`` when it is absent or unreadable.""" 

597 try: 

598 low, high = path.read_text(encoding="utf-8").split()[:2] 

599 return int(low), int(high) 

600 except (OSError, ValueError): 

601 return None 

602 

603 

604def _ephemeral_range() -> tuple[int, int] | None: 

605 """The port range the kernel hands out for unbound sockets, if it says. 

606 

607 ``None`` when neither source answers, which is the signal to fall back to 

608 letting the OS choose. 

609 """ 

610 from_proc = _port_range_from(_PROC_PORT_RANGE) 

611 if from_proc is not None: 

612 return from_proc 

613 try: # macOS and the BSDs, which have no procfs entry for this 

614 out = subprocess.run( 

615 ["/usr/sbin/sysctl", "-n", "net.inet.ip.portrange.first", "net.inet.ip.portrange.last"], 

616 capture_output=True, 

617 text=True, 

618 encoding="utf-8", 

619 errors="replace", 

620 timeout=5, 

621 check=False, 

622 ) 

623 low, high = out.stdout.split()[:2] 

624 return int(low), int(high) 

625 except (OSError, ValueError, subprocess.SubprocessError): 

626 return None 

627 

628 

629# Where lilbee looks for engine ports when the kernel's ephemeral range is known. 

630# Above the registered-service crowd, below every default ephemeral range. 

631_PORT_SEARCH_FLOOR = 20000 

632_PORT_WINDOW_SPAN = 8192 

633# Block width. Each process searches one block, so concurrent lilbees hold 

634# disjoint ranges. A fleet takes one proxy port plus one per member, and embed 

635# and vision replicate per GPU, so 64 covers a 30-GPU host; a wider fleet spills 

636# into the next block. 

637_PORT_BLOCK = 64 

638 

639# Ports handed to a child that has not bound them yet. llama-swap binds a member 

640# port only on that member's first request, so the probe socket is long closed 

641# by then and the port looks free to every later probe. Without this the picker 

642# hands the next group exactly what it gave the last one, every time. 

643_reserved_ports: set[int] = set() 

644_reserved_lock = threading.Lock() 

645 

646 

647def release_reserved_ports(ports: Iterable[int]) -> None: 

648 """Give *ports* back to the picker, once nothing is expected to bind them.""" 

649 with _reserved_lock: 

650 _reserved_ports.difference_update(ports) 

651 

652 

653def _window_span(ceiling: tuple[int, int]) -> int: 

654 """How many ports below the ephemeral floor this host leaves to search.""" 

655 return max(1, min(ceiling[0], _PORT_SEARCH_FLOOR + _PORT_WINDOW_SPAN) - _PORT_SEARCH_FLOOR) 

656 

657 

658def _search_start(ceiling: tuple[int, int]) -> int: 

659 """First port of the block this process owns. 

660 

661 The pid selects a whole block, not an offset: reservation is per-process, and 

662 a fleet takes its ports contiguously, so pid-offset starts one apart overlap 

663 on all but one port. 

664 """ 

665 blocks = max(1, _window_span(ceiling) // _PORT_BLOCK) 

666 return _PORT_SEARCH_FLOOR + (os.getpid() % blocks) * _PORT_BLOCK 

667 

668 

669def _pick_free_ports(count: int) -> list[int]: 

670 """Bind *count* free localhost ports at once and return them. 

671 

672 All sockets stay open until every port is claimed so the OS cannot hand the 

673 same port out twice within one allocation. 

674 

675 Picked from below the kernel's ephemeral range rather than inside it. The 

676 gap between lilbee picking a port and llama-server binding it spans the whole 

677 lazy-spawn wait, and a port inside the ephemeral range can be handed to any 

678 passing outbound connection during that gap; one below it cannot be handed to 

679 anybody, so the only way to lose it is another server binding that exact port 

680 on purpose. Falls back to letting the OS choose when the range is unknown. 

681 """ 

682 ceiling = _ephemeral_range() 

683 sockets = [socket.socket(socket.AF_INET, socket.SOCK_STREAM) for _ in range(count)] 

684 try: 

685 for sock in sockets: 

686 _bind_below_ephemeral(sock, ceiling) 

687 return [int(sock.getsockname()[1]) for sock in sockets] 

688 finally: 

689 for sock in sockets: 

690 sock.close() 

691 

692 

693def _bind_below_ephemeral(sock: socket.socket, ceiling: tuple[int, int] | None) -> None: 

694 """Bind *sock* to a free, unreserved port under the ephemeral floor. 

695 

696 Falls back to letting the OS choose when the range is unknown or the window 

697 is used up, which keeps a fleet start working at the cost of returning to the 

698 ephemeral range for those ports. 

699 """ 

700 if ceiling is not None and ceiling[0] > _PORT_SEARCH_FLOOR: 

701 span = _window_span(ceiling) 

702 start = _search_start(ceiling) 

703 for offset in range(span): 

704 port = _PORT_SEARCH_FLOOR + (start - _PORT_SEARCH_FLOOR + offset) % span 

705 with _reserved_lock: 

706 if port in _reserved_ports: 

707 continue 

708 try: 

709 sock.bind((_HOST, port)) 

710 except OSError: 

711 continue 

712 _reserved_ports.add(port) 

713 return 

714 sock.bind((_HOST, 0)) 

715 

716 

717def _live_children(pid: int) -> list[psutil.Process]: 

718 """The process's current descendants, or none when it already exited.""" 

719 try: 

720 children: list[psutil.Process] = psutil.Process(pid).children(recursive=True) 

721 except psutil.NoSuchProcess: 

722 return [] 

723 return children 

724 

725 

726def _reap_survivors(children: list[psutil.Process]) -> None: 

727 """Terminate then kill any captured child that is still running.""" 

728 survivors = [child for child in children if child.is_running()] 

729 for child in survivors: 

730 with contextlib.suppress(psutil.NoSuchProcess): 

731 child.terminate() 

732 _, alive = psutil.wait_procs(survivors, timeout=_ORPHAN_STOP_TIMEOUT_S) 

733 for child in alive: 

734 with contextlib.suppress(psutil.NoSuchProcess): 

735 child.kill() 

736 _await_killed(alive) 

737 

738 

739def _await_killed(procs: list[psutil.Process]) -> None: 

740 """Wait for SIGKILLed processes to exit so their VRAM is free before any probe.""" 

741 if not procs: 

742 return 

743 _, alive = psutil.wait_procs(procs, timeout=_KILL_WAIT_TIMEOUT_S) 

744 for proc in alive: 

745 log.warning("Process %s survived SIGKILL; its VRAM may still be held.", proc.pid) 

746 

747 

748def _processes_named(needle: str) -> Iterator[psutil.Process]: 

749 """Live processes whose executable name contains *needle*. 

750 

751 ``name()`` is a cheap field (comm/proc_name); ``cmdline()`` reads the full 

752 argument vector and on macOS blocks on entitlement-protected binaries. So the 

753 name is the pre-filter and callers pay for ``cmdline()`` only on a match, 

754 which keeps a full-process-table scan from stalling on an unrelated process. 

755 """ 

756 for proc in psutil.process_iter(["name"]): 

757 # process_iter already skips processes that vanish mid-scan and, per its 

758 # ad_value contract, leaves ``name`` as None where it could not be read. 

759 name = proc.info["name"] or "" 

760 if needle in name: 

761 yield proc 

762 

763 

764def _swaps_for_config(config_path: Path) -> list[psutil.Process]: 

765 """Every live llama-swap (any owner) running against *config_path*. 

766 

767 Identity is the ``-config <path>`` argument, which every llama-swap this 

768 lilbee starts carries and which survives reparenting to init -- so this finds 

769 a leaked duplicate or a swap reparented away from us, neither of which a 

770 tracked Popen handle nor a ``children()`` scan would catch. 

771 """ 

772 target = str(config_path) 

773 swaps: list[psutil.Process] = [] 

774 for proc in _processes_named(_LLAMA_SWAP_PROCESS_NAME): 

775 try: 

776 cmdline = proc.cmdline() 

777 except ( 

778 psutil.NoSuchProcess, 

779 psutil.AccessDenied, 

780 psutil.ZombieProcess, 

781 OSError, 

782 SystemError, 

783 ): 

784 # OSError/SystemError: macOS psutil mishandles entitlement-protected 

785 # binaries (sysctl KERN_PROCARGS2), leaking a raw PermissionError or a 

786 # C-extension SystemError instead of an AccessDenied. 

787 continue 

788 # Identity is the -config path; _processes_named already gated on comm. 

789 if target in cmdline: 

790 swaps.append(proc) 

791 return swaps 

792 

793 

794def find_live_state(data_dir: Path, group: SwapGroup) -> SwapState | None: 

795 """The newest recorded state for *group* at *data_dir* (no liveness check). 

796 

797 A record's presence does not prove the engine is up; callers that need that 

798 probe it with ``state_is_healthy``. The name reflects that a record is written 

799 only for a running engine, not that this function verifies it. 

800 """ 

801 best: SwapState | None = None 

802 for state_path in sorted(data_dir.glob(_STATE_FILE_GLOB)): 

803 if f".{group.value}." not in f".{state_path.name}": 

804 continue 

805 state = _load_state(state_path) 

806 if state is None: 

807 continue 

808 if best is None or (state.created_at or 0) > (best.created_at or 0): 

809 best = state 

810 return best 

811 

812 

813def _running_endpoint_answers(base_url: str) -> bool: 

814 """Whether *base_url* serves llama-swap's ``/running`` endpoint (identity, not 

815 just liveness). 

816 

817 Proxy ports are ephemeral: after an engine dies, any unrelated local service 

818 that later binds the recorded port and returns a 2xx/3xx to an unknown path 

819 would pass a bare status check, so a dead record would look healthy forever and 

820 inference clients would bind to a non-engine endpoint. Requiring the ``running`` 

821 JSON payload shape that only llama-swap produces makes the probe identity-checked. 

822 Total: any transport error or non-conforming body reads as "not our engine". 

823 """ 

824 try: 

825 resp = _probe_client().get(f"{base_url}{_RUNNING_PATH}", timeout=_LIVENESS_TIMEOUT) 

826 except (OSError, httpx.HTTPError): 

827 return False 

828 if resp.status_code >= httpx.codes.BAD_REQUEST: 

829 return False 

830 try: 

831 return isinstance(resp.json().get(_KEY_RUNNING), list) 

832 except (ValueError, AttributeError): 

833 return False 

834 

835 

836def state_is_healthy(state: SwapState) -> bool: 

837 """Whether the engine behind *state* answers on its recorded proxy port.""" 

838 if state.proxy_port is None: 

839 return False 

840 return _running_endpoint_answers(f"http://{_HOST}:{state.proxy_port}") 

841 

842 

843def engine_record_exists(data_dir: Path) -> bool: 

844 """Whether any engine state file is present, without probing proxy health. 

845 

846 A filesystem fact, unlike a proxy HTTP probe: it is true for an engine that 

847 is live but momentarily unprobeable (fd exhaustion, host thrash), so the 

848 ladder can clear a recorded engine before building rather than double-build 

849 beside one an HTTP probe failed to see. 

850 """ 

851 return any(data_dir.glob(_STATE_FILE_GLOB)) 

852 

853 

854def stop_engine(data_dir: Path) -> list[str]: 

855 """Stop every engine the dir's state files record, regardless of liveness. 

856 

857 The unconditional off switch behind ``lilbee engine stop`` and the 

858 last-user-out path: each recorded swap is terminated through its state 

859 record (never a Popen handle, so it works on engines this process did 

860 not build) and its file removed. A record whose llama-swap is already dead 

861 still has its llama-servers (each in its own process group) reaped by 

862 recorded port, exactly as reap_stale does -- otherwise the off switch would 

863 leave those orphans holding VRAM and delete the ports needed to find them. 

864 Stale config files for dead owners are cleaned too, and the persistence 

865 opt-in is dropped with the engine it described, so the dir is left as 

866 clean as a reap leaves it. Unparseable files are left alone, as in 

867 reap_stale: they may be a sibling's in-flight write. Returns the group tokens 

868 whose engine was actually alive, so a caller reports only real stops. 

869 """ 

870 _clean_stale_configs(data_dir) 

871 # The persistence opt-in describes the engine instance being stopped, so it 

872 # dies with it. Cleared here rather than at each call site so no stop path 

873 # can leave a mark that makes the next engine sticky-warm. 

874 clear_keep_warm(data_dir) 

875 stopped: list[str] = [] 

876 for state_path in sorted(data_dir.glob(_STATE_FILE_GLOB)): 

877 state = _load_state(state_path) 

878 if state is None: 

879 continue 

880 if _stop_recorded_engine(state): 

881 group = _state_group(state_path.name) 

882 if group is not None: 

883 stopped.append(group) 

884 state_path.unlink(missing_ok=True) 

885 return stopped 

886 

887 

888def reap_stale(data_dir: Path) -> None: 

889 """Kill every dead or unhealthy recorded engine at *data_dir*. 

890 

891 An OOM-killed lilbee leaves llama-swap (and its servers) holding VRAM, 

892 so planning would otherwise see artificially reduced free memory; the 

893 ladder calls this before its GPU probe. Every state file is scanned 

894 (all groups, including legacy names): an engine that is alive AND 

895 answering on its proxy is spared regardless of who started it (a 

896 reload's own healthy groups, or a bindable engine the ladder skipped); 

897 everything else is stopped through its record and its file removed. An 

898 unparseable file is skipped, never deleted: it may be a sibling's 

899 in-flight write. When the swap itself is dead, its servers (each in 

900 its own process group) may still be alive holding VRAM; they are 

901 matched by name plus recorded member port and stopped before the file 

902 is removed. 

903 

904 Module-level (not a method) because it must run before planning decides 

905 which role groups exist, when no per-group manager has been built yet. 

906 """ 

907 _clean_stale_tmp_files(data_dir) 

908 _clean_stale_configs(data_dir) 

909 for state_path in sorted(data_dir.glob(_STATE_FILE_GLOB)): 

910 state = _load_state(state_path) 

911 if state is None: 

912 continue 

913 if state_is_healthy(state): 

914 # An answering engine is in use (bind accepts on exactly this 

915 # test); reaping must never disagree with binding. 

916 continue 

917 _stop_recorded_engine(state) 

918 state_path.unlink(missing_ok=True) 

919 

920 

921def _clean_stale_tmp_files(data_dir: Path) -> None: 

922 """Remove crash-leftover state and config tmp files whose writer is dead.""" 

923 tmp_glob = f"{_STATE_TMP_PREFIX}*{_STATE_TMP_SUFFIX}" 

924 for tmp_path in data_dir.glob(tmp_glob): 

925 writer_pid = _state_owner_pid(tmp_path.name) 

926 if writer_pid is not None and not psutil.pid_exists(writer_pid): 

927 tmp_path.unlink(missing_ok=True) 

928 

929 

930def _clean_stale_configs(data_dir: Path) -> None: 

931 """Remove per-owner config files whose owner lilbee is gone. 

932 

933 The swaps themselves are reaped from the state files; these leftover config 

934 files are just clutter once their writer pid is dead. A live owner's config 

935 (pid still exists) and a pid-less legacy name are left untouched; skipping on 

936 pid reuse only leaves harmless clutter, never deletes a live owner's config. 

937 """ 

938 for config_path in data_dir.glob(_CONFIG_FILE_GLOB): 

939 owner = _config_owner_pid(config_path.name) 

940 if owner is not None and not psutil.pid_exists(owner): 

941 config_path.unlink(missing_ok=True) 

942 

943 

944def _stop_own_fleet(config_path: Path, member_ports: tuple[int, ...]) -> None: 

945 """Stop every llama-swap this lilbee owns at *config_path* and reap upstreams. 

946 

947 Keyed on config-path identity rather than a tracked Popen or the live process 

948 tree: a warm-up/reload race can leave several llama-swap processes this lilbee 

949 started, any of which may be reparented to init, so no single handle or child 

950 scan finds them all. Every llama-swap running against our config is reaped: 

951 the build lock guarantees one builder per engine dir, so no sibling sparing 

952 applies. Each swap runs each llama-server in its own process group, so the 

953 upstreams are swept separately: captured descendants plus any llama-server 

954 still bound to one of our member ports (a respawned upstream the descendant 

955 snapshot missed), then confirmed gone. 

956 """ 

957 swaps = list(_swaps_for_config(config_path)) 

958 children: list[psutil.Process] = [] 

959 for swap in swaps: 

960 children.extend(_live_children(swap.pid)) 

961 for swap in swaps: 

962 if sys.platform == "win32": 

963 _hard_stop_proc(swap) 

964 else: 

965 _terminate_proc_group(swap) 

966 _reap_survivors(children + _find_orphan_servers(member_ports)) 

967 

968 

969def _terminate_proc_group(proc: psutil.Process) -> None: 

970 """SIGTERM a process's group, escalating to SIGKILL on timeout.""" 

971 try: 

972 pgid = os.getpgid(proc.pid) 

973 except (ProcessLookupError, OSError): # pragma: no cover - exited between checks 

974 return 

975 with contextlib.suppress(ProcessLookupError, PermissionError): 

976 os.killpg(pgid, signal.SIGTERM) 

977 try: 

978 proc.wait(timeout=_STOP_TIMEOUT_S) 

979 except psutil.TimeoutExpired: 

980 with contextlib.suppress(ProcessLookupError, PermissionError): 

981 os.killpg(pgid, _SIGKILL) 

982 _await_killed([proc]) 

983 

984 

985def _hard_stop_proc(proc: psutil.Process) -> None: 

986 """Terminate a process, escalating to a hard kill on timeout (Windows path).""" 

987 with contextlib.suppress(psutil.NoSuchProcess): 

988 proc.terminate() 

989 try: 

990 proc.wait(timeout=_STOP_TIMEOUT_S) 

991 except psutil.TimeoutExpired: 

992 with contextlib.suppress(psutil.NoSuchProcess): 

993 proc.kill() 

994 

995 

996def _load_state(path: Path) -> SwapState | None: 

997 """Parse a state file into a :class:`SwapState`; ``None`` when absent/corrupt.""" 

998 try: 

999 payload = json.loads(path.read_text(encoding="utf-8")) 

1000 raw_pgid = payload.get(_STATE_KEY_PGID) 

1001 raw_created = payload.get(_STATE_KEY_CREATED_AT) 

1002 raw_ports = payload.get(_STATE_KEY_MEMBER_PORTS) or [] 

1003 raw_proxy = payload.get(_STATE_KEY_PROXY_PORT) 

1004 return SwapState( 

1005 pid=int(payload[_STATE_KEY_PID]), 

1006 pgid=int(raw_pgid) if raw_pgid is not None else None, 

1007 created_at=float(raw_created) if raw_created is not None else None, 

1008 member_ports=tuple(int(port) for port in raw_ports), 

1009 proxy_port=int(raw_proxy) if raw_proxy is not None else None, 

1010 launches=tuple(payload.get(_STATE_KEY_LAUNCHES) or ()), 

1011 engine_pin=payload.get(_STATE_KEY_ENGINE_PIN), 

1012 ) 

1013 except (OSError, ValueError, KeyError, TypeError): 

1014 return None 

1015 

1016 

1017def _is_live_llama_swap(state: SwapState) -> bool: 

1018 """True when the recorded pid is alive and is the recorded llama-swap. 

1019 

1020 A recorded create time that differs from the live process's is pid reuse, 

1021 even when the recycled pid runs another instance's llama-swap; a legacy 

1022 state file without one falls back to the cmdline match alone. 

1023 """ 

1024 try: 

1025 proc = psutil.Process(state.pid) 

1026 cmdline = proc.cmdline() 

1027 create_time = proc.create_time() 

1028 except (psutil.NoSuchProcess, psutil.AccessDenied): 

1029 return False 

1030 if state.created_at is not None and abs(create_time - state.created_at) > ( 

1031 _CREATE_TIME_TOLERANCE_S 

1032 ): 

1033 return False 

1034 binary = Path(next(iter(cmdline), "")).name 

1035 return _LLAMA_SWAP_PROCESS_NAME in binary 

1036 

1037 

1038def _stop_stale_swap(state: SwapState) -> None: 

1039 """TERM-then-KILL a stale llama-swap's group and reap the servers it spawned. 

1040 

1041 Swept as wide as ``_stop_own_fleet``: a reparented or respawned server is no 

1042 longer a descendant, and every caller unlinks the record next, so the member 

1043 ports are the last thing that can match it. 

1044 """ 

1045 children = _live_children(state.pid) 

1046 try: 

1047 proc = psutil.Process(state.pid) 

1048 except psutil.NoSuchProcess: 

1049 proc = None 

1050 if proc is not None: 

1051 _signal_stale(state, signal.SIGTERM) 

1052 try: 

1053 proc.wait(timeout=_ORPHAN_STOP_TIMEOUT_S) 

1054 except psutil.TimeoutExpired: 

1055 _signal_stale(state, _SIGKILL) 

1056 _await_killed([proc]) 

1057 _reap_survivors(children + _find_orphan_servers(state.member_ports)) 

1058 

1059 

1060def _signal_stale(state: SwapState, sig: int) -> None: 

1061 """Signal the stale swap's process group, or the pid where groups don't apply.""" 

1062 if state.pgid is not None and sys.platform != "win32": 

1063 with contextlib.suppress(ProcessLookupError, PermissionError): 

1064 os.killpg(state.pgid, sig) 

1065 return 

1066 with contextlib.suppress(psutil.NoSuchProcess): 

1067 psutil.Process(state.pid).send_signal(sig) 

1068 

1069 

1070def _stop_recorded_engine(state: SwapState) -> bool: 

1071 """Terminate a live llama-swap and its servers, or reap the servers a dead one 

1072 orphaned (matched by recorded port, since they run in their own process groups 

1073 and outlive the swap). Returns whether anything was actually alive to stop, so 

1074 the off switch reports a stale record as "nothing stopped" rather than a false 

1075 success. 

1076 """ 

1077 if _is_live_llama_swap(state): 

1078 _stop_stale_swap(state) 

1079 return True 

1080 orphans = _find_orphan_servers(state.member_ports) 

1081 _reap_survivors(orphans) 

1082 return bool(orphans) 

1083 

1084 

1085def _find_orphan_servers(ports: tuple[int, ...]) -> list[psutil.Process]: 

1086 """Live llama-server processes serving one of *ports*. 

1087 

1088 Both the binary name and the ``--port`` value must match, so an unrelated 

1089 process on a recycled port is never killed; a server whose parent is a 

1090 live llama-swap belongs to a current run on a reused port and is spared. 

1091 """ 

1092 if not ports: 

1093 return [] 

1094 targets = {str(port) for port in ports} 

1095 orphans: list[psutil.Process] = [] 

1096 for proc in _processes_named(_LLAMA_SERVER_PROCESS_NAME): 

1097 try: 

1098 cmdline = proc.cmdline() 

1099 except (psutil.NoSuchProcess, psutil.AccessDenied): 

1100 continue 

1101 # Identity is the --port value plus the absence of a live swap parent; 

1102 # _processes_named already gated on the executable name (comm). 

1103 if _port_argument(cmdline) in targets and not _has_live_swap_parent(proc): 

1104 orphans.append(proc) 

1105 return orphans 

1106 

1107 

1108def _state_owner_pid(name: str) -> int | None: 

1109 """Owner pid embedded in a state or state-tmp filename, ``None`` when absent. 

1110 

1111 Handles both the group-qualified form (``llama-swap.state.chat.123.json``) 

1112 and the legacy pre-group form (``llama-swap.state.123.json``): the pid is 

1113 always the last dotted segment of the stem. 

1114 """ 

1115 stem = name.removeprefix(_STATE_TMP_PREFIX).removesuffix(_STATE_TMP_SUFFIX) 

1116 stem = stem.removeprefix(_STATE_FILENAME_PREFIX).removesuffix(_STATE_FILENAME_SUFFIX) 

1117 try: 

1118 return int(stem.rsplit(".", 1)[-1]) 

1119 except ValueError: 

1120 return None 

1121 

1122 

1123def _state_group(name: str) -> str | None: 

1124 """Group token from a group-qualified state filename, ``None`` for legacy names. 

1125 

1126 ``llama-swap.state.chat.123.json`` -> ``chat``; the legacy pre-group form 

1127 ``llama-swap.state.123.json`` has no group token. 

1128 """ 

1129 stem = name.removeprefix(_STATE_FILENAME_PREFIX).removesuffix(_STATE_FILENAME_SUFFIX) 

1130 head, _, _ = stem.rpartition(".") # drop the trailing pid; group is what remains 

1131 return head or None 

1132 

1133 

1134def _config_filename(pid: int, group: str) -> str: 

1135 """This owner's config filename for *group* (``llama-swap-<group>.<pid>.json``).""" 

1136 return _CONFIG_FILENAME_TEMPLATE.format(group=group, pid=pid) 

1137 

1138 

1139def _config_owner_pid(name: str) -> int | None: 

1140 """Owner pid embedded in a config filename, ``None`` for a legacy pid-less name.""" 

1141 stem = name.removeprefix("llama-swap-").removesuffix(".json") 

1142 try: 

1143 return int(stem.rsplit(".", 1)[-1]) 

1144 except ValueError: 

1145 return None 

1146 

1147 

1148def _has_live_swap_parent(proc: psutil.Process) -> bool: 

1149 """True when *proc*'s parent is a live llama-swap (the server is not orphaned).""" 

1150 try: 

1151 parent = proc.parent() 

1152 if parent is None: 

1153 return False 

1154 return _LLAMA_SWAP_PROCESS_NAME in parent.name() 

1155 except (psutil.NoSuchProcess, psutil.AccessDenied): 

1156 return False 

1157 

1158 

1159def _port_argument(cmdline: list[str]) -> str | None: 

1160 """The value following the port flag in *cmdline*, or ``None``.""" 

1161 for flag, value in itertools.pairwise(cmdline): 

1162 if flag == PORT_FLAG: 

1163 return value 

1164 return None