Coverage for src/lilbee/cli/launchers/server.py: 100%

199 statements  

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

1"""Server-lifecycle helpers shared by every ``lilbee launch <client>`` command.""" 

2 

3from __future__ import annotations 

4 

5import contextlib 

6import json 

7import logging 

8import os 

9import shutil 

10import socket 

11import subprocess 

12import sys 

13import time 

14from typing import IO 

15 

16import httpx 

17import typer 

18 

19from lilbee.cli.app import console 

20from lilbee.cli.commands.servers import port_file 

21from lilbee.cli.launchers.warm_render import render_warm 

22from lilbee.core.config import cfg 

23from lilbee.modelhub.registry import ModelRegistry 

24from lilbee.parent_monitor import PARENT_PID_ENV 

25from lilbee.providers.fleet.child_guard import spawn_bound_child 

26from lilbee.providers.fleet.swap_config import cold_load_timeout_s 

27from lilbee.providers.roles import WorkerRole 

28from lilbee.server.auth import server_json_path 

29 

30log = logging.getLogger(__name__) 

31 

32LOOPBACK = "127.0.0.1" 

33"""Loopback address used for launcher-spawned sessions and the URLs we hand to clients.""" 

34 

35_SERVER_BOOT_TIMEOUT_S = 60.0 

36_SERVER_POLL_INTERVAL_S = 0.5 

37# Floor on the cold model-load wait; chat_warm_budget_s() scales it up with the weights. 

38_WARM_TIMEOUT_S = 600.0 

39_HEALTH_PROBE_TIMEOUT_S = 2.0 

40_HTTP_OK = 200 

41_HEALTH_PATH = "/api/health" 

42_TERMINATE_GRACE_S = 10 

43_KILL_GRACE_S = 5 

44# Spawn attempts; free_port()'s released probe port can be stolen before the server binds. 

45_SPAWN_ATTEMPTS = 3 

46 

47 

48def running_server_session() -> tuple[str, int] | None: 

49 """Return ``(token, port)`` for a server already running on this machine, else None.""" 

50 session_path = server_json_path() 

51 port_path = port_file() 

52 if not session_path.exists() or not port_path.exists(): 

53 return None 

54 try: 

55 data = json.loads(session_path.read_text(encoding="utf-8")) 

56 token = data.get("token") 

57 port = int(port_path.read_text(encoding="utf-8").strip()) 

58 except (json.JSONDecodeError, OSError, ValueError): 

59 return None 

60 if not isinstance(token, str) or not token: 

61 return None 

62 return token, port 

63 

64 

65def free_port() -> int: 

66 """Return an unused TCP port on the loopback interface.""" 

67 with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: 

68 s.bind((LOOPBACK, 0)) 

69 return int(s.getsockname()[1]) 

70 

71 

72def _session_token() -> str | None: 

73 """The bearer token from server.json, or None if it is not readable yet.""" 

74 try: 

75 data = json.loads(server_json_path().read_text(encoding="utf-8")) 

76 except (json.JSONDecodeError, UnicodeDecodeError, OSError): 

77 return None 

78 token = data.get("token") 

79 return token if isinstance(token, str) and token else None 

80 

81 

82def _probe_health(port: int) -> dict[str, object] | None: 

83 """GET ``/api/health`` once; return the parsed body on 200, else None. 

84 

85 The single place the probe URL, timeout, error handling, and status check 

86 live, so the three public probes below stay consistent. 

87 

88 Health needs the token like every other route: it reports the chat 

89 engine's last error, which carries model paths and loader failures. The 

90 token is re-read per attempt rather than captured once, because these 

91 probes poll a server that is still starting and server.json does not exist 

92 until its lifespan has run. No token yet means no server yet, which is the 

93 same answer a refused connection gives. 

94 """ 

95 token = _session_token() 

96 if token is None: 

97 return None 

98 try: 

99 resp = httpx.get( 

100 f"http://{LOOPBACK}:{port}{_HEALTH_PATH}", 

101 timeout=_HEALTH_PROBE_TIMEOUT_S, 

102 headers={"Authorization": f"Bearer {token}"}, 

103 ) 

104 except httpx.HTTPError: 

105 return None 

106 if resp.status_code != _HTTP_OK: 

107 return None 

108 try: 

109 body = resp.json() 

110 except ValueError: 

111 return {} 

112 return body if isinstance(body, dict) else {} 

113 

114 

115def health_ok(port: int) -> bool: 

116 """Single-shot ``/api/health`` probe; True iff a 200 comes back fast.""" 

117 return _probe_health(port) is not None 

118 

119 

120def wait_for_health(port: int, timeout_s: float = _SERVER_BOOT_TIMEOUT_S) -> bool: 

121 """Poll ``/api/health`` until it answers 200 or *timeout_s* elapses.""" 

122 deadline = time.monotonic() + timeout_s 

123 while time.monotonic() < deadline: 

124 if health_ok(port): 

125 return True 

126 time.sleep(_SERVER_POLL_INTERVAL_S) 

127 return False 

128 

129 

130def chat_ready(port: int) -> bool: 

131 """Single-shot probe: True iff ``/api/health`` reports the chat engine warm.""" 

132 body = _probe_health(port) 

133 return bool(body and body.get("chat_ready", False)) 

134 

135 

136def served_chat_ctx(port: int) -> int | None: 

137 """The chat window ``/api/health`` reports, or None if unknown/unreachable. 

138 

139 A launcher passes this to the client so it trims history to the model's 

140 actual window instead of overflowing on a long agentic session. 

141 """ 

142 body = _probe_health(port) 

143 if body is None: 

144 return None 

145 ctx = body.get("chat_ctx") 

146 return ctx if isinstance(ctx, int) and ctx > 0 else None 

147 

148 

149def planned_chat_ctx() -> int | None: 

150 """The per-slot window the fleet will serve the configured chat model, or None. 

151 

152 Mirrors the fleet's own single-GPU chat sizing, so it answers before the 

153 engine is up: the same ``cfg.num_ctx`` short-circuit, then the same 

154 :func:`resolve_chat_ctx` against the same budget the fleet sizes with, which 

155 is the memory the GPU reports rather than the host's (see 

156 :func:`lilbee.providers.fleet.planning.plan_sizing_budget`). 

157 

158 A tensor-split chat is sized by the fleet against per-device headroom 

159 instead, so this can over-report there; it is only a fallback for a chat 

160 engine that is not up yet, and the served window wins once it is. A 

161 remote-served chat model has no local window to compute. 

162 """ 

163 from lilbee.providers.base import ProviderError 

164 from lilbee.providers.engine_params import resolve_chat_ctx, resolve_model_path 

165 from lilbee.providers.fleet.planning import plan_sizing_budget 

166 from lilbee.providers.gguf_meta import read_gguf_metadata 

167 from lilbee.providers.model_ref import parse_model_ref 

168 

169 ref = str(cfg.chat_model) 

170 if not ref or not parse_model_ref(ref).is_local: 

171 return None 

172 if cfg.num_ctx is not None: 

173 return cfg.num_ctx 

174 try: 

175 path = resolve_model_path(ref) 

176 return resolve_chat_ctx( 

177 path, read_gguf_metadata(path), available_bytes=plan_sizing_budget() 

178 ) 

179 except (ProviderError, OSError, ValueError): 

180 # Sizing needs the model file and its GGUF header; an absent or unreadable 

181 # one leaves the window unknown rather than failing the launch. 

182 log.debug("planned_chat_ctx failed for %s", ref, exc_info=True) 

183 return None 

184 

185 

186def client_chat_ctx(port: int) -> int | None: 

187 """The chat window to advertise to a launched client, warning when it is small. 

188 

189 The chat role builds lazily, so a launcher that hands off before the engine 

190 is warm gets nothing from ``/api/health``; fall back to the window the fleet 

191 plans to serve rather than leaving the client with no window at all. A window 

192 below ``cfg.chat_n_ctx_target`` means the host could not back what was asked 

193 for, which changes how much history an agent can keep, so say so. 

194 """ 

195 ctx = served_chat_ctx(port) 

196 if ctx is None: 

197 ctx = planned_chat_ctx() 

198 if ctx is not None and ctx < cfg.chat_n_ctx_target: 

199 typer.secho( 

200 f"Warning: the chat model is served with a {ctx:,}-token context, below the " 

201 f"configured chat_n_ctx_target of {cfg.chat_n_ctx_target:,}. Either the model " 

202 "was trained on a smaller window, or its weights leave too little of the " 

203 "memory budget for the KV cache. A longer-context model, a smaller " 

204 "quantization, or a higher gpu_memory_fraction raises it.", 

205 err=True, 

206 fg=typer.colors.YELLOW, 

207 ) 

208 return ctx 

209 

210 

211def chat_warm_budget_s() -> float: 

212 """Warm wait scaled to the chat model's on-disk weights at the engine's cold-load rate.""" 

213 try: 

214 shards = ModelRegistry(cfg.models_dir).shard_paths(str(cfg.chat_model)) 

215 except (KeyError, ValueError): 

216 return _WARM_TIMEOUT_S 

217 total_bytes = sum(shard.stat().st_size for shard in shards) 

218 return max(_WARM_TIMEOUT_S, float(cold_load_timeout_s(total_bytes, WorkerRole.CHAT))) 

219 

220 

221def wait_for_chat_warm(port: int, timeout_s: float | None = None) -> bool: 

222 """Block until the chat model is loaded, showing granular warm progress. 

223 

224 The server warms the chat role on a background thread at startup, so a client 

225 launched the instant the HTTP port binds would otherwise hit an 

226 apparently-dead stream during the cold model load. Streams ``/api/warm/stream`` 

227 to render a real read-phase byte bar then an engine-load spinner; falls back 

228 to a plain readiness poll when that stream can't be opened. 

229 Returns True once the chat engine reports ready, or False if the budget 

230 (weights-scaled via :func:`chat_warm_budget_s` unless given) elapses first; 

231 the caller proceeds either way, so a still-loading model just warms on the 

232 first call. 

233 """ 

234 if timeout_s is None: 

235 timeout_s = chat_warm_budget_s() 

236 if chat_ready(port): 

237 return True 

238 streamed = render_warm(f"http://{LOOPBACK}:{port}", timeout_s) 

239 if streamed is not None: 

240 # The stream ran (ready, error, or its own timeout); don't double-spend 

241 # the budget on a second poll. The caller proceeds on False regardless. 

242 return streamed 

243 return _poll_chat_ready(port, timeout_s) 

244 

245 

246def _poll_chat_ready(port: int, timeout_s: float) -> bool: 

247 """Fallback warm wait when the progress stream is unavailable: poll readiness.""" 

248 deadline = time.monotonic() + timeout_s 

249 with console.status("Warming the chat model..."): 

250 while time.monotonic() < deadline: 

251 if chat_ready(port): 

252 return True 

253 time.sleep(_SERVER_POLL_INTERVAL_S) 

254 return False 

255 

256 

257def spawn_server( 

258 port: int, *, env_overrides: dict[str, str] | None = None 

259) -> subprocess.Popen[bytes]: 

260 """Spawn ``lilbee serve --port <port>`` as a background subprocess. 

261 

262 Prefers the ``lilbee`` binary on PATH so frozen builds (Nuitka standalone) 

263 spawn the binary directly. Falls back to ``sys.executable -m lilbee`` for 

264 pip / editable installs where the entry point shims to the same form. 

265 

266 ``env_overrides`` are layered onto the inherited environment for the child 

267 (e.g. ``LILBEE_CHAT_N_CTX_TARGET`` to size the served window for a launched 

268 agent); ``None`` inherits the parent environment unchanged. 

269 

270 Stdout/stderr go to ``cfg.data_dir / "logs" / "launcher-serve.log"`` (size 

271 capped at 5 MB) so a crash mid-session leaves a trace instead of disappearing. 

272 Set ``LILBEE_LAUNCHER_SERVE_QUIET=1`` to restore the previous DEVNULL behavior. 

273 """ 

274 lilbee_bin = shutil.which("lilbee") 

275 # On Windows, pip/uv may install a ``lilbee.cmd`` wrapper instead of a bare 

276 # executable. Popen(shell=False) raises PermissionError on .cmd files, so 

277 # fall through to the sys.executable -m lilbee form in that case. 

278 _bin_is_cmd = sys.platform == "win32" and ( 

279 lilbee_bin is not None and lilbee_bin.lower().endswith(".cmd") 

280 ) 

281 cmd = ( 

282 [lilbee_bin, "serve", "--port", str(port)] 

283 if lilbee_bin is not None and not _bin_is_cmd 

284 else [sys.executable, "-m", "lilbee", "serve", "--port", str(port)] 

285 ) 

286 

287 log_file: IO[bytes] | None = None 

288 if os.environ.get("LILBEE_LAUNCHER_SERVE_QUIET"): 

289 stdout: int | IO[bytes] = subprocess.DEVNULL 

290 stderr: int | IO[bytes] = subprocess.DEVNULL 

291 else: 

292 log_dir = cfg.data_dir / "logs" 

293 log_dir.mkdir(parents=True, exist_ok=True) 

294 log_path = log_dir / "launcher-serve.log" 

295 # Truncate when the file passes 5 MB so a long-lived session doesn't 

296 # accumulate the chat-completion firehose into the data dir indefinitely. 

297 # On Windows the file may still be held open by a previous session, so 

298 # fall through to append mode when unlink is denied. 

299 if log_path.exists() and log_path.stat().st_size > 5 * 1024 * 1024: 

300 with contextlib.suppress(OSError): 

301 log_path.unlink() 

302 log_file = log_path.open("ab") 

303 stdout = log_file 

304 stderr = subprocess.STDOUT 

305 

306 # LILBEE_PARENT_PID arms serve's parent-death watcher, so a hard-killed 

307 # launcher (whose finally never runs) does not orphan serve holding server_lock. 

308 child_env = {**os.environ, **(env_overrides or {}), PARENT_PID_ENV: str(os.getpid())} 

309 

310 try: 

311 return spawn_bound_child( 

312 cmd, 

313 stdout=stdout, 

314 stderr=stderr, 

315 env=child_env, 

316 ) 

317 finally: 

318 # Popen dups the fd into the child; the parent's handle is no longer 

319 # needed and would otherwise leak for the launcher's whole lifetime. 

320 if log_file is not None: 

321 log_file.close() 

322 

323 

324def stop_spawned_server(proc: subprocess.Popen[bytes]) -> None: 

325 """Terminate *proc* gracefully, escalating to kill if it ignores SIGTERM.""" 

326 if proc.poll() is not None: 

327 return 

328 proc.terminate() 

329 try: 

330 proc.wait(timeout=_TERMINATE_GRACE_S) 

331 except subprocess.TimeoutExpired: 

332 proc.kill() 

333 proc.wait(timeout=_KILL_GRACE_S) 

334 

335 

336def ensure_server_running( 

337 *, env_overrides: dict[str, str] | None = None 

338) -> tuple[tuple[str, int], subprocess.Popen[bytes] | None]: 

339 """Return ``(session, spawned_proc)`` for a usable lilbee server. 

340 

341 Reuses an already-running server when its session files are healthy. 

342 Otherwise spawns a fresh server on a free port. The returned ``spawned_proc`` 

343 is ``None`` when an existing server was reused; the caller is responsible 

344 for stopping a spawned process when it is done with it. 

345 

346 ``env_overrides`` reach a freshly spawned child (e.g. a launcher sizing the 

347 served window); a reused server keeps whatever window it booted with. 

348 """ 

349 existing = running_server_session() 

350 if existing is not None and health_ok(existing[1]): 

351 return existing, None 

352 last_port = 0 

353 for _ in range(_SPAWN_ATTEMPTS): 

354 # Honor a user-pinned port so a persisted agent config keeps a valid URL; 

355 # fall back to a free port when unset (0). 

356 last_port = cfg.server_port or free_port() 

357 spawned = _spawn_and_wait(last_port, env_overrides=env_overrides) 

358 if spawned is not None: 

359 return _session_for_spawned(spawned), spawned 

360 typer.secho( 

361 f"lilbee server failed to start on port {last_port}; check the logs.", 

362 err=True, 

363 fg=typer.colors.RED, 

364 ) 

365 raise typer.Exit(1) 

366 

367 

368def _spawn_and_wait( 

369 port: int, *, env_overrides: dict[str, str] | None = None 

370) -> subprocess.Popen[bytes] | None: 

371 """Spawn a server on *port* and wait for health; None when it never comes up.""" 

372 spawned = spawn_server(port, env_overrides=env_overrides) 

373 with console.status(f"Starting lilbee server on port {port}..."): 

374 healthy = wait_for_health(port) 

375 if healthy: 

376 return spawned 

377 stop_spawned_server(spawned) 

378 return None 

379 

380 

381def _session_for_spawned(spawned: subprocess.Popen[bytes]) -> tuple[str, int]: 

382 """Read the session a freshly-healthy server wrote, stopping it when missing.""" 

383 session = running_server_session() 

384 if session is None: 

385 stop_spawned_server(spawned) 

386 typer.secho( 

387 "lilbee server started but did not write a session file; cannot continue.", 

388 err=True, 

389 fg=typer.colors.RED, 

390 ) 

391 raise typer.Exit(1) 

392 return session