/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/api/dependencies.py
324 строки
11 KB
Alexander Efanov
upd fix
04 авг 2026, 14:09
04 авг 2026, 14:09
f63ab86
Код
Авторство
О чём код?
"""FastAPI dependencies для извлечения контекста и сервисов. Централизованный модуль с FastAPI Depends для: - Database session management - User context extraction (auth + workspace isolation) - Repository factories (с автоматической изоляцией по user_id + workspace_id) - Service factories (бизнес-логика с runtime-интеграцией) Note: используется классический паттерн ``Depends()`` в default-аргументах (стандарт FastAPI). Ruff B008 отключён в pyproject.toml. Usage в endpoints: @router.get("/items") async def list_items( service: AgentService = Depends(get_agent_service), ): return await service.list() """ from __future__ import annotations import uuid from collections.abc import AsyncGenerator from typing import Any import structlog from fastapi import Depends, Request from sqlalchemy.ext.asyncio import AsyncSession from src.db.repositories import ( AgentRepository, ChatMemoryRepository, ChatRepository, FlowRepository, MCPServerRepository, MessageRepository, SkillRepository, TaskRepository, UserRepository, ) from src.db.session import async_session_maker from src.middleware.auth import UserContext, require_user_context from src.runtime import MCPManager, get_mcp_manager from src.services import ( AgentService, ChatService, FlowService, MCPServerService, SkillService, TaskService, ) logger = structlog.get_logger() # ============================================================================ # CORE DEPENDENCIES (Session & User Context) # ============================================================================ async def get_db_session() -> AsyncGenerator[AsyncSession, None]: """ Dependency для получения async DB session. Автоматически: - Создаёт session из pool - Commit'ит при успешном завершении endpoint - Rollback'ит при exception - Возвращает connection в pool """ async with async_session_maker() as session: try: yield session await session.commit() except Exception as e: await session.rollback() logger.warning( "db.session.rollback", error=str(e), error_type=type(e).__name__, ) raise def get_user_ctx(request: Request) -> UserContext: """ Dependency для получения UserContext из request. UserContext извлекается из request.state (помещается AuthMiddleware). Raises: HTTPException(401): если UserContext отсутствует. """ return require_user_context(request) # ============================================================================ # RUNTIME DEPENDENCIES (app.state access) # ============================================================================ def get_tools_manager_dep(request: Request) -> Any: """ Dependency для получения ToolRegistryManager из app.state. ToolRegistryManager создаётся в lifespan (main.py) и хранится в ``app.state.tools_manager``. Используется AgentService для резолвинга имён инструментов в OpenAI tool schemas. Returns: ToolRegistryManager или None (если ещё не инициализирован). """ return getattr(request.app.state, "tools_manager", None) def get_mcp_manager_dep() -> MCPManager: """Singleton MCP менеджер (для реального MCP-взаимодействия).""" return get_mcp_manager() # ============================================================================ # REPOSITORY DEPENDENCIES # ============================================================================ async def get_user_repo( session: AsyncSession = Depends(get_db_session), ) -> UserRepository: """UserRepository (без изоляции — работает с самими пользователями).""" return UserRepository(session) async def get_chat_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> ChatRepository: """ChatRepository с автоматической изоляцией по user + workspace.""" return ChatRepository(session, user_ctx.user_id, user_ctx.workspace_id) async def get_message_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> MessageRepository: """MessageRepository с автоматической изоляцией.""" return MessageRepository(session, user_ctx.user_id, user_ctx.workspace_id) async def get_chat_memory_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> ChatMemoryRepository: """ChatMemoryRepository (долгосрочная память чатов).""" return ChatMemoryRepository(session, user_ctx.user_id, user_ctx.workspace_id) async def get_skill_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> SkillRepository: """SkillRepository с автоматической изоляцией.""" return SkillRepository(session, user_ctx.user_id, user_ctx.workspace_id) async def get_agent_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> AgentRepository: """AgentRepository с автоматической изоляцией.""" return AgentRepository(session, user_ctx.user_id, user_ctx.workspace_id) async def get_task_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> TaskRepository: """TaskRepository с автоматической изоляцией.""" return TaskRepository(session, user_ctx.user_id, user_ctx.workspace_id) async def get_flow_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> FlowRepository: """FlowRepository с автоматической изоляцией.""" return FlowRepository(session, user_ctx.user_id, user_ctx.workspace_id) async def get_mcp_server_repo( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> MCPServerRepository: """MCPServerRepository с автоматической изоляцией.""" return MCPServerRepository(session, user_ctx.user_id, user_ctx.workspace_id) # ============================================================================ # SERVICE DEPENDENCIES (бизнес-логика + runtime) # ============================================================================ async def get_chat_service( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> ChatService: """ChatService (оркестрация RAG + LLM + persistence).""" return ChatService(session, user_ctx.user_id, user_ctx.workspace_id) async def get_agent_service( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), tools_manager: Any = Depends(get_tools_manager_dep), ) -> AgentService: """AgentService (CRUD + запуск через AgentRuntime + tools resolution).""" return AgentService( session, user_ctx.user_id, user_ctx.workspace_id, tools_manager=tools_manager, ) async def get_flow_service( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> FlowService: """FlowService (CRUD + запуск через GraphExecutor).""" return FlowService(session, user_ctx.user_id, user_ctx.workspace_id) async def get_skill_service( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> SkillService: """SkillService (CRUD + запуск: render prompt → LLM).""" return SkillService(session, user_ctx.user_id, user_ctx.workspace_id) async def get_task_service( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> TaskService: """TaskService (CRUD + статусы + runs).""" return TaskService(session, user_ctx.user_id, user_ctx.workspace_id) async def get_mcp_service( user_ctx: UserContext = Depends(get_user_ctx), session: AsyncSession = Depends(get_db_session), ) -> MCPServerService: """ MCPServerService (CRUD + lifecycle + discovery через MCPManager). Graceful degradation: если ``mcp`` SDK не установлен — сервис работает в режиме метаданных (mcp_manager=None), реальное MCP-взаимодействие недоступно, но CRUD/status операции работают. """ try: manager: MCPManager | None = get_mcp_manager() except RuntimeError as e: logger.warning("mcp.manager_unavailable", error=str(e)) manager = None return MCPServerService( session, user_ctx.user_id, user_ctx.workspace_id, mcp_manager=manager, ) # ============================================================================ # UTILITY DEPENDENCIES # ============================================================================ def get_workspace_id( user_ctx: UserContext = Depends(get_user_ctx), ) -> str: """Shortcut для получения только workspace_id.""" return user_ctx.workspace_id def get_user_id( user_ctx: UserContext = Depends(get_user_ctx), ) -> uuid.UUID: """Shortcut для получения только user_id.""" return user_ctx.user_id # ============================================================================ # EXPORTS # ============================================================================ __all__ = [ # noqa: RUF022 # Core "get_db_session", "get_user_ctx", # Repositories "get_user_repo", "get_chat_repo", "get_message_repo", "get_chat_memory_repo", "get_skill_repo", "get_agent_repo", "get_task_repo", "get_flow_repo", "get_mcp_server_repo", # Services "get_chat_service", "get_agent_service", "get_flow_service", "get_skill_service", "get_task_service", "get_mcp_service", # Runtime "get_tools_manager_dep", "get_mcp_manager_dep", # Utilities "get_workspace_id", "get_user_id", ]