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
« 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.
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.
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"""
17from __future__ import annotations
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
33if TYPE_CHECKING:
34 from collections.abc import Iterator
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
53log = logging.getLogger(__name__)
55_SIGNAL_EXIT_BASE = 128
57_HARD_EXIT_THREAD_NAME = "hard-exit-teardown"
60def _default_session_store() -> SessionStore:
61 """Build the file-backed session store, importing it lazily.
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
68 return SessionStore()
71@dataclass
72class CrawlerSyncState:
73 """Process-wide sync coordination state (lock + last-run timestamp)."""
75 lock: threading.Lock = field(default_factory=threading.Lock)
76 last_run: float = 0.0
79@dataclass(frozen=True)
80class Services:
81 """Holds all runtime service instances.
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 """
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)
106 def cancel_inference(self) -> None:
107 """Interrupt any in-flight generation. Idempotent.
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()
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.
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)
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.
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)
141class _ServicesState:
142 """The cached process-global singleton plus the per-task scoped override.
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.
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 """
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
172_state = _ServicesState()
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.
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.
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
212 provider = provider or create_provider(config, hold_warm=interactive)
213 from lilbee.data.extract.backends import sync_xberg_backends
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 )
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()
258def get_services() -> Services:
259 """Return the active container: a scoped override if set, else the cached singleton.
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
272 with _singleton_create_lock:
273 if _state.singleton is not None:
274 return _state.singleton
276 from lilbee.app.settings import reconcile_embedding_dim
277 from lilbee.core.config import cfg
278 from lilbee.modelhub.registry import ModelRegistry
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
293 with suppress(Exception):
294 _state.singleton.provider.warm_up_pool()
295 return _state.singleton
298@contextmanager
299def services_scope(services: Services) -> Iterator[None]:
300 """Bind *services* as the container ``get_services()`` returns for this block.
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)
313def mark_interactive_session() -> None:
314 """Record that this process is an interactive session before services build.
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
326def set_services(services: Services | None) -> None:
327 """Replace the cached Services singleton (for testing)."""
328 _state.singleton = services
331def peek_services() -> Services | None:
332 """Return the cached Services container, or None if not yet initialized.
334 Public read-only accessor for test cleanup helpers that need to
335 inspect the singleton without forcing initialization.
336 """
337 return _state.singleton
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()
346def reset_services() -> None:
347 """Shut down and discard all cached instances.
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()
364def reset_store() -> None:
365 """Close and rebuild only the Store and its dependents; keep providers loaded.
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
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
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()
399class _ServerExit:
400 """Holds the running server's graceful-exit callback while it serves."""
402 def __init__(self) -> None:
403 self._hook: Callable[[], None] | None = None
405 def set(self, hook: Callable[[], None] | None) -> None:
406 self._hook = hook
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
416_server_exit = _ServerExit()
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)
424def request_server_exit() -> bool:
425 """Ask the running server to exit; False when none registered."""
426 return _server_exit.request()
429class _EngineLifecycle:
430 """Owns the hard-exit hooks that stop the engine fleet."""
432 def __init__(self) -> None:
433 self._installed = False
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)
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
453 def reset(self) -> None:
454 """Forget that handlers were installed."""
455 self._installed = False
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.
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)
478_FORCE_QUIT_EXIT_CODE = 130
480_STOPPING_NOTE = "lilbee: stopping the engine. Press Ctrl-C again to force quit.\n"
482_STOPPING_NOTE_GRACE_S = 1.0
483"""Seconds a teardown may run silently; a fast exit stays as quiet as ever."""
485_FORCE_QUIT_NOTE = "lilbee: force quit. The engine stops with it.\n"
487_STRAGGLER_BUDGET_S = 5.0
488"""Seconds every leftover non-daemon thread gets, together, before the exit."""
490_STRAGGLER_NOTE = "lilbee: a background task would not stop; exiting anyway.\n"
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
502def _all_threads() -> list[threading.Thread]:
503 """Seam over threading.enumerate, so tests never patch the stdlib global."""
504 return threading.enumerate()
507def wait_for_hard_exit_teardown() -> None:
508 """Block until any teardown thread (signal-driven or exit-driven) finishes.
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.
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)
531def reset_services_on_exit() -> None:
532 """Tear the container down on a thread no signal reaches, and wait for it.
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()
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 ]
557def exit_when_stragglers_would_hang() -> None:
558 """Exit rather than let interpreter shutdown join a wedged thread forever.
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)
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()
589_lifecycle = _EngineLifecycle()
592def install_engine_lifecycle_hooks() -> None:
593 """Make a terminal close or ``kill`` stop the engine fleet instead of orphaning it."""
594 _lifecycle.install()
597atexit.register(reset_services_on_exit)