/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/gateway/src/proxy/flows.rs
218 строк
7 KB
Alexander Efanov
upd fix
04 авг 2026, 14:09
04 авг 2026, 14:09
f63ab86
Код
Авторство
О чём код?
//! Proxy handlers для Flows API (Gateway → Engine). //! //! Использует generic [`proxy_to_engine`] из [`super::forward`]: //! - Автоопределение SSE streaming (для `POST /flows/{id}/run` и `POST /runs`) //! - Inject `X-User-ID` / `X-Workspace-ID` из AuthContext //! - Circuit breaker + error handling //! //! ## DB-based Flows (UUID): //! - GET /api/v1/flows — список flows //! - POST /api/v1/flows — создать flow //! - GET /api/v1/flows/stats — статистика //! - POST /api/v1/flows/validate — валидировать граф //! - GET /api/v1/flows/{flow_id} — получить flow //! - PATCH /api/v1/flows/{flow_id} — обновить flow //! - DELETE /api/v1/flows/{flow_id} — удалить flow //! - POST /api/v1/flows/{flow_id}/run — запустить flow (SSE) //! //! ## Predefined Flows (string ID из FLOW_REGISTRY): //! - GET /api/v1/runs/flows — список predefined flows //! - POST /api/v1/runs — запустить predefined flow (SSE) //! //! ⚠️ Порядок роутов в main.rs: //! - `/stats`, `/validate` ДО `/{flow_id}` (статические до динамических) //! - `/runs/flows` ДО `/runs` (если бы был catch-all) use axum::{ body::Bytes, extract::{Path, State}, http::{HeaderMap, Method}, response::Response, }; use tracing::info; use super::forward::{ProxyState, proxy_to_engine}; use crate::auth::OptionalAuth; // ============================================================================ // DB FLOWS — COLLECTION ENDPOINTS // ============================================================================ /// GET /api/v1/flows — список flows с фильтрами pub async fn proxy_flows_list( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { proxy_to_engine(&state, Method::GET, "/api/v1/flows", auth, headers, None).await } /// POST /api/v1/flows — создать новый flow pub async fn proxy_flows_create( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { info!("proxy_flows_create"); proxy_to_engine( &state, Method::POST, "/api/v1/flows", auth, headers, Some(body), ) .await } /// GET /api/v1/flows/stats — статистика по flows /// /// ⚠️ Должен быть зарегистрирован ДО `/{flow_id}` в main.rs. pub async fn proxy_flows_stats( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { proxy_to_engine( &state, Method::GET, "/api/v1/flows/stats", auth, headers, None, ) .await } /// POST /api/v1/flows/validate — валидировать граф без создания /// /// ⚠️ Должен быть зарегистрирован ДО `/{flow_id}` в main.rs. pub async fn proxy_flows_validate( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { proxy_to_engine( &state, Method::POST, "/api/v1/flows/validate", auth, headers, Some(body), ) .await } // ============================================================================ // DB FLOWS — ITEM ENDPOINTS // ============================================================================ /// GET /api/v1/flows/{flow_id} — получить flow по ID pub async fn proxy_flows_get( State(state): State<ProxyState>, Path(flow_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { let path = format!("/api/v1/flows/{}", flow_id); proxy_to_engine(&state, Method::GET, &path, auth, headers, None).await } /// PATCH /api/v1/flows/{flow_id} — обновить flow pub async fn proxy_flows_update( State(state): State<ProxyState>, Path(flow_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { info!(flow_id = %flow_id, "proxy_flows_update"); let path = format!("/api/v1/flows/{}", flow_id); proxy_to_engine(&state, Method::PATCH, &path, auth, headers, Some(body)).await } /// DELETE /api/v1/flows/{flow_id} — удалить flow pub async fn proxy_flows_delete( State(state): State<ProxyState>, Path(flow_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { info!(flow_id = %flow_id, "proxy_flows_delete"); let path = format!("/api/v1/flows/{}", flow_id); proxy_to_engine(&state, Method::DELETE, &path, auth, headers, None).await } // ============================================================================ // DB FLOWS — ACTION ENDPOINTS // ============================================================================ /// POST /api/v1/flows/{flow_id}/run — запустить DB flow (sync или SSE) /// /// Engine возвращает SSE (`text/event-stream`) при `stream=true` — /// `proxy_to_engine` автоматически определяет это и стримит ответ /// (события: flow_start, node_start, node_done, agent_message, flow_done). pub async fn proxy_flows_run( State(state): State<ProxyState>, Path(flow_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { info!(flow_id = %flow_id, body_size = body.len(), "proxy_flows_run"); let path = format!("/api/v1/flows/{}/run", flow_id); proxy_to_engine(&state, Method::POST, &path, auth, headers, Some(body)).await } // ============================================================================ // PREDEFINED FLOWS (RUNS) — из FLOW_REGISTRY // ============================================================================ /// GET /api/v1/runs/flows — список predefined flows /// /// Возвращает массив: `[{id, name, description, agents}, ...]` /// Predefined flows: research, review, debate, code_review, consensus, /// compare, brainstorm, analyze. pub async fn proxy_runs_flows( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { proxy_to_engine( &state, Method::GET, "/api/v1/runs/flows", auth, headers, None, ) .await } /// POST /api/v1/runs — запуск predefined flow (SSE при stream=true) /// /// Body: `{"flow_id": "research", "input": {"topic": "..."}, "stream": true}` /// /// Engine возвращает SSE с событиями: /// agent_start, agent_message, agent_done, ideas_generated, flow_done, error. /// /// ⚠️ Это НЕ DB-flow — flow_id здесь строковый ("research"), не UUID. pub async fn proxy_runs( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { info!(body_size = body.len(), "proxy_runs"); proxy_to_engine( &state, Method::POST, "/api/v1/runs", auth, headers, Some(body), ) .await }