diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index 8ae7479b1..e9ec8d41f 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -137,7 +137,7 @@ metrics = { workspace = true } tokio = { workspace = true, features = ["rt-multi-thread", "macros", "test-util"] } criterion = { workspace = true, features = ["html_reports"] } temp-env = { workspace = true, features = ["async_closure"] } -tracing-subscriber = { workspace = true } +tracing-subscriber = { workspace = true, features = ["json"] } serial_test = { workspace = true } opentelemetry_sdk = { workspace = true } diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 50c8dcdec..2eb558c83 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -64,9 +64,14 @@ const EVENT_REPLICATION_WORKER_RESIZE_SKIPPED: &str = "replication_worker_resize const EVENT_REPLICATION_WORKER_RESIZED: &str = "replication_worker_resized"; const EVENT_REPLICATION_BACKPRESSURE: &str = "replication_backpressure"; const EVENT_REPLICATION_RESYNC_LOAD_SKIPPED: &str = "replication_resync_load_skipped"; +const EVENT_REPLICATION_RESYNC_RECOVERED: &str = "replication_resync_recovered"; const EVENT_REPLICATION_CONFIG_LOOKUP_SKIPPED: &str = "replication_config_lookup_skipped"; const EVENT_REPLICATION_MRF_QUEUE_OVERFLOW: &str = "replication_mrf_queue_overflow"; +fn should_auto_resume_resync(status: ResyncStatusType) -> bool { + matches!(status, ResyncStatusType::ResyncPending | ResyncStatusType::ResyncStarted) +} + // Worker limits pub const WORKER_MAX_LIMIT: usize = 500; pub const WORKER_MIN_LIMIT: usize = 50; @@ -1017,6 +1022,11 @@ impl ReplicationPool { // Note: Leader lock implementation would be needed here // let _lock_guard = global_leader_lock.get_lock().await?; + let mut recovered_statuses = Vec::new(); + let mut restart_opts = Vec::new(); + let mut recovered_bucket_count = 0usize; + let mut skipped_failed_target_count = 0usize; + for bucket in buckets { let meta = match load_bucket_resync_metadata(bucket, self.storage.clone()).await { Ok(meta) => meta, @@ -1036,39 +1046,53 @@ impl ReplicationPool { } }; - // Store metadata in resyncer - { - let mut status_map = self.resyncer.status_map.write().await; - status_map.insert(bucket.clone(), meta.clone()); + if meta.targets_map.is_empty() { + continue; } - // Process target statistics - let target_stats = meta.clone_tgt_stats(); - for (arn, stats) in target_stats { - match stats.resync_status { - ResyncStatusType::ResyncFailed | ResyncStatusType::ResyncStarted | ResyncStatusType::ResyncPending => { - // Note: This would spawn a resync task in a real implementation - // For now, we just log the resync request - - let ctx = CancellationToken::new(); - let bucket_clone = bucket.clone(); - let resync = self.resyncer.clone(); - let storage = self.storage.clone(); - let opts = ResyncOpts { - bucket: bucket_clone, - arn, - resync_id: stats.resync_id, - resync_before: stats.resync_before_date, - }; - tokio::spawn(async move { - resync.register_cancel_token(&opts, ctx.clone()).await; - Box::pin(resync.clone().resync_bucket(ctx, storage, true, opts.clone())).await; - resync.clear_cancel_token(&opts).await; - }); - } - _ => {} + recovered_bucket_count += 1; + for (arn, stats) in &meta.targets_map { + if should_auto_resume_resync(stats.resync_status) { + restart_opts.push(ResyncOpts { + bucket: bucket.clone(), + arn: arn.clone(), + resync_id: stats.resync_id.clone(), + resync_before: stats.resync_before_date, + }); + } else if stats.resync_status == ResyncStatusType::ResyncFailed { + skipped_failed_target_count += 1; } } + + recovered_statuses.push((bucket.clone(), meta)); + } + + if !recovered_statuses.is_empty() { + let mut status_map = self.resyncer.status_map.write().await; + status_map.extend(recovered_statuses); + } + + if !restart_opts.is_empty() || skipped_failed_target_count > 0 { + info!( + event = EVENT_REPLICATION_RESYNC_RECOVERED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + recovered_buckets = recovered_bucket_count, + resumed_targets = restart_opts.len(), + skipped_failed_targets = skipped_failed_target_count, + "Recovered replication resync state from persisted metadata; failed targets require manual resync restart" + ); + } + + for opts in restart_opts { + let ctx = CancellationToken::new(); + let resync = self.resyncer.clone(); + let storage = self.storage.clone(); + tokio::spawn(async move { + resync.register_cancel_token(&opts, ctx.clone()).await; + Box::pin(resync.clone().resync_bucket(ctx, storage, true, opts.clone())).await; + resync.clear_cancel_token(&opts).await; + }); } Ok(()) @@ -1498,4 +1522,14 @@ mod tests { admission.merge(ReplicationQueueAdmission::Missed); assert_eq!(admission, ReplicationQueueAdmission::Missed); } + + #[test] + fn auto_resume_resync_only_for_inflight_states() { + assert!(should_auto_resume_resync(ResyncStatusType::ResyncPending)); + assert!(should_auto_resume_resync(ResyncStatusType::ResyncStarted)); + assert!(!should_auto_resume_resync(ResyncStatusType::NoResync)); + assert!(!should_auto_resume_resync(ResyncStatusType::ResyncCanceled)); + assert!(!should_auto_resume_resync(ResyncStatusType::ResyncCompleted)); + assert!(!should_auto_resume_resync(ResyncStatusType::ResyncFailed)); + } } diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index bb6cc99ff..610b51f28 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -87,7 +87,7 @@ use tokio::task::JoinSet; use tokio::time::Duration as TokioDuration; use tokio_util::io::ReaderStream; use tokio_util::sync::CancellationToken; -use tracing::{debug, error, instrument, warn}; +use tracing::{debug, error, instrument, trace, warn}; use uuid::Uuid; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; @@ -920,18 +920,32 @@ impl ReplicationResyncer { (roi.size, None) }; - debug!( - event = EVENT_RESYNC_OBJECT_PROCESSED, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, - reset_id = %reset_id, - bucket = %bucket_name, - object = %roi.name, - version_id = %roi.version_id.unwrap_or_default(), - size, - error = ?err, - "Processed resync object" - ); + if err.is_some() { + debug!( + event = EVENT_RESYNC_OBJECT_PROCESSED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + reset_id = %reset_id, + bucket = %bucket_name, + object = %roi.name, + version_id = %roi.version_id.unwrap_or_default(), + size, + error = ?err, + "Processed resync object with verification error" + ); + } else { + trace!( + event = EVENT_RESYNC_OBJECT_PROCESSED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + reset_id = %reset_id, + bucket = %bucket_name, + object = %roi.name, + version_id = %roi.version_id.unwrap_or_default(), + size, + "Processed resync object" + ); + } if cancel_token.is_cancelled() { return; diff --git a/crates/ecstore/src/client/api_remove.rs b/crates/ecstore/src/client/api_remove.rs index cd2390a49..27e453c9b 100644 --- a/crates/ecstore/src/client/api_remove.rs +++ b/crates/ecstore/src/client/api_remove.rs @@ -34,6 +34,7 @@ use std::{ }; use time::OffsetDateTime; use tokio::sync::mpsc::{self, Receiver, Sender}; +use tracing::Instrument; use crate::client::utils::base64_encode; use crate::client::{ @@ -221,11 +222,14 @@ impl TransitionClient { let self_clone = Arc::clone(&self); let bucket_name_owned = bucket_name.to_string(); - tokio::spawn(async move { - self_clone - .remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts) - .await; - }); + tokio::spawn( + async move { + self_clone + .remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts) + .await; + } + .instrument(tracing::Span::current()), + ); result_rx } @@ -241,26 +245,32 @@ impl TransitionClient { let bucket_name_owned = bucket_name.to_string(); let (result_tx, mut result_rx) = mpsc::channel(1); - tokio::spawn(async move { - self_clone - .remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts) - .await; - }); - tokio::spawn(async move { - while let Some(res) = result_rx.recv().await { - if res.err.is_none() { - continue; - } - error_tx - .send(RemoveObjectError { - object_name: res.object_name, - version_id: res.object_version_id, - err: res.err, - ..Default::default() - }) + tokio::spawn( + async move { + self_clone + .remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts) .await; } - }); + .instrument(tracing::Span::current()), + ); + tokio::spawn( + async move { + while let Some(res) = result_rx.recv().await { + if res.err.is_none() { + continue; + } + error_tx + .send(RemoveObjectError { + object_name: res.object_name, + version_id: res.object_version_id, + err: res.err, + ..Default::default() + }) + .await; + } + } + .instrument(tracing::Span::current()), + ); error_rx } diff --git a/crates/ecstore/src/rpc/peer_rest_client.rs b/crates/ecstore/src/rpc/peer_rest_client.rs index 3dcd4afac..78b13eece 100644 --- a/crates/ecstore/src/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/rpc/peer_rest_client.rs @@ -52,7 +52,7 @@ use tokio::{net::TcpStream, time::Duration}; use tonic::Request; use tonic::service::interceptor::InterceptedService; use tonic::transport::Channel; -use tracing::warn; +use tracing::{Instrument, warn}; pub const PEER_RESTSIGNAL: &str = "signal"; pub const PEER_RESTSUB_SYS: &str = "sub-sys"; @@ -78,6 +78,16 @@ pub struct PeerRestClient { } impl PeerRestClient { + fn recovery_monitor_span(grid_host: &str) -> tracing::Span { + tracing::info_span!( + "recovery-monitor", + component = "ecstore", + subsystem = "peer_rest_client", + kind = "peer_rest", + grid_host = %grid_host + ) + } + pub fn new(host: XHost, grid_host: String) -> Self { Self { host, @@ -172,28 +182,47 @@ impl PeerRestClient { let grid_host = self.grid_host.clone(); let offline = Arc::clone(&self.offline); let recovery_running = Arc::clone(&self.recovery_running); - tokio::spawn(async move { - let mut delay = get_drive_active_check_interval(); - let connect_timeout = get_drive_active_check_timeout(); + let span = Self::recovery_monitor_span(&grid_host); + tokio::spawn( + async move { + let mut delay = get_drive_active_check_interval(); + let connect_timeout = get_drive_active_check_timeout(); - for _ in 0..PEER_REST_RECOVERY_MAX_ATTEMPTS { - tokio::time::sleep(delay).await; - if Self::perform_connectivity_check(&grid_host, connect_timeout).await.is_ok() { - offline.store(false, Ordering::Release); - recovery_running.store(false, Ordering::Release); - return; + for _ in 0..PEER_REST_RECOVERY_MAX_ATTEMPTS { + tokio::time::sleep(delay).await; + if Self::perform_connectivity_check(&grid_host, connect_timeout).await.is_ok() { + offline.store(false, Ordering::Release); + recovery_running.store(false, Ordering::Release); + return; + } + + delay = std::cmp::min(delay.saturating_mul(2), PEER_REST_RECOVERY_MAX_BACKOFF); } - delay = std::cmp::min(delay.saturating_mul(2), PEER_REST_RECOVERY_MAX_BACKOFF); + warn!( + grid_host = %grid_host, + attempts = PEER_REST_RECOVERY_MAX_ATTEMPTS, + "peer recovery monitor reached max attempts; will retry on next request" + ); + recovery_running.store(false, Ordering::Release); } + .instrument(span), + ); + } - warn!( - grid_host = %grid_host, - attempts = PEER_REST_RECOVERY_MAX_ATTEMPTS, - "peer recovery monitor reached max attempts; will retry on next request" - ); - recovery_running.store(false, Ordering::Release); - }); + #[cfg(test)] + fn spawn_recovery_monitor_log_probe_for_test(&self) -> tokio::sync::oneshot::Receiver<()> { + let (tx, rx) = tokio::sync::oneshot::channel(); + let grid_host = self.grid_host.clone(); + let span = Self::recovery_monitor_span(&grid_host); + tokio::spawn( + async move { + warn!(grid_host = %grid_host, "peer recovery monitor log probe"); + let _ = tx.send(()); + } + .instrument(span), + ); + rx } async fn perform_connectivity_check(addr: &str, timeout_duration: Duration) -> Result<()> { @@ -956,6 +985,58 @@ impl PeerRestClient { #[cfg(test)] mod tests { use super::*; + use serde_json::Value; + use std::io::{self, Write}; + use std::sync::{Arc, Mutex}; + use tracing_subscriber::{Registry, fmt::MakeWriter, layer::SubscriberExt}; + + #[derive(Clone, Default)] + struct CapturedLogs { + buffer: Arc>>, + } + + struct CapturedLogWriter { + buffer: Arc>>, + } + + impl CapturedLogs { + fn lines(&self) -> Vec { + let buffer = self + .buffer + .lock() + .expect("captured logs mutex should not be poisoned") + .clone(); + String::from_utf8(buffer) + .expect("captured logs should be valid UTF-8") + .lines() + .map(|line| serde_json::from_str::(line).expect("captured log line should be valid JSON")) + .collect() + } + } + + impl Write for CapturedLogWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.buffer + .lock() + .expect("captured logs mutex should not be poisoned") + .extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + + impl<'a> MakeWriter<'a> for CapturedLogs { + type Writer = CapturedLogWriter; + + fn make_writer(&'a self) -> Self::Writer { + CapturedLogWriter { + buffer: Arc::clone(&self.buffer), + } + } + } fn test_peer_client() -> PeerRestClient { PeerRestClient::new( @@ -1011,4 +1092,40 @@ mod tests { assert!(matches!(err, Error::VolumeNotFound)); assert!(!client.offline.load(Ordering::Acquire)); } + + #[tokio::test(flavor = "current_thread")] + async fn peer_rest_recovery_probe_logs_keep_request_id_span_context() { + let logs = CapturedLogs::default(); + let subscriber = Registry::default().with( + tracing_subscriber::fmt::layer() + .with_writer(logs.clone()) + .with_ansi(false) + .without_time() + .json() + .flatten_event(true) + .with_current_span(true) + .with_span_list(true), + ); + let _guard = tracing::subscriber::set_default(subscriber); + + let client = test_peer_client(); + let span = tracing::info_span!("request-span", request_id = "req-peer-rest"); + let _entered = span.enter(); + let done = client.spawn_recovery_monitor_log_probe_for_test(); + done.await.expect("recovery monitor probe should signal completion"); + + let log = logs + .lines() + .into_iter() + .find(|value| value.get("message").and_then(Value::as_str) == Some("peer recovery monitor log probe")) + .expect("expected peer recovery monitor probe log"); + + assert_eq!(log["span"]["name"], Value::String("recovery-monitor".to_string())); + assert_eq!(log["span"]["kind"], Value::String("peer_rest".to_string())); + let spans = log["spans"].as_array().expect("spans should be present"); + assert!(spans.iter().any(|span| { + span.get("name").and_then(Value::as_str) == Some("request-span") + && span.get("request_id").and_then(Value::as_str) == Some("req-peer-rest") + })); + } } diff --git a/crates/ecstore/src/rpc/peer_s3_client.rs b/crates/ecstore/src/rpc/peer_s3_client.rs index 2dbddc578..83f3dd690 100644 --- a/crates/ecstore/src/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/rpc/peer_s3_client.rs @@ -45,7 +45,7 @@ use tokio_util::sync::CancellationToken; use tonic::Request; use tonic::service::interceptor::InterceptedService; use tonic::transport::Channel; -use tracing::{debug, info, warn}; +use tracing::{Instrument, debug, info, warn}; type Client = Arc>; @@ -596,6 +596,16 @@ pub struct RemotePeerS3Client { } impl RemotePeerS3Client { + fn recovery_monitor_span(addr: &str) -> tracing::Span { + tracing::info_span!( + "recovery-monitor", + component = "ecstore", + subsystem = "peer_s3_client", + kind = "peer_s3", + addr = %addr + ) + } + pub fn new(node: Option, pools: Option>) -> Self { let addr = node.as_ref().map(|v| v.url.to_string()).unwrap_or_default(); let client = Self { @@ -661,10 +671,14 @@ impl RemotePeerS3Client { let health_clone = Arc::clone(&health); let addr_clone = addr.clone(); let cancel_clone = cancel_token.clone(); + let span = Self::recovery_monitor_span(&addr_clone); - tokio::spawn(async move { - Self::monitor_remote_peer_recovery(addr_clone, health_clone, cancel_clone).await; - }); + tokio::spawn( + async move { + Self::monitor_remote_peer_recovery(addr_clone, health_clone, cancel_clone).await; + } + .instrument(span), + ); } } } @@ -767,9 +781,13 @@ impl RemotePeerS3Client { let health = Arc::clone(&self.health); let cancel_token = self.cancel_token.clone(); let addr = self.addr.clone(); - tokio::spawn(async move { - Self::monitor_remote_peer_recovery(addr, health, cancel_token).await; - }); + let span = Self::recovery_monitor_span(&addr); + tokio::spawn( + async move { + Self::monitor_remote_peer_recovery(addr, health, cancel_token).await; + } + .instrument(span), + ); } } } diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index 45f8ecb99..c786f8462 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -65,7 +65,7 @@ use tokio::{ }; use tokio_util::sync::CancellationToken; use tonic::{Request, service::interceptor::InterceptedService, transport::Channel}; -use tracing::{debug, warn}; +use tracing::{Instrument, debug, warn}; use uuid::Uuid; #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -117,6 +117,17 @@ pub struct RemoteDisk { } impl RemoteDisk { + fn recovery_monitor_span(addr: &str, endpoint: &Endpoint) -> tracing::Span { + tracing::info_span!( + "recovery-monitor", + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REMOTE_DISK, + kind = "remote_disk", + endpoint = %endpoint, + addr = %addr + ) + } + fn is_retryable_walk_dir_error(err: &DiskError) -> bool { if is_network_like_disk_error(err) { return true; @@ -239,9 +250,37 @@ impl RemoteDisk { let endpoint = self.endpoint.clone(); let health = Arc::clone(&self.health); let cancel_token = self.cancel_token.clone(); - tokio::spawn(async move { - Self::monitor_remote_disk_recovery(addr, endpoint, health, cancel_token).await; - }); + let span = Self::recovery_monitor_span(&addr, &endpoint); + tokio::spawn( + async move { + Self::monitor_remote_disk_recovery(addr, endpoint, health, cancel_token).await; + } + .instrument(span), + ); + } + + #[cfg(test)] + fn spawn_recovery_monitor_log_probe_for_test(&self) -> tokio::sync::oneshot::Receiver<()> { + let (tx, rx) = tokio::sync::oneshot::channel(); + let endpoint = self.endpoint.clone(); + let addr = self.addr.clone(); + let span = Self::recovery_monitor_span(&addr, &endpoint); + tokio::spawn( + async move { + warn!( + event = EVENT_REMOTE_DISK_HEALTH, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REMOTE_DISK, + endpoint = %endpoint, + addr, + state = "probe", + "remote disk recovery monitor log probe" + ); + let _ = tx.send(()); + } + .instrument(span), + ); + rx } fn mark_suspect_or_offline(&self, reason: &'static str) -> bool { @@ -295,10 +334,14 @@ impl RemoteDisk { let addr_clone = addr.clone(); let endpoint_clone = endpoint.clone(); let cancel_clone = cancel_token.clone(); + let span = Self::recovery_monitor_span(&addr_clone, &endpoint_clone); - tokio::spawn(async move { - Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await; - }); + tokio::spawn( + async move { + Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await; + } + .instrument(span), + ); } loop { @@ -357,10 +400,14 @@ impl RemoteDisk { let addr_clone = addr.clone(); let endpoint_clone = endpoint.clone(); let cancel_clone = cancel_token.clone(); + let span = Self::recovery_monitor_span(&addr_clone, &endpoint_clone); - tokio::spawn(async move { - Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await; - }); + tokio::spawn( + async move { + Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await; + } + .instrument(span), + ); } } } @@ -375,10 +422,28 @@ impl RemoteDisk { cancel_token: CancellationToken, ) { let mut interval = time::interval(get_drive_returning_probe_interval()); + debug!( + event = EVENT_REMOTE_DISK_HEALTH, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REMOTE_DISK, + endpoint = %endpoint, + addr, + state = "recovery_monitor_started", + "Remote disk recovery monitor started" + ); loop { tokio::select! { _ = cancel_token.cancelled() => { + debug!( + event = EVENT_REMOTE_DISK_HEALTH, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REMOTE_DISK, + endpoint = %endpoint, + addr, + state = "recovery_monitor_cancelled", + "Remote disk recovery monitor cancelled" + ); return; } _ = interval.tick() => { @@ -2200,17 +2265,68 @@ mod tests { use super::*; use crate::rpc::{InternodeDataTransportCapabilities, TcpHttpInternodeDataTransport}; use rustfs_common::GLOBAL_CONN_MAP; + use serde_json::Value; + use std::io::{self as std_io, Write}; use std::pin::Pin; - use std::sync::{Mutex as StdMutex, Once}; + use std::sync::{Arc, Mutex, Mutex as StdMutex, Once}; use std::task::{Context, Poll}; use tokio::io::{ReadBuf, duplex}; use tokio::net::TcpListener; use tonic::transport::Endpoint as TonicEndpoint; use tracing::Level; + use tracing_subscriber::{Registry, fmt::MakeWriter, layer::SubscriberExt}; use uuid::Uuid; static INIT: Once = Once::new(); + #[derive(Clone, Default)] + struct CapturedLogs { + buffer: Arc>>, + } + + struct CapturedLogWriter { + buffer: Arc>>, + } + + impl CapturedLogs { + fn lines(&self) -> Vec { + let buffer = self + .buffer + .lock() + .expect("captured logs mutex should not be poisoned") + .clone(); + String::from_utf8(buffer) + .expect("captured logs should be valid UTF-8") + .lines() + .map(|line| serde_json::from_str::(line).expect("captured log line should be valid JSON")) + .collect() + } + } + + impl Write for CapturedLogWriter { + fn write(&mut self, buf: &[u8]) -> std_io::Result { + self.buffer + .lock() + .expect("captured logs mutex should not be poisoned") + .extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std_io::Result<()> { + Ok(()) + } + } + + impl<'a> MakeWriter<'a> for CapturedLogs { + type Writer = CapturedLogWriter; + + fn make_writer(&'a self) -> Self::Writer { + CapturedLogWriter { + buffer: Arc::clone(&self.buffer), + } + } + } + #[derive(Debug, Clone)] enum RecordedTransportCall { Read(ReadStreamRequest), @@ -3500,4 +3616,155 @@ mod tests { assert_eq!(endpoint.set_idx, 2); assert_eq!(endpoint.disk_idx, 3); } + + #[tokio::test(flavor = "current_thread")] + async fn remote_disk_recovery_probe_logs_keep_request_id_span_context() { + let logs = CapturedLogs::default(); + let subscriber = Registry::default().with( + tracing_subscriber::fmt::layer() + .with_writer(logs.clone()) + .with_ansi(false) + .without_time() + .json() + .flatten_event(true) + .with_current_span(true) + .with_span_list(true), + ); + let _guard = tracing::subscriber::set_default(subscriber); + + let endpoint = Endpoint { + url: url::Url::parse("http://127.0.0.1:59996/data").expect("endpoint URL should parse"), + is_local: false, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, + }; + let remote_disk = RemoteDisk::new( + &endpoint, + &DiskOption { + cleanup: false, + health_check: true, + }, + Arc::new(TcpHttpInternodeDataTransport), + ) + .await + .expect("remote disk should construct"); + + let span = tracing::info_span!("request-span", request_id = "req-remote-disk"); + let _entered = span.enter(); + let done = remote_disk.spawn_recovery_monitor_log_probe_for_test(); + done.await + .expect("remote disk recovery monitor probe should signal completion"); + + let log = logs + .lines() + .into_iter() + .find(|value| value.get("message").and_then(Value::as_str) == Some("remote disk recovery monitor log probe")) + .expect("expected remote disk recovery monitor probe log"); + + assert_eq!(log["span"]["name"], Value::String("recovery-monitor".to_string())); + assert_eq!(log["span"]["kind"], Value::String("remote_disk".to_string())); + let spans = log["spans"].as_array().expect("spans should be present"); + assert!(spans.iter().any(|span| { + span.get("name").and_then(Value::as_str) == Some("request-span") + && span.get("request_id").and_then(Value::as_str) == Some("req-remote-disk") + })); + } + + #[tokio::test(flavor = "current_thread")] + async fn remote_disk_network_error_starts_recovery_monitor_with_request_context() { + let logs = CapturedLogs::default(); + let subscriber = Registry::default().with( + tracing_subscriber::fmt::layer() + .with_writer(logs.clone()) + .with_ansi(false) + .without_time() + .json() + .flatten_event(true) + .with_current_span(true) + .with_span_list(true), + ); + let _guard = tracing::subscriber::set_default(subscriber); + + let addr = "http://127.0.0.1:59997".to_string(); + let endpoint = Endpoint { + url: url::Url::parse(&format!("{addr}/data")).expect("endpoint URL should parse"), + is_local: false, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, + }; + + let remote_disk = RemoteDisk::new( + &endpoint, + &DiskOption { + cleanup: false, + health_check: true, + }, + Arc::new(TcpHttpInternodeDataTransport), + ) + .await + .expect("remote disk should construct"); + + let span = tracing::info_span!("request-span", request_id = "req-remote-disk-e2e"); + let _entered = span.enter(); + + let err = remote_disk + .execute_with_timeout( + || async { + Err::<(), Error>(DiskError::Io(std::io::Error::new( + std::io::ErrorKind::ConnectionRefused, + "connection refused", + ))) + }, + Duration::from_secs(1), + ) + .await + .expect_err("network-like operation error should fail"); + assert_eq!( + match &err { + DiskError::Io(io_err) => io_err.kind(), + other => panic!("expected io network error, got {other:?}"), + }, + std::io::ErrorKind::ConnectionRefused + ); + + tokio::task::yield_now().await; + tokio::time::sleep(Duration::from_millis(20)).await; + remote_disk.cancel_token.cancel(); + tokio::task::yield_now().await; + + let lines = logs.lines(); + let marked_suspect = lines + .iter() + .find(|value| value.get("state").and_then(Value::as_str) == Some("marked_suspect")) + .expect("expected marked_suspect log"); + assert!( + marked_suspect["spans"] + .as_array() + .expect("spans should be present") + .iter() + .any(|span| { + span.get("name").and_then(Value::as_str) == Some("request-span") + && span.get("request_id").and_then(Value::as_str) == Some("req-remote-disk-e2e") + }) + ); + + let recovery_started = lines + .iter() + .find(|value| value.get("state").and_then(Value::as_str) == Some("recovery_monitor_started")) + .expect("expected recovery_monitor_started log"); + assert_eq!(recovery_started["span"]["name"], Value::String("recovery-monitor".to_string())); + assert_eq!(recovery_started["span"]["kind"], Value::String("remote_disk".to_string())); + assert!( + recovery_started["spans"] + .as_array() + .expect("spans should be present") + .iter() + .any(|span| { + span.get("name").and_then(Value::as_str) == Some("request-span") + && span.get("request_id").and_then(Value::as_str) == Some("req-remote-disk-e2e") + }) + ); + } } diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index ae07ca9d7..675efd91f 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -127,7 +127,7 @@ use tokio::{ }; use tokio_util::sync::CancellationToken; use tracing::error; -use tracing::{debug, info, warn}; +use tracing::{Instrument, debug, info, warn}; use uuid::Uuid; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; @@ -1491,55 +1491,58 @@ impl SetDisks { if !pending.is_empty() { let cleanup_requests = requests.clone(); let lockers = self.lockers.clone(); - let handle = tokio::spawn(async move { - let mut late_lock_ids_by_client = vec![Vec::new(); lockers.len()]; - let mut pending = pending; - while let Some(join_result) = pending.join_next().await { - match join_result { - Ok((client_idx, Ok(responses))) => { - for (req_idx, request) in cleanup_requests.iter().enumerate() { - if let Some(response) = responses.get(req_idx) - && response.success - { - let lock_id = response - .lock_info - .as_ref() - .map(|lock_info| lock_info.id.clone()) - .unwrap_or_else(|| request.lock_id.clone()); - if let Some(client_locks) = late_lock_ids_by_client.get_mut(client_idx) { - client_locks.push(lock_id); + let handle = tokio::spawn( + async move { + let mut late_lock_ids_by_client = vec![Vec::new(); lockers.len()]; + let mut pending = pending; + while let Some(join_result) = pending.join_next().await { + match join_result { + Ok((client_idx, Ok(responses))) => { + for (req_idx, request) in cleanup_requests.iter().enumerate() { + if let Some(response) = responses.get(req_idx) + && response.success + { + let lock_id = response + .lock_info + .as_ref() + .map(|lock_info| lock_info.id.clone()) + .unwrap_or_else(|| request.lock_id.clone()); + if let Some(client_locks) = late_lock_ids_by_client.get_mut(client_idx) { + client_locks.push(lock_id); + } } } } - } - Ok((_client_idx, Err(err))) => { - tracing::warn!("late distributed delete lock batch request failed: {}", err); - } - Err(err) => { - tracing::warn!("late distributed delete lock batch task join failed: {}", err); - } - } - } - - join_all(lockers.iter().cloned().enumerate().filter_map(|(client_idx, client)| { - let lock_ids = late_lock_ids_by_client.get(client_idx).cloned().unwrap_or_default(); - if lock_ids.is_empty() { - None - } else { - Some(async move { - if let Err(err) = client.release_locks_batch(&lock_ids).await { - tracing::warn!( - client_idx, - lock_count = lock_ids.len(), - "failed to cleanup late distributed delete locks in batch: {}", - err - ); + Ok((_client_idx, Err(err))) => { + tracing::warn!("late distributed delete lock batch request failed: {}", err); } - }) + Err(err) => { + tracing::warn!("late distributed delete lock batch task join failed: {}", err); + } + } } - })) - .await; - }); + + join_all(lockers.iter().cloned().enumerate().filter_map(|(client_idx, client)| { + let lock_ids = late_lock_ids_by_client.get(client_idx).cloned().unwrap_or_default(); + if lock_ids.is_empty() { + None + } else { + Some(async move { + if let Err(err) = client.release_locks_batch(&lock_ids).await { + tracing::warn!( + client_idx, + lock_count = lock_ids.len(), + "failed to cleanup late distributed delete locks in batch: {}", + err + ); + } + }) + } + })) + .await; + } + .instrument(tracing::Span::current()), + ); drop(handle); } diff --git a/crates/ecstore/src/store_list_objects.rs b/crates/ecstore/src/store_list_objects.rs index 034b3ed27..b95374dd1 100644 --- a/crates/ecstore/src/store_list_objects.rs +++ b/crates/ecstore/src/store_list_objects.rs @@ -40,7 +40,7 @@ use std::sync::Arc; use tokio::sync::broadcast::{self}; use tokio::sync::mpsc::{self, Receiver, Sender}; use tokio_util::sync::CancellationToken; -use tracing::{error, info, warn}; +use tracing::{Instrument, error, info, warn}; use uuid::Uuid; const MAX_OBJECT_LIST: i32 = 1000; @@ -710,31 +710,37 @@ impl ECStore { let cancel_rx1 = cancel.clone(); let cancel_rx1_for_err = cancel_rx1.clone(); let err_tx1 = err_tx.clone(); - let job1 = tokio::spawn(async move { - let mut opts = opts; - opts.stop_disk_at_limit = true; - if let Err(err) = store.list_merged(cancel_rx1, opts, sender).await - && !cancel_rx1_for_err.is_cancelled() - { - error!("list_merged err {:?}", err); - let _ = err_tx1.send(Arc::new(err)); + let job1 = tokio::spawn( + async move { + let mut opts = opts; + opts.stop_disk_at_limit = true; + if let Err(err) = store.list_merged(cancel_rx1, opts, sender).await + && !cancel_rx1_for_err.is_cancelled() + { + error!("list_merged err {:?}", err); + let _ = err_tx1.send(Arc::new(err)); + } } - }); + .instrument(tracing::Span::current()), + ); let cancel_rx2 = cancel.clone(); let (result_tx, mut result_rx) = mpsc::channel(1); let err_tx2 = err_tx.clone(); let opts = o.clone(); - let job2 = tokio::spawn(async move { - if let Err(err) = gather_results(cancel_rx2, opts, recv, result_tx).await { - error!("gather_results err {:?}", err); - let _ = err_tx2.send(Arc::new(err)); - } + let job2 = tokio::spawn( + async move { + if let Err(err) = gather_results(cancel_rx2, opts, recv, result_tx).await { + error!("gather_results err {:?}", err); + let _ = err_tx2.send(Arc::new(err)); + } - // cancel call exit spawns - cancel.cancel(); - }); + // cancel call exit spawns + cancel.cancel(); + } + .instrument(tracing::Span::current()), + ); let mut result = { // receiver result @@ -815,11 +821,14 @@ impl ECStore { } } - tokio::spawn(async move { - if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await { - error!("merge_entry_channels err {:?}", err) + tokio::spawn( + async move { + if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await { + error!("merge_entry_channels err {:?}", err) + } } - }); + .instrument(tracing::Span::current()), + ); // let merge_res = merge_entry_channels(rx, inputs, sender.clone(), 1).await; @@ -1007,34 +1016,47 @@ impl ECStore { Err(_) => None, }; - tokio::spawn(async move { - let mut sent_err = false; - while let Some(entry) = merge_rx.recv().await { - if opts.latest_only { - let fi = match entry.to_fileinfo(&bucket_clone) { - Ok(res) => res, - Err(err) => { - if !sent_err { + tokio::spawn( + async move { + let mut sent_err = false; + while let Some(entry) = merge_rx.recv().await { + if opts.latest_only { + let fi = match entry.to_fileinfo(&bucket_clone) { + Ok(res) => res, + Err(err) => { + if !sent_err { + let item = ObjectInfoOrErr { + item: None, + err: Some(err.into()), + }; + + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + + sent_err = true; + + return; + } + + continue; + } + }; + + if let Some(filter) = opts.filter { + if filter(&fi) { let item = ObjectInfoOrErr { - item: None, - err: Some(err.into()), + item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, { + if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } + })), + err: None, }; if let Err(err) = result.send(item).await { error!("walk result send err {:?}", err); } - - sent_err = true; - - return; } - - continue; - } - }; - - if let Some(filter) = opts.filter { - if filter(&fi) { + } else { let item = ObjectInfoOrErr { item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, { if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } @@ -1046,43 +1068,43 @@ impl ECStore { error!("walk result send err {:?}", err); } } - } else { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } + continue; } - continue; - } - let fvs = match entry.file_info_versions(&bucket_clone) { - Ok(res) => res, - Err(err) => { - let item = ObjectInfoOrErr { - item: None, - err: Some(err.into()), - }; + let fvs = match entry.file_info_versions(&bucket_clone) { + Ok(res) => res, + Err(err) => { + let item = ObjectInfoOrErr { + item: None, + err: Some(err.into()), + }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + return; } - return; + }; + + if opts.versions_sort == WalkVersionsSortOrder::Ascending { + //TODO: SORT } - }; - if opts.versions_sort == WalkVersionsSortOrder::Ascending { - //TODO: SORT - } + for fi in fvs.versions.iter() { + if let Some(filter) = opts.filter { + if filter(fi) { + let item = ObjectInfoOrErr { + item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, { + if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } + })), + err: None, + }; - for fi in fvs.versions.iter() { - if let Some(filter) = opts.filter { - if filter(fi) { + if let Err(err) = result.send(item).await { + error!("walk result send err {:?}", err); + } + } + } else { let item = ObjectInfoOrErr { item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, { if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } @@ -1094,23 +1116,13 @@ impl ECStore { error!("walk result send err {:?}", err); } } - } else { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } } } } - }); + .instrument(tracing::Span::current()), + ); - tokio::spawn(async move { merge_entry_channels(rx, inputs, merge_tx, 1).await }); + tokio::spawn(async move { merge_entry_channels(rx, inputs, merge_tx, 1).await }.instrument(tracing::Span::current())); let walk_results = join_all(futures).await; let mut errs = Vec::new(); diff --git a/crates/obs/Cargo.toml b/crates/obs/Cargo.toml index 5c78e3212..33d3eb8b6 100644 --- a/crates/obs/Cargo.toml +++ b/crates/obs/Cargo.toml @@ -61,6 +61,7 @@ opentelemetry-otlp = { workspace = true } opentelemetry-semantic-conventions = { workspace = true } percent-encoding = { workspace = true } serde = { workspace = true } +serde_json = { workspace = true } tracing = { workspace = true, features = ["std", "attributes"] } tracing-appender = { workspace = true } tracing-error = { workspace = true } @@ -83,4 +84,3 @@ jemalloc_pprof = { workspace = true, optional = true } tokio = { workspace = true, features = ["full"] } tempfile = { workspace = true } temp-env = { workspace = true } -serde_json = { workspace = true } diff --git a/crates/obs/src/telemetry/local.rs b/crates/obs/src/telemetry/local.rs index 276741e16..fe431008b 100644 --- a/crates/obs/src/telemetry/local.rs +++ b/crates/obs/src/telemetry/local.rs @@ -48,13 +48,19 @@ use rustfs_config::observability::{ DEFAULT_OBS_LOG_ZSTD_WORKERS, }; use rustfs_config::{APP_NAME, DEFAULT_LOG_KEEP_FILES, DEFAULT_LOG_ROTATION_TIME, DEFAULT_OBS_LOG_STDOUT_ENABLED}; +use serde_json::Value as JsonValue; use std::sync::Arc; +use std::{collections::BTreeMap, fmt}; use std::{fs, io::IsTerminal, time::Duration}; use tracing::Subscriber; use tracing::{info, warn}; use tracing_error::ErrorLayer; use tracing_subscriber::{ - fmt::{format::FmtSpan, time::LocalTime}, + fmt::{ + FormatEvent, + format::{FmtSpan, Format, Json, JsonFields, Writer}, + time::LocalTime, + }, layer::SubscriberExt, registry::LookupSpan, util::SubscriberInitExt, @@ -65,6 +71,73 @@ const LOG_SUBSYSTEM_LOCAL_LOGGING: &str = "local_logging"; const EVENT_LOCAL_LOGGING_STATE: &str = "local_logging_state"; const EVENT_LOG_CLEANER_STATE: &str = "log_cleaner_state"; const STDERR_WARNING_PREFIX: &str = "[WARN]"; +const REQUEST_ID: &str = "request-id"; + +#[derive(Clone, Debug)] +struct RequestIdJsonFormat { + inner: Format, +} + +impl RequestIdJsonFormat { + fn new(inner: Format) -> Self { + Self { inner } + } +} + +impl FormatEvent for RequestIdJsonFormat +where + S: Subscriber + for<'lookup> LookupSpan<'lookup>, + T: tracing_subscriber::fmt::time::FormatTime, +{ + fn format_event( + &self, + ctx: &tracing_subscriber::fmt::FmtContext<'_, S, JsonFields>, + mut writer: Writer<'_>, + event: &tracing::Event<'_>, + ) -> fmt::Result { + let mut buffer = String::new(); + self.inner.format_event(ctx, Writer::new(&mut buffer), event)?; + + let request_id = span_scope_request_id(ctx, event); + if request_id.is_none() { + return writer.write_str(&buffer); + } + + let trimmed = buffer.trim_end(); + let mut payload: JsonValue = serde_json::from_str(trimmed).map_err(|_| fmt::Error)?; + if let Some(object) = payload.as_object_mut() { + object + .entry(REQUEST_ID.to_string()) + .or_insert_with(|| JsonValue::String(request_id.expect("checked is_some"))); + } + + let serialized = serde_json::to_string(&payload).map_err(|_| fmt::Error)?; + writer.write_str(&serialized)?; + writer.write_char('\n') + } +} + +fn span_scope_request_id( + ctx: &tracing_subscriber::fmt::FmtContext<'_, S, JsonFields>, + event: &tracing::Event<'_>, +) -> Option +where + S: Subscriber + for<'lookup> LookupSpan<'lookup>, +{ + let current_span = event.parent().and_then(|id| ctx.span(id)).or_else(|| ctx.lookup_current())?; + let mut request_id = None; + + for span in current_span.scope().from_root() { + let extensions = span.extensions(); + let formatted_fields = extensions.get::>()?; + let fields: BTreeMap = serde_json::from_str(&formatted_fields.fields).ok()?; + if let Some(value) = fields.get(REQUEST_ID).and_then(JsonValue::as_str) { + request_id = Some(value.to_string()); + } + } + + request_id +} pub(super) fn build_json_log_layer(writer: W, enable_ansi: bool, span_events: FmtSpan) -> impl tracing_subscriber::Layer where @@ -81,9 +154,11 @@ where .with_line_number(true) .with_writer(writer) .json() + .flatten_event(true) .with_current_span(true) .with_span_list(true) .with_span_events(span_events) + .map_event_format(RequestIdJsonFormat::new) } /// Initialize local logging (stdout-only or file-rolling). @@ -520,7 +595,7 @@ pub fn spawn_cleanup_task( } Ok(Err(e)) => { counter!(METRIC_LOG_CLEANER_RUN_FAILURES_TOTAL).increment(1); - tracing::warn!( + warn!( event = EVENT_LOG_CLEANER_STATE, component = LOG_COMPONENT_OBS, subsystem = LOG_SUBSYSTEM_LOCAL_LOGGING, @@ -531,7 +606,7 @@ pub fn spawn_cleanup_task( } Err(e) => { counter!(METRIC_LOG_CLEANER_RUN_FAILURES_TOTAL).increment(1); - tracing::warn!( + warn!( event = EVENT_LOG_CLEANER_STATE, component = LOG_COMPONENT_OBS, subsystem = LOG_SUBSYSTEM_LOCAL_LOGGING, @@ -593,7 +668,7 @@ mod tests { let subscriber = Registry::default().with(layer); tracing::subscriber::with_default(subscriber, || { - tracing::info!(component = "obs_contract_test", sink = "stdout", "hello world"); + info!(component = "obs_contract_test", sink = "stdout", "hello world"); }); let raw = String::from_utf8( @@ -609,6 +684,55 @@ mod tests { (raw, parsed) } + fn render_json_log_with_request_span() -> Value { + let writer = SharedWriter::default(); + let layer = build_json_log_layer(writer.clone(), false, FmtSpan::NONE); + let subscriber = Registry::default().with(layer); + + tracing::subscriber::with_default(subscriber, || { + let span = tracing::info_span!("http-request", request_id = "req-123", method = "GET"); + let _guard = span.enter(); + info!(component = "obs_contract_test", message = "inside request span"); + }); + + let raw = String::from_utf8( + writer + .inner + .lock() + .expect("shared writer lock should not be poisoned") + .clone(), + ) + .expect("log output should be valid UTF-8"); + let first_line = raw.lines().next().expect("expected one JSON log line"); + serde_json::from_str(first_line).expect("log line should be valid JSON") + } + + fn render_json_log_with_recovery_monitor_child_span() -> Value { + let writer = SharedWriter::default(); + let layer = build_json_log_layer(writer.clone(), false, FmtSpan::NONE); + let subscriber = Registry::default().with(layer); + + tracing::subscriber::with_default(subscriber, || { + let request_span = tracing::info_span!("request-span", request_id = "req-parent", method = "GET"); + let _request_guard = request_span.enter(); + let recovery_span = + tracing::info_span!("recovery-monitor", kind = "remote_disk", endpoint = "http://127.0.0.1:9000/data"); + let _recovery_guard = recovery_span.enter(); + warn!(component = "obs_contract_test", message = "inside recovery monitor"); + }); + + let raw = String::from_utf8( + writer + .inner + .lock() + .expect("shared writer lock should not be poisoned") + .clone(), + ) + .expect("log output should be valid UTF-8"); + let first_line = raw.lines().next().expect("expected one JSON log line"); + serde_json::from_str(first_line).expect("log line should be valid JSON") + } + #[test] /// Invalid file names should be reported as errors instead of panicking. fn test_init_file_logging_invalid_filename_does_not_panic() { @@ -661,9 +785,9 @@ mod tests { assert!(!raw.contains('\u{1b}'), "JSON log output must not contain ANSI escape sequences: {raw}"); assert_eq!(parsed.get("level").and_then(Value::as_str), Some("INFO")); - assert_eq!(parsed["fields"]["message"], Value::String("hello world".to_string())); - assert_eq!(parsed["fields"]["component"], Value::String("obs_contract_test".to_string())); - assert_eq!(parsed["fields"]["sink"], Value::String("stdout".to_string())); + assert_eq!(parsed["message"], Value::String("hello world".to_string())); + assert_eq!(parsed["component"], Value::String("obs_contract_test".to_string())); + assert_eq!(parsed["sink"], Value::String("stdout".to_string())); assert!(parsed.get("target").is_some(), "expected target field in JSON log output"); } @@ -686,18 +810,54 @@ mod tests { .collect::>(); assert_eq!(plain_keys, color_keys, "top-level JSON log keys must stay stable across sinks"); - let plain_field_keys = plain["fields"] + let plain_field_keys = plain .as_object() - .expect("plain JSON log fields should be an object") + .expect("plain JSON log should be an object") .keys() + .filter(|key| { + !matches!( + key.as_str(), + "timestamp" | "level" | "target" | "filename" | "line_number" | "threadName" | "threadId" + ) + }) .cloned() .collect::>(); - let color_field_keys = color["fields"] + let color_field_keys = color .as_object() - .expect("color JSON log fields should be an object") + .expect("color JSON log should be an object") .keys() + .filter(|key| { + !matches!( + key.as_str(), + "timestamp" | "level" | "target" | "filename" | "line_number" | "threadName" | "threadId" + ) + }) .cloned() .collect::>(); assert_eq!(plain_field_keys, color_field_keys, "structured field keys must stay stable across sinks"); } + + #[test] + fn test_json_log_layer_promotes_request_id_from_current_span() { + let parsed = render_json_log_with_request_span(); + + assert_eq!(parsed["request_id"], Value::String("req-123".to_string())); + assert_eq!(parsed["message"], Value::String("inside request span".to_string())); + assert_eq!(parsed["span"]["request_id"], Value::String("req-123".to_string())); + } + + #[test] + fn test_json_log_layer_promotes_parent_request_id_for_recovery_monitor_child_span() { + let parsed = render_json_log_with_recovery_monitor_child_span(); + + assert_eq!(parsed["request_id"], Value::String("req-parent".to_string())); + assert_eq!(parsed["message"], Value::String("inside recovery monitor".to_string())); + assert_eq!(parsed["span"]["name"], Value::String("recovery-monitor".to_string())); + assert_eq!(parsed["span"]["kind"], Value::String("remote_disk".to_string())); + let spans = parsed["spans"].as_array().expect("spans should be present"); + assert!(spans.iter().any(|span| { + span.get("name").and_then(Value::as_str) == Some("request-span") + && span.get("request_id").and_then(Value::as_str) == Some("req-parent") + })); + } } diff --git a/crates/protocols/src/sftp/driver.rs b/crates/protocols/src/sftp/driver.rs index b1833009d..d4c4d1289 100644 --- a/crates/protocols/src/sftp/driver.rs +++ b/crates/protocols/src/sftp/driver.rs @@ -40,6 +40,7 @@ use std::collections::HashMap; use std::sync::atomic::AtomicU64; use std::sync::{Arc, LazyLock}; use tokio::sync::Semaphore; +use tracing::Instrument; use uuid::Uuid; /// Permits available to the fire-and-forget AbortMultipartUpload tasks @@ -1086,88 +1087,91 @@ impl Drop for SftpDriver { } }; - tokio::spawn(async move { - let _permit = permit; - tracing::warn!( - bucket = %bucket, - key = %key, - upload_id = %upload_id, - peer = %peer, - "aborting orphaned multipart upload on session drop" - ); - // Build AbortMultipartUploadInput inside the spawned - // task so the builder Result is handled in async - // context. The builder only fails on missing required - // fields. bucket, key, and upload_id are all set, so - // log and return on any unexpected failure. - let input = match AbortMultipartUploadInput::builder() - .bucket(bucket.clone()) - .key(key.clone()) - .upload_id(upload_id.clone()) - .build() - { - Ok(input) => input, - Err(e) => { - tracing::error!( - bucket = %bucket, - key = %key, - upload_id = %upload_id, - err = %e, - "failed to build AbortMultipartUploadInput on session drop" - ); - return; - } - }; - match tokio::time::timeout( - std::time::Duration::from_secs(backend_op_timeout_secs), - storage.abort_multipart_upload(input, &access_key, &secret_key), - ) - .await - { - Ok(Ok(_)) => {} - Ok(Err(e)) => { - // close() removes the tombstone only on Ok, so Drop - // retries any abort whose inline attempt caused an - // error. A retried abort can race a concurrent - // successful CompleteMultipartUpload, returning - // NoSuchUpload. Log at debug to keep error-level - // logs reserved for genuine abort failures. - if is_no_such_upload_error(&e) { - tracing::debug!( - bucket = %bucket, - key = %key, - upload_id = %upload_id, - "Drop abort returned NoSuchUpload: upload already completed or aborted", - ); - } else { + tokio::spawn( + async move { + let _permit = permit; + tracing::warn!( + bucket = %bucket, + key = %key, + upload_id = %upload_id, + peer = %peer, + "aborting orphaned multipart upload on session drop" + ); + // Build AbortMultipartUploadInput inside the spawned + // task so the builder Result is handled in async + // context. The builder only fails on missing required + // fields. bucket, key, and upload_id are all set, so + // log and return on any unexpected failure. + let input = match AbortMultipartUploadInput::builder() + .bucket(bucket.clone()) + .key(key.clone()) + .upload_id(upload_id.clone()) + .build() + { + Ok(input) => input, + Err(e) => { tracing::error!( bucket = %bucket, key = %key, upload_id = %upload_id, err = %e, - "failed to abort orphaned multipart upload" + "failed to build AbortMultipartUploadInput on session drop" + ); + return; + } + }; + match tokio::time::timeout( + std::time::Duration::from_secs(backend_op_timeout_secs), + storage.abort_multipart_upload(input, &access_key, &secret_key), + ) + .await + { + Ok(Ok(_)) => {} + Ok(Err(e)) => { + // close() removes the tombstone only on Ok, so Drop + // retries any abort whose inline attempt caused an + // error. A retried abort can race a concurrent + // successful CompleteMultipartUpload, returning + // NoSuchUpload. Log at debug to keep error-level + // logs reserved for genuine abort failures. + if is_no_such_upload_error(&e) { + tracing::debug!( + bucket = %bucket, + key = %key, + upload_id = %upload_id, + "Drop abort returned NoSuchUpload: upload already completed or aborted", + ); + } else { + tracing::error!( + bucket = %bucket, + key = %key, + upload_id = %upload_id, + err = %e, + "failed to abort orphaned multipart upload" + ); + } + } + Err(_elapsed) => { + // Drop's abort task is bounded by the same + // per-call deadline as inline backend calls. + // A timeout here is rare (the runtime drains + // session tasks for SHUTDOWN_DRAIN_TIMEOUT_SECS + // and Drop runs after that), so log at warn so + // operators can correlate the orphaned upload + // with the bucket AbortIncompleteMultipartUpload + // lifecycle rule that will reclaim it. + tracing::warn!( + bucket = %bucket, + key = %key, + upload_id = %upload_id, + timeout_secs = backend_op_timeout_secs, + "Drop abort of orphaned multipart upload timed out; bucket lifecycle rule must reclaim parts", ); } } - Err(_elapsed) => { - // Drop's abort task is bounded by the same - // per-call deadline as inline backend calls. - // A timeout here is rare (the runtime drains - // session tasks for SHUTDOWN_DRAIN_TIMEOUT_SECS - // and Drop runs after that), so log at warn so - // operators can correlate the orphaned upload - // with the bucket AbortIncompleteMultipartUpload - // lifecycle rule that will reclaim it. - tracing::warn!( - bucket = %bucket, - key = %key, - upload_id = %upload_id, - timeout_secs = backend_op_timeout_secs, - "Drop abort of orphaned multipart upload timed out; bucket lifecycle rule must reclaim parts", - ); - } } - }); + .instrument(tracing::Span::current()), + ); } } } diff --git a/crates/protocols/src/webdav/server.rs b/crates/protocols/src/webdav/server.rs index 007ab349d..85b0eb560 100644 --- a/crates/protocols/src/webdav/server.rs +++ b/crates/protocols/src/webdav/server.rs @@ -35,7 +35,7 @@ use std::time::Duration; use tokio::net::TcpListener; use tokio::sync::{broadcast, watch}; use tokio_rustls::TlsAcceptor; -use tracing::{debug, error, info, warn}; +use tracing::{Instrument, debug, error, info, info_span, warn}; const LOG_COMPONENT_PROTOCOLS: &str = "protocols"; const LOG_SUBSYSTEM_WEBDAV_SERVER: &str = "webdav_server"; @@ -144,54 +144,61 @@ where let tls_acceptor = tls_acceptor.clone(); let max_body_size = self.config.max_body_size; - tokio::spawn(async move { - let source_ip: IpAddr = addr.ip(); - - if let Some(acceptor) = tls_acceptor { - match acceptor.accept(stream).await { - Ok(tls_stream) => { - let io = TokioIo::new(tls_stream); - if let Err(e) = Self::handle_connection_impl(io, storage, source_ip, max_body_size).await { + let source_ip: IpAddr = addr.ip(); + let span = info_span!( + "webdav-connection", + peer = %source_ip, + transport = if tls_acceptor.is_some() { "tls" } else { "tcp" }, + ); + tokio::spawn( + async move { + if let Some(acceptor) = tls_acceptor { + match acceptor.accept(stream).await { + Ok(tls_stream) => { + let io = TokioIo::new(tls_stream); + if let Err(e) = Self::handle_connection_impl(io, storage, source_ip, max_body_size).await { + debug!( + event = EVENT_WEBDAV_CONNECTION_STATE, + component = LOG_COMPONENT_PROTOCOLS, + subsystem = LOG_SUBSYSTEM_WEBDAV_SERVER, + result = "error", + peer = %source_ip, + transport = "tls", + error = %e, + "webdav connection ended with error" + ); + } + } + Err(e) => { debug!( event = EVENT_WEBDAV_CONNECTION_STATE, component = LOG_COMPONENT_PROTOCOLS, subsystem = LOG_SUBSYSTEM_WEBDAV_SERVER, - result = "error", + result = "tls_handshake_failed", peer = %source_ip, - transport = "tls", error = %e, "webdav connection ended with error" ); } } - Err(e) => { + } else { + let io = TokioIo::new(stream); + if let Err(e) = Self::handle_connection_impl(io, storage, source_ip, max_body_size).await { debug!( event = EVENT_WEBDAV_CONNECTION_STATE, component = LOG_COMPONENT_PROTOCOLS, subsystem = LOG_SUBSYSTEM_WEBDAV_SERVER, - result = "tls_handshake_failed", + result = "error", peer = %source_ip, + transport = "tcp", error = %e, "webdav connection ended with error" ); } } - } else { - let io = TokioIo::new(stream); - if let Err(e) = Self::handle_connection_impl(io, storage, source_ip, max_body_size).await { - debug!( - event = EVENT_WEBDAV_CONNECTION_STATE, - component = LOG_COMPONENT_PROTOCOLS, - subsystem = LOG_SUBSYSTEM_WEBDAV_SERVER, - result = "error", - peer = %source_ip, - transport = "tcp", - error = %e, - "webdav connection ended with error" - ); - } } - }); + .instrument(span), + ); } Err(e) => { error!( diff --git a/rustfs/src/admin/console.rs b/rustfs/src/admin/console.rs index 500f39c39..55314341d 100644 --- a/rustfs/src/admin/console.rs +++ b/rustfs/src/admin/console.rs @@ -15,7 +15,11 @@ use crate::admin::handlers::health::{HealthProbe, build_health_response_parts, collect_dependency_readiness}; use crate::license::has_valid_license; use crate::server::has_path_prefix; -use crate::server::{CONSOLE_PREFIX, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, LICENSE, RUSTFS_ADMIN_PREFIX, VERSION}; +use crate::server::{ + CONSOLE_PREFIX, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, HeaderMapCarrier, LICENSE, RUSTFS_ADMIN_PREFIX, + RequestContextLayer, VERSION, +}; +use crate::storage::request_context::RequestContext; use crate::version::build; use axum::{ Json, Router, @@ -27,6 +31,7 @@ use axum::{ }; use http::{HeaderMap, HeaderName, HeaderValue, Method, StatusCode, Uri}; use mime_guess::from_path; +use opentelemetry::global; use rust_embed::RustEmbed; use serde::Serialize; use std::{ @@ -41,7 +46,8 @@ 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}; +use tracing::{Span, debug, error, info, instrument, warn}; +use tracing_opentelemetry::OpenTelemetrySpanExt; #[derive(RustEmbed)] #[folder = "$CARGO_MANIFEST_DIR/static"] @@ -491,9 +497,45 @@ fn setup_console_middleware_stack( // Add comprehensive middleware layers using tower-http features app = app .layer(CatchPanicLayer::new()) + .layer( + TraceLayer::new_for_http() + .make_span_with(|request: &Request| { + let request_context = request.extensions().get::(); + let request_id = request_context.map(|ctx| ctx.request_id.as_str()).unwrap_or("unknown"); + let trace_id = request_context.and_then(|ctx| ctx.trace_id.as_deref()).unwrap_or("unknown"); + let span_id = request_context.and_then(|ctx| ctx.span_id.as_deref()).unwrap_or("unknown"); + + let parent_context = global::get_text_map_propagator(|propagator| { + propagator.extract(&HeaderMapCarrier::new(request.headers())) + }); + + let span = tracing::info_span!( + "console-request", + request_id = %request_id, + trace_id = %trace_id, + span_id = %span_id, + method = %request.method(), + uri = %request.uri(), + status_code = tracing::field::Empty, + ); + + if span.is_disabled() { + return span; + } + + if let Err(err) = span.set_parent(parent_context) { + debug!(error = ?err, "Failed to propagate tracing context for console request"); + } + + span + }) + .on_response(|response: &Response, _latency: Duration, span: &Span| { + span.record("status_code", tracing::field::display(response.status())); + }), + ) .layer(PropagateRequestIdLayer::x_request_id()) + .layer(RequestContextLayer) .layer(SetRequestIdLayer::x_request_id(MakeRequestUuid)) - .layer(TraceLayer::new_for_http()) // Compress responses .layer(CompressionLayer::new()) .layer(middleware::from_fn(console_logging_middleware)) @@ -684,12 +726,47 @@ pub(crate) fn make_console_server() -> Router { mod tests { use super::*; use axum::body::Body; + use axum::routing::get; use http::{Request, StatusCode}; use http_body_util::BodyExt; use serial_test::serial; + use std::io; use std::net::{IpAddr, Ipv4Addr}; + use std::sync::{Arc, Mutex}; use temp_env::async_with_vars; use tower::ServiceExt; + use tracing_subscriber::{Registry, layer::SubscriberExt}; + + #[derive(Clone, Default)] + struct SharedWriter { + inner: Arc>>, + } + + struct SharedWriterGuard { + inner: Arc>>, + } + + impl<'writer> tracing_subscriber::fmt::MakeWriter<'writer> for SharedWriter { + type Writer = SharedWriterGuard; + + fn make_writer(&'writer self) -> Self::Writer { + SharedWriterGuard { + inner: Arc::clone(&self.inner), + } + } + } + + impl io::Write for SharedWriterGuard { + fn write(&mut self, buf: &[u8]) -> io::Result { + let mut inner = self.inner.lock().expect("shared writer lock should not be poisoned"); + inner.extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } #[test] fn console_api_base_url_keeps_rustfs_admin_prefix() { @@ -739,6 +816,112 @@ mod tests { ); } + #[tokio::test(flavor = "current_thread")] + async fn console_trace_layer_records_request_id_on_current_span() { + let writer = SharedWriter::default(); + let subscriber = Registry::default().with( + tracing_subscriber::fmt::layer() + .without_time() + .with_target(false) + .with_level(false) + .with_ansi(false) + .json() + .flatten_event(true) + .with_current_span(true) + .with_span_list(true) + .with_writer(writer.clone()), + ); + let _guard = tracing::subscriber::set_default(subscriber); + + let app = Router::new() + .route( + "/trace-test", + get(|| async { + tracing::info!("console handler log"); + StatusCode::OK + }), + ) + .layer( + TraceLayer::new_for_http() + .make_span_with(|request: &http::Request<_>| { + let request_context = request.extensions().get::(); + let request_id = request_context.map(|ctx| ctx.request_id.as_str()).unwrap_or("unknown"); + let trace_id = request_context.and_then(|ctx| ctx.trace_id.as_deref()).unwrap_or("unknown"); + let span_id = request_context.and_then(|ctx| ctx.span_id.as_deref()).unwrap_or("unknown"); + + let parent_context = global::get_text_map_propagator(|propagator| { + propagator.extract(&HeaderMapCarrier::new(request.headers())) + }); + + let span = tracing::info_span!( + "console-request", + request_id = %request_id, + trace_id = %trace_id, + span_id = %span_id, + method = %request.method(), + uri = %request.uri(), + status_code = tracing::field::Empty, + ); + + if span.is_disabled() { + return span; + } + + if let Err(err) = span.set_parent(parent_context) { + debug!(error = ?err, "Failed to propagate tracing context for console request"); + } + + span + }) + .on_response(|response: &Response, _latency: Duration, span: &Span| { + span.record("status_code", tracing::field::display(response.status())); + }), + ) + .layer(RequestContextLayer) + .layer(SetRequestIdLayer::x_request_id(MakeRequestUuid)); + + app.oneshot( + Request::builder() + .uri("/trace-test") + .body(Body::empty()) + .expect("failed to build trace test request"), + ) + .await + .expect("trace test request should complete"); + + let output = String::from_utf8( + writer + .inner + .lock() + .expect("shared writer lock should not be poisoned") + .clone(), + ) + .expect("console trace log output should be valid UTF-8"); + + let log = output + .lines() + .map(|line| serde_json::from_str::(line).expect("console log line should be valid JSON")) + .find(|value| value.get("message").and_then(serde_json::Value::as_str) == Some("console handler log")) + .expect("expected console handler log entry"); + + assert_eq!(log["span"]["method"], serde_json::Value::String("GET".to_string())); + assert_eq!(log["span"]["name"], serde_json::Value::String("console-request".to_string())); + assert!( + log.get("span") + .and_then(|span| span.get("request_id")) + .and_then(serde_json::Value::as_str) + .is_some(), + "{output}" + ); + assert_ne!( + log.get("span") + .and_then(|span| span.get("request_id")) + .and_then(serde_json::Value::as_str), + Some("unknown"), + "{output}" + ); + } + /// Regression: when no console CORS origins are configured (the new /// default), the layer must NOT emit `Access-Control-Allow-Origin`, so /// browsers treat responses as same-origin only. diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index 431c59fdf..9d0054721 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -17,6 +17,7 @@ use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::app::context::resolve_object_store_handle; use crate::server::ADMIN_PREFIX; use crate::server::RemoteAddr; +use crate::storage::request_context::spawn_traced; use bytes::Bytes; use http::{HeaderMap, HeaderValue, Uri}; use hyper::{Method, StatusCode}; @@ -34,7 +35,6 @@ use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; use serde::{Deserialize, Serialize}; use std::path::PathBuf; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; -use tokio::spawn; use tokio::sync::mpsc; use tracing::{info, warn}; @@ -395,7 +395,7 @@ impl Operation for HealHandler { let tx_clone = tx.clone(); let heal_path_str = heal_path.to_str().unwrap_or_default().to_string(); let client_token = hip.client_token.clone(); - spawn(async move { + spawn_traced(async move { match rustfs_common::heal_channel::query_heal_status(heal_path_str, client_token).await { Ok(response) if response.success => { let (summary, items) = heal_channel_response_status(&response); @@ -445,7 +445,7 @@ impl Operation for HealHandler { let client_token = hip.client_token.clone(); let client_address = client_address.clone(); let heal_settings = hip.hs; - spawn(async move { + spawn_traced(async move { match rustfs_common::heal_channel::cancel_heal_task(heal_path_str, client_token.clone()).await { Ok(response) if response.success => { let resp_bytes = if client_token.is_empty() { @@ -495,7 +495,7 @@ impl Operation for HealHandler { // Use new heal channel mechanism let tx_clone = tx.clone(); let client_address = client_address.clone(); - spawn(async move { + spawn_traced(async move { // Create heal request through channel let heal_request = build_heal_channel_request(&hip); let client_token = heal_request.id.clone(); diff --git a/rustfs/src/admin/handlers/metrics.rs b/rustfs/src/admin/handlers/metrics.rs index b0af4211f..0ae43236f 100644 --- a/rustfs/src/admin/handlers/metrics.rs +++ b/rustfs/src/admin/handlers/metrics.rs @@ -19,6 +19,7 @@ //! exposition endpoint. use crate::admin::router::Operation; +use crate::storage::request_context::spawn_traced; use bytes::Bytes; use futures::{Stream, StreamExt}; use http::{HeaderMap, HeaderValue, Uri}; @@ -35,9 +36,9 @@ use std::collections::{HashMap, HashSet}; use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Duration as StdDuration; +use tokio::select; use tokio::sync::mpsc; use tokio::time::interval; -use tokio::{select, spawn}; use tokio_stream::wrappers::ReceiverStream; use tracing::{debug, error, warn}; @@ -217,7 +218,7 @@ impl Operation for MetricsHandler { }); let body = Body::from(in_stream); - spawn(async move { + spawn_traced(async move { while n > 0 { let mut metrics = RealtimeMetrics::default(); let local_metrics = collect_local_metrics(types, &opts).await; diff --git a/rustfs/src/admin/handlers/tier.rs b/rustfs/src/admin/handlers/tier.rs index e1ccc9039..d6720516c 100644 --- a/rustfs/src/admin/handlers/tier.rs +++ b/rustfs/src/admin/handlers/tier.rs @@ -21,6 +21,7 @@ use crate::{ app::context::{resolve_object_store_handle, resolve_tier_config_handle}, auth::{check_key_valid, get_session_token}, server::{ADMIN_PREFIX, RemoteAddr}, + storage::request_context::spawn_traced, }; use http::Uri; use http::{HeaderMap, StatusCode}; @@ -54,7 +55,6 @@ use s3s::{ use serde_urlencoded::from_bytes; use std::collections::HashMap; use time::OffsetDateTime; -use tokio::spawn; use tracing::{debug, warn}; const LOG_COMPONENT_ADMIN: &str = "admin"; @@ -91,7 +91,7 @@ pub struct AddTier {} fn spawn_transition_tier_config_propagation(action: &'static str) { if let Some(notification_sys) = get_global_notification_sys() { - spawn(async move { + spawn_traced(async move { for peer_result in notification_sys.load_transition_tier_config().await { if let Some(err) = peer_result.err { warn!( diff --git a/rustfs/src/admin/router.rs b/rustfs/src/admin/router.rs index b2211c6c0..6aa2a5da1 100644 --- a/rustfs/src/admin/router.rs +++ b/rustfs/src/admin/router.rs @@ -23,6 +23,7 @@ use crate::server::{ ADMIN_PREFIX, HEALTH_PREFIX, HEALTH_READY_PATH, MINIO_ADMIN_PREFIX, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH, is_admin_path, }; use crate::storage::access::{ReqInfo, authorize_request}; +use crate::storage::request_context::spawn_traced; use aws_sdk_s3::primitives::ByteStream as AwsByteStream; use bytes::Bytes; use futures::{Stream, StreamExt}; @@ -1310,7 +1311,7 @@ fn build_listen_notification_response(uri: &Uri, bucket: Option<&str>) -> S3Resu inner: ReceiverStream::new(rx), }); - tokio::spawn(async move { + spawn_traced(async move { let mut ticker = tokio::time::interval(interval_duration); let mut peer_ticker = tokio::time::interval(LISTEN_NOTIFICATION_PEER_POLL_INTERVAL); // Skip the immediate first tick so behavior starts after interval duration. diff --git a/rustfs/src/server/mod.rs b/rustfs/src/server/mod.rs index f60edc87d..34836aeae 100644 --- a/rustfs/src/server/mod.rs +++ b/rustfs/src/server/mod.rs @@ -41,7 +41,9 @@ pub use service_state::ShutdownSignal; pub use service_state::wait_for_shutdown; // Items only used within the library crate (admin handlers, server/http.rs, etc.). +pub(crate) use http::HeaderMapCarrier; pub(crate) use http::active_http_requests; +pub(crate) use layer::RequestContextLayer; pub(crate) use module_switch::{ ModuleSwitchSnapshot, ModuleSwitchSource, PersistedModuleSwitches, current_module_switch_snapshot, refresh_persisted_module_switches_from_store, save_persisted_module_switches_to_store, validate_module_switch_update, diff --git a/rustfs/src/storage/rpc/http_service.rs b/rustfs/src/storage/rpc/http_service.rs index 2e230fb94..ddb58f4f2 100644 --- a/rustfs/src/storage/rpc/http_service.rs +++ b/rustfs/src/storage/rpc/http_service.rs @@ -13,6 +13,7 @@ // limitations under the License. use crate::server::RPC_PREFIX; +use crate::storage::request_context::spawn_traced; use bytes::{Bytes, BytesMut}; use futures_util::TryStreamExt; use http::{HeaderMap, Method, Request, Response, StatusCode, Uri}; @@ -243,7 +244,7 @@ async fn handle_walk_dir(req: Request) -> Response { }; let (rd, mut wd) = tokio::io::duplex(DEFAULT_READ_BUFFER_SIZE); - tokio::spawn(async move { + spawn_traced(async move { if let Err(e) = disk.walk_dir(args, &mut wd).await { warn!(error = %e, "walk_dir failed"); } diff --git a/scripts/run.sh b/scripts/run.sh index 4e798ada2..987d6afee 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -77,7 +77,7 @@ fi # 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