exercises

Unnamed repository; edit this file 'description' to name the repository.
Log | Files | Refs | README

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 }