Files
bruno f8acb906e5
FlowDeck CI / lint (push) Successful in 2m14s
FlowDeck CI / test (push) Failing after 17m51s
FlowDeck CI / docker (push) Skipped
feat: connecteurs Discord/Telegram/MCP + outils MCP dynamiques au registre (v7.57.0)
- 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_<serveur>_<outil> 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
2026-10-07 09:09:19 -04:00

337 lines
14 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Connecteurs de l'agent (v7.55.0) — socle : catalogue, statut, fetch gardé SSRF.
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
(Google, M365, Discord/Telegram, MCP) arriveront avec leur propre module.
"""
from __future__ import annotations
import json
import logging
from app.config import settings
from app.db import get_conn
from app.services import oauth_connectors
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)
# ── Natifs (config, sans réseau) ─────────────────────────────────────────────
def native_state(kind: str, user_id: int | None = None) -> tuple[str, str]:
"""(status, detail) d'un connecteur natif — dérivé de la config, sans réseau."""
if kind in OAUTH_KINDS:
try:
oauth_connectors._client(kind) # lève si client_id/secret absents
except ValueError as exc:
return ("missing", str(exc))
tok = oauth_connectors.tokens(kind, user_id)
if tok.get("access_token"):
scope = (tok.get("scope") or "").split()
return ("ok", f"connecté · {len(scope)} scope(s)")
return ("missing", "non connecté — lancer la connexion OAuth")
if kind == "gitea":
ok = bool(settings.gitea_url) and settings.gitea_token not in ("", "change-me")
return ("ok", settings.gitea_url if ok else "GITEA_TOKEN non configuré")
if kind == "github":
if settings.github_token:
return ("ok", "PAT configuré (30 req/min)")
return ("missing", "aucun GITHUB_TOKEN (10 req/min anonyme)")
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] = []
for kind in NATIVE_KINDS:
status, detail = native_state(kind, user_id)
out.append({
"id": None, "builtin": True, "kind": kind,
"name": (oauth_connectors.PROVIDERS[kind]["name"] if kind in OAUTH_KINDS
else kind.capitalize()),
"url": "", "enabled": True, "status": status, "detail": detail,
"has_secret": False, "oauth": kind in OAUTH_KINDS,
})
with get_conn() as conn:
rows = conn.execute(
"SELECT id, name, url, secret_encrypted, enabled, status, detail, "
"kind, auth, tools_json FROM agent_connectors ORDER BY id"
).fetchall()
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, "
"kind, auth, tools_json FROM agent_connectors WHERE id=?",
(connector_id,),
).fetchone()
return _row(r) if r else None
# ── CRUD (personnels) ────────────────────────────────────────────────────────
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")
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, kind, auth) VALUES (?,?,?,?,?,?,?)",
(name, url, encrypt_secret(secret or ""), "unknown", "", kind, auth),
)
conn.commit()
cid = cur.lastrowid
return get(cid)
def update_connector(connector_id: int, *, name=None, enabled=None, secret=None) -> dict | None:
row = get(connector_id)
if row is None:
return None
sets, params = [], []
if name is not None:
name = str(name).strip()
if not name:
raise ValueError("name est requis")
sets.append("name=?")
params.append(name)
if enabled is not None:
sets.append("enabled=?")
params.append(1 if enabled else 0)
if secret is not None:
sets.append("secret_encrypted=?")
params.append(encrypt_secret(str(secret)))
if sets:
with get_conn() as conn:
params.append(connector_id)
conn.execute(f"UPDATE agent_connectors SET {', '.join(sets)} WHERE id=?", params)
conn.commit()
return get(connector_id)
def delete_connector(connector_id: int) -> bool:
with get_conn() as conn:
cur = conn.execute("DELETE FROM agent_connectors WHERE id=?", (connector_id,))
conn.commit()
return cur.rowcount > 0
# ── 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.
``auth`` = bearer (défaut) | bot (Discord, jeton préfixé) | none (jeton ailleurs).
"""
headers = {"User-Agent": "FlowDeck-Connectors/1.0"}
secret, auth = "", "bearer"
if connector_id:
with get_conn() as conn:
r = conn.execute(
"SELECT secret_encrypted, auth FROM agent_connectors WHERE id=?",
(connector_id,),
).fetchone()
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).
Point d'injection des tests : on monkeypatche ``connectors._get``.
"""
from app.services.http_client import shared_client
from app.services.importers.url_fetch import _is_public_host, _validate_url
safe = _validate_url(url)
async with shared_client(timeout=10, follow_redirects=True, headers=headers or {}) as client:
resp = await client.get(safe)
if resp.url.host and not _is_public_host(resp.url.host):
raise ValueError("Redirection vers un hôte non autorisé")
return resp.status_code, (resp.text or "")[:FETCH_LIMIT]
def _resolve(connector: str) -> tuple[str, str | int]:
"""« 3 », « Ma clé API » ou « gitea » → (kind, id_custom|kind)."""
token = str(connector or "").strip()
if not token:
raise ValueError("connector est requis")
if token.isdigit():
row = get(int(token))
if row is None:
raise ValueError(f"Connecteur inconnu: {token}")
return (row["kind"], row["id"])
low = token.lower()
if low in NATIVE_KINDS:
return (low, low)
with get_conn() as conn:
row = conn.execute(
"SELECT id FROM agent_connectors WHERE lower(name)=?", (low,)
).fetchone()
if row:
return (get(row["id"])["kind"], row["id"])
raise ValueError(f"Connecteur inconnu: {token}")
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 == "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]
elif kind == "gitea":
code, _ = await _get(f"{settings.gitea_url.rstrip('/')}/api/v1/version",
{"Authorization": f"token {settings.gitea_token}"})
status, detail = ("ok" if code < 400 else "error"), f"HTTP {code}"
elif kind == "github":
code, _ = await _get("https://api.github.com/rate_limit", _gh_headers())
status, detail = ("ok" if code < 400 else "error"), f"HTTP {code}"
elif kind == "web":
status, detail = "ok", native_state("web")[1]
else:
row = get(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 not in NATIVE_KINDS:
with get_conn() as conn:
conn.execute(
"UPDATE agent_connectors SET status=?, detail=? WHERE id=?",
(status, detail, int(ref)),
)
conn.commit()
return {"connector": connector, "status": status, "detail": detail}
def _gh_headers() -> dict:
headers = {"User-Agent": "FlowDeck-Connectors/1.0",
"Accept": "application/vnd.github+json"}
if settings.github_token:
headers["Authorization"] = f"Bearer {settings.github_token}"
return headers
async def connector_fetch(connector: str, path: str = "", query: str = "",
user_id: int | None = None) -> dict:
"""Lit un connecteur pour l'LLM : natif = API, web = recherche, perso = GET."""
kind, ref = _resolve(connector)
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"
code, text = await _get(url, {"Authorization": f"token {settings.gitea_token}"})
elif kind == "github":
url = "https://api.github.com/" + path.lstrip("/") if path else "https://api.github.com/rate_limit"
code, text = await _get(url, _gh_headers())
elif kind == "web":
from app.services.web_search import search_web
if not (query or "").strip():
return {"status": "error", "text": "query est requis pour le connecteur web"}
results, provider = await search_web(query.strip(), 5)
text = "\n".join(f"- {r.get('title', '')} — {r.get('url', '')}\n {r.get('snippet', '')}"
for r in results) or "Aucun résultat."
return {"status": "ok", "text": f"provider: {provider}\n{text}"[:FETCH_LIMIT]}
else:
row = get(int(ref))
if row is None:
return {"status": "error", "text": "Connecteur supprimé"}
if not row["enabled"]:
return {"status": "error", "text": "Connecteur désactivé"}
url = _with_secret(row["url"], int(ref)).rstrip("/")
if path:
url += "/" + path.lstrip("/")
code, text = await _get(url, _headers(int(ref)))
if code >= 400:
return {"status": "error", "text": f"HTTP {code} — {text[:400]}"}
return {"status": "ok", "text": text[:FETCH_LIMIT]}