From 33eff4c3c4514ee299c4895cbdcdeb20e38cfcce Mon Sep 17 00:00:00 2001 From: houseme Date: Sun, 16 Aug 2026 23:36:46 +0800 Subject: [PATCH] test(ecstore): add metadata slow-tail fault hook (#6150) Add a diagnostic metadata-only read_version delay hook for GET data-read fanout so bounded/default behavior can be compared under controlled slow-tail metadata responses. Co-authored-by: heihutu --- .../src/set_disk/core/io_primitives.rs | 134 ++++++++++++++++++ crates/ecstore/src/set_disk/mod.rs | 97 ++++++++++++- 2 files changed, 229 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 299299971..4e0e0e759 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -2446,6 +2446,7 @@ impl SetDisks { let bucket: Arc = Arc::from(bucket); let object: Arc = Arc::from(object); let version_id: Arc = Arc::from(version_id); + let slowtail_fault = get_metadata_slowtail_fault_request(bucket.as_ref(), object.as_ref(), read_data); let futures = disks.iter().enumerate().map(|(disk_index, disk)| { let disk = disk.clone(); let task_opts = opts; @@ -2453,10 +2454,14 @@ impl SetDisks { let bucket = bucket.clone(); let object = object.clone(); let version_id = version_id.clone(); + let slowtail_fault = slowtail_fault.clone(); tokio::spawn(async move { let response_start = observe.then(Instant::now); let result = if let Some(disk) = disk { Self::record_read_version_call(&object, disk_index); + if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(disk_index)) { + tokio::time::sleep(delay).await; + } disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts) .await } else { @@ -2552,6 +2557,7 @@ impl SetDisks { let mut scheduled_count = 0usize; let mut force_full_wait = false; let mut final_miss_reason_override = None; + let slowtail_fault = get_metadata_slowtail_fault_request(bucket.as_ref(), object.as_ref(), read_data); let spawn_read_version = |join_set: &mut JoinSet<(usize, disk::error::Result, Duration)>, index: usize, disk: Option| { let task_opts = opts; @@ -2559,6 +2565,7 @@ impl SetDisks { let bucket = bucket.clone(); let object = object.clone(); let version_id = version_id.clone(); + let slowtail_fault = slowtail_fault.clone(); join_set.spawn(async move { let response_start = Instant::now(); let result = if let Some(disk) = disk { @@ -2567,6 +2574,9 @@ impl SetDisks { Self::record_read_version_call(&object, index); #[cfg(test)] Self::read_version_fanout_barrier(&object, index).await; + if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(index)) { + tokio::time::sleep(delay).await; + } disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts) .await } else { @@ -5744,6 +5754,130 @@ mod tests { (dirs, disks) } + #[test] + fn metadata_slowtail_fault_delay_parses_and_filters_request() { + temp_env::with_vars( + [ + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DELAY_MS, Some("25")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DISKS, Some("1,3")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_BUCKET, Some("bench-bucket")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_OBJECT_PREFIX, Some("objects/")), + ], + || { + assert_eq!( + get_metadata_slowtail_fault_delay("bench-bucket", "objects/000001", 3, true), + Some(Duration::from_millis(25)) + ); + assert!(get_metadata_slowtail_fault_delay("bench-bucket", "objects/000001", 2, true).is_none()); + assert!(get_metadata_slowtail_fault_delay("other-bucket", "objects/000001", 3, true).is_none()); + assert!(get_metadata_slowtail_fault_delay("bench-bucket", "other/000001", 3, true).is_none()); + assert!(get_metadata_slowtail_fault_delay("bench-bucket", "objects/000001", 3, false).is_none()); + }, + ); + } + + #[test] + fn metadata_slowtail_fault_delay_disables_invalid_disk_list() { + temp_env::with_vars( + [ + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DELAY_MS, Some("25")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DISKS, Some("1,nope")), + ], + || { + assert!(get_metadata_slowtail_fault_delay("bucket", "object", 1, true).is_none()); + }, + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn metadata_slowtail_fault_delays_only_data_read_metadata_task() { + const DISKS: usize = 4; + let bucket = "metadata-slowtail-fault-bucket"; + let object = "objects/metadata-slowtail-fault-object"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_metadata_fanout_fileinfo(&disks, bucket, object, None).await; + + temp_env::async_with_vars( + [ + (ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE, Some("false")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DELAY_MS, Some("150")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DISKS, Some("3")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_BUCKET, Some(bucket)), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_OBJECT_PREFIX, Some("objects/")), + ], + async { + let read_without_data = + SetDisks::read_all_fileinfo_observed(&disks, bucket, bucket, object, "", false, false, false, true, 2); + tokio::time::timeout(Duration::from_millis(100), read_without_data) + .await + .expect("non-data metadata fanout must not be delayed by the data-read slowtail hook") + .expect("metadata fanout without read_data should resolve"); + + let mut read_with_data = Box::pin(SetDisks::read_all_fileinfo_observed( + &disks, bucket, bucket, object, "", true, false, false, true, 2, + )); + assert!( + tokio::time::timeout(Duration::from_millis(40), &mut read_with_data) + .await + .is_err(), + "data-read metadata fanout must wait for the injected slow read_version response" + ); + let (parts_metadata, errs, diagnostics) = tokio::time::timeout(Duration::from_secs(2), read_with_data) + .await + .expect("injected slowtail should eventually complete") + .expect("data-read metadata fanout should resolve"); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + assert_eq!(diagnostics.total_responses(), DISKS); + }, + ) + .await; + + drop(dirs); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn metadata_slowtail_fault_delays_early_stop_metadata_task() { + const DISKS: usize = 4; + let bucket = "metadata-slowtail-early-stop-bucket"; + let object = "objects/metadata-slowtail-early-stop-object"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_metadata_fanout_fileinfo(&disks, bucket, object, None).await; + + temp_env::async_with_vars( + [ + (ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE, Some("true")), + (ENV_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE, Some("true")), + (ENV_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT, Some("false")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DELAY_MS, Some("150")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DISKS, Some("3")), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_BUCKET, Some(bucket)), + (ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_OBJECT_PREFIX, Some("objects/")), + ], + async { + let mut read_with_data = Box::pin(SetDisks::read_all_fileinfo_observed( + &disks, bucket, bucket, object, "", true, false, false, true, 2, + )); + assert!( + tokio::time::timeout(Duration::from_millis(40), &mut read_with_data) + .await + .is_err(), + "early-stop metadata fanout must still wait for the injected slow response after fallback to full wait" + ); + let (parts_metadata, errs, diagnostics) = tokio::time::timeout(Duration::from_secs(2), read_with_data) + .await + .expect("injected early-stop slowtail should eventually complete") + .expect("early-stop metadata fanout should resolve"); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + assert_eq!(diagnostics.total_responses(), DISKS); + }, + ) + .await; + + drop(dirs); + } + /// Demo / regression guard for the backlog#1325 per-disk call counters. /// /// The metadata fan-out issues each `read_version` inside its own diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index ad5db566b..a6a875a7b 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -174,15 +174,14 @@ use std::future::Future; use std::hash::{BuildHasher, Hash, Hasher}; use std::mem::{self}; use std::pin::Pin; -use std::sync::OnceLock; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::{Arc, OnceLock}; use std::task::{Context, Poll}; use std::time::{Instant, SystemTime, UNIX_EPOCH}; use std::{ collections::{HashMap, HashSet}, io::{Cursor, Write}, path::Path, - sync::Arc, time::Duration, }; use time::OffsetDateTime; @@ -717,6 +716,11 @@ const DEFAULT_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE: bool = true; const ENV_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT: &str = "RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT"; const DEFAULT_RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT: bool = false; +const ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DELAY_MS: &str = "RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DELAY_MS"; +const ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DISKS: &str = "RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DISKS"; +const ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_BUCKET: &str = "RUSTFS_GET_METADATA_SLOWTAIL_FAULT_BUCKET"; +const ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_OBJECT_PREFIX: &str = "RUSTFS_GET_METADATA_SLOWTAIL_FAULT_OBJECT_PREFIX"; + // --- Multipart Reader-Setup Prefetch Configuration (backlog#870) --- const ENV_RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH: &str = "RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH"; @@ -1695,6 +1699,95 @@ fn is_get_metadata_early_stop_bounded_fanout_enabled() -> bool { } } +#[derive(Debug)] +struct GetMetadataSlowtailFaultConfig { + delay: Duration, + disks: Arc<[usize]>, + bucket: Option, + object_prefix: Option, +} + +#[derive(Clone, Debug)] +struct GetMetadataSlowtailFaultRequest { + delay: Duration, + disks: Arc<[usize]>, +} + +impl GetMetadataSlowtailFaultRequest { + fn delay_for_disk(&self, disk_index: usize) -> Option { + self.disks.contains(&disk_index).then_some(self.delay) + } +} + +fn parse_get_metadata_slowtail_fault_disks(raw: &str) -> Option> { + let mut disks = Vec::new(); + for item in raw.split(',').map(str::trim).filter(|item| !item.is_empty()) { + let Ok(index) = item.parse::() else { + return None; + }; + if !disks.contains(&index) { + disks.push(index); + } + } + (!disks.is_empty()).then_some(disks) +} + +fn load_get_metadata_slowtail_fault_config() -> Option { + let delay_ms = rustfs_utils::get_env_u64(ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DELAY_MS, 0); + if delay_ms == 0 { + return None; + } + let disks = parse_get_metadata_slowtail_fault_disks(&std::env::var(ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_DISKS).ok()?)?; + let bucket = std::env::var(ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_BUCKET) + .ok() + .filter(|value| !value.is_empty()); + let object_prefix = std::env::var(ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_OBJECT_PREFIX) + .ok() + .filter(|value| !value.is_empty()); + Some(GetMetadataSlowtailFaultConfig { + delay: Duration::from_millis(delay_ms), + disks: Arc::from(disks.into_boxed_slice()), + bucket, + object_prefix, + }) +} + +fn get_metadata_slowtail_fault_request(bucket: &str, object: &str, read_data: bool) -> Option { + if !read_data { + return None; + } + + #[cfg(test)] + let config = load_get_metadata_slowtail_fault_config(); + #[cfg(test)] + let config = config.as_ref()?; + #[cfg(not(test))] + let config = ({ + static CACHED: OnceLock> = OnceLock::new(); + CACHED.get_or_init(load_get_metadata_slowtail_fault_config).as_ref() + })?; + + if let Some(expected_bucket) = &config.bucket + && expected_bucket != bucket + { + return None; + } + if let Some(expected_prefix) = &config.object_prefix + && !object.starts_with(expected_prefix) + { + return None; + } + Some(GetMetadataSlowtailFaultRequest { + delay: config.delay, + disks: config.disks.clone(), + }) +} + +#[cfg(test)] +fn get_metadata_slowtail_fault_delay(bucket: &str, object: &str, disk_index: usize, read_data: bool) -> Option { + get_metadata_slowtail_fault_request(bucket, object, read_data)?.delay_for_disk(disk_index) +} + /// Check if multipart reads prefetch the next part's bitrot reader setup /// while the current part decodes (backlog#870). ///