/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/gateway/src/proxy/forward.rs
228 строк
7 KB
Alexander Efanov
upd fix
04 авг 2026, 14:09
04 авг 2026, 14:09
f63ab86
Код
Авторство
О чём код?
//! Generic proxy helper — единая логика проксирования на engine. //! //! Содержит: //! - [`ProxyState`] — shared state (client + circuit breaker) //! - [`proxy_to_engine`] — generic proxy (используется всеми proxy-модулями) //! - [`read_body`] / [`error_response`] — вспомогательные функции use axum::{ Json, body::Body, http::{HeaderMap, HeaderValue, Method, StatusCode, header}, response::{IntoResponse, Response}, }; use bytes::Bytes; use futures_util::StreamExt; use reqwest::header::CONTENT_TYPE; use serde_json::json; use tracing::{debug, error, warn}; use super::circuit_breaker::CircuitBreaker; use super::client::EngineClient; use crate::auth::AuthContext; // ============================================================================ // Proxy State // ============================================================================ /// Shared proxy state (client + circuit breaker). #[derive(Clone)] pub struct ProxyState { pub client: EngineClient, pub circuit_breaker: CircuitBreaker, } // ============================================================================ // Generic Proxy Helper // ============================================================================ /// Generic proxy запрос к engine. /// /// - Circuit breaker check /// - Копирование headers (кроме hop-by-hop) /// - Inject `X-User-ID` / `X-Workspace-ID` из [`AuthContext`] /// - Автоопределение SSE streaming (по `Content-Type` ответа) pub async fn proxy_to_engine( state: &ProxyState, method: Method, path: &str, auth: Option<AuthContext>, headers: HeaderMap, body: Option<Bytes>, ) -> Response { if !state.circuit_breaker.allow_request() { warn!(path = %path, "Circuit breaker open, rejecting request"); return error_response( StatusCode::SERVICE_UNAVAILABLE, "service_unavailable", "Engine service is temporarily unavailable. Please try again later.", ); } let url = state.client.url(path); debug!(url = %url, method = %method, "Proxying request to engine"); let mut req_builder = state.client.inner().request(method, &url); 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, "Engine request failed"); return error_response( StatusCode::BAD_GATEWAY, "bad_gateway", &format!("Engine request failed: {}", e), ); } }; let status = response.status(); debug!(status = %status, path = %path, "Engine responded"); // SSE auto-detect 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 && status.is_success() { return stream_response(response).await; } // Обычный response 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() { 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 engine response"); error_response( StatusCode::BAD_GATEWAY, "bad_gateway", &format!("Failed to read response: {}", e), ) } } } /// Inject `X-User-ID` и `X-Workspace-ID` из AuthContext. fn inject_auth_headers( mut req_builder: reqwest::RequestBuilder, auth: &Option<AuthContext>, ) -> reqwest::RequestBuilder { let Some(ctx) = auth else { return req_builder; }; if let Ok(value) = HeaderValue::from_str(ctx.workspace_id.as_str()) { req_builder = req_builder.header("X-Workspace-ID", value); } let user_id = ctx .user_id .clone() .unwrap_or_else(|| "00000000-0000-0000-0000-000000000000".to_string()); if let Ok(value) = HeaderValue::from_str(&user_id) { req_builder = req_builder.header("X-User-ID", value); } req_builder } // ============================================================================ // Helpers // ============================================================================ /// Прочитать body запроса (с лимитом). pub async fn read_body(body: Body, limit: usize) -> Result<Bytes, Response> { axum::body::to_bytes(body, limit).await.map_err(|e| { error_response( StatusCode::BAD_REQUEST, "bad_request", &format!("Failed to read request body: {}", e), ) }) } /// Стандартизированный error response. pub fn error_response(status: StatusCode, error_type: &str, message: &str) -> Response { ( status, Json(json!({ "error": { "type": error_type, "message": message, "status": status.as_u16() } })), ) .into_response() } /// SSE streaming — прозрачно проксируем байты от engine. 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 engine stream"); None } } }); let mut builder = Response::builder() .status(status.as_u16()) .header(CONTENT_TYPE, "text/event-stream") .header("Cache-Control", "no-cache") .header("Connection", "keep-alive"); 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", ) }) }