mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-17 10:17:55 +00:00
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 <heihutu@gmail.com>
This commit is contained in:
@@ -2446,6 +2446,7 @@ impl SetDisks {
|
|||||||
let bucket: Arc<str> = Arc::from(bucket);
|
let bucket: Arc<str> = Arc::from(bucket);
|
||||||
let object: Arc<str> = Arc::from(object);
|
let object: Arc<str> = Arc::from(object);
|
||||||
let version_id: Arc<str> = Arc::from(version_id);
|
let version_id: Arc<str> = 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 futures = disks.iter().enumerate().map(|(disk_index, disk)| {
|
||||||
let disk = disk.clone();
|
let disk = disk.clone();
|
||||||
let task_opts = opts;
|
let task_opts = opts;
|
||||||
@@ -2453,10 +2454,14 @@ impl SetDisks {
|
|||||||
let bucket = bucket.clone();
|
let bucket = bucket.clone();
|
||||||
let object = object.clone();
|
let object = object.clone();
|
||||||
let version_id = version_id.clone();
|
let version_id = version_id.clone();
|
||||||
|
let slowtail_fault = slowtail_fault.clone();
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let response_start = observe.then(Instant::now);
|
let response_start = observe.then(Instant::now);
|
||||||
let result = if let Some(disk) = disk {
|
let result = if let Some(disk) = disk {
|
||||||
Self::record_read_version_call(&object, disk_index);
|
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)
|
disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts)
|
||||||
.await
|
.await
|
||||||
} else {
|
} else {
|
||||||
@@ -2552,6 +2557,7 @@ impl SetDisks {
|
|||||||
let mut scheduled_count = 0usize;
|
let mut scheduled_count = 0usize;
|
||||||
let mut force_full_wait = false;
|
let mut force_full_wait = false;
|
||||||
let mut final_miss_reason_override = None;
|
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 =
|
let spawn_read_version =
|
||||||
|join_set: &mut JoinSet<(usize, disk::error::Result<FileInfo>, Duration)>, index: usize, disk: Option<DiskStore>| {
|
|join_set: &mut JoinSet<(usize, disk::error::Result<FileInfo>, Duration)>, index: usize, disk: Option<DiskStore>| {
|
||||||
let task_opts = opts;
|
let task_opts = opts;
|
||||||
@@ -2559,6 +2565,7 @@ impl SetDisks {
|
|||||||
let bucket = bucket.clone();
|
let bucket = bucket.clone();
|
||||||
let object = object.clone();
|
let object = object.clone();
|
||||||
let version_id = version_id.clone();
|
let version_id = version_id.clone();
|
||||||
|
let slowtail_fault = slowtail_fault.clone();
|
||||||
join_set.spawn(async move {
|
join_set.spawn(async move {
|
||||||
let response_start = Instant::now();
|
let response_start = Instant::now();
|
||||||
let result = if let Some(disk) = disk {
|
let result = if let Some(disk) = disk {
|
||||||
@@ -2567,6 +2574,9 @@ impl SetDisks {
|
|||||||
Self::record_read_version_call(&object, index);
|
Self::record_read_version_call(&object, index);
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
Self::read_version_fanout_barrier(&object, index).await;
|
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)
|
disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts)
|
||||||
.await
|
.await
|
||||||
} else {
|
} else {
|
||||||
@@ -5744,6 +5754,130 @@ mod tests {
|
|||||||
(dirs, disks)
|
(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.
|
/// Demo / regression guard for the backlog#1325 per-disk call counters.
|
||||||
///
|
///
|
||||||
/// The metadata fan-out issues each `read_version` inside its own
|
/// The metadata fan-out issues each `read_version` inside its own
|
||||||
|
|||||||
@@ -174,15 +174,14 @@ use std::future::Future;
|
|||||||
use std::hash::{BuildHasher, Hash, Hasher};
|
use std::hash::{BuildHasher, Hash, Hasher};
|
||||||
use std::mem::{self};
|
use std::mem::{self};
|
||||||
use std::pin::Pin;
|
use std::pin::Pin;
|
||||||
use std::sync::OnceLock;
|
|
||||||
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||||
|
use std::sync::{Arc, OnceLock};
|
||||||
use std::task::{Context, Poll};
|
use std::task::{Context, Poll};
|
||||||
use std::time::{Instant, SystemTime, UNIX_EPOCH};
|
use std::time::{Instant, SystemTime, UNIX_EPOCH};
|
||||||
use std::{
|
use std::{
|
||||||
collections::{HashMap, HashSet},
|
collections::{HashMap, HashSet},
|
||||||
io::{Cursor, Write},
|
io::{Cursor, Write},
|
||||||
path::Path,
|
path::Path,
|
||||||
sync::Arc,
|
|
||||||
time::Duration,
|
time::Duration,
|
||||||
};
|
};
|
||||||
use time::OffsetDateTime;
|
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 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 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) ---
|
// --- Multipart Reader-Setup Prefetch Configuration (backlog#870) ---
|
||||||
|
|
||||||
const ENV_RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH: &str = "RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH";
|
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<String>,
|
||||||
|
object_prefix: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Debug)]
|
||||||
|
struct GetMetadataSlowtailFaultRequest {
|
||||||
|
delay: Duration,
|
||||||
|
disks: Arc<[usize]>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl GetMetadataSlowtailFaultRequest {
|
||||||
|
fn delay_for_disk(&self, disk_index: usize) -> Option<Duration> {
|
||||||
|
self.disks.contains(&disk_index).then_some(self.delay)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn parse_get_metadata_slowtail_fault_disks(raw: &str) -> Option<Vec<usize>> {
|
||||||
|
let mut disks = Vec::new();
|
||||||
|
for item in raw.split(',').map(str::trim).filter(|item| !item.is_empty()) {
|
||||||
|
let Ok(index) = item.parse::<usize>() else {
|
||||||
|
return None;
|
||||||
|
};
|
||||||
|
if !disks.contains(&index) {
|
||||||
|
disks.push(index);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
(!disks.is_empty()).then_some(disks)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn load_get_metadata_slowtail_fault_config() -> Option<GetMetadataSlowtailFaultConfig> {
|
||||||
|
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<GetMetadataSlowtailFaultRequest> {
|
||||||
|
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<Option<GetMetadataSlowtailFaultConfig>> = 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<Duration> {
|
||||||
|
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
|
/// Check if multipart reads prefetch the next part's bitrot reader setup
|
||||||
/// while the current part decodes (backlog#870).
|
/// while the current part decodes (backlog#870).
|
||||||
///
|
///
|
||||||
|
|||||||
Reference in New Issue
Block a user