Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-mcp/src/lexigram/ai/mcp/di/provider.py: 20%
232 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 07:19 +0800
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 07:19 +0800
1"""MCP provider for Lexigram dependency injection."""
3from __future__ import annotations
5from typing import TYPE_CHECKING, Any, cast
7from lexigram.ai.mcp.config import MCPConfig
8from lexigram.ai.mcp.exceptions import MCPInitializationError
9from lexigram.ai.mcp.server import MCPServer
10from lexigram.ai.mcp.server.handlers import (
11 LoggingHandler,
12 PromptHandler,
13 ResourceHandler,
14 SamplingHandler,
15 ToolHandler,
16)
17from lexigram.ai.mcp.transport import SSETransport, StdioTransport
18from lexigram.contracts import (
19 ProviderPriority,
20)
21from lexigram.contracts.core.health import HealthCheckResult, HealthStatus
22from lexigram.contracts.exceptions import UnresolvableDependencyError
23from lexigram.contracts.mcp.protocols import MCPAuthorizerProtocol
24from lexigram.di.provider import Provider
25from lexigram.logging import (
26 get_logger,
27)
29if TYPE_CHECKING:
30 from lexigram.contracts.core.di import (
31 BootContainerProtocol,
32 ContainerRegistrarProtocol,
33 )
35logger = get_logger(__name__)
38class _ConnectorToolBundle:
39 """Sentinel DI type: combined tool provider for all built-in connectors."""
42class _ConnectorResourceBundle:
43 """Sentinel DI type: combined resource provider for all built-in connectors."""
46class MCPProvider(Provider):
47 name = "mcp"
48 priority = ProviderPriority.PRESENTATION
49 config_key: str | None = "ai_mcp"
50 config_model: type | None = MCPConfig
52 """Provider for MCP server components.
54 Registers MCP server, handlers, and transports with the DI container.
55 When ``MCPConfig.client_stdio_command`` or ``MCPConfig.client_url`` is
56 set, an :class:`~lexigram.ai.mcp.client.MCPClient` is also registered.
58 When ``controllers`` is provided (typically set by ``MCPModule``), the
59 controller instances are registered as tool/resource/prompt providers
60 before the handlers are wired up.
61 """
63 def __init__(
64 self,
65 config: MCPConfig | None = None,
66 controllers: list[type] | None = None,
67 services: list[type] | None = None,
68 include_methods: list[str] | None = None,
69 enable_streaming: bool = True,
70 **provider_kwargs: Any,
71 ) -> None:
72 super().__init__()
73 self._requested_config = config
74 self._config = config or MCPConfig()
75 self._controllers: list[type] = controllers or []
76 self._services: list[type] = services or []
77 self._include_methods: list[str] | None = include_methods
78 self._enable_streaming = enable_streaming
79 self._managed_transports: list[StdioTransport | SSETransport] = []
80 self._managed_client: Any | None = None
82 async def register(self, container: ContainerRegistrarProtocol) -> None:
83 """Register MCPConfig with the DI container."""
84 self._config = self._requested_config or (
85 self.config
86 if isinstance(getattr(self, "config", None), MCPConfig)
87 else self._config
88 )
89 container.singleton(MCPConfig, self._config)
91 async def boot(self, container: BootContainerProtocol) -> None:
92 """Wire up all MCP components using the resolved container."""
93 if not self._config.enabled:
94 return
96 if self._controllers:
97 await self._boot_controller_providers(container)
98 if self._services:
99 await self._boot_service_providers(container)
100 await self._boot_skill_bridge(container)
101 await self._boot_connectors(container)
102 await self._boot_handlers(container)
103 await self._boot_server(container)
104 self._boot_transports(container)
105 await self._boot_client(container)
107 async def _boot_controller_providers(
108 self, container: BootContainerProtocol
109 ) -> None:
110 """Instantiate MCPController classes and register compound providers.
112 Controllers are resolved via the DI container when possible so that
113 constructor dependencies are injected automatically.
114 """
115 from lexigram.ai.mcp.controllers import (
116 ControllerPromptProvider,
117 ControllerResourceProvider,
118 ControllerToolProvider,
119 MCPController,
120 )
122 instances: list[MCPController] = []
123 for ctrl_class in self._controllers:
124 try:
125 instance: MCPController = await container.resolve(ctrl_class)
126 except (UnresolvableDependencyError, LookupError) as exc:
127 raise MCPInitializationError(
128 message=f"Failed to resolve MCP controller: {ctrl_class.__name__}"
129 ) from exc
130 instances.append(instance)
131 logger.debug("mcp_controller_registered", controller=ctrl_class.__name__)
133 container.singleton(ControllerToolProvider, ControllerToolProvider(instances))
134 container.singleton(
135 ControllerResourceProvider, ControllerResourceProvider(instances)
136 )
137 container.singleton(
138 ControllerPromptProvider, ControllerPromptProvider(instances)
139 )
141 async def _boot_service_providers(self, container: BootContainerProtocol) -> None:
142 """Resolve service classes and register a ServiceToolProvider."""
143 from lexigram.ai.mcp.controllers import ServiceToolProvider
145 instances: list[object] = []
146 for svc_class in self._services:
147 try:
148 instance: object = await container.resolve(svc_class)
149 except (UnresolvableDependencyError, LookupError) as exc:
150 raise MCPInitializationError(
151 message=f"Failed to resolve MCP service: {svc_class.__name__}"
152 ) from exc
153 instances.append(instance)
154 logger.debug("mcp_service_registered", service=svc_class.__name__)
156 container.singleton(
157 ServiceToolProvider,
158 ServiceToolProvider.from_services(instances, self._include_methods),
159 )
161 async def _boot_skill_bridge(self, container: BootContainerProtocol) -> None:
162 """Register SkillToolAdapter if lexigram-ai-skills is available.
164 When ``SkillRegistryProtocol`` and ``SkillExecutorProtocol`` are
165 registered in the container (by ``lexigram-ai-skills``), a
166 :class:`~lexigram.ai.mcp.adapters.skill_adapter.SkillToolAdapter` is
167 created and registered so that skills are exposed as MCP tools.
168 """
169 skill_registry = None
170 skill_executor = None
171 try:
172 from lexigram.contracts.ai.skills import (
173 SkillExecutorProtocol,
174 SkillRegistryProtocol,
175 )
177 skill_registry = await container.resolve(SkillRegistryProtocol)
178 skill_executor = await container.resolve(SkillExecutorProtocol)
179 except (
180 LookupError,
181 RuntimeError,
182 AttributeError,
183 ImportError,
184 UnresolvableDependencyError,
185 ):
186 logger.debug("mcp_skills_not_available")
187 return
189 from lexigram.ai.mcp.adapters.skill_adapter import SkillToolAdapter
191 adapter = SkillToolAdapter(
192 skill_registry=skill_registry,
193 skill_executor=skill_executor,
194 )
195 container.singleton(SkillToolAdapter, adapter)
196 logger.info("mcp_skill_bridge_registered")
198 async def _boot_connectors(self, container: BootContainerProtocol) -> None:
199 """Instantiate and register built-in connectors from ``MCPConfig.connectors``.
201 Each connector that has sufficient configuration is instantiated and
202 collected into :class:`~lexigram.ai.mcp.controllers._ConnectorToolBundle`
203 and :class:`~lexigram.ai.mcp.controllers._ConnectorResourceBundle`
204 sentinel types. :meth:`_boot_handlers` then picks them up and merges
205 them with the controller/service providers.
206 """
207 from lexigram.ai.mcp.connectors import (
208 FilesystemConnector,
209 GitHubConnector,
210 GoogleDriveConnector,
211 SlackConnector,
212 SQLConnector,
213 WebFetchConnector,
214 WebSearchConnector,
215 )
216 from lexigram.ai.mcp.controllers import (
217 _CombinedResourceProvider,
218 _CombinedToolProvider,
219 )
221 cfg = self._config.connectors
222 connectors: list[object] = []
224 if cfg.filesystem.root_dir:
225 connectors.append(
226 FilesystemConnector(
227 root_dir=cfg.filesystem.root_dir,
228 read_only=cfg.filesystem.read_only,
229 )
230 )
231 logger.info("mcp_connector_registered", connector="filesystem")
233 if cfg.github.token:
234 connectors.append(GitHubConnector(token=cfg.github.token))
235 logger.info("mcp_connector_registered", connector="github")
237 if cfg.web_fetch.enabled:
238 connectors.append(
239 WebFetchConnector(
240 max_content_bytes=cfg.web_fetch.max_content_bytes,
241 )
242 )
243 logger.info("mcp_connector_registered", connector="web_fetch")
245 if cfg.web_search.api_key:
246 connectors.append(
247 WebSearchConnector(
248 provider=cfg.web_search.provider,
249 api_key=cfg.web_search.api_key,
250 max_results=cfg.web_search.max_results,
251 )
252 )
253 logger.info("mcp_connector_registered", connector="web_search")
255 if cfg.slack.bot_token:
256 connectors.append(
257 SlackConnector(
258 bot_token=cfg.slack.bot_token,
259 max_messages=cfg.slack.max_messages,
260 )
261 )
262 logger.info("mcp_connector_registered", connector="slack")
264 if cfg.google_drive.service_account_json:
265 connectors.append(
266 GoogleDriveConnector(
267 service_account_json=cfg.google_drive.service_account_json,
268 impersonated_email=cfg.google_drive.impersonated_email,
269 )
270 )
271 logger.info("mcp_connector_registered", connector="google_drive")
273 if cfg.sql.dsn:
274 try:
275 from lexigram.contracts.data import DatabaseProviderProtocol
277 db_provider = await container.resolve_optional(DatabaseProviderProtocol)
278 if db_provider is not None:
279 connectors.append(
280 SQLConnector(
281 db=db_provider,
282 allowed_tables=cfg.sql.allowed_tables,
283 read_only=cfg.sql.read_only,
284 )
285 )
286 logger.info("mcp_connector_registered", connector="sql")
287 except (ImportError, AttributeError, LookupError):
288 logger.debug("mcp_sql_connector_skipped_no_db_provider")
290 if not connectors:
291 return
293 tool_bundle = _CombinedToolProvider(connectors)
294 resource_bundle = _CombinedResourceProvider(connectors)
295 container.singleton(_ConnectorToolBundle, tool_bundle)
296 container.singleton(_ConnectorResourceBundle, resource_bundle)
298 async def _boot_handlers(self, container: BootContainerProtocol) -> None:
299 """Register MCP handlers, wiring in controller providers when available.
301 When ``MCPModule(controllers=[...])`` has already registered
302 ``ControllerToolProvider`` / ``ControllerResourceProvider`` /
303 ``ControllerPromptProvider`` in the container, those are passed to the
304 handlers so that ``@tool``/``@resource``/``@prompt`` decorated methods
305 are automatically exposed. Falls back to empty (no-op) handlers when
306 no controllers are configured.
307 """
308 from lexigram.ai.mcp.controllers import (
309 ControllerPromptProvider,
310 ControllerResourceProvider,
311 ControllerToolProvider,
312 )
314 tool_provider: Any | None = await container.resolve_optional(
315 ControllerToolProvider
316 )
317 resource_provider: Any | None = await container.resolve_optional(
318 ControllerResourceProvider
319 )
320 prompt_provider: Any | None = await container.resolve_optional(
321 ControllerPromptProvider
322 )
324 # Merge service tools with controller tools when both exist
325 service_tool_provider = None
326 try:
327 from lexigram.ai.mcp.controllers import ServiceToolProvider
329 service_tool_provider = await container.resolve_optional(
330 ServiceToolProvider
331 )
332 except (ImportError, AttributeError):
333 pass
335 if tool_provider is not None and service_tool_provider is not None:
336 from lexigram.ai.mcp.controllers import _CombinedToolProvider
338 tool_provider = _CombinedToolProvider(
339 [tool_provider, service_tool_provider]
340 )
341 elif service_tool_provider is not None:
342 tool_provider = service_tool_provider
344 # Merge built-in connector tools / resources when any connectors are configured
345 connector_tool_bundle = await container.resolve_optional(_ConnectorToolBundle)
346 connector_resource_bundle = await container.resolve_optional(
347 _ConnectorResourceBundle
348 )
349 if connector_tool_bundle is not None:
350 from lexigram.ai.mcp.controllers import _CombinedToolProvider
352 tool_provider = _CombinedToolProvider(
353 [p for p in [tool_provider, connector_tool_bundle] if p is not None]
354 )
355 if connector_resource_bundle is not None:
356 from lexigram.ai.mcp.controllers import _CombinedResourceProvider
358 resource_provider = _CombinedResourceProvider(
359 [
360 p
361 for p in [resource_provider, connector_resource_bundle]
362 if p is not None
363 ]
364 )
366 tool_handler = ToolHandler(tool_provider)
367 resource_handler = ResourceHandler(resource_provider)
368 prompt_handler = PromptHandler(prompt_provider)
370 container.singleton(ToolHandler, tool_handler)
371 container.singleton(ResourceHandler, resource_handler)
372 container.singleton(PromptHandler, prompt_handler)
374 # Optional: sampling handler (requires an LLM client in the container)
375 try:
376 from lexigram.contracts.ai.llm import LLMClientProtocol
378 llm_client = await container.resolve_optional(LLMClientProtocol)
379 if llm_client is not None:
380 sampling_handler = SamplingHandler(llm=llm_client)
381 container.singleton(SamplingHandler, sampling_handler)
382 logger.info("mcp_sampling_handler_registered")
383 except (ImportError, AttributeError):
384 pass
386 # Always register a logging handler
387 logging_handler = LoggingHandler()
388 container.singleton(LoggingHandler, logging_handler)
390 # G41: Optional observability wrapping
391 try:
392 from lexigram.ai.mcp.server.observability import (
393 ObservableResourceHandler,
394 ObservableToolHandler,
395 )
396 from lexigram.contracts.observability.ai import AIMetricsProtocol
398 metrics = await container.resolve_optional(AIMetricsProtocol)
399 if metrics is not None:
400 observable_tool = ObservableToolHandler(tool_handler, metrics)
401 observable_resource = ObservableResourceHandler(
402 resource_handler, metrics
403 )
404 container.singleton(ToolHandler, observable_tool)
405 container.singleton(ResourceHandler, observable_resource)
406 logger.info("mcp_observability_wrapping_registered")
407 except (ImportError, AttributeError):
408 pass
410 self._handlers_registered = True
412 async def _boot_server(self, container: BootContainerProtocol) -> None:
413 """Register MCP server."""
414 tool_handler = await container.resolve(ToolHandler)
415 resource_handler = await container.resolve(ResourceHandler)
416 prompt_handler = await container.resolve(PromptHandler)
417 sampling_handler = await container.resolve_optional(SamplingHandler)
418 logging_handler = await container.resolve_optional(LoggingHandler)
419 authorizer = await container.resolve_optional(MCPAuthorizerProtocol)
421 server = MCPServer(
422 name=self._config.server_name,
423 version=self._config.server_version,
424 tool_handler=tool_handler,
425 resource_handler=resource_handler,
426 prompt_handler=prompt_handler,
427 sampling_handler=sampling_handler,
428 logging_handler=logging_handler,
429 authorizer=authorizer,
430 allow_unauthenticated=self._config.allow_unauthenticated,
431 )
433 registrar = cast("ContainerRegistrarProtocol", container)
434 registrar.singleton(MCPServer, server)
436 def _boot_transports(self, container: ContainerRegistrarProtocol) -> None:
437 """Register MCP transport singletons."""
438 stdio_transport = StdioTransport()
439 sse_transport = SSETransport()
441 container.singleton(StdioTransport, stdio_transport)
442 container.singleton(SSETransport, sse_transport)
443 self._managed_transports = [stdio_transport, sse_transport]
445 async def shutdown(self) -> None:
446 """Clean up MCP provider resources."""
447 for transport in self._managed_transports:
448 try:
449 await transport.stop()
450 except (RuntimeError, AttributeError, TypeError, OSError) as exc:
451 logger.warning(
452 "mcp_transport_stop_failed",
453 transport=type(transport).__name__,
454 error=str(exc),
455 )
457 self._managed_transports.clear()
459 if self._managed_client is not None:
460 try:
461 await self._managed_client.disconnect()
462 except (
463 RuntimeError,
464 AttributeError,
465 TypeError,
466 OSError,
467 ConnectionError,
468 TimeoutError,
469 ) as exc:
470 logger.warning("mcp_client_disconnect_failed", error=str(exc))
471 self._managed_client = None
473 logger.info("mcp_provider_shutdown_complete")
475 async def health_check(self, timeout: float = 5.0) -> HealthCheckResult:
476 """Return a health check result for this provider.
478 Args:
479 timeout: Maximum seconds to wait for health response.
481 Returns:
482 A :class:`~lexigram.contracts.core.health.HealthCheckResult`
483 reflecting the current MCP provider state.
484 """
485 if self._managed_client is not None:
486 is_connected = getattr(self._managed_client, "is_connected", None)
487 if (callable(is_connected) and not is_connected()) or (
488 isinstance(is_connected, bool) and not is_connected
489 ):
490 return HealthCheckResult(
491 component=self.name,
492 status=HealthStatus.DEGRADED,
493 details={"reason": "MCP client disconnected"},
494 )
495 return HealthCheckResult(
496 component=self.name,
497 status=HealthStatus.HEALTHY,
498 details={"status": "operational"},
499 )
501 async def _boot_client(self, container: BootContainerProtocol) -> None:
502 """Register MCPClient if client configuration is present."""
503 from lexigram.ai.mcp.client.core import (
504 MCPClient,
505 SSEClientTransport,
506 StdioClientTransport,
507 )
509 self._managed_client = None
511 if self._config.client_stdio_command:
512 stdio_client_transport = StdioClientTransport(
513 self._config.client_stdio_command,
514 startup_timeout=self._config.request_timeout,
515 )
516 client = MCPClient(
517 stdio_client_transport, request_timeout=self._config.request_timeout
518 )
519 registrar = cast("ContainerRegistrarProtocol", container)
520 registrar.singleton(MCPClient, client)
521 self._managed_client = client
522 logger.info(
523 "mcp_client_registered",
524 mode="stdio",
525 command=self._config.client_stdio_command[0],
526 )
527 elif self._config.client_url:
528 sse_client_transport = SSEClientTransport(
529 self._config.client_url,
530 request_timeout=self._config.request_timeout,
531 )
532 client = MCPClient(
533 sse_client_transport, request_timeout=self._config.request_timeout
534 )
535 registrar = cast("ContainerRegistrarProtocol", container)
536 registrar.singleton(MCPClient, client)
537 self._managed_client = client
538 logger.info(
539 "mcp_client_registered",
540 mode="sse",
541 url=self._config.client_url,
542 )
545__all__ = ["MCPProvider"]