File size: 19,421 Bytes
c1a62a8
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
"""
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]


@dataclass
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)

    @staticmethod
    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]