Spaces:
Paused
Paused
| """ | |
| External Tool Servers — MCP klient (Model Context Protocol, Streamable HTTP). | |
| Umožňuje za běhu připojit libovolný MCP server (veřejný hostovaný, vlastní, | |
| nebo z registru), stáhnout jeho nástroje a nabídnout je hlavnímu agentovi | |
| jako OpenAI function-calling nástroje s prefixem `ext_<server>_<tool>`. | |
| - Transport: Streamable HTTP (spec 2025-06-18) — jediný POST endpoint, | |
| odpověď buď application/json, nebo text/event-stream (SSE) s jednou | |
| message událostí. Obojí je podporováno (veřejné servery používají obě). | |
| - Session: server může vrátit hlavičku Mcp-Session-Id při initialize; | |
| posílá se ve všech dalších požadavcích. Při 4xx se session obnoví. | |
| - Bezpečnost: výstupy externích nástrojů jsou NEDŮVĚRYHODNÁ data — popisy | |
| se agentovi označují prefixem serveru, výsledky ořezává tool_result_max_chars | |
| v agent smyčce. Přidávej jen servery, kterým věříš. | |
| - Persistence: /data/toolservers.json (bucket) nebo .agent/toolservers.json — | |
| stejná logika jako settings; změny platí okamžitě, bez restartu. | |
| Objevování veřejných serverů: search_registry() dotazuje OFICIÁLNÍ registr | |
| registry.modelcontextprotocol.io (metaregistr podporovaný Anthropic/GitHub/ | |
| Microsoft/PulseMCP) a vrací servery s remote streamable-http endpointem. | |
| CATALOG níže obsahuje ručně ověřené veřejné servery a adresáře. | |
| """ | |
| from __future__ import annotations | |
| import json | |
| import logging | |
| import os | |
| import re | |
| import tempfile | |
| import time | |
| from dataclasses import asdict, dataclass, field | |
| from pathlib import Path | |
| from threading import RLock | |
| import httpx | |
| from settings import _persist_dir | |
| logger = logging.getLogger("codeagent.toolservers") | |
| MCP_PROTOCOL = "2025-06-18" | |
| REGISTRY_URL = "https://registry.modelcontextprotocol.io/v0/servers" | |
| DEFAULT_TIMEOUT = 60 | |
| _MASK = "********" | |
| # Ručně ověřené veřejné zdroje (2026-07): hostované Streamable HTTP servery | |
| # připojitelné na jeden klik + registry/adresáře pro hledání dalších. | |
| CATALOG = { | |
| "servers": [ | |
| {"name": "context7", "url": "https://mcp.context7.com/mcp", | |
| "desc": "Aktuální dokumentace a příklady kódu pro tisíce knihoven (Upstash)."}, | |
| {"name": "deepwiki", "url": "https://mcp.deepwiki.com/mcp", | |
| "desc": "Dokumentace a Q&A nad veřejnými GitHub repozitáři (Devin/Cognition)."}, | |
| {"name": "huggingface", "url": "https://huggingface.co/mcp", | |
| "desc": "HF Hub: hledání modelů/datasetů/paperů; s HF_TOKEN i Spaces nástroje."}, | |
| {"name": "microsoft-learn", "url": "https://learn.microsoft.com/api/mcp", | |
| "desc": "Oficiální Microsoft/Azure dokumentace."}, | |
| {"name": "cloudflare-docs", "url": "https://docs.mcp.cloudflare.com/mcp", | |
| "desc": "Cloudflare dokumentace."}, | |
| ], | |
| "registries": [ | |
| {"name": "Oficiální MCP registr", | |
| "url": "https://registry.modelcontextprotocol.io", | |
| "desc": "Metaregistr (Anthropic, GitHub, Microsoft, PulseMCP) — " | |
| "prohledatelný přímo z aplikace (tab Externí nástroje)."}, | |
| {"name": "PulseMCP", "url": "https://www.pulsemcp.com/servers", | |
| "desc": "Největší adresář, 10 000+ serverů, denní aktualizace."}, | |
| {"name": "Smithery", "url": "https://smithery.ai", | |
| "desc": "Registr + hosting: servery běží jako Streamable HTTP endpointy " | |
| "(server.smithery.ai/<jméno>/mcp, vyžaduje API klíč)."}, | |
| {"name": "Glama", "url": "https://glama.ai/mcp/servers", | |
| "desc": "Adresář s hodnocením kvality a bezpečnosti."}, | |
| {"name": "mcp.so", "url": "https://mcp.so", | |
| "desc": "Komunitní katalog."}, | |
| {"name": "Awesome MCP Servers", | |
| "url": "https://github.com/punkpeye/awesome-mcp-servers", | |
| "desc": "Kurátorovaný GitHub seznam."}, | |
| ], | |
| } | |
| # ---------------------------------------------------------------- JSON-RPC/SSE | |
| def parse_mcp_response(content_type: str, body: str) -> dict: | |
| """Vytáhne JSON-RPC odpověď z application/json i text/event-stream těla.""" | |
| if "text/event-stream" in (content_type or ""): | |
| last = None | |
| for line in body.splitlines(): | |
| if not line.startswith("data:"): | |
| continue | |
| try: | |
| obj = json.loads(line[5:].strip()) | |
| except json.JSONDecodeError: | |
| continue | |
| if isinstance(obj, dict) and ("result" in obj or "error" in obj): | |
| last = obj | |
| if last is None: | |
| raise ValueError("SSE odpověď neobsahuje JSON-RPC message") | |
| return last | |
| return json.loads(body) | |
| class MCPError(RuntimeError): | |
| pass | |
| class MCPClient: | |
| """Minimální klient MCP Streamable HTTP (initialize/tools list+call).""" | |
| def __init__(self, url: str, token: str = "", timeout: int = DEFAULT_TIMEOUT): | |
| self.url = url | |
| self.token = token | |
| self.timeout = timeout | |
| self.session_id: str | None = None | |
| self.server_info: dict = {} | |
| self._id = 0 | |
| def _headers(self) -> dict: | |
| headers = { | |
| "Content-Type": "application/json", | |
| "Accept": "application/json, text/event-stream", | |
| "MCP-Protocol-Version": MCP_PROTOCOL, | |
| } | |
| if self.token: | |
| headers["Authorization"] = f"Bearer {self.token}" | |
| if self.session_id: | |
| headers["Mcp-Session-Id"] = self.session_id | |
| return headers | |
| def _rpc(self, method: str, params: dict | None = None, | |
| notification: bool = False) -> dict | None: | |
| payload: dict = {"jsonrpc": "2.0", "method": method} | |
| if params is not None: | |
| payload["params"] = params | |
| if not notification: | |
| self._id += 1 | |
| payload["id"] = self._id | |
| resp = httpx.post(self.url, json=payload, headers=self._headers(), | |
| timeout=self.timeout, follow_redirects=True) | |
| if notification: | |
| return None | |
| if resp.status_code >= 400: | |
| raise MCPError(f"HTTP {resp.status_code}: {resp.text[:300]}") | |
| sid = resp.headers.get("Mcp-Session-Id") or resp.headers.get("mcp-session-id") | |
| if sid: | |
| self.session_id = sid | |
| data = parse_mcp_response(resp.headers.get("Content-Type", ""), resp.text) | |
| if "error" in data: | |
| err = data["error"] | |
| raise MCPError(f"{err.get('code')}: {err.get('message')}") | |
| return data.get("result", {}) | |
| def initialize(self) -> dict: | |
| result = self._rpc("initialize", { | |
| "protocolVersion": MCP_PROTOCOL, | |
| "capabilities": {}, | |
| "clientInfo": {"name": "codeagent", "version": "5.0"}, | |
| }) | |
| self.server_info = result.get("serverInfo", {}) | |
| try: | |
| self._rpc("notifications/initialized", {}, notification=True) | |
| except httpx.HTTPError: | |
| pass # některé servery notifikaci nevyžadují | |
| return result | |
| def list_tools(self) -> list[dict]: | |
| result = self._rpc("tools/list", {}) | |
| return result.get("tools", []) | |
| def call_tool(self, name: str, arguments: dict) -> dict: | |
| return self._rpc("tools/call", {"name": name, "arguments": arguments}) | |
| # ---------------------------------------------------------------- manager | |
| def _slug(text: str) -> str: | |
| return re.sub(r"[^a-z0-9_]+", "_", text.lower()).strip("_")[:24] | |
| class ToolServer: | |
| name: str | |
| url: str | |
| token: str = "" | |
| enabled: bool = True | |
| # čárkami oddělený seznam povolených nástrojů; prázdné = všechny | |
| allowed_tools: str = "" | |
| timeout: int = DEFAULT_TIMEOUT | |
| # runtime stav (persistuje se jen konfigurace + poslední známé nástroje) | |
| tools: list = field(default_factory=list) | |
| status: str = "new" # new | ok | error | |
| error: str = "" | |
| server_info: dict = field(default_factory=dict) | |
| refreshed_at: float = 0.0 | |
| class ToolServerManager: | |
| """Thread-safe registr externích MCP serverů s persistencí.""" | |
| def __init__(self, path: Path | None = None): | |
| self._lock = RLock() | |
| self.path = path or (_persist_dir() / "toolservers.json") | |
| self._servers: dict[str, ToolServer] = {} | |
| self._clients: dict[str, MCPClient] = {} | |
| self._routing: dict[str, tuple[str, str]] = {} # ext_name -> (server, tool) | |
| self._load() | |
| # ------------------------------------------------------------- CRUD | |
| def add_or_update(self, name: str, url: str, token: str = "", | |
| enabled: bool = True, allowed_tools: str = "", | |
| timeout: int = DEFAULT_TIMEOUT, | |
| refresh: bool = True) -> ToolServer: | |
| name = _slug(name) | |
| if not name: | |
| raise ValueError("Neplatné jméno serveru") | |
| if not url.startswith(("http://", "https://")): | |
| raise ValueError("URL musí začínat http(s)://") | |
| with self._lock: | |
| existing = self._servers.get(name) | |
| if token == _MASK and existing: | |
| token = existing.token | |
| server = ToolServer(name=name, url=url.strip(), token=token, | |
| enabled=bool(enabled), | |
| allowed_tools=allowed_tools.strip(), | |
| timeout=max(5, int(timeout or DEFAULT_TIMEOUT))) | |
| if existing: | |
| server.tools = existing.tools | |
| server.status = existing.status | |
| self._servers[name] = server | |
| self._clients.pop(name, None) | |
| self._persist() | |
| if refresh: | |
| self.refresh(name) | |
| return self._servers[name] | |
| def remove(self, name: str) -> bool: | |
| with self._lock: | |
| existed = self._servers.pop(name, None) is not None | |
| self._clients.pop(name, None) | |
| if existed: | |
| self._rebuild_routing() | |
| self._persist() | |
| return existed | |
| # ------------------------------------------------------------- discovery | |
| def refresh(self, name: str | None = None) -> dict: | |
| """Znovu načte nástroje serveru (None = všech zapnutých).""" | |
| names = [name] if name else [s.name for s in self._servers.values() | |
| if s.enabled] | |
| results = {} | |
| for n in names: | |
| server = self._servers.get(n) | |
| if server is None: | |
| results[n] = "neznámý server" | |
| continue | |
| try: | |
| client = MCPClient(server.url, server.token, server.timeout) | |
| client.initialize() | |
| tools = client.list_tools() | |
| with self._lock: | |
| server.tools = tools | |
| server.status = "ok" | |
| server.error = "" | |
| server.server_info = client.server_info | |
| server.refreshed_at = time.time() | |
| self._clients[n] = client | |
| results[n] = f"ok ({len(tools)} nástrojů)" | |
| logger.info("Tool server %s: %s nástrojů", n, len(tools)) | |
| except (httpx.HTTPError, MCPError, ValueError) as e: | |
| with self._lock: | |
| server.status = "error" | |
| server.error = str(e)[:300] | |
| results[n] = f"error: {e}" | |
| logger.warning("Tool server %s selhal: %s", n, e) | |
| with self._lock: | |
| self._rebuild_routing() | |
| self._persist() | |
| return results | |
| def _allowed(self, server: ToolServer, tool_name: str) -> bool: | |
| if not server.allowed_tools: | |
| return True | |
| allowed = {t.strip() for t in server.allowed_tools.split(",") if t.strip()} | |
| return tool_name in allowed | |
| def _rebuild_routing(self): | |
| self._routing = {} | |
| for server in self._servers.values(): | |
| for tool in server.tools: | |
| ext = f"ext_{server.name}_{_slug(tool.get('name', ''))}"[:64] | |
| self._routing[ext] = (server.name, tool.get("name", "")) | |
| # ------------------------------------------------------------- agent API | |
| def openai_tools(self) -> list[dict]: | |
| """OpenAI function-calling definice všech zapnutých externích nástrojů.""" | |
| out = [] | |
| with self._lock: | |
| for server in self._servers.values(): | |
| if not server.enabled or server.status != "ok": | |
| continue | |
| for tool in server.tools: | |
| tname = tool.get("name", "") | |
| if not tname or not self._allowed(server, tname): | |
| continue | |
| ext = f"ext_{server.name}_{_slug(tname)}"[:64] | |
| schema = tool.get("inputSchema") or {"type": "object", "properties": {}} | |
| desc = (f"[externí nástroj · server {server.name}] " | |
| f"{tool.get('description', '')}")[:1024] | |
| out.append({"type": "function", "function": { | |
| "name": ext, "description": desc, "parameters": schema}}) | |
| return out | |
| def call(self, ext_name: str, arguments: dict) -> dict: | |
| """Zavolá externí nástroj. Výsledek je NEDŮVĚRYHODNÝ obsah.""" | |
| with self._lock: | |
| route = self._routing.get(ext_name) | |
| if route is None: | |
| return {"error": f"Neznámý externí nástroj {ext_name}"} | |
| server_name, tool_name = route | |
| server = self._servers.get(server_name) | |
| if server is None or not server.enabled: | |
| return {"error": f"Server {server_name} není dostupný/zapnutý"} | |
| if not self._allowed(server, tool_name): | |
| return {"error": f"Nástroj {tool_name} není na serveru povolen"} | |
| client = self._clients.get(server_name) | |
| if client is None: | |
| client = MCPClient(server.url, server.token, server.timeout) | |
| try: | |
| client.initialize() | |
| except (httpx.HTTPError, MCPError) as e: | |
| return {"error": f"Připojení k {server_name} selhalo: {e}"} | |
| self._clients[server_name] = client | |
| try: | |
| result = client.call_tool(tool_name, arguments or {}) | |
| except (httpx.HTTPError, MCPError) as first_error: | |
| # session mohla vypršet — jedna reinicializace a retry | |
| try: | |
| client = MCPClient(server.url, server.token, server.timeout) | |
| client.initialize() | |
| self._clients[server_name] = client | |
| result = client.call_tool(tool_name, arguments or {}) | |
| except (httpx.HTTPError, MCPError): | |
| return {"error": f"Volání {tool_name}@{server_name} selhalo: " | |
| f"{first_error}"} | |
| return self._normalize_result(result) | |
| def _normalize_result(result: dict) -> dict: | |
| texts = [] | |
| for block in result.get("content") or []: | |
| if isinstance(block, dict) and block.get("type") == "text": | |
| texts.append(block.get("text", "")) | |
| out: dict = {} | |
| if texts: | |
| out["result"] = "\n".join(texts) | |
| if result.get("structuredContent") is not None: | |
| out["structured"] = result["structuredContent"] | |
| if result.get("isError"): | |
| return {"error": out.get("result") or "nástroj vrátil chybu"} | |
| return out or {"result": "(prázdný výsledek)"} | |
| # ------------------------------------------------------------- introspekce | |
| def list(self, mask_secrets: bool = True) -> list[dict]: | |
| with self._lock: | |
| out = [] | |
| for server in self._servers.values(): | |
| data = asdict(server) | |
| if mask_secrets and data.get("token"): | |
| data["token"] = _MASK | |
| data["tool_names"] = [t.get("name") for t in server.tools] | |
| data["tools"] = len(server.tools) | |
| out.append(data) | |
| return out | |
| def summary(self) -> dict: | |
| with self._lock: | |
| enabled = [s for s in self._servers.values() if s.enabled] | |
| return { | |
| "servers": len(self._servers), | |
| "enabled": len(enabled), | |
| "tools": sum(len(s.tools) for s in enabled if s.status == "ok"), | |
| } | |
| # ------------------------------------------------------------- persistence | |
| def _persist(self): | |
| try: | |
| data = [] | |
| for server in self._servers.values(): | |
| item = asdict(server) | |
| item.pop("server_info", None) | |
| data.append(item) | |
| self.path.parent.mkdir(parents=True, exist_ok=True) | |
| fd, tmp = tempfile.mkstemp(dir=str(self.path.parent), suffix=".tmp") | |
| with os.fdopen(fd, "w", encoding="utf-8") as f: | |
| f.write(json.dumps(data, ensure_ascii=False, indent=2)) | |
| os.replace(tmp, self.path) | |
| except OSError as e: | |
| logger.warning("Persistence tool serverů selhala: %s", e) | |
| def _load(self): | |
| if not self.path.exists(): | |
| return | |
| try: | |
| data = json.loads(self.path.read_text(encoding="utf-8")) | |
| except (OSError, json.JSONDecodeError) as e: | |
| logger.warning("Načtení %s selhalo: %s", self.path, e) | |
| return | |
| for item in data: | |
| try: | |
| server = ToolServer( | |
| name=item["name"], url=item["url"], | |
| token=item.get("token", ""), | |
| enabled=bool(item.get("enabled", True)), | |
| allowed_tools=item.get("allowed_tools", ""), | |
| timeout=int(item.get("timeout", DEFAULT_TIMEOUT))) | |
| server.tools = item.get("tools", []) | |
| server.status = item.get("status", "new") | |
| self._servers[server.name] = server | |
| except (KeyError, TypeError, ValueError): | |
| continue | |
| self._rebuild_routing() | |
| logger.info("Načteno %s tool serverů z %s", len(self._servers), self.path) | |
| # ---------------------------------------------------------------- registr | |
| def parse_registry_payload(data: dict) -> list[dict]: | |
| """Z odpovědi oficiálního registru vytáhne servery s HTTP endpointem.""" | |
| out = [] | |
| for item in data.get("servers", []): | |
| server = item.get("server", item) | |
| remotes = [r for r in server.get("remotes", []) | |
| if r.get("type") in ("streamable-http", "http") and r.get("url")] | |
| out.append({ | |
| "name": server.get("name", "?"), | |
| "description": (server.get("description") or "")[:200], | |
| "version": server.get("version", ""), | |
| "urls": [r["url"] for r in remotes], | |
| "remote": bool(remotes), | |
| }) | |
| return out | |
| def search_registry(query: str, limit: int = 12, | |
| remote_only: bool = True) -> list[dict]: | |
| """Prohledá oficiální MCP registr (registry.modelcontextprotocol.io).""" | |
| resp = httpx.get(REGISTRY_URL, params={"search": query, "limit": max(limit, 30)}, | |
| timeout=20) | |
| resp.raise_for_status() | |
| results = parse_registry_payload(resp.json()) | |
| if remote_only: | |
| results = [r for r in results if r["remote"]] | |
| return results[:limit] | |