Coverage for src/lilbee/mcp_server.py: 100%
749 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"""MCP server exposing lilbee as tools for AI agents."""
3from __future__ import annotations
5import asyncio
6import concurrent.futures
7import functools
8import inspect
9import json
10import logging
11import os
12import re
13import threading
14import uuid
15from collections.abc import Callable, Iterator
16from contextlib import contextmanager
17from contextvars import ContextVar
18from copy import deepcopy
19from dataclasses import asdict
20from pathlib import Path
21from typing import TYPE_CHECKING, Any, TypeVar, cast
22from weakref import WeakKeyDictionary
24import anyio
25from mcp.server.mcpserver import Context, MCPServer
26from mcp.types import Tool as MCPTool
28from lilbee.app.memory import (
29 MEMORY_DISABLED_HINT,
30 forget,
31 list_memories,
32 memory_enabled,
33 recall,
34 remember,
35)
36from lilbee.app.placement import (
37 PlacementView,
38 get_placement,
39 placement_refused_message,
40 preview_placement,
41 set_placement,
42)
43from lilbee.app.search import clean_result
44from lilbee.app.services import get_services, reset_services, reset_store
45from lilbee.app.settings import (
46 SettingInfo,
47 apply_settings_update,
48 config_write_failure_message,
49 get_setting,
50 list_settings,
51 provider_reset_refused_message,
52 requires_services_reset,
53 reset_settings,
54)
55from lilbee.catalog.types import ModelSource
56from lilbee.core.config import cfg, validate_ocr_timeout
57from lilbee.core.config.enums import CrawlRenderMode
58from lilbee.core.settings import overlay_persisted_settings
59from lilbee.core.system import LOCAL_ROOT_DIRNAME, canonical_data_root
60from lilbee.crawler import crawler_available, is_url, require_valid_crawl_url
61from lilbee.crawler.task import get_task, start_crawl
62from lilbee.data.store import (
63 EmbeddingModelMismatchError,
64 MemoryKind,
65 MemorySource,
66 SearchScope,
67 agent_owner,
68 scope_to_chunk_type,
69)
70from lilbee.runtime.cancellation import TaskCancelledError
71from lilbee.runtime.hardware import FitLevel, available_memory_for_fit, make_fit_filter
72from lilbee.runtime.lock import ResetRefusedError
73from lilbee.sessions import (
74 AGENT_SESSIONS_DISABLED_HINT,
75 MessageRole,
76 Session,
77 SessionForkRangeError,
78 SessionMessage,
79 SessionNotFoundError,
80 SessionOrigin,
81 SessionOwnershipError,
82 TitleSource,
83 agent_sessions_enabled,
84)
85from lilbee.wiki.shared import (
86 INVALID_DRAFT_SLUG_ERROR,
87 WIKI_DISABLED_ERROR,
88 WikiSubdir,
89 total_wiki_pages,
90)
92if TYPE_CHECKING:
93 from lilbee.providers.fleet.placement_spec import PlacementSpec
95log = logging.getLogger(__name__)
97_INSTRUCTIONS = (
98 "Local search engine over the user's files, code, and crawled pages. "
99 "For any question about the user's own documents or codebase -- a lookup, a "
100 "find-in-docs, 'where is X', 'how does Y work here' -- call lilbee_search first "
101 "and answer from its cited chunks. Prefer it over web-fetch or file-read tools: "
102 "those cannot see the indexed corpus."
103)
106class _TransportState:
107 """Process-level MCP transport facts.
109 ``http_mounted`` is True only when MCP is exposed over the shared
110 streamable-http daemon (set by build_mcp_mount). On that transport multiple
111 agents share one process and one global cfg/Services singleton, so
112 vault-switching (init) and factory reset are refused: switching or tearing
113 down the store under concurrent in-flight handlers is a use-after-close /
114 identity race. stdio (one agent per process) keeps both.
115 """
117 http_mounted: bool = False
120_transport = _TransportState()
123def set_http_mounted(value: bool) -> None:
124 """Mark whether this process serves MCP over the shared HTTP daemon."""
125 _transport.http_mounted = value
128_F = TypeVar("_F", bound=Callable[..., Any])
130# Set by the offload for the duration of one sync handler and read by the long
131# ones. A parameter would reach the generated tool schema and invite an agent to
132# pass it; the handlers keep the plain signature that in-process callers use.
133_CANCEL: ContextVar[threading.Event | None] = ContextVar("lilbee_mcp_cancel", default=None)
136def _caller_cancelled() -> threading.Event | None:
137 """The stop token for the handler running on this thread, if the offload set one."""
138 return _CANCEL.get()
141def _offload_sync(fn: _F) -> _F:
142 """Run a sync tool handler off the event loop; async handlers pass through.
144 The bundled mcp SDK calls sync tool handlers directly on the event loop, so
145 under the shared streamable-http daemon one slow handler would stall every
146 connected agent. ``functools.wraps`` preserves the wrapped signature so the
147 generated tool schema is unchanged.
149 ``abandon_on_cancel`` releases the caller the moment it cancels. A worker
150 thread cannot be interrupted, so the handler still runs to completion; the
151 default would additionally hold the cancelling task until it did, which on
152 the shared mount keeps a connection open for a request nobody is waiting
153 for. Work that must actually stop takes a cancel token instead, as the
154 ingest tools do.
155 """
156 if inspect.iscoroutinefunction(fn):
157 return fn
159 @functools.wraps(fn)
160 async def _runner(*args: Any, **kwargs: Any) -> Any:
161 with _cancel_token() as token:
162 _CANCEL.set(token)
163 return await anyio.to_thread.run_sync(
164 functools.partial(fn, *args, **kwargs), abandon_on_cancel=True
165 )
167 return cast("_F", _runner)
170# (handler, wire name override, gate) applied by build_mcp_server per instance.
171_REGISTRATIONS: list[tuple[Callable[..., Any], str | None, Callable[[], bool] | None]] = []
174def _tool(fn: _F) -> _F:
175 """Register *fn* as an MCP tool with sync handlers offloaded off the loop.
177 Returns the original callable so in-process callers (tests, the stdio
178 fallback) keep the synchronous API while the schema sees the offloaded form.
179 """
180 _REGISTRATIONS.append((_offload_sync(fn), None, None))
181 return fn
184def _tool_named(name: str) -> Callable[[_F], _F]:
185 """Register an MCP tool under an explicit wire *name* (sync handlers offloaded)."""
187 def deco(fn: _F) -> _F:
188 _REGISTRATIONS.append((_offload_sync(fn), name, None))
189 return fn
191 return deco
194def _tool_if(when: Callable[[], bool]) -> Callable[[_F], _F]:
195 """Register an MCP tool gated on *when*, evaluated at server-build time.
197 The function stays importable so direct callers (tests, in-process
198 fallback) can still reach it. A server built after a config change
199 carries the current tool surface; live servers keep theirs.
200 """
201 if not callable(when):
202 raise TypeError("_tool_if takes a zero-arg callable, evaluated per build")
204 def deco(fn: _F) -> _F:
205 _REGISTRATIONS.append((_offload_sync(fn), None, when))
206 return fn
208 return deco
211# Settings read once per server build to decide which tools register. Writing
212# one over ``settings_set`` persists it, but the tool list only changes on the
213# next connection: MCP sends no ``tools/list_changed`` from here. The settings
214# reference documents that, and generates the list from this set.
215TOOL_GATE_SETTINGS: frozenset[str] = frozenset({"wiki", "memory_enabled", "mcp_sessions_enabled"})
218def _wiki_enabled() -> bool:
219 return cfg.wiki
222def _error(msg: str) -> dict[str, Any]:
223 """Uniform error envelope MCP tool handlers return on a failure path.
225 Typed as ``dict[str, Any]`` rather than a TypedDict so it composes
226 with the success-side returns under the existing handler signatures
227 without forcing every caller to widen its return type.
228 """
229 return {"error": msg}
232def _wiki_tool(fn: _F) -> _F:
233 """Register a wiki MCP tool: gated at build time, re-checked per call.
235 A live server keeps the tool surface it was built with, so every handler
236 re-reads ``cfg.wiki`` and refuses once the setting is turned off at runtime.
237 """
239 @functools.wraps(fn)
240 def guarded(*args: Any, **kwargs: Any) -> Any:
241 if not cfg.wiki:
242 return _error(WIKI_DISABLED_ERROR)
243 return fn(*args, **kwargs)
245 checked = cast("_F", guarded)
246 _tool_if(_wiki_enabled)(checked)
247 return checked
250@_tool
251def search(
252 query: str, top_k: int | None = None, scope: str = SearchScope.BOTH.value
253) -> list[dict[str, Any]] | dict[str, Any]:
254 """Search the user's indexed documents, code, and crawled pages; prefer it over web-fetch or
255 file-read tools. Returns chunks with citations. ``scope``: "both" (default) / "raw" / "wiki"."""
256 if not query or not query.strip():
257 return _error("query must not be empty")
258 try:
259 chunk_type = scope_to_chunk_type(scope)
260 except ValueError:
261 # Smaller models routinely echo prose like "indexed docs" or "all"
262 # back as the scope value. Treat unrecognised scopes as the default
263 # "both" rather than a hard failure so the request still does work.
264 log.warning("lilbee_search: unknown scope %r, falling back to %r", scope, SearchScope.BOTH)
265 chunk_type = scope_to_chunk_type(SearchScope.BOTH.value)
266 if top_k is not None and top_k < 1:
267 # Same lenient stance as the scope fallback above: do the search with
268 # the configured default rather than hard-failing the agent's call.
269 log.warning("lilbee_search: top_k %d is not positive, using the default", top_k)
270 top_k = None
271 effective_top_k = top_k if top_k is not None else cfg.top_k
272 try:
273 results = get_services().searcher.search(
274 query, top_k=effective_top_k, chunk_type=chunk_type
275 )
276 return [clean_result(r) for r in results]
277 except EmbeddingModelMismatchError as exc:
278 # Structured so an agent can offer to adopt the index's embedder
279 # rather than parse prose out of a generic error. Names the embedder,
280 # which the HTTP search route keeps out of its generic 503.
281 return {
282 "error": str(exc),
283 "code": "INDEX_EMBEDDER_MISMATCH",
284 "persisted_model": exc.persisted_model,
285 "persisted_dim": exc.persisted_dim,
286 "adoptable": exc.dims_match,
287 }
288 except Exception as exc:
289 return _error(str(exc))
292@_tool
293def status() -> dict[str, Any]:
294 """Show indexed documents, configuration, and chunk counts."""
295 from lilbee.app.status import gather_status
297 return gather_status().model_dump()
300@contextmanager
301def _cancel_token() -> Iterator[threading.Event]:
302 """A stop token for a long operation, set on any abnormal exit.
304 An agent cancelling a tool call gets what the HTTP surface gets on a client
305 disconnect: the work stops at its next boundary and keeps what it finished,
306 rather than running the rest of a corpus, download or wiki build for a
307 request that has gone away. The work runs on a thread asyncio cannot
308 interrupt, which is why cancellation alone does not reach it and the token
309 has to be polled.
311 Set on every exception, not only cancellation: a pass that raised is over
312 either way, and a caller vanishing does not always surface as a cancel.
313 """
314 token = threading.Event()
315 try:
316 yield token
317 except BaseException:
318 token.set()
319 raise
322@_tool
323async def sync(
324 force_rebuild: bool = False,
325 retry_skipped: bool = False,
326 prune_ignored: bool = False,
327 enable_ocr: bool | None = None,
328 ocr_timeout: float | None = None,
329) -> dict[str, Any]:
330 """Sync the documents directory into the vector store.
332 ``force_rebuild`` drops every table and re-ingests. ``retry_skipped``
333 clears failed-file skip markers. ``prune_ignored`` also drops sources a
334 ``.lilbeeignore`` now excludes.
335 """
336 from lilbee.app.ingest import temporary_ocr_config
337 from lilbee.data.ingest import sync as run_sync
339 try:
340 validate_ocr_timeout(ocr_timeout)
341 except ValueError as exc:
342 return _error(str(exc))
343 with temporary_ocr_config(enable_ocr, ocr_timeout), _cancel_token() as cancel:
344 return (
345 await run_sync(
346 quiet=True,
347 force_rebuild=force_rebuild,
348 retry_skipped=retry_skipped,
349 prune_ignored=prune_ignored,
350 cancel=cancel,
351 )
352 ).model_dump()
355async def _sync_after_add(
356 reached_corpus: bool, enable_ocr: bool | None, ocr_timeout: float | None
357) -> dict[str, Any] | None:
358 """The sync that follows an add, or None when nothing named reached the corpus."""
359 from lilbee.app.ingest import temporary_ocr_config
360 from lilbee.data.ingest import sync as run_sync
362 if not reached_corpus:
363 return None
364 with temporary_ocr_config(enable_ocr, ocr_timeout), _cancel_token() as cancel:
365 return (await run_sync(quiet=True, cancel=cancel)).model_dump()
368async def _crawl_add_urls(
369 urls: list[str], render_mode: CrawlRenderMode | None
370) -> tuple[int, list[str]]:
371 """Crawl each URL as a single page for ``add``; returns (crawled_count, url_errors)."""
372 from lilbee.crawler import crawl_and_save
374 crawled_count = 0
375 errors: list[str] = []
376 for url in urls:
377 try:
378 # URL validation resolves the host (blocking DNS); run it off the
379 # event loop like the sibling crawl tool does.
380 await anyio.to_thread.run_sync(require_valid_crawl_url, url)
381 except ValueError as exc:
382 errors.append(f"{url}: {exc}")
383 continue
384 # add fetches single pages (depth=0); site crawls go through the crawl tool
385 crawled_paths = await crawl_and_save(url, depth=0, render_mode=render_mode)
386 crawled_count += len(crawled_paths)
387 return crawled_count, errors
390@_tool
391async def add(
392 paths: list[str],
393 force: bool = False,
394 enable_ocr: bool | None = None,
395 ocr_timeout: float | None = None,
396 render_mode: CrawlRenderMode | None = None,
397) -> dict[str, Any]:
398 """Add files, directories, or URLs to the knowledge base, then sync.
399 Paths are absolute and resolve on the machine running lilbee, which is not
400 the caller's machine when the server is remote. URLs are fetched as single
401 pages; use ``crawl`` for sites."""
402 from lilbee.app.ingest import register_sources
404 try:
405 validate_ocr_timeout(ocr_timeout)
406 except ValueError as exc:
407 return _error(str(exc))
409 errors: list[str] = []
410 valid: list[Path] = []
411 urls: list[str] = []
412 for p_str in paths:
413 if is_url(p_str):
414 urls.append(p_str)
415 else:
416 p = Path(p_str)
417 if not p.exists():
418 # Name the machine and the root. A remote caller sends paths
419 # from its own filesystem, and a bare path in `errors` reads as
420 # a transient hiccup: agents took it as success and moved on.
421 errors.append(
422 f"{p_str}: no such file or directory on the lilbee server. "
423 f"Paths resolve on the server, whose documents root is "
424 f"{cfg.documents_dir}."
425 )
426 else:
427 valid.append(p)
429 # Crawl URLs
430 crawled_count = 0
431 if urls:
432 from lilbee.crawler import crawler_available
434 if not crawler_available():
435 return _error("Web crawling requires: pip install 'lilbee[crawler]'")
436 crawled_count, url_errors = await _crawl_add_urls(urls, render_mode)
437 errors.extend(url_errors)
439 # Registration touches config.toml (a locked read-modify-write); keep the
440 # blocking disk I/O off the event loop.
441 reg_result = await anyio.to_thread.run_sync(
442 functools.partial(register_sources, valid, force=force)
443 )
444 errors.extend(reg_result.refused)
445 reached = reg_result.reached_corpus or bool(crawled_count)
446 sync_result = await _sync_after_add(reached, enable_ocr, ocr_timeout)
447 result: dict[str, Any] = {
448 "command": "add",
449 "copied": reg_result.registered,
450 "name_taken": reg_result.name_taken,
451 "overlapping": reg_result.overlapping,
452 "tracked": reg_result.tracked,
453 "crawled": crawled_count,
454 "errors": errors,
455 "sync": sync_result,
456 }
457 if errors and not reached:
458 # Nothing was added. Returning the success shape with a warning let a
459 # caller report the add as done over an untouched index.
460 return _error("add indexed nothing. " + " ".join(errors))
461 if errors or (sync_result is not None and sync_result.get("failed")):
462 result["warning"] = "some files could not be processed"
463 return result
466@_tool_if(crawler_available)
467async def crawl(
468 url: str,
469 depth: int | None = 0,
470 max_pages: int | None = None,
471 render_mode: CrawlRenderMode | None = None,
472 include_subdomains: bool = False,
473) -> dict[str, Any]:
474 """Start a non-blocking crawl; poll via ``crawl_status(task_id)``.
475 ``depth=0`` (default) = single URL, ``N`` = follow links N levels,
476 ``null`` = whole site. ``render_mode``: "http"/"browser"."""
477 from lilbee.crawler import crawler_available
479 if not crawler_available():
480 return _error("Web crawling requires: pip install 'lilbee[crawler]'")
481 # Mirror the REST CrawlRequest bounds so a negative value is a clean error,
482 # not an unbounded crawl.
483 if depth is not None and depth < 0:
484 return _error("depth must be 0 or greater (pass depth=null to crawl the whole site)")
485 if max_pages is not None and max_pages < 0:
486 return _error("max_pages must be 0 or greater (0 = unlimited, omit for the safety cap)")
487 try:
488 # URL validation resolves the host (blocking DNS), so it runs off the loop.
489 # The crawl itself must be scheduled ON the loop: start_crawl uses
490 # asyncio.create_task, which requires a running event loop.
491 await anyio.to_thread.run_sync(require_valid_crawl_url, url)
492 except ValueError as exc:
493 return _error(str(exc))
495 task_id = start_crawl(
496 url,
497 depth=depth,
498 max_pages=max_pages,
499 render_mode=render_mode,
500 include_subdomains=include_subdomains,
501 )
502 return {"status": "started", "task_id": task_id, "url": url}
505@_tool_if(crawler_available)
506def crawl_status(task_id: str) -> dict[str, Any]:
507 """Poll a crawl task by id; returns ``{status, pages, error}``."""
508 task = get_task(task_id)
509 if task is None:
510 return _error(f"No task found with id: {task_id}")
511 return {
512 "task_id": task.task_id,
513 "url": task.url,
514 "status": task.status.value,
515 "pages_crawled": task.pages_crawled,
516 "pages_total": task.pages_total,
517 "pages_failed": task.pages_failed,
518 "failure_reasons": task.failure_reasons,
519 "error": task.error,
520 "started_at": task.started_at,
521 "finished_at": task.finished_at,
522 }
525@_tool
526def crawl_cancel(task_id: str) -> dict[str, Any]:
527 """Stop a running crawl started by ``crawl``. Pages already saved are kept."""
528 from lilbee.crawler.task import cancel_crawl
530 if get_task(task_id) is None:
531 return _error(f"No task found with id: {task_id}")
532 return {"command": "crawl_cancel", "task_id": task_id, "cancelling": cancel_crawl(task_id)}
535@_tool
536def init(path: str = "") -> dict[str, Any]:
537 """Initialize a local ``.lilbee/`` knowledge base; empty path = cwd.
539 Switches the MCP session to use it for subsequent calls.
540 """
541 if _transport.http_mounted:
542 return _error(
543 "init is unavailable on the HTTP server: it is bound to one vault and "
544 "shared by every connected client. Start a separate server for another vault."
545 )
546 # Canonical so this vault keys the same lock paths a CLI or server process
547 # would derive for the same directory.
548 base = canonical_data_root(path) if path else canonical_data_root(Path.cwd())
549 root = base / LOCAL_ROOT_DIRNAME
551 created = False
552 if not root.is_dir():
553 (root / "documents").mkdir(parents=True)
554 (root / "data").mkdir(parents=True)
555 (root / ".gitignore").write_text("data/\n", encoding="utf-8")
556 created = True
558 # Switch MCP session to this project's KB. Overlay any persisted
559 # config.toml so per-vault model / generation settings take effect,
560 # matching the CLI's --data-dir behaviour. Env export mirrors
561 # cli/app.py::_apply_data_root for worker-log parity.
562 cfg.data_root = base
563 cfg.documents_dir = root / "documents"
564 cfg.data_dir = root / "data"
565 cfg.lancedb_dir = root / "data" / "lancedb"
566 os.environ["LILBEE_DATA"] = str(base)
567 overlay_persisted_settings(base)
568 reset_services()
570 return {"command": "init", "path": str(root), "created": created}
573@_tool
574def remove(names: list[str]) -> dict[str, Any]:
575 """Remove documents from the index by source name, folder, or glob pattern.
577 Source files are never deleted. A folder name removes every document indexed
578 beneath it; a glob (``*``/``?``/``[]``) removes every matching source."""
579 from lilbee.app.ingest import remove_documents_durably
581 result = remove_documents_durably(names)
582 return {"command": "remove", "removed": result.removed, "not_found": result.not_found}
585@_tool
586def list_documents() -> dict[str, Any]:
587 """List all indexed documents with their chunk counts."""
588 sources = get_services().store.get_sources()
589 return {
590 "documents": [
591 {"filename": s["filename"], "chunk_count": s.get("chunk_count", 0)} for s in sources
592 ],
593 "total": len(sources),
594 }
597_AGENT_ORIGINS = frozenset({SessionOrigin.MCP})
600def _require_agent_session(session_id: str) -> Session:
601 """Human conversations are private: foreign ids answer not-found, so an
602 agent cannot even probe which ones exist."""
603 session = get_services().session_store.get(session_id)
604 if session.meta.origin is not SessionOrigin.MCP:
605 raise SessionNotFoundError(session_id)
606 return session
609@_tool_if(agent_sessions_enabled)
610def sessions_list() -> dict[str, Any]:
611 """List the agent's sessions, newest first."""
612 if not agent_sessions_enabled():
613 return _error(AGENT_SESSIONS_DISABLED_HINT)
614 metas = get_services().session_store.list(origins=_AGENT_ORIGINS)
615 return {"sessions": [asdict(meta) for meta in metas], "total": len(metas)}
618@_tool_if(agent_sessions_enabled)
619def session_get(session_id: str) -> dict[str, Any]:
620 """Return one agent session: metadata, transcript, summary."""
621 if not agent_sessions_enabled():
622 return _error(AGENT_SESSIONS_DISABLED_HINT)
623 try:
624 session = _require_agent_session(session_id)
625 except SessionNotFoundError as exc:
626 return _error(str(exc))
627 return _session_payload(session)
630def _session_payload(session: Session) -> dict[str, Any]:
631 """An agent session as the session tools return it: meta, transcript, summary."""
632 return {
633 "meta": asdict(session.meta),
634 "messages": [
635 {
636 "role": message.role.value,
637 "content": message.content,
638 "sources": list(message.sources),
639 "ts": message.ts,
640 }
641 for message in session.messages
642 ],
643 # What compaction folded the older turns into (empty if never compacted).
644 # An agent that resumes and continues the conversation needs it, or it
645 # rebuilds history without what was already condensed.
646 "summary": session.summary,
647 }
650@_tool_if(agent_sessions_enabled)
651def session_fork(session_id: str, message_count: int | None = None) -> dict[str, Any]:
652 """Copy the first message_count messages (all when omitted) into a new session."""
653 if not agent_sessions_enabled():
654 return _error(AGENT_SESSIONS_DISABLED_HINT)
655 store = get_services().session_store
656 try:
657 _require_agent_session(session_id)
658 fork_id = store.fork(session_id, message_count=message_count, origin=SessionOrigin.MCP)
659 except (SessionNotFoundError, SessionForkRangeError) as exc:
660 return _error(str(exc))
661 return _session_payload(store.get(fork_id))
664@_tool_if(agent_sessions_enabled)
665def session_create(model_ref: str, scope: str = "both") -> dict[str, Any]:
666 """Start a saved chat session; returns its id."""
667 if not agent_sessions_enabled():
668 return _error(AGENT_SESSIONS_DISABLED_HINT)
669 session_id = get_services().session_store.create(
670 model_ref=model_ref, scope=scope, origin=SessionOrigin.MCP
671 )
672 return {"id": session_id, "model_ref": model_ref, "scope": scope}
675@_tool_if(agent_sessions_enabled)
676def session_add_message(
677 session_id: str,
678 role: MessageRole,
679 content: str,
680 sources: list[str] | None = None,
681 claim: bool = False,
682) -> dict[str, Any]:
683 """Append one turn; a foreign session errors unless claim=True (ask first)."""
684 if not agent_sessions_enabled():
685 return _error(AGENT_SESSIONS_DISABLED_HINT)
686 store = get_services().session_store
687 try:
688 # Re-coerce: the MCP layer passes the enum, but a raw string still
689 # arrives via direct library calls, and a bad one must error cleanly.
690 message = SessionMessage(
691 role=MessageRole(role), content=content, sources=tuple(sources or ())
692 )
693 if claim:
694 store.transfer(session_id, SessionOrigin.MCP)
695 store.add_message(session_id, message, surface=SessionOrigin.MCP)
696 except (SessionNotFoundError, SessionOwnershipError) as exc:
697 return _error(str(exc))
698 except ValueError as exc:
699 return _error(f"invalid role {role!r}: {exc}")
700 return {"id": session_id, "added": True}
703@_tool_if(agent_sessions_enabled)
704def session_set_summary(session_id: str, summary: str) -> dict[str, Any]:
705 """Replace an agent session's compaction summary."""
706 if not agent_sessions_enabled():
707 return _error(AGENT_SESSIONS_DISABLED_HINT)
708 try:
709 _require_agent_session(session_id)
710 get_services().session_store.set_summary(session_id, summary)
711 except SessionNotFoundError as exc:
712 return _error(str(exc))
713 return {"id": session_id, "summary": summary}
716@_tool_if(agent_sessions_enabled)
717def session_rename(session_id: str, title: str) -> dict[str, Any]:
718 """Rename an agent session."""
719 if not agent_sessions_enabled():
720 return _error(AGENT_SESSIONS_DISABLED_HINT)
721 try:
722 _require_agent_session(session_id)
723 get_services().session_store.set_title(session_id, title, TitleSource.CUSTOM)
724 except SessionNotFoundError as exc:
725 return _error(str(exc))
726 return {"id": session_id, "title": title}
729@_tool_if(agent_sessions_enabled)
730def session_delete(session_id: str) -> dict[str, Any]:
731 """Delete an agent session."""
732 if not agent_sessions_enabled():
733 return _error(AGENT_SESSIONS_DISABLED_HINT)
734 try:
735 _require_agent_session(session_id)
736 get_services().session_store.delete(session_id)
737 except SessionNotFoundError as exc:
738 return _error(str(exc))
739 return {"id": session_id, "deleted": True}
742@_tool
743def export_dataset(output: str, fmt: str = "", source: str = "") -> dict[str, Any]:
744 """Write the per-page {source, page, text} dataset to a file (no vectors).
746 ``fmt`` is parquet/jsonl (empty infers from the suffix); ``source`` limits to one file.
747 """
748 from lilbee.app.dataset import DatasetError, export_to_path
750 try:
751 summary = export_to_path(Path(output), fmt, source or None, cancel=_caller_cancelled())
752 except DatasetError as exc:
753 return _error(str(exc))
754 return summary.model_dump()
757@_tool
758async def import_dataset(dataset: str, fmt: str = "", ctx: Context | None = None) -> dict[str, Any]:
759 """Import a per-page text dataset, re-embedding under the current model.
761 Replaces existing copies; imported sources are detached so sync won't delete them.
762 """
763 from lilbee.app.dataset import DatasetError, import_from_path
764 from lilbee.runtime.progress import EmbedEvent, EventType, ProgressEvent
766 loop = asyncio.get_running_loop()
768 with _cancel_token() as cancel:
770 def on_progress(event_type: EventType, data: ProgressEvent) -> None:
771 # Raise rather than return: the embed work runs off the loop, so a
772 # quiet return would re-embed the whole dataset for a caller that
773 # has gone. Checked before the isinstance filter so the stop lands
774 # on every event, not only the ones that map to a percent.
775 if cancel.is_set():
776 raise TaskCancelledError
777 # EMBED events carry chunk/total_chunks; other event types don't map to a percent.
778 if ctx is None or not isinstance(data, EmbedEvent):
779 return
780 future = asyncio.run_coroutine_threadsafe(
781 ctx.report_progress(
782 progress=float(data.chunk), total=float(data.total_chunks), message=data.file
783 ),
784 loop,
785 )
786 future.add_done_callback(_log_progress_failure)
788 try:
789 summary = await import_from_path(Path(dataset), fmt, on_progress=on_progress)
790 except DatasetError as exc:
791 return _error(str(exc))
792 return summary.model_dump()
795@_tool
796def reset(confirm: bool = False) -> dict[str, Any]:
797 """Factory reset: delete all documents and indexed data. Requires ``confirm=true``."""
798 if _transport.http_mounted:
799 return _error(
800 "reset is unavailable on the HTTP server: it would wipe the shared index for "
801 "every connected client. Run it from the CLI or the stdio MCP server."
802 )
803 if not confirm:
804 return _error("pass confirm=true to confirm deletion")
805 from lilbee.app.reset import perform_reset
807 try:
808 result = perform_reset().model_dump()
809 except ResetRefusedError as exc:
810 return _error(str(exc))
811 # Reopen LanceDB against the empty data dir; keep providers loaded.
812 reset_store()
813 return result
816@_wiki_tool
817def wiki_lint(wiki_source: str = "") -> dict[str, Any]:
818 """Lint wiki pages; empty ``wiki_source`` lints all."""
819 from lilbee.wiki.lint import LintReport, lint_all, lint_wiki_page
821 store = get_services().store
822 report = (
823 LintReport(issues=lint_wiki_page(wiki_source, store))
824 if wiki_source
825 else lint_all(store, cancel=_caller_cancelled())
826 )
827 return {
828 "command": "wiki_lint",
829 "issues": [i.to_dict() for i in report.issues],
830 "total": len(report.issues),
831 "errors": report.error_count,
832 "warnings": report.warning_count,
833 }
836@_wiki_tool
837def wiki_citations(wiki_source: str = "", source: str = "") -> dict[str, Any]:
838 """List a wiki page's citations, or with ``source``, the pages citing that document.
840 Pass exactly one of ``wiki_source`` (forward) or ``source`` (reverse).
841 """
842 if bool(wiki_source) == bool(source):
843 return _error("pass either wiki_source or source, not both")
844 store = get_services().store
845 if source:
846 records = store.get_citations_for_source(source)
847 return {
848 "command": "wiki_citations",
849 "source": source,
850 "citations": [dict(r) for r in records],
851 "total": len(records),
852 }
853 records = store.get_citations_for_wiki(wiki_source)
854 return {
855 "command": "wiki_citations",
856 "wiki_source": wiki_source,
857 "citations": [dict(r) for r in records],
858 "total": len(records),
859 }
862@_tool
863def wiki_status() -> dict[str, Any]:
864 """Show wiki layer status: page counts, recent lint issues.
866 Registered even when the wiki is disabled, like the HTTP status route, so a
867 caller can read the disabled state instead of finding no tool at all.
868 """
869 from lilbee.wiki.lint import lint_all
871 wiki_root = cfg.data_root / cfg.wiki_dir
872 if not cfg.wiki or not wiki_root.exists():
873 # Same keys as the enabled arm so a client can key on lint_errors in
874 # either state; a disabled wiki reports zeros rather than being linted.
875 return {
876 "wiki_enabled": cfg.wiki,
877 WikiSubdir.SUMMARIES: 0,
878 WikiSubdir.DRAFTS: 0,
879 "pages": 0,
880 "lint_errors": 0,
881 "lint_warnings": 0,
882 }
884 summaries_dir = wiki_root / WikiSubdir.SUMMARIES
885 drafts_dir = wiki_root / WikiSubdir.DRAFTS
886 summaries = list(summaries_dir.rglob("*.md")) if summaries_dir.exists() else []
887 drafts = list(drafts_dir.rglob("*.md")) if drafts_dir.exists() else []
889 # Read-only status: lint for counts without appending to the audit log.
890 report = lint_all(get_services().store, record_log=False, cancel=_caller_cancelled())
891 return {
892 "wiki_enabled": cfg.wiki,
893 WikiSubdir.SUMMARIES: len(summaries),
894 WikiSubdir.DRAFTS: len(drafts),
895 "pages": total_wiki_pages(wiki_root),
896 "lint_errors": report.error_count,
897 "lint_warnings": report.warning_count,
898 }
901@_wiki_tool
902def wiki_list() -> dict[str, Any]:
903 """List wiki pages with metadata."""
904 from dataclasses import asdict
906 from lilbee.wiki.browse import list_pages
908 wiki_root = cfg.data_root / cfg.wiki_dir
909 pages = list_pages(wiki_root)
910 return {
911 "command": "wiki_list",
912 "pages": [asdict(p) for p in pages],
913 "total": len(pages),
914 }
917@_wiki_tool
918def wiki_read(slug: str) -> dict[str, Any]:
919 """Read a wiki page's content + frontmatter by slug."""
920 from dataclasses import asdict
922 from lilbee.wiki.browse import read_page
924 wiki_root = cfg.data_root / cfg.wiki_dir
925 result = read_page(wiki_root, slug)
926 if result is None:
927 return _error(f"wiki page not found: {slug}")
928 return {"command": "wiki_read", **asdict(result)}
931@_wiki_tool
932def wiki_build(dry_run: bool = False) -> dict[str, Any]:
933 """Build the concept and entity wiki across all ingested sources. Blocks until done.
935 ``dry_run=True`` returns the NER entity candidates a build would cover and
936 makes no LLM call.
937 """
938 from lilbee.wiki import run_full_build
939 from lilbee.wiki.generation import DRY_RUN_CONCEPT_NOTE, preview_build_entities
941 if dry_run:
942 rows = preview_build_entities(cfg)
943 return {
944 "command": "wiki_build",
945 "dry_run": True,
946 "entities": rows,
947 "count": len(rows),
948 "note": DRY_RUN_CONCEPT_NOTE,
949 }
950 return {"command": "wiki_build", **run_full_build(cfg, cancel=_caller_cancelled())}
953@_wiki_tool
954def wiki_update() -> dict[str, Any]:
955 """Refresh the concept and entity wiki after an ingest. A full rebuild; blocks until done."""
956 from lilbee.wiki import run_full_build
958 return {"command": "wiki_update", **run_full_build(cfg, cancel=_caller_cancelled())}
961@_wiki_tool
962def wiki_synthesize() -> dict[str, Any]:
963 """Generate synthesis pages for concept clusters with three or more sources."""
964 from lilbee.wiki import run_full_synthesize
966 return {"command": "wiki_synthesize", **run_full_synthesize(cfg, cancel=_caller_cancelled())}
969@_wiki_tool
970def wiki_prune() -> dict[str, Any]:
971 """Prune stale and orphaned wiki pages."""
972 from lilbee.wiki.prune import prune_wiki
974 report = prune_wiki(get_services().store, cancel=_caller_cancelled())
975 return {
976 "command": "wiki_prune",
977 "records": [r.to_dict() for r in report.records],
978 "archived": report.archived_count,
979 "flagged": report.flagged_count,
980 "reconciled": report.reconciled_count,
981 }
984@_wiki_tool
985def wiki_index() -> dict[str, Any]:
986 """Rebuild the browse index of pages the corpus could have. No LLM call."""
987 from lilbee.wiki.stubs import refresh_stub_index
989 stubs = refresh_stub_index(get_services().store)
990 return {"command": "wiki_index", "entries": len(stubs)}
993@_wiki_tool
994def wiki_generate(slug: str) -> dict[str, Any]:
995 """Generate one indexed wiki page. Costs a single LLM call and is GPU-heavy."""
996 from lilbee.wiki.browse import page_slug
997 from lilbee.wiki.lazy import UnknownStubError, generate_stub_page
999 try:
1000 path = generate_stub_page(slug, get_services().store)
1001 except UnknownStubError as exc:
1002 return _error(str(exc))
1003 if path is None:
1004 return _error(f"index entry for {slug} is stale; its sources are gone")
1005 # The read surfaces address pages by section, so answer with that slug.
1006 read_slug = page_slug(path, cfg.data_root / cfg.wiki_dir)
1007 return {"command": "wiki_generate", "slug": read_slug, "path": path.as_posix()}
1010@_tool
1011def wiki_wipe(confirm: bool = False) -> dict[str, Any]:
1012 """Delete every generated wiki page and its indexed rows. Pass ``confirm=true``.
1014 Registered even with the wiki disabled, because turning the setting off
1015 leaves the pages generated earlier in place.
1016 """
1017 if not confirm:
1018 return _error("pass confirm=true to delete the wiki; this cannot be undone")
1019 from lilbee.wiki.wipe import wipe_wiki
1021 report = wipe_wiki(get_services().store)
1022 if not report.rows_deleted:
1023 return _error(report.summary())
1024 return {
1025 "command": "wiki_wipe",
1026 "pages_removed": report.pages_removed,
1027 "sources_cleared": report.sources_cleared,
1028 }
1031def _setting_info_to_dict(info: SettingInfo) -> dict[str, Any]:
1032 """Render a SettingInfo as a JSON-safe dict for the MCP wire format."""
1033 return {
1034 "key": info.key,
1035 "value": _json_safe(info.value),
1036 "default": _json_safe(info.default),
1037 "type": info.type,
1038 "nullable": info.nullable,
1039 "group": info.group.value,
1040 "help": info.help_text,
1041 "choices": list(info.choices) if info.choices else None,
1042 "reindex_required": info.reindex_required,
1043 }
1046def _json_safe(value: Any) -> Any:
1047 """Coerce Path / frozenset / tuple to JSON-friendly primitives."""
1048 if isinstance(value, str | int | float | bool | list | type(None)):
1049 return value
1050 return str(value)
1053@_tool
1054def settings_list(group: str = "") -> dict[str, Any]:
1055 """List writable lilbee settings (each with value, default, type, help, choices).
1057 ``group`` filters by group name (case-insensitive); empty returns all.
1058 """
1060 try:
1061 infos = list_settings(group or None)
1062 except ValueError as exc:
1063 return _error(str(exc))
1064 return {
1065 "command": "settings_list",
1066 "settings": [_setting_info_to_dict(info) for info in infos],
1067 "total": len(infos),
1068 }
1071@_tool
1072def settings_get(key: str) -> dict[str, Any]:
1073 """Get a single setting's current value + metadata."""
1075 try:
1076 info = get_setting(key)
1077 except KeyError as exc:
1078 return _error(str(exc))
1079 return {"command": "settings_get", "setting": _setting_info_to_dict(info)}
1082@_tool
1083def settings_set(updates: dict[str, Any]) -> dict[str, Any]:
1084 """Atomically update writable settings; rolls back on validation error.
1085 Persists to config.toml; returns updated, reindex_required, warnings."""
1086 if _transport.http_mounted and requires_services_reset(updates):
1087 return _error(provider_reset_refused_message("Switching"))
1088 try:
1089 result = apply_settings_update(updates)
1090 except (ValueError, TypeError) as exc:
1091 return _error(str(exc))
1092 except OSError as exc:
1093 return _error(config_write_failure_message(exc))
1094 return {
1095 "command": "settings_set",
1096 "updated": result.updated,
1097 "reindex_required": result.reindex_required,
1098 "warnings": list(result.warnings),
1099 }
1102@_tool
1103def settings_reset(keys: list[str]) -> dict[str, Any]:
1104 """Reset writable settings to their built-in defaults."""
1105 if _transport.http_mounted and requires_services_reset(dict.fromkeys(keys)):
1106 return _error(provider_reset_refused_message("Resetting"))
1107 try:
1108 result = reset_settings(keys)
1109 except (ValueError, TypeError) as exc:
1110 return _error(str(exc))
1111 except OSError as exc:
1112 return _error(config_write_failure_message(exc))
1113 return {
1114 "command": "settings_reset",
1115 "updated": result.updated,
1116 "reindex_required": result.reindex_required,
1117 }
1120@_tool
1121def model_list(source: str = "", task: str = "") -> dict[str, Any]:
1122 """List installed models. ``source`` is ``native`` / ``remote``; ``task`` filters by role."""
1123 from lilbee.app.models import list_models_data
1124 from lilbee.catalog.types import ModelTask
1126 try:
1127 src = ModelSource.parse(source)
1128 except ValueError as exc:
1129 return _error(str(exc))
1130 try:
1131 parsed_task = ModelTask(task) if task else None
1132 except ValueError as exc:
1133 return _error(str(exc))
1134 return list_models_data(source=src, task=parsed_task).model_dump()
1137@_tool
1138def catalog_browse(
1139 task: str = "",
1140 search: str = "",
1141 size: str = "",
1142 installed: bool | None = None,
1143 featured: bool | None = None,
1144 max_fit: str = "",
1145 sort: str = "featured",
1146 limit: int = 20,
1147 offset: int = 0,
1148) -> dict[str, Any]:
1149 """Browse the model catalog. ``task``: chat/embedding/vision/rerank.
1150 ``size``: small/medium/large/huge, by parameter count.
1151 ``max_fit``: fits/tight/wont_run, worst fit to keep.
1152 ``sort``: featured/downloads/name/size_asc/size_desc."""
1153 from lilbee.catalog.query import get_catalog
1154 from lilbee.catalog.types import CatalogSize, CatalogSort, ModelTask
1156 try:
1157 parsed_task = ModelTask(task) if task else None
1158 parsed_size = CatalogSize(size) if size else None
1159 parsed_sort = CatalogSort(sort)
1160 parsed_max_fit = FitLevel(max_fit) if max_fit else None
1161 except ValueError as exc:
1162 return _error(str(exc))
1163 # The rows here carry no fit chip; only the filter needs the uncached GPU probe.
1164 available_bytes = available_memory_for_fit() if parsed_max_fit is not None else None
1165 try:
1166 result = get_catalog(
1167 task=parsed_task,
1168 search=search,
1169 size=parsed_size,
1170 installed=installed,
1171 featured=featured,
1172 fit_filter=make_fit_filter(parsed_max_fit, available_bytes),
1173 sort=parsed_sort,
1174 limit=limit,
1175 offset=offset,
1176 model_manager=get_services().model_manager,
1177 )
1178 except ValueError as exc:
1179 return _error(str(exc))
1180 return {
1181 "command": "catalog_browse",
1182 "total": result.total,
1183 "limit": result.limit,
1184 "offset": result.offset,
1185 "has_more": result.has_more,
1186 "truncated": result.truncated,
1187 "models": [
1188 {
1189 "ref": m.hf_repo,
1190 "display_name": m.display_name,
1191 "task": m.task.value,
1192 "size_gb": m.size_gb,
1193 "min_ram_gb": m.min_ram_gb,
1194 "downloads": m.downloads,
1195 "featured": m.featured,
1196 "description": m.description,
1197 "architecture": m.architecture,
1198 "compat": m.compat.value,
1199 "safety_stripped": m.safety_stripped,
1200 }
1201 for m in result.models
1202 ],
1203 }
1206@_tool
1207def model_show(model: str) -> dict[str, Any]:
1208 """Show catalog and installed metadata for a model ref."""
1209 from lilbee.app.models import show_model_data
1210 from lilbee.modelhub.model_manager import ModelNotFoundError
1212 try:
1213 return show_model_data(model).model_dump()
1214 except ModelNotFoundError as exc:
1215 return _error(str(exc))
1218def _log_progress_failure(future: concurrent.futures.Future[None]) -> None:
1219 """Log report_progress failures without raising.
1221 Progress notifications are best-effort: a failure should not abort
1222 an in-flight pull.
1223 """
1224 try:
1225 future.result()
1226 except Exception:
1227 log.warning("MCP report_progress failed", exc_info=True)
1230@_tool
1231async def model_pull(
1232 model: str,
1233 source: str = ModelSource.NATIVE.value,
1234 allow_unsupported: bool = False,
1235 ctx: Context | None = None,
1236) -> dict[str, Any]:
1237 """Download a model and stream progress.
1239 ``source`` is ``native`` (GGUF) or ``remote`` (SDK).
1240 ``allow_unsupported`` overrides the supported-architecture refusal.
1241 """
1242 from lilbee.app.models import pull_model_data
1243 from lilbee.catalog import DownloadProgress
1244 from lilbee.catalog.compat import SUPPORTED_ARCHS, UnsupportedArchError
1246 try:
1247 src = ModelSource.parse(source) or ModelSource.NATIVE
1248 except ValueError as exc:
1249 return _error(str(exc))
1251 loop = asyncio.get_running_loop()
1253 with _cancel_token() as cancel:
1255 def on_update(p: DownloadProgress) -> None:
1256 # Raise rather than return: the download runs on a thread asyncio
1257 # cannot interrupt, so returning would leave a multi-GB pull going
1258 # for a caller that has gone. Same idiom the HTTP pull uses.
1259 if cancel.is_set():
1260 raise TaskCancelledError
1261 if ctx is None:
1262 return
1263 future = asyncio.run_coroutine_threadsafe(
1264 ctx.report_progress(progress=float(p.percent), total=100.0, message=p.detail),
1265 loop,
1266 )
1267 future.add_done_callback(_log_progress_failure)
1269 try:
1270 result = await asyncio.to_thread(
1271 pull_model_data,
1272 model,
1273 src,
1274 on_update=on_update,
1275 allow_unsupported=allow_unsupported,
1276 cancel=cancel,
1277 )
1278 except UnsupportedArchError as exc:
1279 return {
1280 "ok": False,
1281 "command": "model_pull",
1282 "error": {
1283 "code": "unsupported_arch",
1284 "arch": exc.architecture,
1285 "ref": exc.ref,
1286 "supported_examples": sorted(SUPPORTED_ARCHS)[:5],
1287 "total_supported": len(SUPPORTED_ARCHS),
1288 },
1289 }
1290 except (RuntimeError, PermissionError) as exc:
1291 return _error(str(exc))
1292 return result.model_dump()
1295@_tool
1296def model_rm(model: str, source: str = "") -> dict[str, Any]:
1297 """Remove an installed model. Only native GGUF models lilbee downloaded;
1298 Ollama/LM Studio are read-only."""
1299 from lilbee.app.models import remove_model_data
1301 try:
1302 src = ModelSource.parse(source)
1303 return remove_model_data(model, source=src).model_dump()
1304 except ValueError as exc:
1305 return _error(str(exc))
1308@_wiki_tool
1309def wiki_drafts_list() -> dict[str, Any]:
1310 """List pending wiki drafts. Read-only: promotion is reserved for the human
1311 surfaces (CLI, TUI, and the authenticated HTTP API)."""
1312 from lilbee.wiki.drafts import list_drafts
1314 wiki_root = cfg.data_root / cfg.wiki_dir
1315 drafts = list_drafts(wiki_root)
1316 return {
1317 "command": "wiki_drafts_list",
1318 "drafts": [d.to_dict() for d in drafts],
1319 "total": len(drafts),
1320 }
1323@_wiki_tool
1324def wiki_drafts_diff(slug: str) -> dict[str, Any]:
1325 """Unified diff of a draft against its published counterpart."""
1326 from lilbee.core.security import PathTraversalError
1327 from lilbee.wiki.drafts import diff_draft
1329 wiki_root = cfg.data_root / cfg.wiki_dir
1330 try:
1331 diff = diff_draft(slug, wiki_root)
1332 except FileNotFoundError as exc:
1333 return _error(str(exc))
1334 except PathTraversalError:
1335 return _error(INVALID_DRAFT_SLUG_ERROR)
1336 return {"command": "wiki_drafts_diff", "slug": slug, "diff": diff}
1339def _collapse_nullable_anyof(prop: dict[str, Any]) -> None:
1340 """Collapse ``anyOf: [{type: X}, {type: null}]`` to ``{type: X}`` in place.
1342 Pydantic emits ``T | None`` parameters as a two-arm anyOf with a null
1343 branch. The null branch carries no information the model needs to pick
1344 or shape its call, but it costs tokens at every dispatch. Drop it.
1345 """
1346 arms = prop.get("anyOf")
1347 if not isinstance(arms, list):
1348 return
1349 non_null = [a for a in arms if isinstance(a, dict) and a.get("type") != "null"]
1350 if len(non_null) == 1 and len(non_null) < len(arms):
1351 prop.pop("anyOf", None)
1352 for key, value in non_null[0].items():
1353 prop.setdefault(key, value)
1356def _strip_property_noise(prop: dict[str, Any]) -> None:
1357 """Drop tokens that don't change the model's behavior."""
1358 prop.pop("title", None)
1359 prop.pop("default", None)
1360 _collapse_nullable_anyof(prop)
1361 if prop.get("additionalProperties") is True:
1362 prop.pop("additionalProperties", None)
1365def _flatten_tool_description(text: str) -> str:
1366 """Flatten a triple-quoted tool docstring for the tools wire.
1368 The summary line carries no indent while continuation lines are indented to
1369 the source, so ``textwrap.dedent`` alone is a no-op (the common prefix is the
1370 empty string) and leaves source indentation on every body line -- including
1371 deeper-indented Args lines. Strip each line so the model sees flat text;
1372 blank lines are kept so paragraph breaks survive.
1373 """
1374 return "\n".join(line.strip() for line in text.strip().splitlines())
1377def _strip_schema(schema: dict[str, Any]) -> dict[str, Any]:
1378 """Trim auto-generated noise from a tool's input schema, on a copy.
1380 Drops:
1381 - SDK/Pydantic ``title`` keys (per-schema + per-property). Tools the
1382 model picks by name don't need a separate display title.
1383 - ``default`` values on properties: clients send what they want and
1384 omitted fields fall back server-side.
1385 - ``additionalProperties: true`` on dict params: Pydantic emits it for
1386 every ``dict[str, Any]`` but it's the JSON Schema default behavior.
1387 - The ``null`` arm of ``anyOf: [{type: X}, {type: null}]`` unions for
1388 ``T | None`` defaults; the null branch is implicit.
1390 A roughly 25-35% reduction in the serialized tools payload, which matters
1391 most for small-context (16K) chat models where the tools surface was
1392 previously eating ~60% of the budget.
1393 """
1394 schema = deepcopy(schema)
1395 schema.pop("title", None)
1396 properties = schema.get("properties")
1397 if isinstance(properties, dict):
1398 for prop in properties.values():
1399 if isinstance(prop, dict):
1400 _strip_property_noise(prop)
1401 return schema
1404_NO_WIKI_SCOPE_HINT = ' No wiki layer here: use scope "raw" or "both".'
1407class LilbeeMCP(MCPServer):
1408 """MCP server that trims its tools wire and keeps it current with config."""
1410 async def list_tools(self) -> list[MCPTool]:
1411 """The registered tools with schema noise stripped and flat descriptions.
1413 The transforms run on the wire representation per request, never on the
1414 stored registrations, so they cannot drift out of sync with config. The
1415 ``search`` description advertises only the scopes this corpus has: when
1416 wiki generation is off, a model that guesses ``scope="wiki"`` gets a
1417 silent fallback to the full pool, so raw/both only.
1418 """
1419 tools = await super().list_tools()
1420 for tool in tools:
1421 tool.input_schema = _strip_schema(tool.input_schema)
1422 if isinstance(tool.description, str):
1423 description = _flatten_tool_description(tool.description)
1424 if tool.name == "search" and not cfg.wiki:
1425 description += _NO_WIKI_SCOPE_HINT
1426 tool.description = description
1427 return tools
1430def _client_name(ctx: Context | None) -> str:
1431 """The MCP client's self-reported name from the initialize handshake, or empty."""
1432 if ctx is None:
1433 return ""
1434 params = ctx.session.client_params
1435 return params.clientInfo.name if params is not None else ""
1438def _slug(value: str) -> str:
1439 """Lowercase, hyphenated id fragment; falls back to ``generic`` when empty."""
1440 slug = re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-")
1441 return slug or "generic"
1444# Per-connection fallback ids for agents that report no identity. Keyed by the
1445# live MCP session so each connection gets a distinct, stable namespace instead
1446# of every unidentified agent colliding on a shared one. WeakKeyDictionary drops
1447# entries once the session is collected, so this does not grow unbounded. The
1448# lock guards the get-or-create because sync tool handlers run on the offload
1449# threadpool, so concurrent connections (and weakref-removal callbacks) can
1450# touch the mapping from different threads.
1451_ANON_OWNER_IDS: WeakKeyDictionary[object, str] = WeakKeyDictionary()
1452_ANON_OWNER_LOCK = threading.Lock()
1455def _anon_owner_id(ctx: Context | None) -> str:
1456 """A stable per-connection id for an agent that reported no identity.
1458 Without this, two unidentified agents would both slug to ``generic`` and
1459 share a memory namespace; keying on the session keeps them isolated.
1460 """
1461 if ctx is None:
1462 return "anonymous"
1463 session = ctx.session
1464 with _ANON_OWNER_LOCK:
1465 existing = _ANON_OWNER_IDS.get(session)
1466 if existing is None:
1467 existing = f"anon-{uuid.uuid4().hex[:12]}"
1468 _ANON_OWNER_IDS[session] = existing
1469 return existing
1472def _derive_owner(agent_id: str, ctx: Context | None) -> str:
1473 """Resolve the calling agent's stable owner namespace.
1475 Precedence: explicit ``agent_id`` argument, then the ``LILBEE_AGENT_ID`` env var
1476 (pinned in the client's MCP config), then the MCP client name, then a stable
1477 per-connection fallback so unidentified agents never share a namespace.
1478 """
1479 explicit = agent_id or os.environ.get("LILBEE_AGENT_ID", "")
1480 resolved = explicit or _client_name(ctx)
1481 if resolved:
1482 return agent_owner(_slug(resolved))
1483 return agent_owner(_slug(_anon_owner_id(ctx)))
1486@_tool_if(memory_enabled)
1487def memory_remember(
1488 text: str,
1489 kind: MemoryKind = MemoryKind.FACT,
1490 shared: bool = False,
1491 agent_id: str = "",
1492 ctx: Context | None = None,
1493) -> dict[str, Any]:
1494 """Store a durable memory. ``kind``: "fact" (similarity-recalled) or "preference" (always on).
1495 ``shared`` exposes it to the human's TUI/CLI."""
1496 if not memory_enabled():
1497 return _error(MEMORY_DISABLED_HINT)
1498 owner = _derive_owner(agent_id, ctx)
1499 memory_id = remember(text, owner=owner, kind=kind, source=MemorySource.AGENT, shared=shared)
1500 return {"ok": True, "id": memory_id, "owner": owner}
1503@_tool_if(memory_enabled)
1504def memory_recall(
1505 query: str, limit: int = 0, agent_id: str = "", ctx: Context | None = None
1506) -> dict[str, Any]:
1507 """Recall this agent's memories (plus any the human shared) relevant to *query*."""
1508 if not memory_enabled():
1509 return _error(MEMORY_DISABLED_HINT)
1510 owner = _derive_owner(agent_id, ctx)
1511 memories = recall(query, owner, top_k=limit if limit > 0 else None)
1512 return {
1513 "memories": [
1514 {"id": m.id, "text": m.text, "kind": m.kind.value, "owner": m.owner} for m in memories
1515 ]
1516 }
1519@_tool_if(memory_enabled)
1520def memory_list(agent_id: str = "", ctx: Context | None = None) -> dict[str, Any]:
1521 """List every memory in this agent's namespace (any kind, newest first)."""
1522 if not memory_enabled():
1523 return _error(MEMORY_DISABLED_HINT)
1524 owner = _derive_owner(agent_id, ctx)
1525 memories = list_memories(owner)
1526 return {
1527 "memories": [
1528 {"id": m.id, "text": m.text, "kind": m.kind.value, "shared": m.shared} for m in memories
1529 ]
1530 }
1533@_tool_if(memory_enabled)
1534def memory_forget(memory_id: str, agent_id: str = "", ctx: Context | None = None) -> dict[str, Any]:
1535 """Delete one of this agent's own memories by id (agent_id scopes the namespace)."""
1536 if not memory_enabled():
1537 return _error(MEMORY_DISABLED_HINT)
1538 owner = _derive_owner(agent_id, ctx)
1539 if not forget(memory_id, owner=owner):
1540 return _error(f"No memory '{memory_id}' in this agent's namespace.")
1541 return {"ok": True, "id": memory_id}
1544def _placement_dict(view: PlacementView) -> dict[str, Any]:
1545 from lilbee.server.models import PlacementResponse
1547 return PlacementResponse.from_view(view).model_dump(mode="json")
1550def _placement_guard(serialize: Callable[[], dict[str, Any]]) -> dict[str, Any]:
1551 """Run a placement query and serialize it, returning a structured error on failure."""
1552 from lilbee.providers.base import ProviderError
1553 from lilbee.providers.fleet.placement_spec import PlacementError
1555 try:
1556 return serialize()
1557 except (PlacementError, ProviderError) as exc:
1558 return _error(str(exc))
1561def _placement_result(action: Callable[[], PlacementView]) -> dict[str, Any]:
1562 """Run a placement action and serialize its view, returning a structured error on failure."""
1563 return _placement_guard(lambda: _placement_dict(action()))
1566def _parse_spec(spec: dict[str, Any] | None) -> PlacementSpec | None:
1567 from lilbee.providers.fleet.placement_spec import PlacementSpec
1569 return PlacementSpec.from_json(json.dumps(spec)) if spec else None
1572@_tool_named("get_gpus")
1573def get_gpus_tool() -> dict[str, Any]:
1574 """List detected GPUs with free/total VRAM (the placement HTTP /api/gpus equivalent)."""
1576 def _body() -> dict[str, Any]:
1577 from lilbee.cli.tui import messages as msg
1578 from lilbee.providers.fleet.gpu_stats import probe_intel_util_hint
1580 view = get_placement()
1581 hint = probe_intel_util_hint(view.gpus)
1582 return {
1583 "gpus": _placement_dict(view)["gpus"],
1584 "notice": msg.intel_util_hint_text(hint) if hint else None,
1585 }
1587 return _placement_guard(_body)
1590@_tool_named("get_placement")
1591def get_placement_tool() -> dict[str, Any]:
1592 """Show the current effective multi-GPU model placement."""
1593 return _placement_result(get_placement)
1596@_tool_named("preview_placement")
1597def preview_placement_tool(spec: dict[str, Any] | None = None) -> dict[str, Any]:
1598 """Preview what a placement spec (or auto, when omitted) would place. No changes made."""
1599 return _placement_result(lambda: preview_placement(_parse_spec(spec)))
1602@_tool_named("set_placement")
1603def set_placement_tool(spec: dict[str, Any]) -> dict[str, Any]:
1604 """Set and apply a manual multi-GPU placement spec (persists to config).
1606 The spec maps a role ("chat"/"embed"/"rerank"/"vision") to a placement, e.g.
1607 ``{"chat": {"devices": [0, 1], "tensor_split": [1, 1]}}``. ``devices`` is the
1608 GPU indices (get_gpus lists them); ``tensor_split`` is optional per-device
1609 weights (omit for an even split). Omit a role to leave it auto-placed.
1610 """
1611 from lilbee.providers.fleet.placement_spec import PlacementSpec
1613 # set_placement restarts the shared fleet's moved roles: gate it on the
1614 # shared HTTP transport exactly like the REST PUT/DELETE placement routes.
1615 if _transport.http_mounted and not cfg.allow_http_placement:
1616 return _error(placement_refused_message())
1617 # Always build a spec (even {}) so an empty/invalid one is rejected, not cleared.
1618 return _placement_result(lambda: set_placement(PlacementSpec.from_json(json.dumps(spec))))
1621@_tool_named("clear_placement")
1622def clear_placement_tool() -> dict[str, Any]:
1623 """Clear the manual placement and return to automatic placement."""
1624 if _transport.http_mounted and not cfg.allow_http_placement:
1625 return _error(placement_refused_message())
1626 return _placement_result(lambda: set_placement(None))
1629def build_mcp_server() -> LilbeeMCP:
1630 """Build an MCP server carrying every tool registered in this module.
1632 Each transport builds its own instance: the SDK caches one
1633 ``StreamableHTTPSessionManager`` per server and its ``run()`` is single-use,
1634 so a shared server cannot back two apps in one process. Gates registered
1635 via ``_tool_if`` are evaluated here, against current config.
1636 """
1637 server = LilbeeMCP("lilbee", instructions=_INSTRUCTIONS)
1638 for fn, name, gate in _REGISTRATIONS:
1639 if gate is None or gate():
1640 server.add_tool(fn, name=name)
1641 return server
1644_PARENT_DEATH_CLEANUP_S = 5.0
1647def _exit_on_parent_death() -> None:
1648 """Release engine membership best-effort, then hard-exit promptly.
1650 ``os._exit`` skips atexit, so an explicit release stops this process's engine
1651 when it was the last user and keeps the machine clean. But this watchdog's one
1652 contract is to exit promptly on parent death, so the release runs on a daemon
1653 thread joined with a short deadline: a peer holding an engine build lock can
1654 never keep the orphaned process alive with its models resident. The kernel
1655 releases this process's user lock on exit regardless, so a skipped release only
1656 defers the engine stop to the peers' reap and the idle TTL.
1657 """
1658 cleanup = threading.Thread(target=reset_services, name="parent-death-cleanup", daemon=True)
1659 cleanup.start()
1660 cleanup.join(timeout=_PARENT_DEATH_CLEANUP_S)
1661 os._exit(0)
1664def main() -> None:
1665 """Entry point for the stdio MCP server."""
1666 # Preload so the first tool call doesn't pay the cold-start cost
1667 # of provider/embedder/store init. Failures (missing model, bad
1668 # config) still surface on the first tool call rather than crashing
1669 # the server before it attaches to stdio.
1670 try:
1671 get_services()
1672 except Exception:
1673 log.debug("MCP pre-warm failed; services will init on first call", exc_info=True)
1675 from lilbee.parent_monitor import parse_parent_pid, watch_parent_thread
1677 parent_pid = parse_parent_pid()
1678 if parent_pid is not None:
1679 watch_parent_thread(parent_pid, _exit_on_parent_death)
1681 build_mcp_server().run()