/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/api/flows.py
329 строк
11 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""API endpoints для управления flows (мульти-агентные workflow). Предоставляет: - CRUD операции для flows (графы nodes/edges) - Запуск flows (sync и SSE streaming) - Валидацию графа (START/END, связность, безопасные conditions) - Статистику Эндпоинты: - GET /api/v1/flows — список flows - POST /api/v1/flows — создать flow - GET /api/v1/flows/stats — статистика - POST /api/v1/flows/validate — валидировать граф - GET /api/v1/flows/{id} — получить flow - PATCH /api/v1/flows/{id} — обновить flow - DELETE /api/v1/flows/{id} — удалить flow - POST /api/v1/flows/{id}/run — запустить flow (sync/stream) """ from __future__ import annotations import json import uuid from collections.abc import AsyncIterator from typing import Any import structlog from fastapi import APIRouter, Depends, HTTPException, Query, status from pydantic import BaseModel from sse_starlette.sse import EventSourceResponse from src.api.dependencies import get_flow_repo, get_flow_service, get_user_ctx from src.db.models import Flow as FlowModel from src.db.repositories import FlowRepository from src.middleware.auth import UserContext from src.primitives import Task from src.schemas.flow import ( FlowCreate, FlowResponse, FlowRunRequest, FlowRunResponse, FlowUpdate, ) from src.services import FlowService logger = structlog.get_logger() router = APIRouter(prefix="/api/v1/flows", tags=["flows"]) # ============================================================================ # RESPONSE MODELS (специфичные для API) # ============================================================================ class FlowStatsResponse(BaseModel): """Статистика по flows.""" total_flows: int by_status: dict[str, int] class FlowValidateResponse(BaseModel): """Результат валидации графа flow.""" valid: bool nodes_count: int edges_count: int # ============================================================================ # HELPER FUNCTIONS # ============================================================================ def _build_flow_response(flow: FlowModel) -> FlowResponse: """Собрать FlowResponse из DB-модели (через from_attributes).""" return FlowResponse.model_validate(flow) def _build_run_response(flow_id: uuid.UUID, task: Task) -> FlowRunResponse: """Собрать FlowRunResponse из результата выполнения (primitives.Task).""" output = task.output if isinstance(task.output, dict) else {"output": task.output} return FlowRunResponse( flow_id=flow_id, status=task.status.value, output=output, steps=output.get("steps", []), tokens_used=output.get("tokens", 0), duration_ms=task.duration_ms, error=task.error, ) # ============================================================================ # CRUD ENDPOINTS # ============================================================================ @router.get("", response_model=list[FlowResponse]) async def list_flows( status_filter: str | None = Query(None, alias="status", description="Фильтр по статусу"), limit: int = Query(100, ge=1, le=500), offset: int = Query(0, ge=0), user_ctx: UserContext = Depends(get_user_ctx), flow_repo: FlowRepository = Depends(get_flow_repo), ) -> list[FlowResponse]: """Список flows в текущем workspace с фильтрацией и пагинацией.""" flows = await flow_repo.list(status=status_filter, limit=limit, offset=offset) logger.info( "flows.listed", count=len(flows), status=status_filter, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return [_build_flow_response(flow) for flow in flows] @router.post("", response_model=FlowResponse, status_code=status.HTTP_201_CREATED) async def create_flow( data: FlowCreate, user_ctx: UserContext = Depends(get_user_ctx), flow_service: FlowService = Depends(get_flow_service), ) -> FlowResponse: """ Создать новый flow. Граф автоматически валидируется (START/END nodes, связность, безопасные conditions) через Pydantic model_validator в FlowCreate. """ flow = await flow_service.create(data) logger.info( "flow.created", flow_id=str(flow.id), name=flow.name, nodes_count=len(flow.nodes or []), edges_count=len(flow.edges or []), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_flow_response(flow) @router.get("/stats", response_model=FlowStatsResponse) async def get_flows_stats( user_ctx: UserContext = Depends(get_user_ctx), flow_repo: FlowRepository = Depends(get_flow_repo), ) -> FlowStatsResponse: """Статистика по flows (количество по статусам).""" statuses = ["draft", "active", "paused", "archived"] by_status = {s: await flow_repo.count(status=s) for s in statuses} total = await flow_repo.count() logger.info( "flows.stats", total=total, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return FlowStatsResponse(total_flows=total, by_status=by_status) @router.post("/validate", response_model=FlowValidateResponse) async def validate_flow( data: FlowCreate, user_ctx: UserContext = Depends(get_user_ctx), ) -> FlowValidateResponse: """ Валидировать граф flow без создания. Note: FlowCreate уже валидирует граф через model_validator. Если запрос дошёл сюда — граф валиден (иначе FastAPI вернул бы 422). """ return FlowValidateResponse( valid=True, nodes_count=len(data.nodes), edges_count=len(data.edges), ) @router.get("/{flow_id}", response_model=FlowResponse) async def get_flow( flow_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), flow_repo: FlowRepository = Depends(get_flow_repo), ) -> FlowResponse: """Получить flow по ID.""" flow = await flow_repo.get(flow_id) if flow is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Flow {flow_id} not found", ) return _build_flow_response(flow) @router.patch("/{flow_id}", response_model=FlowResponse) async def update_flow( flow_id: uuid.UUID, data: FlowUpdate, user_ctx: UserContext = Depends(get_user_ctx), flow_service: FlowService = Depends(get_flow_service), ) -> FlowResponse: """ Обновить flow (частичное обновление). Если переданы nodes И edges — граф повторно валидируется. """ flow = await flow_service.update(flow_id, data) if flow is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Flow {flow_id} not found", ) logger.info( "flow.updated", flow_id=str(flow_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_flow_response(flow) @router.delete("/{flow_id}", status_code=status.HTTP_204_NO_CONTENT) async def delete_flow( flow_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), flow_service: FlowService = Depends(get_flow_service), ) -> None: """Удалить flow.""" deleted = await flow_service.delete(flow_id) if not deleted: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Flow {flow_id} not found", ) logger.info( "flow.deleted", flow_id=str(flow_id), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) # ============================================================================ # RUN ENDPOINTS # ============================================================================ # ✅ response_model=None — функция возвращает Union (FlowRunResponse | EventSourceResponse), # FastAPI не может вывести единую Pydantic-схему из EventSourceResponse @router.post("/{flow_id}/run", response_model=None) async def run_flow( flow_id: uuid.UUID, data: FlowRunRequest, user_ctx: UserContext = Depends(get_user_ctx), flow_service: FlowService = Depends(get_flow_service), ) -> FlowRunResponse | EventSourceResponse: """ Запустить flow с указанными входными данными. Два режима: - **stream=true** (default): SSE поток событий (flow_start, node_start, node_done, agent_message, flow_done, error) - **stream=false**: JSON с полным результатом """ logger.info( "flow.run.requested", flow_id=str(flow_id), stream=data.stream, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) if data.stream: return EventSourceResponse(_stream_flow_run(flow_id, data.input, flow_service)) # Non-streaming mode task = await flow_service.run(flow_id, data.input) return _build_run_response(flow_id, task) async def _stream_flow_run( flow_id: uuid.UUID, input_data: dict[str, Any], flow_service: FlowService, ) -> AsyncIterator[dict[str, str]]: """Генератор SSE событий для streaming выполнения flow.""" try: async for event in flow_service.run_stream(flow_id, input_data): event_type = event.get("type", "message") yield { "event": event_type, "data": json.dumps(event, ensure_ascii=False, default=str), } except HTTPException as e: yield { "event": "error", "data": json.dumps({"type": "error", "error": str(e.detail)}, ensure_ascii=False), } except Exception as e: logger.exception( "flow.run.stream.error", flow_id=str(flow_id), error=str(e), error_type=type(e).__name__, ) yield { "event": "error", "data": json.dumps( { "type": "error", "error": str(e), "error_type": type(e).__name__, }, ensure_ascii=False, ), }