diff --git a/Cargo.lock b/Cargo.lock index 988712641..2d4e3adf8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7870,6 +7870,7 @@ dependencies = [ "mimalloc", "mime_guess", "opentelemetry", + "opentelemetry_sdk", "percent-encoding", "pin-project-lite", "pprof-pyroscope-fork", @@ -7937,6 +7938,7 @@ dependencies = [ "tower-http", "tracing", "tracing-opentelemetry", + "tracing-subscriber", "url", "urlencoding", "uuid", @@ -8100,6 +8102,8 @@ dependencies = [ "memmap2 0.9.10", "metrics", "num_cpus", + "opentelemetry", + "opentelemetry_sdk", "parking_lot 0.12.5", "path-absolutize", "pin-project-lite", @@ -8146,6 +8150,7 @@ dependencies = [ "tonic", "tower", "tracing", + "tracing-opentelemetry", "tracing-subscriber", "url", "urlencoding", diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index 76d3e38c0..2a19d4488 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -55,6 +55,7 @@ flatbuffers.workspace = true futures.workspace = true futures-util.workspace = true tracing.workspace = true +tracing-opentelemetry.workspace = true serde.workspace = true time.workspace = true bytesize.workspace = true @@ -62,6 +63,7 @@ serde_json.workspace = true quick-xml = { workspace = true, features = ["serialize", "async-tokio"] } s3s.workspace = true http.workspace = true +opentelemetry.workspace = true http-body = { workspace = true } http-body-util.workspace = true url.workspace = true @@ -127,6 +129,7 @@ criterion = { workspace = true, features = ["html_reports"] } temp-env = { workspace = true } tracing-subscriber = { workspace = true } serial_test = { workspace = true } +opentelemetry_sdk = { workspace = true } [build-dependencies] shadow-rs = { workspace = true, features = ["build", "metadata"] } diff --git a/crates/ecstore/src/rpc/client.rs b/crates/ecstore/src/rpc/client.rs index a9830ef85..6314c6571 100644 --- a/crates/ecstore/src/rpc/client.rs +++ b/crates/ecstore/src/rpc/client.rs @@ -20,6 +20,8 @@ use std::error::Error; use tonic::{service::interceptor::InterceptedService, transport::Channel}; use tracing::debug; +use super::context_propagation::{inject_request_id_into_metadata, inject_trace_context_into_metadata}; + /// 3. Subsequent calls will attempt fresh connections /// 4. If node is still down, connection will fail fast (3s timeout) pub async fn node_service_time_out_client( @@ -55,6 +57,8 @@ impl tonic::service::Interceptor for TonicSignatureInterceptor { fn call(&mut self, mut req: tonic::Request<()>) -> Result, tonic::Status> { let headers = gen_signature_headers(TONIC_RPC_PREFIX, &Method::GET); req.metadata_mut().as_mut().extend(headers); + inject_trace_context_into_metadata(req.metadata_mut()); + inject_request_id_into_metadata(req.metadata_mut()); Ok(req) } } @@ -84,3 +88,34 @@ impl tonic::service::Interceptor for TonicInterceptor { } } } + +#[cfg(test)] +mod tests { + use super::*; + use tonic::service::Interceptor; + + #[test] + fn test_signature_interceptor_keeps_auth_headers() { + let mut interceptor = TonicSignatureInterceptor; + let req = tonic::Request::new(()); + + let req = interceptor.call(req).expect("interceptor call should succeed"); + + assert!(req.metadata().contains_key("x-rustfs-signature")); + assert!(req.metadata().contains_key("x-rustfs-timestamp")); + } + + #[test] + fn test_signature_interceptor_may_inject_request_id() { + let mut interceptor = TonicSignatureInterceptor; + let req = tonic::Request::new(()); + + let span = tracing::info_span!("grpc-rpc-test-span"); + let _guard = span.enter(); + let req = interceptor.call(req).expect("interceptor call should succeed"); + + if let Some(v) = req.metadata().get("x-request-id") { + assert!(!v.as_encoded_bytes().is_empty()); + } + } +} diff --git a/crates/ecstore/src/rpc/context_propagation.rs b/crates/ecstore/src/rpc/context_propagation.rs new file mode 100644 index 000000000..3ee1e806b --- /dev/null +++ b/crates/ecstore/src/rpc/context_propagation.rs @@ -0,0 +1,223 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use http::{HeaderMap, HeaderValue}; +use opentelemetry::{global, propagation::Injector, trace::TraceContextExt}; +use tracing::Span; +use tracing_opentelemetry::OpenTelemetrySpanExt; + +pub(crate) const REQUEST_ID_HEADER: &str = "x-request-id"; + +struct HttpHeaderInjector<'a> { + headers: &'a mut HeaderMap, +} + +impl Injector for HttpHeaderInjector<'_> { + fn set(&mut self, key: &str, value: String) { + let Ok(name) = http::header::HeaderName::from_bytes(key.as_bytes()) else { + return; + }; + let Ok(val) = HeaderValue::from_str(&value) else { + return; + }; + self.headers.insert(name, val); + } +} + +struct MetadataInjector<'a> { + metadata: &'a mut tonic::metadata::MetadataMap, +} + +impl Injector for MetadataInjector<'_> { + fn set(&mut self, key: &str, value: String) { + let Ok(meta_key) = tonic::metadata::MetadataKey::from_bytes(key.as_bytes()) else { + return; + }; + let Ok(meta_value) = tonic::metadata::MetadataValue::try_from(value.as_str()) else { + return; + }; + self.metadata.insert(meta_key, meta_value); + } +} + +fn current_trace_id() -> Option { + let current_context = Span::current().context(); + let current_span = current_context.span(); + let span_context = current_span.span_context(); + if !span_context.is_valid() { + return None; + } + Some(span_context.trace_id().to_string()) +} + +fn fallback_request_id() -> String { + format!("req-{}", &uuid::Uuid::new_v4().to_string()[..8]) +} + +fn propagated_request_id() -> String { + current_trace_id() + .map(|trace_id| format!("trace-{trace_id}")) + .unwrap_or_else(fallback_request_id) +} + +pub(crate) fn inject_trace_context_into_http_headers(headers: &mut HeaderMap) { + let current_context = Span::current().context(); + global::get_text_map_propagator(|propagator| { + let mut injector = HttpHeaderInjector { headers }; + propagator.inject_context(¤t_context, &mut injector); + }); +} + +pub(crate) fn inject_request_id_into_http_headers(headers: &mut HeaderMap) { + if headers.contains_key(REQUEST_ID_HEADER) { + return; + } + let request_id = propagated_request_id(); + if let Ok(value) = HeaderValue::from_str(&request_id) { + headers.insert(REQUEST_ID_HEADER, value); + } +} + +pub(crate) fn inject_trace_context_into_metadata(metadata: &mut tonic::metadata::MetadataMap) { + let current_context = Span::current().context(); + global::get_text_map_propagator(|propagator| { + let mut injector = MetadataInjector { metadata }; + propagator.inject_context(¤t_context, &mut injector); + }); +} + +pub(crate) fn inject_request_id_into_metadata(metadata: &mut tonic::metadata::MetadataMap) { + let request_id_key = tonic::metadata::MetadataKey::from_static(REQUEST_ID_HEADER); + if metadata.contains_key(&request_id_key) { + return; + } + let request_id = propagated_request_id(); + let Ok(value) = tonic::metadata::MetadataValue::try_from(request_id.as_str()) else { + return; + }; + metadata.insert(request_id_key, value); +} + +#[cfg(test)] +mod tests { + use super::*; + use opentelemetry::trace::{SpanContext, TraceContextExt, TraceFlags, TraceId, TraceState, TracerProvider as _}; + use opentelemetry_sdk::trace::SdkTracerProvider; + use tracing_opentelemetry::OpenTelemetrySpanExt; + use tracing_subscriber::{Registry, layer::SubscriberExt}; + + fn with_trace_parent(trace_id_hex: &str, f: F) + where + F: FnOnce(), + { + let provider = SdkTracerProvider::builder().build(); + let tracer = provider.tracer("context-propagation-tests"); + let subscriber = Registry::default().with(tracing_opentelemetry::layer().with_tracer(tracer)); + + tracing::subscriber::with_default(subscriber, || { + let span = tracing::info_span!("context-propagation-test-span"); + + let trace_id = TraceId::from_hex(trace_id_hex).expect("trace id should be valid hex"); + let span_id = opentelemetry::trace::SpanId::from_hex("0102030405060708").expect("span id should be valid hex"); + let parent = SpanContext::new(trace_id, span_id, TraceFlags::SAMPLED, true, TraceState::default()); + span.set_parent(opentelemetry::Context::new().with_remote_span_context(parent)) + .expect("failed to set parent context"); + let _guard = span.enter(); + + f(); + }); + let _ = provider.shutdown(); + } + + #[test] + fn test_inject_request_id_into_http_headers_preserves_existing_value() { + let mut headers = HeaderMap::new(); + headers.insert(REQUEST_ID_HEADER, HeaderValue::from_static("req-upstream-123")); + + with_trace_parent("0123456789abcdef0123456789abcdef", || { + inject_request_id_into_http_headers(&mut headers); + }); + + assert_eq!(headers.get(REQUEST_ID_HEADER).and_then(|v| v.to_str().ok()), Some("req-upstream-123")); + } + + #[test] + fn test_inject_request_id_into_http_headers_uses_trace_id_when_missing() { + let trace_id = "abcdefabcdefabcdefabcdefabcdefab"; + let mut headers = HeaderMap::new(); + + with_trace_parent(trace_id, || { + inject_request_id_into_http_headers(&mut headers); + }); + + assert_eq!( + headers.get(REQUEST_ID_HEADER).and_then(|v| v.to_str().ok()), + Some(format!("trace-{trace_id}").as_str()) + ); + } + + #[test] + fn test_inject_request_id_into_metadata_preserves_existing_value() { + let mut metadata = tonic::metadata::MetadataMap::new(); + metadata.insert( + tonic::metadata::MetadataKey::from_static(REQUEST_ID_HEADER), + tonic::metadata::MetadataValue::from_static("req-upstream-456"), + ); + + with_trace_parent("fedcba9876543210fedcba9876543210", || { + inject_request_id_into_metadata(&mut metadata); + }); + + assert_eq!(metadata.get(REQUEST_ID_HEADER).and_then(|v| v.to_str().ok()), Some("req-upstream-456")); + } + + #[test] + fn test_inject_request_id_into_metadata_uses_trace_id_when_missing() { + let trace_id = "1234567890abcdef1234567890abcdef"; + let mut metadata = tonic::metadata::MetadataMap::new(); + + with_trace_parent(trace_id, || { + inject_request_id_into_metadata(&mut metadata); + }); + + assert_eq!( + metadata.get(REQUEST_ID_HEADER).and_then(|v| v.to_str().ok()), + Some(format!("trace-{trace_id}").as_str()) + ); + } + + #[test] + fn test_inject_request_id_into_http_headers_uses_req_fallback_when_trace_missing() { + let mut headers = HeaderMap::new(); + inject_request_id_into_http_headers(&mut headers); + + let request_id = headers + .get(REQUEST_ID_HEADER) + .and_then(|v| v.to_str().ok()) + .expect("request id should be injected"); + assert!(request_id.starts_with("req-"), "expected req- fallback, got: {request_id}"); + } + + #[test] + fn test_inject_request_id_into_metadata_uses_req_fallback_when_trace_missing() { + let mut metadata = tonic::metadata::MetadataMap::new(); + inject_request_id_into_metadata(&mut metadata); + + let request_id = metadata + .get(REQUEST_ID_HEADER) + .and_then(|v| v.to_str().ok()) + .expect("request id should be injected"); + assert!(request_id.starts_with("req-"), "expected req- fallback, got: {request_id}"); + } +} diff --git a/crates/ecstore/src/rpc/http_auth.rs b/crates/ecstore/src/rpc/http_auth.rs index 5d69e2803..9f267d248 100644 --- a/crates/ecstore/src/rpc/http_auth.rs +++ b/crates/ecstore/src/rpc/http_auth.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::rpc::context_propagation::{inject_request_id_into_http_headers, inject_trace_context_into_http_headers}; use base64::Engine as _; use base64::engine::general_purpose; use hmac::{Hmac, KeyInit, Mac}; @@ -62,6 +63,8 @@ pub fn build_auth_headers(url: &str, method: &Method, headers: &mut HeaderMap) { let auth_headers = gen_signature_headers(url, method); headers.extend(auth_headers); + inject_trace_context_into_http_headers(headers); + inject_request_id_into_http_headers(headers); } pub fn gen_signature_headers(url: &str, method: &Method) -> HeaderMap { @@ -132,6 +135,7 @@ pub fn verify_rpc_signature(url: &str, method: &Method, headers: &HeaderMap) -> #[cfg(test)] mod tests { use super::*; + use crate::rpc::context_propagation::REQUEST_ID_HEADER; use http::{HeaderMap, Method}; use time::OffsetDateTime; @@ -210,6 +214,33 @@ mod tests { assert!((current_time - timestamp).abs() <= 1, "Timestamp should be close to current time"); } + #[test] + fn test_build_auth_headers_preserves_existing_request_id() { + let url = "http://example.com/api/test"; + let method = Method::GET; + let mut headers = HeaderMap::new(); + headers.insert(REQUEST_ID_HEADER, HeaderValue::from_static("req-upstream-123")); + + build_auth_headers(url, &method, &mut headers); + + assert_eq!(headers.get(REQUEST_ID_HEADER).and_then(|v| v.to_str().ok()), Some("req-upstream-123")); + } + + #[test] + fn test_build_auth_headers_may_set_request_id_from_trace_id() { + let url = "http://example.com/api/test"; + let method = Method::GET; + let mut headers = HeaderMap::new(); + + let span = tracing::info_span!("rpc-test-span"); + let _guard = span.enter(); + build_auth_headers(url, &method, &mut headers); + + if let Some(value) = headers.get(REQUEST_ID_HEADER).and_then(|v| v.to_str().ok()) { + assert!(!value.is_empty(), "request id should not be empty"); + } + } + #[test] fn test_verify_rpc_signature_success() { let url = "http://example.com/api/test"; diff --git a/crates/ecstore/src/rpc/mod.rs b/crates/ecstore/src/rpc/mod.rs index a59935534..f2c0c77c4 100644 --- a/crates/ecstore/src/rpc/mod.rs +++ b/crates/ecstore/src/rpc/mod.rs @@ -13,6 +13,7 @@ // limitations under the License. mod client; +mod context_propagation; mod http_auth; mod peer_rest_client; mod peer_s3_client; diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 16b008ee2..7edf316a8 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -195,6 +195,8 @@ aws-config = { workspace = true } anyhow = { workspace = true } tokio = { workspace = true, features = ["test-util"] } temp-env = { workspace = true, features = ["async_closure"] } +tracing-subscriber = { workspace = true } +opentelemetry_sdk = { workspace = true } [build-dependencies] http.workspace = true diff --git a/rustfs/src/admin/console.rs b/rustfs/src/admin/console.rs index 23c39392e..64cce7b7e 100644 --- a/rustfs/src/admin/console.rs +++ b/rustfs/src/admin/console.rs @@ -38,6 +38,7 @@ use tower_http::catch_panic::CatchPanicLayer; use tower_http::compression::CompressionLayer; use tower_http::cors::{AllowOrigin, Any, CorsLayer}; use tower_http::limit::RequestBodyLimitLayer; +use tower_http::request_id::{MakeRequestUuid, PropagateRequestIdLayer, SetRequestIdLayer}; use tower_http::timeout::TimeoutLayer; use tower_http::trace::TraceLayer; use tracing::{debug, error, info, instrument, warn}; @@ -379,9 +380,16 @@ async fn console_logging_middleware(req: Request, next: middleware::Next) -> Res let start = std::time::Instant::now(); let response = next.run(req).await; let duration = start.elapsed(); + let request_id = response + .headers() + .get("x-request-id") + .and_then(|v| v.to_str().ok()) + .unwrap_or("unknown") + .to_string(); info!( target: "rustfs::console::access", + request_id = %request_id, method = %method, uri = %uri, status = %response.status(), @@ -463,6 +471,8 @@ fn setup_console_middleware_stack( // Add comprehensive middleware layers using tower-http features app = app .layer(CatchPanicLayer::new()) + .layer(PropagateRequestIdLayer::x_request_id()) + .layer(SetRequestIdLayer::x_request_id(MakeRequestUuid)) .layer(TraceLayer::new_for_http()) // Compress responses .layer(CompressionLayer::new()) @@ -654,7 +664,10 @@ pub(crate) fn make_console_server() -> Router { #[cfg(test)] mod tests { use super::*; + use axum::body::Body; + use http::{Request, StatusCode}; use std::net::{IpAddr, Ipv4Addr}; + use tower::ServiceExt; #[test] fn console_api_base_url_keeps_rustfs_admin_prefix() { @@ -684,4 +697,20 @@ mod tests { assert!(!is_console_path("/minio/admin/v3/info")); assert!(!is_console_path("/rustfs/admin/v3/info")); } + + #[tokio::test] + async fn console_middleware_stack_propagates_request_id_header() { + let app = setup_console_middleware_stack(parse_cors_origins(None), false, 0, 30); + let request = Request::builder() + .uri(format!("{CONSOLE_PREFIX}{HEALTH_PREFIX}")) + .body(Body::empty()) + .expect("failed to build request"); + + let response = app.oneshot(request).await.expect("request should succeed"); + assert_eq!(response.status(), StatusCode::OK); + assert!( + response.headers().contains_key("x-request-id"), + "console response should include propagated x-request-id header" + ); + } } diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index b4e584bfa..bfe053c15 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -23,7 +23,7 @@ use crate::storage::concurrency::{ }; use crate::storage::ecfs::*; use crate::storage::head_prefix::{head_prefix_not_found_message, probe_prefix_has_children}; -use crate::storage::helper::{OperationHelper, spawn_background, spawn_background_with_context}; +use crate::storage::helper::{OperationHelper, spawn_background_with_context}; use crate::storage::options::{ copy_dst_opts, copy_src_opts, del_opts, extract_metadata, extract_metadata_from_mime_with_object_name, filter_object_metadata, get_content_sha256_with_query, get_opts, normalize_content_encoding_for_storage, put_opts, @@ -3105,7 +3105,11 @@ impl DefaultObjectUsecase { .as_ref() .map(|context| context.notify()) .unwrap_or_else(default_notify_interface); - spawn_background(async move { + let request_context = req + .extensions + .get::() + .cloned(); + spawn_background_with_context(request_context, async move { for res in delete_results { if let Some(dobj) = res.delete_object { let event_name = if dobj.delete_marker { diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 437b099f1..d008af389 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -522,6 +522,45 @@ impl<'a> opentelemetry::propagation::Extractor for HeaderMapCarrier<'a> { } } +/// Adapter that implements the OpenTelemetry [`Extractor`] trait for gRPC +/// metadata maps so internode gRPC requests can continue distributed traces. +struct MetadataMapCarrier<'a> { + metadata: &'a tonic::metadata::MetadataMap, +} + +impl<'a> MetadataMapCarrier<'a> { + fn new(metadata: &'a tonic::metadata::MetadataMap) -> Self { + Self { metadata } + } +} + +impl<'a> opentelemetry::propagation::Extractor for MetadataMapCarrier<'a> { + fn get(&self, key: &str) -> Option<&str> { + self.metadata.get(key).and_then(|v| v.to_str().ok()) + } + + fn keys(&self) -> Vec<&str> { + self.metadata + .keys() + .filter_map(|key| match key { + tonic::metadata::KeyRef::Ascii(v) => Some(v.as_str()), + tonic::metadata::KeyRef::Binary(_) => None, + }) + .collect() + } + + fn get_all(&self, key: &str) -> Option> { + let values = self + .metadata + .get_all(key) + .iter() + .filter_map(|value| value.to_str().ok()) + .collect::>(); + + if values.is_empty() { None } else { Some(values) } + } +} + /// Process a single incoming TCP connection. /// /// This function is executed in a new Tokio task, and it will: @@ -845,6 +884,21 @@ fn check_auth(req: Request<()>) -> std::result::Result, Status> { error!("RPC signature verification failed: {}", e); Status::unauthenticated("No valid auth token") })?; + + let parent_context = + global::get_text_map_propagator(|propagator| propagator.extract(&MetadataMapCarrier::new(req.metadata()))); + if parent_context.has_active_span() { + let span_ref = parent_context.span(); + debug!( + otel_trace_id = %span_ref.span_context().trace_id(), + otel_parent_span_id = %span_ref.span_context().span_id(), + sampled = span_ref.span_context().is_sampled(), + "Extracted trace context from incoming gRPC metadata" + ); + if let Err(e) = tracing::Span::current().set_parent(parent_context) { + warn!("Failed to propagate tracing context from gRPC metadata: `{:?}`", e); + } + } Ok(req) } diff --git a/rustfs/src/storage/helper.rs b/rustfs/src/storage/helper.rs index 53bf7fbb2..f1593897e 100644 --- a/rustfs/src/storage/helper.rs +++ b/rustfs/src/storage/helper.rs @@ -16,6 +16,7 @@ use crate::storage::access::{ReqInfo, request_context_from_req}; use crate::storage::request_context::{RequestContext, extract_request_id_from_headers}; use hashbrown::HashMap; use http::StatusCode; +use metrics::counter; use rustfs_audit::{ entity::{ApiDetails, ApiDetailsBuilder, AuditEntryBuilder}, global::AuditLogger, @@ -80,6 +81,7 @@ pub struct OperationHelper { impl OperationHelper { /// Create a new OperationHelper for S3 requests. pub fn new(req: &S3Request, event: EventName, op: S3Operation) -> Self { + counter!("rustfs.log.chain.audit.total").increment(1); // Parse path -> bucket/object let path = req.uri.path().trim_start_matches('/'); let mut segs = path.splitn(2, '/'); @@ -119,8 +121,11 @@ impl OperationHelper { } // Audit builder // Resolve canonical request context and request_id in a single pass: - // RequestContext.request_id > extract_request_id_from_headers() > "unknown" + // RequestContext.request_id > extract_request_id_from_headers() > generated fallback id let request_context = request_context_from_req(req); + if request_context.is_none() { + counter!("rustfs.log.chain.orphan.total", "component" => "operation_helper").increment(1); + } let request_id = request_context .as_ref() .map(|ctx| ctx.request_id.clone()) diff --git a/rustfs/src/storage/request_context.rs b/rustfs/src/storage/request_context.rs index 9e4514282..a2a088cf8 100644 --- a/rustfs/src/storage/request_context.rs +++ b/rustfs/src/storage/request_context.rs @@ -43,6 +43,7 @@ //! - Canonical source: HTTP ingress `x-request-id` header (set by `SetRequestIdLayer`) //! - `x-amz-request_id` is an alias for S3 compatibility, always equal to `request_id` //! - Internal modules MUST NOT generate a second request-id under the name `request_id` +//! except for orphan/non-ingress fallback paths where no canonical request-id exists. //! - Internal identifiers for sub-operations should use `operation_id` or `subtask_id` //! //! ## tokio::spawn usage @@ -56,8 +57,14 @@ //! - NEVER use bare `tokio::spawn` in request-handling code paths use http::HeaderMap; +use metrics::counter; +use opentelemetry::trace::TraceContextExt; use rustfs_utils::http::headers::AMZ_REQUEST_ID; use std::time::Instant; +use tracing::Span; +use tracing_opentelemetry::OpenTelemetrySpanExt; + +const REQUEST_ID_HEADER: &str = "x-request-id"; /// Canonical request context carried through the entire request lifecycle. /// @@ -79,32 +86,62 @@ pub struct RequestContext { impl RequestContext { /// Create a fallback `RequestContext` for paths that bypass HTTP ingress. - /// Generates a `req-{uuid}` format request-id. + /// Generates a `trace-{trace_id}` or `req-{uuid}` format request-id. pub fn fallback() -> Self { - let id = format!("req-{}", &uuid::Uuid::new_v4().to_string()[..8]); + let trace_ctx = current_trace_context_ids(); + let id = build_fallback_request_id(trace_ctx.as_ref()); + counter!("rustfs.log.chain.fallback_request_id.total", "source" => "request_context_fallback").increment(1); Self { request_id: id.clone(), x_amz_request_id: id, - trace_id: None, - span_id: None, + trace_id: trace_ctx.as_ref().map(|(trace_id, _)| trace_id.clone()), + span_id: trace_ctx.as_ref().map(|(_, span_id)| span_id.clone()), start_time: Instant::now(), } } } +fn current_trace_context_ids() -> Option<(String, String)> { + let current_context = Span::current().context(); + let current_span = current_context.span(); + let span_context = current_span.span_context(); + if !span_context.is_valid() { + return None; + } + + Some((span_context.trace_id().to_string(), span_context.span_id().to_string())) +} + +fn build_fallback_request_id(trace_ctx: Option<&(String, String)>) -> String { + trace_ctx + .map(|(trace_id, _)| format!("trace-{trace_id}")) + .unwrap_or_else(|| format!("req-{}", &uuid::Uuid::new_v4().to_string()[..8])) +} + +fn generate_fallback_request_id() -> String { + let trace_ctx = current_trace_context_ids(); + build_fallback_request_id(trace_ctx.as_ref()) +} + /// Extract the canonical request ID from HTTP headers. /// /// Priority: /// 1. `x-request-id` (primary, set by `SetRequestIdLayer`) /// 2. `x-amz-request-id` (fallback, from S3 client forwarding) -/// 3. `"unknown"` (no header present) +/// 3. generated fallback id (`trace-{trace_id}` or `req-{uuid}`) pub fn extract_request_id_from_headers(headers: &HeaderMap) -> String { - headers - .get("x-request-id") + let request_id = headers + .get(REQUEST_ID_HEADER) .and_then(|v| v.to_str().ok()) .map(String::from) .or_else(|| headers.get(AMZ_REQUEST_ID).and_then(|v| v.to_str().ok()).map(String::from)) - .unwrap_or_else(|| "unknown".to_string()) + .unwrap_or_else(generate_fallback_request_id); + + if !headers.contains_key(REQUEST_ID_HEADER) && !headers.contains_key(AMZ_REQUEST_ID) { + counter!("rustfs.log.chain.fallback_request_id.total", "source" => "headers_missing").increment(1); + } + + request_id } /// Spawn a request-internal task that inherits the current tracing span. @@ -126,6 +163,33 @@ where #[cfg(test)] mod tests { use super::*; + use opentelemetry::trace::{SpanContext, TraceContextExt, TraceFlags, TraceId, TraceState, TracerProvider as _}; + use opentelemetry_sdk::trace::SdkTracerProvider; + use tracing_opentelemetry::OpenTelemetrySpanExt; + use tracing_subscriber::{Registry, layer::SubscriberExt}; + + fn with_trace_parent(trace_id_hex: &str, f: F) + where + F: FnOnce(), + { + let provider = SdkTracerProvider::builder().build(); + let tracer = provider.tracer("request-context-tests"); + let subscriber = Registry::default().with(tracing_opentelemetry::layer().with_tracer(tracer)); + + tracing::subscriber::with_default(subscriber, || { + let span = tracing::info_span!("request-context-test-span"); + + let trace_id = TraceId::from_hex(trace_id_hex).expect("trace id should be valid hex"); + let span_id = opentelemetry::trace::SpanId::from_hex("0102030405060708").expect("span id should be valid hex"); + let parent = SpanContext::new(trace_id, span_id, TraceFlags::SAMPLED, true, TraceState::default()); + span.set_parent(opentelemetry::Context::new().with_remote_span_context(parent)) + .expect("failed to set parent context"); + let _guard = span.enter(); + + f(); + }); + let _ = provider.shutdown(); + } #[test] fn test_request_context_clone_send_sync() { @@ -142,6 +206,17 @@ mod tests { assert!(ctx.span_id.is_none()); } + #[test] + fn test_request_context_fallback_uses_trace_prefix_when_span_context_valid() { + let trace_id = "70f5f77e2f0a4f24be343b59f8b66f8f"; + with_trace_parent(trace_id, || { + let ctx = RequestContext::fallback(); + assert_eq!(ctx.request_id, format!("trace-{trace_id}")); + assert_eq!(ctx.trace_id.as_deref(), Some(trace_id)); + assert!(ctx.span_id.is_some()); + }); + } + #[test] fn test_extract_request_id_from_x_request_id() { let mut headers = HeaderMap::new(); @@ -171,6 +246,20 @@ mod tests { fn test_extract_request_id_no_headers() { let headers = HeaderMap::new(); let id = extract_request_id_from_headers(&headers); - assert_eq!(id, "unknown"); + assert!( + id.starts_with("req-") || id.starts_with("trace-"), + "fallback request id should use req-/trace- prefix, got: {}", + id + ); + } + + #[test] + fn test_extract_request_id_no_headers_uses_trace_prefix_when_span_context_valid() { + let trace_id = "8d8b7d58055d45f793b8ca7fcb91bc17"; + with_trace_parent(trace_id, || { + let headers = HeaderMap::new(); + let id = extract_request_id_from_headers(&headers); + assert_eq!(id, format!("trace-{trace_id}")); + }); } } diff --git a/scripts/run.sh b/scripts/run.sh index e392d2e98..2000c444a 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -55,7 +55,7 @@ export RUSTFS_CONSOLE_ADDRESS=":9001" # export RUSTFS_TLS_PATH="./deploy/certs" # Observability related configuration -#export RUSTFS_OBS_ENDPOINT=http://localhost:4318 # OpenTelemetry Collector address +export RUSTFS_OBS_ENDPOINT=http://localhost:4318 # OpenTelemetry Collector address # RustFS OR OTEL exporter configuration #export RUSTFS_OBS_TRACE_ENDPOINT=http://localhost:4318/v1/traces # OpenTelemetry Collector trace address http://localhost:4318/v1/traces #export OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://localhost:14318/v1/traces