/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/tools/registry.py
936 строк
33 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
""" Реестр инструментов (Tool Registry). Объединяет все tools в единый managed registry: - Встроенные tools: calculator, browser, code_execution, diagram_gen, file_ops, python_repl, web_search - MCP server tools: confluence, github, gitlab, google_workspace, jira, linear, postgres, custom — автоматически адаптируются через `MCPToolAdapter` в стандартный `Tool` interface (MCP 2024-11-05) - Custom tools: пользовательские tools, зарегистрированные через `register_tool()` Архитектура: - `MCPToolAdapter` — превращает MCP tool definition в `Tool` subclass (динамический type()) - `MCPBundle` — lazy-loading wrapper для MCP server с lifecycle management - `ToolRegistryManager` — объединяет `ToolRegistry` + MCP bundles + lifecycle hooks Lifecycle (production): async with registry_lifecycle( file_ops_base_dir="/workspace", enable_mcp_bundles=["github", "postgres"], ) as manager: tool = manager.get_tool("calculator") result = await tool.execute(expression="2 + 2") Simple API (singleton): from src.tools.registry import get_tool, list_tools, register_tool tool = get_tool("calculator") """ from __future__ import annotations import json import logging from collections.abc import AsyncIterator, Callable from contextlib import asynccontextmanager from typing import Any, TYPE_CHECKING from src.tools.base import ( TextContent, Tool, ToolRegistry, ToolResult, ) # ============================================================================ # Встроенные tools (builtins) — всегда доступны # ============================================================================ from src.tools.calculator import CalculatorTool from src.tools.code_execution import CodeExecutionConfig, CodeExecutionTool from src.tools.diagram_gen import DiagramConfig, DiagramGenTool from src.tools.file_ops import FileOpsConfig, FileOpsTool from src.tools.python_repl import PythonReplConfig, PythonReplTool # Опциональные builtins — импортируем с runtime проверкой BrowserTool: type[Tool] | None = None _HAS_BROWSER = False try: from src.tools.browser import BrowserTool as _BrowserTool # type: ignore[import-not-found] BrowserTool = _BrowserTool _HAS_BROWSER = True except ImportError: pass WebSearchTool: type[Tool] | None = None _HAS_WEB_SEARCH = False try: from src.tools.web_search import WebSearchTool as _WebSearchTool # type: ignore[import-not-found] WebSearchTool = _WebSearchTool _HAS_WEB_SEARCH = True except ImportError: pass # ============================================================================ # MCP Servers — graceful degradation при отсутствии зависимостей # ============================================================================ logger = logging.getLogger(__name__) def _try_import_mcp(module_name: str) -> Any: """ Безопасный импорт MCP модуля. Returns: Модуль если импорт успешен, None если модуль недоступен. """ try: import importlib return importlib.import_module(f"src.tools.mcp.{module_name}") except ImportError as e: logger.debug(f"MCP module '{module_name}' not available: {e}") return None except Exception as e: logger.warning(f"Failed to import MCP module '{module_name}': {e}") return None _mcp_confluence = _try_import_mcp("confluence") _mcp_github = _try_import_mcp("github") _mcp_gitlab = _try_import_mcp("gitlab") _mcp_google_workspace = _try_import_mcp("google_workspace") _mcp_jira = _try_import_mcp("jira") _mcp_linear = _try_import_mcp("linear") _mcp_postgres = _try_import_mcp("postgres") _mcp_custom = _try_import_mcp("custom") # ============================================================================ # MCP Tool Adapter # ============================================================================ class MCPToolAdapter: """ Адаптер MCP tool definition → `Tool` subclass. Интегрирует MCP tools (tools/list → name + description + inputSchema, tools/call → result) в единый Tool registry согласно MCP 2024-11-05. Архитектурные принципы: - Динамическое создание subclass через `type()` — type-safe для Pylance - MCP result конвертируется в `ToolResult` с TextContent + metadata - MCP ошибки пробрасываются как `ToolResult.failure` - isError флаг из MCP protocol сохраняется в metadata """ @staticmethod def create_tool_class( server: Any, tool_name: str, tool_description: str, input_schema: dict[str, Any] | None, tags: list[str] | None = None, is_read_only: bool = True, ) -> type[Tool]: """ Создать `Tool` subclass, оборачивающий MCP tool. Args: server: MCP server instance (имеет метод `call_tool(name, arguments)`) tool_name: Имя MCP tool (например, "github_create_issue") tool_description: Описание tool input_schema: JSON Schema параметров (MCP inputSchema) tags: Теги для категоризации is_read_only: Read-only flag Returns: Новый `Tool` subclass, готовый к регистрации. """ async def execute_impl(self: Tool, **kwargs: Any) -> ToolResult: """Делегирует выполнение MCP server и конвертирует результат.""" try: mcp_result = await server.call_tool(tool_name, kwargs) except Exception as e: logger.exception(f"MCP tool '{tool_name}' raised exception: {e}") return ToolResult.failure( f"MCP tool '{tool_name}' failed: {type(e).__name__}: {e}", metadata={ "mcp_tool": tool_name, "error_type": "exception", }, ) if not isinstance(mcp_result, dict): return ToolResult.failure( f"MCP tool '{tool_name}' returned invalid result", metadata={"mcp_tool": tool_name, "raw": str(mcp_result)}, ) success = bool(mcp_result.get("success", False)) metadata: dict[str, Any] = { "mcp_tool": tool_name, "mcp_result": mcp_result, } if success: data = mcp_result.get("result") text = MCPToolAdapter._format_result_text(tool_name, data) return ToolResult.success_result( [TextContent(text=text)], metadata=metadata, ) else: error_msg = str(mcp_result.get("error", "Unknown MCP error")) error_type = mcp_result.get("error_type", "mcp_error") metadata["error_type"] = error_type text_lines = [ f"# MCP Tool Error: `{tool_name}`\n", f"**Error type:** {error_type}", f"**Message:** {error_msg}", ] return ToolResult.failure( error=error_msg, content=[TextContent(text="\n".join(text_lines))], metadata=metadata, ) safe_class_name = tool_name.replace("-", "_").replace(".", "_") cls_name = f"MCP_{safe_class_name}_Tool" final_tags = list(tags or []) if "mcp" not in final_tags: final_tags.append("mcp") cls = type( cls_name, (Tool,), { "name": tool_name, "description": tool_description, "input_schema": input_schema, "parameters_schema": input_schema, "execute": execute_impl, "is_read_only": is_read_only, "requires_confirmation": False, "tags": final_tags, }, ) return cls # type: ignore[return-value] @staticmethod def _format_result_text(tool_name: str, data: Any) -> str: """Форматировать MCP result в читаемый markdown текст. Чистая функция.""" lines: list[str] = [ f"# MCP Tool: `{tool_name}`\n", "**Status:** ✅ Success\n", ] if data is None: lines.append("*(no data returned)*") return "\n".join(lines) if isinstance(data, (dict, list)): try: json_str = json.dumps(data, indent=2, ensure_ascii=False, default=str) if len(json_str) > 50000: json_str = json_str[:50000] + "\n... <truncated>" lines.append("## Result\n```json\n" + json_str + "\n```") if isinstance(data, list): lines.append(f"\n*{len(data)} items*") except (TypeError, ValueError): lines.append("## Result\n```\n" + str(data) + "\n```") elif isinstance(data, str): text = data if len(text) > 50000: text = text[:50000] + "\n... <truncated>" lines.append("## Result\n```\n" + text + "\n```") else: lines.append("## Result\n```\n" + str(data) + "\n```") return "\n".join(lines) # ============================================================================ # MCP Bundle # ============================================================================ class MCPBundle: """ Bundle для MCP server с lazy initialization и lifecycle. Управляет: - Подключением к MCP server (через `server_factory` — async context manager) - Регистрацией MCP tools в ToolRegistry - Graceful shutdown при уничтожении bundle """ def __init__( self, name: str, server_factory: Callable[[], Any], tags: list[str] | None = None, enabled: bool = False, description: str = "", ): self.name = name self.server_factory = server_factory self.tags = list(tags or []) self.enabled = enabled self.description = description self._context_manager: Any = None self._server: Any = None self._tools: list[Tool] = [] self._initialized = False self._initialization_error: str | None = None @property def is_initialized(self) -> bool: """Проверить, инициализирован ли bundle.""" return self._initialized and self._server is not None @property def tools_count(self) -> int: """Количество зарегистрированных tools.""" return len(self._tools) async def initialize(self) -> bool: """ Инициализировать bundle: подключиться к MCP server. Returns: True если успешно, False если ошибка. """ if not self.enabled: logger.debug(f"MCP bundle '{self.name}' is disabled, skipping") return False if self._initialized: return True try: self._context_manager = self.server_factory() self._server = await self._context_manager.__aenter__() self._initialized = True logger.info(f"MCP bundle '{self.name}' initialized successfully") return True except Exception as e: self._initialization_error = f"{type(e).__name__}: {e}" logger.warning( f"MCP bundle '{self.name}' failed to initialize: " f"{self._initialization_error}" ) self.enabled = False if self._context_manager is not None: try: await self._context_manager.__aexit__(None, None, None) except Exception: pass self._context_manager = None self._server = None return False async def register_tools(self, registry: ToolRegistry) -> int: """ Зарегистрировать все MCP tools bundle в registry. Args: registry: Target ToolRegistry instance. Returns: Количество зарегистрированных tools. """ if not self.is_initialized: logger.warning( f"Cannot register tools from bundle '{self.name}': " "bundle not initialized" ) return 0 try: tool_defs = self._server.get_tools() except Exception as e: logger.warning( f"Failed to get tools from MCP bundle '{self.name}': {e}" ) return 0 if not isinstance(tool_defs, list): logger.warning( f"MCP bundle '{self.name}' returned invalid tools list: {type(tool_defs)}" ) return 0 registered_count = 0 for tool_def in tool_defs: if not isinstance(tool_def, dict): continue tool_name = tool_def.get("name") if not tool_name or not isinstance(tool_name, str): continue tool_description = str(tool_def.get("description", "")) # MCP 2024-11-05 spec: inputSchema (также поддерживаем старые варианты) input_schema = ( tool_def.get("inputSchema") or tool_def.get("parameters") or tool_def.get("parameters_schema") ) try: tool_cls = MCPToolAdapter.create_tool_class( server=self._server, tool_name=tool_name, tool_description=tool_description, input_schema=input_schema, tags=list(self.tags), ) tool_instance = tool_cls() registry.register_tool(tool_instance) self._tools.append(tool_instance) registered_count += 1 except Exception as e: logger.warning( f"Failed to create tool '{tool_name}' from " f"bundle '{self.name}': {e}" ) logger.info( f"Registered {registered_count} tools from MCP bundle '{self.name}'" ) return registered_count async def shutdown(self) -> None: """Graceful shutdown bundle.""" if self._context_manager is not None: try: await self._context_manager.__aexit__(None, None, None) except Exception as e: logger.warning(f"MCP bundle '{self.name}' shutdown error: {e}") finally: self._context_manager = None self._server = None self._initialized = False self._tools.clear() def get_info(self) -> dict[str, Any]: """Получить информацию о bundle. Чистая функция.""" return { "name": self.name, "enabled": self.enabled, "initialized": self._initialized, "description": self.description, "tags": self.tags, "tools_count": len(self._tools), "tools": [t.name for t in self._tools], "initialization_error": self._initialization_error, } # ============================================================================ # Tool Registry Manager # ============================================================================ class ToolRegistryManager: """ Главный менеджер реестра инструментов. Объединяет: - `ToolRegistry` (из base.py) для storage - `MCPBundle`s для внешних MCP servers (MCP 2024-11-05) - Lifecycle hooks (initialize / shutdown) - Special handling для stateful tools (python_repl) """ def __init__(self) -> None: self.registry = ToolRegistry() self._bundles: dict[str, MCPBundle] = {} self._python_repl: PythonReplTool | None = None self._initialized = False # ---------- Builtins ---------- def register_builtin_tools( self, file_ops_base_dir: str = ".", file_ops_read_only: bool = False, ) -> None: """ Зарегистрировать все встроенные (non-MCP) tools. Args: file_ops_base_dir: Корневая директория для file_ops sandbox file_ops_read_only: Read-only mode для file_ops """ self.registry.register_tool(CalculatorTool()) self.registry.register_tool(CodeExecutionTool(CodeExecutionConfig())) self.registry.register_tool(DiagramGenTool(DiagramConfig())) file_ops_config = FileOpsConfig( base_dir=file_ops_base_dir, read_only=file_ops_read_only, ) self.registry.register_tool(FileOpsTool(file_ops_config)) self._python_repl = PythonReplTool(PythonReplConfig()) self.registry.register_tool(self._python_repl) # Browser tool — опциональный if _HAS_BROWSER and BrowserTool is not None: try: self.registry.register_tool(BrowserTool()) except Exception as e: logger.warning(f"Failed to register BrowserTool: {e}") # Web search tool — опциональный if _HAS_WEB_SEARCH and WebSearchTool is not None: try: self.registry.register_tool(WebSearchTool()) except Exception as e: logger.warning(f"Failed to register WebSearchTool: {e}") logger.info( f"Registered {len(self.registry.list_tools())} builtin tools" ) # ---------- MCP Bundles ---------- def register_mcp_bundle( self, name: str, server_factory: Callable[[], Any], tags: list[str] | None = None, enabled: bool = False, description: str = "", ) -> None: """Зарегистрировать MCP bundle (без инициализации).""" if name in self._bundles: logger.warning(f"MCP bundle '{name}' already registered, replacing") self._bundles[name] = MCPBundle( name=name, server_factory=server_factory, tags=tags, enabled=enabled, description=description, ) def register_default_mcp_bundles(self) -> None: """ Зарегистрировать все известные MCP bundles как disabled по умолчанию. Безопасный дефолт: MCP servers требуют credentials и не должны подключаться без явного запроса. """ if _mcp_confluence is not None: self.register_mcp_bundle( name="confluence", server_factory=_mcp_confluence.create_confluence_mcp, tags=["mcp", "confluence", "wiki", "atlassian"], enabled=False, description="Atlassian Confluence (pages, spaces)", ) if _mcp_github is not None: self.register_mcp_bundle( name="github", server_factory=_mcp_github.create_github_mcp, tags=["mcp", "github", "vcs", "git"], enabled=False, description="GitHub (repos, issues, PRs, code search)", ) if _mcp_gitlab is not None: self.register_mcp_bundle( name="gitlab", server_factory=_mcp_gitlab.create_gitlab_mcp, tags=["mcp", "gitlab", "vcs", "git"], enabled=False, description="GitLab (projects, issues, MRs, sprints)", ) if _mcp_google_workspace is not None: self.register_mcp_bundle( name="google_workspace", server_factory=_mcp_google_workspace.create_google_workspace_mcp, tags=["mcp", "google", "gmail", "drive", "calendar", "docs", "sheets"], enabled=False, description="Google Workspace (Gmail, Drive, Calendar, Docs, Sheets)", ) if _mcp_jira is not None: self.register_mcp_bundle( name="jira", server_factory=_mcp_jira.create_jira_mcp, tags=["mcp", "jira", "atlassian", "issues", "agile"], enabled=False, description="Atlassian Jira (issues, boards, sprints)", ) if _mcp_linear is not None: self.register_mcp_bundle( name="linear", server_factory=_mcp_linear.create_linear_mcp, tags=["mcp", "linear", "issues", "agile"], enabled=False, description="Linear (issues, projects, cycles)", ) if _mcp_postgres is not None: self.register_mcp_bundle( name="postgres", server_factory=_mcp_postgres.create_postgres_mcp, tags=["mcp", "postgres", "database", "sql"], enabled=False, description="PostgreSQL (queries, schema, stats)", ) # ---------- Bundle control ---------- def enable_mcp_bundle(self, name: str) -> bool: """Включить MCP bundle (будет инициализирован при `initialize()`).""" bundle = self._bundles.get(name) if bundle is None: logger.warning(f"MCP bundle '{name}' not found") return False bundle.enabled = True return True def disable_mcp_bundle(self, name: str) -> bool: """Отключить MCP bundle.""" bundle = self._bundles.get(name) if bundle is None: return False bundle.enabled = False return True def list_bundles(self) -> list[dict[str, Any]]: """Получить информацию о всех MCP bundles. Чистая функция.""" return [bundle.get_info() for bundle in self._bundles.values()] def get_bundle(self, name: str) -> MCPBundle | None: """Получить bundle по имени.""" return self._bundles.get(name) # ---------- Lifecycle ---------- async def initialize(self) -> None: """ Инициализировать все enabled MCP bundles. Idempotent — повторный вызов ничего не делает. Bundles с ошибками инициализации пропускаются с warning. """ if self._initialized: logger.debug("ToolRegistryManager already initialized") return for bundle in self._bundles.values(): if bundle.enabled: success = await bundle.initialize() if success: await bundle.register_tools(self.registry) self._initialized = True total_tools = len(self.registry.list_tools()) total_mcp_tools = sum(b.tools_count for b in self._bundles.values()) logger.info( f"Registry initialized: {total_tools} total tools " f"({total_mcp_tools} from MCP bundles)" ) async def shutdown(self) -> None: """ Graceful shutdown всех ресурсов. - Python REPL — закрывает все worker subprocesses - MCP bundles — закрывает соединения с серверами """ if self._python_repl is not None: try: await self._python_repl.shutdown() except Exception as e: logger.warning(f"Python REPL shutdown error: {e}") for bundle in self._bundles.values(): try: await bundle.shutdown() except Exception as e: logger.warning(f"MCP bundle '{bundle.name}' shutdown error: {e}") self._initialized = False logger.info("ToolRegistryManager shutdown complete") # ---------- Tool access shortcuts ---------- def get_tool(self, name: str) -> Tool | None: """Получить tool по имени.""" if self.registry.has_tool(name): try: return self.registry.get_tool(name) except Exception: return None return None def list_tools(self) -> list[Tool]: """Получить список всех tools.""" return self.registry.list_tools() def list_tool_names(self) -> list[str]: """Получить список имён всех tools.""" return self.registry.list_tool_names() def register_tool(self, tool: Tool) -> None: """Зарегистрировать кастомный tool.""" self.registry.register_tool(tool) def unregister_tool(self, name: str) -> bool: """Удалить tool из реестра.""" return self.registry.unregister_tool(name) def has_tool(self, name: str) -> bool: """Проверить наличие tool.""" return self.registry.has_tool(name) def list_tools_by_tag(self, tag: str) -> list[Tool]: """Получить список tools по тегу.""" return [ tool for tool in self.registry.list_tools() if tag in getattr(tool, "tags", []) ] def stats(self) -> dict[str, Any]: """Получить статистику registry. Чистая функция.""" builtin_count = 0 mcp_count = 0 for tool in self.registry.list_tools(): tags = getattr(tool, "tags", []) if "mcp" in tags: mcp_count += 1 else: builtin_count += 1 return { "total_tools": len(self.registry.list_tools()), "builtin_tools": builtin_count, "mcp_tools": mcp_count, "mcp_bundles": len(self._bundles), "mcp_bundles_enabled": sum( 1 for b in self._bundles.values() if b.enabled ), "mcp_bundles_initialized": sum( 1 for b in self._bundles.values() if b.is_initialized ), "initialized": self._initialized, "bundles": [b.get_info() for b in self._bundles.values()], } # ============================================================================ # Singleton Manager # ============================================================================ _manager: ToolRegistryManager | None = None def get_manager() -> ToolRegistryManager: """ Получить singleton ToolRegistryManager. Автоматически регистрирует builtin tools и default MCP bundles (disabled) при первом вызове. Для production используйте `registry_lifecycle()` для proper async initialization/shutdown. """ global _manager if _manager is None: _manager = ToolRegistryManager() _manager.register_builtin_tools() _manager.register_default_mcp_bundles() return _manager def reset_manager() -> None: """Сбросить singleton (для тестов). Не вызывает shutdown.""" global _manager _manager = None # ============================================================================ # Module-level API (singleton-based) # ============================================================================ def get_tool(name: str) -> Tool | None: """Получить tool по имени.""" return get_manager().get_tool(name) def list_tools() -> list[Tool]: """Получить список всех зарегистрированных tools.""" return get_manager().list_tools() def list_tool_names() -> list[str]: """Получить список имён всех tools.""" return get_manager().list_tool_names() def register_tool(tool: Tool) -> None: """Зарегистрировать кастомный tool в реестре.""" get_manager().register_tool(tool) def unregister_tool(name: str) -> bool: """Удалить tool из реестра.""" return get_manager().unregister_tool(name) def has_tool(name: str) -> bool: """Проверить наличие tool.""" return get_manager().has_tool(name) def list_tools_by_tag(tag: str) -> list[Tool]: """Получить tools по тегу.""" return get_manager().list_tools_by_tag(tag) def enable_mcp_bundle(name: str) -> bool: """Включить MCP bundle (до `initialize()`).""" return get_manager().enable_mcp_bundle(name) def disable_mcp_bundle(name: str) -> bool: """Отключить MCP bundle.""" return get_manager().disable_mcp_bundle(name) def list_bundles() -> list[dict[str, Any]]: """Получить информацию о всех MCP bundles.""" return get_manager().list_bundles() def stats() -> dict[str, Any]: """Получить статистику реестра.""" return get_manager().stats() async def initialize() -> None: """Инициализировать все enabled MCP bundles.""" await get_manager().initialize() async def shutdown() -> None: """Graceful shutdown всех ресурсов.""" await get_manager().shutdown() # ============================================================================ # Context Manager: Lifecycle # ============================================================================ @asynccontextmanager async def registry_lifecycle( file_ops_base_dir: str = ".", file_ops_read_only: bool = False, enable_mcp_bundles: list[str] | None = None, disable_mcp_bundles: list[str] | None = None, ) -> AsyncIterator[ToolRegistryManager]: """ Async context manager для полного lifecycle registry. Гарантирует proper initialization и shutdown даже при exceptions. Args: file_ops_base_dir: Корневая директория для file_ops sandbox file_ops_read_only: Read-only mode для file_ops enable_mcp_bundles: Список MCP bundles для включения disable_mcp_bundles: Список MCP bundles для принудительного отключения Usage: async with registry_lifecycle( file_ops_base_dir="/workspace", enable_mcp_bundles=["github", "postgres"], ) as manager: calculator = manager.get_tool("calculator") result = await calculator.execute(expression="2 + 2") """ manager = ToolRegistryManager() manager.register_builtin_tools( file_ops_base_dir=file_ops_base_dir, file_ops_read_only=file_ops_read_only, ) manager.register_default_mcp_bundles() if enable_mcp_bundles: for name in enable_mcp_bundles: manager.enable_mcp_bundle(name) if disable_mcp_bundles: for name in disable_mcp_bundles: manager.disable_mcp_bundle(name) try: await manager.initialize() yield manager finally: await manager.shutdown() # ============================================================================ # Helper: Quick tool invocation # ============================================================================ async def call_tool(name: str, **kwargs: Any) -> ToolResult: """ Быстрый способ вызвать tool по имени. Shortcut для: tool = get_tool(name) if tool: result = await tool.execute(**kwargs) Returns: ToolResult (failure если tool не найден) """ tool = get_tool(name) if tool is None: available = list_tool_names() return ToolResult.failure( f"Tool '{name}' not found. Available: {', '.join(sorted(available))}", metadata={"requested_tool": name, "available_tools": available}, ) return await tool.safe_execute(**kwargs) # ============================================================================ # Exports # ============================================================================ __all__ = [ "ToolRegistryManager", "MCPBundle", "MCPToolAdapter", "get_manager", "reset_manager", "get_tool", "list_tools", "list_tool_names", "register_tool", "unregister_tool", "has_tool", "list_tools_by_tag", "enable_mcp_bundle", "disable_mcp_bundle", "list_bundles", "stats", "initialize", "shutdown", "call_tool", "registry_lifecycle", ]