Coverage for src/lilbee/server/handlers/__init__.py: 100%

120 statements  

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

1"""Framework-agnostic route handlers for the lilbee HTTP server. 

2 

3Every public function is a plain async callable; no framework imports. 

4Return types are dicts (JSON responses), lists, or async generators of SSE strings. 

5 

6Handlers are grouped by concern (sse, rag, models, ingest, config, documents, 

7crawl) under sibling submodules. The names re-exported below are the public 

8API consumed by ``server/routes/*.py``. 

9""" 

10 

11from __future__ import annotations 

12 

13import asyncio 

14import dataclasses 

15import logging 

16import time 

17from collections.abc import AsyncGenerator, Callable, Sequence 

18from typing import TYPE_CHECKING, Literal 

19 

20from lilbee.app.services import get_services 

21from lilbee.app.status import gather_status 

22from lilbee.app.version import get_version 

23from lilbee.core.config import cfg 

24from lilbee.providers.roles import WorkerRole 

25from lilbee.providers.warm_progress import WarmPhase, WarmProgress, is_active_warm 

26from lilbee.runtime.progress import SseEvent 

27from lilbee.server.handlers.agent_config import agent_config, agent_config_index 

28from lilbee.server.handlers.config import ( 

29 get_config, 

30 get_config_defaults, 

31 get_config_schema, 

32 update_config, 

33) 

34from lilbee.server.handlers.crawl import crawl_stream 

35from lilbee.server.handlers.documents import ( 

36 delete_documents, 

37 get_source_content, 

38 list_documents, 

39) 

40from lilbee.server.handlers.ingest import ( 

41 add_files_stream, 

42 add_uploads_stream, 

43 import_stream, 

44 sync_stream, 

45 validate_add_paths, 

46 validate_upload_names, 

47) 

48from lilbee.server.handlers.models import ( 

49 TASK_ENDPOINT_PATH, 

50 ModelCatalogSection, 

51 ModelsResponse, 

52 enforce_pull_arch_compat, 

53 format_task_mismatch, 

54 list_external_models, 

55 list_models, 

56 models_catalog, 

57 models_delete, 

58 models_installed, 

59 models_pull, 

60 models_show, 

61 set_chat_model, 

62 set_embedding_model, 

63 set_reranker_model, 

64 set_vision_model, 

65) 

66from lilbee.server.handlers.rag import ( 

67 ask, 

68 ask_stream, 

69 chat, 

70 chat_stream, 

71 search, 

72) 

73from lilbee.server.handlers.sse import ( 

74 SseStream, 

75 classify_load_error, 

76 sse_done, 

77 sse_error, 

78 sse_event, 

79) 

80from lilbee.server.handlers.wiki import ( 

81 wiki_build_stream, 

82 wiki_generate_stream, 

83 wiki_synthesize_stream, 

84) 

85from lilbee.server.models import ( 

86 GpusResponse, 

87 HealthResponse, 

88 PlacementResponse, 

89 ShutdownResponse, 

90 StatusResponse, 

91) 

92 

93if TYPE_CHECKING: 

94 from lilbee.app.placement import GpuInfo, PlacementView 

95 from lilbee.providers.base import LLMProvider 

96 

97log = logging.getLogger(__name__) 

98 

99# How often the warm stream re-snapshots provider state; sub-second so the read 

100# bar advances smoothly without busy-spinning. 

101_WARM_POLL_INTERVAL_S = 0.25 

102# Upper bound on the warm stream; the launcher hands off when this elapses, so a 

103# model still loading past it just warms on the client's first call. Generous to 

104# cover a cold tensor-split giant off a slow filesystem. 

105_WARM_STREAM_TIMEOUT_S = 1800.0 

106 

107 

108def _chat_status( 

109 provider: LLMProvider, 

110) -> tuple[Literal["ready", "loading", "not_started", "error"], str | None]: 

111 """Classify the chat engine's readiness for /api/health, with the error reason. 

112 

113 ``ready`` once the role serves; ``error`` when warm-up failed (paired with the 

114 warm tracker's reason); ``loading`` while a warm is in flight; ``not_started`` 

115 when nothing is warming and the role isn't up (no chat model planned, or chat 

116 is swapped out for its co-tenant; the next chat request loads it). 

117 """ 

118 if provider.role_ready(WorkerRole.CHAT): 

119 return "ready", None 

120 snapshot = provider.warm_progress() 

121 if snapshot is not None and snapshot.phase is WarmPhase.ERROR: 

122 return "error", snapshot.error 

123 if is_active_warm(snapshot): 

124 return "loading", None 

125 return "not_started", None 

126 

127 

128async def health() -> HealthResponse: 

129 """Return service health, version, and whether the chat engine is warm.""" 

130 services = get_services() 

131 provider = services.provider 

132 chat_status, chat_error = _chat_status(provider) 

133 prefill = provider.chat_prefill_progress() 

134 return HealthResponse( 

135 status="ok", 

136 version=get_version(), 

137 chat_ready=provider.role_ready(WorkerRole.CHAT), 

138 chat_status=chat_status, 

139 chat_error=chat_error, 

140 chat_ctx=provider.served_chat_ctx(), 

141 chat_slots=provider.served_chat_slots(), 

142 chat_prefill_processed=prefill[0] if prefill else None, 

143 chat_prefill_total=prefill[1] if prefill else None, 

144 embed_token_cap=provider.embed_token_cap(), 

145 warnings=[*services.store.health_warnings(), *provider.health_warnings()], 

146 ) 

147 

148 

149async def shutdown() -> ShutdownResponse: 

150 """Accept an API-requested stop; the route's background task stops the server. 

151 

152 Litestar runs that task only after the response has been handed to the 

153 transport, so the stop cannot beat the 202 out and no wall-clock delay 

154 has to be guessed. Stopping through the serving loop keeps teardown 

155 ordered; without a loop the task falls back to SIGTERM. 

156 """ 

157 log.info("Shutdown requested via the API") 

158 return ShutdownResponse(status="shutting_down") 

159 

160 

161async def warm_stream() -> AsyncGenerator[str, None]: 

162 """Stream chat-model cold-load progress as SSE until the engine is ready. 

163 

164 A launcher subscribes to render granular warm feedback. Each 

165 :data:`SseEvent.WARM` event carries a :class:`WarmProgress` snapshot; a 

166 terminal :data:`SseEvent.DONE` closes the stream once the chat role is ready 

167 or has failed, or when the budget elapses (the caller proceeds either way, so 

168 a still-loading model just warms on its first call). When nothing is loading 

169 because the engine is already warm, a single ready snapshot is emitted. 

170 """ 

171 provider = get_services().provider 

172 deadline = time.monotonic() + _WARM_STREAM_TIMEOUT_S 

173 while time.monotonic() < deadline: 

174 snapshot = provider.warm_progress() 

175 if snapshot is None: 

176 if provider.role_ready(WorkerRole.CHAT): 

177 yield sse_event(SseEvent.WARM, WarmProgress(phase=WarmPhase.READY).model_dump()) 

178 break 

179 yield sse_event(SseEvent.WARM, WarmProgress(phase=WarmPhase.STARTING).model_dump()) 

180 else: 

181 yield sse_event(SseEvent.WARM, snapshot.model_dump()) 

182 # A stale READY after eviction must not short-circuit; wait for the role. 

183 if snapshot.phase is WarmPhase.READY and provider.role_ready(WorkerRole.CHAT): 

184 break 

185 if snapshot.phase is WarmPhase.ERROR: 

186 break 

187 await asyncio.sleep(_WARM_POLL_INTERVAL_S) 

188 yield sse_done({}) 

189 

190 

191_GPU_STATS_INTERVAL_S = 1.0 

192 

193 

194async def gpu_stats_stream( 

195 devices: Sequence[GpuInfo], 

196 interval_s: float = _GPU_STATS_INTERVAL_S, 

197 max_ticks: int | None = None, 

198) -> AsyncGenerator[str, None]: 

199 """Stream live per-GPU utilization + free memory as SSE for the placement view. 

200 

201 Devices are resolved by the caller before the stream starts so a ProviderError 

202 surfaces as a 503 at route time, not mid-stream. The client keeps the stream 

203 open while visible; ``max_ticks`` bounds it for tests. A heartbeat is emitted 

204 every ``cfg.sse_heartbeat_interval`` seconds of idle so clients don't time out. 

205 

206 The per-vendor probe runs on a worker thread, not here. It is not light: every 

207 backend shells out to an SMI tool with a five-second timeout, and the Intel 

208 paths sleep and scan /proc on top of that. Driven inline it held the event 

209 loop for the whole subprocess on every tick, once per connected client, which 

210 stalls chat, search and embedding requests along with it. 

211 """ 

212 from lilbee.cli.tui import messages as msg 

213 from lilbee.providers.fleet.gpu_stats import intel_util_hint, probe_gpu_stats_shared 

214 

215 last_heartbeat = time.monotonic() 

216 tick = 0 

217 while max_ticks is None or tick < max_ticks: 

218 stats = await asyncio.to_thread(probe_gpu_stats_shared, devices) 

219 payload: dict[str, object] = {"gpus": [dataclasses.asdict(s) for s in stats.values()]} 

220 hint = intel_util_hint(devices, stats) 

221 if hint: 

222 payload["notice"] = msg.intel_util_hint_text(hint) 

223 yield sse_event(SseEvent.GPU_STATS, payload) 

224 tick += 1 

225 if max_ticks is None or tick < max_ticks: 

226 await asyncio.sleep(interval_s) 

227 now = time.monotonic() 

228 heartbeat_interval = cfg.sse_heartbeat_interval 

229 if heartbeat_interval > 0 and now - last_heartbeat >= heartbeat_interval: 

230 last_heartbeat = now 

231 yield sse_event(SseEvent.HEARTBEAT, {"ts": time.time()}) 

232 

233 

234async def status() -> StatusResponse: 

235 """Return config, sources, and chunk counts.""" 

236 raw = gather_status() 

237 return StatusResponse(**raw.model_dump(exclude_none=True)) 

238 

239 

240async def placement() -> PlacementResponse: 

241 """Current effective placement.""" 

242 from lilbee.app.placement import get_placement 

243 

244 return await _placement_response_off_loop(get_placement) 

245 

246 

247async def placement_preview(spec_json: str | None) -> PlacementResponse: 

248 """Preview a candidate spec (or auto when no spec). No persistence.""" 

249 from lilbee.app.placement import preview_placement 

250 from lilbee.providers.fleet.placement_spec import PlacementSpec 

251 

252 spec = PlacementSpec.from_json(spec_json) if spec_json else None 

253 return await _placement_response_off_loop(lambda: preview_placement(spec)) 

254 

255 

256async def placement_set(spec_json: str) -> PlacementResponse: 

257 """Apply a manual placement spec; persists and rebuilds the fleet.""" 

258 from lilbee.app.placement import set_placement 

259 from lilbee.providers.fleet.placement_spec import PlacementSpec 

260 

261 spec = PlacementSpec.from_json(spec_json) 

262 return await _placement_response_off_loop(lambda: set_placement(spec)) 

263 

264 

265async def placement_clear() -> PlacementResponse: 

266 """Clear manual placement; returns to the auto planner and rebuilds the fleet.""" 

267 from lilbee.app.placement import set_placement 

268 

269 return await _placement_response_off_loop(lambda: set_placement(None)) 

270 

271 

272async def _placement_response_off_loop(action: Callable[[], PlacementView]) -> PlacementResponse: 

273 """Run a placement action and serialize it off the event loop. 

274 

275 Placement actions and the Intel util notice both shell out to GPU probes, 

276 so neither may run on the loop. 

277 """ 

278 return await asyncio.to_thread(lambda: _placement_response(action())) 

279 

280 

281def _placement_response(view: PlacementView) -> PlacementResponse: 

282 """Serialize a placement view with the host-level Intel util notice attached.""" 

283 resp = PlacementResponse.from_view(view) 

284 resp.notice = _intel_notice_text(view.gpus) 

285 return resp 

286 

287 

288def _intel_notice_text(devices: Sequence[GpuInfo]) -> str | None: 

289 """Formatted Intel util fix for the JSON surfaces, or None when util reads fine.""" 

290 from lilbee.cli.tui import messages as msg 

291 from lilbee.providers.fleet.gpu_stats import probe_intel_util_hint 

292 

293 hint = probe_intel_util_hint(devices) 

294 return msg.intel_util_hint_text(hint) if hint else None 

295 

296 

297async def gpus() -> GpusResponse: 

298 """Detected GPUs with free/total VRAM, plus the host-level Intel util notice.""" 

299 from lilbee.app.placement import get_placement 

300 

301 def _body() -> GpusResponse: 

302 view = get_placement() 

303 return GpusResponse( 

304 gpus=PlacementResponse.from_view(view).gpus, 

305 notice=_intel_notice_text(view.gpus), 

306 ) 

307 

308 return await asyncio.to_thread(_body) 

309 

310 

311__all__ = [ 

312 "TASK_ENDPOINT_PATH", 

313 "ModelCatalogSection", 

314 "ModelsResponse", 

315 "SseStream", 

316 "add_files_stream", 

317 "add_uploads_stream", 

318 "agent_config", 

319 "agent_config_index", 

320 "ask", 

321 "ask_stream", 

322 "chat", 

323 "chat_stream", 

324 "classify_load_error", 

325 "crawl_stream", 

326 "delete_documents", 

327 "enforce_pull_arch_compat", 

328 "format_task_mismatch", 

329 "get_config", 

330 "get_config_defaults", 

331 "get_config_schema", 

332 "get_source_content", 

333 "gpu_stats_stream", 

334 "gpus", 

335 "health", 

336 "import_stream", 

337 "list_documents", 

338 "list_external_models", 

339 "list_models", 

340 "models_catalog", 

341 "models_delete", 

342 "models_installed", 

343 "models_pull", 

344 "models_show", 

345 "placement", 

346 "placement_clear", 

347 "placement_preview", 

348 "placement_set", 

349 "search", 

350 "set_chat_model", 

351 "set_embedding_model", 

352 "set_reranker_model", 

353 "set_vision_model", 

354 "sse_done", 

355 "sse_error", 

356 "sse_event", 

357 "status", 

358 "sync_stream", 

359 "update_config", 

360 "validate_add_paths", 

361 "validate_upload_names", 

362 "warm_stream", 

363 "wiki_build_stream", 

364 "wiki_generate_stream", 

365 "wiki_synthesize_stream", 

366]