Coverage for src/lilbee/app/services.py: 100%

232 statements  

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

1"""Typed service container: single point of access for all singletons. 

2 

3All runtime dependencies (provider, store, embedder, reranker, concepts, 

4clusterer, searcher, worker pool) are created lazily on first call to 

5``get_services()`` and cached for the process lifetime. Tests call 

6``reset_services()`` between runs. 

7 

8``build_services(config)`` is the construction seam: it builds a full container 

9against an arbitrary Config without touching the process-global singleton. The 

10library API (:class:`lilbee.Lilbee`) builds one per instance and installs it for 

11the duration of each call via :func:`services_scope`, so ingest code reaching for 

12``get_services()`` resolves the caller's container. The override is a ContextVar, 

13so ``reset_services`` / ``set_services`` / ``peek_services`` (which operate only 

14on the global singleton) never see it. 

15""" 

16 

17from __future__ import annotations 

18 

19import asyncio 

20import atexit 

21import logging 

22import os 

23import signal 

24import sys 

25import threading 

26import time 

27from collections.abc import Callable 

28from contextlib import contextmanager 

29from contextvars import ContextVar 

30from dataclasses import dataclass, field 

31from typing import TYPE_CHECKING 

32 

33if TYPE_CHECKING: 

34 from collections.abc import Iterator 

35 

36 from lilbee.catalog.hf_client import HfClient 

37 from lilbee.core.config import Config 

38 from lilbee.data.store import Store 

39 from lilbee.modelhub.model_manager import ModelManager 

40 from lilbee.modelhub.model_manager.discovery import KnownModelCache 

41 from lilbee.modelhub.registry import ModelRegistry 

42 from lilbee.providers.base import LLMProvider 

43 from lilbee.providers.roles import WorkerRole 

44 from lilbee.retrieval.clustering import Clusterer 

45 from lilbee.retrieval.concepts import ConceptGraph 

46 from lilbee.retrieval.embedder import Embedder 

47 from lilbee.retrieval.query import Searcher 

48 from lilbee.retrieval.reranker import Reranker 

49 from lilbee.runtime.ingest_lock import IngestLockRegistry 

50 from lilbee.sessions import SessionStore 

51 

52 

53log = logging.getLogger(__name__) 

54 

55_SIGNAL_EXIT_BASE = 128 

56 

57_HARD_EXIT_THREAD_NAME = "hard-exit-teardown" 

58 

59 

60def _default_session_store() -> SessionStore: 

61 """Build the file-backed session store, importing it lazily. 

62 

63 ``lilbee.sessions`` pulls in the config/catalog import chain, so importing it 

64 at this module's top would form a cycle during CLI config load. 

65 """ 

66 from lilbee.sessions import SessionStore 

67 

68 return SessionStore() 

69 

70 

71@dataclass 

72class CrawlerSyncState: 

73 """Process-wide sync coordination state (lock + last-run timestamp).""" 

74 

75 lock: threading.Lock = field(default_factory=threading.Lock) 

76 last_run: float = 0.0 

77 

78 

79@dataclass(frozen=True) 

80class Services: 

81 """Holds all runtime service instances. 

82 

83 Inference lifecycle (cancel, per-role reload, spawn notifications) is owned 

84 by the provider, which manages the llama-server fleet. Services exposes thin 

85 pass-throughs so callers (Ctrl+C, the chat-stream cancel action, the settings 

86 and model-bar pickers, the TUI task bar) need not reach into the provider's 

87 API. ``cancel_inference()`` is the canonical cancel entry point. 

88 """ 

89 

90 provider: LLMProvider 

91 store: Store 

92 embedder: Embedder 

93 reranker: Reranker 

94 concepts: ConceptGraph 

95 clusterer: Clusterer 

96 searcher: Searcher 

97 registry: ModelRegistry 

98 hf_client: HfClient 

99 ingest_lock_registry: IngestLockRegistry 

100 model_manager: ModelManager 

101 crawler_semaphore: asyncio.Semaphore | None 

102 crawler_sync_state: CrawlerSyncState 

103 known_models: KnownModelCache 

104 session_store: SessionStore = field(default_factory=_default_session_store) 

105 

106 def cancel_inference(self) -> None: 

107 """Interrupt any in-flight generation. Idempotent. 

108 

109 The fleet engine severs its live chat streams (llama-server stops 

110 generating when the connection drops); providers with nothing in 

111 flight treat this as a no-op. 

112 """ 

113 self.provider.cancel_inference() 

114 

115 def reload_role(self, role_name: WorkerRole, *, wait: bool = False) -> None: 

116 """Respawn only *role_name*'s model server so it picks up changed cfg. 

117 

118 Other roles' servers and any in-flight stream they own are untouched. Use 

119 when one role-bound model setting changed (e.g. embedding_model). The 

120 respawn runs off the caller's thread, so this returns immediately, unless 

121 ``wait=True`` (the caller is already off the event loop and wants to block 

122 until the new model has loaded). 

123 """ 

124 self.provider.reload_role(role_name, wait=wait) 

125 

126 def add_pool_listener( 

127 self, 

128 *, 

129 on_spawning: Callable[[WorkerRole], None] | None = None, 

130 on_spawned: Callable[[WorkerRole], None] | None = None, 

131 ) -> None: 

132 """Subscribe to server spawn lifecycle events. 

133 

134 Forwards to :meth:`LLMProvider.add_spawn_listener`. The TUI uses this to 

135 surface "Starting <role>..." / "<role> ready" notifications when a role's 

136 server (re)spawns (cold start after a non-eager boot, or a reload). 

137 """ 

138 self.provider.add_spawn_listener(on_spawning=on_spawning, on_spawned=on_spawned) 

139 

140 

141class _ServicesState: 

142 """The cached process-global singleton plus the per-task scoped override. 

143 

144 ``singleton`` is set on first ``get_services()`` call. Concurrency 

145 contract: creation is serialized by ``_singleton_create_lock`` (several 

146 worker threads can first-touch services at once, and a duplicate build 

147 would collide in xberg's process-global backend registry), and the 

148 Services dataclass is logically immutable post-construction, so concurrent 

149 reads are safe without a lock. Tests that need a custom container call 

150 ``set_services(make_mock_services(...))``; ``peek_services()`` is the 

151 read-only inspector for cleanup fixtures. 

152 

153 ``override`` shadows the singleton for the entering task: set by 

154 :func:`services_scope` (the library API's per-call binding), read by 

155 :func:`get_services`, and invisible to ``reset_services`` / 

156 ``set_services`` / ``peek_services``, which only touch the singleton. 

157 """ 

158 

159 def __init__(self) -> None: 

160 self.singleton: Services | None = None 

161 self.override: ContextVar[Services | None] = ContextVar( 

162 "lilbee_services_override", default=None 

163 ) 

164 # Whether this process is an interactive session (the TUI). Recorded by 

165 # the interactive entry point before anything builds the container, and 

166 # read once at build so the provider it creates holds its fleet resident 

167 # for the session. Build-time intent only; the state that matters after 

168 # that lives on the provider itself. 

169 self.interactive: bool = False 

170 

171 

172_state = _ServicesState() 

173 

174 

175def build_services( 

176 config: Config, 

177 *, 

178 provider: LLMProvider | None = None, 

179 registry: ModelRegistry | None = None, 

180 interactive: bool = False, 

181) -> Services: 

182 """Build a full Services container bound to *config*, without caching it. 

183 

184 ``get_services()`` calls this with the process-global cfg; the library API 

185 calls it per instance. Service modules are imported inside the function to 

186 keep CLI startup fast (they transitively pull in lancedb / xberg). Pass 

187 *provider* to reuse a caller-supplied one; otherwise it is built from 

188 *config* via the provider factory. Pass *registry* to reuse one already built 

189 (get_services builds it for embedding-dim reconciliation). Embedding-dim 

190 reconciliation is a global-cfg concern owned by :func:`get_services`, not 

191 done here. 

192 

193 Side effect: binds *provider* into xberg's process-global OCR and embedding 

194 backends (the chunker binds the tokenizer on demand). The registry is one 

195 process-wide slot, so every container (singleton or per-instance library) 

196 binds its own, not only get_services. 

197 """ 

198 from lilbee.catalog.hf_client import HfClient 

199 from lilbee.data.store import Store 

200 from lilbee.modelhub.model_manager import ModelManager 

201 from lilbee.modelhub.model_manager.discovery import KnownModelCache 

202 from lilbee.modelhub.registry import ModelRegistry 

203 from lilbee.modelhub.role_validator import warn_unregistered_role_refs 

204 from lilbee.providers.factory import create_provider 

205 from lilbee.retrieval.clustering import Clusterer 

206 from lilbee.retrieval.concepts import ConceptGraph 

207 from lilbee.retrieval.embedder import Embedder 

208 from lilbee.retrieval.query import Searcher 

209 from lilbee.retrieval.reranker import Reranker 

210 from lilbee.runtime.ingest_lock import IngestLockRegistry 

211 

212 provider = provider or create_provider(config, hold_warm=interactive) 

213 from lilbee.data.extract.backends import sync_xberg_backends 

214 

215 sync_xberg_backends(provider) 

216 registry = registry or ModelRegistry(config.models_dir) 

217 # The first point that sees env vars, --model and config.toml together. 

218 warn_unregistered_role_refs(config, registry) 

219 store = Store(config) 

220 embedder = Embedder(config, provider) 

221 reranker = Reranker(config) 

222 concepts = ConceptGraph(config, store) 

223 clusterer = Clusterer(config, store) 

224 searcher = Searcher(config, provider, store, embedder, reranker, concepts) 

225 hf_client = HfClient() 

226 ingest_lock_registry = IngestLockRegistry() 

227 model_manager = ModelManager(config.models_dir) 

228 crawler_semaphore = ( 

229 asyncio.Semaphore(config.crawl_max_concurrent) if config.crawl_max_concurrent > 0 else None 

230 ) 

231 crawler_sync_state = CrawlerSyncState() 

232 known_models = KnownModelCache() 

233 return Services( 

234 provider=provider, 

235 store=store, 

236 embedder=embedder, 

237 reranker=reranker, 

238 concepts=concepts, 

239 clusterer=clusterer, 

240 searcher=searcher, 

241 registry=registry, 

242 hf_client=hf_client, 

243 ingest_lock_registry=ingest_lock_registry, 

244 model_manager=model_manager, 

245 crawler_semaphore=crawler_semaphore, 

246 crawler_sync_state=crawler_sync_state, 

247 known_models=known_models, 

248 ) 

249 

250 

251# Serializes first-touch singleton creation: several worker threads (e.g. 

252# concurrent downloads) can call get_services() before the singleton exists, 

253# and a duplicate build re-registers xberg's process-global backends mid-flight, 

254# raising 'already registered' in the losing thread. 

255_singleton_create_lock = threading.Lock() 

256 

257 

258def get_services() -> Services: 

259 """Return the active container: a scoped override if set, else the cached singleton. 

260 

261 Creates the singleton on first call (against the process-global cfg). A 

262 config-file embedding_model with no embedding_dim would otherwise build the 

263 store at the stale 768 default, so the width is pinned to the embedder before 

264 the store is built. 

265 """ 

266 override = _state.override.get() 

267 if override is not None: 

268 return override 

269 if _state.singleton is not None: 

270 return _state.singleton 

271 

272 with _singleton_create_lock: 

273 if _state.singleton is not None: 

274 return _state.singleton 

275 

276 from lilbee.app.settings import reconcile_embedding_dim 

277 from lilbee.core.config import cfg 

278 from lilbee.modelhub.registry import ModelRegistry 

279 

280 registry = ModelRegistry(cfg.models_dir) 

281 # Pin the store width to the embedder before Store(); pass the registry so 

282 # resolution doesn't re-enter this half-built get_services. 

283 reconcile_embedding_dim(registry) 

284 _state.singleton = build_services(cfg, registry=registry, interactive=_state.interactive) 

285 # Eager start is the default: pay the spawn cost per role server at TUI mount 

286 # so the first user action lands on a warm fleet. Roles whose model is unset 

287 # are skipped, so a setup with only chat + embed never spawns rerank or 

288 # vision. Set ``cfg.worker_pool_eager_start = false`` for headless scripts 

289 # where mount time matters more than first-call latency. 

290 if cfg.worker_pool_eager_start: 

291 from contextlib import suppress 

292 

293 with suppress(Exception): 

294 _state.singleton.provider.warm_up_pool() 

295 return _state.singleton 

296 

297 

298@contextmanager 

299def services_scope(services: Services) -> Iterator[None]: 

300 """Bind *services* as the container ``get_services()`` returns for this block. 

301 

302 Isolated to the entering task via a ContextVar (it propagates into 

303 ``to_ingest_thread`` workers), and never affects the global singleton, so 

304 ``reset_services`` is unnecessary and unused around a scoped call. 

305 """ 

306 token = _state.override.set(services) 

307 try: 

308 yield 

309 finally: 

310 _state.override.reset(token) 

311 

312 

313def mark_interactive_session() -> None: 

314 """Record that this process is an interactive session before services build. 

315 

316 The TUI owns the process for its whole lifetime, so an engine bound to it 

317 keeps its weights resident rather than idle-unloading under a user who is 

318 still in the app; closing lilbee releases it. A keep-warm engine outlives the 

319 session and keeps its idle window instead. Called by the interactive entry point 

320 before anything touches ``get_services``, so the provider is constructed with 

321 that intent; a one-shot CLI or the MCP server never calls it. 

322 """ 

323 _state.interactive = True 

324 

325 

326def set_services(services: Services | None) -> None: 

327 """Replace the cached Services singleton (for testing).""" 

328 _state.singleton = services 

329 

330 

331def peek_services() -> Services | None: 

332 """Return the cached Services container, or None if not yet initialized. 

333 

334 Public read-only accessor for test cleanup helpers that need to 

335 inspect the singleton without forcing initialization. 

336 """ 

337 return _state.singleton 

338 

339 

340# Serializes the singleton swap in reset_services: a signal's teardown thread 

341# and an exiting caller's can race, and both tearing down the same container 

342# would double-close the store. 

343_reset_swap_lock = threading.Lock() 

344 

345 

346def reset_services() -> None: 

347 """Shut down and discard all cached instances. 

348 

349 Swap the module reference to ``None`` *before* tearing the old instances 

350 down, so a new caller never observes a half-closed container. The swap is 

351 locked so concurrent callers (a signal's teardown thread plus an exiting 

352 one) tear the container down exactly once. On the shared HTTP daemon every 

353 entry point that would call this mid-flight is refused, so it only ever 

354 runs single-client (CLI, TUI, stdio MCP). 

355 """ 

356 with _reset_swap_lock: 

357 old = _state.singleton 

358 _state.singleton = None 

359 if old is not None: 

360 old.provider.shutdown() 

361 old.store.close() 

362 

363 

364def reset_store() -> None: 

365 """Close and rebuild only the Store and its dependents; keep providers loaded. 

366 

367 Used after a data-dir wipe (``/reset``) where the LanceDB handle is invalid 

368 but the running provider/embedder/reranker are still good. Avoids the 

369 multi-second reload cost of ``reset_services()``. 

370 """ 

371 svc = _state.singleton 

372 if svc is None: 

373 return 

374 from dataclasses import replace 

375 

376 from lilbee.core.config import cfg 

377 from lilbee.data.store import Store 

378 from lilbee.retrieval.clustering import Clusterer 

379 from lilbee.retrieval.concepts import ConceptGraph 

380 from lilbee.retrieval.query import Searcher 

381 

382 # Build the replacement, swap the reference, then close the old store last so 

383 # a new caller never observes a closed handle mid-swap. 

384 old_store = svc.store 

385 store = Store(cfg) 

386 concepts = ConceptGraph(cfg, store) 

387 clusterer = Clusterer(cfg, store) 

388 searcher = Searcher(cfg, svc.provider, store, svc.embedder, svc.reranker, concepts) 

389 _state.singleton = replace( 

390 svc, 

391 store=store, 

392 concepts=concepts, 

393 clusterer=clusterer, 

394 searcher=searcher, 

395 ) 

396 old_store.close() 

397 

398 

399class _ServerExit: 

400 """Holds the running server's graceful-exit callback while it serves.""" 

401 

402 def __init__(self) -> None: 

403 self._hook: Callable[[], None] | None = None 

404 

405 def set(self, hook: Callable[[], None] | None) -> None: 

406 self._hook = hook 

407 

408 def request(self) -> bool: 

409 """Run the hook; False when no server registered one.""" 

410 if self._hook is None: 

411 return False 

412 self._hook() 

413 return True 

414 

415 

416_server_exit = _ServerExit() 

417 

418 

419def set_server_exit_hook(hook: Callable[[], None] | None) -> None: 

420 """Register the running server's graceful-exit callback (None clears it).""" 

421 _server_exit.set(hook) 

422 

423 

424def request_server_exit() -> bool: 

425 """Ask the running server to exit; False when none registered.""" 

426 return _server_exit.request() 

427 

428 

429class _EngineLifecycle: 

430 """Owns the hard-exit hooks that stop the engine fleet.""" 

431 

432 def __init__(self) -> None: 

433 self._installed = False 

434 

435 @staticmethod 

436 def _hard_exit_signals() -> tuple[signal.Signals, ...]: # pragma: no cover - platform split 

437 """Signals whose default disposition kills us without running atexit.""" 

438 if sys.platform == "win32": # Windows has no SIGHUP 

439 return (signal.SIGTERM,) 

440 return (signal.SIGTERM, signal.SIGHUP) 

441 

442 def install(self) -> None: 

443 """Route hard-exit signals through teardown. Idempotent; no-op off the main thread.""" 

444 if self._installed: 

445 return 

446 try: 

447 for sig in self._hard_exit_signals(): 

448 signal.signal(sig, self._on_hard_exit) 

449 except ValueError: 

450 return 

451 self._installed = True 

452 

453 def reset(self) -> None: 

454 """Forget that handlers were installed.""" 

455 self._installed = False 

456 

457 def _on_hard_exit(self, signum: int, frame: object) -> None: 

458 """Stop the fleet on its own thread, then exit with the signal status. 

459 

460 The serving loop is asked to stop first: a SystemExit raised while the 

461 main thread runs a request dies inside that request's ASGI wrapper, so 

462 the flag stops the loop when the raise cannot reach it. Signal handlers 

463 all run on the main thread, and a second signal can interrupt this one 

464 mid-teardown: the kernel pairs SIGCONT with SIGHUP for an orphaned 

465 process group, and Textual's SIGCONT handler raises once the event loop 

466 is gone, which aborted the reap half-done and orphaned a loaded fleet. 

467 A dedicated non-daemon thread cannot be interrupted by signals, and the 

468 interpreter waits for it even as the SystemExit unwinds the main thread. 

469 """ 

470 del frame 

471 request_server_exit() 

472 threading.Thread( 

473 target=_teardown_for_signal, args=(signum,), name=_HARD_EXIT_THREAD_NAME 

474 ).start() 

475 raise SystemExit(_SIGNAL_EXIT_BASE + signum) 

476 

477 

478_FORCE_QUIT_EXIT_CODE = 130 

479 

480_STOPPING_NOTE = "lilbee: stopping the engine. Press Ctrl-C again to force quit.\n" 

481 

482_STOPPING_NOTE_GRACE_S = 1.0 

483"""Seconds a teardown may run silently; a fast exit stays as quiet as ever.""" 

484 

485_FORCE_QUIT_NOTE = "lilbee: force quit. The engine stops with it.\n" 

486 

487_STRAGGLER_BUDGET_S = 5.0 

488"""Seconds every leftover non-daemon thread gets, together, before the exit.""" 

489 

490_STRAGGLER_NOTE = "lilbee: a background task would not stop; exiting anyway.\n" 

491 

492 

493def _write_exit_note(note: str) -> None: 

494 """Print an exit-path status line; stderr may already be gone at exit.""" 

495 try: 

496 sys.stderr.write(note) 

497 sys.stderr.flush() 

498 except (OSError, ValueError): 

499 pass 

500 

501 

502def _all_threads() -> list[threading.Thread]: 

503 """Seam over threading.enumerate, so tests never patch the stdlib global.""" 

504 return threading.enumerate() 

505 

506 

507def wait_for_hard_exit_teardown() -> None: 

508 """Block until any teardown thread (signal-driven or exit-driven) finishes. 

509 

510 Lets ``serve`` hold its OS locks through the fleet stop, so a successor 

511 cannot acquire them while this server's models still occupy memory. 

512 

513 An interrupt while waiting force-quits the process: waiting is the only 

514 thing left, the interpreter would otherwise still join the non-daemon 

515 teardown thread, and every platform's child guard reaps the engine when 

516 the process dies (kill-on-close job object, pdeathsig, death pipe). 

517 """ 

518 try: 

519 for thread in _all_threads(): 

520 if thread.name != _HARD_EXIT_THREAD_NAME: 

521 continue 

522 thread.join(_STOPPING_NOTE_GRACE_S) 

523 if thread.is_alive(): 

524 _write_exit_note(_STOPPING_NOTE) 

525 thread.join() 

526 except KeyboardInterrupt: 

527 _write_exit_note(_FORCE_QUIT_NOTE) 

528 os._exit(_FORCE_QUIT_EXIT_CODE) 

529 

530 

531def reset_services_on_exit() -> None: 

532 """Tear the container down on a thread no signal reaches, and wait for it. 

533 

534 Engine release waits on the fleet build lock before it releases anything, so 

535 a Ctrl-C on the main thread skips the release, and atexit cannot retry: the 

536 singleton is already cleared. The teardown thread is non-daemon and takes 

537 no signals; a wait that outlives the grace names what is happening. 

538 """ 

539 if peek_services() is None: 

540 return 

541 threading.Thread(target=reset_services, name=_HARD_EXIT_THREAD_NAME).start() 

542 wait_for_hard_exit_teardown() 

543 

544 

545def _straggler_threads() -> list[threading.Thread]: 

546 """Non-daemon threads the interpreter would join at shutdown, besides us.""" 

547 return [ 

548 thread 

549 for thread in _all_threads() 

550 if thread.is_alive() 

551 and not thread.daemon 

552 and thread is not threading.main_thread() 

553 and thread is not threading.current_thread() 

554 ] 

555 

556 

557def exit_when_stragglers_would_hang() -> None: 

558 """Exit rather than let interpreter shutdown join a wedged thread forever. 

559 

560 threading joins every non-daemon thread after atexit, so one worker 

561 blocked in an unbounded network read hangs the process after all real 

562 work is done. Stragglers share a short budget; whatever remains cannot 

563 finish, and the child guard reaps the engine when the process dies. 

564 Called from the TUI's own exit path and inert elsewhere: a test process 

565 or embedding host owns its threads, and an exit would take them with it. 

566 """ 

567 if not _state.interactive: 

568 return 

569 try: 

570 deadline = time.monotonic() + _STRAGGLER_BUDGET_S 

571 for thread in _straggler_threads(): 

572 thread.join(max(0.0, deadline - time.monotonic())) 

573 if _straggler_threads(): 

574 _write_exit_note(_STRAGGLER_NOTE) 

575 os._exit(0) 

576 except KeyboardInterrupt: 

577 _write_exit_note(_FORCE_QUIT_NOTE) 

578 os._exit(_FORCE_QUIT_EXIT_CODE) 

579 

580 

581def _teardown_for_signal(signum: int) -> None: 

582 """Log the fatal signal, then stop services; runs off the signal handler's thread.""" 

583 log.info( 

584 "Received signal %s; stopping the engine fleet before exit", signal.Signals(signum).name 

585 ) 

586 reset_services() 

587 

588 

589_lifecycle = _EngineLifecycle() 

590 

591 

592def install_engine_lifecycle_hooks() -> None: 

593 """Make a terminal close or ``kill`` stop the engine fleet instead of orphaning it.""" 

594 _lifecycle.install() 

595 

596 

597atexit.register(reset_services_on_exit)