Files
mcphub/core/user_endpoints.py
airano-ir 788439e377
Some checks failed
Release / Test before release (push) Has been cancelled
Release / Publish to PyPI (push) Has been cancelled
Release / Publish to Docker Hub (push) Has been cancelled
Release / Create GitHub Release (push) Has been cancelled
feat(F.7+F.17): v3.11.0 — Coolify plugin (67 tools) + tool access overhaul
Catch-up sync spanning v3.7.0 → v3.11.0 of the internal repo.

Platform
- Total tools: 565 → 633 (+68) across 10 plugins (Coolify added)
- Tests: 481 → 828 passing

New plugin: Coolify (67 tools, Track F.17)
- Applications (17): CRUD, lifecycle, logs, env vars
- Deployments (5): list/get/cancel/deploy, app history
- Servers (8): CRUD, resources, domains, validation
- Projects (8), Databases (16, 6 DB types + backups), Services (13)

Tool access system (Track F.7 → F.7d)
- Scope → category mapping with per-tool `category` + `sensitivity`
- Schema v7: `site_tool_toggles(site_id)` + `sites.tool_scope` column
- Schema v8: per-site API keys (`api_keys.site_id`)
- Plugin-specific access-level presets (WP / WC / Gitea / OpenPanel /
  Coolify 5-tier)
- Credential-requirement notice tailored per plugin and tier
- Admin Tools count card on service page
- Dropped redundant `write` tier on WP / WP Advanced / WooCommerce
  (admin-scope tool count = 0 → identical to admin tier)

Dashboard
- Unified site manage page (Connection / Tool Access / Connect)
- /dashboard/keys unified (was /api-keys and /connect)
- CSRF interceptor via meta-tag; removed conflicting cookie reader
- Tailwind: pre-built CSS (scripts/build-css.sh) replaces CDN

Docs
- README / DOCKER_README / CLAUDE updated to 633 tools / 10 plugins
- CHANGELOG entries for v3.7.0 → v3.11.0
- FastMCP compatibility note updated to 3.x (post-v3.5 upgrade)

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-04-14 21:05:07 +02:00

492 lines
17 KiB
Python

"""Per-user MCP endpoint handler (Track E.3).
Handles MCP JSON-RPC requests for user-owned sites at
``/u/{user_id}/{alias}/mcp``. Implements the Streamable HTTP transport
protocol directly (no per-user FastMCP instances).
Flow:
1. Validate user API key (Bearer token)
2. Look up user's site from SQLite
3. Decrypt credentials
4. For tools/list: return plugin tools (without ``site`` param)
5. For tools/call: create plugin instance, call method, return result
Usage:
# In server.py route registration:
from core.user_endpoints import user_mcp_handler
Route("/u/{user_id}/{alias}/mcp", endpoint=user_mcp_handler, methods=["POST"])
"""
import json
import logging
import os
import time
from copy import deepcopy
from typing import Any
from starlette.requests import Request
from starlette.responses import JSONResponse, Response
from core.tool_registry import ToolDefinition
logger = logging.getLogger(__name__)
# Per-user rate limiting defaults
USER_RATE_LIMIT_PER_MIN = int(os.getenv("USER_RATE_LIMIT_PER_MIN", "30"))
USER_RATE_LIMIT_PER_HR = int(os.getenv("USER_RATE_LIMIT_PER_HR", "500"))
# In-memory rate limit tracking: user_id -> list of timestamps
_rate_limits: dict[str, list[float]] = {}
def _check_user_rate_limit(user_id: str) -> tuple[bool, str]:
"""Check per-user rate limits.
Args:
user_id: User UUID.
Returns:
Tuple of (allowed, error_message). error_message is empty if allowed.
"""
now = time.time()
timestamps = _rate_limits.setdefault(user_id, [])
# Prune old entries (older than 1 hour)
cutoff_hr = now - 3600
_rate_limits[user_id] = [t for t in timestamps if t > cutoff_hr]
timestamps = _rate_limits[user_id]
# Check per-minute limit
cutoff_min = now - 60
recent_min = sum(1 for t in timestamps if t > cutoff_min)
if recent_min >= USER_RATE_LIMIT_PER_MIN:
return False, f"Rate limit exceeded: {USER_RATE_LIMIT_PER_MIN} requests/minute"
# Check per-hour limit
if len(timestamps) >= USER_RATE_LIMIT_PER_HR:
return False, f"Rate limit exceeded: {USER_RATE_LIMIT_PER_HR} requests/hour"
timestamps.append(now)
return True, ""
def _tools_to_mcp_schema(tools: list[ToolDefinition]) -> list[dict[str, Any]]:
"""Convert ToolDefinition objects into MCP ``tools/list`` response shape.
Strips the auto-injected ``site`` parameter, since user endpoints bind a
single site per alias.
"""
result = []
for tool_def in tools:
schema = deepcopy(tool_def.input_schema)
if "properties" in schema:
schema["properties"].pop("site", None)
if "required" in schema and "site" in schema["required"]:
schema["required"] = [r for r in schema["required"] if r != "site"]
result.append(
{
"name": tool_def.name,
"description": tool_def.description,
"inputSchema": schema,
}
)
return result
async def _get_visible_tools_for_site(
site_id: str,
key_scopes: list[str],
plugin_type: str,
) -> list[dict[str, Any]]:
"""Return tools/list payload filtered by key scope + site scope + toggles (F.7b)."""
from core.tool_access import get_tool_access_manager
access = get_tool_access_manager()
tools = await access.get_visible_tools(
site_id=site_id, key_scopes=key_scopes, plugin_type=plugin_type
)
return _tools_to_mcp_schema(tools)
async def _execute_tool(
tool_name: str,
arguments: dict[str, Any],
plugin_type: str,
config_dict: dict[str, Any],
) -> Any:
"""Execute a tool by creating a plugin instance and calling the method.
Uses the same pattern as unified_handler in tool_generator.py.
"""
from plugins import registry as plugin_registry
if not plugin_registry.is_registered(plugin_type):
return {"type": "text", "text": f"Error: Unknown plugin type '{plugin_type}'"}
method_name = tool_name
# Strip plugin_type prefix: "wordpress_list_posts" → "list_posts"
prefix = f"{plugin_type}_"
if method_name.startswith(prefix):
method_name = method_name[len(prefix) :]
try:
plugin_instance = plugin_registry.create_instance(
plugin_type,
project_id=f"user_{config_dict.get('alias', 'unknown')}",
config=config_dict,
)
if not hasattr(plugin_instance, method_name):
return {"type": "text", "text": f"Error: Method '{method_name}' not found"}
method = getattr(plugin_instance, method_name)
# Process arguments (parse JSON strings, filter None/empty)
filtered_args = {}
for key, value in arguments.items():
if value is None:
continue
if isinstance(value, str):
stripped = value.strip()
if stripped == "":
continue
if (stripped.startswith("{") and stripped.endswith("}")) or (
stripped.startswith("[") and stripped.endswith("]")
):
try:
value = json.loads(value)
except json.JSONDecodeError:
pass
filtered_args[key] = value
result = await method(**filtered_args)
return result
except Exception as e:
logger.error("Tool execution error %s: %s", tool_name, e, exc_info=True)
return {"type": "text", "text": f"Error: {type(e).__name__}: {e}"}
def _jsonrpc_error(req_id: Any, code: int, message: str) -> dict[str, Any]:
"""Build a JSON-RPC error response."""
return {
"jsonrpc": "2.0",
"id": req_id,
"error": {"code": code, "message": message},
}
def _jsonrpc_result(req_id: Any, result: Any) -> dict[str, Any]:
"""Build a JSON-RPC success response."""
return {
"jsonrpc": "2.0",
"id": req_id,
"result": result,
}
async def user_mcp_handler(request: Request) -> Response:
"""Handle MCP JSON-RPC requests for user endpoints.
Route: POST /u/{user_id}/{alias}/mcp
Validates user API key, looks up the site, and handles MCP protocol
methods (initialize, tools/list, tools/call).
"""
user_id = request.path_params.get("user_id", "")
alias = request.path_params.get("alias", "")
if not user_id or not alias:
return JSONResponse(
_jsonrpc_error(None, -32600, "Invalid endpoint path"),
status_code=400,
)
# --- Authentication ---
auth_header = request.headers.get("authorization", "")
if not auth_header.startswith("Bearer "):
return JSONResponse(
_jsonrpc_error(None, -32600, "Missing Authorization header"),
status_code=401,
)
api_key = auth_header[7:] # Strip "Bearer "
# Shared scope tracking — set by whichever auth path succeeds
key_scopes: list[str] = []
# Try mhu_ API key first, then fall back to OAuth JWT token
if api_key.startswith("mhu_"):
try:
from core.user_keys import get_user_key_manager
key_mgr = get_user_key_manager()
key_info = await key_mgr.validate_key(api_key)
except RuntimeError:
return JSONResponse(
_jsonrpc_error(None, -32600, "Authentication service unavailable"),
status_code=503,
)
if key_info is None:
return JSONResponse(
_jsonrpc_error(None, -32600, "Invalid API key"),
status_code=401,
)
if key_info["user_id"] != user_id:
return JSONResponse(
_jsonrpc_error(None, -32600, "API key does not match user"),
status_code=403,
)
# Check site-scoped key: if key is scoped to a site, it can only
# be used for that specific site's alias.
key_site_id = key_info.get("site_id")
if key_site_id:
try:
from core.database import get_database
db = get_database()
site = await db.get_site(key_site_id, user_id)
if site is None or site.get("alias") != alias:
return JSONResponse(
_jsonrpc_error(
None,
-32600,
"API key is scoped to a different site",
),
status_code=403,
)
except RuntimeError:
pass # DB unavailable — allow through
key_scopes = key_info.get("scopes", "read").split()
else:
# Try OAuth JWT token (issued after consent flow via GitHub/Google login)
try:
import jwt as pyjwt
from core.oauth import get_token_manager
token_manager = get_token_manager()
jwt_payload = token_manager.validate_access_token(api_key)
# sub = "user:{uuid}" — extract actual user_id
sub = jwt_payload.get("sub", "")
if not sub.startswith("user:"):
return JSONResponse(
_jsonrpc_error(None, -32600, "Token not authorized for user endpoint"),
status_code=403,
)
jwt_user_id = sub[len("user:") :]
if jwt_user_id != user_id:
return JSONResponse(
_jsonrpc_error(None, -32600, "Token user mismatch"),
status_code=403,
)
key_scopes = jwt_payload.get("scope", "read").split()
except pyjwt.ExpiredSignatureError:
return JSONResponse(
_jsonrpc_error(None, -32600, "Token expired"),
status_code=401,
)
except Exception:
return JSONResponse(
_jsonrpc_error(None, -32600, "Invalid token"),
status_code=401,
)
# --- Rate Limiting ---
allowed, rate_msg = _check_user_rate_limit(user_id)
if not allowed:
return JSONResponse(
_jsonrpc_error(None, -32600, rate_msg),
status_code=429,
)
# --- Look up site ---
try:
from core.database import get_database
db = get_database()
site = await db.get_site_by_alias(user_id, alias)
except RuntimeError:
return JSONResponse(
_jsonrpc_error(None, -32600, "Database unavailable"),
status_code=503,
)
if site is None:
return JSONResponse(
_jsonrpc_error(None, -32600, f"Site '{alias}' not found"),
status_code=404,
)
if site["status"] == "disabled":
return JSONResponse(
_jsonrpc_error(None, -32600, "Site is disabled"),
status_code=403,
)
# --- Plugin visibility check ---
from core.plugin_visibility import is_plugin_public
if not is_plugin_public(site["plugin_type"]):
return JSONResponse(
_jsonrpc_error(
None, -32600, f"Plugin '{site['plugin_type']}' is not currently available"
),
status_code=403,
)
# --- Parse JSON-RPC body ---
try:
body = await request.json()
except Exception:
return JSONResponse(
_jsonrpc_error(None, -32700, "Parse error"),
status_code=400,
)
req_id = body.get("id")
method = body.get("method", "")
params = body.get("params", {})
# --- Handle MCP methods ---
if method == "initialize":
version = "3.1.0"
return JSONResponse(
_jsonrpc_result(
req_id,
{
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}},
"serverInfo": {
"name": f"mcphub-{alias}",
"version": version,
},
},
)
)
elif method == "notifications/initialized":
# Notification — no response needed
return Response(status_code=204)
elif method == "tools/list":
tools = await _get_visible_tools_for_site(site["id"], key_scopes, site["plugin_type"])
return JSONResponse(_jsonrpc_result(req_id, {"tools": tools}))
elif method == "tools/call":
tool_name = params.get("name", "")
arguments = params.get("arguments", {})
# Verify tool belongs to this plugin type
plugin_prefix = f"{site['plugin_type']}_"
if not tool_name.startswith(plugin_prefix):
return JSONResponse(
_jsonrpc_error(req_id, -32601, f"Tool '{tool_name}' not available for this site")
)
# Check required scope
from core.tool_registry import get_tool_registry
registry = get_tool_registry()
tool_def = registry.get_by_name(tool_name)
if not tool_def:
return JSONResponse(_jsonrpc_error(req_id, -32601, f"Tool '{tool_name}' not found"))
required_scope = tool_def.required_scope
# key_scopes is set during authentication (both mhu_ and JWT paths)
# F.7b: enforce category-based scope allowlist in addition to the
# legacy read/write/admin hierarchy. A tool is allowed only if BOTH
# (a) the legacy hierarchy grants it, AND
# (b) the tool's category is in BOTH the key-scope set AND the
# site's stored tool_scope set (the narrower layer wins).
from core.tool_access import KNOWN_CATEGORIES, SCOPE_CUSTOM, scopes_to_categories
scope_hierarchy = {"read": 1, "write": 2, "admin": 3}
required_level = scope_hierarchy.get(required_scope, 0)
key_level = max([scope_hierarchy.get(s, 0) for s in key_scopes] + [0])
legacy_ok = key_level >= required_level
key_cats = scopes_to_categories(key_scopes)
key_category_ok = tool_def.category not in KNOWN_CATEGORIES or tool_def.category in key_cats
# Site-level scope check (skipped for "custom" preset).
site_scope = site.get("tool_scope") or "admin"
if site_scope and site_scope != SCOPE_CUSTOM:
site_cats = scopes_to_categories([site_scope])
site_category_ok = (
tool_def.category not in KNOWN_CATEGORIES or tool_def.category in site_cats
)
else:
site_category_ok = True
if not (legacy_ok and key_category_ok and site_category_ok):
return JSONResponse(
_jsonrpc_error(
req_id,
-32600,
f"Insufficient scope. Tool '{tool_name}' requires "
f"scope '{required_scope}' (category '{tool_def.category}').",
)
)
# F.7b: honour per-site tool toggles — a disabled tool cannot be called
# even if scopes would otherwise allow it.
try:
toggles = await db.get_site_tool_toggles(site["id"])
if not toggles.get(tool_name, True):
return JSONResponse(
_jsonrpc_error(
req_id,
-32600,
f"Tool '{tool_name}' is disabled for this site.",
)
)
except Exception as exc: # non-fatal — fall through on DB errors
logger.warning("Failed to check site tool toggles for %s: %s", site["id"], exc)
# Decrypt credentials
try:
from core.encryption import get_credential_encryption
encryptor = get_credential_encryption()
credentials = encryptor.decrypt_credentials(site["credentials"], site["id"])
except Exception as e:
logger.error("Credential decryption failed for site %s: %s", site["id"], e)
return JSONResponse(
_jsonrpc_error(req_id, -32603, "Failed to decrypt site credentials")
)
# Build config dict for plugin instantiation
config_dict = {
"site_url": site["url"],
"url": site["url"],
"alias": alias,
**credentials,
}
result = await _execute_tool(tool_name, arguments, site["plugin_type"], config_dict)
# Format result as MCP content
if isinstance(result, str):
content = [{"type": "text", "text": result}]
elif isinstance(result, dict) and "type" in result:
content = [result]
elif isinstance(result, list):
content = result
else:
content = [{"type": "text", "text": json.dumps(result, default=str)}]
return JSONResponse(_jsonrpc_result(req_id, {"content": content}))
else:
return JSONResponse(_jsonrpc_error(req_id, -32601, f"Method '{method}' not supported"))