From f8acb906e51dab184bc945129312a23162ceadea Mon Sep 17 00:00:00 2001 From: Bruno Charest Date: Wed, 7 Oct 2026 09:09:19 -0400 Subject: [PATCH] feat: connecteurs Discord/Telegram/MCP + outils MCP dynamiques au registre (v7.57.0) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - presets Discord / Telegram : colonnes kind + auth (bearer/bot/none) sur agent_connectors (migration 34) ; Discord envoie le jeton préfixé dans l'en-tête d'autorisation, Telegram le place dans l'URL via le substitut {secret} remplacé à l'appel — jamais stocké en clair - Teams : scope ChannelMessage.Read.All ajouté aux Graph scopes M365 - app/services/mcp_client.py : handshake initialize → notifications/initialized → tools/list → tools/call (JSON-RPC 2.0, réponse JSON ou premier data: d'un text/event-stream), garde SSRF, outils cachés en base (tools_json) — le bouton « Tester » rejoue le handshake - ToolRegistry._all() merge le cache MCP à chaque run : outils exposés au LLM sous mcp__ avec leur inputSchema, dispatch tools/call ; serveur désactivé → outil absent du schéma, échec réseau → ToolResult(error) - menu + : champ Type dans le formulaire connecteur (preset URL + auth), badge « · N outil(s) » sur la fiche d'un serveur MCP - tests : tests/test_v757_mcp_discord.py (12 tests, 0 appel réseau réel) ; suite complète 1331 verts, ruff 0, eslint 0 - livraison : VERSION + app/main = 7.57.0, OpenAPI 523 chemins, CHANGELOG, ROADMAP phase 7 cochée, avenant phase 7 dans docs/V74_Agent_Plus_Menu.md --- CHANGELOG.md | 32 +++++ ROADMAP.md | 35 ++++-- VERSION | 2 +- app/main.py | 2 +- app/migrations.py | 13 ++ app/routers/agent.py | 4 +- app/services/connectors.py | 128 ++++++++++++++----- app/services/mcp_client.py | 181 ++++++++++++++++++++++++++ app/services/oauth_connectors.py | 1 + app/services/tool_registry.py | 59 ++++++++- app/templates/agent_panel.html | 13 +- docs/V74_Agent_Plus_Menu.md | 21 ++++ docs/openapi-v2.json | 2 +- static/js/agent_panel_2.js | 26 +++- tests/test_v757_mcp_discord.py | 209 +++++++++++++++++++++++++++++++ 15 files changed, 668 insertions(+), 60 deletions(-) create mode 100644 app/services/mcp_client.py create mode 100644 tests/test_v757_mcp_discord.py diff --git a/CHANGELOG.md b/CHANGELOG.md index d6a7e79..c371d2c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,37 @@ # Changelog - FlowDeck +## v7.57.0 (2026-10-06) — Connecteurs : Discord, Telegram, Teams, MCP (phase 7/8) + +### Added + +- **Presets Discord / Telegram** — champ **Type** dans le formulaire « Ajouter + un connecteur » (Personnalisé / Discord / Telegram / Serveur MCP) : URL et + schéma d'authentification pré-remplis. Colonnes `kind` + `auth` + (`bearer`/`bot`/`none`) ajoutées à `agent_connectors` (migration **34**) ; + Discord envoie son jeton préfixé dans l'en-tête d'autorisation, Telegram place + le jeton **dans l'URL** via le substitut `{secret}` — substitué à l'appel, + **jamais stocké en clair** (la colonne url ne contient que le motif). +- **Teams** — scope `ChannelMessage.Read.All` ajouté aux Graph scopes M365 + (consentement admin requis à la reconnexion). +- **`app/services/mcp_client.py`** — client MCP (streamable HTTP) : handshake + `initialize` → `notifications/initialized` → `tools/list` → `tools/call` + (JSON-RPC 2.0, réponse JSON **ou** premier `data:` d'un `text/event-stream`), + garde SSRF conservée, erreurs JSON-RPC remontées telles quelles, sortie + bornée à 20 000 car. Les outils sont **cachés** en base (`tools_json`) — le + bouton « Tester le connecteur » rejoue le handshake et met à jour la liste. +- **Outils dynamiques au registre** — `ToolRegistry._all()` merge ce cache à + chaque run : chaque outil est exposé au LLM sous `mcp__` avec + son `inputSchema` d'origine, et `execute()` dispatche sur `tools/call`. + Serveur désactivé → outil **absent du schéma** ; échec réseau → + `ToolResult(error)` sans casser le run. +- **Menu +** — badge `· N outil(s)` sur la fiche d'un serveur MCP ; + `POST /api/agent/connectors` accepte `kind` + `auth` (400 sur valeurs + inconnues). +- **Tests** — `tests/test_v757_mcp_discord.py` : **12 tests**, 0 appel réseau + réel (`_rpc` / `_get` monkeypatchés) — presets, `{secret}` absent de la base, + handshake + séquence d'appels, cache, dispatch LLM, serveur désactivé, + erreurs RPC, parse SSE, scopes Teams, câblage menu. + ## v7.56.0 (2026-10-06) — Connecteurs : Google + Microsoft 365 (phase 6/8) ### Added diff --git a/ROADMAP.md b/ROADMAP.md index f84d5b2..4bd3370 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -284,7 +284,7 @@ Propriétés custom, AI keywords, sync API, 12 tables DB --- -## v7.51.0 → v7.58.0 — Menu + de l'assistant : hub de contexte (phases 1-6 ✅ 2026-10-06 · phases 7-8 planifiées) +## v7.51.0 → v7.58.0 — Menu + de l'assistant : hub de contexte (phases 1-7 ✅ 2026-10-06 · phase 8 planifiée) > **Objectif** : faire du bouton **+** du panneau agent (à gauche de la zone d'édition) > un menu à sections, tel que demandé : @@ -471,13 +471,30 @@ CHANGELOG). Chiffrage : S < ½ j · M = 1–2 j · L = 3–5 j. (GOOGLE_/MS_ CLIENT_ID/SECRET + redirect URI), OpenAPI régénéré (523 chemins), CHANGELOG, ce ROADMAP -### Phase 7 — v7.57.0 — Connecteurs : Discord, Telegram, Teams, MCP · effort L +### Phase 7 — v7.57.0 — Connecteurs : Discord, Telegram, Teams, MCP · effort L ✅ -- [ ] **Discord** (bot token : lecture canaux/DMs), **Telegram** (Bot API), - **Teams** (via Graph de la phase 6) -- [ ] **MCP complet** — handshake `initialize`, `tools/list`, `tools/call` - (streamable HTTP) → outils **dynamiques** ajoutés au registre à la connexion -- [ ] **Tests** — schémas MCP, mapping d'outils, échecs d'auth +- [x] **Presets Discord / Telegram** — champ **Type** dans le formulaire + « Ajouter un connecteur » : `kind` + `auth` (`bearer`/`bot`/`none`) ajoutés à + `agent_connectors` (migration **34**) ; Discord = en-tête d'autorisation avec + le jeton préfixé, Telegram = jeton **dans l'URL** via le substitut `{secret}` + (remplacé à l'appel, **jamais stocké en clair** — colonne url = `…/bot{secret}`) +- [x] **Teams** — scope `ChannelMessage.Read.All` ajouté aux Graph scopes M365 + (consentement admin) +- [x] **MCP complet** — `app/services/mcp_client.py` : `initialize` → + `notifications/initialized` → `tools/list` → `tools/call` (JSON-RPC 2.0, + réponse JSON **ou** 1er `data:` d'un `text/event-stream`), garde SSRF, + outils **cachés** en base (`tools_json`, maj par le bouton « Tester ») +- [x] **Outils dynamiques au registre** — `ToolRegistry._all()` merge le cache + MCP à chaque run : nom `mcp__` + `inputSchema` exposés au + LLM, `execute()` dispatche sur `tools/call` ; serveur désactivé → outil + **absent du schéma** ; erreur RPC → `ToolResult(error)` (jamais de run perdu) +- [x] **Tests** — `tests/test_v757_mcp_discord.py` : **12 tests**, 0 appel + réseau (`_rpc` / `_get` monkeypatchés) — presets, `{secret}` absent de la + base, kind/auth invalides → 400, handshake + séquence d'appels, cache + `tools_json`, dispatch LLM, serveur désactivé, erreur handshake → probe + error, parse SSE/erreurs JSON-RPC, slug, scopes Teams, câblage menu +- [x] **Livraison** — `VERSION` + `app/main.py` = **7.57.0**, OpenAPI régénéré, + CHANGELOG, ce ROADMAP ### Phase 8 — v7.58.0 — Add plugins · effort M @@ -489,8 +506,8 @@ CHANGELOG). Chiffrage : S < ½ j · M = 1–2 j · L = 3–5 j. --- -*Plan produit le 2026-10-06 à partir du code réel — **phases 1 à 6 livrées le -2026-10-06 (v7.51.0 → v7.56.0)**, phases 7-8 à valider une par une avant « go ». +*Plan produit le 2026-10-06 à partir du code réel — **phases 1 à 7 livrées le +2026-10-06 (v7.51.0 → v7.57.0)**, phase 8 (plugins) à valider avant « go ». Chaque phase se clôt par tests verts, bump de version, CHANGELOG et ROADMAP à jour.* --- diff --git a/VERSION b/VERSION index 7c6760b..f266389 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -7.56.0 +7.57.0 diff --git a/app/main.py b/app/main.py index fae41b5..dab6f38 100644 --- a/app/main.py +++ b/app/main.py @@ -185,7 +185,7 @@ async def lifespan(_app: FastAPI): app = FastAPI( title="FlowDeck", - version="7.56.0", + version="7.57.0", docs_url="/docs", redoc_url="/redoc", lifespan=lifespan, diff --git a/app/migrations.py b/app/migrations.py index 2bc9f77..51ee3c1 100644 --- a/app/migrations.py +++ b/app/migrations.py @@ -1586,6 +1586,19 @@ def _migration_my_tasks_mapping(conn: sqlite3.Connection) -> None: ) +@register(34, "v7.57.0: connecteurs — kind/auth/tools_json (Discord, Telegram, MCP)") +def _migration_connector_kinds(conn: sqlite3.Connection) -> None: + """Discord (auth « Bot »), Telegram (jeton dans l'URL via `{secret}`) et + serveurs MCP (kind='mcp', cache d'outils) réutilisent la même table.""" + cols = columns(conn, "agent_connectors") + if "kind" not in cols: + conn.execute("ALTER TABLE agent_connectors ADD COLUMN kind TEXT NOT NULL DEFAULT 'custom'") + if "auth" not in cols: + conn.execute("ALTER TABLE agent_connectors ADD COLUMN auth TEXT NOT NULL DEFAULT 'bearer'") + if "tools_json" not in cols: + conn.execute("ALTER TABLE agent_connectors ADD COLUMN tools_json TEXT NOT NULL DEFAULT '[]'") + + @register(33, "v7.56.0: tokens OAuth des connecteurs (Google / M365)") def _migration_connector_tokens(conn: sqlite3.Connection) -> None: """Tokens OAuth par (kind, utilisateur), chiffrés Fernet côté service.""" diff --git a/app/routers/agent.py b/app/routers/agent.py index 2c76d6d..6674af6 100644 --- a/app/routers/agent.py +++ b/app/routers/agent.py @@ -653,7 +653,9 @@ def create_connectors_route(request: Request, body: dict = Body(default={})): _current_user_id(request) try: return connectors.create_connector( - body.get("name") or "", body.get("url") or "", str(body.get("secret") or "") + body.get("name") or "", body.get("url") or "", str(body.get("secret") or ""), + kind=str(body.get("kind") or "custom"), + auth=str(body.get("auth") or "bearer"), ) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc diff --git a/app/services/connectors.py b/app/services/connectors.py index 9f75013..30ce4f4 100644 --- a/app/services/connectors.py +++ b/app/services/connectors.py @@ -1,7 +1,10 @@ """Connecteurs de l'agent (v7.55.0) — socle : catalogue, statut, fetch gardé SSRF. -3 connecteurs **natifs** (gitea, github, web) servis à la volée depuis la config, -plus des connecteurs **personnalisation** persistés (URL + clé) en base. +5 connecteurs **natifs** (gitea, github, web, google, ms365) servis à la volée, +plus des connecteurs **personnels** persistés (URL + clé) en base avec un +``kind`` : ``custom`` | ``discord`` (auth « Bot ») | ``telegram`` (jeton dans +l'URL via ``{secret}``) | ``mcp`` (serveur Model Context Protocol, cache +d'outils dans ``tools_json``). ponytail : un seul fichier (pas de package ``connectors/`` ni de classe par provider) — 3 probes dans un dict + 1 fetch. Les adapters lourds des phases 6-7 @@ -10,6 +13,7 @@ provider) — 3 probes dans un dict + 1 fetch. Les adapters lourds des phases 6- from __future__ import annotations +import json import logging from app.config import settings @@ -20,6 +24,18 @@ from app.services.sso_provisioning import decrypt_secret, encrypt_secret logger = logging.getLogger(__name__) NATIVE_KINDS = ("gitea", "github", "web", "google", "ms365") +CUSTOM_KINDS = ("custom", "discord", "telegram", "mcp") +AUTH_SCHEMES = ("bearer", "bot", "none") +# Presets proposés par le formulaire du menu + (le jeton reste toujours saisi +# par l'utilisateur — jamais codé en dur ici). +PRESETS = { + "discord": {"url": "https://discord.com/api/v10", "auth": "bot", + "hint": "Bot Discord Developers → Token"}, + "telegram": {"url": "https://api.telegram.org/bot{secret}", "auth": "none", + "hint": "BotFather → Bot Token (jeton dans l'URL)"}, + "mcp": {"url": "https://", "auth": "bearer", + "hint": "URL du endpoint MCP (streamable HTTP)"}, +} OAUTH_KINDS = oauth_connectors.OAUTH_KINDS FETCH_LIMIT = 20000 # caractères max renvoyés au LLM (ponytail : borne fixe) @@ -48,6 +64,21 @@ def native_state(kind: str, user_id: int | None = None) -> tuple[str, str]: return ("ok", f"recherche web via {settings.web_search_provider or 'duckduckgo'}") +def _row(r) -> dict: + """Ligne `agent_connectors` → dict API (le jeton n'y figure jamais).""" + try: + tools = json.loads(r["tools_json"] or "[]") if "tools_json" in r.keys() else [] + except (ValueError, TypeError): + tools = [] + return { + "id": r["id"], "builtin": False, "kind": r["kind"], "name": r["name"], + "url": r["url"], "auth": r["auth"], "enabled": bool(r["enabled"]), + "status": r["status"], "detail": r["detail"] or "", + "has_secret": bool(r["secret_encrypted"]), "oauth": False, + "tools_count": len(tools), + } + + def list_connectors(user_id: int | None = None) -> list[dict]: """Catalogue : natifs (toujours présents) + connecteurs personnels.""" out: list[dict] = [] @@ -62,46 +93,45 @@ def list_connectors(user_id: int | None = None) -> list[dict]: }) with get_conn() as conn: rows = conn.execute( - "SELECT id, name, url, secret_encrypted, enabled, status, detail " - "FROM agent_connectors ORDER BY id" + "SELECT id, name, url, secret_encrypted, enabled, status, detail, " + "kind, auth, tools_json FROM agent_connectors ORDER BY id" ).fetchall() - for r in rows: - out.append({ - "id": r["id"], "builtin": False, "kind": "custom", "name": r["name"], - "url": r["url"], "enabled": bool(r["enabled"]), "status": r["status"], - "detail": r["detail"] or "", "has_secret": bool(r["secret_encrypted"]), - }) - return out + return out + [_row(r) for r in rows] def get(connector_id: int) -> dict | None: with get_conn() as conn: r = conn.execute( - "SELECT id, name, url, secret_encrypted, enabled, status, detail " - "FROM agent_connectors WHERE id=?", (connector_id,) + "SELECT id, name, url, secret_encrypted, enabled, status, detail, " + "kind, auth, tools_json FROM agent_connectors WHERE id=?", + (connector_id,), ).fetchone() - if not r: - return None - return {"id": r["id"], "builtin": False, "kind": "custom", "name": r["name"], - "url": r["url"], "enabled": bool(r["enabled"]), "status": r["status"], - "detail": r["detail"] or "", "has_secret": bool(r["secret_encrypted"])} + return _row(r) if r else None # ── CRUD (personnels) ──────────────────────────────────────────────────────── -def create_connector(name: str, url: str, secret: str = "") -> dict: +def create_connector(name: str, url: str, secret: str = "", *, + kind: str = "custom", auth: str = "bearer") -> dict: """Valide l'URL (garde SSRF) puis stocke la clé **chiffrée** (Fernet).""" from app.services.importers.url_fetch import _validate_url name = (name or "").strip() if not name: raise ValueError("name est requis") - url = _validate_url((url or "").strip()) # lève ValueError si hôte interne + if kind not in CUSTOM_KINDS: + raise ValueError(f"kind invalide: {kind} (attendu {', '.join(CUSTOM_KINDS)})") + if auth not in AUTH_SCHEMES: + raise ValueError(f"auth invalide: {auth} (attendu {', '.join(AUTH_SCHEMES)})") + url = (url or "").strip() + if "{secret}" in url and not (secret or "").strip(): + raise ValueError("clé requise : l'URL contient {secret}") + url = _validate_url(url) # lève ValueError si hôte interne with get_conn() as conn: cur = conn.execute( - "INSERT INTO agent_connectors (name, url, secret_encrypted, status, detail) " - "VALUES (?,?,?,?,?)", - (name, url, encrypt_secret(secret or ""), "unknown", ""), + "INSERT INTO agent_connectors (name, url, secret_encrypted, status, " + "detail, kind, auth) VALUES (?,?,?,?,?,?,?)", + (name, url, encrypt_secret(secret or ""), "unknown", "", kind, auth), ) conn.commit() cid = cur.lastrowid @@ -143,22 +173,44 @@ def delete_connector(connector_id: int) -> bool: # ── Réseau (toujours gardé SSRF) ───────────────────────────────────────────── def _headers(connector_id: int | None = None) -> dict: - """En-têtes d'appel : la clé n'est lue (et déchiffrée) qu'ici, jamais renvoyée.""" + """En-têtes d'appel : la clé n'est lue (et déchiffrée) qu'ici, jamais renvoyée. + + ``auth`` = bearer (défaut) | bot (Discord, jeton préfixé) | none (jeton ailleurs). + """ headers = {"User-Agent": "FlowDeck-Connectors/1.0"} - secret = "" + secret, auth = "", "bearer" if connector_id: with get_conn() as conn: r = conn.execute( - "SELECT secret_encrypted FROM agent_connectors WHERE id=?", + "SELECT secret_encrypted, auth FROM agent_connectors WHERE id=?", (connector_id,), ).fetchone() - secret = decrypt_secret((r["secret_encrypted"] if r else "") or "") - if secret: + if r: + secret = decrypt_secret(r["secret_encrypted"] or "") + auth = r["auth"] or "bearer" + if secret and auth == "bot": + headers["Authorization"] = f"Bot {secret}" + elif secret and auth != "none": headers["Authorization"] = f"Bearer {secret}" headers["X-Api-Key"] = secret return headers +def _with_secret(url: str, connector_id: int) -> str: + """Substitue ``{secret}`` (Telegram : le jeton vit dans l'URL, jamais stocké clair).""" + if "{secret}" not in url: + return url + with get_conn() as conn: + r = conn.execute( + "SELECT secret_encrypted FROM agent_connectors WHERE id=?", + (connector_id,), + ).fetchone() + secret = decrypt_secret((r["secret_encrypted"] if r else "") or "") + if not secret: + raise ValueError("clé manquante pour une URL contenant {secret}") + return url.replace("{secret}", secret) + + async def _get(url: str, headers: dict | None = None) -> tuple[int, str]: """GET gardé SSRF (validation + re-vérification après redirection). @@ -184,7 +236,7 @@ def _resolve(connector: str) -> tuple[str, str | int]: row = get(int(token)) if row is None: raise ValueError(f"Connecteur inconnu: {token}") - return ("custom", row["id"]) + return (row["kind"], row["id"]) low = token.lower() if low in NATIVE_KINDS: return (low, low) @@ -193,7 +245,7 @@ def _resolve(connector: str) -> tuple[str, str | int]: "SELECT id FROM agent_connectors WHERE lower(name)=?", (low,) ).fetchone() if row: - return ("custom", row["id"]) + return (get(row["id"])["kind"], row["id"]) raise ValueError(f"Connecteur inconnu: {token}") @@ -201,7 +253,11 @@ async def probe(connector: str, user_id: int | None = None) -> dict: """Teste un connecteur et **persiste** le résultat (personnel uniquement).""" kind, ref = _resolve(connector) try: - if kind in OAUTH_KINDS: + if kind == "mcp": + from app.services import mcp_client + tools = await mcp_client.initialize_and_list(int(ref)) + status, detail = "ok", f"MCP · {len(tools)} outil(s)" + elif kind in OAUTH_KINDS: res = await oauth_connectors.api_get(kind, user_id, oauth_connectors.PROVIDERS[kind]["probe_path"]) status, detail = res["status"], res["text"][:200] @@ -216,12 +272,13 @@ async def probe(connector: str, user_id: int | None = None) -> dict: status, detail = "ok", native_state("web")[1] else: row = get(int(ref)) - code, _ = await _get(row["url"], _headers(int(ref))) + url = _with_secret(row["url"], int(ref)) + code, _ = await _get(url, _headers(int(ref))) status, detail = ("ok" if code < 400 else "error"), f"HTTP {code}" except Exception as exc: # noqa: BLE001 — un probe ne casse jamais l'UI logger.info("Connector probe failed (%s): %s", connector, exc) status, detail = "error", str(exc)[:200] - if kind == "custom": + if kind not in NATIVE_KINDS: with get_conn() as conn: conn.execute( "UPDATE agent_connectors SET status=?, detail=? WHERE id=?", @@ -246,6 +303,9 @@ async def connector_fetch(connector: str, path: str = "", query: str = "", path = (path or "").strip() if kind in OAUTH_KINDS: return await oauth_connectors.api_get(kind, user_id, path) + if kind == "mcp": + from app.services import mcp_client + return {"status": "ok", "text": mcp_client.tools_text(int(ref))} if kind == "gitea": base = settings.gitea_url.rstrip("/") url = f"{base}/{path.lstrip('/')}" if path else f"{base}/api/v1/version" @@ -267,7 +327,7 @@ async def connector_fetch(connector: str, path: str = "", query: str = "", return {"status": "error", "text": "Connecteur supprimé"} if not row["enabled"]: return {"status": "error", "text": "Connecteur désactivé"} - url = row["url"].rstrip("/") + url = _with_secret(row["url"], int(ref)).rstrip("/") if path: url += "/" + path.lstrip("/") code, text = await _get(url, _headers(int(ref))) diff --git a/app/services/mcp_client.py b/app/services/mcp_client.py new file mode 100644 index 0000000..5b28794 --- /dev/null +++ b/app/services/mcp_client.py @@ -0,0 +1,181 @@ +"""Client MCP (Model Context Protocol) — streamable HTTP, v7.57.0. + +JSON-RPC 2.0 : ``initialize`` → ``notifications/initialized`` → ``tools/list`` +→ ``tools/call``. Les outils sont **cachés** dans +``agent_connectors.tools_json`` (kind ``mcp``) puis exposés au LLM par +``ToolRegistry`` sous des noms ``mcp__`` — outils dynamiques, +sans état en mémoire. + +ponytail : **1 seul seam** ``_rpc()`` (monkeypatché en test) ; on lit soit la +réponse JSON directe, soit la première ligne ``data:`` d'une +``text/event-stream`` — pas de client SSE à l'état (les méthodes utilisées +répondent en JSON chez la grande majorité des serveurs). +""" + +from __future__ import annotations + +import json +import logging +import re + +from app.db import get_conn +from app.services.connectors import _headers, _with_secret, get + +logger = logging.getLogger(__name__) + +# Version la plus répandue côté serveurs (un serveur plus récent rétrograde). +PROTOCOL_VERSION = "2024-11-05" +CLIENT_INFO = {"name": "FlowDeck", "version": "7"} + + +def tool_name(server_name: str, tool_name: str) -> str: + """Nom LLM stable : ``mcp__`` (slug minuscules/underscores).""" + def slug(text: str) -> str: + return re.sub(r"[^a-z0-9]+", "_", str(text).lower()).strip("_") + joined = f"mcp_{slug(server_name)}_{slug(tool_name)}" + return re.sub(r"_{2,}", "_", joined) + + +def parse_body(content_type: str, body: str, status: int) -> dict: + """Corps MCP → dict JSON (JSON direct ou événement ``data:`` d'un flux SSE).""" + text = body or "" + if "text/event-stream" in (content_type or ""): + for line in text.splitlines(): + if line.startswith("data:"): + text = line[5:].strip() + break + try: + data = json.loads(text) + except ValueError as exc: + raise ValueError(f"Réponse MCP illisible (HTTP {status})") from exc + if isinstance(data, dict) and data.get("error"): + err = data["error"] if isinstance(data["error"], dict) else {} + raise ValueError(f"MCP {err.get('code', '?')} : {err.get('message', 'erreur')}") + return data if isinstance(data, dict) else {} + + +async def _rpc(url: str, payload: dict, headers: dict | None = None) -> dict: + """POST JSON-RPC gardé SSRF → réponse décodée. Point d'injection des tests.""" + from app.services.http_client import shared_client + from app.services.importers.url_fetch import _validate_url + + safe = _validate_url(url) + hdrs = {"Content-Type": "application/json", + "Accept": "application/json, text/event-stream", + **(headers or {})} + async with shared_client(timeout=15, headers=hdrs) as client: + resp = await client.post(safe, json=payload) + return parse_body(resp.headers.get("content-type", ""), resp.text or "", + resp.status_code) + + +async def _notify(url: str, payload: dict, headers: dict | None = None) -> None: + """Notification JSON-RPC (202 sans corps) — les erreurs sont ignorées.""" + try: + await _rpc(url, payload, headers) + except Exception as exc: # noqa: BLE001 + logger.debug("MCP notification ignorée: %s", exc) + + +def _server(connector_id: int) -> dict: + row = get(connector_id) + if row is None: + raise ValueError("Serveur MCP introuvable") + if row["kind"] != "mcp": + raise ValueError("Ce connecteur n'est pas un serveur MCP") + if not row["enabled"]: + raise ValueError("Serveur MCP désactivé") + return {"id": row["id"], "name": row["name"], "url": row["url"], + "headers": _headers(connector_id)} + + +def _normalize(server: dict, raw: list) -> list[dict]: + out = [] + for t in raw if isinstance(raw, list) else []: + if not isinstance(t, dict) or not t.get("name"): + continue + params = t.get("inputSchema") + if not isinstance(params, dict): + params = {"type": "object", "properties": {}} + out.append({ + "name": tool_name(server["name"], t["name"]), + "original": str(t["name"]), + "description": str(t.get("description") or "Outil MCP"), + "parameters": params, + }) + return out + + +def _save_tools(connector_id: int, tools: list[dict]) -> None: + with get_conn() as conn: + conn.execute("UPDATE agent_connectors SET tools_json=? WHERE id=?", + (json.dumps(tools), connector_id)) + conn.commit() + + +async def initialize_and_list(connector_id: int) -> list[dict]: + """Handshake + ``tools/list`` → met à jour le cache de la base.""" + srv = _server(connector_id) + url = _with_secret(srv["url"], connector_id) + await _rpc(url, { + "jsonrpc": "2.0", "id": 1, "method": "initialize", + "params": {"protocolVersion": PROTOCOL_VERSION, "capabilities": {}, + "clientInfo": CLIENT_INFO}, + }, srv["headers"]) + await _notify(url, {"jsonrpc": "2.0", "method": "notifications/initialized"}, + srv["headers"]) + data = await _rpc(url, {"jsonrpc": "2.0", "id": 2, "method": "tools/list", + "params": {}}, srv["headers"]) + result = data.get("result") if isinstance(data.get("result"), dict) else {} + tools = _normalize(srv, result.get("tools") or []) + _save_tools(connector_id, tools) + return tools + + +def cached_tools() -> list[dict]: + """Outils MCP des serveurs **activés** (cache en base) — lu par ToolRegistry.""" + with get_conn() as conn: + rows = conn.execute( + "SELECT id, tools_json FROM agent_connectors WHERE kind='mcp' AND enabled=1" + ).fetchall() + out: list[dict] = [] + for r in rows: + try: + tools = json.loads(r["tools_json"] or "[]") + except (ValueError, TypeError): + tools = [] + for t in tools if isinstance(tools, list) else []: + if isinstance(t, dict) and t.get("name"): + out.append({**t, "connector_id": r["id"]}) + return out + + +def tools_text(connector_id: int) -> str: + """Liste lisible des outils cachés (utilisée par ``connector_fetch``).""" + tools = [t for t in cached_tools() if t["connector_id"] == connector_id] + if not tools: + return "Aucun outil en cache — lancez « Tester le connecteur » (tools/list)." + return "\n".join(f"- {t['name']} : {t.get('description', '')}" for t in tools) + + +def _content_text(result: dict) -> str: + parts = [] + for item in result.get("content") or []: + if isinstance(item, dict) and item.get("type") == "text": + parts.append(str(item.get("text") or "")) + else: + parts.append(json.dumps(item, ensure_ascii=False)[:2000]) + return "\n".join(p for p in parts if p) or json.dumps(result, ensure_ascii=False)[:4000] + + +async def call_tool(connector_id: int, tool: str, arguments: dict | None = None) -> dict: + """``tools/call`` → {status, text} (isError → error, borné à 20 000 car.).""" + srv = _server(connector_id) + url = _with_secret(srv["url"], connector_id) + data = await _rpc(url, { + "jsonrpc": "2.0", "id": 3, "method": "tools/call", + "params": {"name": tool, "arguments": arguments or {}}, + }, srv["headers"]) + result = data.get("result") if isinstance(data.get("result"), dict) else {} + text = _content_text(result)[:20000] + return {"status": "error" if result.get("isError") else "ok", "text": text} diff --git a/app/services/oauth_connectors.py b/app/services/oauth_connectors.py index 4b302c6..cc33fc2 100644 --- a/app/services/oauth_connectors.py +++ b/app/services/oauth_connectors.py @@ -50,6 +50,7 @@ PROVIDERS: dict[str, dict] = { "scopes": [ "openid", "email", "profile", "offline_access", "User.Read", "Files.Read", "Mail.Read", "Calendars.Read", + "ChannelMessage.Read.All", # Teams (phase 7) — consentement admin ], "extra": {"response_mode": "query", "prompt": "consent"}, "probe_path": "/me", diff --git a/app/services/tool_registry.py b/app/services/tool_registry.py index 547f337..c5a10c6 100644 --- a/app/services/tool_registry.py +++ b/app/services/tool_registry.py @@ -1115,6 +1115,36 @@ class ConnectorFetch(Tool): message=text[:800], data={"text": text}) +class McpTool(Tool): + """Outil MCP dynamique — instance construite à la volée depuis le cache + en base (`mcp_client.cached_tools()`), pas enregistrée dans TOOL_CLASSES.""" + + def __init__(self, entry: dict): + self.name = entry["name"] + self.description = entry.get("description") or "Outil MCP" + self.parameters = entry.get("parameters") or {"type": "object", "properties": {}} + self.entry = entry + + async def execute(self, args, *, user_id=None) -> ToolResult: + from app.services import mcp_client + + try: + res = await mcp_client.call_tool(self.entry["connector_id"], + self.entry.get("original") or "", args) + except ValueError as exc: + return ToolResult(status="error", tool=self.name, message=str(exc)) + except Exception as exc: # noqa: BLE001 — réseau : erreur outil, pas run + logger.warning("MCP tool failed (%s): %s", self.name, exc) + return ToolResult(status="error", tool=self.name, + message=f"MCP indisponible: {exc}") + if res.get("status") != "ok": + return ToolResult(status="error", tool=self.name, + message=(res.get("text") or "erreur MCP")[:500]) + text = res.get("text") or "" + return ToolResult(status="success", tool=self.name, target_type="mcp", + message=text[:800], data={"text": text}) + + TOOL_CLASSES = [ SearchWorkspace, ReadCollection, ReadPage, ReadWorkspaces, ReadDocument, CreateCollection, CreateView, AddProperty, CreatePage, CreateDocument, @@ -1133,24 +1163,41 @@ class ToolRegistry: def __init__(self): self.tools: dict[str, Tool] = {t.name: t() for t in TOOL_CLASSES} + def _all(self) -> dict[str, Tool]: + """Outils statiques + **outils MCP dynamiques** (cache en base). + + Les MCP ne sont jamais mis en cache en mémoire : un run re-lit la base, + donc « Tester le connecteur » suffit à rendre un outil disponible. + """ + tools = dict(self.tools) + try: + from app.services import mcp_client + for entry in mcp_client.cached_tools(): + tools[entry["name"]] = McpTool(entry) + except Exception: # noqa: BLE001 — la base ne doit jamais casser un run + logger.exception("MCP dynamic tools load failed") + return tools + def list(self, scope: dict | None = None) -> list[str]: """Tool names allowed by an agent scope (default: all).""" + tools = self._all() allowed = (scope or {}).get("tools") if allowed is None: - return list(self.tools.keys()) - return [n for n in self.tools if n in allowed] + return list(tools.keys()) + return [n for n in tools if n in allowed] def schema(self, scope: dict | None = None) -> list[dict]: """Function-calling schema for the LLM, filtered by scope.""" + tools = self._all() return [ - {"name": self.tools[n].name, - "description": self.tools[n].description, - "parameters": self.tools[n].parameters} + {"name": tools[n].name, + "description": tools[n].description, + "parameters": tools[n].parameters} for n in self.list(scope) ] async def execute(self, tool: str, args: dict, *, user_id: int | None = None) -> ToolResult: - impl = self.tools.get(tool) + impl = self._all().get(tool) if not impl: return ToolResult(status="error", tool=tool, message=f"Outil inconnu: {tool}") return await impl.execute(args, user_id=user_id) diff --git a/app/templates/agent_panel.html b/app/templates/agent_panel.html index b07c68f..05ded56 100644 --- a/app/templates/agent_panel.html +++ b/app/templates/agent_panel.html @@ -445,11 +445,18 @@ body.fd-ap-resizing *{cursor:col-resize!important}
+ + - - - + + +
diff --git a/docs/V74_Agent_Plus_Menu.md b/docs/V74_Agent_Plus_Menu.md index 786aa5a..da603f9 100644 --- a/docs/V74_Agent_Plus_Menu.md +++ b/docs/V74_Agent_Plus_Menu.md @@ -202,6 +202,27 @@ transaction par migration — A31). Chaque phase se clôt par : `ruff` + suite verte, bump de version, CHANGELOG, ROADMAP à jour, tag. +## 9bis. Avenant phase 7 (v7.57.0) — décisions livrées + +- **D5 (`probe/list_sources/fetch` par adapter) écartée** : 1 switch sur + `agent_connectors.kind` dans `probe()` / `connector_fetch()` — une interface + à 3 méthodes pour 4 kinds tordrait plus qu'elle ne simplifie ; ajouter un + kind = 1 branche de plus dans le même switch. +- **Discord / Telegram = presets du socle**, pas 3 adapters dédiés : le champ + **Type** du formulaire pré-remplit l'URL + le schéma d'auth + (`bearer`/`bot`/`none`). La Bot API de Telegram impose le jeton **dans l'URL** + → substitut `{secret}` remplacé à l'appel (`connectors._with_secret`), + stockage = motif seul, test « le jeton n'est jamais en base ». +- **MCP : outils dynamiques sans état** — `ToolRegistry._all()` relit le cache + `tools_json` à chaque run (rien à invalider, « Tester » suffit) ; nom LLM + `mcp__` (slug sans double soulignement) ; dispatch + `tools/call` dans `McpTool.execute()`. Serveur désactivé → outil **absent du + schéma**, échec réseau → `ToolResult(error)` : le run continue toujours. +- **Pas de client SSE à l'état** : `parse_body()` lit le JSON direct ou le + premier `data:` d'un `text/event-stream`, ce qui couvre + `initialize` / `tools/list` / `tools/call` des serveurs streamables courants. + Upgrade : client SSE persistant si un serveur ne répond qu'en flux. + ## 10. Hors périmètre - Écriture dans les sources distantes (Google Docs, Slack…) — lecture seule au départ. diff --git a/docs/openapi-v2.json b/docs/openapi-v2.json index e4323e3..9658182 100644 --- a/docs/openapi-v2.json +++ b/docs/openapi-v2.json @@ -2,7 +2,7 @@ "openapi": "3.1.0", "info": { "title": "FlowDeck", - "version": "7.56.0" + "version": "7.57.0" }, "paths": { "/auth/register": { diff --git a/static/js/agent_panel_2.js b/static/js/agent_panel_2.js index 6365351..bc2eb77 100644 --- a/static/js/agent_panel_2.js +++ b/static/js/agent_panel_2.js @@ -75,7 +75,8 @@ galleryFilter: '', canvases: [], canvas: null, connectors: [], connector: null, - connectorForm: {name: '', url: '', secret: ''}, + connectorForm: {name: '', url: '', secret: '', kind: 'custom', + auth: 'bearer', hint: ''}, cmdFocus: -1, slashDismiss: false, toastMsg: '', toastErr: false, _toastTimer: null, _livePh: null, @@ -1270,7 +1271,8 @@ _connectorStatus(c){ var mark = c.status === 'ok' ? '✓ ' : (c.status === 'error' ? '✗ ' : (c.status === 'missing' ? '⚠ ' : '? ')); - return mark + (c.detail || c.status) + (c.enabled ? '' : ' · désactivé'); + var tools = (c.kind === 'mcp' && c.tools_count) ? ' · ' + c.tools_count + ' outil(s)' : ''; + return mark + (c.detail || c.status) + tools + (c.enabled ? '' : ' · désactivé'); }, _connectorItems(){ var self = this; @@ -1377,6 +1379,20 @@ self.fetchConnectors(); }).catch(function(){ self.toast('Suppression impossible (réseau).', true); }); }, + // Type choisi dans le formulaire → URL / schéma d'auth / aide pré-remplis. + connectorPreset(){ + var f = this.connectorForm; + var presets = { + custom: {url:'', auth:'bearer', hint:'URL publique de votre API (hôtes privés refusés)'}, + discord: {url:'https://discord.com/api/v10', auth:'bot', + hint:'Discord Developers → Bot Token (jeton en Authorization)'}, + telegram:{url:'https://api.telegram.org/bot{secret}', auth:'none', + hint:'BotFather → Bot Token ({secret} complète l\u2019URL)'}, + mcp: {url:'', auth:'bearer', hint:'URL du endpoint MCP (streamable HTTP)'} + }; + var p = presets[f.kind] || presets.custom; + f.url = p.url; f.auth = p.auth; f.hint = p.hint; + }, connectorCreate(){ var self = this, f = this.connectorForm; if(!(f.name || '').trim() || !(f.url || '').trim()){ @@ -1384,7 +1400,8 @@ } fetch('/api/agent/connectors', { method:'POST', headers:{'X-CSRF-Token': getCsrf(), 'Content-Type':'application/json'}, - body: JSON.stringify({name: f.name, url: f.url, secret: f.secret || ''}) + body: JSON.stringify({name: f.name, url: f.url, secret: f.secret || '', + kind: f.kind || 'custom', auth: f.auth || 'bearer'}) }).then(function(r){ return r.json().then(function(d){ return {ok:r.ok, d:d}; }); }) .then(function(res){ if(!res.ok){ self.toast((res.d && res.d.detail) || 'Création impossible.', true); return; } @@ -1558,7 +1575,8 @@ this.connector = it.connector; this.plusSection = 'connector'; this.plusFocus = 0; return; } if(it.action === 'connector-new'){ - this.connectorForm = {name:'', url:'', secret:''}; + this.connectorForm = {name:'', url:'', secret:'', kind:'custom', + auth:'bearer', hint:''}; this.plusSection = 'connector-new'; this.plusFocus = -1; return; } if(it.action === 'connector-probe'){ this.connectorProbe(); return; } diff --git a/tests/test_v757_mcp_discord.py b/tests/test_v757_mcp_discord.py new file mode 100644 index 0000000..dbf1714 --- /dev/null +++ b/tests/test_v757_mcp_discord.py @@ -0,0 +1,209 @@ +"""v7.57.0 — Discord / Telegram (presets) + serveurs MCP (outils dynamiques).""" + +import asyncio +from pathlib import Path + +import pytest +from conftest import anon_csrf + +from app.db import get_conn +from app.services import connectors, mcp_client +from app.services.tool_registry import ToolRegistry + +ROOT = Path(__file__).resolve().parent.parent +JS_PATH = ROOT / "static" / "js" / "agent_panel_2.js" +PANEL_HTML = ROOT / "app" / "templates" / "agent_panel.html" + + +def _create(client, **kw): + body = {"name": "Conn", "url": "https://example.com/v1", "secret": "sk-1", **kw} + r = client.post("/api/agent/connectors", json=body) + assert r.status_code == 200, r.text + return r.json() + + +# ── Presets Discord / Telegram ────────────────────────────────────────────── + +def test_discord_preset_uses_bot_auth(client): + row = _create(client, name="Discord", url="https://discord.com/api/v10", + auth="bot", kind="discord") + assert row["kind"] == "discord" and row["auth"] == "bot" + assert connectors._headers(row["id"])["Authorization"].startswith("Bot ") + + +def test_telegram_secret_lives_in_url_not_in_db(client, monkeypatch): + row = _create(client, name="Telegram", kind="telegram", auth="none", + url="https://api.telegram.org/bot{secret}", secret="123456:abc") + with get_conn() as conn: + stored = conn.execute("SELECT url FROM agent_connectors WHERE id=?", + (row["id"],)).fetchone()["url"] + assert "{secret}" in stored and "123456" not in stored # jamais en clair + + seen = [] + + async def fake_get(url, headers=None): + seen.append((url, headers)) + return 200, '{"ok": true}' + + monkeypatch.setattr(connectors, "_get", fake_get) + asyncio.run(connectors.connector_fetch(str(row["id"]), path="getMe")) + assert seen[0][0] == "https://api.telegram.org/bot123456:abc/getMe" + assert "Authorization" not in seen[0][1] # auth=none + + +def test_secret_placeholder_requires_key_and_valid_kinds(client): + no_key = client.post("/api/agent/connectors", + json={"name": "TG", "url": "https://api.telegram.org/bot{secret}"}) + assert no_key.status_code == 400 and "secret" in no_key.json()["detail"] + bad_kind = client.post("/api/agent/connectors", + json={"name": "X", "url": "https://example.com/x", "kind": "slack"}) + assert bad_kind.status_code == 400 and "kind" in bad_kind.json()["detail"] + bad_auth = client.post("/api/agent/connectors", + json={"name": "X", "url": "https://example.com/x", "auth": "digest"}) + assert bad_auth.status_code == 400 and "auth" in bad_auth.json()["detail"] + + +# ── MCP : handshake, cache, outils dynamiques ─────────────────────────────── + +def _mcp_server(client) -> int: + return _create(client, name="Demo Server", kind="mcp", url="https://example.com/mcp")["id"] + + +def _fake_rpc_factory(calls): + async def fake_rpc(url, payload, headers=None): + assert url.startswith("https://example.com/mcp") # SSRF: hôte public fixe + calls.append(payload.get("method")) + method = payload.get("method") + if method == "initialize": + return {"result": {"protocolVersion": "2024-11-05", + "serverInfo": {"name": "demo"}}} + if method == "tools/list": + return {"result": {"tools": [ + {"name": "search_docs", "description": "Cherche dans les notes", + "inputSchema": {"type": "object", "properties": {"q": {"type": "string"}}, + "required": ["q"]}}, + {"name": "Echo", "description": "Répète"}, + ]}} + if method == "tools/call": + return {"result": {"content": [{"type": "text", "text": "résultat utile"}]}} + return {} + return fake_rpc + + +def test_mcp_probe_handshakes_and_caches_tools(client, monkeypatch): + sid = _mcp_server(client) + calls = [] + monkeypatch.setattr(mcp_client, "_rpc", _fake_rpc_factory(calls)) + + res = asyncio.run(connectors.probe(str(sid))) + assert res["status"] == "ok", res + assert "2 outil(s)" in res["detail"] + assert calls == ["initialize", "notifications/initialized", "tools/list"] + + with get_conn() as conn: + tools = conn.execute("SELECT tools_json FROM agent_connectors WHERE id=?", + (sid,)).fetchone()["tools_json"] + assert "mcp_demo_server_search_docs" in tools # nom LLM calculé + rows = client.get("/api/agent/connectors").json()["connectors"] + me = {r["id"]: r for r in rows if r["id"] == sid}[sid] + assert me["kind"] == "mcp" and me["tools_count"] == 2 + + +def test_dynamic_tools_reach_the_llm_schema_and_execute(client, monkeypatch): + sid = _mcp_server(client) + monkeypatch.setattr(mcp_client, "_rpc", _fake_rpc_factory([])) + asyncio.run(connectors.probe(str(sid))) + + reg = ToolRegistry() + schema = {t["name"]: t for t in reg.schema()} + assert "mcp_demo_server_search_docs" in schema + assert schema["mcp_demo_server_search_docs"]["parameters"]["required"] == ["q"] + assert "mcp_demo_server_echo" in schema + + res = asyncio.run(reg.execute("mcp_demo_server_search_docs", {"q": "notes"})) + assert res.status == "success" + assert "résultat utile" in res.data["text"] + + +def test_disabled_mcp_server_hides_its_tools(client, monkeypatch): + sid = _mcp_server(client) + monkeypatch.setattr(mcp_client, "_rpc", _fake_rpc_factory([])) + asyncio.run(connectors.probe(str(sid))) + assert "mcp_demo_server_echo" in {t["name"] for t in ToolRegistry().schema()} + + client.patch(f"/api/agent/connectors/{sid}", json={"enabled": False}) + assert "mcp_demo_server_echo" not in {t["name"] for t in ToolRegistry().schema()} + # hors du registre → refus explicite plutôt qu'un appel réseau raté + res = asyncio.run(ToolRegistry().execute("mcp_demo_server_echo", {})) + assert res.status == "error" and "inconnu" in res.message.lower() + + +def test_mcp_rpc_errors_become_probe_errors(client, monkeypatch): + sid = _mcp_server(client) + + async def bad_rpc(url, payload, headers=None): + if payload.get("method") == "initialize": + raise ValueError("MCP -32000 : handshake refusé") + return {} + + monkeypatch.setattr(mcp_client, "_rpc", bad_rpc) + res = asyncio.run(connectors.probe(str(sid))) + assert res["status"] == "error" and "handshake" in res["detail"] + # les outils restent absents : pas de cache partiel + with get_conn() as conn: + tools = conn.execute("SELECT tools_json FROM agent_connectors WHERE id=?", + (sid,)).fetchone()["tools_json"] + assert tools == "[]" + + +def test_parse_body_handles_sse_and_rpc_errors(): + sse = mcp_client.parse_body( + "text/event-stream; charset=utf-8", + ": keep-alive\nevent: message\ndata: {\"result\": {\"ok\": 1}}\n\n", 200) + assert sse == {"result": {"ok": 1}} + assert mcp_client.parse_body("application/json", '{"result": {}}', 200) == {"result": {}} + with pytest.raises(ValueError, match="illisible"): + mcp_client.parse_body("text/html", "gateway", 502) + with pytest.raises(ValueError, match="-32601"): + mcp_client.parse_body("application/json", + '{"error": {"code": -32601, "message": "Method not found"}}', 200) + + +def test_tool_name_slug(): + assert mcp_client.tool_name("Demo Server", "search-docs") == "mcp_demo_server_search_docs" + assert mcp_client.tool_name("!!", "Echo") == "mcp_echo" + + +# ── Teams + câblage ───────────────────────────────────────────────────────── + +def test_ms365_scopes_include_teams_messages(): + from app.services.oauth_connectors import PROVIDERS + assert "ChannelMessage.Read.All" in PROVIDERS["ms365"]["scopes"] + assert "Files.Read" in PROVIDERS["ms365"]["scopes"] + + +def test_connectors_menu_wired_for_presets(client): + src = JS_PATH.read_text(encoding="utf-8") + for token in ("connectorPreset()", "kind: f.kind || 'custom'", + "auth: f.auth || 'bearer'", + "https://discord.com/api/v10", + "https://api.telegram.org/bot{secret}", "presets[f.kind]"): + assert token in src, token + html = PANEL_HTML.read_text(encoding="utf-8") + assert "fd-plus-cn-kind" in html + for token in ('value="discord"', 'value="telegram"', 'value="mcp"', + 'connectorForm.kind'): + assert token in html, token + # 1 seule section « bientôt » restante (plugins) + root = src.split("var FD_PLUS_MENU = [", 1)[1].split("];", 1)[0] + assert root.count("disabled:true") == 1 + resp = client.get("/accounts") + assert resp.status_code == 200 and "connectorPreset" in resp.text + + +def test_mcp_routes_require_session(client): + anon_csrf(client) + assert client.get("/api/agent/connectors").status_code == 401 + assert client.post("/api/agent/connectors", + json={"name": "x", "url": "https://example.com/x", + "kind": "mcp"}).status_code == 401