From eabbea46d384b3ca172409c122818d08e7f5c362 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 4 Apr 2026 12:20:59 +0800 Subject: [PATCH] refactor(tracing): unify request-context propagation and fix tracing chain breaks (#2394) Co-authored-by: heihutu --- rustfs/src/admin/router.rs | 2 + rustfs/src/app/bucket_usecase.rs | 8 +- rustfs/src/app/multipart_usecase.rs | 3 +- rustfs/src/app/object_usecase.rs | 38 +++--- rustfs/src/protocols/client.rs | 1 + rustfs/src/server/http.rs | 39 ++++-- rustfs/src/server/layer.rs | 115 +++++++++++++++++ rustfs/src/storage/access.rs | 14 ++ rustfs/src/storage/helper.rs | 135 ++++++++++++++++++-- rustfs/src/storage/mod.rs | 1 + rustfs/src/storage/request_context.rs | 176 ++++++++++++++++++++++++++ rustfs/src/storage/timeout_wrapper.rs | 17 ++- 12 files changed, 499 insertions(+), 50 deletions(-) create mode 100644 rustfs/src/storage/request_context.rs diff --git a/rustfs/src/admin/router.rs b/rustfs/src/admin/router.rs index 3551c4365..f25441187 100644 --- a/rustfs/src/admin/router.rs +++ b/rustfs/src/admin/router.rs @@ -1442,6 +1442,7 @@ async fn authorize_replication_extension_request(req: &mut S3Request, ext_ object: None, version_id: None, region: get_global_region(), + ..Default::default() }); license_check().map_err(|er| match er.kind() { @@ -2163,6 +2164,7 @@ async fn authorize_misc_extension_request(req: &mut S3Request, route: &Mis object, version_id: None, region: get_global_region(), + ..Default::default() }); license_check().map_err(|er| match er.kind() { diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index 64e963f82..e23e5ef91 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -22,7 +22,7 @@ use crate::auth::get_condition_values; use crate::error::ApiError; use crate::server::RemoteAddr; use crate::storage::access::{ReqInfo, authorize_request, req_info_ref}; -use crate::storage::helper::OperationHelper; +use crate::storage::helper::{OperationHelper, spawn_background_with_context}; use crate::storage::s3_api::bucket::{build_list_buckets_output, build_list_objects_v2_output}; use crate::storage::s3_api::common::rustfs_owner; use crate::storage::s3_api::{acl, encryption, replication, tagging}; @@ -1494,7 +1494,11 @@ impl DefaultBucketUsecase { && let Some(store) = new_object_layer_fn() { let bucket_name = bucket.clone(); - tokio::spawn(async move { + let request_context = req + .extensions + .get::() + .cloned(); + spawn_background_with_context(request_context, async move { if let Err(err) = enqueue_transition_for_existing_objects(store, &bucket_name).await { warn!(bucket = %bucket_name, error = ?err, "failed to enqueue transition for existing objects"); } diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index fc1e069d4..0f97f1fe1 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -25,6 +25,7 @@ use crate::storage::options::{ copy_src_opts, extract_metadata, get_complete_multipart_upload_opts, get_content_sha256_with_query, get_opts, parse_copy_source_range, put_opts, }; +use crate::storage::request_context::spawn_traced; use crate::storage::s3_api::multipart::build_list_parts_output; use crate::storage::*; use bytes::Bytes; @@ -406,7 +407,7 @@ impl DefaultMultipartUsecase { }; let mpu_version_clone = mpu_version.clone(); let mpu_version_for_event = mpu_version.clone(); - tokio::spawn(async move { + spawn_traced(async move { manager .invalidate_cache_versioned(&mpu_bucket, &mpu_key, mpu_version_clone.as_deref()) .await; diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index f97b3519c..efebe94e4 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -24,11 +24,12 @@ 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}; +use crate::storage::helper::{OperationHelper, spawn_background, 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, }; +use crate::storage::request_context::spawn_traced; use crate::storage::s3_api::multipart::parse_list_parts_params; use crate::storage::s3_api::{acl, restore, select}; use crate::storage::timeout_wrapper::{RequestTimeoutWrapper, TimeoutConfig}; @@ -928,7 +929,7 @@ impl DefaultObjectUsecase { fn spawn_cache_invalidation(bucket: String, key: String, version_id: Option) { let manager = get_concurrency_manager(); - tokio::spawn(async move { + spawn_traced(async move { manager.invalidate_cache_versioned(&bucket, &key, version_id.as_deref()).await; }); } @@ -1015,9 +1016,9 @@ impl DefaultObjectUsecase { ))) } - fn init_get_object_bootstrap(bucket: &str, key: &str) -> S3Result { + fn init_get_object_bootstrap(bucket: &str, key: &str, request_id: &str) -> S3Result { let timeout_config = TimeoutConfig::from_env(); - let wrapper = RequestTimeoutWrapper::with_request_id(timeout_config.clone(), format!("get-{bucket}-{key}")); + let wrapper = RequestTimeoutWrapper::with_request_id(timeout_config.clone(), request_id.to_string()); let request_start = std::time::Instant::now(); let request_guard = ConcurrencyManager::track_request(); let concurrent_requests = GetObjectGuard::concurrent_requests(); @@ -1535,7 +1536,7 @@ impl DefaultObjectUsecase { .with_last_modified(last_modified_str.unwrap_or_default()); let cache_key_clone = cache_key.to_string(); - tokio::spawn(async move { + spawn_traced(async move { let manager = get_concurrency_manager(); manager.put_cached_object(cache_key_clone.clone(), cached_response).await; debug!("Object cached successfully with metadata: {}", cache_key_clone); @@ -2369,7 +2370,7 @@ impl DefaultObjectUsecase { let cache_key = ConcurrencyManager::make_cache_key(&bucket, &object, version_id.clone().as_deref()); let cache_bucket = bucket.clone(); let cache_object = object.clone(); - tokio::spawn(async move { + spawn_traced(async move { manager .invalidate_cache_versioned(&cache_bucket, &cache_object, version_id.as_deref()) .await; @@ -2599,7 +2600,12 @@ impl DefaultObjectUsecase { let _ = context.object_store(); } - let bootstrap = Self::init_get_object_bootstrap(&req.input.bucket, &req.input.key)?; + let request_id = req + .extensions + .get::() + .map(|ctx| ctx.request_id.clone()) + .unwrap_or_else(|| crate::storage::request_context::RequestContext::fallback().request_id); + let bootstrap = Self::init_get_object_bootstrap(&req.input.bucket, &req.input.key, &request_id)?; let timeout_config = bootstrap.timeout_config; let wrapper = bootstrap.wrapper; let request_start = bootstrap.request_start; @@ -3711,7 +3717,7 @@ impl DefaultObjectUsecase { let manager = get_concurrency_manager(); let bucket_clone = bucket.clone(); let deleted_objects = dobjs.clone(); - tokio::spawn(async move { + spawn_traced(async move { for dobj in deleted_objects { manager .invalidate_cache_versioned( @@ -4114,7 +4120,7 @@ impl DefaultObjectUsecase { let version_id_clone = version_id.clone(); let cache_bucket = bucket.clone(); let cache_object = object.clone(); - tokio::spawn(async move { + spawn_traced(async move { manager .invalidate_cache_versioned(&cache_bucket, &cache_object, version_id_clone.as_deref()) .await; @@ -4626,7 +4632,7 @@ impl DefaultObjectUsecase { let rreq_clone = rreq.clone(); let version_id_clone = version_id.clone(); - tokio::spawn(async move { + spawn_traced(async move { let opts = ObjectOptions { transition: TransitionOptions { restore_request: rreq_clone, @@ -4647,8 +4653,6 @@ impl DefaultObjectUsecase { object_clone, err.to_string() ); - // Note: Errors from background tasks cannot be returned to client - // Consider adding to monitoring/metrics system } else { info!("successfully restored transitioned object: {}/{}", bucket_clone, object_clone); } @@ -4721,7 +4725,7 @@ impl DefaultObjectUsecase { let (tx, rx) = mpsc::channel::>(2); let stream = ReceiverStream::new(rx); - tokio::spawn(async move { + spawn_traced(async move { let _ = tx .send(Ok(SelectObjectContentEvent::Cont(ContinuationEvent::default()))) .await; @@ -5078,7 +5082,7 @@ impl DefaultObjectUsecase { let manager = get_concurrency_manager(); let fpath_clone = fpath.clone(); let bucket_clone = bucket.clone(); - tokio::spawn(async move { + spawn_traced(async move { manager.invalidate_cache_versioned(&bucket_clone, &fpath_clone, None).await; }); @@ -5102,7 +5106,11 @@ impl DefaultObjectUsecase { }; let notify = notify.clone(); - tokio::spawn(async move { + let request_context = req + .extensions + .get::() + .cloned(); + spawn_background_with_context(request_context, async move { notify.notify(event_args).await; }); } diff --git a/rustfs/src/protocols/client.rs b/rustfs/src/protocols/client.rs index 65347d84f..766d3c147 100644 --- a/rustfs/src/protocols/client.rs +++ b/rustfs/src/protocols/client.rs @@ -81,6 +81,7 @@ impl ProtocolStorageClient { object: params.object, version_id: None, region: None, + request_context: Some(crate::storage::request_context::RequestContext::fallback()), }); let req = S3Request { diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 0fdcf3d1d..58a5b7af6 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -21,7 +21,10 @@ use crate::server::{ ReadinessGateLayer, RemoteAddr, ServiceState, ServiceStateManager, compress::{CompressionConfig, PathAwareCompressionPredicate, PathCategoryInjectionLayer}, hybrid::hybrid, - layer::{AdminChunkedContentLengthCompatLayer, ConditionalCorsLayer, ObjectAttributesEtagFixLayer, RedirectLayer}, + layer::{ + AdminChunkedContentLengthCompatLayer, ConditionalCorsLayer, ObjectAttributesEtagFixLayer, RedirectLayer, + RequestContextLayer, + }, tls_material::{TlsAcceptorHolder, TlsHandshakeFailureKind, TlsMaterialSnapshot, spawn_reload_loop}, }; use crate::storage; @@ -593,17 +596,18 @@ fn process_connection( // 2. AddExtensionLayer — per-connection raw socket addr (TrustedProxy) // 3. TrustedProxyLayer — conditional, parses X-Forwarded-For // 4. SetRequestIdLayer — generates X-Request-ID - // 5. AdminChunkedContentLengthCompatLayer — admin API compat - // 6. CatchPanicLayer — panic → 500 - // 7. ReadinessGateLayer — blocks until ready - // 8. KeystoneAuthLayer — X-Auth-Token validation - // 9. TraceLayer — request/response tracing + metrics - // 10. PropagateRequestIdLayer — X-Request-ID → response - // 11. PathCategoryInjectionLayer — injects path category for compression - // 12. CompressionLayer — response compression (whitelist, path-aware) - // 13. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes - // 14. ConditionalCorsLayer — S3 API CORS - // 15. RedirectLayer — console redirect (conditional) + // 5. RequestContextLayer — creates RequestContext in extensions + // 6. AdminChunkedContentLengthCompatLayer — admin API compat + // 7. CatchPanicLayer — panic → 500 + // 8. ReadinessGateLayer — blocks until ready + // 9. KeystoneAuthLayer — X-Auth-Token validation + // 10. TraceLayer — request/response tracing + metrics + // 11. PropagateRequestIdLayer — X-Request-ID → response + // 12. PathCategoryInjectionLayer — injects path category for compression + // 13. CompressionLayer — response compression (whitelist, path-aware) + // 14. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes + // 15. ConditionalCorsLayer — S3 API CORS + // 16. RedirectLayer — console redirect (conditional) // ───────────────────────────────────────────────────────────── let hybrid_service = ServiceBuilder::new() // NOTE: Both extension types are intentionally inserted to maintain compatibility: @@ -619,6 +623,7 @@ fn process_connection( // Pre-computed in ConnectionContext to avoid per-connection is_enabled() check. .option_layer(trusted_proxy_layer) .layer(SetRequestIdLayer::x_request_id(MakeRequestUuid)) + .layer(RequestContextLayer) .layer(AdminChunkedContentLengthCompatLayer) .layer(CatchPanicLayer::new()) // CRITICAL: Insert ReadinessGateLayer before business logic @@ -687,6 +692,12 @@ fn process_connection( debug!("http started method: {}, url path: {}", request.method(), request.uri().path()); let labels = [("key_request_method", request.method().to_string())]; counter!("rustfs.api.requests.total", &labels).increment(1); + // Aggregate request body size for throughput monitoring (lightweight) + if let Some(cl) = request.headers().get("content-length") + && let Some(len) = cl.to_str().ok().and_then(|s| s.parse::().ok()) + { + counter!("rustfs.request.body.bytes_total", "direction" => "request").increment(len); + } }) .on_response(|response: &Response<_>, latency: Duration, span: &Span| { span.record("status_code", tracing::field::display(response.status())); @@ -695,6 +706,8 @@ fn process_connection( debug!("http response generated in {:?}", latency) }) .on_body_chunk(|chunk: &Bytes, latency: Duration, span: &Span| { + // Always track aggregate body bytes (lightweight counter, no debug logging) + counter!("rustfs.request.body.bytes_total", "direction" => "response").increment(chunk.len() as u64); #[cfg(feature = "tracing-chunk-debug")] { let _enter = span.enter(); @@ -703,7 +716,7 @@ fn process_connection( } #[cfg(not(feature = "tracing-chunk-debug"))] { - let _ = (chunk, latency, span); + let _ = (latency, span); } }) .on_eos(|_trailers: Option<&HeaderMap>, stream_duration: Duration, span: &Span| { diff --git a/rustfs/src/server/layer.rs b/rustfs/src/server/layer.rs index da2c95046..13a198141 100644 --- a/rustfs/src/server/layer.rs +++ b/rustfs/src/server/layer.rs @@ -17,19 +17,134 @@ use crate::server::cors; use crate::server::hybrid::HybridBody; use crate::server::{ADMIN_PREFIX, CONSOLE_PREFIX, MINIO_ADMIN_PREFIX, MINIO_ADMIN_V3_PREFIX, RPC_PREFIX, RUSTFS_ADMIN_PREFIX}; use crate::storage::apply_cors_headers; +use crate::storage::request_context::{RequestContext, extract_request_id_from_headers}; use bytes::Bytes; use http::{HeaderMap, HeaderValue, Method, Request as HttpRequest, Response, StatusCode}; use http_body::Body; use http_body_util::BodyExt; use hyper::body::Incoming; +use opentelemetry::global; +use opentelemetry::trace::TraceContextExt; use rustfs_utils::get_env_opt_str; +use rustfs_utils::http::headers::AMZ_REQUEST_ID; use std::future::Future; use std::pin::Pin; use std::sync::Arc; use std::task::{Context, Poll}; +use std::time::Instant; use tower::{Layer, Service}; use tracing::debug; +/// A carrier that adapts [`HeaderMap`] for OpenTelemetry trace context propagation. +struct HeaderMapCarrier<'a>(&'a HeaderMap); + +impl<'a> opentelemetry::propagation::Extractor for HeaderMapCarrier<'a> { + fn get(&self, key: &str) -> Option<&str> { + self.0.get(key).and_then(|v| v.to_str().ok()) + } + + fn keys(&self) -> Vec<&str> { + self.0.keys().map(|k| k.as_str()).collect() + } + + fn get_all(&self, key: &str) -> Option> { + let headers = self + .0 + .get_all(key) + .iter() + .filter_map(|value| value.to_str().ok()) + .collect::>(); + + if headers.is_empty() { None } else { Some(headers) } + } +} + +/// Tower middleware layer that creates a canonical [`RequestContext`] from HTTP headers +/// and injects it into `request.extensions()`. +/// +/// This layer must be placed after `SetRequestIdLayer` in the middleware stack, +/// as it reads the `x-request-id` header that `SetRequestIdLayer` generates. +/// +/// Additionally, it sets the `x-amz-request-id` request header for S3 compatibility +/// if not already present. +#[derive(Clone, Default)] +pub struct RequestContextLayer; + +impl Layer for RequestContextLayer { + type Service = RequestContextService; + + fn layer(&self, inner: S) -> Self::Service { + RequestContextService { inner } + } +} + +/// Service that injects [`RequestContext`] into every request. +#[derive(Clone)] +pub struct RequestContextService { + inner: S, +} + +impl Service> for RequestContextService +where + S: Service>, +{ + type Response = S::Response; + type Error = S::Error; + type Future = S::Future; + + fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll> { + self.inner.poll_ready(cx) + } + + fn call(&mut self, mut req: HttpRequest) -> Self::Future { + let request_id = extract_request_id_from_headers(req.headers()); + + // Extract OpenTelemetry trace/span context from incoming headers + let parent_cx = global::get_text_map_propagator(|propagator| propagator.extract(&HeaderMapCarrier(req.headers()))); + let span_ref = parent_cx.span(); + let span_context = span_ref.span_context(); + let trace_id = if span_context.is_valid() { + Some(span_context.trace_id().to_string()) + } else { + None + }; + let span_id = if span_context.is_valid() { + Some(span_context.span_id().to_string()) + } else { + None + }; + + // Preserve the upstream x-amz-request-id if present (S3 client forwarding), + // otherwise fall back to the canonical request_id. + let x_amz_request_id = req + .headers() + .get(AMZ_REQUEST_ID) + .and_then(|v| v.to_str().ok()) + .map(String::from) + .unwrap_or_else(|| request_id.clone()); + + let ctx = RequestContext { + request_id: request_id.clone(), + x_amz_request_id, + trace_id, + span_id, + start_time: Instant::now(), + }; + + req.extensions_mut().insert(ctx); + + // Set x-amz-request-id for S3 compatibility downstream + if !req.headers().contains_key(AMZ_REQUEST_ID) + && let Ok(val) = HeaderValue::from_str(&request_id) + { + req.headers_mut() + .insert(http::header::HeaderName::from_static(AMZ_REQUEST_ID), val); + } + + self.inner.call(req) + } +} + /// Redirect layer that redirects browser requests to the console #[derive(Clone)] pub struct RedirectLayer; diff --git a/rustfs/src/storage/access.rs b/rustfs/src/storage/access.rs index a681dda4e..86937d650 100644 --- a/rustfs/src/storage/access.rs +++ b/rustfs/src/storage/access.rs @@ -17,6 +17,7 @@ use crate::auth::{check_key_valid, get_condition_values_with_query, get_session_ use crate::error::ApiError; use crate::license::license_check; use crate::server::RemoteAddr; +use crate::storage::request_context::RequestContext; use metrics::counter; use rustfs_ecstore::bucket::metadata_sys; use rustfs_ecstore::bucket::policy_sys::PolicySys; @@ -45,6 +46,7 @@ pub(crate) struct ReqInfo { pub version_id: Option, #[allow(dead_code)] pub region: Option, + pub request_context: Option, } #[derive(Clone, Debug)] @@ -67,6 +69,15 @@ fn ext_req_info_mut(ext: &mut http::Extensions) -> S3Result<&mut ReqInfo> { .ok_or_else(|| s3_error!(InternalError, "ReqInfo not found in request extensions")) } +/// Extract the canonical `RequestContext` from a request, checking both +/// the request extensions directly and the `ReqInfo.request_context` field. +pub(crate) fn request_context_from_req(req: &S3Request) -> Option { + req.extensions + .get::() + .cloned() + .or_else(|| req.extensions.get::().and_then(|ri| ri.request_context.clone())) +} + #[derive(Clone, Debug)] pub(crate) struct ObjectTagConditions { bucket: String, @@ -731,10 +742,13 @@ impl S3Access for FS { (None, false) }; + let request_context = cx.extensions_mut().get::().cloned(); + let req_info = ReqInfo { cred, is_owner, region: rustfs_ecstore::global::get_global_region(), + request_context, ..Default::default() }; diff --git a/rustfs/src/storage/helper.rs b/rustfs/src/storage/helper.rs index 4a2bd7352..968dfb363 100644 --- a/rustfs/src/storage/helper.rs +++ b/rustfs/src/storage/helper.rs @@ -12,7 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::storage::access::ReqInfo; +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 rustfs_audit::{ entity::{ApiDetails, ApiDetailsBuilder, AuditEntryBuilder}, @@ -24,10 +26,13 @@ use rustfs_s3_common::record_s3_op; use rustfs_s3_common::{EventName, S3Operation}; use rustfs_utils::{ extract_params_header, extract_req_params, extract_resp_elements, get_request_host, get_request_port, get_request_user_agent, + http::headers::AMZ_REQUEST_ID, }; use s3s::{S3Request, S3Response, S3Result}; +use serde_json::Value; use std::future::Future; use tokio::runtime::{Builder, Handle}; +use tracing::{Instrument, info_span}; /// Schedules an asynchronous task on the current runtime; /// if there is no runtime, creates a minimal runtime execution on a new thread. @@ -46,12 +51,30 @@ where } } +/// Spawn a background task with request context correlation. +/// Creates a child span with the request_id for tracing continuity, +/// ensuring audit/notify tasks can be traced back to the original request. +pub(crate) fn spawn_background_with_context(request_context: Option, fut: F) +where + F: Future + Send + 'static, +{ + match request_context { + Some(ctx) => { + let request_id = ctx.request_id.clone(); + let span = info_span!("background-task", request_id = %request_id); + spawn_background(Instrument::instrument(fut, span)); + } + None => spawn_background(fut), + } +} + /// A unified helper structure for building and distributing audit logs and event notifications via RAII mode at the end of an S3 operation scope. pub struct OperationHelper { audit_builder: Option, api_builder: ApiDetailsBuilder, event_builder: Option, start_time: std::time::Instant, + request_context: Option, } impl OperationHelper { @@ -95,18 +118,21 @@ impl OperationHelper { api_builder = api_builder.object(&object_key); } // Audit builder - let mut audit_builder = AuditEntryBuilder::new("1.0", event, trigger, ApiDetails::default()) + // Resolve canonical request context and request_id in a single pass: + // RequestContext.request_id > extract_request_id_from_headers() > "unknown" + let request_context = request_context_from_req(req); + let request_id = request_context + .as_ref() + .map(|ctx| ctx.request_id.clone()) + .unwrap_or_else(|| extract_request_id_from_headers(&req.headers)); + + let audit_builder = AuditEntryBuilder::new("1.0", event, trigger, ApiDetails::default()) .remote_host(remote_host) .user_agent(get_request_user_agent(&req.headers)) .req_host(get_request_host(&req.headers)) .req_path(req.uri.path().to_string()) - .req_query(extract_req_params(req)); - - if let Some(req_id) = req.headers.get("x-amz-request-id") - && let Ok(id_str) = req_id.to_str() - { - audit_builder = audit_builder.request_id(id_str); - } + .req_query(extract_req_params(req)) + .request_id(&request_id); let event_object = ObjectInfo { bucket: bucket.clone(), @@ -115,6 +141,12 @@ impl OperationHelper { }; let mut req_params = extract_params_header(&req.headers); + // Inject x-amz-request-id from RequestContext into req_params for event correlation + if let Some(ref ctx) = request_context { + req_params + .entry(AMZ_REQUEST_ID.to_string()) + .or_insert_with(|| ctx.x_amz_request_id.clone()); + } if let Some(principal_id) = req_info .and_then(|info| info.cred.as_ref()) .map(|cred| cred.access_key.clone()) @@ -141,7 +173,11 @@ impl OperationHelper { audit_builder: Some(audit_builder), api_builder, event_builder: Some(event_builder), - start_time: std::time::Instant::now(), + start_time: request_context + .as_ref() + .map(|ctx| ctx.start_time) + .unwrap_or_else(std::time::Instant::now), + request_context, } } @@ -211,6 +247,20 @@ impl OperationHelper { final_builder = final_builder.access_key(&sk); } + // Inject OpenTelemetry trace context into audit tags for distributed tracing correlation + if let Some(ref ctx) = self.request_context + && (ctx.trace_id.is_some() || ctx.span_id.is_some()) + { + let mut tags = HashMap::new(); + if let Some(ref tid) = ctx.trace_id { + tags.insert("traceId".to_string(), Value::String(tid.clone())); + } + if let Some(ref sid) = ctx.span_id { + tags.insert("spanId".to_string(), Value::String(sid.clone())); + } + final_builder = final_builder.tags(tags); + } + self.audit_builder = Some(final_builder); self.api_builder = ApiDetailsBuilder(api_details); // Store final details for Drop use } @@ -234,7 +284,8 @@ impl Drop for OperationHelper { fn drop(&mut self) { // Distribute audit logs if let Some(builder) = self.audit_builder.take() { - spawn_background(async move { + let ctx = self.request_context.clone(); + spawn_background_with_context(ctx, async move { AuditLogger::log(builder.build()).await; }); } @@ -246,7 +297,8 @@ impl Drop for OperationHelper { let event_args = builder.build(); // Avoid generating notifications for copy requests if !event_args.is_replication_request() { - spawn_background(async move { + let ctx = self.request_context.clone(); + spawn_background_with_context(ctx, async move { notifier_global::notify(event_args).await; }); } @@ -305,4 +357,63 @@ mod tests { assert_eq!(event_args.version_id, "version-123"); assert_eq!(event_args.req_params.get("principalId").map(String::as_str), Some("notifyTag")); } + + #[test] + fn operation_helper_prioritizes_request_context_for_request_id() { + let input = DeleteObjectTaggingInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + let mut req = build_request(input, Method::DELETE, Uri::from_static("/test-bucket/test-key")); + req.headers.insert("host", HeaderValue::from_static("example.com")); + req.headers.insert("user-agent", HeaderValue::from_static("rustfs-test")); + + // Insert RequestContext (set by ingress layer) with a specific request_id + req.extensions.insert(RequestContext { + request_id: "ingress-canonical-uuid".to_string(), + x_amz_request_id: "ingress-canonical-uuid".to_string(), + trace_id: None, + span_id: None, + start_time: std::time::Instant::now(), + }); + + req.extensions.insert(ReqInfo { + bucket: Some("test-bucket".to_string()), + object: Some("test-key".to_string()), + ..Default::default() + }); + + let helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject); + + // Verify the helper stored the RequestContext + assert!(helper.request_context.is_some()); + assert_eq!(helper.request_context.as_ref().unwrap().request_id, "ingress-canonical-uuid"); + } + + #[test] + fn operation_helper_no_request_context_when_absent() { + let input = DeleteObjectTaggingInput::builder() + .bucket("test-bucket".to_string()) + .key("test-key".to_string()) + .build() + .unwrap(); + let mut req = build_request(input, Method::DELETE, Uri::from_static("/test-bucket/test-key")); + req.headers.insert("host", HeaderValue::from_static("example.com")); + req.headers.insert("user-agent", HeaderValue::from_static("rustfs-test")); + req.headers + .insert("x-amz-request-id", HeaderValue::from_static("amz-header-uuid")); + + // No RequestContext inserted + req.extensions.insert(ReqInfo { + bucket: Some("test-bucket".to_string()), + object: Some("test-key".to_string()), + ..Default::default() + }); + + let helper = OperationHelper::new(&req, EventName::ObjectAccessedGet, S3Operation::GetObject); + + // Verify the helper has no RequestContext + assert!(helper.request_context.is_none()); + } } diff --git a/rustfs/src/storage/mod.rs b/rustfs/src/storage/mod.rs index 52de62dae..8b514ce33 100644 --- a/rustfs/src/storage/mod.rs +++ b/rustfs/src/storage/mod.rs @@ -21,6 +21,7 @@ pub(crate) mod entity; pub(crate) mod helper; pub mod lock_optimizer; pub mod options; +pub mod request_context; pub mod rpc; pub(crate) mod s3_api; mod sse; diff --git a/rustfs/src/storage/request_context.rs b/rustfs/src/storage/request_context.rs new file mode 100644 index 000000000..9e4514282 --- /dev/null +++ b/rustfs/src/storage/request_context.rs @@ -0,0 +1,176 @@ +// 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. + +//! Canonical request context carried through the entire request lifecycle. +//! +//! # Architecture +//! +//! ```text +//! HTTP Ingress (SetRequestIdLayer) +//! → generates x-request-id UUID +//! → RequestContextLayer creates RequestContext +//! → stores in request.extensions() +//! → sets x-amz-request-id header +//! Auth (FS::check) +//! → copies RequestContext into ReqInfo.request_context +//! Storage (FS methods) +//! → reads ReqInfo for bucket/object/version +//! → reads RequestContext for request_id/trace_id/span_id +//! Timeout Wrapper +//! → receives canonical request_id from caller +//! → passes to deadlock_detector.register_request() +//! OperationHelper +//! → reads RequestContext.request_id for audit log +//! → spawn_background_with_context() for audit/notify +//! tokio::spawn (request-internal) +//! → spawn_traced() = tokio::spawn + .instrument(Span::current()) +//! ``` +//! +//! # Frozen Rules (T00 Guardrails) +//! +//! ## request-id +//! - 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` +//! - Internal identifiers for sub-operations should use `operation_id` or `subtask_id` +//! +//! ## tokio::spawn usage +//! - **Request-internal tasks** (cache invalidation, metrics, read/write subtasks): +//! Use `spawn_traced()` which wraps `tokio::spawn` with `.instrument(Span::current())` +//! - **Post-request side effects** (audit flush, notify, replication enqueue): +//! Use `spawn_background_with_context()` which creates a correlated child span +//! with explicit `request_id` +//! - **Infrastructure tasks** (server loop, TLS reload, deadlock detection): +//! Plain `tokio::spawn` is acceptable; these are not request-scoped +//! - NEVER use bare `tokio::spawn` in request-handling code paths + +use http::HeaderMap; +use rustfs_utils::http::headers::AMZ_REQUEST_ID; +use std::time::Instant; + +/// Canonical request context carried through the entire request lifecycle. +/// +/// Created exactly once at HTTP ingress. Cloned by value; never mutated after creation. +#[derive(Clone, Debug)] +pub struct RequestContext { + /// Canonical request ID (from `x-request-id` header, set by `SetRequestIdLayer`). + pub request_id: String, + /// S3-compatible request ID alias (preserves upstream `x-amz-request-id` if present, + /// otherwise equals `request_id`). + pub x_amz_request_id: String, + /// OpenTelemetry trace ID (if present from upstream propagation). + pub trace_id: Option, + /// OpenTelemetry span ID (if present from upstream propagation). + pub span_id: Option, + /// Request ingress timestamp. + pub start_time: Instant, +} + +impl RequestContext { + /// Create a fallback `RequestContext` for paths that bypass HTTP ingress. + /// Generates a `req-{uuid}` format request-id. + pub fn fallback() -> Self { + let id = format!("req-{}", &uuid::Uuid::new_v4().to_string()[..8]); + Self { + request_id: id.clone(), + x_amz_request_id: id, + trace_id: None, + span_id: None, + start_time: Instant::now(), + } + } +} + +/// 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) +pub fn extract_request_id_from_headers(headers: &HeaderMap) -> String { + headers + .get("x-request-id") + .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()) +} + +/// Spawn a request-internal task that inherits the current tracing span. +/// +/// Use this for tasks that are part of the request processing pipeline +/// (e.g., cache invalidation, metrics recording, read/write subtasks). +/// +/// # Rules +/// - Do NOT use this for post-request side effects (audit, notify). +/// Use `crate::storage::helper::spawn_background_with_context` instead. +/// - Do NOT use bare `tokio::spawn` in request-handling code paths. +pub fn spawn_traced(fut: F) +where + F: std::future::Future + Send + 'static, +{ + tokio::spawn(tracing::Instrument::instrument(fut, tracing::Span::current())); +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_request_context_clone_send_sync() { + fn assert_clone_send_sync() {} + assert_clone_send_sync::(); + } + + #[test] + fn test_request_context_fallback_generates_id() { + let ctx = RequestContext::fallback(); + assert!(ctx.request_id.starts_with("req-")); + assert_eq!(ctx.request_id, ctx.x_amz_request_id); + assert!(ctx.trace_id.is_none()); + assert!(ctx.span_id.is_none()); + } + + #[test] + fn test_extract_request_id_from_x_request_id() { + let mut headers = HeaderMap::new(); + headers.insert("x-request-id", "test-uuid-123".parse().unwrap()); + let id = extract_request_id_from_headers(&headers); + assert_eq!(id, "test-uuid-123"); + } + + #[test] + fn test_extract_request_id_fallback_to_amz() { + let mut headers = HeaderMap::new(); + headers.insert("x-amz-request-id", "amz-uuid-456".parse().unwrap()); + let id = extract_request_id_from_headers(&headers); + assert_eq!(id, "amz-uuid-456"); + } + + #[test] + fn test_extract_request_id_priority() { + let mut headers = HeaderMap::new(); + headers.insert("x-request-id", "x-req-789".parse().unwrap()); + headers.insert("x-amz-request-id", "amz-req-000".parse().unwrap()); + let id = extract_request_id_from_headers(&headers); + assert_eq!(id, "x-req-789"); + } + + #[test] + fn test_extract_request_id_no_headers() { + let headers = HeaderMap::new(); + let id = extract_request_id_from_headers(&headers); + assert_eq!(id, "unknown"); + } +} diff --git a/rustfs/src/storage/timeout_wrapper.rs b/rustfs/src/storage/timeout_wrapper.rs index 47d711854..017239096 100644 --- a/rustfs/src/storage/timeout_wrapper.rs +++ b/rustfs/src/storage/timeout_wrapper.rs @@ -234,12 +234,15 @@ pub struct RequestTimeoutWrapper { impl RequestTimeoutWrapper { /// Create a new timeout wrapper with the given configuration. + /// + /// Note: This uses a sentinel request_id. Prefer `with_request_id()` to pass + /// the canonical request-id from `RequestContext`. pub fn new(config: TimeoutConfig) -> Self { Self { config, start_time: Instant::now(), cancel_token: CancellationToken::new(), - request_id: format!("req-{}", &uuid::Uuid::new_v4().to_string()[..8]), + request_id: "no-request-id".to_string(), } } @@ -253,17 +256,17 @@ impl RequestTimeoutWrapper { } } - /// Create a new timeout wrapper with operation size for dynamic timeout calculation + /// Create a new timeout wrapper with operation size for dynamic timeout calculation. + /// + /// Note: This uses a sentinel request_id. Prefer `with_request_id()` to pass + /// the canonical request-id from `RequestContext`. pub fn with_operation_size(config: TimeoutConfig, operation_size: Option) -> Self { - // Store operation size in config for later use - // Note: Currently we don't store the size in the wrapper itself, - // but the config can be used to calculate appropriate timeout - let _ = operation_size; // Suppress unused warning for now + let _ = operation_size; Self { config, start_time: Instant::now(), cancel_token: CancellationToken::new(), - request_id: format!("req-{}", &uuid::Uuid::new_v4().to_string()[..8]), + request_id: "no-request-id".to_string(), } }