Coverage for src/local_deep_research/web_search_engines/search_engines_config.py: 92%
204 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-10-02 13:53 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-10-02 13:53 +0000
1"""
2Configuration file for search engines.
3Loads search engine definitions from the user's configuration.
4"""
6import copy
7import threading
8import time
9from typing import Any, Dict, Optional
10from cachetools import LRUCache
11from sqlalchemy.orm import Session
13from ..security.secure_logging import logger
15from ..config.thread_settings import get_setting_from_snapshot
16from ..security.log_sanitizer import scrub_error
17from ..utilities.db_utils import get_settings_manager
18from .search_engine_base import _is_api_key_placeholder
20_COLLECTION_CACHE_LOCK = threading.Lock()
21_COLLECTION_CACHE_MAXSIZE = 50
22_COLLECTION_CACHE_TTL = 30.0 # seconds
23_COLLECTION_CACHE_MAX_STALE = 300.0 # seconds (5 minutes stale age cap)
24# username -> (timestamp, {engine_id: engine_dict})
25_COLLECTION_ENGINES_CACHE: LRUCache[str, tuple[float, Dict[str, Any]]] = (
26 LRUCache(maxsize=_COLLECTION_CACHE_MAXSIZE)
27)
30def invalidate_collection_engines_cache(
31 username: Optional[str] = None,
32) -> None:
33 """Invalidate cached collection search engine configurations.
35 Args:
36 username: If provided, invalidate cache for this username only.
37 If None, clear the entire collection engines cache.
38 """
39 with _COLLECTION_CACHE_LOCK:
40 if username:
41 _COLLECTION_ENGINES_CACHE.pop(username, None)
42 else:
43 _COLLECTION_ENGINES_CACHE.clear()
46def _get_setting(
47 key: str,
48 default_value: Any = None,
49 db_session: Optional[Session] = None,
50 settings_snapshot: Optional[Dict[str, Any]] = None,
51 username: Optional[str] = None,
52) -> Any:
53 """
54 Get a setting from either a database session or settings snapshot.
56 Args:
57 key: The setting key
58 default_value: Default value if setting not found
59 db_session: Database session for direct access
60 settings_snapshot: Settings snapshot for thread context
61 username: Username for backward compatibility
63 Returns:
64 The setting value or default_value if not found
65 """
66 # Try settings snapshot first (thread context)
67 if settings_snapshot:
68 try:
69 return get_setting_from_snapshot(
70 key, default_value, settings_snapshot=settings_snapshot
71 )
72 except Exception as e:
73 logger.debug(
74 f"Could not get setting {key} from snapshot: {scrub_error(e)}"
75 )
77 # Try database session if available
78 if db_session:
79 try:
80 settings_manager = get_settings_manager(db_session, username)
81 return settings_manager.get_setting(key, default_value)
82 except Exception as e:
83 logger.debug(
84 f"Could not get setting {key} from db_session: {scrub_error(e)}"
85 )
87 # Return default if all methods fail
88 logger.warning(
89 f"Could not retrieve setting '{key}', returning default: {default_value}"
90 )
91 return default_value
94def _extract_per_engine_config(
95 raw_config: Dict[str, Any],
96) -> Dict[str, Dict[str, Any]]:
97 """
98 Converts the "flat" configuration loaded from the settings database into
99 individual settings dictionaries for each engine.
101 Args:
102 raw_config: The raw "flat" configuration.
104 Returns:
105 Configuration dictionaries indexed by engine name.
107 """
108 nested_config: dict[str, Any] = {}
109 for key, value in raw_config.items():
110 if "." in key:
111 # This is a higher-level key.
112 top_level_key = key.split(".")[0]
113 lower_keys = ".".join(key.split(".")[1:])
114 nested_config.setdefault(top_level_key, {})[lower_keys] = value
115 else:
116 # This is a low-level key.
117 nested_config[key] = value
119 # Expand all the lower-level keys.
120 for key, value in nested_config.items():
121 if isinstance(value, dict):
122 # Expand the child keys.
123 nested_config[key] = _extract_per_engine_config(value)
125 return nested_config
128def search_config(
129 username: Optional[str] = None,
130 db_session: Optional[Session] = None,
131 settings_snapshot: Optional[Dict[str, Any]] = None,
132) -> Dict[str, Any]:
133 """
134 Returns the search engine configuration loaded from the database or settings snapshot.
136 Args:
137 username: Username used to scope the per-user retriever listing
138 (own registrations plus shared ones) so one user's retrievers
139 never leak into another user's engine list. Falls back to the
140 ``_username`` carried in ``settings_snapshot`` when omitted.
141 db_session: Database session for direct access (preferred for web routes)
142 settings_snapshot: Settings snapshot for thread context (preferred for background threads)
144 Returns:
145 The search engine configuration loaded from the database or snapshot.
146 """
147 # Extract search engine definitions
148 config_data = _get_setting(
149 "search.engine.web",
150 {},
151 db_session=db_session,
152 settings_snapshot=settings_snapshot,
153 username=username,
154 )
156 search_engines = _extract_per_engine_config(config_data)
158 # Inject module/class from the hardcoded engine registry.
159 # This is the single source of truth for which Python module implements
160 # each engine — these values are never read from the settings DB.
161 from .engine_registry import ENGINE_REGISTRY
163 for name, entry in ENGINE_REGISTRY.items():
164 if name in search_engines:
165 search_engines[name]["module_path"] = entry.module_path
166 search_engines[name]["class_name"] = entry.class_name
167 if entry.full_search_module:
168 search_engines[name]["full_search_module"] = (
169 entry.full_search_module
170 )
171 search_engines[name]["full_search_class"] = (
172 entry.full_search_class
173 )
175 # Add registered retrievers as available search engines. Scope the
176 # listing to the requesting user (their own registrations plus shared
177 # ones) so one user's retriever names/metadata never leak into another
178 # user's engine list. Prefer an explicit username, else the ``_username``
179 # carried in the snapshot; None lists shared retrievers only.
180 from ..search_system import username_from_snapshot
182 effective_username = username or username_from_snapshot(settings_snapshot)
183 from .retriever_registry import retriever_registry
185 for name in retriever_registry.list_registered(username=effective_username):
186 search_engines[name] = {
187 "module_path": ".engines.search_engine_retriever",
188 "class_name": "RetrieverSearchEngine",
189 "requires_api_key": False,
190 "requires_llm": False,
191 "description": f"LangChain retriever: {name}",
192 "strengths": [
193 "Domain-specific knowledge",
194 "No rate limits",
195 "Fast retrieval",
196 ],
197 "weaknesses": ["Limited to indexed content"],
198 "supports_full_search": True,
199 "is_retriever": True, # Mark as retriever for identification
200 }
202 logger.info(
203 f"Loaded {len(search_engines)} search engines from configuration file"
204 )
205 logger.info(f"\n {', '.join(sorted(search_engines.keys()))} \n")
207 # Register Library RAG as a search engine
208 library_enabled = _get_setting(
209 "search.engine.library.enabled",
210 True,
211 db_session=db_session,
212 settings_snapshot=settings_snapshot,
213 username=username,
214 )
216 if library_enabled:
217 search_engines["library"] = {
218 "module_path": ".engines.search_engine_library",
219 "class_name": "LibraryRAGSearchEngine",
220 "requires_llm": True,
221 "display_name": "Search All Collections",
222 "default_params": {},
223 "description": "Search across all your document collections using semantic search",
224 "strengths": [
225 "Searches all your curated collections of research papers and documents",
226 "Uses semantic search for better relevance",
227 "Returns documents you've already saved and reviewed",
228 ],
229 "weaknesses": [
230 "Limited to documents already in your collections",
231 "Requires documents to be indexed first",
232 ],
233 "reliability": "High - searches all your collections",
234 }
235 logger.info("Registered Library RAG as search engine")
237 # Register document collections as individual search engines
238 if library_enabled:
239 try:
240 from ..database.models.library import Collection
241 from ..database.session_context import get_user_db_session
242 from ..search_system import username_from_snapshot
244 # Get username from settings_snapshot if available
245 collection_username = (
246 username_from_snapshot(settings_snapshot)
247 if settings_snapshot
248 else username
249 )
251 if collection_username:
252 now = time.monotonic()
253 cached_engines = None
254 with _COLLECTION_CACHE_LOCK:
255 cache_entry = _COLLECTION_ENGINES_CACHE.get(
256 collection_username
257 )
258 if cache_entry and (
259 now - cache_entry[0] < _COLLECTION_CACHE_TTL
260 ):
261 cached_engines = cache_entry[1]
263 if cached_engines is not None:
264 search_engines.update(copy.deepcopy(cached_engines))
265 logger.debug(
266 f"Reused {len(cached_engines)} cached document collection engines for {collection_username}"
267 )
268 else:
269 try:
270 with get_user_db_session(
271 collection_username
272 ) as session:
273 collections = session.query(Collection).all()
275 user_collection_engines = {}
276 for collection in collections:
277 engine_id = f"collection_{collection.id}"
278 # Add suffix to distinguish from the all-collections search
279 display_name = f"{collection.name} (Collection)"
280 # Egress classification follows the per-collection
281 # public/private flag (default private). A "public"
282 # collection counts as a public engine (allowed under
283 # PUBLIC_ONLY); a private one is local-only. NULL
284 # (pre-migration rows) reads as private — the safe
285 # default.
286 collection_is_public = bool(
287 getattr(collection, "is_public", False)
288 )
289 # Usability flag (NOT egress): whether the LangGraph
290 # research agent offers this collection as a tool. NULL
291 # (pre-migration rows) reads as available (True). Uses
292 # the same `is not False` idiom as the rag_routes
293 # serializers so all call sites share one NULL→available
294 # default and can't drift.
295 collection_agent_enabled = (
296 getattr(collection, "agent_enabled", True)
297 is not False
298 )
299 user_collection_engines[engine_id] = {
300 "module_path": ".engines.search_engine_collection",
301 "class_name": "CollectionSearchEngine",
302 "requires_llm": True,
303 "is_local": True,
304 "is_public": collection_is_public,
305 "agent_enabled": collection_agent_enabled,
306 "display_name": display_name,
307 "default_params": {
308 "collection_id": collection.id,
309 "collection_name": collection.name,
310 },
311 "description": (
312 collection.description
313 if collection.description
314 else f"Search documents in {collection.name} collection only"
315 ),
316 "strengths": [
317 f"Searches only documents in {collection.name}",
318 "Focused semantic search within specific topic area",
319 "Returns documents from a curated collection",
320 ],
321 "weaknesses": [
322 "Limited to documents in this collection",
323 "Smaller result pool than full library search",
324 ],
325 "reliability": "High - searches a specific collection",
326 }
328 with _COLLECTION_CACHE_LOCK:
329 _COLLECTION_ENGINES_CACHE[
330 collection_username
331 ] = (
332 now,
333 copy.deepcopy(user_collection_engines),
334 )
336 search_engines.update(
337 copy.deepcopy(user_collection_engines)
338 )
339 logger.info(
340 f"Registered {len(user_collection_engines)} document collections as search engines"
341 )
342 except Exception as exc:
343 safe_msg = scrub_error(exc)
344 logger.warning(
345 f"Could not register document collections for {collection_username}: {safe_msg}"
346 )
347 logger.debug(
348 f"Traceback for collection registration failure for {collection_username}",
349 exc_info=True,
350 )
351 with _COLLECTION_CACHE_LOCK:
352 stale_entry = _COLLECTION_ENGINES_CACHE.get(
353 collection_username
354 )
355 if stale_entry and (
356 now - stale_entry[0] <= _COLLECTION_CACHE_MAX_STALE
357 ):
358 search_engines.update(copy.deepcopy(stale_entry[1]))
359 logger.warning(
360 f"Using {len(stale_entry[1])} stale cached collection engines for {collection_username} after DB error"
361 )
362 elif stale_entry:
363 logger.warning(
364 f"Stale cached collection engines for {collection_username} expired "
365 f"(age {now - stale_entry[0]:.1f}s > {_COLLECTION_CACHE_MAX_STALE}s), skipping"
366 )
367 else:
368 logger.debug(
369 "No username available for collection registration"
370 )
371 except Exception as exc:
372 safe_msg = scrub_error(exc)
373 logger.warning(
374 f"Could not register document collections: {safe_msg}"
375 )
376 logger.debug(
377 "Traceback for document collections registration failure",
378 exc_info=True,
379 )
381 return search_engines
384def list_eligible_engine_configs(
385 settings_snapshot: Optional[Dict[str, Any]] = None,
386 *,
387 use_api_key_services: bool = True,
388 egress_context: Optional[Any] = None,
389 check_agent_enabled: bool = False,
390 check_auto_search: bool = False,
391) -> Dict[str, Any]:
392 """Return every engine from ``search_config`` that can be instantiated.
394 This is the candidate pool for surfaces that decide visibility
395 through their own flag (``agent_enabled`` for the langgraph research
396 agent) — NOT for surfaces governed by ``use_in_auto_search``. Engines
397 that need a key are dropped when no key is set (or when
398 ``use_api_key_services`` is False), so callers don't have to repeat the
399 credential check; everything else (retrievers, the built-in ``library``
400 engine, dynamic ``collection_*`` entries) flows through.
402 The name keys in the returned dict are the canonical engine names from
403 ``search_config`` (e.g. ``searxng``, ``elasticsearch``) — matching the
404 registry in ``engine_registry.py`` 1:1, so callers can compare against
405 ``search.tool`` directly.
407 Engines whose ``is_available(settings_snapshot)`` returns ``False`` are
408 excluded — this catches the case where the engine is configured but its
409 backing service is unreachable (e.g. Elasticsearch with no running
410 cluster on ``localhost:9200``). Without this filter the agent still
411 SEES the broken engine as a tool in its heartbeat and the factory logs
412 ``Failed to create search engine '<name>' (ConnectionError)`` per step.
413 Subclasses of ``BaseSearchEngine`` override ``is_available`` to opt in;
414 the default is True so existing engines are unaffected.
416 Filters for disabled state (``agent_enabled`` / ``use_in_auto_search``) and
417 egress policy are evaluated BEFORE calling ``is_available`` so unreachable or
418 disabled engines do not trigger unnecessary network probes.
420 Args:
421 settings_snapshot: Thread-safe settings snapshot.
422 use_api_key_services: When False, engines that require an API key
423 are excluded even when the key is present.
424 egress_context: Optional EgressContext to evaluate policy before probing.
425 check_agent_enabled: When True, excludes engines disabled for agent.
426 check_auto_search: When True, excludes engines disabled for auto-search.
428 Returns:
429 Dict of engine_name → config for engines that passed all checks.
430 """
431 if not settings_snapshot:
432 logger.warning(
433 "list_eligible_engine_configs called without settings_snapshot, "
434 "returning empty dict"
435 )
436 return {}
438 all_engines = search_config(settings_snapshot=settings_snapshot)
440 if egress_context is not None:
441 from ..security.egress.policy import (
442 evaluate_engine,
443 evaluate_retriever,
444 )
446 eligible: Dict[str, Any] = {}
447 for name, config in all_engines.items():
448 requires_key = config.get("requires_api_key", False)
450 if requires_key and not use_api_key_services:
451 continue
452 if requires_key:
453 api_key = _resolve_api_key(name, config, settings_snapshot)
454 if not api_key:
455 logger.debug(
456 f"Skipping {name} — requires API key but none configured"
457 )
458 continue
460 # Optional auto-search flag check prior to probe
461 if check_auto_search:
462 use_in_auto = config.get("use_in_auto_search")
463 if use_in_auto is None: 463 ↛ 468line 463 didn't jump to line 468 because the condition on line 463 was always true
464 auto_search_key = f"search.engine.web.{name}.use_in_auto_search"
465 use_in_auto = get_setting_from_snapshot(
466 auto_search_key, False, settings_snapshot=settings_snapshot
467 )
468 if not use_in_auto:
469 continue
471 # Optional agent_enabled flag check prior to probe
472 if check_agent_enabled:
473 agent_enabled = config.get("agent_enabled")
474 if agent_enabled is None:
475 agent_enabled = get_setting_from_snapshot(
476 f"search.engine.web.{name}.agent_enabled",
477 True,
478 settings_snapshot=settings_snapshot,
479 )
480 if not agent_enabled:
481 logger.debug(f"Skipping {name} — disabled for research agent")
482 continue
484 # Egress policy check prior to network probe
485 if egress_context is not None:
486 if config.get("is_retriever"): 486 ↛ 487line 486 didn't jump to line 487 because the condition on line 486 was never true
487 try:
488 from .retriever_registry import retriever_registry
490 meta = retriever_registry.get_metadata(
491 name,
492 username=getattr(egress_context, "username", None),
493 )
494 except Exception:
495 meta = None
496 decision = evaluate_retriever(
497 name, egress_context, metadata=meta
498 )
499 else:
500 decision = evaluate_engine(
501 name,
502 egress_context,
503 settings_snapshot=settings_snapshot,
504 metadata=config,
505 )
506 if not decision.allowed:
507 logger.debug(
508 f"Skipping {name} — disallowed by egress policy ({decision.reason})"
509 )
510 continue
512 # Runtime availability probe. Cheap (cached, single TCP connect at
513 # worst) and excludes engines whose backing service isn't reachable
514 # so the agent doesn't waste steps advertising broken tools.
515 # ``is_retriever`` entries go through the retriever registry's own
516 # ``is_available`` path — skip them here so we don't double-probe
517 # or fall back to the base class default for custom retriever
518 # classes that never subclass ``BaseSearchEngine``.
519 if not config.get("is_retriever", False): 519 ↛ 526line 519 didn't jump to line 526 because the condition on line 519 was always true
520 if not _engine_class_is_available(name, config, settings_snapshot):
521 logger.debug(
522 f"Skipping {name} — runtime availability probe failed"
523 )
524 continue
526 eligible[name] = config
528 return eligible
531def _engine_class_is_available(
532 name: str,
533 config: Dict[str, Any],
534 settings_snapshot: Dict[str, Any],
535) -> bool:
536 """Load the engine class for ``config`` and call its ``is_available``.
538 Returns True (fail-open) on any error — the class load can legitimately
539 fail (bad module_path, missing dependency, sandbox import restrictions)
540 and we don't want ``list_eligible_engine_configs`` to mask real config
541 problems with an unrelated "engine unavailable" log line. The factory
542 will surface the real ``Failed to create search engine`` error when
543 (and only when) the caller actually tries to instantiate it.
544 """
545 from .search_engine_base import BaseSearchEngine
547 module_path = config.get("module_path")
548 class_name = config.get("class_name")
549 if not module_path or not class_name:
550 return True
552 try:
553 from ..security.module_whitelist import get_safe_module_class
555 engine_class = get_safe_module_class(module_path, class_name)
556 except Exception as exc:
557 logger.debug(
558 "Skipping availability probe for {} — could not load engine "
559 "class ({}: {})",
560 name,
561 type(exc).__name__,
562 exc,
563 )
564 return True
566 # Only subclasses of BaseSearchEngine get the ``is_available`` contract.
567 # Library / Collection / Retriever engines either subclass it (in which
568 # case the default ``True`` applies) or don't (also default True).
569 if not isinstance(engine_class, type) or not issubclass( 569 ↛ 572line 569 didn't jump to line 572 because the condition on line 569 was never true
570 engine_class, BaseSearchEngine
571 ):
572 return True
574 try:
575 return bool(
576 engine_class.is_available(settings_snapshot=settings_snapshot)
577 )
578 except Exception as exc:
579 logger.debug(
580 "is_available raised for {} — treating as available ({})",
581 name,
582 type(exc).__name__,
583 )
584 return True
587def get_available_engines(
588 settings_snapshot: Optional[Dict[str, Any]] = None,
589 use_api_key_services: bool = True,
590 exclude_engines: Optional[set] = None,
591) -> Dict[str, Any]:
592 """
593 Return search engines that are actually usable: enabled for auto-search
594 and with valid API keys when required.
596 This is the single shared filter used by ``auto``-mode search so the
597 LangGraph agent agrees with the rest of the system on which engines
598 show up in non-agent searches. For the agent's specialised tool list,
599 use :func:`list_eligible_engine_configs` instead — the agent's surface
600 is governed by per-engine ``agent_enabled``, not by this auto-search
601 filter (see #5015).
603 Args:
604 settings_snapshot: Thread-safe settings snapshot.
605 use_api_key_services: If False, engines that require an API key are
606 excluded even when the key is present.
607 exclude_engines: Additional engine names to skip (e.g. the caller's
608 own name).
610 Returns:
611 Dict of engine_name → config for engines that passed all checks.
612 """
613 eligible = list_eligible_engine_configs(
614 settings_snapshot=settings_snapshot,
615 use_api_key_services=use_api_key_services,
616 check_auto_search=True,
617 )
618 excluded = set(exclude_engines) if exclude_engines else set()
620 return {
621 name: config
622 for name, config in eligible.items()
623 if name not in excluded
624 }
627def _resolve_api_key(
628 engine_name: str,
629 engine_config: Dict[str, Any],
630 settings_snapshot: Dict[str, Any],
631) -> Optional[str]:
632 """
633 Try to find a valid API key for *engine_name*.
635 Resolution order (mirrors ``create_search_engine``):
636 1. ``search.engine.web.<name>.api_key`` in the snapshot
637 2. ``api_key`` inside the engine config dict
639 Returns the key string or None.
640 """
641 api_key = None
642 api_key_path = f"search.engine.web.{engine_name}.api_key"
644 api_key_setting = settings_snapshot.get(api_key_path)
645 if api_key_setting:
646 api_key = (
647 api_key_setting.get("value")
648 if isinstance(api_key_setting, dict)
649 else api_key_setting
650 )
652 if not api_key:
653 api_key = engine_config.get("api_key")
655 if not api_key:
656 return None
658 # Reject common placeholder values
659 api_key_str = str(api_key).strip()
660 if _is_api_key_placeholder(api_key_str): 660 ↛ 661line 660 didn't jump to line 661 because the condition on line 660 was never true
661 return None
663 return api_key_str
666def default_search_engine(
667 username: Optional[str] = None,
668 db_session: Optional[Session] = None,
669 settings_snapshot: Optional[Dict[str, Any]] = None,
670) -> str:
671 """
672 Returns the configured default search engine.
674 Args:
675 username: Username for backward compatibility (deprecated)
676 db_session: Database session for direct access (preferred for web routes)
677 settings_snapshot: Settings snapshot for thread context (preferred for background threads)
679 Returns:
680 The configured default search engine.
681 """
682 return str(
683 _get_setting(
684 "search.engine.DEFAULT_SEARCH_ENGINE",
685 "wikipedia",
686 db_session=db_session,
687 settings_snapshot=settings_snapshot,
688 username=username,
689 )
690 )