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

1"""MCP provider for Lexigram dependency injection.""" 

2 

3from __future__ import annotations 

4 

5from typing import TYPE_CHECKING, Any, cast 

6 

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) 

28 

29if TYPE_CHECKING: 

30 from lexigram.contracts.core.di import ( 

31 BootContainerProtocol, 

32 ContainerRegistrarProtocol, 

33 ) 

34 

35logger = get_logger(__name__) 

36 

37 

38class _ConnectorToolBundle: 

39 """Sentinel DI type: combined tool provider for all built-in connectors.""" 

40 

41 

42class _ConnectorResourceBundle: 

43 """Sentinel DI type: combined resource provider for all built-in connectors.""" 

44 

45 

46class MCPProvider(Provider): 

47 name = "mcp" 

48 priority = ProviderPriority.PRESENTATION 

49 config_key: str | None = "ai_mcp" 

50 config_model: type | None = MCPConfig 

51 

52 """Provider for MCP server components. 

53 

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. 

57 

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

62 

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 

81 

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) 

90 

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 

95 

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) 

106 

107 async def _boot_controller_providers( 

108 self, container: BootContainerProtocol 

109 ) -> None: 

110 """Instantiate MCPController classes and register compound providers. 

111 

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 ) 

121 

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

132 

133 container.singleton(ControllerToolProvider, ControllerToolProvider(instances)) 

134 container.singleton( 

135 ControllerResourceProvider, ControllerResourceProvider(instances) 

136 ) 

137 container.singleton( 

138 ControllerPromptProvider, ControllerPromptProvider(instances) 

139 ) 

140 

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 

144 

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

155 

156 container.singleton( 

157 ServiceToolProvider, 

158 ServiceToolProvider.from_services(instances, self._include_methods), 

159 ) 

160 

161 async def _boot_skill_bridge(self, container: BootContainerProtocol) -> None: 

162 """Register SkillToolAdapter if lexigram-ai-skills is available. 

163 

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 ) 

176 

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 

188 

189 from lexigram.ai.mcp.adapters.skill_adapter import SkillToolAdapter 

190 

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

197 

198 async def _boot_connectors(self, container: BootContainerProtocol) -> None: 

199 """Instantiate and register built-in connectors from ``MCPConfig.connectors``. 

200 

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 ) 

220 

221 cfg = self._config.connectors 

222 connectors: list[object] = [] 

223 

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

232 

233 if cfg.github.token: 

234 connectors.append(GitHubConnector(token=cfg.github.token)) 

235 logger.info("mcp_connector_registered", connector="github") 

236 

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

244 

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

254 

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

263 

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

272 

273 if cfg.sql.dsn: 

274 try: 

275 from lexigram.contracts.data import DatabaseProviderProtocol 

276 

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

289 

290 if not connectors: 

291 return 

292 

293 tool_bundle = _CombinedToolProvider(connectors) 

294 resource_bundle = _CombinedResourceProvider(connectors) 

295 container.singleton(_ConnectorToolBundle, tool_bundle) 

296 container.singleton(_ConnectorResourceBundle, resource_bundle) 

297 

298 async def _boot_handlers(self, container: BootContainerProtocol) -> None: 

299 """Register MCP handlers, wiring in controller providers when available. 

300 

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 ) 

313 

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 ) 

323 

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 

328 

329 service_tool_provider = await container.resolve_optional( 

330 ServiceToolProvider 

331 ) 

332 except (ImportError, AttributeError): 

333 pass 

334 

335 if tool_provider is not None and service_tool_provider is not None: 

336 from lexigram.ai.mcp.controllers import _CombinedToolProvider 

337 

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 

343 

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 

351 

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 

357 

358 resource_provider = _CombinedResourceProvider( 

359 [ 

360 p 

361 for p in [resource_provider, connector_resource_bundle] 

362 if p is not None 

363 ] 

364 ) 

365 

366 tool_handler = ToolHandler(tool_provider) 

367 resource_handler = ResourceHandler(resource_provider) 

368 prompt_handler = PromptHandler(prompt_provider) 

369 

370 container.singleton(ToolHandler, tool_handler) 

371 container.singleton(ResourceHandler, resource_handler) 

372 container.singleton(PromptHandler, prompt_handler) 

373 

374 # Optional: sampling handler (requires an LLM client in the container) 

375 try: 

376 from lexigram.contracts.ai.llm import LLMClientProtocol 

377 

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 

385 

386 # Always register a logging handler 

387 logging_handler = LoggingHandler() 

388 container.singleton(LoggingHandler, logging_handler) 

389 

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 

397 

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 

409 

410 self._handlers_registered = True 

411 

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) 

420 

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 ) 

432 

433 registrar = cast("ContainerRegistrarProtocol", container) 

434 registrar.singleton(MCPServer, server) 

435 

436 def _boot_transports(self, container: ContainerRegistrarProtocol) -> None: 

437 """Register MCP transport singletons.""" 

438 stdio_transport = StdioTransport() 

439 sse_transport = SSETransport() 

440 

441 container.singleton(StdioTransport, stdio_transport) 

442 container.singleton(SSETransport, sse_transport) 

443 self._managed_transports = [stdio_transport, sse_transport] 

444 

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 ) 

456 

457 self._managed_transports.clear() 

458 

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 

472 

473 logger.info("mcp_provider_shutdown_complete") 

474 

475 async def health_check(self, timeout: float = 5.0) -> HealthCheckResult: 

476 """Return a health check result for this provider. 

477 

478 Args: 

479 timeout: Maximum seconds to wait for health response. 

480 

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 ) 

500 

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 ) 

508 

509 self._managed_client = None 

510 

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 ) 

543 

544 

545__all__ = ["MCPProvider"]