/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/gateway/src/proxy/rag_forward.rs
294 строки
10 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
//! Forward handlers для проксирования запросов к RAG сервису. //! //! Использует generic [`proxy_to_rag`] (аналог [`super::forward::proxy_to_engine`]): //! - Circuit breaker //! - Inject `X-User-ID` / `X-Workspace-ID` из AuthContext //! - Автоопределение SSE streaming (по `Content-Type` ответа) //! - URL нормализация через [`RagClient::url`] //! //! # Body лимиты //! //! Глобальный лимит 50MB установлен в `main.rs` (`DefaultBodyLimit`) — //! покрывает ingest больших документов. Search body обычно < 1MB. use axum::{ body::{Body, Bytes}, extract::{Path, RawQuery, State}, http::{HeaderMap, HeaderValue, Method, StatusCode, header}, response::Response, }; use futures_util::StreamExt; use reqwest::header::CONTENT_TYPE; use tracing::{debug, error, warn}; use super::circuit_breaker::CircuitBreaker; use super::forward::error_response; use super::rag_client::RagClient; use crate::auth::workspace::WorkspaceId; use crate::auth::{AuthContext, AuthMethod, OptionalAuth}; // ============================================================================ // Proxy State // ============================================================================ /// Shared proxy state для RAG. #[derive(Clone)] pub struct RagProxyState { pub client: RagClient, pub circuit_breaker: CircuitBreaker, } // ============================================================================ // Generic Proxy Helper // ============================================================================ /// Generic proxy запрос к RAG сервису. /// /// Единая логика (аналог [`super::forward::proxy_to_engine`]): /// - Circuit breaker check /// - URL нормализация через [`RagClient::url`] /// - Копирование headers клиента (кроме hop-by-hop) /// - Inject auth headers /// - Автоопределение SSE streaming (по `Content-Type`) async fn proxy_to_rag( state: &RagProxyState, method: Method, path: &str, auth: Option<AuthContext>, headers: HeaderMap, body: Option<Bytes>, ) -> Response { // Circuit breaker if !state.circuit_breaker.allow_request() { warn!(path = %path, "Circuit breaker open, rejecting RAG request"); return error_response( StatusCode::SERVICE_UNAVAILABLE, "service_unavailable", "RAG service is temporarily unavailable. Please try again later.", ); } let url = state.client.url(path); debug!(url = %url, method = %method, "Proxying request to RAG"); let mut req_builder = state.client.inner().request(method, &url); // Body (если есть) if let Some(bytes) = body { req_builder = req_builder.body(bytes); } // Копируем headers клиента (кроме hop-by-hop) for (name, value) in headers.iter() { if name != header::HOST && name != header::CONTENT_LENGTH { req_builder = req_builder.header(name, value); } } // Inject auth headers req_builder = inject_auth_headers(req_builder, &auth); // Отправка запроса let response = match req_builder.send().await { Ok(resp) => { state.circuit_breaker.record_success(); resp } Err(e) => { state.circuit_breaker.record_failure(); error!(error = %e, url = %url, "RAG request failed"); return error_response( StatusCode::BAD_GATEWAY, "bad_gateway", &format!("RAG request failed: {}", e), ); } }; // ✅ Автоопределение SSE streaming по Content-Type let is_stream = response .headers() .get(CONTENT_TYPE) .and_then(|v| v.to_str().ok()) .map(|v| v.contains("text/event-stream")) .unwrap_or(false); if is_stream && response.status().is_success() { return stream_response(response).await; } build_proxy_response(response).await } /// Inject `X-User-ID` и `X-Workspace-ID` из AuthContext. /// /// Для anonymous — используется default workspace (`WorkspaceId::default_dev()`). fn inject_auth_headers( mut req_builder: reqwest::RequestBuilder, auth: &Option<AuthContext>, ) -> reqwest::RequestBuilder { // Workspace — из auth или default let workspace_id = match auth { Some(ctx) if ctx.auth_method != AuthMethod::Anonymous => ctx.workspace_id.clone(), _ => WorkspaceId::default_dev(), }; if let Ok(value) = HeaderValue::from_str(workspace_id.as_str()) { req_builder = req_builder.header("X-Workspace-ID", value); } // User ID — из auth (если есть) if let Some(ctx) = auth { if let Some(user_id) = &ctx.user_id { if let Ok(value) = HeaderValue::from_str(user_id) { req_builder = req_builder.header("X-User-ID", value); } } } req_builder } // ============================================================================ // Response Builders // ============================================================================ /// SSE streaming — прозрачно проксируем байты от RAG к клиенту. async fn stream_response(response: reqwest::Response) -> Response { let status = response.status(); let response_headers = response.headers().clone(); let stream = response.bytes_stream().filter_map(|result| async move { match result { Ok(bytes) => Some(Ok::<Bytes, std::convert::Infallible>(bytes)), Err(e) => { warn!(error = %e, "Error reading from RAG stream"); None } } }); let mut builder = Response::builder() .status(status.as_u16()) .header(CONTENT_TYPE, "text/event-stream") .header("Cache-Control", "no-cache, no-transform") .header("Connection", "keep-alive") .header("X-Accel-Buffering", "no"); // nginx: отключить буферизацию SSE for (name, value) in response_headers.iter() { if name != CONTENT_TYPE && name != header::CONTENT_LENGTH { builder = builder.header(name, value); } } builder.body(Body::from_stream(stream)).unwrap_or_else(|_| { error_response( StatusCode::INTERNAL_SERVER_ERROR, "internal_error", "Failed to build SSE response", ) }) } /// Построить прокси ответ из RAG response (буферизованный). async fn build_proxy_response(response: reqwest::Response) -> Response { let status = response.status(); let response_headers = response.headers().clone(); match response.bytes().await { Ok(bytes) => { let mut builder = Response::builder().status(status.as_u16()); for (name, value) in response_headers.iter() { let name_lower = name.as_str().to_lowercase(); // Пропускаем hop-by-hop headers if name_lower == "content-length" || name_lower == "transfer-encoding" { continue; } builder = builder.header(name, value); } builder.body(Body::from(bytes)).unwrap_or_else(|_| { error_response( StatusCode::INTERNAL_SERVER_ERROR, "internal_error", "Failed to build response", ) }) } Err(e) => { error!(error = %e, "Failed to read RAG response"); error_response( StatusCode::BAD_GATEWAY, "bad_gateway", &format!("Failed to read RAG response: {}", e), ) } } } // ============================================================================ // RAG PROXY HANDLERS // ============================================================================ /// POST /api/v1/rag/ingest → RAG POST /ingest (загрузка документа, до 50MB) pub async fn proxy_rag_ingest( State(state): State<RagProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { proxy_to_rag(&state, Method::POST, "/ingest", auth, headers, Some(body)).await } /// POST /api/v1/rag/search → RAG POST /search (семантический поиск, SSE при stream=true) pub async fn proxy_rag_search( State(state): State<RagProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { proxy_to_rag(&state, Method::POST, "/search", auth, headers, Some(body)).await } /// GET /api/v1/rag/stats → RAG GET /stats pub async fn proxy_rag_stats( State(state): State<RagProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { proxy_to_rag(&state, Method::GET, "/stats", auth, headers, None).await } /// GET /api/v1/rag/documents → RAG GET /documents (с query params) pub async fn proxy_rag_list_documents( State(state): State<RagProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, RawQuery(query): RawQuery, ) -> Response { let path = match query { Some(q) if !q.is_empty() => format!("/documents?{}", q), _ => "/documents".to_string(), }; proxy_to_rag(&state, Method::GET, &path, auth, headers, None).await } /// DELETE /api/v1/rag/documents/{id} → RAG DELETE /documents/{id} pub async fn proxy_rag_delete_document( State(state): State<RagProxyState>, Path(document_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { let path = format!("/documents/{}", document_id); proxy_to_rag(&state, Method::DELETE, &path, auth, headers, None).await } /// DELETE /api/v1/rag/workspaces/{id} → RAG DELETE /workspaces/{id} pub async fn proxy_rag_delete_workspace( State(state): State<RagProxyState>, Path(workspace_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { let path = format!("/workspaces/{}", workspace_id); proxy_to_rag(&state, Method::DELETE, &path, auth, headers, None).await }