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