diff --git a/README.md b/README.md index cd85754..6f2c586 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ > 本文件用于团队开发期间快速配置环境和启动项目,不是正式的项目 README。 -> 当前基线:2026-09-02。第一阶段 Web 联调前后端已经完成;第二阶段已完成 Workspace 去 Mock、Agent Trace 持久化与 SSE 恢复、stdio MCP Bridge、隔离 Plugin Host,以及 Plugin Command/Settings 后端 Contract 和前端 Service。真实音频、Provider 协议增强、Benchmark、导出、主题包、Trace 可视化、Mermaid 与函数图像仍在后续开发;Tauri Host、Stronghold、原生多 Vault 文件系统和 Sync Server 尚未接入。 +> 当前基线:2026-09-03。第一阶段 Web 联调前后端已经完成;第二阶段已完成 Workspace 去 Mock、Agent Trace 持久化与 SSE 恢复、stdio MCP Bridge、隔离 Plugin Host、Plugin Command/Settings,以及独立 MCP Server 配置中心 C.1/P0。Streamable HTTP MCP、真实音频、Provider 协议增强、Benchmark、导出、主题包、Trace 可视化、Mermaid 与函数图像仍在后续开发;Tauri Host、Stronghold、原生多 Vault 文件系统和 Sync Server 尚未接入。 ## 当前目录 diff --git a/backend/app/container.py b/backend/app/container.py index 62b21be..ee05cd8 100644 --- a/backend/app/container.py +++ b/backend/app/container.py @@ -5,6 +5,7 @@ from app.agent.builtin_tools import register_builtin_tools from app.contracts import ModelCapability, ProviderConfig, ProviderType from app.config import BACKEND_DIR, get_settings from app.extensions import PluginRuntime, SkillRuntime +from app.extensions.mcp_registry import McpServerRegistry from app.providers import MockProvider, ProviderFactory, ProviderRegistry from app.providers.credentials import ( ChainedCredentialResolver, @@ -22,6 +23,7 @@ class ApplicationContainer: permissions: PermissionManager skills: SkillRuntime plugins: PluginRuntime + mcp_servers: McpServerRegistry agent: AgentRuntime @@ -61,6 +63,14 @@ def build_container() -> ApplicationContainer: plugins.install(BACKEND_DIR / "extensions" / "plugins" / "text-tools") plugins.enable("text-tools") + mcp_servers = McpServerRegistry( + tools, + credentials, + settings.data_dir, + allow_process_launch=settings.environment == "development", + ) + mcp_servers.restore_enabled() + skills = SkillRuntime(tools) skills.install(BACKEND_DIR / "extensions" / "skills" / "knowledge-assistant") skills.enable("knowledge-assistant") @@ -81,6 +91,7 @@ def build_container() -> ApplicationContainer: permissions=permissions, skills=skills, plugins=plugins, + mcp_servers=mcp_servers, agent=agent, ) diff --git a/backend/app/contracts.py b/backend/app/contracts.py index f95e890..ee8f3ae 100644 --- a/backend/app/contracts.py +++ b/backend/app/contracts.py @@ -199,7 +199,7 @@ class ToolDefinition(Contract): description: str parameters: dict[str, Any] = Field(default_factory=dict) permission: str | None = None - source: Literal["builtin", "plugin"] = "builtin" + source: Literal["builtin", "plugin", "mcp_server"] = "builtin" class ToolCall(Contract): @@ -486,6 +486,72 @@ class PluginHostStatus(Contract): error: str | None = None +# Independent user-managed MCP Server Registry. This is deliberately separate +# from Plugin manifests: a server can contribute tools without being a Plugin. +class McpServerTransport(str, Enum): + stdio = "stdio" + streamable_http = "streamable_http" + sse = "sse" + + +class McpServerCreateRequest(Contract): + name: str = Field(min_length=1, max_length=80) + transport: McpServerTransport = McpServerTransport.stdio + command: str = Field(min_length=1, max_length=1024) + args: list[str] = Field(default_factory=list, max_length=64) + environment: dict[str, str] = Field(default_factory=dict) + secret_environment_keys: list[str] = Field(default_factory=list) + permissions: list[str] = Field(default_factory=list) + startup_timeout_seconds: float = Field(default=15, ge=1, le=120) + tool_timeout_seconds: float = Field(default=30, ge=1, le=300) + + +class McpServerUpdateRequest(McpServerCreateRequest): + pass + + +class McpServerSecretWriteRequest(Contract): + secret: SecretStr = Field(min_length=1, max_length=32768) + + +class McpServerSecretStatus(Contract): + key: str + configured: bool + + +class McpServerTrustRequest(Contract): + command_digest: str = Field(min_length=64, max_length=64) + + +class McpServer(Contract): + server_id: str + name: str + transport: McpServerTransport + command: str + args: list[str] = Field(default_factory=list) + environment: dict[str, str] = Field(default_factory=dict) + secret_environment: dict[str, bool] = Field(default_factory=dict) + permissions: list[str] = Field(default_factory=list) + startup_timeout_seconds: float + tool_timeout_seconds: float + enabled: bool = False + trusted: bool = False + command_digest: str + command_summary: str + status: PluginHostState = PluginHostState.stopped + tools_count: int = 0 + protocol_version: str | None = None + remote_server_name: str | None = None + remote_server_version: str | None = None + error: str | None = None + last_tested_at: datetime | None = None + last_test_succeeded: bool | None = None + + +class McpServerListResponse(Contract): + items: list[McpServer] = Field(default_factory=list) + + class PluginCommandLocation(str, Enum): command_palette = "command_palette" context_menu = "context_menu" diff --git a/backend/app/extensions/mcp.py b/backend/app/extensions/mcp.py index dc3e62f..834ef65 100644 --- a/backend/app/extensions/mcp.py +++ b/backend/app/extensions/mcp.py @@ -74,12 +74,14 @@ class McpStdioClient: command: list[str], *, cwd: Path, + environment: dict[str, str] | None = None, on_seen: Callable[[], None], on_broken: Callable[[str], None], on_tools_changed: Callable[[], None], ) -> None: self.command = command self.cwd = cwd + self.environment = environment or {} self.on_seen = on_seen self.on_broken = on_broken self.on_tools_changed = on_tools_changed @@ -99,6 +101,7 @@ class McpStdioClient: # 平台级沙箱启动器;uvx 只隔离 Python 依赖,不能替代系统权限限制。 creation_flags = getattr(subprocess, "CREATE_NO_WINDOW", 0) if os.name == "nt" else 0 environment = _subprocess_environment() + environment.update(self.environment) environment.setdefault("PYTHONUNBUFFERED", "1") try: self.process = subprocess.Popen( @@ -383,6 +386,10 @@ class McpBridge: package_path: Path, declared_permissions: list[str], on_unavailable: Callable[[str, str], None], + *, + command_override: list[str] | None = None, + environment: dict[str, str] | None = None, + tool_source: str = "plugin", ) -> list[McpDiscoveredTool]: if backend.transport != "stdio": raise McpBridgeError( @@ -390,7 +397,7 @@ class McpBridge: "Phase C only supports the MCP stdio transport.", status_code=501, ) - command = self._resolve_command(package_path, backend) + command = command_override or self._resolve_command(package_path, backend) now = datetime.now(timezone.utc) status = PluginHostStatus( plugin_id=plugin_id, @@ -420,6 +427,7 @@ class McpBridge: client = McpStdioClient( command, cwd=package_path, + environment=environment, on_seen=seen, on_broken=broken, on_tools_changed=tools_changed, @@ -470,7 +478,7 @@ class McpBridge: status.server_version = _optional_string(server_info.get("version")) client.notify("notifications/initialized") discovered = self._discover_tools( - plugin_id, client, backend, declared_permissions + plugin_id, client, backend, declared_permissions, tool_source ) status.status = PluginHostState.ready status.tools_count = len(discovered) @@ -601,6 +609,7 @@ class McpBridge: client: McpStdioClient, backend: PluginBackend, declared_permissions: list[str], + tool_source: str, ) -> list[McpDiscoveredTool]: discovered: list[McpDiscoveredTool] = [] cursor: str | None = None @@ -620,7 +629,7 @@ class McpBridge: ) for raw in raw_tools: discovered.append( - self._map_tool(plugin_id, raw, declared_permissions) + self._map_tool(plugin_id, raw, declared_permissions, tool_source) ) if len(discovered) > MAX_MCP_TOOLS: raise McpBridgeError( @@ -648,7 +657,10 @@ class McpBridge: @staticmethod def _map_tool( - plugin_id: str, raw: Any, declared_permissions: list[str] + plugin_id: str, + raw: Any, + declared_permissions: list[str], + tool_source: str = "plugin", ) -> McpDiscoveredTool: if not isinstance(raw, dict): raise McpBridgeError( @@ -712,7 +724,7 @@ class McpBridge: description=description if isinstance(description, str) else remote_name, parameters=schema, permission=permission, - source="plugin", + source=tool_source, ), ) diff --git a/backend/app/extensions/mcp_registry.py b/backend/app/extensions/mcp_registry.py new file mode 100644 index 0000000..e576f7b --- /dev/null +++ b/backend/app/extensions/mcp_registry.py @@ -0,0 +1,554 @@ +"""Independent, user-managed MCP server registry for development builds.""" + +from __future__ import annotations + +import hashlib +import json +import re +import threading +from datetime import UTC, datetime +from pathlib import Path +from typing import Any +from uuid import uuid4 + +from pydantic import BaseModel, ConfigDict, create_model + +from app.agent.permissions import KNOWN_PERMISSIONS +from app.agent.tools import ToolExecutionContext, ToolRegistry +from app.contracts import ( + McpServer, + McpServerCreateRequest, + McpServerSecretStatus, + McpServerTransport, + McpServerUpdateRequest, + PluginBackend, + PluginHostState, +) +from app.extensions.mcp import McpBridge, McpBridgeError, McpDiscoveredTool +from app.providers.credentials import CredentialStoreError, EncryptedCredentialStore + +_ENVIRONMENT_KEY = re.compile(r"^[A-Za-z_][A-Za-z0-9_]{0,127}$") + + +class McpRegistryError(RuntimeError): + def __init__(self, code: str, message: str, *, status_code: int = 422) -> None: + super().__init__(message) + self.code = code + self.message = message + self.status_code = status_code + + +class McpServerRegistry: + """Persists configuration and owns stdio host/tool lifecycles.""" + + def __init__( + self, + registry: ToolRegistry, + credentials: EncryptedCredentialStore, + data_dir: Path, + *, + allow_process_launch: bool, + bridge: McpBridge | None = None, + ) -> None: + self.tools = registry + self.credentials = credentials + self.data_dir = data_dir + self.allow_process_launch = allow_process_launch + self.bridge = bridge or McpBridge() + self._lock = threading.RLock() + self._records = self._read() + self._registered: dict[str, list[str]] = {} + self._last_status: dict[str, dict[str, Any]] = {} + + def list(self) -> list[McpServer]: + with self._lock: + return [ + self._public(server_id, record) + for server_id, record in self._records.items() + ] + + def get(self, server_id: str) -> McpServer: + with self._lock: + return self._public(server_id, self._record(server_id)) + + def create(self, request: McpServerCreateRequest) -> McpServer: + self._validate(request) + server_id = uuid4().hex[:12] + record = request.model_dump(mode="json") + record["name"] = request.name.strip() + record["command"] = request.command.strip() + record.update(enabled=False, approved_digest=None) + with self._lock: + updated = {**self._records, server_id: record} + self._write(updated) + self._records = updated + return self.get(server_id) + + def update(self, server_id: str, request: McpServerUpdateRequest) -> McpServer: + self._validate(request) + self.disable(server_id) + with self._lock: + previous = self._record(server_id) + removed = set(previous.get("secret_environment_keys", [])) - set( + request.secret_environment_keys + ) + record = request.model_dump(mode="json") + record["name"] = request.name.strip() + record["command"] = request.command.strip() + record.update(enabled=False, approved_digest=None) + updated = {**self._records, server_id: record} + self._write(updated) + self._records = updated + self._last_status.pop(server_id, None) + for key in removed: + try: + self.credentials.delete(self._secret_id(server_id, key)) + except CredentialStoreError as exc: + raise McpRegistryError( + "MCP_SECRET_STORE_ERROR", str(exc), status_code=500 + ) from exc + return self.get(server_id) + + def delete(self, server_id: str) -> None: + self.disable(server_id) + with self._lock: + record = self._record(server_id) + secret_ids = [ + self._secret_id(server_id, key) + for key in record.get("secret_environment_keys", []) + ] + updated = dict(self._records) + del updated[server_id] + self._write(updated) + self._records = updated + self._last_status.pop(server_id, None) + try: + self.credentials.delete_many(secret_ids) + except CredentialStoreError as exc: + raise McpRegistryError( + "MCP_SECRET_STORE_ERROR", str(exc), status_code=500 + ) from exc + self.bridge.remove(self._host_id(server_id)) + + def trust(self, server_id: str, command_digest: str) -> McpServer: + with self._lock: + record = self._record(server_id) + current = self._digest(record) + if command_digest != current: + raise McpRegistryError( + "MCP_TRUST_DIGEST_STALE", + "MCP server configuration changed; review it again.", + status_code=409, + ) + approved = {**record, "approved_digest": current} + updated = {**self._records, server_id: approved} + self._write(updated) + self._records = updated + return self.get(server_id) + + def put_secret( + self, server_id: str, key: str, secret: str + ) -> McpServerSecretStatus: + with self._lock: + record = self._record(server_id) + self._validate_environment_key(key) + if key not in record.get("secret_environment_keys", []): + raise McpRegistryError( + "MCP_SECRET_NOT_DECLARED", + "Secret environment key is not declared in this server configuration.", + ) + try: + self.credentials.put(self._secret_id(server_id, key), secret) + except CredentialStoreError as exc: + raise McpRegistryError( + "MCP_SECRET_STORE_ERROR", str(exc), status_code=500 + ) from exc + return McpServerSecretStatus(key=key, configured=True) + + def delete_secret(self, server_id: str, key: str) -> McpServerSecretStatus: + record = self._record(server_id) + if key not in record.get("secret_environment_keys", []): + raise McpRegistryError( + "MCP_SECRET_NOT_DECLARED", + "Secret environment key is not declared in this server configuration.", + ) + try: + self.credentials.delete(self._secret_id(server_id, key)) + except CredentialStoreError as exc: + raise McpRegistryError( + "MCP_SECRET_STORE_ERROR", str(exc), status_code=500 + ) from exc + return McpServerSecretStatus(key=key, configured=False) + + def test(self, server_id: str) -> McpServer: + record = self._record(server_id) + if record.get("enabled"): + raise McpRegistryError( + "MCP_SERVER_ALREADY_ENABLED", + "Disable the MCP server before running an isolated connection test.", + status_code=409, + ) + self._require_launch_allowed(record) + try: + discovered = self._start(server_id, record) + except Exception as exc: + self._last_status[server_id] = { + "status": PluginHostState.error, + "error": str(exc), + "last_tested_at": datetime.now(UTC), + "last_test_succeeded": False, + } + raise + status = self.bridge.status(self._host_id(server_id), self._backend(record)) + self._last_status[server_id] = { + "status": PluginHostState.stopped, + "tools_count": len(discovered), + "protocol_version": status.protocol_version, + "remote_server_name": status.server_name, + "remote_server_version": status.server_version, + "error": None, + "last_tested_at": datetime.now(UTC), + "last_test_succeeded": True, + } + self.bridge.stop(self._host_id(server_id)) + return self.get(server_id) + + def enable(self, server_id: str) -> McpServer: + record = self._record(server_id) + if server_id in self._registered: + return self.get(server_id) + self._require_launch_allowed(record) + discovered = self._start(server_id, record) + registered: list[str] = [] + try: + for item in discovered: + self._register(server_id, item) + registered.append(item.definition.name) + except Exception: + for name in registered: + self.tools.unregister(name) + self.bridge.stop(self._host_id(server_id)) + raise + try: + with self._lock: + enabled_record = {**record, "enabled": True} + updated = {**self._records, server_id: enabled_record} + self._write(updated) + self._records = updated + self._registered[server_id] = registered + except McpRegistryError: + for name in registered: + self.tools.unregister(name) + self.bridge.stop(self._host_id(server_id)) + raise + return self.get(server_id) + + def disable(self, server_id: str) -> McpServer: + with self._lock: + record = self._record(server_id) + disabled_record = {**record, "enabled": False} + updated = {**self._records, server_id: disabled_record} + self._write(updated) + self._records = updated + for name in self._registered.pop(server_id, []): + self.tools.unregister(name) + self.bridge.stop(self._host_id(server_id)) + return self.get(server_id) + + def restore_enabled(self) -> None: + if not self._records: + return + for server_id, record in list(self._records.items()): + if record.get("enabled"): + try: + self.enable(server_id) + except (McpRegistryError, ValueError, OSError) as exc: + self._records[server_id] = {**record, "enabled": False} + self._last_status[server_id] = { + "status": PluginHostState.error, + "error": str(exc), + } + self._write() + + def shutdown(self) -> None: + for server_id in list(self._records): + for name in self._registered.pop(server_id, []): + self.tools.unregister(name) + self.bridge.stop(self._host_id(server_id)) + + def _start(self, server_id: str, record: dict[str, Any]) -> list[McpDiscoveredTool]: + environment = dict(record.get("environment", {})) + for key in record.get("secret_environment_keys", []): + try: + value = self.credentials.resolve(self._secret_id(server_id, key)) + except CredentialStoreError as exc: + raise McpRegistryError( + "MCP_SECRET_STORE_ERROR", str(exc), status_code=500 + ) from exc + if value is None: + raise McpRegistryError( + "MCP_SECRET_REQUIRED", + f"Secret environment variable is not configured: {key}", + status_code=409, + ) + environment[key] = value + host_id = self._host_id(server_id) + self.bridge.remove(host_id) + try: + return self.bridge.start( + host_id, + self._backend(record), + self._server_dir(server_id), + list(record.get("permissions", [])), + lambda _host, message: self._unavailable(server_id, message), + command_override=[record["command"], *record.get("args", [])], + environment=environment, + tool_source="mcp_server", + ) + except McpBridgeError as exc: + raise McpRegistryError( + exc.code, exc.message, status_code=exc.status_code + ) from exc + + def _register(self, server_id: str, discovered: McpDiscoveredTool) -> None: + definition = discovered.definition + model_name = "McpArgs_" + re.sub(r"\W+", "_", definition.name) + arguments_model = create_model(model_name, __config__=ConfigDict(extra="allow")) + + async def executor(arguments: BaseModel, context: ToolExecutionContext) -> Any: + return await self.bridge.call_tool( + self._host_id(server_id), + discovered.remote_name, + arguments.model_dump(exclude_unset=True), + request_id=context.tool_call_id + or f"{context.run_id}:{definition.name}", + ) + + self.tools.register(definition, arguments_model, executor) + + def _unavailable(self, server_id: str, message: str) -> None: + with self._lock: + for name in self._registered.pop(server_id, []): + self.tools.unregister(name) + record = self._records.get(server_id) + if record is not None: + self._records[server_id] = {**record, "enabled": False} + self._last_status[server_id] = { + "status": PluginHostState.unhealthy, + "error": message, + } + self._write() + + def _require_launch_allowed(self, record: dict[str, Any]) -> None: + if record.get("transport") != McpServerTransport.stdio.value: + raise McpRegistryError( + "MCP_TRANSPORT_UNSUPPORTED", + "C.1 currently supports stdio; Streamable HTTP and SSE are reserved for a later increment.", + status_code=501, + ) + if not self.allow_process_launch: + raise McpRegistryError( + "MCP_SANDBOX_REQUIRED", + "Python process launch is disabled outside development until the desktop sandbox is available.", + status_code=403, + ) + if record.get("approved_digest") != self._digest(record): + raise McpRegistryError( + "MCP_TRUST_APPROVAL_REQUIRED", + "Review and approve the current MCP command before testing or enabling it.", + status_code=409, + ) + + def _public(self, server_id: str, record: dict[str, Any]) -> McpServer: + digest = self._digest(record) + backend = self._backend(record) + status = self.bridge.status(self._host_id(server_id), backend) + cached = self._last_status.get(server_id, {}) + return McpServer( + server_id=server_id, + name=record["name"], + transport=record["transport"], + command=record["command"], + args=list(record.get("args", [])), + environment=dict(record.get("environment", {})), + secret_environment={ + key: self._secret_configured(server_id, key) + for key in record.get("secret_environment_keys", []) + }, + permissions=list(record.get("permissions", [])), + startup_timeout_seconds=backend.startup_timeout_seconds, + tool_timeout_seconds=backend.tool_timeout_seconds, + enabled=bool(record.get("enabled")), + trusted=record.get("approved_digest") == digest, + command_digest=digest, + command_summary=self._summary(record), + status=status.status + if record.get("enabled") + else cached.get("status", PluginHostState.stopped), + tools_count=status.tools_count + if record.get("enabled") + else cached.get("tools_count", 0), + protocol_version=status.protocol_version + if record.get("enabled") + else cached.get("protocol_version"), + remote_server_name=status.server_name + if record.get("enabled") + else cached.get("remote_server_name"), + remote_server_version=status.server_version + if record.get("enabled") + else cached.get("remote_server_version"), + error=status.error if record.get("enabled") else cached.get("error"), + last_tested_at=cached.get("last_tested_at"), + last_test_succeeded=cached.get("last_test_succeeded"), + ) + + def _validate(self, request: McpServerCreateRequest) -> None: + if not request.name.strip(): + raise McpRegistryError( + "MCP_SERVER_NAME_INVALID", "MCP server name cannot be blank." + ) + if not request.command.strip() or "\x00" in request.command: + raise McpRegistryError("MCP_COMMAND_INVALID", "MCP executable is invalid.") + if any("\x00" in arg for arg in request.args): + raise McpRegistryError( + "MCP_COMMAND_INVALID", "MCP argument contains a null byte." + ) + for key in [*request.environment, *request.secret_environment_keys]: + self._validate_environment_key(key) + if set(request.environment) & set(request.secret_environment_keys): + raise McpRegistryError( + "MCP_ENVIRONMENT_INVALID", + "An environment key cannot be both plain and secret.", + ) + unknown_permissions = set(request.permissions) - KNOWN_PERMISSIONS + if unknown_permissions: + raise McpRegistryError( + "MCP_PERMISSION_INVALID", + f"Unknown MCP permission: {min(unknown_permissions)}", + ) + + @staticmethod + def _validate_environment_key(key: str) -> None: + if not _ENVIRONMENT_KEY.fullmatch(key): + raise McpRegistryError( + "MCP_ENVIRONMENT_INVALID", f"Invalid environment variable name: {key}" + ) + + @staticmethod + def _backend(record: dict[str, Any]) -> PluginBackend: + return PluginBackend( + type="mcp", + transport="stdio", + command=record["command"], + args=record.get("args", []), + startup_timeout_seconds=record.get("startup_timeout_seconds", 15), + tool_timeout_seconds=record.get("tool_timeout_seconds", 30), + ) + + @staticmethod + def _host_id(server_id: str) -> str: + return f"mcp.{server_id}" + + def _server_dir(self, server_id: str) -> Path: + path = self.data_dir / "mcp" / "workdirs" / server_id + path.mkdir(parents=True, exist_ok=True) + return path + + @staticmethod + def _digest(record: dict[str, Any]) -> str: + executable = { + key: record.get(key) + for key in ( + "transport", + "command", + "args", + "environment", + "secret_environment_keys", + "permissions", + ) + } + return hashlib.sha256( + json.dumps( + executable, sort_keys=True, ensure_ascii=False, separators=(",", ":") + ).encode() + ).hexdigest() + + @staticmethod + def _summary(record: dict[str, Any]) -> str: + return " ".join( + [ + record["command"], + *[ + json.dumps(arg, ensure_ascii=False) + for arg in record.get("args", []) + ], + ] + ) + + @staticmethod + def _secret_id(server_id: str, key: str) -> str: + suffix = hashlib.sha256(key.encode()).hexdigest()[:20] + return f"mcp.{server_id}.{suffix}" + + def _secret_configured(self, server_id: str, key: str) -> bool: + try: + return self.credentials.has(self._secret_id(server_id, key)) + except CredentialStoreError as exc: + raise McpRegistryError( + "MCP_SECRET_STORE_ERROR", str(exc), status_code=500 + ) from exc + + def _record(self, server_id: str) -> dict[str, Any]: + try: + return self._records[server_id] + except KeyError as exc: + raise McpRegistryError( + "MCP_SERVER_NOT_FOUND", + "MCP server configuration was not found.", + status_code=404, + ) from exc + + @property + def _path(self) -> Path: + return self.data_dir / "mcp" / "servers.json" + + def _read(self) -> dict[str, dict[str, Any]]: + if not self._path.exists(): + return {} + try: + value = json.loads(self._path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as exc: + raise McpRegistryError( + "MCP_REGISTRY_INVALID", + "MCP server registry cannot be loaded.", + status_code=500, + ) from exc + if not isinstance(value, dict): + raise McpRegistryError( + "MCP_REGISTRY_INVALID", + "MCP server registry has an invalid format.", + status_code=500, + ) + return value + + def _write(self, records: dict[str, dict[str, Any]] | None = None) -> None: + temporary = self._path.with_suffix(".tmp") + try: + self._path.parent.mkdir(parents=True, exist_ok=True) + temporary.write_text( + json.dumps( + records if records is not None else self._records, + ensure_ascii=False, + indent=2, + sort_keys=True, + ), + encoding="utf-8", + ) + temporary.replace(self._path) + except OSError as exc: + temporary.unlink(missing_ok=True) + raise McpRegistryError( + "MCP_REGISTRY_WRITE_FAILED", + "MCP server registry cannot be written.", + status_code=500, + ) from exc diff --git a/backend/app/main.py b/backend/app/main.py index 8a3aae5..bfaaf91 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -19,6 +19,7 @@ async def lifespan(_: FastAPI): yield # 第三方 MCP Server 必须跟随 AI Core 退出,不能遗留孤儿进程。 container.plugins.shutdown() + container.mcp_servers.shutdown() app = FastAPI( diff --git a/backend/app/providers/credentials.py b/backend/app/providers/credentials.py index 88d9485..accdacb 100644 --- a/backend/app/providers/credentials.py +++ b/backend/app/providers/credentials.py @@ -14,6 +14,7 @@ from app.config import get_settings _CREDENTIAL_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$") _PLUGIN_CREDENTIAL_PREFIX = "plugin." +_MCP_CREDENTIAL_PREFIX = "mcp." class CredentialStoreError(RuntimeError): @@ -27,10 +28,10 @@ class CredentialResolver(Protocol): def validate_provider_credential_id(credential_id: str | None) -> None: """阻止 Provider 和通用凭据 API 跨入 Plugin 私有命名空间。""" - if credential_id and credential_id.casefold().startswith( - _PLUGIN_CREDENTIAL_PREFIX - ): + if credential_id and credential_id.casefold().startswith(_PLUGIN_CREDENTIAL_PREFIX): raise CredentialStoreError("Credential namespace is reserved for Plugin settings.") + if credential_id and credential_id.casefold().startswith(_MCP_CREDENTIAL_PREFIX): + raise CredentialStoreError("Credential namespace is reserved for MCP settings.") class EnvironmentCredentialResolver: diff --git a/backend/app/routes.py b/backend/app/routes.py index 962bec7..438d41b 100644 --- a/backend/app/routes.py +++ b/backend/app/routes.py @@ -21,6 +21,13 @@ from app.contracts import ( IndexJob, IndexRebuildRequest, IndexStatus, + McpServer, + McpServerCreateRequest, + McpServerListResponse, + McpServerSecretStatus, + McpServerSecretWriteRequest, + McpServerTrustRequest, + McpServerUpdateRequest, ModelEvent, ModelEventType, Note, @@ -72,6 +79,7 @@ from app.agent import AgentCapacityError, AgentRunNotFoundError from app.container import container from app.errors import ApiError from app.extensions import ExtensionError +from app.extensions.mcp_registry import McpRegistryError from app.providers.registry import ProviderNotFoundError from app.providers.factory import UnsupportedProviderError from app.providers.base import ProviderError @@ -91,6 +99,21 @@ from app.services import ( router = APIRouter(prefix="/api") +def mcp_call(operation): + try: + return operation() + except McpRegistryError as exc: + raise ApiError(exc.status_code, exc.code, exc.message) from exc + + +async def mcp_call_async(operation): + """MCP process operations wait on stdio and must not block the API event loop.""" + try: + return await asyncio.to_thread(operation) + except McpRegistryError as exc: + raise ApiError(exc.status_code, exc.code, exc.message) from exc + + def utc_now() -> datetime: return datetime.now(timezone.utc) @@ -481,6 +504,63 @@ async def uninstall_skill(skill_id: str) -> OperationResponse: return OperationResponse(status="completed", resource_id=skill_id, message="uninstalled") +# Independent MCP Server Registry +@router.get("/mcp/servers", response_model=McpServerListResponse, tags=["MCP Servers"]) +async def list_mcp_servers() -> McpServerListResponse: + return McpServerListResponse(items=mcp_call(container.mcp_servers.list)) + + +@router.post("/mcp/servers", response_model=McpServer, status_code=201, tags=["MCP Servers"]) +async def create_mcp_server(request: McpServerCreateRequest) -> McpServer: + return mcp_call(lambda: container.mcp_servers.create(request)) + + +@router.get("/mcp/servers/{server_id}", response_model=McpServer, tags=["MCP Servers"]) +async def get_mcp_server(server_id: str) -> McpServer: + return mcp_call(lambda: container.mcp_servers.get(server_id)) + + +@router.put("/mcp/servers/{server_id}", response_model=McpServer, tags=["MCP Servers"]) +async def update_mcp_server(server_id: str, request: McpServerUpdateRequest) -> McpServer: + return await mcp_call_async(lambda: container.mcp_servers.update(server_id, request)) + + +@router.delete("/mcp/servers/{server_id}", response_model=OperationResponse, tags=["MCP Servers"]) +async def delete_mcp_server(server_id: str) -> OperationResponse: + await mcp_call_async(lambda: container.mcp_servers.delete(server_id)) + return OperationResponse(status="completed", resource_id=server_id, message="deleted") + + +@router.post("/mcp/servers/{server_id}/trust", response_model=McpServer, tags=["MCP Servers"]) +async def trust_mcp_server(server_id: str, request: McpServerTrustRequest) -> McpServer: + return mcp_call(lambda: container.mcp_servers.trust(server_id, request.command_digest)) + + +@router.post("/mcp/servers/{server_id}/test", response_model=McpServer, tags=["MCP Servers"]) +async def test_mcp_server(server_id: str) -> McpServer: + return await mcp_call_async(lambda: container.mcp_servers.test(server_id)) + + +@router.post("/mcp/servers/{server_id}/enable", response_model=McpServer, tags=["MCP Servers"]) +async def enable_mcp_server(server_id: str) -> McpServer: + return await mcp_call_async(lambda: container.mcp_servers.enable(server_id)) + + +@router.post("/mcp/servers/{server_id}/disable", response_model=McpServer, tags=["MCP Servers"]) +async def disable_mcp_server(server_id: str) -> McpServer: + return await mcp_call_async(lambda: container.mcp_servers.disable(server_id)) + + +@router.put("/mcp/servers/{server_id}/secrets/{key}", response_model=McpServerSecretStatus, tags=["MCP Servers"]) +async def put_mcp_server_secret(server_id: str, key: str, request: McpServerSecretWriteRequest) -> McpServerSecretStatus: + return mcp_call(lambda: container.mcp_servers.put_secret(server_id, key, request.secret.get_secret_value())) + + +@router.delete("/mcp/servers/{server_id}/secrets/{key}", response_model=McpServerSecretStatus, tags=["MCP Servers"]) +async def delete_mcp_server_secret(server_id: str, key: str) -> McpServerSecretStatus: + return mcp_call(lambda: container.mcp_servers.delete_secret(server_id, key)) + + # Plugins @router.get("/plugins", response_model=PluginListResponse, tags=["Plugins"]) async def list_plugins() -> PluginListResponse: diff --git a/backend/tests/test_mcp_registry.py b/backend/tests/test_mcp_registry.py new file mode 100644 index 0000000..a4eb658 --- /dev/null +++ b/backend/tests/test_mcp_registry.py @@ -0,0 +1,123 @@ +import sys + +import pytest + +from app.agent.tools import ToolRegistry +from app.config import BACKEND_DIR, get_settings +from app.contracts import McpServerCreateRequest, McpServerUpdateRequest +from app.extensions.mcp_registry import McpRegistryError, McpServerRegistry +from app.providers.credentials import EncryptedCredentialStore + +SERVER = BACKEND_DIR / "extensions" / "fixtures" / "mcp-echo" / "server.py" + + +def request(**overrides) -> McpServerCreateRequest: + values = { + "name": "Echo MCP", + "command": sys.executable, + "args": [str(SERVER)], + "permissions": ["notes.read", "secrets.use"], + "secret_environment_keys": ["TEST_MCP_SECRET"], + } + values.update(overrides) + return McpServerCreateRequest(**values) + + +def registry(*, launch: bool = True) -> McpServerRegistry: + return McpServerRegistry( + ToolRegistry(), + EncryptedCredentialStore(), + get_settings().data_dir, + allow_process_launch=launch, + ) + + +def test_registry_requires_current_trust_and_never_returns_secret() -> None: + service = registry() + created = service.create(request()) + assert created.trusted is False + assert created.secret_environment == {"TEST_MCP_SECRET": False} + + service.put_secret(created.server_id, "TEST_MCP_SECRET", "do-not-return") + configured = service.get(created.server_id) + assert configured.secret_environment == {"TEST_MCP_SECRET": True} + assert "do-not-return" not in configured.model_dump_json() + + with pytest.raises(McpRegistryError, match="approve"): + service.test(created.server_id) + + service.trust(created.server_id, created.command_digest) + tested = service.test(created.server_id) + assert tested.status == "stopped" + assert tested.last_test_succeeded is True + assert tested.tools_count > 0 + service.shutdown() + + +def test_update_disables_server_and_revokes_command_trust() -> None: + service = registry() + created = service.create(request(secret_environment_keys=[])) + service.trust(created.server_id, created.command_digest) + enabled = service.enable(created.server_id) + assert enabled.enabled is True + assert any( + item.name.startswith(f"mcp.{created.server_id}.") + for item in service.tools.definitions() + ) + + updated = service.update( + created.server_id, + McpServerUpdateRequest( + **request(name="Changed", secret_environment_keys=[]).model_dump() + ), + ) + assert updated.enabled is False + assert updated.trusted is False + assert not any( + item.name.startswith(f"mcp.{created.server_id}.") + for item in service.tools.definitions() + ) + service.shutdown() + + +def test_production_rejects_process_launch_even_after_approval() -> None: + service = registry(launch=False) + created = service.create(request(secret_environment_keys=[])) + service.trust(created.server_id, created.command_digest) + with pytest.raises(McpRegistryError) as error: + service.enable(created.server_id) + assert error.value.code == "MCP_SANDBOX_REQUIRED" + + +def test_non_stdio_transport_is_explicitly_reserved() -> None: + service = registry() + created = service.create( + request( + transport="streamable_http", + command="https://example.invalid/mcp", + secret_environment_keys=[], + ) + ) + service.trust(created.server_id, created.command_digest) + with pytest.raises(McpRegistryError) as error: + service.test(created.server_id) + assert error.value.code == "MCP_TRANSPORT_UNSUPPORTED" + + +def test_enabled_server_is_restored_from_persisted_registry() -> None: + first = registry() + created = first.create(request(secret_environment_keys=[])) + first.trust(created.server_id, created.command_digest) + first.enable(created.server_id) + first.shutdown() + + restored = registry() + restored.restore_enabled() + current = restored.get(created.server_id) + assert current.enabled is True + assert current.status == "ready" + assert any( + item.name.startswith(f"mcp.{created.server_id}.") + for item in restored.tools.definitions() + ) + restored.shutdown() diff --git a/docs/README.md b/docs/README.md index de9938b..eb9ff2e 100644 --- a/docs/README.md +++ b/docs/README.md @@ -32,6 +32,7 @@ - [Knowledge 与 Retrieval Core 开发说明](development/Knowledge与Retrieval-Core开发说明.md) - [模型提供商与模型发现开发说明](development/模型提供商与模型发现开发说明.md) - [MCP Bridge 与 Plugin Host 开发说明](development/MCP-Bridge与Plugin-Host开发说明.md) +- [独立 MCP Server 配置中心开发说明](development/独立MCP-Server配置中心开发说明.md) - [Plugin Command 与 Settings 开发说明](development/Plugin-Command与Settings开发说明.md) - [前端壳子与接口层开发说明](development/前端壳子与接口层开发说明.md) - [前端写作体验优化开发说明](development/前端写作体验优化开发说明.md) diff --git a/docs/contracts/第二阶段接口契约-开发版.md b/docs/contracts/第二阶段接口契约-开发版.md index 87aceb4..a600282 100644 --- a/docs/contracts/第二阶段接口契约-开发版.md +++ b/docs/contracts/第二阶段接口契约-开发版.md @@ -47,6 +47,12 @@ | Agent Trace | GET | `/api/agent/runs/{run_id}/trace` | 已实现 | 分页读取可回放 Trace 快照 | | Plugin Host | GET | `/api/plugins/{plugin_id}/host` | 已实现 | 获取 MCP Host 健康状态 | | Plugin Host | POST | `/api/plugins/{plugin_id}/host/restart` | 已实现 | 重启异常 Host 并重新发现 Tool | +| MCP Server | GET/POST | `/api/mcp/servers` | 已实现(C.1/P0) | 列出、创建独立 MCP Server 配置 | +| MCP Server | GET/PUT/DELETE | `/api/mcp/servers/{server_id}` | 已实现(C.1/P0) | 读取、修改、删除独立配置 | +| MCP Server | POST | `/api/mcp/servers/{server_id}/trust` | 已实现(C.1/P0) | 确认当前可执行配置摘要 | +| MCP Server | POST | `/api/mcp/servers/{server_id}/test` | 已实现(C.1/P0) | 隔离启动、握手、发现工具后退出 | +| MCP Server | POST | `/api/mcp/servers/{server_id}/enable`、`disable` | 已实现(C.1/P0) | 控制 Host 与动态 Tool 生命周期 | +| MCP Server | PUT/DELETE | `/api/mcp/servers/{server_id}/secrets/{key}` | 已实现(C.1/P0) | 写入或删除加密环境变量 | | Plugin Command | GET | `/api/plugin-contributions/commands` | 已实现 | 获取前端可展示的 Command | | Plugin Command | POST | `/api/plugin-contributions/commands/{command_id}/execute` | 已实现 | 受控执行 Command | | Plugin Settings | GET | `/api/plugins/{plugin_id}/settings` | 已实现 | 获取 Schema 与非敏感配置 | @@ -647,6 +653,18 @@ MCP_TOOL_SCHEMA_INVALID MCP_TOOL_CALL_FAILED MCP_TOOL_RESULT_TOO_LARGE MCP_TRUST_APPROVAL_REQUIRED +MCP_TRUST_DIGEST_STALE +MCP_SANDBOX_REQUIRED +MCP_TRANSPORT_UNSUPPORTED +MCP_SERVER_NOT_FOUND +MCP_SERVER_NAME_INVALID +MCP_SERVER_ALREADY_ENABLED +MCP_REGISTRY_WRITE_FAILED +MCP_SECRET_REQUIRED +MCP_SECRET_NOT_DECLARED +MCP_SECRET_STORE_ERROR +MCP_ENVIRONMENT_INVALID +MCP_PERMISSION_INVALID PLUGIN_COMMAND_NOT_FOUND PLUGIN_COMMAND_CONFLICT PLUGIN_COMMAND_INVALID @@ -670,6 +688,14 @@ PLUGIN_STORAGE_ERROR CREDENTIAL_NAMESPACE_RESERVED ``` +### 7.8 独立 MCP Server Registry(C.1) + +独立 Server 不依附 Plugin Manifest,配置持久化于 `APP_DATA_DIR/mcp/servers.json`,敏感环境变量只以 `mcp.*` 引用进入加密凭据存储。响应仅返回每个 Secret 是否已配置,不返回明文。动态工具使用 `mcp.{server_id}.{remote_tool}` 命名空间,来源标记为 `mcp_server`,仍通过统一 Tool Registry、Permission Manager 与 Agent Trace。 + +P0 只真实支持 `stdio`。`streamable_http` 与 `sse` 已作为后续 Contract 枚举保留,但测试或启用会返回 `501 MCP_TRANSPORT_UNSUPPORTED`,前端不可伪装为可用。命令始终以 executable 与 args 数组通过 `shell=False` 启动;普通环境变量和加密 Secret 显式注入,不继承 Provider Key、数据库或 Vault 路径。 + +创建或编辑配置后,调用方必须向 `/trust` 回传服务端计算的 `command_digest`。后端只接受与当前 transport、command、args、环境变量键值及权限完全一致的摘要;配置变化会撤销旧信任。测试连接同样会实际启动进程,因此也要求确认。当前 Python Host 仅在 `APP_ENVIRONMENT=development` 时允许启动;其他环境返回 `403 MCP_SANDBOX_REQUIRED`,等待第三阶段桌面端沙箱接管。 + --- ## 8. Provider Adapter 扩展 diff --git a/docs/development/独立MCP-Server配置中心开发说明.md b/docs/development/独立MCP-Server配置中心开发说明.md new file mode 100644 index 0000000..26d68f0 --- /dev/null +++ b/docs/development/独立MCP-Server配置中心开发说明.md @@ -0,0 +1,47 @@ +# 独立 MCP Server 配置中心开发说明 + +> 更新日期:2026-09-03。本文记录第二阶段 C.1 的 P0 实现;它与 Plugin 自带的 MCP Host 是两个并列入口。 + +## 1. 已实现范围 + +- 独立 Server 的创建、读取、编辑和删除; +- stdio 命令、参数、普通环境变量、加密环境变量及超时配置; +- 命令摘要确认、测试连接、启用、停用和异常状态展示; +- MCP initialize、`tools/list` 与动态 Tool 注册,工具命名为 `mcp.{server_id}.{tool}`; +- 启用状态持久化与开发服务重启恢复; +- 前端独立“MCP”导航与配置弹窗,提供 stdio/uvx 模板; +- Streamable HTTP 与旧 SSE 仅作为后续选项展示为禁用,不属于本次完成范围。 + +## 2. 数据与 Secret + +普通配置原子写入 `APP_DATA_DIR/mcp/servers.json`。Secret 使用 `mcp.{server_id}.{key_hash}` 作为内部引用写入现有 Fernet 凭据存储;API 和前端只看到 `configured: true/false`。删除 Server 或移除 Secret 键时同步清理密文。 + +前端 Secret 输入使用密码框,提交后立即清空,不写入 localStorage、普通配置 JSON 或日志。`mcp.*` 同 `plugin.*` 一样属于保留凭据命名空间,Provider 配置与通用凭据 API 无权读取。 + +## 3. 启动安全边界 + +命令不经过 Shell,管道、重定向和拼接字符串不会被解释。Host 只继承启动所需的系统变量,再叠加用户显式配置;`uvx` 可隔离 Python 依赖,但不能限制文件、网络和系统调用。 + +后端会对影响执行的配置计算 SHA-256 摘要。测试或启用前,用户必须确认并回传当前摘要;修改配置会立即撤销旧确认。由于 Python Host 尚无 OS 沙箱,非开发环境硬拒绝启动。第三阶段由 Tauri/Rust Host 提供平台级隔离后再替换这道临时门禁。 + +## 4. 测试与启动 + +```powershell +cd backend +uv run pytest -q tests/test_mcp_registry.py +uv run uvicorn app.main:app --reload + +cd ../frontend +npm run type-check +npm test +npm run dev +``` + +打开知识库后进入左侧“MCP”。保存配置,按提示确认命令,先执行“测试连接”;成功后再启用。默认模板 `uvx mcp-server-fetch` 仅为配置示例,首次下载是否联网由本机 uv 缓存与网络环境决定。 + +## 5. 后续增量 + +- P1:Streamable HTTP 连接、认证 Header 与重连策略; +- 兼容项:仅在确有旧服务需求时增加 SSE; +- 第三阶段前:把命令确认与进程创建迁移至 Tauri/Rust 沙箱; +- 增加面向真实第三方 Server 的兼容矩阵,不用单一 Fixture 代表协议全兼容。 diff --git a/frontend/src/components/common/PrimarySidebar.vue b/frontend/src/components/common/PrimarySidebar.vue index 052f0b4..9561166 100644 --- a/frontend/src/components/common/PrimarySidebar.vue +++ b/frontend/src/components/common/PrimarySidebar.vue @@ -1,7 +1,7 @@ + + + + diff --git a/frontend/src/router/index.ts b/frontend/src/router/index.ts index 2116443..4655ccb 100644 --- a/frontend/src/router/index.ts +++ b/frontend/src/router/index.ts @@ -44,6 +44,12 @@ const routes = [ component: () => import('@/features/skills/SkillsView.vue'), meta: { title: 'Skill 管理', requiresVault: true }, }, + { + path: '/extensions/mcp', + name: 'mcp-servers', + component: () => import('@/features/mcp/McpServersView.vue'), + meta: { title: 'MCP 服务器', requiresVault: true }, + }, { path: '/extensions/plugins', name: 'plugins', diff --git a/frontend/src/services/index.ts b/frontend/src/services/index.ts index 968f5b3..849e204 100644 --- a/frontend/src/services/index.ts +++ b/frontend/src/services/index.ts @@ -8,6 +8,7 @@ export * as chatService from './chatService' export * as agentService from './agentService' export * as skillService from './skillService' export * as pluginService from './pluginService' +export * as mcpServerService from './mcpServerService' export * as providerService from './providerService' export * as taskService from './taskService' export * as indexService from './indexService' diff --git a/frontend/src/services/mcpServerService.ts b/frontend/src/services/mcpServerService.ts new file mode 100644 index 0000000..a146a6b --- /dev/null +++ b/frontend/src/services/mcpServerService.ts @@ -0,0 +1,17 @@ +import apiClient from './apiClient' +import type { McpServer, McpServerInput, OperationResponse } from '@/contracts' + +const base = '/api/mcp/servers' + +export async function listMcpServers(): Promise { + return (await apiClient.get<{ items: McpServer[] }>(base)).items +} +export const createMcpServer = (input: McpServerInput) => apiClient.post(base, input) +export const updateMcpServer = (id: string, input: McpServerInput) => apiClient.put(`${base}/${id}`, input) +export const deleteMcpServer = (id: string) => apiClient.delete(`${base}/${id}`) +export const trustMcpServer = (server: McpServer) => apiClient.post(`${base}/${server.server_id}/trust`, { command_digest: server.command_digest }) +export const testMcpServer = (id: string) => apiClient.post(`${base}/${id}/test`) +export const enableMcpServer = (id: string) => apiClient.post(`${base}/${id}/enable`) +export const disableMcpServer = (id: string) => apiClient.post(`${base}/${id}/disable`) +export const putMcpServerSecret = (id: string, key: string, secret: string) => apiClient.put(`${base}/${id}/secrets/${encodeURIComponent(key)}`, { secret }) +export const deleteMcpServerSecret = (id: string, key: string) => apiClient.delete(`${base}/${id}/secrets/${encodeURIComponent(key)}`)