/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/rag/src/api/mod.rs
316 строк
11 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
//! HTTP API module for RAG service. //! //! Определяет главный роутер API со всеми эндпоинтами и middleware: //! //! ## Эндпоинты //! //! - `GET /health` — liveness probe //! - `GET /ready` — readiness probe //! - `GET /stats` — статистика хранилища //! - `POST /ingest` — загрузка документа //! - `POST /ingest/batch` — пакетная загрузка //! - `POST /search` — семантический поиск //! - `POST /search/batch` — пакетный поиск //! - `GET /documents` — список документов (с фильтром по workspace) //! - `GET /documents/all` — список документов всех workspaces //! - `DELETE /documents/{document_id}` — удаление документа //! - `DELETE /workspaces/{workspace_id}` — удаление workspace //! //! ## Middleware //! //! - **CORS** — разрешает запросы с любых источников //! - **TraceLayer** — трейсинг запросов/ответов //! - **TimeoutLayer** — ограничение времени обработки запроса //! - **Request ID** — уникальный `X-Request-ID` в каждом response //! - **DefaultBodyLimit** — лимит размера тела запроса (50MB) pub mod handlers; pub mod models; use std::sync::Arc; use std::time::Duration; use axum::{ extract::{DefaultBodyLimit, Request}, middleware::{self, Next}, response::Response, routing::{delete, get, post}, Router, }; use tower_http::{ cors::{Any, CorsLayer}, timeout::TimeoutLayer, trace::TraceLayer, }; use crate::pipeline::RagPipeline; /// Create the main API router with all endpoints and middleware. /// /// # Arguments /// /// * `pipeline` — общий RAG pipeline (Arc для шаринга между хендлерами) /// * `request_timeout` — максимальное время обработки запроса /// /// # Middleware порядок (применяются в обратном порядке к запросу) /// /// 1. **CORS** — первый обрабатывает preflight OPTIONS /// 2. **TraceLayer** — логирует запрос/ответ /// 3. **Timeout** — прерывает долгий запрос /// 4. **Request ID** — добавляет X-Request-ID /// 5. **DefaultBodyLimit** — ограничивает размер тела (50MB) #[allow(deprecated)] pub fn create_router(pipeline: Arc<RagPipeline>, request_timeout: Duration) -> Router { let cors = CorsLayer::new() .allow_origin(Any) .allow_methods(Any) .allow_headers(Any); let timeout = TimeoutLayer::new(request_timeout); Router::new() // ==================================================================== // Health & stats // ==================================================================== .route("/health", get(handlers::health)) .route("/ready", get(handlers::ready)) .route("/stats", get(handlers::stats)) // ==================================================================== // Ingest // ==================================================================== .route("/ingest", post(handlers::ingest)) .route("/ingest/batch", post(handlers::ingest_batch)) // ==================================================================== // Search // ==================================================================== .route("/search", post(handlers::search)) .route("/search/batch", post(handlers::search_batch)) // ==================================================================== // Documents — listing (статические роуты ДО динамических!) // ==================================================================== // ⚠️ ВАЖНО: эти роуты должны идти ДО `/documents/{document_id}`, // иначе axum будет пытаться матчить "all" как document_id .route("/documents", get(handlers::list_documents)) .route("/documents/all", get(handlers::list_all_documents)) // ==================================================================== // Documents & Workspaces — delete (динамические роуты) // ==================================================================== .route( "/documents/{document_id}", delete(handlers::delete_document), ) .route( "/workspaces/{workspace_id}", delete(handlers::delete_workspace), ) // ==================================================================== // Middleware (applied in reverse order to the request) // ==================================================================== // 5. Ограничение размера тела запроса (50MB для ingest) .layer(DefaultBodyLimit::max(50 * 1024 * 1024)) // 4. Request ID .layer(middleware::from_fn(add_request_id)) // 3. Timeout .layer(timeout) // 2. Trace layer .layer(TraceLayer::new_for_http()) // 1. CORS .layer(cors) // Shared state .with_state(pipeline) } /// Middleware: добавить уникальный X-Request-ID в каждый response. /// /// Позволяет коррелировать логи одного запроса между разными сервисами. /// ID генерируется как UUID v4 и добавляется в заголовок `X-Request-ID`. async fn add_request_id(request: Request, next: Next) -> Response { let request_id = uuid::Uuid::new_v4().to_string(); let mut response = next.run(request).await; response .headers_mut() .insert("X-Request-ID", request_id.parse().unwrap()); response } // ============================================================================ // Tests // ============================================================================ #[cfg(test)] mod tests { use super::*; use axum::body::Body; use axum::http::{Method, Request as HttpRequest, StatusCode}; use tower::ServiceExt; #[tokio::test] async fn test_router_creation() { let pipeline = Arc::new(RagPipeline::for_testing().unwrap()); let timeout = Duration::from_secs(30); let _router = create_router(pipeline, timeout); } #[tokio::test] async fn test_health_endpoint() { let pipeline = Arc::new(RagPipeline::for_testing().unwrap()); let router = create_router(pipeline, Duration::from_secs(30)); let response = router .oneshot( HttpRequest::builder() .method(Method::GET) .uri("/health") .body(Body::empty()) .unwrap(), ) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); assert!(response.headers().get("X-Request-ID").is_some()); } #[tokio::test] async fn test_stats_endpoint() { let pipeline = Arc::new(RagPipeline::for_testing().unwrap()); let router = create_router(pipeline, Duration::from_secs(30)); let response = router .oneshot( HttpRequest::builder() .method(Method::GET) .uri("/stats") .body(Body::empty()) .unwrap(), ) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); } #[tokio::test] async fn test_documents_endpoint() { let pipeline = Arc::new(RagPipeline::for_testing().unwrap()); let router = create_router(pipeline, Duration::from_secs(30)); // Без query параметра (должен использовать default workspace) let response = router .clone() .oneshot( HttpRequest::builder() .method(Method::GET) .uri("/documents") .body(Body::empty()) .unwrap(), ) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); // С query параметром let response = router .oneshot( HttpRequest::builder() .method(Method::GET) .uri("/documents?workspace_id=test") .body(Body::empty()) .unwrap(), ) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); } #[tokio::test] async fn test_documents_all_endpoint() { let pipeline = Arc::new(RagPipeline::for_testing().unwrap()); let router = create_router(pipeline, Duration::from_secs(30)); let response = router .oneshot( HttpRequest::builder() .method(Method::GET) .uri("/documents/all") .body(Body::empty()) .unwrap(), ) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); } #[tokio::test] async fn test_request_id_header() { let pipeline = Arc::new(RagPipeline::for_testing().unwrap()); let router = create_router(pipeline, Duration::from_secs(30)); let response1 = router .clone() .oneshot( HttpRequest::builder() .method(Method::GET) .uri("/health") .body(Body::empty()) .unwrap(), ) .await .unwrap(); let response2 = router .oneshot( HttpRequest::builder() .method(Method::GET) .uri("/health") .body(Body::empty()) .unwrap(), ) .await .unwrap(); let id1 = response1 .headers() .get("X-Request-ID") .unwrap() .to_str() .unwrap(); let id2 = response2 .headers() .get("X-Request-ID") .unwrap() .to_str() .unwrap(); // Request ID должны быть уникальными для каждого запроса assert_ne!(id1, id2); } #[tokio::test] async fn test_cors_headers() { let pipeline = Arc::new(RagPipeline::for_testing().unwrap()); let router = create_router(pipeline, Duration::from_secs(30)); let response = router .oneshot( HttpRequest::builder() .method(Method::OPTIONS) .uri("/health") .header("Origin", "http://localhost:3000") .header("Access-Control-Request-Method", "GET") .body(Body::empty()) .unwrap(), ) .await .unwrap(); assert_eq!(response.status(), StatusCode::OK); assert!(response .headers() .get("access-control-allow-origin") .is_some()); } }