/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/tools/mcp/custom.py
1 549 строк
50 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
""" Custom MCP Server — универсальный MCP клиент для подключения к любым внешним MCP серверам. Реализует Model Context Protocol (MCP) спецификацию: https://modelcontextprotocol.io/specification Поддерживаемые транспорты: - stdio: запуск процесса и общение через stdin/stdout - sse: Server-Sent Events over HTTP - streamable-http: Streamable HTTP (новый стандарт MCP) Архитектурные принципы: - Чистые функции для преобразования JSON-RPC сообщений - Явная обработка ошибок через типизированные исключения - Transport abstraction — единый интерфейс для разных транспортов - Async-first дизайн с context managers - Делегирование через явные _require_* методы для type safety Использование: async with create_custom_mcp( type="stdio", command="npx", args=["-y", "@modelcontextprotocol/server-github"], env={"GITHUB_PERSONAL_ACCESS_TOKEN": "..."}, ) as mcp: # Динамически получить tools tools = await mcp.list_tools() # Вызвать tool result = await mcp.call_tool("create_issue", {"title": "Bug"}) """ from __future__ import annotations import asyncio import json import logging import os from collections.abc import AsyncIterator from contextlib import asynccontextmanager from dataclasses import dataclass, field from datetime import datetime from enum import Enum from typing import Any, TYPE_CHECKING from src.primitives.context import MCPContext if TYPE_CHECKING: import aiohttp # type: ignore[import-not-found,import-untyped] logger = logging.getLogger(__name__) # ============================================================================ # Exceptions # ============================================================================ class MCPError(Exception): """Базовое исключение MCP.""" pass class MCPConnectionError(MCPError): """Ошибка подключения к MCP серверу.""" pass class MCPTransportError(MCPError): """Ошибка транспорта (stdio/sse/http).""" pass class MCPProtocolError(MCPError): """Ошибка MCP протокола (JSON-RPC).""" def __init__( self, message: str, code: int = -32000, data: Any = None, ): super().__init__(message) self.code = code self.data = data class MCPToolNotFoundError(MCPError): """Tool не найден на сервере.""" pass class MCPTimeoutError(MCPError): """Превышен таймаут ожидания ответа.""" pass # ============================================================================ # Enums # ============================================================================ class MCPTransportType(str, Enum): """Типы транспортов MCP.""" STDIO = "stdio" SSE = "sse" STREAMABLE_HTTP = "streamable-http" # ============================================================================ # Configuration # ============================================================================ @dataclass class MCPConfig: """ Конфигурация подключения к MCP серверу. Чистая структура данных — валидация и преобразования вынесены в отдельные функции. Примеры: # stdio MCPConfig( type=MCPTransportType.STDIO, command="npx", args=["-y", "@modelcontextprotocol/server-github"], env={"GITHUB_PERSONAL_ACCESS_TOKEN": "..."}, ) # SSE MCPConfig( type=MCPTransportType.SSE, url="https://mcp.example.com/sse", headers={"Authorization": "Bearer ..."}, ) # Streamable HTTP MCPConfig( type=MCPTransportType.STREAMABLE_HTTP, url="https://mcp.example.com/mcp", headers={"Authorization": "Bearer ..."}, ) """ # Тип транспорта type: MCPTransportType = MCPTransportType.STDIO # Для stdio command: str | None = None args: list[str] = field(default_factory=list) env: dict[str, str] = field(default_factory=dict) cwd: str | None = None # Для sse и streamable-http url: str | None = None headers: dict[str, str] = field(default_factory=dict) # Общие настройки timeout_seconds: float = 30.0 connect_timeout_seconds: float = 10.0 max_retries: int = 3 retry_delay_seconds: float = 1.0 # MCP protocol client_name: str = "flowstack-mcp-client" client_version: str = "1.0.0" protocol_version: str = "2024-11-05" # Идентификатор server_name: str | None = None def validate(self) -> list[str]: """ Валидировать конфигурацию. Чистая функция — возвращает список ошибок. """ errors: list[str] = [] if self.type == MCPTransportType.STDIO: if not self.command: errors.append("command is required for stdio transport") elif self.type in (MCPTransportType.SSE, MCPTransportType.STREAMABLE_HTTP): if not self.url: errors.append("url is required for HTTP-based transports") else: errors.append(f"Unknown transport type: {self.type}") if self.timeout_seconds <= 0: errors.append("timeout_seconds must be positive") return errors @classmethod def from_dict(cls, data: dict[str, Any]) -> MCPConfig: """Создать из словаря. Чистая функция.""" transport_type_raw = data.get("type", "stdio") try: transport_type = MCPTransportType(transport_type_raw) except ValueError: transport_type = MCPTransportType.STDIO return cls( type=transport_type, command=data.get("command"), args=list(data.get("args") or []), env=dict(data.get("env") or {}), cwd=data.get("cwd"), url=data.get("url"), headers=dict(data.get("headers") or {}), timeout_seconds=float(data.get("timeout_seconds", 30.0)), connect_timeout_seconds=float(data.get("connect_timeout_seconds", 10.0)), max_retries=int(data.get("max_retries", 3)), retry_delay_seconds=float(data.get("retry_delay_seconds", 1.0)), client_name=str(data.get("client_name", "flowstack-mcp-client")), client_version=str(data.get("client_version", "1.0.0")), protocol_version=str(data.get("protocol_version", "2024-11-05")), server_name=data.get("server_name"), ) # ============================================================================ # JSON-RPC Messages # ============================================================================ @dataclass class JSONRPCRequest: """JSON-RPC 2.0 Request.""" id: int | str method: str params: dict[str, Any] | None = None def to_dict(self) -> dict[str, Any]: """Преобразовать в JSON-RPC dict. Чистая функция.""" message: dict[str, Any] = { "jsonrpc": "2.0", "id": self.id, "method": self.method, } if self.params is not None: message["params"] = self.params return message def to_json(self) -> str: """Сериализовать в JSON строку. Чистая функция.""" return json.dumps(self.to_dict(), ensure_ascii=False) @dataclass class JSONRPCNotification: """JSON-RPC 2.0 Notification (без id).""" method: str params: dict[str, Any] | None = None def to_dict(self) -> dict[str, Any]: """Преобразовать в JSON-RPC dict. Чистая функция.""" message: dict[str, Any] = { "jsonrpc": "2.0", "method": self.method, } if self.params is not None: message["params"] = self.params return message def to_json(self) -> str: """Сериализовать в JSON строку. Чистая функция.""" return json.dumps(self.to_dict(), ensure_ascii=False) @dataclass class JSONRPCError: """JSON-RPC 2.0 Error.""" code: int message: str data: Any = None def to_dict(self) -> dict[str, Any]: """Преобразовать в dict. Чистая функция.""" result: dict[str, Any] = { "code": self.code, "message": self.message, } if self.data is not None: result["data"] = self.data return result @dataclass class JSONRPCResponse: """JSON-RPC 2.0 Response.""" id: int | str | None result: Any = None error: JSONRPCError | None = None @classmethod def from_dict(cls, data: dict[str, Any]) -> JSONRPCResponse: """ Создать из dict. Чистая функция. Обрабатывает как success responses (с result), так и error responses (с error). """ error_data = data.get("error") error: JSONRPCError | None = None if isinstance(error_data, dict): error = JSONRPCError( code=int(error_data.get("code", -32000)), message=str(error_data.get("message", "Unknown error")), data=error_data.get("data"), ) return cls( id=data.get("id"), result=data.get("result"), error=error, ) def is_error(self) -> bool: """Проверить, является ли ответ ошибкой. Чистая функция.""" return self.error is not None def parse_jsonrpc_message(raw: str) -> JSONRPCResponse | JSONRPCNotification | None: """ Парсить JSON-RPC сообщение из строки. Чистая функция — возвращает None при невалидном JSON. """ try: data = json.loads(raw) except (json.JSONDecodeError, ValueError): return None if not isinstance(data, dict): return None # Response — есть id и (result или error) if "id" in data and ("result" in data or "error" in data): return JSONRPCResponse.from_dict(data) # Notification — есть method, но нет id if "method" in data and "id" not in data: return JSONRPCNotification( method=str(data["method"]), params=data.get("params"), ) return None # ============================================================================ # MCP Domain Models # ============================================================================ @dataclass class MCPTool: """ Tool, предоставляемый MCP сервером. Соответствует MCP спецификации: https://modelcontextprotocol.io/specification/2024-11-05/server/tools """ name: str description: str = "" input_schema: dict[str, Any] = field(default_factory=dict) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь. Чистая функция.""" return { "name": self.name, "description": self.description, "inputSchema": self.input_schema, } @classmethod def from_dict(cls, data: dict[str, Any]) -> MCPTool: """Создать из словаря. Чистая функция.""" return cls( name=str(data.get("name", "")), description=str(data.get("description", "")), input_schema=data.get("inputSchema") or data.get("input_schema") or {}, ) @dataclass class MCPResource: """ Resource, предоставляемый MCP сервером. Соответствует MCP спецификации: https://modelcontextprotocol.io/specification/2024-11-05/server/resources """ uri: str name: str = "" description: str = "" mime_type: str | None = None def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь. Чистая функция.""" result: dict[str, Any] = { "uri": self.uri, "name": self.name, "description": self.description, } if self.mime_type: result["mimeType"] = self.mime_type return result @classmethod def from_dict(cls, data: dict[str, Any]) -> MCPResource: """Создать из словаря. Чистая функция.""" return cls( uri=str(data.get("uri", "")), name=str(data.get("name", "")), description=str(data.get("description", "")), mime_type=data.get("mimeType") or data.get("mime_type"), ) @dataclass class MCPPrompt: """ Prompt, предоставляемый MCP сервером. Соответствует MCP спецификации: https://modelcontextprotocol.io/specification/2024-11-05/server/prompts """ name: str description: str = "" arguments: list[dict[str, Any]] = field(default_factory=list) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь. Чистая функция.""" return { "name": self.name, "description": self.description, "arguments": self.arguments, } @classmethod def from_dict(cls, data: dict[str, Any]) -> MCPPrompt: """Создать из словаря. Чистая функция.""" args_raw = data.get("arguments") arguments: list[dict[str, Any]] = ( [a for a in args_raw if isinstance(a, dict)] if isinstance(args_raw, list) else [] ) return cls( name=str(data.get("name", "")), description=str(data.get("description", "")), arguments=arguments, ) @dataclass class MCPServerInfo: """ Информация о MCP сервере (из initialize ответа). """ name: str version: str capabilities: dict[str, Any] = field(default_factory=dict) def has_tools(self) -> bool: """Проверить, поддерживает ли сервер tools.""" return "tools" in self.capabilities def has_resources(self) -> bool: """Проверить, поддерживает ли сервер resources.""" return "resources" in self.capabilities def has_prompts(self) -> bool: """Проверить, поддерживает ли сервер prompts.""" return "prompts" in self.capabilities def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь. Чистая функция.""" return { "name": self.name, "version": self.version, "capabilities": self.capabilities, } @classmethod def from_dict(cls, data: dict[str, Any]) -> MCPServerInfo: """Создать из словаря. Чистая функция.""" server_info_raw = data.get("serverInfo") or data.get("server_info") or {} server_info: dict[str, Any] = ( server_info_raw if isinstance(server_info_raw, dict) else {} ) capabilities_raw = data.get("capabilities") capabilities: dict[str, Any] = ( capabilities_raw if isinstance(capabilities_raw, dict) else {} ) return cls( name=str(server_info.get("name", "")), version=str(server_info.get("version", "")), capabilities=capabilities, ) @dataclass class MCPToolCallResult: """Результат вызова MCP tool.""" content: list[dict[str, Any]] = field(default_factory=list) is_error: bool = False duration_ms: int = 0 def get_text(self) -> str: """ Извлечь текстовое содержимое. Объединяет все content элементы с типом "text". """ texts: list[str] = [] for item in self.content: if isinstance(item, dict) and item.get("type") == "text": text = item.get("text") if isinstance(text, str): texts.append(text) return "\n\n".join(texts) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь. Чистая функция.""" return { "content": self.content, "isError": self.is_error, "duration_ms": self.duration_ms, } # ============================================================================ # Transport Layer # ============================================================================ class MCPTransport: """ Базовый класс для MCP транспортов. Определяет интерфейс для отправки/получения JSON-RPC сообщений. """ async def connect(self) -> None: """Подключиться к серверу.""" raise NotImplementedError async def close(self) -> None: """Закрыть соединение.""" raise NotImplementedError async def send(self, message: str) -> None: """Отправить JSON-RPC сообщение.""" raise NotImplementedError async def receive(self) -> str | None: """ Получить JSON-RPC сообщение. Возвращает None при закрытии соединения. """ raise NotImplementedError @property def is_connected(self) -> bool: """Проверить, активно ли соединение.""" return False class StdioTransport(MCPTransport): """ Transport через stdio (stdin/stdout) процесса. Запускает MCP сервер как subprocess и общается через JSON-RPC. Используется для локальных MCP серверов (например, `npx @mcp/server-github`). MCP спецификация: - Сервер запускается как subprocess - Сообщения delimited by newlines - stderr для логирования """ def __init__(self, config: MCPConfig): self.config = config self._process: asyncio.subprocess.Process | None = None self._connected: bool = False async def connect(self) -> None: """Запустить процесс MCP сервера.""" if self._process is not None: return command = self.config.command if not command: raise MCPConnectionError("stdio transport requires 'command'") # Merge environment env = os.environ.copy() env.update(self.config.env) try: self._process = await asyncio.create_subprocess_exec( command, *self.config.args, stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, env=env, cwd=self.config.cwd, ) self._connected = True logger.info(f"Started MCP process: {command} (PID: {self._process.pid})") except FileNotFoundError as e: raise MCPConnectionError(f"Command not found: {command}") from e except Exception as e: raise MCPConnectionError(f"Failed to start process: {e}") from e async def close(self) -> None: """Завершить процесс.""" if self._process is None: return try: if self._process.stdin is not None: self._process.stdin.close() # Подождать graceful shutdown try: await asyncio.wait_for( self._process.wait(), timeout=5.0, ) except asyncio.TimeoutError: # Force kill self._process.terminate() try: await asyncio.wait_for(self._process.wait(), timeout=2.0) except asyncio.TimeoutError: self._process.kill() await self._process.wait() except Exception as e: logger.warning(f"Error closing MCP process: {e}") finally: self._process = None self._connected = False async def send(self, message: str) -> None: """Отправить JSON-RPC сообщение через stdin.""" process = self._require_process() if process.stdin is None: raise MCPTransportError("Process stdin is not available") try: data = (message + "\n").encode("utf-8") process.stdin.write(data) await process.stdin.drain() except Exception as e: raise MCPTransportError(f"Failed to send message: {e}") from e async def receive(self) -> str | None: """Получить JSON-RPC сообщение из stdout.""" process = self._require_process() if process.stdout is None: raise MCPTransportError("Process stdout is not available") try: line = await process.stdout.readline() if not line: # EOF — process завершился self._connected = False return None return line.decode("utf-8").strip() except asyncio.CancelledError: raise except Exception as e: raise MCPTransportError(f"Failed to receive message: {e}") from e @property def is_connected(self) -> bool: """Проверить, активен ли процесс.""" return ( self._connected and self._process is not None and self._process.returncode is None ) def _require_process(self) -> asyncio.subprocess.Process: """Получить активный процесс или выбросить ошибку.""" if self._process is None: raise MCPTransportError("Transport is not connected") return self._process class HTTPTransport(MCPTransport): """ Transport через HTTP (SSE или Streamable HTTP). Реализует MCP Streamable HTTP протокол: https://modelcontextprotocol.io/specification/2024-11-05/basic/transports Использует POST для отправки запросов и опционально GET с SSE для получения уведомлений. """ def __init__(self, config: MCPConfig): self.config = config self._session: aiohttp.ClientSession | None = None self._connected: bool = False self._session_id: str | None = None async def connect(self) -> None: """Установить HTTP соединение.""" if self._session is not None: return try: import aiohttp # type: ignore[import-not-found,import-untyped] except ImportError as e: raise MCPConnectionError( "aiohttp is required for HTTP transport. " "Install it with: pip install aiohttp" ) from e if not self.config.url: raise MCPConnectionError("HTTP transport requires 'url'") timeout = aiohttp.ClientTimeout( total=self.config.timeout_seconds, connect=self.config.connect_timeout_seconds, ) self._session = aiohttp.ClientSession( timeout=timeout, headers=self.config.headers, ) self._connected = True logger.info(f"Connected to MCP server: {self.config.url}") async def close(self) -> None: """Закрыть HTTP соединение.""" if self._session is not None: await self._session.close() self._session = None self._connected = False self._session_id = None async def send(self, message: str) -> None: """ Отправить JSON-RPC сообщение через POST. Для HTTP транспорта send используется только для уведомлений. """ session = self._require_session() url = self.config.url if not url: raise MCPTransportError("URL is not configured") headers = {"Content-Type": "application/json"} if self._session_id: headers["Mcp-Session-Id"] = self._session_id try: async with session.post( url, data=message, headers=headers, ) as response: if response.status >= 400: text = await response.text() raise MCPTransportError( f"HTTP error {response.status}: {text}" ) # Сохранить session ID если сервер его прислал new_session_id = response.headers.get("Mcp-Session-Id") if new_session_id: self._session_id = new_session_id except Exception as e: if "aiohttp" in type(e).__module__: raise MCPTransportError(f"HTTP request failed: {e}") from e raise async def send_and_receive(self, message: str) -> str | None: """ Отправить сообщение и получить ответ (для HTTP). Возвращает None для уведомлений. """ session = self._require_session() url = self.config.url if not url: raise MCPTransportError("URL is not configured") headers = {"Content-Type": "application/json", "Accept": "application/json"} if self._session_id: headers["Mcp-Session-Id"] = self._session_id try: async with session.post( url, data=message, headers=headers, ) as response: if response.status >= 400: text = await response.text() raise MCPTransportError( f"HTTP error {response.status}: {text}" ) # Сохранить session ID new_session_id = response.headers.get("Mcp-Session-Id") if new_session_id: self._session_id = new_session_id # 202 Accepted — уведомление, нет тела if response.status == 202: return None text = await response.text() return text if text.strip() else None except Exception as e: if "aiohttp" in type(e).__module__: raise MCPTransportError(f"HTTP request failed: {e}") from e raise async def receive(self) -> str | None: """ Получить сообщение (для HTTP это не основной путь). HTTP transport использует request-response, поэтому receive возвращает None. """ return None @property def is_connected(self) -> bool: """Проверить, активна ли сессия.""" return self._connected and self._session is not None def _require_session(self) -> aiohttp.ClientSession: """Получить активную сессию или выбросить ошибку.""" if self._session is None: raise MCPTransportError("Transport is not connected") return self._session def create_transport(config: MCPConfig) -> MCPTransport: """ Factory для создания транспорта по конфигу. Чистая функция — не мутирует конфиг, возвращает новый transport. """ if config.type == MCPTransportType.STDIO: return StdioTransport(config) elif config.type in (MCPTransportType.SSE, MCPTransportType.STREAMABLE_HTTP): return HTTPTransport(config) else: raise MCPError(f"Unknown transport type: {config.type}") # ============================================================================ # MCP Client # ============================================================================ class MCPClient: """ Клиент для взаимодействия с MCP сервером. Управляет: - Подключением через transport - Инициализацией MCP handshake - JSON-RPC request/response циклом - Динамическим discovery tools/resources/prompts """ def __init__(self, config: MCPConfig): self.config = config self._transport = create_transport(config) self._request_id: int = 0 self._server_info: MCPServerInfo | None = None async def __aenter__(self) -> MCPClient: """Async context manager entry.""" await self.connect() return self async def __aexit__(self, exc_type, exc_val, exc_tb) -> None: """Async context manager exit.""" await self.close() async def connect(self) -> None: """ Подключиться и инициализировать MCP handshake. 1. Устанавливает транспорт 2. Отправляет initialize запрос 3. Отправляет initialized уведомление """ errors = self.config.validate() if errors: raise MCPConnectionError( f"Invalid configuration: {', '.join(errors)}" ) await self._transport.connect() # MCP handshake try: init_result = await self._request("initialize", { "protocolVersion": self.config.protocol_version, "capabilities": {}, "clientInfo": { "name": self.config.client_name, "version": self.config.client_version, }, }) if isinstance(init_result, dict): self._server_info = MCPServerInfo.from_dict(init_result) # Отправить initialized уведомление await self._notify("notifications/initialized", {}) except Exception as e: await self._transport.close() raise MCPConnectionError(f"MCP handshake failed: {e}") from e async def close(self) -> None: """Закрыть соединение.""" await self._transport.close() self._server_info = None def get_server_info(self) -> MCPServerInfo | None: """Получить информацию о сервере.""" return self._server_info def _next_request_id(self) -> int: """Получить следующий request ID.""" self._request_id += 1 return self._request_id async def _request( self, method: str, params: dict[str, Any] | None = None, ) -> Any: """ Отправить JSON-RPC запрос и получить ответ. Args: method: JSON-RPC метод params: Параметры запроса Returns: result из JSON-RPC response Raises: MCPProtocolError: если сервер вернул ошибку MCPTimeoutError: если превышен таймаут MCPTransportError: если проблема с транспортом """ request_id = self._next_request_id() request = JSONRPCRequest( id=request_id, method=method, params=params, ) # Для HTTP транспорта используем send_and_receive if isinstance(self._transport, HTTPTransport): response_text = await asyncio.wait_for( self._transport.send_and_receive(request.to_json()), timeout=self.config.timeout_seconds, ) if response_text is None: # Уведомление без ответа (не должно быть для request) raise MCPProtocolError("No response from server") message = parse_jsonrpc_message(response_text) if not isinstance(message, JSONRPCResponse): raise MCPProtocolError("Invalid response format") if message.error is not None: raise MCPProtocolError( message.error.message, code=message.error.code, data=message.error.data, ) return message.result # Для stdio транспорта отправляем и читаем await self._transport.send(request.to_json()) # Читать до получения ответа с нашим id try: async with asyncio.timeout(self.config.timeout_seconds): while True: response_text = await self._transport.receive() if response_text is None: raise MCPTransportError("Connection closed by server") message = parse_jsonrpc_message(response_text) # Игнорируем уведомления if isinstance(message, JSONRPCNotification): logger.debug( f"Received notification: {message.method}" ) continue if not isinstance(message, JSONRPCResponse): continue # Наш ответ? if message.id == request_id: if message.error is not None: raise MCPProtocolError( message.error.message, code=message.error.code, data=message.error.data, ) return message.result # Чужой ответ — логируем logger.warning( f"Received response with unexpected id: {message.id}" ) except TimeoutError: raise MCPTimeoutError( f"Request {method} timed out after {self.config.timeout_seconds}s" ) async def _notify( self, method: str, params: dict[str, Any] | None = None, ) -> None: """ Отправить JSON-RPC уведомление (без ожидания ответа). """ notification = JSONRPCNotification(method=method, params=params) await self._transport.send(notification.to_json()) # ==================================================================== # MCP High-Level API # ==================================================================== async def list_tools(self) -> list[MCPTool]: """ Получить список tools, предоставляемых сервером. MCP метод: tools/list """ result = await self._request("tools/list", {}) tools_raw = result.get("tools") if isinstance(result, dict) else None tools_data: list[Any] = tools_raw if isinstance(tools_raw, list) else [] return [ MCPTool.from_dict(t) for t in tools_data if isinstance(t, dict) ] async def call_tool( self, name: str, arguments: dict[str, Any] | None = None, ) -> MCPToolCallResult: """ Вызвать tool на сервере. MCP метод: tools/call """ start_time = datetime.utcnow() try: result = await self._request("tools/call", { "name": name, "arguments": arguments or {}, }) except MCPProtocolError as e: # Tool not found — проверяем текст ошибки через str() error_text = str(e).lower() if "not found" in error_text or "unknown tool" in error_text: raise MCPToolNotFoundError(f"Tool '{name}' not found: {e}") from e raise duration_ms = int((datetime.utcnow() - start_time).total_seconds() * 1000) if not isinstance(result, dict): return MCPToolCallResult( content=[], is_error=False, duration_ms=duration_ms, ) return MCPToolCallResult( content=list(result.get("content") or []), is_error=bool(result.get("isError", False)), duration_ms=duration_ms, ) async def list_resources(self) -> list[MCPResource]: """ Получить список resources. MCP метод: resources/list """ result = await self._request("resources/list", {}) resources_raw = result.get("resources") if isinstance(result, dict) else None resources_data: list[Any] = ( resources_raw if isinstance(resources_raw, list) else [] ) return [ MCPResource.from_dict(r) for r in resources_data if isinstance(r, dict) ] async def read_resource(self, uri: str) -> dict[str, Any]: """ Прочитать resource по URI. MCP метод: resources/read """ result = await self._request("resources/read", {"uri": uri}) return result if isinstance(result, dict) else {} async def list_prompts(self) -> list[MCPPrompt]: """ Получить список prompts. MCP метод: prompts/list """ result = await self._request("prompts/list", {}) prompts_raw = result.get("prompts") if isinstance(result, dict) else None prompts_data: list[Any] = ( prompts_raw if isinstance(prompts_raw, list) else [] ) return [ MCPPrompt.from_dict(p) for p in prompts_data if isinstance(p, dict) ] async def get_prompt( self, name: str, arguments: dict[str, Any] | None = None, ) -> dict[str, Any]: """ Получить prompt по имени с аргументами. MCP метод: prompts/get """ result = await self._request("prompts/get", { "name": name, "arguments": arguments or {}, }) return result if isinstance(result, dict) else {} # ============================================================================ # Custom MCP Server (FlowStack обёртка) # ============================================================================ class CustomMCPServer: """ Обёртка над MCPClient для интеграции в FlowStack. Предоставляет унифицированный интерфейс для работы с любыми внешними MCP серверами через FlowStack tools API. Следует принципу "You Might Not Need an Effect": - Состояние клиента управляется через context manager - Нет скрытых side effects - Явная инициализация и очистка ресурсов - Все методы делегируют через _require_client() """ def __init__(self, config: MCPConfig | None = None): self.config = config or MCPConfig() self._client: MCPClient | None = None self._cached_tools: list[MCPTool] | None = None async def __aenter__(self) -> CustomMCPServer: """Async context manager entry.""" self._client = MCPClient(self.config) await self._client.connect() return self async def __aexit__(self, exc_type, exc_val, exc_tb) -> None: """Async context manager exit.""" if self._client is not None: await self._client.close() self._client = None self._cached_tools = None def _require_client(self) -> MCPClient: """Получить активный клиент или выбросить ошибку.""" if self._client is None: raise MCPError( "Server is not connected. Use 'async with' to manage connection." ) return self._client def get_server_info(self) -> MCPServerInfo | None: """Получить информацию о подключённом сервере.""" if self._client is None: return None return self._client.get_server_info() async def list_tools(self, use_cache: bool = True) -> list[MCPTool]: """ Получить список tools сервера. Args: use_cache: Использовать кэш (по умолчанию True) """ if use_cache and self._cached_tools is not None: return list(self._cached_tools) client = self._require_client() tools = await client.list_tools() self._cached_tools = tools return tools async def get_tools_schema(self) -> list[dict[str, Any]]: """ Получить tools в формате для LLM (OpenAI-like function calling). Преобразует MCP tools в формат, совместимый с OpenAI function calling / Anthropic tool_use. """ tools = await self.list_tools() return [ { "type": "function", "function": { "name": tool.name, "description": tool.description, "parameters": tool.input_schema or {"type": "object"}, }, } for tool in tools ] async def call_tool( self, name: str, arguments: dict[str, Any], context: MCPContext | None = None, ) -> dict[str, Any]: """ Вызвать tool по имени. Args: name: Имя tool arguments: Аргументы context: MCP контекст для метрик Returns: Результат в формате {"success": bool, "result"|"error": ..., "error_type"?} """ try: client = self._require_client() except MCPError as e: return { "success": False, "error": str(e), "error_type": "not_connected", } try: result = await client.call_tool(name, arguments) if context is not None: context.add_tool_call() return { "success": not result.is_error, "result": result.to_dict(), "duration_ms": result.duration_ms, } except MCPToolNotFoundError as e: return { "success": False, "error": str(e), "error_type": "tool_not_found", } except MCPTimeoutError as e: return { "success": False, "error": str(e), "error_type": "timeout", } except MCPProtocolError as e: return { "success": False, "error": f"Protocol error: {e}", "error_type": "protocol_error", "code": e.code, "data": e.data, } except MCPError as e: return { "success": False, "error": str(e), "error_type": "mcp_error", } except Exception as e: logger.exception(f"Unexpected error calling tool {name}") return { "success": False, "error": f"Unexpected error: {e}", "error_type": "unexpected", } async def list_resources(self) -> list[MCPResource]: """Получить список resources.""" client = self._require_client() return await client.list_resources() async def read_resource(self, uri: str) -> dict[str, Any]: """Прочитать resource по URI.""" client = self._require_client() return await client.read_resource(uri) async def list_prompts(self) -> list[MCPPrompt]: """Получить список prompts.""" client = self._require_client() return await client.list_prompts() async def get_prompt( self, name: str, arguments: dict[str, Any] | None = None, ) -> dict[str, Any]: """Получить prompt по имени.""" client = self._require_client() return await client.get_prompt(name, arguments) def invalidate_cache(self) -> None: """Очистить кэш tools (например, при обновлении сервера).""" self._cached_tools = None # ============================================================================ # Helper Functions (public API) # ============================================================================ @asynccontextmanager async def create_custom_mcp( type: MCPTransportType | str = MCPTransportType.STDIO, command: str | None = None, args: list[str] | None = None, env: dict[str, str] | None = None, url: str | None = None, headers: dict[str, str] | None = None, timeout_seconds: float = 30.0, **kwargs: Any, ) -> AsyncIterator[CustomMCPServer]: """ Создать Custom MCP Server с автоматическим управлением подключением. Args: type: Тип транспорта (stdio, sse, streamable-http) command: Команда для запуска (для stdio) args: Аргументы команды (для stdio) env: Переменные окружения (для stdio) url: URL сервера (для sse/streamable-http) headers: HTTP headers (для sse/streamable-http) timeout_seconds: Таймаут запросов **kwargs: Дополнительные параметры MCPConfig Usage: # stdio (локальный MCP server) async with create_custom_mcp( type="stdio", command="npx", args=["-y", "@modelcontextprotocol/server-github"], env={"GITHUB_PERSONAL_ACCESS_TOKEN": "ghp_..."}, ) as mcp: tools = await mcp.list_tools() result = await mcp.call_tool("create_issue", {...}) # HTTP (удалённый MCP server) async with create_custom_mcp( type="streamable-http", url="https://mcp.example.com/mcp", headers={"Authorization": "Bearer ..."}, ) as mcp: tools = await mcp.list_tools() """ try: transport_type = MCPTransportType(type) except ValueError: transport_type = MCPTransportType.STDIO config = MCPConfig( type=transport_type, command=command, args=list(args or []), env=dict(env or {}), url=url, headers=dict(headers or {}), timeout_seconds=timeout_seconds, **{k: v for k, v in kwargs.items() if hasattr(MCPConfig, k)}, ) async with CustomMCPServer(config) as server: yield server def mcp_config_from_env(prefix: str = "MCP_") -> MCPConfig: """ Создать MCPConfig из переменных окружения. Читает переменные: - {prefix}TYPE (stdio|sse|streamable-http) - {prefix}COMMAND - {prefix}ARGS (JSON array) - {prefix}ENV (JSON object) - {prefix}URL - {prefix}HEADERS (JSON object) - {prefix}TIMEOUT """ transport_type = os.getenv(f"{prefix}TYPE", "stdio") args_raw = os.getenv(f"{prefix}ARGS") env_raw = os.getenv(f"{prefix}ENV") headers_raw = os.getenv(f"{prefix}HEADERS") timeout_raw = os.getenv(f"{prefix}TIMEOUT") args: list[str] = [] if args_raw: try: parsed = json.loads(args_raw) if isinstance(parsed, list): args = [str(x) for x in parsed] except (json.JSONDecodeError, ValueError): pass env_dict: dict[str, str] = {} if env_raw: try: parsed = json.loads(env_raw) if isinstance(parsed, dict): env_dict = {str(k): str(v) for k, v in parsed.items()} except (json.JSONDecodeError, ValueError): pass headers_dict: dict[str, str] = {} if headers_raw: try: parsed = json.loads(headers_raw) if isinstance(parsed, dict): headers_dict = {str(k): str(v) for k, v in parsed.items()} except (json.JSONDecodeError, ValueError): pass timeout = 30.0 if timeout_raw: try: timeout = float(timeout_raw) except ValueError: pass try: type_enum = MCPTransportType(transport_type) except ValueError: type_enum = MCPTransportType.STDIO return MCPConfig( type=type_enum, command=os.getenv(f"{prefix}COMMAND"), args=args, env=env_dict, url=os.getenv(f"{prefix}URL"), headers=headers_dict, timeout_seconds=timeout, ) # ============================================================================ # Exports # ============================================================================ __all__ = [ # Exceptions "MCPError", "MCPConnectionError", "MCPTransportError", "MCPProtocolError", "MCPToolNotFoundError", "MCPTimeoutError", # Enums "MCPTransportType", # Config "MCPConfig", # JSON-RPC "JSONRPCRequest", "JSONRPCNotification", "JSONRPCResponse", "JSONRPCError", "parse_jsonrpc_message", # Domain models "MCPTool", "MCPResource", "MCPPrompt", "MCPServerInfo", "MCPToolCallResult", # Transport "MCPTransport", "StdioTransport", "HTTPTransport", "create_transport", # Client "MCPClient", "CustomMCPServer", # Helpers "create_custom_mcp", "mcp_config_from_env", ]