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
« 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."""
3from __future__ import annotations
5import contextlib
6import json
7import logging
8import os
9import shutil
10import socket
11import subprocess
12import sys
13import time
14from typing import IO
16import httpx
17import typer
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
30log = logging.getLogger(__name__)
32LOOPBACK = "127.0.0.1"
33"""Loopback address used for launcher-spawned sessions and the URLs we hand to clients."""
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
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
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])
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
82def _probe_health(port: int) -> dict[str, object] | None:
83 """GET ``/api/health`` once; return the parsed body on 200, else None.
85 The single place the probe URL, timeout, error handling, and status check
86 live, so the three public probes below stay consistent.
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 {}
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
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
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))
136def served_chat_ctx(port: int) -> int | None:
137 """The chat window ``/api/health`` reports, or None if unknown/unreachable.
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
149def planned_chat_ctx() -> int | None:
150 """The per-slot window the fleet will serve the configured chat model, or None.
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`).
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
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
186def client_chat_ctx(port: int) -> int | None:
187 """The chat window to advertise to a launched client, warning when it is small.
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
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)))
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.
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)
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
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.
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.
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.
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 )
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
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())}
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()
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)
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.
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.
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)
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
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