/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/runtime/mcp_manager.py
378 строк
15 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""MCP Manager — реальное взаимодействие с MCP серверами. Использует официальный ``mcp`` SDK (pip install mcp): - **stdio**: запуск subprocess, JSON-RPC через stdin/stdout - **sse**: Server-Sent Events подключение - **http**: Streamable HTTP (с fallback на SSE) Управляет долгоживущими подключениями (start/stop) и ad-hoc тестами. """ from __future__ import annotations import asyncio import time from contextlib import AsyncExitStack from dataclasses import dataclass, field from typing import Any import structlog logger = structlog.get_logger() # ============================================================================ # MCP SDK import (опциональная зависимость) # ============================================================================ try: from mcp import ClientSession, StdioServerParameters from mcp.client.sse import sse_client try: from mcp.client.streamable_http import streamablehttp_client except ImportError: # старые версии SDK streamablehttp_client = None # type: ignore[assignment] MCP_AVAILABLE = True except ImportError: MCP_AVAILABLE = False ClientSession = None # type: ignore[assignment] StdioServerParameters = None # type: ignore[assignment] sse_client = None # type: ignore[assignment] streamablehttp_client = None # type: ignore[assignment] def _require_mcp() -> None: """Проверить что mcp SDK установлен.""" if not MCP_AVAILABLE: raise RuntimeError("Пакет 'mcp' не установлен. Выполните: pip install mcp") # ============================================================================ # Connection Info # ============================================================================ @dataclass class ConnectionInfo: """Информация об активном подключении к MCP серверу.""" server_id: str session: Any # ClientSession exit_stack: AsyncExitStack config: dict[str, Any] connected_at: float = field(default_factory=time.time) last_used_at: float = field(default_factory=time.time) def touch(self) -> None: """Обновить время последнего использования.""" self.last_used_at = time.time() # ============================================================================ # MCP Manager # ============================================================================ class MCPManager: """ Менеджер подключений к MCP серверам. Реализует protocol из ``mcp_service.MCPManager``. Хранит активные сессии (server_id → ConnectionInfo). """ def __init__(self, default_timeout: float = 30.0) -> None: _require_mcp() self._connections: dict[str, ConnectionInfo] = {} self._lock = asyncio.Lock() self.default_timeout = default_timeout logger.info("mcp_manager.initialized") # ======================================================================== # Transport helpers # ======================================================================== async def _open_transport( self, exit_stack: AsyncExitStack, config: dict[str, Any] ) -> tuple[Any, Any]: """ Открыть transport (stdio/sse/http) через exit_stack. Returns: (read_stream, write_stream) для ClientSession. """ server_type = config.get("type", "stdio") timeout = config.get("timeout_seconds", self.default_timeout) if server_type == "stdio": command = config.get("command") if not command: raise ValueError("stdio сервер требует 'command'") params = StdioServerParameters( command=command, args=config.get("args", []), env=config.get("env") or {}, ) return await exit_stack.enter_async_context(stdio_client(params)) # http / sse url = config.get("url") if not url: raise ValueError(f"{server_type} сервер требует 'url'") if server_type == "http" and streamablehttp_client is not None: return await exit_stack.enter_async_context(streamablehttp_client(url, timeout=timeout)) # SSE (или fallback для http) return await exit_stack.enter_async_context(sse_client(url, timeout=timeout)) async def _create_session(self, exit_stack: AsyncExitStack, config: dict[str, Any]) -> Any: """Открыть transport + создать и инициализировать ClientSession.""" read_stream, write_stream = await self._open_transport(exit_stack, config) session = await exit_stack.enter_async_context(ClientSession(read_stream, write_stream)) await session.initialize() return session def _get_connection(self, server_id: str) -> ConnectionInfo: """Получить активное подключение (или бросить ошибку).""" conn = self._connections.get(server_id) if conn is None: raise ValueError( f"MCP server '{server_id}' не подключён. Сначала вызовите start_server()." ) return conn # ======================================================================== # Lifecycle (start / stop) # ======================================================================== async def start_server(self, server_id: str, config: dict[str, Any]) -> bool: """ Запустить MCP сервер и установить сессию. Для stdio — запускает subprocess. Для http/sse — подключается к URL. """ async with self._lock: if server_id in self._connections: logger.info("mcp_manager.already_started", server_id=server_id) return True exit_stack = AsyncExitStack() try: session = await self._create_session(exit_stack, config) self._connections[server_id] = ConnectionInfo( server_id=server_id, session=session, exit_stack=exit_stack, config=config, ) logger.info( "mcp_manager.server_started", server_id=server_id, type=config.get("type"), ) return True except Exception as e: await exit_stack.aclose() logger.error("mcp_manager.start_failed", server_id=server_id, error=str(e)) raise async def stop_server(self, server_id: str) -> bool: """Остановить MCP сервер (закрыть сессию и subprocess).""" async with self._lock: conn = self._connections.pop(server_id, None) if conn is None: return False try: await conn.exit_stack.aclose() logger.info("mcp_manager.server_stopped", server_id=server_id) return True except Exception as e: logger.warning("mcp_manager.stop_error", server_id=server_id, error=str(e)) return False async def shutdown(self) -> None: """Закрыть все активные подключения (вызывать при shutdown приложения).""" async with self._lock: server_ids = list(self._connections.keys()) for server_id in server_ids: conn = self._connections.pop(server_id, None) if conn is not None: try: await conn.exit_stack.aclose() except Exception as e: logger.warning( "mcp_manager.shutdown_error", server_id=server_id, error=str(e) ) logger.info("mcp_manager.shutdown_complete", closed=len(server_ids)) # ======================================================================== # Discovery (tools / resources) # ======================================================================== async def list_tools(self, server_id: str) -> list[dict[str, Any]]: """Получить список tools сервера.""" conn = self._get_connection(server_id) conn.touch() timeout = conn.config.get("timeout_seconds", self.default_timeout) result = await asyncio.wait_for(conn.session.list_tools(), timeout=timeout) return [ { "name": tool.name, "description": getattr(tool, "description", None), "input_schema": getattr(tool, "inputSchema", {}), } for tool in result.tools ] async def list_resources(self, server_id: str) -> list[dict[str, Any]]: """Получить список resources сервера.""" conn = self._get_connection(server_id) conn.touch() timeout = conn.config.get("timeout_seconds", self.default_timeout) try: result = await asyncio.wait_for(conn.session.list_resources(), timeout=timeout) return [ { "uri": str(resource.uri), "name": getattr(resource, "name", None), "description": getattr(resource, "description", None), "mime_type": getattr(resource, "mimeType", None), } for resource in result.resources ] except Exception as e: # Некоторые серверы не поддерживают resources logger.debug( "mcp_manager.list_resources_unsupported", server_id=server_id, error=str(e) ) return [] # ======================================================================== # Tool call # ======================================================================== async def call_tool(self, server_id: str, tool_name: str, arguments: dict[str, Any]) -> Any: """Вызвать инструмент MCP сервера.""" conn = self._get_connection(server_id) conn.touch() timeout = conn.config.get("timeout_seconds", self.default_timeout) result = await asyncio.wait_for( conn.session.call_tool(tool_name, arguments), timeout=timeout ) if getattr(result, "isError", False): error_text = self._extract_content_text(result.content) raise RuntimeError(f"Tool '{tool_name}' error: {error_text}") return self._extract_content(result.content) @staticmethod def _extract_content(content: Any) -> Any: """Извлечь контент из результата вызова tool.""" if not content: return None extracted = [] for block in content: if hasattr(block, "text"): extracted.append(block.text) elif hasattr(block, "model_dump"): extracted.append(block.model_dump()) else: extracted.append(str(block)) # Один текстовый блок — возвращаем строкой, иначе список if len(extracted) == 1 and isinstance(extracted[0], str): return extracted[0] return extracted @staticmethod def _extract_content_text(content: Any) -> str: """Извлечь текст из контент-блоков (для ошибок).""" if not content: return "Unknown error" parts = [] for block in content: if hasattr(block, "text"): parts.append(block.text) else: parts.append(str(block)) return "\n".join(parts) # ======================================================================== # Test connection # ======================================================================== async def test_connection(self, server_id: str, config: dict[str, Any]) -> bool: """ Проверить соединение (ad-hoc, без сохранения в пул). Открывает временную сессию, отправляет ping, закрывает. """ exit_stack = AsyncExitStack() try: session = await self._create_session(exit_stack, config) # ping для проверки живости соединения if hasattr(session, "send_ping"): await session.send_ping() return True except Exception as e: logger.warning("mcp_manager.test_failed", server_id=server_id, error=str(e)) return False finally: await exit_stack.aclose() # ======================================================================== # Monitoring # ======================================================================== def is_connected(self, server_id: str) -> bool: """Проверить, подключён ли сервер.""" return server_id in self._connections def list_active(self) -> list[dict[str, Any]]: """Список активных подключений.""" return [ { "server_id": conn.server_id, "type": conn.config.get("type"), "connected_at": conn.connected_at, "last_used_at": conn.last_used_at, "uptime_seconds": round(time.time() - conn.connected_at, 1), } for conn in self._connections.values() ] async def get_status(self) -> dict[str, Any]: """Статус менеджера.""" return { "mcp_available": MCP_AVAILABLE, "active_connections": len(self._connections), "connections": self.list_active(), } # ============================================================================ # Singleton # ============================================================================ _mcp_manager: MCPManager | None = None def get_mcp_manager() -> MCPManager: """Получить singleton MCP менеджер.""" global _mcp_manager # noqa: PLW0603 if _mcp_manager is None: _mcp_manager = MCPManager() return _mcp_manager async def shutdown_mcp_manager() -> None: """Закрыть singleton MCP менеджер (для lifespan shutdown).""" global _mcp_manager # noqa: PLW0603 if _mcp_manager is not None: await _mcp_manager.shutdown() _mcp_manager = None