mod.rs (5890B)
1 pub mod inventory; 2 pub mod orders; 3 pub mod products; 4 5 use std::time::Instant; 6 7 use axum::{ 8 body::Body, 9 extract::Request, 10 http::{HeaderValue, StatusCode}, 11 response::{IntoResponse, Response}, 12 }; 13 use opentelemetry::KeyValue; 14 15 use crate::metrics::metrics; 16 17 /// Generic proxy handler that forwards requests to a backend service 18 /// 19 /// This function handles: 20 /// - Path and query string forwarding 21 /// - Header forwarding (except Host header) 22 /// - Request body forwarding 23 /// - Response status, headers, and body forwarding 24 /// 25 /// # Arguments 26 /// * `service_url` - Base URL of the backend service (e.g., "http://products:3001") 27 /// * `req` - The incoming request to proxy 28 /// * `http_client` - HTTP client for making the proxy request 29 pub async fn proxy_request( 30 service_url: &str, 31 req: Request, 32 http_client: &reqwest_middleware::ClientWithMiddleware, 33 ) -> impl IntoResponse { 34 // Extract the path after /api prefix 35 let path = req.uri().path(); 36 let forwarded_path = path.strip_prefix("/api").unwrap_or(path); 37 let query = req 38 .uri() 39 .query() 40 .map(|q| format!("?{}", q)) 41 .unwrap_or_default(); 42 43 // Build target URL 44 let target_url = format!("{}{}{}", service_url, forwarded_path, query); 45 46 // Extract method and headers 47 let method = req.method().clone(); 48 let method_str = method.to_string(); 49 let headers = req.headers().clone(); 50 51 // Extract body 52 let body_bytes = match axum::body::to_bytes(req.into_body(), usize::MAX).await { 53 Ok(bytes) => bytes, 54 Err(_) => { 55 return Response::builder() 56 .status(StatusCode::BAD_REQUEST) 57 .body(Body::from("Failed to read request body")) 58 .unwrap(); 59 } 60 }; 61 62 // Build request to backend service 63 let mut client_req = http_client.request(method, &target_url); 64 65 // Forward headers (except Host header) 66 for (name, value) in headers.iter() { 67 if name != "host" { 68 if let Ok(val) = HeaderValue::from_bytes(value.as_bytes()) { 69 client_req = client_req.header(name.as_str(), val); 70 } 71 } 72 } 73 74 // Add body if present 75 if !body_bytes.is_empty() { 76 client_req = client_req.body(body_bytes.to_vec()); 77 } 78 79 // Send request to backend service and measure duration 80 // The TracingMiddleware automatically: 81 // 1. Creates a client span with HTTP semantic convention attributes 82 // 2. Injects trace context (traceparent/tracestate) headers 83 // 3. Records response status on span completion 84 let start = Instant::now(); 85 let response = client_req.send().await; 86 let duration = start.elapsed().as_secs_f64(); 87 88 // Derive the upstream service name from the target URL 89 let upstream_service = derive_upstream_service(service_url); 90 91 // Handle request failure 92 let response = match response { 93 Ok(resp) => resp, 94 Err(e) => { 95 // Record upstream request failure metric 96 metrics().upstream_request_duration.record( 97 duration, 98 &[ 99 KeyValue::new("upstream.service", upstream_service.clone()), 100 KeyValue::new("http.request.method", method_str), 101 KeyValue::new("error.type", "connection_error"), 102 ], 103 ); 104 105 // Log the downstream service failure with structured fields 106 tracing::error!( 107 error.r#type = "http_client", 108 error.message = %e, 109 downstream.service = %upstream_service, 110 downstream.url = %target_url, 111 downstream.duration_ms = format!("{:.1}", duration * 1000.0), 112 "Downstream service unavailable" 113 ); 114 115 return Response::builder() 116 .status(StatusCode::SERVICE_UNAVAILABLE) 117 .header("content-type", "application/json") 118 .body(Body::from(r#"{"error":"Service unavailable"}"#)) 119 .unwrap(); 120 } 121 }; 122 123 // Extract status and headers before consuming the response 124 let status = response.status(); 125 let resp_headers = response.headers().clone(); 126 127 // Record successful upstream request duration metric 128 metrics().upstream_request_duration.record( 129 duration, 130 &[ 131 KeyValue::new("upstream.service", upstream_service.clone()), 132 KeyValue::new("http.request.method", method_str), 133 KeyValue::new("http.response.status_code", i64::from(status.as_u16())), 134 ], 135 ); 136 137 // Log completed downstream call with response status 138 tracing::info!( 139 downstream.service = %upstream_service, 140 downstream.status = status.as_u16(), 141 downstream.duration_ms = format!("{:.1}", duration * 1000.0), 142 "Downstream call completed" 143 ); 144 145 // Get response body (this consumes the response) 146 let body = match response.bytes().await { 147 Ok(bytes) => bytes, 148 Err(_) => { 149 return Response::builder() 150 .status(StatusCode::INTERNAL_SERVER_ERROR) 151 .body(Body::from("Failed to read response body")) 152 .unwrap(); 153 } 154 }; 155 156 // Build response with forwarded headers 157 let mut builder = Response::builder().status(status); 158 for (name, value) in resp_headers.iter() { 159 if let Ok(val) = HeaderValue::from_bytes(value.as_bytes()) { 160 builder = builder.header(name.as_str(), val); 161 } 162 } 163 164 builder.body(Body::from(body)).unwrap() 165 } 166 167 /// Derives a short service name from a backend URL 168 /// (e.g. "http://products:3001" → "products"). 169 fn derive_upstream_service(url: &str) -> String { 170 url.trim_start_matches("http://") 171 .trim_start_matches("https://") 172 .split(":") 173 .next() 174 .unwrap_or("unknown") 175 .to_string() 176 }