diff --git a/Cargo.lock b/Cargo.lock index 96ae63e0c..c5504c8d0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6832,6 +6832,7 @@ dependencies = [ "mimalloc", "mime_guess", "moka", + "opentelemetry", "pin-project-lite", "pprof", "rand 0.10.0-rc.6", @@ -6888,6 +6889,7 @@ dependencies = [ "tower", "tower-http", "tracing", + "tracing-opentelemetry", "url", "urlencoding", "uuid", diff --git a/crates/obs/src/telemetry.rs b/crates/obs/src/telemetry.rs index 54780f7c3..af187fcae 100644 --- a/crates/obs/src/telemetry.rs +++ b/crates/obs/src/telemetry.rs @@ -21,6 +21,7 @@ use nu_ansi_term::Color; use opentelemetry::{KeyValue, global, trace::TracerProvider}; use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge; use opentelemetry_otlp::{Compression, Protocol, WithExportConfig, WithHttpConfig}; +use opentelemetry_sdk::propagation::TraceContextPropagator; use opentelemetry_sdk::{ Resource, logs::SdkLoggerProvider, @@ -439,6 +440,7 @@ fn init_observability_http(config: &OtelConfig, logger_level: &str, is_productio let provider = builder.build(); global::set_tracer_provider(provider.clone()); + global::set_text_map_propagator(TraceContextPropagator::new()); Some(provider) } }; diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 68ce4479b..6c2ebd94b 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -133,6 +133,8 @@ aes-gcm = { workspace = true } # Observability and Metrics metrics = { workspace = true } +opentelemetry = { workspace = true } +tracing-opentelemetry = { workspace = true } [target.'cfg(target_os = "linux")'.dependencies] libsystemd.workspace = true diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index cf586e661..8551e95a3 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -33,6 +33,7 @@ use hyper_util::{ service::TowerToHyperService, }; use metrics::{counter, histogram}; +use opentelemetry::global; use rustfs_common::GlobalReadiness; use rustfs_config::{RUSTFS_TLS_CERT, RUSTFS_TLS_KEY}; use rustfs_ecstore::rpc::{TONIC_RPC_PREFIX, verify_rpc_signature}; @@ -56,6 +57,7 @@ use tower_http::compression::CompressionLayer; use tower_http::request_id::{MakeRequestUuid, PropagateRequestIdLayer, SetRequestIdLayer}; use tower_http::trace::TraceLayer; use tracing::{Span, debug, error, info, instrument, warn}; +use tracing_opentelemetry::OpenTelemetrySpanExt; pub async fn start_http_server( opt: &config::Opt, @@ -524,6 +526,40 @@ struct ConnectionContext { readiness: Arc, } +/// Adapter that implements the OpenTelemetry [`Extractor`] trait for Hyper's +/// [`HeaderMap`], enabling trace context propagation by extracting +/// OpenTelemetry headers from incoming HTTP requests. +pub struct HeaderMapCarrier<'a> { + headers: &'a HeaderMap, +} + +impl<'a> HeaderMapCarrier<'a> { + pub fn new(headers: &'a HeaderMap) -> Self { + Self { headers } + } +} + +impl<'a> opentelemetry::propagation::Extractor for HeaderMapCarrier<'a> { + fn get(&self, key: &str) -> Option<&str> { + self.headers.get(key).and_then(|v| v.to_str().ok()) + } + + fn keys(&self) -> Vec<&str> { + self.headers.keys().map(|k| k.as_str()).collect() + } + + fn get_all(&self, key: &str) -> Option> { + let headers = self + .headers + .get_all(key) + .iter() + .filter_map(|value| value.to_str().ok()) + .collect::>(); + + if headers.is_empty() { None } else { Some(headers) } + } +} + /// Process a single incoming TCP connection. /// /// This function is executed in a new Tokio task, and it will: @@ -593,6 +629,10 @@ fn process_connection( .and_then(|v| v.to_str().ok()) .unwrap_or("unknown"); + let parent_context = global::get_text_map_propagator(|propagator| { + propagator.extract(&HeaderMapCarrier::new(request.headers())) + }); + // Extract real client IP from trusted proxy middleware if available let client_info = request.extensions().get::(); let real_ip = client_info @@ -607,6 +647,9 @@ fn process_connection( uri = %request.uri(), version = ?request.version(), ); + if let Err(e) = span.set_parent(parent_context) { + warn!("Failed to propagate tracing context: `{:?}`", e); + } for (header_name, header_value) in request.headers() { if header_name == "user-agent" || header_name == "content-type" || header_name == "content-length" { span.record(header_name.as_str(), header_value.to_str().unwrap_or("invalid")); @@ -835,3 +878,84 @@ fn get_default_tcp_keepalive() -> TcpKeepalive { .with_retries(3) } } + +#[cfg(test)] +mod tests { + use super::*; + use http::HeaderMap; + use opentelemetry::propagation::Extractor; + + #[test] + fn test_headermap_carrier_new() { + let headers = HeaderMap::new(); + let carrier = HeaderMapCarrier::new(&headers); + assert_eq!(carrier.keys().len(), 0); + } + + #[test] + fn test_headermap_carrier_get() { + let mut headers = HeaderMap::new(); + headers.insert("user-agent", "test-agent".parse().unwrap()); + headers.insert("x-request-id", "12345".parse().unwrap()); + + let carrier = HeaderMapCarrier::new(&headers); + + assert_eq!(carrier.get("user-agent"), Some("test-agent")); + assert_eq!(carrier.get("x-request-id"), Some("12345")); + assert_eq!(carrier.get("content-type"), None); + } + + #[test] + fn test_headermap_carrier_keys() { + let mut headers = HeaderMap::new(); + headers.insert("user-agent", "test-agent".parse().unwrap()); + headers.insert("content-type", "application/json".parse().unwrap()); + + let carrier = HeaderMapCarrier::new(&headers); + let keys = carrier.keys(); + + assert_eq!(keys.len(), 2); + assert!(keys.contains(&"user-agent")); + assert!(keys.contains(&"content-type")); + } + + #[test] + fn test_headermap_carrier_get_all() { + let mut headers = HeaderMap::new(); + headers.append("x-custom-header", "value1".parse().unwrap()); + headers.append("x-custom-header", "value2".parse().unwrap()); + headers.insert("user-agent", "test-agent".parse().unwrap()); + + let carrier = HeaderMapCarrier::new(&headers); + + // Test multi-value header + let values = carrier.get_all("x-custom-header"); + assert!(values.is_some()); + let v = values.unwrap(); + assert_eq!(v.len(), 2); + assert!(v.contains(&"value1")); + assert!(v.contains(&"value2")); + + // Test single value header + let values = carrier.get_all("user-agent"); + assert!(values.is_some()); + let v = values.unwrap(); + assert_eq!(v.len(), 1); + assert_eq!(v[0], "test-agent"); + + // Test missing header + assert_eq!(carrier.get_all("missing-header"), None); + } + + #[test] + fn test_headermap_carrier_case_insensitivity() { + let mut headers = HeaderMap::new(); + headers.insert("content-type", "application/json".parse().unwrap()); + + let carrier = HeaderMapCarrier::new(&headers); + + // HeaderMap::get is case insensitive + assert_eq!(carrier.get("Content-Type"), Some("application/json")); + assert_eq!(carrier.get("CONTENT-TYPE"), Some("application/json")); + } +}