/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/api/mcp_servers.py
516 строк
18 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""API endpoints для управления MCP (Model Context Protocol) серверами. Использует MCPServerService + MCPManager для реального lifecycle (запуск/остановка subprocess, discovery tools, вызов инструментов). Эндпоинты: - GET /api/v1/mcp-servers — список серверов - POST /api/v1/mcp-servers — создать сервер - GET /api/v1/mcp-servers/stats — статистика - GET /api/v1/mcp-servers/{id} — получить сервер - PATCH /api/v1/mcp-servers/{id} — обновить сервер - DELETE /api/v1/mcp-servers/{id} — удалить сервер - POST /api/v1/mcp-servers/{id}/start — запустить - POST /api/v1/mcp-servers/{id}/stop — остановить - POST /api/v1/mcp-servers/{id}/restart — перезапустить - POST /api/v1/mcp-servers/{id}/test — тест соединения - GET /api/v1/mcp-servers/{id}/tools — список tools - GET /api/v1/mcp-servers/{id}/resources — список resources - POST /api/v1/mcp-servers/{id}/refresh-tools — обновить tools - POST /api/v1/mcp-servers/{id}/tools/call — вызвать tool """ from __future__ import annotations import uuid from typing import Any import structlog from fastapi import APIRouter, Depends, HTTPException, Query, status from pydantic import BaseModel from src.api.dependencies import get_mcp_service, get_user_ctx from src.db.models import MCPServer from src.middleware.auth import UserContext from src.schemas.mcp import ( MCPServerCreate, MCPServerResponse, MCPServerUpdate, MCPTestResult, MCPToolCallRequest, MCPToolCallResponse, ) from src.services import MCPServerService logger = structlog.get_logger() router = APIRouter(prefix="/api/v1/mcp-servers", tags=["mcp-servers"]) # ============================================================================ # RESPONSE MODELS (специфичные для API) # ============================================================================ class MCPServerStatsResponse(BaseModel): """Статистика по MCP серверам.""" total: int running: int active: int stopped: int # ============================================================================ # HELPER FUNCTIONS # ============================================================================ def _build_server_response(server: MCPServer) -> MCPServerResponse: """ Собрать MCPServerResponse из DB-модели. ⚠️ ``auth`` НЕ возвращается (security — секреты не в response). """ return MCPServerResponse( id=server.id, name=server.name, description=server.description, icon=server.icon, type=server.type, url=server.url, command=server.command, args=list(server.args or []), env=dict(server.env or {}), # auth намеренно исключён (секреты) tools=list(server.tools or []), resources=list(server.resources or []), is_active=server.is_active, auto_start=server.auto_start, timeout_seconds=server.timeout_seconds, tags=list(server.tags or []), is_public=server.is_public, version=server.version, usage_count=server.usage_count, last_used_at=server.last_used_at.isoformat() if server.last_used_at else None, error_count=server.error_count, status=server.status, last_error=server.last_error, created_at=server.created_at.isoformat(), updated_at=server.updated_at.isoformat(), ) # ============================================================================ # CRUD ENDPOINTS # ============================================================================ @router.get("", response_model=list[MCPServerResponse]) async def list_mcp_servers( is_active: bool | None = Query(None, description="Фильтр по активности"), status_filter: str | None = Query(None, alias="status", description="Фильтр по статусу"), type_filter: str | None = Query(None, alias="type", description="Фильтр по типу"), limit: int = Query(100, ge=1, le=500), offset: int = Query(0, ge=0), user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> list[MCPServerResponse]: """Список MCP серверов с фильтрами.""" servers = await mcp_service.list( is_active=is_active, status=status_filter, type=type_filter, limit=limit, offset=offset, ) logger.info( "mcp_servers.listed", count=len(servers), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return [_build_server_response(server) for server in servers] @router.post("", response_model=MCPServerResponse, status_code=status.HTTP_201_CREATED) async def create_mcp_server( data: MCPServerCreate, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPServerResponse: """ Создать новый MCP сервер. Валидация конфигурации (stdio→command, http/sse→url) через model_validator в MCPServerCreate. """ server = await mcp_service.create(data) logger.info( "mcp_server.created", server_id=str(server.id), name=server.name, type=server.type, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_server_response(server) @router.get("/stats", response_model=MCPServerStatsResponse) async def get_mcp_servers_stats( user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPServerStatsResponse: """Статистика по MCP серверам.""" stats = await mcp_service.get_stats() logger.info( "mcp_servers.stats", total=stats["total"], user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return MCPServerStatsResponse( total=stats["total"], running=stats["running"], active=stats["active"], stopped=stats["stopped"], ) @router.get("/{server_id}", response_model=MCPServerResponse) async def get_mcp_server( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPServerResponse: """Получить MCP сервер по ID.""" server = await mcp_service.get(server_id) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) return _build_server_response(server) @router.patch("/{server_id}", response_model=MCPServerResponse) async def update_mcp_server( server_id: uuid.UUID, data: MCPServerUpdate, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPServerResponse: """Обновить MCP сервер (частичное обновление).""" server = await mcp_service.update(server_id, data) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) logger.info( "mcp_server.updated", server_id=str(server_id), updated_fields=list(data.model_dump(exclude_unset=True).keys()), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_server_response(server) @router.delete("/{server_id}", status_code=status.HTTP_204_NO_CONTENT) async def delete_mcp_server( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> None: """Удалить MCP сервер (автоматически останавливает перед удалением).""" server = await mcp_service.get(server_id) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) # Нельзя удалить запущенный сервер if server.status == "running": raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail=f"Cannot delete MCP server '{server.name}' while running. Stop it first.", ) deleted = await mcp_service.delete(server_id) if not deleted: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) logger.info( "mcp_server.deleted", server_id=str(server_id), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) # ============================================================================ # LIFECYCLE ENDPOINTS (реальные, через MCPManager) # ============================================================================ @router.post("/{server_id}/start", response_model=MCPServerResponse) async def start_mcp_server( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPServerResponse: """Запустить MCP сервер (stdio subprocess / HTTP подключение).""" try: server = await mcp_service.start(server_id) except ValueError as e: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) except RuntimeError as e: # MCPManager недоступен или ошибка запуска raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(e)) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) logger.info( "mcp_server.started", server_id=str(server_id), status=server.status, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_server_response(server) @router.post("/{server_id}/stop", response_model=MCPServerResponse) async def stop_mcp_server( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPServerResponse: """Остановить MCP сервер.""" try: server = await mcp_service.stop(server_id) except ValueError as e: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) except RuntimeError as e: raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(e)) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) logger.info( "mcp_server.stopped", server_id=str(server_id), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_server_response(server) @router.post("/{server_id}/restart", response_model=MCPServerResponse) async def restart_mcp_server( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPServerResponse: """Перезапустить MCP сервер.""" try: server = await mcp_service.restart(server_id) except ValueError as e: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) except RuntimeError as e: raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(e)) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) logger.info( "mcp_server.restarted", server_id=str(server_id), status=server.status, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_server_response(server) # ============================================================================ # TEST CONNECTION # ============================================================================ @router.post("/{server_id}/test", response_model=MCPTestResult) async def test_mcp_connection( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPTestResult: """Проверить соединение с MCP сервером (с latency и tools_count).""" try: result = await mcp_service.test_connection(server_id) except ValueError as e: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) except RuntimeError as e: raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(e)) logger.info( "mcp_server.tested", server_id=str(server_id), success=result["success"], user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return MCPTestResult( server_id=server_id, success=result["success"], latency_ms=result.get("latency_ms"), message=result.get("message"), tools_count=result.get("tools_count", 0), ) # ============================================================================ # TOOLS & RESOURCES ENDPOINTS # ============================================================================ @router.get("/{server_id}/tools") async def list_mcp_server_tools( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> dict[str, Any]: """Список tools, предоставляемых MCP сервером.""" server = await mcp_service.get(server_id) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) return { "server_id": str(server_id), "server_name": server.name, "tools": server.tools or [], "tools_count": len(server.tools or []), } @router.get("/{server_id}/resources") async def list_mcp_server_resources( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> dict[str, Any]: """Список resources, предоставляемых MCP сервером.""" server = await mcp_service.get(server_id) if server is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"MCP server {server_id} not found", ) return { "server_id": str(server_id), "server_name": server.name, "resources": server.resources or [], "resources_count": len(server.resources or []), } @router.post("/{server_id}/refresh-tools") async def refresh_mcp_server_tools( server_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> dict[str, Any]: """ Обновить список tools из MCP сервера (реальный discovery). Требует запущенный сервер. """ try: tools = await mcp_service.refresh_tools(server_id) except ValueError as e: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) except RuntimeError as e: # Сервер не запущен или MCPManager недоступен raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(e)) logger.info( "mcp_server.tools.refreshed", server_id=str(server_id), tools_count=len(tools), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return { "status": "refreshed", "server_id": str(server_id), "tools": tools, "tools_count": len(tools), } @router.post("/{server_id}/tools/call", response_model=MCPToolCallResponse) async def call_mcp_tool( server_id: uuid.UUID, data: MCPToolCallRequest, user_ctx: UserContext = Depends(get_user_ctx), mcp_service: MCPServerService = Depends(get_mcp_service), ) -> MCPToolCallResponse: """Вызвать инструмент MCP сервера.""" try: result = await mcp_service.call_tool(server_id, data.tool_name, data.arguments) except ValueError as e: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(e)) except RuntimeError as e: # MCPManager недоступен или ошибка вызова raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(e)) logger.info( "mcp_server.tool.called", server_id=str(server_id), tool_name=data.tool_name, success=result["success"], user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return MCPToolCallResponse( server_id=server_id, tool_name=data.tool_name, success=result["success"], result=result.get("result"), error=result.get("error"), duration_ms=result.get("duration_ms", 0.0), )