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

1"""MCP server exposing lilbee as tools for AI agents.""" 

2 

3from __future__ import annotations 

4 

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 

23 

24import anyio 

25from mcp.server.mcpserver import Context, MCPServer 

26from mcp.types import Tool as MCPTool 

27 

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) 

91 

92if TYPE_CHECKING: 

93 from lilbee.providers.fleet.placement_spec import PlacementSpec 

94 

95log = logging.getLogger(__name__) 

96 

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) 

104 

105 

106class _TransportState: 

107 """Process-level MCP transport facts. 

108 

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 """ 

116 

117 http_mounted: bool = False 

118 

119 

120_transport = _TransportState() 

121 

122 

123def set_http_mounted(value: bool) -> None: 

124 """Mark whether this process serves MCP over the shared HTTP daemon.""" 

125 _transport.http_mounted = value 

126 

127 

128_F = TypeVar("_F", bound=Callable[..., Any]) 

129 

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) 

134 

135 

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() 

139 

140 

141def _offload_sync(fn: _F) -> _F: 

142 """Run a sync tool handler off the event loop; async handlers pass through. 

143 

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. 

148 

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 

158 

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 ) 

166 

167 return cast("_F", _runner) 

168 

169 

170# (handler, wire name override, gate) applied by build_mcp_server per instance. 

171_REGISTRATIONS: list[tuple[Callable[..., Any], str | None, Callable[[], bool] | None]] = [] 

172 

173 

174def _tool(fn: _F) -> _F: 

175 """Register *fn* as an MCP tool with sync handlers offloaded off the loop. 

176 

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 

182 

183 

184def _tool_named(name: str) -> Callable[[_F], _F]: 

185 """Register an MCP tool under an explicit wire *name* (sync handlers offloaded).""" 

186 

187 def deco(fn: _F) -> _F: 

188 _REGISTRATIONS.append((_offload_sync(fn), name, None)) 

189 return fn 

190 

191 return deco 

192 

193 

194def _tool_if(when: Callable[[], bool]) -> Callable[[_F], _F]: 

195 """Register an MCP tool gated on *when*, evaluated at server-build time. 

196 

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") 

203 

204 def deco(fn: _F) -> _F: 

205 _REGISTRATIONS.append((_offload_sync(fn), None, when)) 

206 return fn 

207 

208 return deco 

209 

210 

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"}) 

216 

217 

218def _wiki_enabled() -> bool: 

219 return cfg.wiki 

220 

221 

222def _error(msg: str) -> dict[str, Any]: 

223 """Uniform error envelope MCP tool handlers return on a failure path. 

224 

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} 

230 

231 

232def _wiki_tool(fn: _F) -> _F: 

233 """Register a wiki MCP tool: gated at build time, re-checked per call. 

234 

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 """ 

238 

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) 

244 

245 checked = cast("_F", guarded) 

246 _tool_if(_wiki_enabled)(checked) 

247 return checked 

248 

249 

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)) 

290 

291 

292@_tool 

293def status() -> dict[str, Any]: 

294 """Show indexed documents, configuration, and chunk counts.""" 

295 from lilbee.app.status import gather_status 

296 

297 return gather_status().model_dump() 

298 

299 

300@contextmanager 

301def _cancel_token() -> Iterator[threading.Event]: 

302 """A stop token for a long operation, set on any abnormal exit. 

303 

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. 

310 

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 

320 

321 

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. 

331 

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 

338 

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() 

353 

354 

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 

361 

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() 

366 

367 

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 

373 

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 

388 

389 

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 

403 

404 try: 

405 validate_ocr_timeout(ocr_timeout) 

406 except ValueError as exc: 

407 return _error(str(exc)) 

408 

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) 

428 

429 # Crawl URLs 

430 crawled_count = 0 

431 if urls: 

432 from lilbee.crawler import crawler_available 

433 

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) 

438 

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 

464 

465 

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 

478 

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)) 

494 

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} 

503 

504 

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 } 

523 

524 

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 

529 

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)} 

533 

534 

535@_tool 

536def init(path: str = "") -> dict[str, Any]: 

537 """Initialize a local ``.lilbee/`` knowledge base; empty path = cwd. 

538 

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 

550 

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 

557 

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() 

569 

570 return {"command": "init", "path": str(root), "created": created} 

571 

572 

573@_tool 

574def remove(names: list[str]) -> dict[str, Any]: 

575 """Remove documents from the index by source name, folder, or glob pattern. 

576 

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 

580 

581 result = remove_documents_durably(names) 

582 return {"command": "remove", "removed": result.removed, "not_found": result.not_found} 

583 

584 

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 } 

595 

596 

597_AGENT_ORIGINS = frozenset({SessionOrigin.MCP}) 

598 

599 

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 

607 

608 

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)} 

616 

617 

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) 

628 

629 

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 } 

648 

649 

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)) 

662 

663 

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} 

673 

674 

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} 

701 

702 

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} 

714 

715 

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} 

727 

728 

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} 

740 

741 

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). 

745 

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 

749 

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() 

755 

756 

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. 

760 

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 

765 

766 loop = asyncio.get_running_loop() 

767 

768 with _cancel_token() as cancel: 

769 

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) 

787 

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() 

793 

794 

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 

806 

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 

814 

815 

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 

820 

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 } 

834 

835 

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. 

839 

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 } 

860 

861 

862@_tool 

863def wiki_status() -> dict[str, Any]: 

864 """Show wiki layer status: page counts, recent lint issues. 

865 

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 

870 

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 } 

883 

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 [] 

888 

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 } 

899 

900 

901@_wiki_tool 

902def wiki_list() -> dict[str, Any]: 

903 """List wiki pages with metadata.""" 

904 from dataclasses import asdict 

905 

906 from lilbee.wiki.browse import list_pages 

907 

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 } 

915 

916 

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 

921 

922 from lilbee.wiki.browse import read_page 

923 

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)} 

929 

930 

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. 

934 

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 

940 

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())} 

951 

952 

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 

957 

958 return {"command": "wiki_update", **run_full_build(cfg, cancel=_caller_cancelled())} 

959 

960 

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 

965 

966 return {"command": "wiki_synthesize", **run_full_synthesize(cfg, cancel=_caller_cancelled())} 

967 

968 

969@_wiki_tool 

970def wiki_prune() -> dict[str, Any]: 

971 """Prune stale and orphaned wiki pages.""" 

972 from lilbee.wiki.prune import prune_wiki 

973 

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 } 

982 

983 

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 

988 

989 stubs = refresh_stub_index(get_services().store) 

990 return {"command": "wiki_index", "entries": len(stubs)} 

991 

992 

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 

998 

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()} 

1008 

1009 

1010@_tool 

1011def wiki_wipe(confirm: bool = False) -> dict[str, Any]: 

1012 """Delete every generated wiki page and its indexed rows. Pass ``confirm=true``. 

1013 

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 

1020 

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 } 

1029 

1030 

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 } 

1044 

1045 

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) 

1051 

1052 

1053@_tool 

1054def settings_list(group: str = "") -> dict[str, Any]: 

1055 """List writable lilbee settings (each with value, default, type, help, choices). 

1056 

1057 ``group`` filters by group name (case-insensitive); empty returns all. 

1058 """ 

1059 

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 } 

1069 

1070 

1071@_tool 

1072def settings_get(key: str) -> dict[str, Any]: 

1073 """Get a single setting's current value + metadata.""" 

1074 

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)} 

1080 

1081 

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 } 

1100 

1101 

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 } 

1118 

1119 

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 

1125 

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() 

1135 

1136 

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 

1155 

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 } 

1204 

1205 

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 

1211 

1212 try: 

1213 return show_model_data(model).model_dump() 

1214 except ModelNotFoundError as exc: 

1215 return _error(str(exc)) 

1216 

1217 

1218def _log_progress_failure(future: concurrent.futures.Future[None]) -> None: 

1219 """Log report_progress failures without raising. 

1220 

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) 

1228 

1229 

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. 

1238 

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 

1245 

1246 try: 

1247 src = ModelSource.parse(source) or ModelSource.NATIVE 

1248 except ValueError as exc: 

1249 return _error(str(exc)) 

1250 

1251 loop = asyncio.get_running_loop() 

1252 

1253 with _cancel_token() as cancel: 

1254 

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) 

1268 

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() 

1293 

1294 

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 

1300 

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)) 

1306 

1307 

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 

1313 

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 } 

1321 

1322 

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 

1328 

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} 

1337 

1338 

1339def _collapse_nullable_anyof(prop: dict[str, Any]) -> None: 

1340 """Collapse ``anyOf: [{type: X}, {type: null}]`` to ``{type: X}`` in place. 

1341 

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) 

1354 

1355 

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) 

1363 

1364 

1365def _flatten_tool_description(text: str) -> str: 

1366 """Flatten a triple-quoted tool docstring for the tools wire. 

1367 

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()) 

1375 

1376 

1377def _strip_schema(schema: dict[str, Any]) -> dict[str, Any]: 

1378 """Trim auto-generated noise from a tool's input schema, on a copy. 

1379 

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. 

1389 

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 

1402 

1403 

1404_NO_WIKI_SCOPE_HINT = ' No wiki layer here: use scope "raw" or "both".' 

1405 

1406 

1407class LilbeeMCP(MCPServer): 

1408 """MCP server that trims its tools wire and keeps it current with config.""" 

1409 

1410 async def list_tools(self) -> list[MCPTool]: 

1411 """The registered tools with schema noise stripped and flat descriptions. 

1412 

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 

1428 

1429 

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 "" 

1436 

1437 

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" 

1442 

1443 

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() 

1453 

1454 

1455def _anon_owner_id(ctx: Context | None) -> str: 

1456 """A stable per-connection id for an agent that reported no identity. 

1457 

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 

1470 

1471 

1472def _derive_owner(agent_id: str, ctx: Context | None) -> str: 

1473 """Resolve the calling agent's stable owner namespace. 

1474 

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))) 

1484 

1485 

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} 

1501 

1502 

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 } 

1517 

1518 

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 } 

1531 

1532 

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} 

1542 

1543 

1544def _placement_dict(view: PlacementView) -> dict[str, Any]: 

1545 from lilbee.server.models import PlacementResponse 

1546 

1547 return PlacementResponse.from_view(view).model_dump(mode="json") 

1548 

1549 

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 

1554 

1555 try: 

1556 return serialize() 

1557 except (PlacementError, ProviderError) as exc: 

1558 return _error(str(exc)) 

1559 

1560 

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())) 

1564 

1565 

1566def _parse_spec(spec: dict[str, Any] | None) -> PlacementSpec | None: 

1567 from lilbee.providers.fleet.placement_spec import PlacementSpec 

1568 

1569 return PlacementSpec.from_json(json.dumps(spec)) if spec else None 

1570 

1571 

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).""" 

1575 

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 

1579 

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 } 

1586 

1587 return _placement_guard(_body) 

1588 

1589 

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) 

1594 

1595 

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))) 

1600 

1601 

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). 

1605 

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 

1612 

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)))) 

1619 

1620 

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)) 

1627 

1628 

1629def build_mcp_server() -> LilbeeMCP: 

1630 """Build an MCP server carrying every tool registered in this module. 

1631 

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 

1642 

1643 

1644_PARENT_DEATH_CLEANUP_S = 5.0 

1645 

1646 

1647def _exit_on_parent_death() -> None: 

1648 """Release engine membership best-effort, then hard-exit promptly. 

1649 

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) 

1662 

1663 

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) 

1674 

1675 from lilbee.parent_monitor import parse_parent_pid, watch_parent_thread 

1676 

1677 parent_pid = parse_parent_pid() 

1678 if parent_pid is not None: 

1679 watch_parent_thread(parent_pid, _exit_on_parent_death) 

1680 

1681 build_mcp_server().run()