/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/api/runs.py
195 строк
6 KB
Alexander Efanov
upd fix
04 авг 2026, 14:09
04 авг 2026, 14:09
f63ab86
Код
Авторство
О чём код?
"""API endpoints для запуска predefined flows (мультиагентные workflow из FLOW_REGISTRY). В отличие от DB-based flows (/api/v1/flows/{uuid}/run), predefined flows идентифицируются строковым ID ("research", "brainstorm", etc.) и исполняются через flow.runner() — async generator событий. Эндпоинты: - GET /api/v1/runs/flows — список predefined flows - POST /api/v1/runs — запустить predefined flow (sync/stream) """ from __future__ import annotations import json from collections.abc import AsyncIterator from typing import Any import structlog from fastapi import APIRouter, HTTPException, status from pydantic import BaseModel, Field from sse_starlette.sse import EventSourceResponse from src.flows import FLOW_REGISTRY, get_flow, list_flows logger = structlog.get_logger() router = APIRouter(prefix="/api/v1/runs", tags=["runs"]) # ============================================================================ # SCHEMAS # ============================================================================ class PredefinedFlowInfo(BaseModel): """Информация о predefined flow.""" id: str name: str description: str agents: list[str] class RunRequest(BaseModel): """Запрос на запуск predefined flow.""" flow_id: str = Field(..., description="ID predefined flow (research, brainstorm, etc.)") input: dict[str, Any] = Field(default_factory=dict, description="Входные данные") stream: bool = Field(default=True, description="SSE streaming") class RunResponse(BaseModel): """Ответ non-streaming запуска.""" flow_id: str status: str output: str tokens: int stages: dict[str, Any] = Field(default_factory=dict) metadata: dict[str, Any] = Field(default_factory=dict) # ============================================================================ # ENDPOINTS # ============================================================================ @router.get("/flows", response_model=list[PredefinedFlowInfo]) async def list_predefined_flows() -> list[PredefinedFlowInfo]: """Список всех доступных predefined flows.""" result = [] for flow in list_flows(): result.append( PredefinedFlowInfo( id=flow.id, name=flow.name, description=flow.description, agents=flow.agents, ) ) return result @router.post("", response_model=None) async def run_predefined_flow(data: RunRequest) -> RunResponse | EventSourceResponse: """ Запустить predefined flow. - **stream=true** (default): SSE поток событий - **stream=false**: JSON с полным результатом """ flow = get_flow(data.flow_id) if flow is None: available = sorted(FLOW_REGISTRY.keys()) raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={ "error": f"Predefined flow '{data.flow_id}' not found", "available_flows": available, }, ) if not flow.has_runner(): raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, detail=f"Flow '{data.flow_id}' has no runner", ) logger.info( "runs.predefined_flow.requested", flow_id=data.flow_id, stream=data.stream, input_keys=list(data.input.keys()), ) if data.stream: return EventSourceResponse(_stream_run(flow.id, flow.runner, data.input)) # Non-streaming: собираем все события, возвращаем финальный return await _sync_run(flow.id, flow.runner, data.input) # ============================================================================ # HELPERS # ============================================================================ async def _stream_run( flow_id: str, runner, input_data: dict[str, Any], ) -> AsyncIterator[dict[str, str]]: """SSE генератор для predefined flow.""" try: async for event in runner(input_data): event_type = event.get("type", "message") yield { "event": event_type, "data": json.dumps(event, ensure_ascii=False, default=str), } except Exception as e: logger.exception( "runs.predefined_flow.stream.error", flow_id=flow_id, error=str(e), ) yield { "event": "error", "data": json.dumps( {"type": "error", "error": str(e)}, ensure_ascii=False, ), } async def _sync_run( flow_id: str, runner, input_data: dict[str, Any], ) -> RunResponse: """Non-streaming запуск: собираем все события, возвращаем flow_done.""" final_output = "" total_tokens = 0 stages: dict[str, Any] = {} metadata: dict[str, Any] = {} error: str | None = None try: async for event in runner(input_data): event_type = event.get("type") if event_type == "flow_done": final_output = event.get("output", "") total_tokens = event.get("tokens", 0) stages = event.get("stages", {}) metadata = event.get("metadata", {}) elif event_type == "error": error = event.get("error", "Unknown error") break except Exception as e: error = str(e) if error: raise HTTPException( status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail={"error": error, "flow_id": flow_id}, ) return RunResponse( flow_id=flow_id, status="completed", output=final_output, tokens=total_tokens, stages=stages, metadata=metadata, )