import asyncio import hashlib import json import os import re import logging import secrets import shutil import sys import time from contextlib import asynccontextmanager from typing import Any, Optional from urllib.parse import urlencode import httpx from dotenv import load_dotenv from fastapi import HTTPException, Query from fastapi.responses import HTMLResponse from backend.config.Apps import SubApp from backend.apps.tools_lib.models import ToolDefinition, ToolCreate, ToolUpdate, BUILTIN_TOOLS logger = logging.getLogger(__name__) # Default Google OAuth credentials for the OpenSwarm project. # These are public credentials for a desktop/web OAuth client (safe to embed per Google's docs). # Users can override via GOOGLE_OAUTH_CLIENT_ID / GOOGLE_OAUTH_CLIENT_SECRET env vars. _DEFAULT_GOOGLE_CLIENT_ID = "6741219524-8vpt07arcc5rvkdb4j1b6v9g53469ugq.apps.googleusercontent.com" _DEFAULT_GOOGLE_CLIENT_SECRET = "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" os.environ.setdefault("GOOGLE_OAUTH_CLIENT_ID", _DEFAULT_GOOGLE_CLIENT_ID) os.environ.setdefault("GOOGLE_OAUTH_CLIENT_SECRET", _DEFAULT_GOOGLE_CLIENT_SECRET) from backend.config.paths import BACKEND_DIR, DATA_ROOT, TOOLS_DIR as DATA_DIR, BUILTIN_PERMISSIONS_PATH as BUILTIN_PERMS_PATH load_dotenv(os.path.join(BACKEND_DIR, ".env")) if os.environ.get("OPENSWARM_PACKAGED") == "1": load_dotenv(os.path.join(os.path.dirname(DATA_ROOT), ".env"), override=True) @asynccontextmanager async def tools_lib_lifespan(): os.makedirs(DATA_DIR, exist_ok=True) yield tools_lib = SubApp("tools", tools_lib_lifespan) GOOGLE_AUTH_URL = "https://accounts.google.com/o/oauth2/v2/auth" GOOGLE_TOKEN_URL = "https://oauth2.googleapis.com/token" GOOGLE_USERINFO_URL = "https://www.googleapis.com/oauth2/v2/userinfo" GOOGLE_SCOPES = [ "openid", "https://www.googleapis.com/auth/userinfo.email", "https://www.googleapis.com/auth/gmail.modify", "https://www.googleapis.com/auth/calendar", "https://www.googleapis.com/auth/drive", "https://www.googleapis.com/auth/contacts.readonly", ] AIRTABLE_AUTH_URL = "https://airtable.com/oauth2/v1/authorize" AIRTABLE_TOKEN_URL = "https://airtable.com/oauth2/v1/token" AIRTABLE_SCOPES = [ "data.records:read", "data.records:write", "data.recordComments:read", "data.recordComments:write", "schema.bases:read", "schema.bases:write", "user.email:read", ] HUBSPOT_AUTH_URL = "https://mcp-na2.hubspot.com/oauth/authorize/user" HUBSPOT_TOKEN_URL = "https://api.hubapi.com/oauth/v1/token" DISCORD_AUTH_URL = "https://discord.com/oauth2/authorize" DISCORD_TOKEN_URL = "https://discord.com/api/oauth2/token" # Maps state -> {tool_id, code_verifier (for PKCE flows)} _pending_oauth: dict[str, dict] = {} def _load_all() -> list[ToolDefinition]: result = [] if not os.path.exists(DATA_DIR): return result for fname in os.listdir(DATA_DIR): if fname.endswith(".json"): with open(os.path.join(DATA_DIR, fname)) as f: result.append(ToolDefinition(**json.load(f))) return result def _save(tool: ToolDefinition): with open(os.path.join(DATA_DIR, f"{tool.id}.json"), "w") as f: json.dump(tool.model_dump(), f, indent=2) def _load(tool_id: str) -> ToolDefinition: path = os.path.join(DATA_DIR, f"{tool_id}.json") if not os.path.exists(path): raise HTTPException(status_code=404, detail="Tool not found") with open(path) as f: return ToolDefinition(**json.load(f)) @tools_lib.router.get("/builtin") async def list_builtin_tools(): return {"tools": [t.model_dump() for t in BUILTIN_TOOLS]} def load_builtin_permissions() -> dict[str, str]: if not os.path.exists(BUILTIN_PERMS_PATH): return {} with open(BUILTIN_PERMS_PATH) as f: return json.load(f) def save_builtin_permissions(perms: dict[str, str]): os.makedirs(os.path.dirname(BUILTIN_PERMS_PATH), exist_ok=True) with open(BUILTIN_PERMS_PATH, "w") as f: json.dump(perms, f, indent=2) @tools_lib.router.get("/builtin/permissions") async def get_builtin_permissions(): return {"permissions": load_builtin_permissions()} @tools_lib.router.put("/builtin/permissions") async def update_builtin_permissions(body: dict): valid_tools = {t.name for t in BUILTIN_TOOLS} valid_policies = {"always_allow", "ask", "deny"} perms = load_builtin_permissions() for name, policy in body.get("permissions", {}).items(): if name in valid_tools and policy in valid_policies: perms[name] = policy save_builtin_permissions(perms) return {"permissions": perms} @tools_lib.router.get("/list") async def list_tools(): return {"tools": [t.model_dump() for t in _load_all()]} @tools_lib.router.get("/oauth/callback") async def oauth_callback(code: str = Query(...), state: str = Query("")): pending = _pending_oauth.pop(state, None) if not pending: return HTMLResponse("
{resp.text}", status_code=400)
tokens = resp.json()
tool.oauth_tokens = {
"access_token": tokens.get("access_token", ""),
"refresh_token": tokens.get("refresh_token", ""),
"token_expiry": time.time() + tokens.get("expires_in", 7200),
}
tool.auth_type = "oauth2"
tool.auth_status = "connected"
tool.connected_account_email = "Airtable account"
elif tool.name.lower() == "hubspot":
# HubSpot OAuth 2.1: PKCE flow
client_id = os.environ.get("HUBSPOT_OAUTH_CLIENT_ID", "")
client_secret = os.environ.get("HUBSPOT_OAUTH_CLIENT_SECRET", "")
async with httpx.AsyncClient(timeout=15.0) as client:
resp = await client.post(HUBSPOT_TOKEN_URL, data={
"grant_type": "authorization_code",
"code": code,
"redirect_uri": redirect_uri,
"client_id": client_id,
"client_secret": client_secret,
"code_verifier": code_verifier or "",
}, headers={
"Content-Type": "application/x-www-form-urlencoded",
})
if resp.status_code != 200:
logger.warning(f"HubSpot OAuth token exchange failed: {resp.text}")
return HTMLResponse(f"{resp.text}", status_code=400)
tokens = resp.json()
tool.oauth_tokens = {
"access_token": tokens.get("access_token", ""),
"refresh_token": tokens.get("refresh_token", ""),
"token_expiry": time.time() + tokens.get("expires_in", 1800),
}
tool.auth_type = "oauth2"
tool.auth_status = "connected"
tool.connected_account_email = "HubSpot account"
elif tool.name.lower() == "discord":
# Discord bot install OAuth: exchange code, capture guild_id of the
# server the user added the bot to. Multiple connect calls APPEND
# additional guild_ids so users can authorize multiple servers.
client_id = os.environ.get("DISCORD_OAUTH_CLIENT_ID", "")
client_secret = os.environ.get("DISCORD_OAUTH_CLIENT_SECRET", "")
async with httpx.AsyncClient(timeout=15.0) as client:
resp = await client.post(DISCORD_TOKEN_URL, data={
"grant_type": "authorization_code",
"code": code,
"redirect_uri": redirect_uri,
"client_id": client_id,
"client_secret": client_secret,
}, headers={
"Content-Type": "application/x-www-form-urlencoded",
})
if resp.status_code != 200:
logger.warning(f"Discord OAuth token exchange failed: {resp.text}")
return HTMLResponse(f"{resp.text}", status_code=400)
tokens = resp.json()
guild = tokens.get("guild") or {}
new_guild_id = guild.get("id", "")
new_guild_name = guild.get("name", "")
existing = tool.oauth_tokens.get("guilds") or []
# Append unless this guild was already authorized
if new_guild_id and not any(g.get("id") == new_guild_id for g in existing):
existing.append({"id": new_guild_id, "name": new_guild_name})
tool.oauth_tokens = {
# Bot token lives in .env, NEVER stored on the tool. We only
# track the list of authorized guilds for scope enforcement.
"guilds": existing,
}
tool.auth_type = "oauth2"
tool.auth_status = "connected"
names = ", ".join(g.get("name", "") for g in existing if g.get("name"))
tool.connected_account_email = f"{len(existing)} server{'s' if len(existing) != 1 else ''}" + (f" · {names}" if names else "")
elif tool.name.lower() == "notion":
# Notion OAuth: Basic auth with client_id:secret
notion_client_id = os.environ.get("NOTION_OAUTH_CLIENT_ID", "")
notion_client_secret = os.environ.get("NOTION_OAUTH_CLIENT_SECRET", "")
import base64
credentials = base64.b64encode(f"{notion_client_id}:{notion_client_secret}".encode()).decode()
async with httpx.AsyncClient(timeout=15.0) as client:
resp = await client.post("https://api.notion.com/v1/oauth/token", json={
"grant_type": "authorization_code",
"code": code,
"redirect_uri": redirect_uri,
}, headers={
"Authorization": f"Basic {credentials}",
"Content-Type": "application/json",
})
if resp.status_code != 200:
logger.warning(f"Notion OAuth token exchange failed: {resp.text}")
return HTMLResponse(f"{resp.text}", status_code=400)
tokens = resp.json()
tool.oauth_tokens = {
"access_token": tokens.get("access_token", ""),
}
tool.auth_type = "oauth2"
tool.auth_status = "connected"
tool.connected_account_email = tokens.get("workspace_name", "Notion workspace")
else:
# Google OAuth
client_id = os.environ.get("GOOGLE_OAUTH_CLIENT_ID", "")
client_secret = os.environ.get("GOOGLE_OAUTH_CLIENT_SECRET", "")
async with httpx.AsyncClient(timeout=15.0) as client:
resp = await client.post(GOOGLE_TOKEN_URL, data={
"code": code,
"client_id": client_id,
"client_secret": client_secret,
"redirect_uri": redirect_uri,
"grant_type": "authorization_code",
})
if resp.status_code != 200:
logger.warning(f"OAuth token exchange failed: {resp.text}")
return HTMLResponse(f"{resp.text}", status_code=400)
tokens = resp.json()
access_token = tokens.get("access_token", "")
tool.oauth_tokens = {
"access_token": access_token,
"refresh_token": tokens.get("refresh_token", ""),
"token_expiry": time.time() + tokens.get("expires_in", 3600),
}
tool.auth_type = "oauth2"
tool.auth_status = "connected"
if access_token:
try:
async with httpx.AsyncClient(timeout=10.0) as info_client:
info_resp = await info_client.get(
GOOGLE_USERINFO_URL,
headers={"Authorization": f"Bearer {access_token}"},
)
if info_resp.status_code == 200:
tool.connected_account_email = info_resp.json().get("email")
except Exception as e:
logger.warning(f"Failed to fetch Google userinfo: {e}")
_save(tool)
return HTMLResponse("""
You can close this window.
""") @tools_lib.router.get("/{tool_id}") async def get_tool(tool_id: str): return _load(tool_id).model_dump() @tools_lib.router.post("/create") async def create_tool(body: ToolCreate): tool = ToolDefinition( name=body.name, description=body.description, command=body.command, mcp_config=body.mcp_config, credentials=body.credentials, auth_type=body.auth_type, auth_status=body.auth_status, ) _save(tool) return {"ok": True, "tool": tool.model_dump()} @tools_lib.router.put("/{tool_id}") async def update_tool(tool_id: str, body: ToolUpdate): tool = _load(tool_id) for k, v in body.model_dump(exclude_none=True).items(): setattr(tool, k, v) _save(tool) return {"ok": True, "tool": tool.model_dump()} @tools_lib.router.delete("/{tool_id}") async def delete_tool(tool_id: str): path = os.path.join(DATA_DIR, f"{tool_id}.json") if os.path.exists(path): os.remove(path) return {"ok": True} # --------------------------------------------------------------------------- # MCP config derivation # --------------------------------------------------------------------------- def _sanitize_server_name(name: str) -> str: """Convert a tool name into a valid MCP server identifier (alphanumeric + hyphens).""" return re.sub(r"[^a-z0-9]+", "-", name.lower()).strip("-") def _extra_bin_dirs() -> list[str]: """Well-known user-local bin directories that may not be on PATH in packaged apps.""" home = os.path.expanduser("~") # Bundled uv-bin (ships uvx for non-dev users) _backend = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) dirs = [ os.path.join(_backend, "uv-bin"), os.path.join(home, ".bun", "bin"), os.path.join(home, ".cargo", "bin"), os.path.join(home, ".local", "bin"), os.path.join(home, ".volta", "bin"), "/opt/homebrew/bin", "/usr/local/bin", ] # nvm: pick the newest installed node version nvm_node = os.path.join(home, ".nvm", "versions", "node") try: if os.path.isdir(nvm_node): versions = sorted(os.listdir(nvm_node), reverse=True) if versions: dirs.insert(0, os.path.join(nvm_node, versions[0], "bin")) except OSError: pass # fnm fnm_bin = os.path.join(home, "Library", "Application Support", "fnm", "aliases", "default", "bin") if os.path.isdir(fnm_bin): dirs.insert(0, fnm_bin) return dirs def _resolve_command(command: str) -> str | None: """Find a command on PATH, falling back to common user-local bin directories and bundled binaries (uv-bin for uvx/uv).""" found = shutil.which(command) if found: return found # Windows binaries need an extension. shutil.which() handles PATHEXT for # PATH lookups, but we manually scan _extra_bin_dirs below — replicate # the suffix probing here so `uvx` finds `uvx.exe`, etc. if sys.platform == "win32": suffixes = [""] + os.environ.get("PATHEXT", ".COM;.EXE;.BAT;.CMD").lower().split(os.pathsep) else: suffixes = [""] def _probe(directory: str) -> str | None: for suffix in suffixes: candidate = os.path.join(directory, command + suffix) if os.path.isfile(candidate) and os.access(candidate, os.X_OK): return candidate return None for d in _extra_bin_dirs(): hit = _probe(d) if hit: return hit # Check bundled uv-bin directory (ships uv/uvx for non-dev users) _backend = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) return _probe(os.path.join(_backend, "uv-bin")) def _augmented_path() -> str: """Return PATH with extra bin dirs prepended (for child process environments).""" extra = [d for d in _extra_bin_dirs() if os.path.isdir(d)] current = os.environ.get("PATH", "") seen: set[str] = set() parts: list[str] = [] for p in extra + current.split(os.pathsep): if p and p not in seen: seen.add(p) parts.append(p) return os.pathsep.join(parts) def derive_mcp_config(tool: ToolDefinition) -> Optional[dict]: """Build the claude_agent_sdk mcp_servers config entry for a tool. Returns None if the tool cannot be configured (e.g. missing data). """ if not tool.mcp_config: return None config: dict = dict(tool.mcp_config) if tool.credentials: if config.get("type") in ("http", "sse"): headers = config.setdefault("headers", {}) for key, val in tool.credentials.items(): if key.lower() in ("authorization", "api_key", "api-key"): headers.setdefault("Authorization", f"Bearer {val}") else: env = config.setdefault("env", {}) env.update(tool.credentials) if tool.oauth_tokens.get("access_token"): if config.get("type") in ("http", "sse"): headers = config.setdefault("headers", {}) headers["Authorization"] = f"Bearer {tool.oauth_tokens['access_token']}" else: env = config.setdefault("env", {}) env["OAUTH_ACCESS_TOKEN"] = tool.oauth_tokens["access_token"] if tool.name.lower() == "notion": env["NOTION_TOKEN"] = tool.oauth_tokens["access_token"] if tool.name.lower() == "hubspot": env["PRIVATE_APP_ACCESS_TOKEN"] = tool.oauth_tokens["access_token"] if tool.oauth_tokens.get("refresh_token"): env["GOOGLE_WORKSPACE_REFRESH_TOKEN"] = tool.oauth_tokens["refresh_token"] client_id = os.environ.get("GOOGLE_OAUTH_CLIENT_ID", "") client_secret = os.environ.get("GOOGLE_OAUTH_CLIENT_SECRET", "") if client_id: env["GOOGLE_WORKSPACE_CLIENT_ID"] = client_id if client_secret: env["GOOGLE_WORKSPACE_CLIENT_SECRET"] = client_secret # Discord: bot token is loaded from .env at MCP launch time. It is NEVER # stored on the tool definition or exposed to the frontend. The tool only # tracks the list of authorized guild IDs (in oauth_tokens.guilds) which # are used by the agent system prompt to scope what the agent may access. if tool.name.lower() == "discord" and config.get("type") == "stdio": bot_token = os.environ.get("DISCORD_BOT_TOKEN", "") if bot_token: env = config.setdefault("env", {}) env["DISCORD_TOKEN"] = bot_token # Microsoft 365 MCP: use a stable token cache path shared across process spawns if tool.name.lower() == "microsoft 365" and config.get("type") == "stdio": env = config.setdefault("env", {}) cache_dir = os.path.join(os.path.expanduser("~"), ".openswarm") os.makedirs(cache_dir, exist_ok=True) env["MS365_MCP_TOKEN_CACHE_PATH"] = os.path.join(cache_dir, "ms365-token-cache.json") env["MS365_MCP_SELECTED_ACCOUNT_PATH"] = os.path.join(cache_dir, "ms365-selected-account.json") if config.get("type") == "stdio": if config.get("command"): # Check for bundled npm MCP servers — use Electron's Node.js instead of npx if config["command"] in ("npx", "bunx"): pkg_name = next((a for a in (config.get("args") or []) if not a.startswith("-")), None) if pkg_name: _backend = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) electron_path = os.environ.get("OPENSWARM_ELECTRON_PATH") # Check for single-file bundle first (e.g. reddit-mcp-buddy) bundle_path = os.path.join(_backend, "mcp-bundles", f"{pkg_name}.js") if os.path.isfile(bundle_path) and electron_path: config["command"] = electron_path config["args"] = [bundle_path] config.setdefault("env", {})["ELECTRON_RUN_AS_NODE"] = "1" logger.info(f"Using bundled MCP server for {pkg_name}") else: # Check for pre-installed npm package (works in both dev and packaged modes) safe_dir = pkg_name.replace("/", "-").replace("@", "") npm_dir = os.path.join(_backend, "npm-servers", safe_dir) pkg_json_path = os.path.join(npm_dir, "node_modules", pkg_name, "package.json") if os.path.isfile(pkg_json_path): import json as _json with open(pkg_json_path) as f: pkg_meta = _json.load(f) bin_field = pkg_meta.get("bin", {}) entry = list(bin_field.values())[0] if isinstance(bin_field, dict) else bin_field node_cmd = electron_path or shutil.which("node") if node_cmd: config["command"] = node_cmd config["args"] = [os.path.join(npm_dir, "node_modules", pkg_name, entry)] if electron_path: config.setdefault("env", {})["ELECTRON_RUN_AS_NODE"] = "1" logger.info(f"Using pre-installed npm MCP server for {pkg_name}") if not os.path.isabs(config.get("command", "")): resolved = _resolve_command(config["command"]) if resolved: config["command"] = resolved else: logger.warning(f"Command '{config['command']}' not found on PATH or bundled directories") env = config.setdefault("env", {}) env.setdefault("PATH", _augmented_path()) env.setdefault("PYTHONPATH", "") # Point uv/uvx at our bundled Python — avoids macOS CLT popup on fresh Macs # and avoids downloading Python at runtime _is_packaged = os.environ.get("OPENSWARM_PACKAGED") == "1" _is_windows = sys.platform == "win32" if _is_packaged: _resources = os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))) if _is_windows: _bundled_python = os.path.join(_resources, "python-env", "python.exe") else: _bundled_python = os.path.join(_resources, "python-env", "bin", "python3") if os.path.exists(_bundled_python): env.setdefault("UV_PYTHON", _bundled_python) else: _backend = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) if _is_windows: _venv_python = os.path.join(_backend, ".venv", "Scripts", "python.exe") else: _venv_python = os.path.join(_backend, ".venv", "bin", "python3") if os.path.exists(_venv_python): env.setdefault("UV_PYTHON", _venv_python) return config # --------------------------------------------------------------------------- # OAuth2 flow for Google Workspace (and other OAuth providers) # --------------------------------------------------------------------------- # --------------------------------------------------------------------------- # MCP tool discovery # --------------------------------------------------------------------------- _READ_PREFIXES = ("get", "list", "read", "search", "fetch", "find", "query", "count", "check", "describe", "show", "download", "browse", "analy", "explain") _WRITE_PREFIXES = ("create", "write", "delete", "update", "send", "remove", "modify", "add", "set", "put", "post", "patch", "insert", "move", "copy", "rename", "archive", "trash", "publish", "approve", "reject") _SERVICE_RULES: list[tuple[list[str], str, str]] = [ # (keywords, service_name, group) # Google Workspace (["gmail"], "Gmail", "Google"), (["drive"], "Drive", "Google"), (["calendar", "event", "freebusy"], "Calendar", "Google"), (["spreadsheet", "sheet"], "Sheets", "Google"), (["doc", "paragraph", "table"], "Docs", "Google"), (["chat", "space", "reaction", "message"], "Chat", "Google"), (["form", "publish_settings"], "Forms", "Google"), (["presentation", "slide", "page"], "Slides", "Google"), (["task_list", "task"], "Tasks", "Google"), (["contact"], "Contacts", "Google"), (["script", "deployment", "version", "trigger"], "Apps Script", "Google"), (["search_custom", "search_engine"], "Search", "Google"), # YouTube (["transcript", "caption"], "Transcripts", "YouTube"), (["video_detail", "video_comment", "video_categor", "video_engagement"], "Videos", "YouTube"), (["search_video", "trending_video"], "Search", "YouTube"), (["channel_stat", "channel_top"], "Channels", "YouTube"), # Reddit (before Twitter so "search_reddit" etc. don't mis-match) (["subreddit"], "Subreddits", "Reddit"), (["search_reddit"], "Search", "Reddit"), (["post_detail"], "Posts", "Reddit"), (["user_analysis"], "Users", "Reddit"), (["reddit_explain"], "Reference", "Reddit"), ] def _categorize_tool(name: str) -> str: lower = name.lower().replace("_", " ").replace("-", " ").strip() for word in lower.split(): for prefix in _READ_PREFIXES: if word.startswith(prefix): return "read" for prefix in _WRITE_PREFIXES: if word.startswith(prefix): return "write" return "write" def _extract_service(name: str) -> tuple[str, str]: """Extract the service and group from a tool name (e.g. 'search_gmail_messages' -> ('Gmail', 'Google')).""" lower = name.lower() for keywords, display, group in _SERVICE_RULES: for kw in keywords: if kw in lower: return display, group return "Other", "" def _parse_sse_json(text: str) -> dict | None: """Extract JSON from an SSE response body (handles `data: {...}` lines).""" for line in text.splitlines(): stripped = line.strip() if stripped.startswith("data:"): payload = stripped[len("data:"):].strip() if payload: try: return json.loads(payload) except json.JSONDecodeError: continue try: return json.loads(text) except json.JSONDecodeError: return None async def _discover_mcp_tools_http(url: str, headers: dict | None = None) -> list[dict]: """Connect to a Streamable HTTP MCP server and call tools/list via JSON-RPC POST.""" h = { "Content-Type": "application/json", "Accept": "application/json, text/event-stream", **(headers or {}), } async with httpx.AsyncClient(timeout=30.0) as client: init_resp = await client.post(url, headers=h, json={ "jsonrpc": "2.0", "id": 1, "method": "initialize", "params": {"protocolVersion": "2025-03-26", "capabilities": {}, "clientInfo": {"name": "self-swarm", "version": "0.1.0"}}, }) if init_resp.status_code not in (200, 201): raise HTTPException(status_code=502, detail=f"MCP initialize failed: {init_resp.status_code}") session_id = init_resp.headers.get("mcp-session-id", "") if session_id: h["mcp-session-id"] = session_id await client.post(url, headers=h, json={ "jsonrpc": "2.0", "method": "notifications/initialized", }) list_resp = await client.post(url, headers=h, json={ "jsonrpc": "2.0", "id": 2, "method": "tools/list", "params": {}, }) if list_resp.status_code not in (200, 201): raise HTTPException(status_code=502, detail=f"MCP tools/list failed: {list_resp.status_code}") ct = list_resp.headers.get("content-type", "") if "text/event-stream" in ct: data = _parse_sse_json(list_resp.text) else: data = list_resp.json() if not data: raise HTTPException(status_code=502, detail="Empty response from MCP server") tools_list = data.get("result", {}).get("tools", []) return [{"name": t.get("name", ""), "description": t.get("description", ""), "inputSchema": t.get("inputSchema")} for t in tools_list] async def _discover_mcp_tools_sse(url: str, headers: dict | None = None) -> list[dict]: """Connect to a legacy SSE MCP server (GET event-stream + POST messages) and call tools/list.""" from mcp.client.sse import sse_client from mcp import ClientSession from mcp.types import Implementation try: async with sse_client( url=url, headers=headers, timeout=30, sse_read_timeout=30, ) as (read_stream, write_stream): async with ClientSession( read_stream, write_stream, client_info=Implementation(name="self-swarm", version="0.1.0"), ) as session: await session.initialize() result = await session.list_tools() return [{"name": t.name, "description": t.description or "", "inputSchema": t.inputSchema if t.inputSchema else None} for t in result.tools] except BaseExceptionGroup as eg: first = eg.exceptions[0] if eg.exceptions else eg raise HTTPException(status_code=502, detail=f"SSE discovery failed: {first}") from first _NPX_CACHE_RE = re.compile(r"_npx[/\\]([0-9a-f]{8,})[/\\]") def _try_heal_npx_cache(stderr: str) -> str | None: """On `ERR_MODULE_NOT_FOUND` pointing into `~/.npm/_npx/