diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 196d563ff..1fd1b02d5 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -39,8 +39,8 @@ use crate::diagnostics::get::{ GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND, GET_METADATA_RESPONSE_CORRUPT, GET_METADATA_RESPONSE_DISK_NOT_FOUND, GET_METADATA_RESPONSE_ERROR, GET_METADATA_RESPONSE_IGNORED, GET_METADATA_RESPONSE_NOT_FOUND, GET_METADATA_RESPONSE_TIMEOUT, GET_METADATA_RESPONSE_VALID, GET_METADATA_RESPONSE_VERSION_NOT_FOUND, GET_OBJECT_PATH_CODEC_STREAMING, - GET_OBJECT_PATH_DIRECT_MEMORY, GET_OBJECT_PATH_LEGACY_DUPLEX, GET_OBJECT_PATH_SET_DISK, GET_STAGE_DECODE, - GET_STAGE_METADATA_CACHE_LOOKUP, GET_STAGE_METADATA_RESOLVE, GET_STAGE_RANGE, GET_STAGE_READER_SETUP, + GET_OBJECT_PATH_DIRECT_MEMORY, GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX, GET_OBJECT_PATH_SET_DISK, + GET_STAGE_DECODE, GET_STAGE_METADATA_CACHE_LOOKUP, GET_STAGE_METADATA_RESOLVE, GET_STAGE_RANGE, GET_STAGE_READER_SETUP, GET_STAGE_READER_SETUP_DROP_PENDING, GET_STAGE_READER_SETUP_SCHEDULE, GET_STAGE_READER_SETUP_WAIT_QUORUM, GET_STAGE_READER_TASK_BITROT_READER_INIT, GET_STAGE_READER_TASK_FILE_OPEN, GET_STAGE_READER_TASK_READER_CONSTRUCTION, GetObjectFailureReason, classify_disk_error, get_stage_timer_if_enabled, record_get_object_pipeline_failure, @@ -70,6 +70,57 @@ use std::{ task::{Context, Poll}, time::{Duration, Instant}, }; + +fn metadata_metrics_path(bucket: &str) -> &'static str { + if crate::bucket::utils::is_meta_bucketname(bucket) { + GET_OBJECT_PATH_INTERNAL_META + } else { + GET_OBJECT_PATH_LEGACY_DUPLEX + } +} + +fn metadata_distribution_key(bucket: &str, object: &str) -> String { + [bucket, object].join("/") +} + +pub(in crate::set_disk) fn bounded_metadata_fanout_order( + bucket: &str, + object: &str, + total_disks: usize, + default_parity_count: usize, +) -> Vec { + let fallback_order = || (0..total_disks).collect::>(); + if default_parity_count == 0 || default_parity_count >= total_disks { + return fallback_order(); + } + + let data_blocks = total_disks - default_parity_count; + let distribution_key = metadata_distribution_key(bucket, object); + let distribution = FileInfo::new(&distribution_key, data_blocks, default_parity_count) + .erasure + .distribution; + if distribution.len() != total_disks { + return fallback_order(); + } + + let mut order = Vec::with_capacity(total_disks); + for block_index in 1..=data_blocks { + let Some(disk_index) = distribution + .iter() + .position(|distributed_block| *distributed_block == block_index) + else { + return fallback_order(); + }; + order.push(disk_index); + } + order.extend( + distribution + .iter() + .enumerate() + .filter_map(|(disk_index, block_index)| (*block_index > data_blocks).then_some(disk_index)), + ); + order +} use tokio::io::{AsyncRead, ReadBuf}; use tokio::sync::RwLock; use tokio::task::JoinSet; @@ -600,6 +651,83 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo && left.erasure.distribution == right.erasure.distribution } +pub(in crate::set_disk) async fn data_read_early_stop_inline_body_verified( + bucket: &str, + object: &str, + candidate: &FileInfo, + parts_metadata: &[FileInfo], + disks: &[Option], +) -> bool { + if !candidate.inline_data() + || candidate.is_compressed() + || candidate + .metadata + .keys() + .any(|key| rustfs_utils::http::is_object_encryption_marker(key)) + || candidate.is_remote() + || candidate.deleted + || candidate.size <= 0 + || candidate.parts.len() != 1 + || !candidate.has_valid_erasure_geometry() + { + return false; + } + + let Ok(object_size) = usize::try_from(candidate.size) else { + return false; + }; + if candidate.parts.first().is_none_or(|part| part.size != object_size) { + return false; + } + if !can_try_inline_data_shards_direct(object_size, candidate.erasure.block_size) { + return false; + } + + let Ok(erasure) = coding::Erasure::try_new_with_options( + candidate.erasure.data_blocks, + candidate.erasure.parity_blocks, + candidate.erasure.block_size, + candidate.uses_legacy_checksum, + ) else { + return false; + }; + let Some(data_files) = + collect_inline_data_shard_fileinfos_by_index(parts_metadata, candidate, erasure.data_shards, |index| { + disks.get(index).is_some_and(Option::is_some) + }) + else { + return false; + }; + + let Some(part) = candidate.parts.first() else { + return false; + }; + let checksum_info = candidate.erasure.get_checksum_info(part.number); + let checksum_algo = if candidate.uses_legacy_checksum && checksum_info.algorithm == HashAlgorithm::HighwayHash256S { + HashAlgorithm::HighwayHash256SLegacy + } else { + checksum_info.algorithm + }; + let read_length = inline_erasure_shard_file_offset( + 0, + object_size, + object_size, + candidate.erasure.block_size, + erasure.data_shards, + candidate.uses_legacy_checksum, + ); + let shard_size = inline_erasure_shard_size(candidate.erasure.block_size, erasure.data_shards, candidate.uses_legacy_checksum); + let Ok(mut readers) = + build_inline_bitrot_readers_from_refs(&data_files, bucket, object, read_length, shard_size, &checksum_algo, false).await + else { + return false; + }; + + try_read_inline_data_shards_direct(&mut readers, erasure.data_shards, read_length, object_size) + .await + .is_some_and(|body| body.len() == object_size) +} + pub(in crate::set_disk) fn classify_metadata_response_error(err: &DiskError) -> &'static str { match err { DiskError::FileNotFound | DiskError::VolumeNotFound => GET_METADATA_RESPONSE_NOT_FOUND, @@ -2187,11 +2315,12 @@ impl SetDisks { .await; } if early_stop_enabled { + let metrics_path = metadata_metrics_path(bucket); rustfs_io_metrics::record_get_object_metadata_early_stop_miss( - GET_OBJECT_PATH_LEGACY_DUPLEX, + metrics_path, GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, ); - rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0); + rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(metrics_path, 0); } Self::read_all_fileinfo_full_wait( @@ -2224,6 +2353,7 @@ impl SetDisks { let mut ress = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len()); let mut observations = observe.then(|| Vec::with_capacity(disks.len())); + let scheduled_count = disks.len(); let opts = ReadOptions { incl_free_versions, read_data, @@ -2289,6 +2419,14 @@ impl SetDisks { (Some(fanout_start), Some(observations)) => MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations), _ => MetadataFanoutDiagnostics::default(), }; + if observe { + rustfs_io_metrics::record_get_object_metadata_fanout_lifecycle( + metadata_metrics_path(bucket.as_ref()), + scheduled_count, + scheduled_count, + 0, + ); + } Ok((ress, errors, diagnostics)) } @@ -2319,9 +2457,17 @@ impl SetDisks { let bucket: Arc = Arc::from(bucket); let object: Arc = Arc::from(object); let version_id: Arc = Arc::from(version_id); + let metrics_path = metadata_metrics_path(bucket.as_ref()); let mut join_set = JoinSet::new(); let bounded_fanout = is_get_metadata_early_stop_bounded_fanout_enabled(); - let mut next_disk_index = 0usize; + let fanout_order = if bounded_fanout { + bounded_metadata_fanout_order(bucket.as_ref(), object.as_ref(), disks.len(), default_parity_count) + } else { + Vec::new() + }; + let mut next_fanout_index = 0usize; + let mut scheduled_count = 0usize; + let mut force_full_wait = false; let spawn_read_version = |join_set: &mut JoinSet<(usize, disk::error::Result, Duration)>, index: usize, disk: Option| { let task_opts = opts; @@ -2332,6 +2478,8 @@ impl SetDisks { join_set.spawn(async move { let response_start = Instant::now(); let result = if let Some(disk) = disk { + #[allow(clippy::let_unit_value)] + let _fanout_task_guard = Self::rename_fanout_task_guard(&object); Self::record_read_version_call(&object, index); #[cfg(test)] Self::read_version_fanout_barrier(&object, index).await; @@ -2346,15 +2494,18 @@ impl SetDisks { if bounded_fanout { let initial_target = accumulator.default_write_quorum().min(disks.len()); - while next_disk_index < initial_target { - if let Some(disk) = disks.get(next_disk_index).cloned() { - spawn_read_version(&mut join_set, next_disk_index, disk); + while next_fanout_index < initial_target { + let disk_index = fanout_order[next_fanout_index]; + if let Some(disk) = disks.get(disk_index).cloned() { + spawn_read_version(&mut join_set, disk_index, disk); + scheduled_count = scheduled_count.saturating_add(1); } - next_disk_index = next_disk_index.saturating_add(1); + next_fanout_index = next_fanout_index.saturating_add(1); } } else { for (index, disk) in disks.iter().cloned().enumerate() { spawn_read_version(&mut join_set, index, disk); + scheduled_count = scheduled_count.saturating_add(1); } } @@ -2383,46 +2534,87 @@ impl SetDisks { } } - if let Some(decision) = accumulator - .early_stop_decision() - .or_else(|| accumulator.version_early_stop_decision()) + if !force_full_wait + && let Some(decision) = accumulator + .early_stop_decision() + .or_else(|| accumulator.version_early_stop_decision()) { - let saved_responses = if bounded_fanout { - disks.len().saturating_sub(observations.len()) + let should_return_early = if read_data { + let allow_data_read_early_stop = match accumulator.candidate.as_ref() { + Some(candidate) => { + data_read_early_stop_inline_body_verified(bucket.as_ref(), object.as_ref(), candidate, &ress, disks) + .await + } + None => false, + }; + if !allow_data_read_early_stop { + force_full_wait = true; + } + allow_data_read_early_stop } else { - join_set.len() + true }; - join_set.abort_all(); - rustfs_io_metrics::record_get_object_metadata_early_stop_hit(GET_OBJECT_PATH_LEGACY_DUPLEX, decision.reason); - rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses( - GET_OBJECT_PATH_LEGACY_DUPLEX, - saved_responses, - ); - while join_set.join_next().await.is_some() {} - let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); - return Ok((ress, errors, diagnostics)); + + if should_return_early { + let saved_responses = if bounded_fanout { + disks.len().saturating_sub(observations.len()) + } else { + join_set.len() + }; + join_set.abort_all(); + rustfs_io_metrics::record_get_object_metadata_early_stop_hit(metrics_path, decision.reason); + rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(metrics_path, saved_responses); + let mut cancelled_count = 0usize; + while let Some(join_result) = join_set.join_next().await { + match join_result { + Err(join_error) if join_error.is_cancelled() => { + cancelled_count = cancelled_count.saturating_add(1); + } + _ => {} + } + } + rustfs_io_metrics::record_get_object_metadata_fanout_lifecycle( + metrics_path, + scheduled_count, + scheduled_count.saturating_sub(cancelled_count), + cancelled_count, + ); + let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); + return Ok((ress, errors, diagnostics)); + } } let pending_responses = join_set.len(); - let should_hedge_single_pending_data_read = - read_data && pending_responses == 1 && accumulator.can_still_reach_early_stop_with_pending(pending_responses); - if bounded_fanout - && next_disk_index < disks.len() + let should_hedge_single_pending_data_read = read_data + && !force_full_wait + && pending_responses == 1 + && accumulator.can_still_reach_early_stop_with_pending(pending_responses); + if bounded_fanout && force_full_wait { + while next_fanout_index < disks.len() { + let disk_index = fanout_order[next_fanout_index]; + if let Some(disk) = disks.get(disk_index).cloned() { + spawn_read_version(&mut join_set, disk_index, disk); + scheduled_count = scheduled_count.saturating_add(1); + } + next_fanout_index = next_fanout_index.saturating_add(1); + } + } else if bounded_fanout + && next_fanout_index < disks.len() && (!accumulator.can_still_reach_early_stop_with_pending(pending_responses) || should_hedge_single_pending_data_read) { - if let Some(disk) = disks.get(next_disk_index).cloned() { - spawn_read_version(&mut join_set, next_disk_index, disk); + let disk_index = fanout_order[next_fanout_index]; + if let Some(disk) = disks.get(disk_index).cloned() { + spawn_read_version(&mut join_set, disk_index, disk); + scheduled_count = scheduled_count.saturating_add(1); } - next_disk_index = next_disk_index.saturating_add(1); + next_fanout_index = next_fanout_index.saturating_add(1); } } - rustfs_io_metrics::record_get_object_metadata_early_stop_miss( - GET_OBJECT_PATH_LEGACY_DUPLEX, - accumulator.final_miss_reason(), - ); - rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0); + rustfs_io_metrics::record_get_object_metadata_early_stop_miss(metrics_path, accumulator.final_miss_reason()); + rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(metrics_path, 0); + rustfs_io_metrics::record_get_object_metadata_fanout_lifecycle(metrics_path, scheduled_count, scheduled_count, 0); let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); Ok((ress, errors, diagnostics)) } @@ -5081,7 +5273,7 @@ mod tests { use super::*; use std::io::Cursor; use tempfile::TempDir; - use tokio::io::AsyncReadExt; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; #[test] fn write_precondition_lookup_errors_fail_closed_unless_absence_is_known() { @@ -5411,39 +5603,598 @@ mod tests { } } + async fn inline_metadata_fanout_fileinfos_with_mode( + bucket: &str, + object: &str, + payload: &[u8], + uses_legacy_checksum: bool, + ) -> Vec { + let distribution_key = metadata_distribution_key(bucket, object); + let mut base = FileInfo::new(&distribution_key, 2, 2); + base.volume = bucket.to_string(); + base.name = object.to_string(); + base.size = i64::try_from(payload.len()).expect("test payload should fit i64"); + base.is_latest = true; + base.version_id = Some(Uuid::new_v4()); + base.data_dir = Some(Uuid::new_v4()); + base.mod_time = Some(OffsetDateTime::now_utc()); + base.metadata.insert("etag".to_string(), "etag-inline".to_string()); + base.add_object_part(1, "part-etag-inline".to_string(), payload.len(), base.mod_time, base.size, None, None); + base.set_inline_data(); + base.uses_legacy_checksum = uses_legacy_checksum; + + let erasure = coding::Erasure::new_with_options( + base.erasure.data_blocks, + base.erasure.parity_blocks, + base.erasure.block_size, + base.uses_legacy_checksum, + ); + let shards = erasure.encode_data(payload).expect("inline payload should encode"); + let checksum_algo = if base.uses_legacy_checksum { + HashAlgorithm::HighwayHash256SLegacy + } else { + HashAlgorithm::HighwayHash256S + }; + + let mut files = Vec::with_capacity(shards.len()); + for (index, shard) in shards.into_iter().enumerate() { + let mut writer = coding::BitrotWriterWrapper::new( + coding::CustomWriter::new_inline_buffer(), + erasure.shard_size(), + checksum_algo.clone(), + ); + writer.write(&shard).await.expect("inline shard should write"); + writer.shutdown().await.expect("inline writer should shutdown"); + let mut fi = base.clone(); + fi.erasure.index = index + 1; + fi.data = Some(Bytes::from( + writer + .into_inline_data() + .expect("inline bitrot writer should retain encoded data"), + )); + files.push(fi); + } + files + } + + async fn inline_metadata_fanout_fileinfos(bucket: &str, object: &str, payload: &[u8]) -> Vec { + inline_metadata_fanout_fileinfos_with_mode(bucket, object, payload, false).await + } + + async fn install_inline_metadata_fanout_fileinfo( + disks: &[Option], + bucket: &str, + object: &str, + payload: &[u8], + mutate: impl FnOnce(&mut [FileInfo]), + ) { + let mut files = inline_metadata_fanout_fileinfos(bucket, object, payload).await; + mutate(&mut files); + install_inline_metadata_fanout_files(disks, bucket, object, files).await; + } + + async fn install_inline_metadata_fanout_files(disks: &[Option], bucket: &str, object: &str, files: Vec) { + let distribution = files + .first() + .map(|file| file.erasure.distribution.clone()) + .expect("inline metadata fixture should include shards"); + for (disk_index, disk) in disks + .iter() + .enumerate() + .filter_map(|(disk_index, disk)| disk.as_ref().map(|disk| (disk_index, disk))) + { + let block_index = distribution + .get(disk_index) + .copied() + .expect("inline metadata fixture should cover every disk"); + let file_info = files + .get(block_index.checked_sub(1).expect("erasure block indexes are one-based")) + .expect("inline metadata fixture should include every distributed shard") + .clone(); + disk.write_metadata(bucket, bucket, object, file_info) + .await + .expect("inline metadata should be installed on every disk"); + } + } + + fn object_with_initial_data_shards(bucket: &str, prefix: &str, data_shards: usize, initial_fanout: usize) -> String { + (0..1000) + .map(|index| format!("{prefix}-{index}")) + .find(|name| { + let order = bounded_metadata_fanout_order(bucket, name, data_shards + 2, 2); + let distribution_key = metadata_distribution_key(bucket, name); + let distribution = FileInfo::new(&distribution_key, data_shards, 2).erasure.distribution; + order + .iter() + .take(initial_fanout) + .filter_map(|disk_index| distribution.get(*disk_index).copied()) + .filter(|block_index| (1..=data_shards).contains(block_index)) + .collect::>() + .len() + == data_shards + }) + .expect("test should find an object whose initial fanout covers every data shard") + } + + fn initial_data_shard_indexes(bucket: &str, object: &str, data_shards: usize, initial_fanout: usize) -> Vec { + let order = bounded_metadata_fanout_order(bucket, object, data_shards + 2, 2); + let distribution_key = metadata_distribution_key(bucket, object); + let distribution = FileInfo::new(&distribution_key, data_shards, 2).erasure.distribution; + order + .iter() + .take(initial_fanout) + .filter_map(|disk_index| distribution.get(*disk_index).copied()) + .filter(|block_index| (1..=data_shards).contains(block_index)) + .collect() + } + + fn bounded_spare_disk_index(bucket: &str, object: &str, data_shards: usize, parity_shards: usize) -> usize { + let total_disks = data_shards + parity_shards; + let initial_fanout = if data_shards == parity_shards { + data_shards + 1 + } else { + data_shards + }; + *bounded_metadata_fanout_order(bucket, object, total_disks, parity_shards) + .get(initial_fanout) + .expect("test geometry should leave one bounded spare disk") + } + + #[test] + fn bounded_metadata_fanout_order_prioritizes_default_data_shards() { + let bucket = "bounded-order-bucket"; + let object = "bounded-order-data-shards-first"; + let order = bounded_metadata_fanout_order(bucket, object, 16, 4); + let distribution_key = [bucket, object].join("/"); + let distribution = FileInfo::new(&distribution_key, 12, 4).erasure.distribution; + let initial_blocks: HashSet<_> = order + .iter() + .take(12) + .filter_map(|disk_index| distribution.get(*disk_index).copied()) + .collect(); + + assert_eq!(order.len(), 16); + assert_eq!(order.iter().copied().collect::>().len(), 16); + assert_eq!(initial_blocks, (1..=12).collect()); + } + + #[test] + fn bounded_metadata_fanout_order_uses_written_bucket_object_key() { + let bucket = "bounded-order-rotation-bucket"; + let (object, bare_distribution, stored_distribution) = (0..1000) + .map(|index| format!("bounded-order-rotation-object-{index}")) + .find_map(|object| { + let bare_distribution = FileInfo::new(&object, 12, 4).erasure.distribution; + let stored_key = [bucket, object.as_str()].join("/"); + let stored_distribution = FileInfo::new(&stored_key, 12, 4).erasure.distribution; + (bare_distribution != stored_distribution).then_some((object, bare_distribution, stored_distribution)) + }) + .expect("test should find a bucket/object pair with a different distribution rotation"); + + let order = bounded_metadata_fanout_order(bucket, &object, 16, 4); + let initial_stored_blocks: HashSet<_> = order + .iter() + .take(12) + .filter_map(|disk_index| stored_distribution.get(*disk_index).copied()) + .collect(); + let initial_bare_blocks: HashSet<_> = order + .iter() + .take(12) + .filter_map(|disk_index| bare_distribution.get(*disk_index).copied()) + .collect(); + + assert_eq!(initial_stored_blocks, (1..=12).collect()); + assert_ne!( + initial_bare_blocks, + (1..=12).collect(), + "test fixture must prove the bare object distribution would pick the wrong initial data-shard set" + ); + } + #[tokio::test] - async fn bounded_metadata_early_stop_ab_hedges_data_get_read_version_fanout() { - const DISKS: usize = 4; - let bucket = "bounded-data-get-fanout-bucket"; - let control_object = "bounded-data-get-control-object"; - let treatment_object = "bounded-data-get-treatment-object"; + async fn bounded_metadata_early_stop_non_inline_fallback_schedules_beyond_initial_quorum() { + const DISKS: usize = 6; + let bucket = "bounded-data-get-six-disk-fanout-bucket"; + let object = "bounded-data-get-six-disk-object"; let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; - install_metadata_fanout_fileinfo(&disks, bucket, control_object, None).await; - install_metadata_fanout_fileinfo(&disks, bucket, treatment_object, None).await; + install_metadata_fanout_fileinfo(&disks, bucket, object, None).await; temp_env::async_with_vars( [ ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), - ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("false")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), ], async { - let calls = disk_call_counters::observe(control_object); - let (_, _, diagnostics) = - SetDisks::read_all_fileinfo_observed(&disks, bucket, bucket, control_object, "", true, false, false, true, 2) + let calls = disk_call_counters::observe(object); + let (parts_metadata, errs, diagnostics) = + SetDisks::read_all_fileinfo_observed(&disks, bucket, bucket, object, "", true, false, false, true, 3) .await - .expect("control metadata should resolve"); + .expect("non-inline metadata should resolve after force-full-wait fallback"); assert_eq!( calls.total(disk_call_counters::KIND_READ_VERSION), DISKS as u64, - "control path should keep full fanout when data-read early stop is explicitly disabled" + "force-full-wait fallback must schedule disks beyond the initial write quorum" ); assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); }, ) .await; + drop(dirs); + } + + #[tokio::test] + async fn bounded_metadata_early_stop_allows_verified_inline_data_get_quorum() { + const DISKS: usize = 4; + let bucket = "bounded-inline-data-get-fanout-bucket"; + let object = object_with_initial_data_shards(bucket, "bounded-inline-data-get-object", 2, 3); + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_inline_metadata_fanout_fileinfo(&disks, bucket, &object, b"verified inline payload", |_| {}).await; + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let spare_disk = bounded_spare_disk_index(bucket, &object, 2, 2); + let barrier = rename_fanout_barrier::arm(&object, spare_disk, rename_fanout_barrier::PHASE_READ_VERSION); + let tracker = rename_fanout_barrier::observe_tasks(&object); + let calls = disk_call_counters::observe(&object); + let disks_for_read = disks.clone(); + let object_for_read = object.clone(); + let mut read = tokio::spawn(async move { + SetDisks::read_all_fileinfo_observed( + &disks_for_read, + bucket, + bucket, + &object_for_read, + "", + true, + false, + false, + true, + 2, + ) + .await + }); + + tokio::time::timeout(BARRIER_PAUSE_GUARD, barrier.wait_until_paused()) + .await + .expect("bounded fanout should hedge and pause the spare disk"); + let (parts_metadata, errs, diagnostics) = tokio::time::timeout(BARRIER_PAUSE_GUARD, &mut read) + .await + .expect("verified inline quorum should return without the paused spare") + .expect("read task should join") + .expect("verified inline metadata should reach early-stop quorum"); + + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + DISKS as u64, + "bounded fanout may schedule one spare before the verified inline quorum returns" + ); + assert_eq!( + tracker.running(), + 0, + "early-stop should drain spawned read_version tasks before returning" + ); + assert_eq!(diagnostics.total_responses(), 3); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), 3); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + + drop(dirs); + } + + #[tokio::test] + async fn data_read_early_stop_verifies_legacy_inline_checksum_payload() { + let bucket = "legacy-inline-data-get-fanout-bucket"; + let object = "legacy-inline-data-get-object"; + let payload = b"legacy inline payload whose size is not divisible by the data shard count"; + let (_dirs, disks) = call_counter_local_disks(bucket, 4).await; + let files = inline_metadata_fanout_fileinfos_with_mode(bucket, object, payload, true).await; + let distribution = files + .first() + .map(|file| file.erasure.distribution.clone()) + .expect("legacy fixture should include metadata"); + let order = bounded_metadata_fanout_order(bucket, object, 4, 2); + let mut parts_metadata = vec![FileInfo::default(); 4]; + for disk_index in order.into_iter().take(3) { + let block_index = distribution + .get(disk_index) + .copied() + .expect("legacy fixture distribution should cover every disk"); + parts_metadata[disk_index] = files + .get(block_index.checked_sub(1).expect("erasure block indexes are one-based")) + .expect("legacy fixture should include every distributed shard") + .clone(); + } + let candidate = parts_metadata + .iter() + .find(|file| file.name == object) + .expect("legacy fixture should include observed metadata") + .clone(); + + assert!( + data_read_early_stop_inline_body_verified(bucket, object, &candidate, &parts_metadata, &disks).await, + "legacy inline metadata must use the legacy bitrot shard sizing and checksum algorithm" + ); + } + + #[test] + #[serial_test::serial] + fn metadata_fanout_lifecycle_records_real_early_stop_abort() { + assert_metadata_fanout_lifecycle_records_real_early_stop_abort( + "lifecycle-inline-data-get-fanout-bucket", + "lifecycle-inline-data-get-object", + GET_OBJECT_PATH_LEGACY_DUPLEX, + None, + ); + } + + #[test] + #[serial_test::serial] + fn metadata_fanout_lifecycle_records_bounded_early_stop_abort() { + assert_metadata_fanout_lifecycle_records_real_early_stop_abort( + "lifecycle-bounded-inline-data-get-fanout-bucket", + "lifecycle-bounded-inline-data-get-object", + GET_OBJECT_PATH_LEGACY_DUPLEX, + Some("true"), + ); + } + + #[test] + #[serial_test::serial] + fn metadata_fanout_lifecycle_records_internal_meta_early_stop_abort_path() { + assert_metadata_fanout_lifecycle_records_real_early_stop_abort( + RUSTFS_META_BUCKET, + "buckets/.usage-cache/lifecycle-inline-data-get-object", + GET_OBJECT_PATH_INTERNAL_META, + None, + ); + } + + #[test] + #[serial_test::serial] + fn metadata_fanout_records_internal_meta_final_miss_path() { + const DISKS: usize = 4; + let bucket = RUSTFS_META_BUCKET; + let object = object_with_initial_data_shards(bucket, "buckets/.usage-cache/final-miss-inline-data-get-object", 2, 3); + let corrupt_shard = initial_data_shard_indexes(bucket, &object, 2, 3)[0]; + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("test runtime should build"); + let recorder = crate::test_metrics::CapturingRecorder::default(); + let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled(); + rustfs_io_metrics::set_get_stage_metrics_enabled(true); + + metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_inline_metadata_fanout_fileinfo(&disks, bucket, &object, b"verified inline payload", |files| { + if let Some(data) = files.get_mut(corrupt_shard - 1).and_then(|file| file.data.as_mut()) { + let mut corrupt = data.to_vec(); + corrupt[0] ^= 0x01; + *data = Bytes::from(corrupt); + } + }) + .await; + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let (parts_metadata, errs, diagnostics) = SetDisks::read_all_fileinfo_observed( + &disks, bucket, bucket, &object, "", true, false, false, true, 2, + ) + .await + .expect("corrupt internal inline metadata should fall back to full fanout"); + + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + drop(dirs); + }); + }); + rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate); + + assert_eq!( + recorder.counter_value( + "rustfs_io_get_object_metadata_early_stop_total", + &[ + ("path", GET_OBJECT_PATH_INTERNAL_META), + ("decision", "miss"), + ("reason", GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM), + ], + ), + 1, + "internal metadata final early-stop miss must retain its path label" + ); + assert_eq!( + recorder.counter_value( + "rustfs_io_get_object_metadata_early_stop_total", + &[ + ("path", GET_OBJECT_PATH_LEGACY_DUPLEX), + ("decision", "miss"), + ("reason", GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM), + ], + ), + 0, + "internal metadata final early-stop miss must not leak into legacy_duplex" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_early_stop_saved_responses", + &[("path", GET_OBJECT_PATH_INTERNAL_META)] + ), + vec![0.0], + "internal metadata final miss must record zero saved responses on internal_meta" + ); + assert!( + recorder + .histogram_values( + "rustfs_io_get_object_metadata_early_stop_saved_responses", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ) + .is_empty(), + "internal metadata final miss saved responses must not leak into legacy_duplex" + ); + } + + fn assert_metadata_fanout_lifecycle_records_real_early_stop_abort( + bucket: &'static str, + object_prefix: &str, + expected_path: &'static str, + bounded_fanout_env: Option<&'static str>, + ) { + const DISKS: usize = 4; + let object = object_with_initial_data_shards(bucket, object_prefix, 2, 3); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("test runtime should build"); + let recorder = crate::test_metrics::CapturingRecorder::default(); + let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled(); + rustfs_io_metrics::set_get_stage_metrics_enabled(true); + + metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_inline_metadata_fanout_fileinfo(&disks, bucket, &object, b"verified inline payload", |_| {}).await; + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", bounded_fanout_env), + ], + async { + let barrier_disk = if bounded_fanout_env.is_some() { + bounded_spare_disk_index(bucket, &object, 2, 2) + } else { + 3 + }; + let barrier = + rename_fanout_barrier::arm(&object, barrier_disk, rename_fanout_barrier::PHASE_READ_VERSION); + let disks_for_read = disks.clone(); + let object_for_read = object.clone(); + let mut read = tokio::spawn(async move { + SetDisks::read_all_fileinfo_observed( + &disks_for_read, + bucket, + bucket, + &object_for_read, + "", + true, + false, + false, + true, + 2, + ) + .await + }); + + tokio::time::timeout(BARRIER_PAUSE_GUARD, barrier.wait_until_paused()) + .await + .expect("metadata fanout should pause the spare metadata task"); + let (parts_metadata, errs, diagnostics) = tokio::time::timeout(BARRIER_PAUSE_GUARD, &mut read) + .await + .expect("verified inline quorum should abort the paused eager fanout") + .expect("read task should join") + .expect("verified inline metadata should reach early-stop quorum"); + + assert_eq!(diagnostics.total_responses(), 3); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), 3); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + drop(dirs); + }); + }); + rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate); + + assert_eq!( + recorder.histogram_values("rustfs_io_get_object_metadata_fanout_scheduled", &[("path", expected_path)]), + vec![4.0] + ); + assert_eq!( + recorder.histogram_values("rustfs_io_get_object_metadata_fanout_completed", &[("path", expected_path)]), + vec![3.0] + ); + assert_eq!( + recorder.histogram_values("rustfs_io_get_object_metadata_fanout_cancelled", &[("path", expected_path)]), + vec![1.0] + ); + assert_eq!( + recorder.counter_value( + "rustfs_io_get_object_metadata_early_stop_total", + &[ + ("path", expected_path), + ("decision", "hit"), + ("reason", GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM), + ], + ), + 1, + "early-stop hit must retain its path label" + ); + let unexpected_path = if expected_path == GET_OBJECT_PATH_INTERNAL_META { + GET_OBJECT_PATH_LEGACY_DUPLEX + } else { + GET_OBJECT_PATH_INTERNAL_META + }; + assert_eq!( + recorder.counter_value( + "rustfs_io_get_object_metadata_early_stop_total", + &[ + ("path", unexpected_path), + ("decision", "hit"), + ("reason", GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM), + ], + ), + 0, + "early-stop hit must not leak into the other metadata path" + ); + assert_eq!( + recorder.histogram_values("rustfs_io_get_object_metadata_early_stop_saved_responses", &[("path", expected_path)]), + vec![1.0] + ); + assert!( + recorder + .histogram_values("rustfs_io_get_object_metadata_early_stop_saved_responses", &[("path", unexpected_path)]) + .is_empty(), + "early-stop saved responses must not leak into the other metadata path" + ); + } + + #[tokio::test] + async fn bounded_metadata_early_stop_full_waits_when_inline_data_is_corrupt() { + const DISKS: usize = 4; + let bucket = "bounded-inline-data-get-corrupt-bucket"; + let object = object_with_initial_data_shards(bucket, "bounded-inline-data-get-corrupt-object", 2, 3); + let corrupt_shard = initial_data_shard_indexes(bucket, &object, 2, 3)[0]; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_inline_metadata_fanout_fileinfo(&disks, bucket, &object, b"verified inline payload", |files| { + if let Some(data) = files.get_mut(corrupt_shard - 1).and_then(|file| file.data.as_mut()) { + let mut corrupt = data.to_vec(); + corrupt[0] ^= 0x01; + *data = Bytes::from(corrupt); + } + }) + .await; + temp_env::async_with_vars( [ ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), @@ -5451,31 +6202,359 @@ mod tests { ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), ], async { - let calls = disk_call_counters::observe(treatment_object); - let (parts_metadata, errs, diagnostics) = SetDisks::read_all_fileinfo_observed( - &disks, - bucket, - bucket, - treatment_object, - "", - true, - false, - false, - true, - 2, - ) - .await - .expect("healthy object metadata should reach early-stop quorum"); + let calls = disk_call_counters::observe(&object); + let (parts_metadata, errs, diagnostics) = + SetDisks::read_all_fileinfo_observed(&disks, bucket, bucket, &object, "", true, false, false, true, 2) + .await + .expect("corrupt inline metadata should fall back to full fanout"); - assert!( - (3..=DISKS as u64).contains(&calls.total(disk_call_counters::KIND_READ_VERSION)), - "healthy 2+2 bounded data-read fanout may finish at quorum before a spare hedge is needed" + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + DISKS as u64, + "inline bitrot failure must keep metadata fanout open instead of aborting spares" ); - assert!( - (3..=DISKS).contains(&diagnostics.total_responses()), - "treatment path should return after reaching quorum, with at most the spare hedge response observed" + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + + drop(dirs); + } + + #[tokio::test] + async fn bounded_metadata_early_stop_full_waits_when_inline_generation_differs() { + const DISKS: usize = 4; + let bucket = "bounded-inline-data-get-generation-bucket"; + let object = object_with_initial_data_shards(bucket, "bounded-inline-data-get-generation-object", 2, 3); + let stale_shard = initial_data_shard_indexes(bucket, &object, 2, 3)[0]; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_inline_metadata_fanout_fileinfo(&disks, bucket, &object, b"verified inline payload", |files| { + if let Some(file) = files.get_mut(stale_shard - 1) { + file.version_id = Some(Uuid::new_v4()); + } + }) + .await; + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let calls = disk_call_counters::observe(&object); + let (parts_metadata, errs, diagnostics) = + SetDisks::read_all_fileinfo_observed(&disks, bucket, bucket, &object, "", true, false, false, true, 2) + .await + .expect("mixed-generation inline metadata should fall back to full fanout"); + + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + DISKS as u64, + "inline data from a different metadata generation must not satisfy the data-read gate" ); - assert!(parts_metadata.iter().filter(|fi| fi.name == treatment_object).count() >= 3); + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + + drop(dirs); + } + + #[tokio::test] + async fn bounded_metadata_early_stop_full_waits_when_inline_shard_identity_is_copied() { + const DISKS: usize = 4; + let bucket = "bounded-inline-data-get-copied-shard-bucket"; + let object = object_with_initial_data_shards(bucket, "bounded-inline-data-get-copied-shard-object", 2, 3); + let data_shards = initial_data_shard_indexes(bucket, &object, 2, 3); + let source_shard = data_shards[0]; + let target_shard = data_shards[1]; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_inline_metadata_fanout_fileinfo(&disks, bucket, &object, b"verified inline payload", |files| { + let copied = files[source_shard - 1].clone(); + files[target_shard - 1].erasure.index = copied.erasure.index; + files[target_shard - 1].data = copied.data; + }) + .await; + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let spare_disk = bounded_spare_disk_index(bucket, &object, 2, 2); + let barrier = rename_fanout_barrier::arm(&object, spare_disk, rename_fanout_barrier::PHASE_READ_VERSION); + let calls = disk_call_counters::observe(&object); + let disks_for_read = disks.clone(); + let object_for_read = object.clone(); + let mut read = tokio::spawn(async move { + SetDisks::read_all_fileinfo_observed( + &disks_for_read, + bucket, + bucket, + &object_for_read, + "", + true, + false, + false, + true, + 2, + ) + .await + }); + + tokio::time::timeout(BARRIER_PAUSE_GUARD, barrier.wait_until_paused()) + .await + .expect("full-wait fallback should schedule the spare disk"); + assert!( + tokio::time::timeout(BARRIER_PAUSE_GUARD, &mut read).await.is_err(), + "copied inline shard identity must not satisfy the early-stop data gate" + ); + barrier.release(); + + let (parts_metadata, errs, diagnostics) = read + .await + .expect("metadata read task should not panic") + .expect("copied inline shard identity should fall back to full fanout"); + + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + DISKS as u64, + "copied shard identity must keep metadata fanout open instead of aborting spares" + ); + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + + drop(dirs); + } + + #[tokio::test] + async fn bounded_metadata_early_stop_full_waits_for_transformed_inline_metadata() { + const DISKS: usize = 4; + let bucket = "bounded-inline-data-get-transform-bucket"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + + async fn assert_full_wait(disks: &[Option], bucket: &str, object: &str, mutate: impl FnOnce(&mut [FileInfo])) { + const DISKS: usize = 4; + install_inline_metadata_fanout_fileinfo(disks, bucket, object, b"verified inline payload", mutate).await; + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let calls = disk_call_counters::observe(object); + let (parts_metadata, errs, diagnostics) = + SetDisks::read_all_fileinfo_observed(disks, bucket, bucket, object, "", true, false, false, true, 2) + .await + .expect("transformed inline metadata should fall back to full fanout"); + + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + DISKS as u64, + "transformed inline metadata must not pass the plaintext data-read gate" + ); + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + } + + assert_full_wait(&disks, bucket, "bounded-inline-data-get-compressed-object", |files| { + for file in files { + rustfs_utils::http::insert_str(&mut file.metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string()); + } + }) + .await; + + assert_full_wait(&disks, bucket, "bounded-inline-data-get-encrypted-object", |files| { + for file in files { + file.metadata + .insert(rustfs_utils::http::headers::SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string()); + } + }) + .await; + + drop(dirs); + } + + #[tokio::test] + async fn bounded_metadata_early_stop_full_waits_for_unsafe_inline_metadata_shapes() { + const DISKS: usize = 4; + let bucket = "bounded-inline-data-get-unsafe-shape-bucket"; + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + + async fn assert_full_wait(disks: &[Option], bucket: &str, object: &str, mutate: impl FnOnce(&mut [FileInfo])) { + const DISKS: usize = 4; + install_inline_metadata_fanout_fileinfo(disks, bucket, object, b"verified inline payload", mutate).await; + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let calls = disk_call_counters::observe(object); + let (parts_metadata, errs, diagnostics) = + SetDisks::read_all_fileinfo_observed(disks, bucket, bucket, object, "", true, false, false, true, 2) + .await + .expect("unsafe inline metadata shape should fall back to full fanout"); + + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + DISKS as u64, + "unsafe inline metadata shape {object} must not abort remaining metadata responses" + ); + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); + assert!(errs.iter().all(Option::is_none)); + }, + ) + .await; + } + + assert_full_wait(&disks, bucket, "bounded-inline-data-get-remote-object", |files| { + for file in files { + file.transition_status = TRANSITION_COMPLETE.to_string(); + } + }) + .await; + + assert_full_wait(&disks, bucket, "bounded-inline-data-get-zero-size-object", |files| { + for file in files { + file.size = 0; + if let Some(part) = file.parts.first_mut() { + part.size = 0; + } + } + }) + .await; + + assert_full_wait(&disks, bucket, "bounded-inline-data-get-multipart-object", |files| { + for file in files { + let mut second_part = file.parts[0].clone(); + second_part.number = 2; + file.parts.push(second_part); + } + }) + .await; + + assert_full_wait(&disks, bucket, "bounded-inline-data-get-part-size-mismatch-object", |files| { + for file in files { + if let Some(part) = file.parts.first_mut() { + part.size = part.size.saturating_add(1); + } + } + }) + .await; + + assert_full_wait(&disks, bucket, "bounded-inline-data-get-oversize-object", |files| { + for file in files { + let oversize = file.erasure.block_size.saturating_add(1); + file.size = i64::try_from(oversize).expect("test block size should fit i64"); + if let Some(part) = file.parts.first_mut() { + part.size = oversize; + } + } + }) + .await; + + drop(dirs); + } + + #[tokio::test] + async fn bounded_metadata_early_stop_full_waits_for_purge_pending_inline_payload() { + const DISKS: usize = 4; + let bucket = "bounded-inline-data-get-purge-pending-bucket"; + let object = object_with_initial_data_shards(bucket, "bounded-inline-data-get-purge-pending-object", 2, 3); + let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await; + install_inline_metadata_fanout_fileinfo(&disks, bucket, &object, b"verified inline payload", |files| { + for file in files { + rustfs_utils::http::insert_str( + &mut file.metadata, + rustfs_utils::http::SUFFIX_PURGESTATUS, + "target=PENDING;".to_string(), + ); + let replication_state = crate::bucket::replication::ReplicationState { + version_purge_status_internal: Some("target=PENDING;".to_string()), + purge_targets: crate::bucket::replication::version_purge_statuses_map("target=PENDING;"), + ..Default::default() + }; + file.replication_state_internal = + Some(crate::bucket::replication::replication_state_to_filemeta(&replication_state)); + file.deleted = true; + assert!( + !file.is_canonical_delete_marker(), + "test fixture must remain an erasure-backed purge-pending payload" + ); + } + }) + .await; + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let spare_disk = bounded_spare_disk_index(bucket, &object, 2, 2); + let barrier = rename_fanout_barrier::arm(&object, spare_disk, rename_fanout_barrier::PHASE_READ_VERSION); + let calls = disk_call_counters::observe(&object); + let disks_for_read = disks.clone(); + let object_for_read = object.clone(); + let mut read = tokio::spawn(async move { + SetDisks::read_all_fileinfo_observed( + &disks_for_read, + bucket, + bucket, + &object_for_read, + "", + true, + false, + false, + true, + 2, + ) + .await + }); + + tokio::time::timeout(BARRIER_PAUSE_GUARD, barrier.wait_until_paused()) + .await + .expect("purge-pending fallback should schedule the spare disk"); + assert!( + tokio::time::timeout(BARRIER_PAUSE_GUARD, &mut read).await.is_err(), + "purge-pending payload metadata must not satisfy the inline data-read gate" + ); + barrier.release(); + + let (parts_metadata, errs, diagnostics) = read + .await + .expect("metadata read task should not panic") + .expect("purge-pending inline payload should fall back to full fanout"); + + assert_eq!( + calls.total(disk_call_counters::KIND_READ_VERSION), + DISKS as u64, + "purge-pending payload must keep metadata fanout open instead of aborting spares" + ); + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); assert!(errs.iter().all(Option::is_none)); }, ) @@ -5485,7 +6564,7 @@ mod tests { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] - async fn bounded_data_get_hedges_single_pending_read_version() { + async fn bounded_non_inline_data_get_hedges_then_waits_for_full_fanout() { const DISKS: usize = 4; let bucket = "bounded-data-get-hedge-bucket"; let object = "bounded-data-get-hedge-object"; @@ -5518,12 +6597,14 @@ mod tests { .await .expect("bounded data-read fanout should hedge by starting the spare disk"); - let completed = tokio::time::timeout(BARRIER_PAUSE_GUARD, &mut read).await; - if completed.is_err() { - barrier.release(); - } - let (parts_metadata, errs, diagnostics) = completed - .expect("spare metadata should allow early-stop without waiting for the paused disk") + let pending = tokio::time::timeout(BARRIER_PAUSE_GUARD, &mut read).await; + assert!( + pending.is_err(), + "non-inline data reads must not return before the paused metadata response" + ); + barrier.release(); + let (parts_metadata, errs, diagnostics) = read + .await .expect("metadata read task should not panic") .expect("healthy spare metadata should resolve"); @@ -5532,8 +6613,8 @@ mod tests { DISKS as u64, "bounded data-read fanout should issue the paused disk plus one spare hedge" ); - assert_eq!(diagnostics.total_responses(), 3); - assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), 3); + assert_eq!(diagnostics.total_responses(), DISKS); + assert_eq!(parts_metadata.iter().filter(|fi| fi.name == object).count(), DISKS); assert!(errs.iter().all(Option::is_none)); }, ) diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 20e481c2d..35cec2f28 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -880,12 +880,44 @@ mod prepared_get_object_metadata_tests { use super::*; use crate::ecstore_validation_blackbox::make_local_set_disks; use crate::object_api::{BLOCK_SIZE_V2, PutObjReader}; - use crate::set_disk::core::io_primitives::disk_call_counters; + use crate::set_disk::core::io_primitives::{bounded_metadata_fanout_order, disk_call_counters, rename_fanout_barrier}; use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; + use crate::test_metrics::CapturingRecorder; use http::HeaderMap; use tokio::io::AsyncReadExt; + const READ_VERSION_BARRIER_GUARD: std::time::Duration = std::time::Duration::from_secs(10); + + fn object_with_initial_data_shards(bucket: &str, prefix: &str) -> String { + (0..1000) + .map(|index| format!("{prefix}-{index}.bin")) + .find(|name| { + let order = bounded_metadata_fanout_order(bucket, name, 4, 2); + let distribution = FileInfo::new(&[bucket, name].join("/"), 2, 2).erasure.distribution; + let mut seen = [false; 2]; + for disk_index in order.into_iter().take(3) { + if let Some(block_index @ 1..=2) = distribution.get(disk_index).copied() { + seen[block_index - 1] = true; + } + } + seen.into_iter().all(|seen| seen) + }) + .expect("test should find an object whose initial fanout covers both data shards") + } + + fn bounded_spare_disk_index(bucket: &str, object: &str) -> usize { + *bounded_metadata_fanout_order(bucket, object, 4, 2) + .get(3) + .expect("4-disk test geometry should leave one bounded spare disk") + } + + fn bounded_slow_initial_disk_index(bucket: &str, object: &str) -> usize { + *bounded_metadata_fanout_order(bucket, object, 4, 2) + .get(2) + .expect("4-disk test geometry should include a third initial metadata disk") + } + #[tokio::test] async fn prepared_metadata_is_consumed_exactly_once() { let snapshot = GetObjectFileInfo::owned(FileInfo::default(), Vec::new(), Vec::new()); @@ -1002,6 +1034,307 @@ mod prepared_get_object_metadata_tests { ); } + #[test] + #[serial_test::serial(body_cache_hook)] + fn inline_data_read_early_stop_reader_returns_exact_body() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current-thread runtime should build"); + let bucket = "inline-data-read-early-stop-reader"; + let object = object_with_initial_data_shards(bucket, "inline-data-read-early-stop-reader-object"); + let payload = b"inline early-stop reader payload".repeat(256); + let recorder = CapturingRecorder::default(); + let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled(); + rustfs_io_metrics::set_get_stage_metrics_enabled(true); + + let (restored, object_size, calls_total) = metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut put_reader = PutObjReader::from_vec(payload.clone()); + set_disks + .put_object(bucket, &object, &mut put_reader, &opts) + .await + .expect("inline object should be written"); + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let slow_initial_disk = bounded_slow_initial_disk_index(bucket, &object); + let barrier = + rename_fanout_barrier::arm(&object, slow_initial_disk, rename_fanout_barrier::PHASE_READ_VERSION); + let calls = disk_call_counters::observe(&object); + let set_disks_for_read = Arc::clone(&set_disks); + let opts_for_read = opts.clone(); + let object_for_read = object.clone(); + let mut open_reader = tokio::spawn(async move { + set_disks_for_read + .get_object_reader(bucket, &object_for_read, None, HeaderMap::new(), &opts_for_read) + .await + }); + + tokio::time::timeout(READ_VERSION_BARRIER_GUARD, barrier.wait_until_paused()) + .await + .expect("bounded inline GET should pause a slow initial metadata read"); + let mut reader = tokio::time::timeout(READ_VERSION_BARRIER_GUARD, &mut open_reader) + .await + .expect("production inline GET should return before the paused metadata response") + .expect("inline GET reader task should not panic") + .expect("inline GET reader should open"); + let object_size = reader.object_info.size; + let mut restored = Vec::new(); + reader + .stream + .read_to_end(&mut restored) + .await + .expect("inline GET body should stream"); + + (restored, object_size, calls.total(disk_call_counters::KIND_READ_VERSION)) + }, + ) + .await + }) + }); + rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate); + + assert_eq!(object_size, payload.len() as i64); + assert_eq!(restored, payload); + assert_eq!(calls_total, 4, "bounded production GET should schedule the initial quorum plus one spare"); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_scheduled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![4.0], + "bounded production GET should record all scheduled metadata tasks" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_completed", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![3.0], + "bounded production GET should record only observed metadata responses as completed" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_cancelled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![1.0], + "bounded production GET should record the aborted slow metadata task" + ); + } + + #[test] + #[serial_test::serial(body_cache_hook)] + fn prepared_metadata_uses_full_fanout_even_when_data_read_early_stop_is_enabled() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current-thread runtime should build"); + let bucket = "prepared-metadata-early-stop-enabled"; + let object = object_with_initial_data_shards(bucket, "prepared-metadata-early-stop-enabled-object"); + let payload = b"prepared metadata early-stop enabled payload".repeat(16); + let recorder = CapturingRecorder::default(); + let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled(); + rustfs_io_metrics::set_get_stage_metrics_enabled(true); + + let (restored, calls_total) = metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut put_reader = PutObjReader::from_vec(payload.clone()); + set_disks + .put_object(bucket, &object, &mut put_reader, &opts) + .await + .expect("object should be written"); + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let calls = disk_call_counters::observe(&object); + let metadata = set_disks + .prepare_get_object_metadata(bucket, &object, &opts) + .await + .expect("prepared metadata should resolve"); + let calls_total = calls.total(disk_call_counters::KIND_READ_VERSION); + + let mut reader = set_disks + .get_object_reader_with_prepared_metadata(bucket, &object, None, HeaderMap::new(), &opts, metadata) + .await + .expect("prepared body reader should open"); + let mut restored = Vec::new(); + reader + .stream + .read_to_end(&mut restored) + .await + .expect("prepared body should stream"); + (restored, calls_total) + }, + ) + .await + }) + }); + rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate); + + assert_eq!(restored, payload); + assert_eq!( + calls_total, 4, + "prepared metadata must opt out of data-read early-stop until the read shape is known" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_scheduled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![4.0], + "prepared metadata should schedule the full metadata fanout" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_completed", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![4.0], + "prepared metadata must wait for every scheduled metadata response" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_cancelled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![0.0], + "prepared metadata must not cancel metadata responses" + ); + } + + #[test] + #[serial_test::serial(body_cache_hook)] + fn data_read_early_stop_request_shapes_full_wait_in_production_reader() { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current-thread runtime should build"); + let bucket = "data-read-early-stop-shape-reader"; + let payload = b"shape-gated inline reader payload".repeat(256); + + for (object_prefix, range, configure_opts, expected_body) in [ + ( + "data-read-early-stop-range-reader-object", + Some(HTTPRangeSpec { + start: 0, + end: 3, + is_suffix_length: false, + }), + None, + payload[..4].to_vec(), + ), + ("data-read-early-stop-part-reader-object", None, Some(1), payload.clone()), + ] { + let recorder = CapturingRecorder::default(); + let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled(); + rustfs_io_metrics::set_get_stage_metrics_enabled(true); + let (restored, calls_total) = metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let (_dirs, set_disks) = make_local_set_disks(4, 2).await; + let object = object_with_initial_data_shards(bucket, object_prefix); + let mut opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + opts.part_number = configure_opts; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + let mut put_reader = PutObjReader::from_vec(payload.clone()); + set_disks + .put_object(bucket, &object, &mut put_reader, &opts) + .await + .expect("inline object should be written"); + + temp_env::async_with_vars( + [ + ("RUSTFS_GET_METADATA_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE", Some("true")), + ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), + ], + async { + let calls = disk_call_counters::observe(&object); + let mut reader = set_disks + .get_object_reader(bucket, &object, range, HeaderMap::new(), &opts) + .await + .expect("shape-gated GET reader should open"); + let mut restored = Vec::new(); + reader + .stream + .read_to_end(&mut restored) + .await + .expect("shape-gated GET body should stream"); + (restored, calls.total(disk_call_counters::KIND_READ_VERSION)) + }, + ) + .await + }) + }); + rustfs_io_metrics::set_get_stage_metrics_enabled(previous_gate); + + assert_eq!(restored, expected_body); + assert_eq!(calls_total, 4, "shape-gated production GET should keep full metadata fanout"); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_scheduled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![4.0], + "shape-gated production GET should schedule the full metadata fanout" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_completed", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![4.0], + "shape-gated production GET must wait for every scheduled metadata response" + ); + assert_eq!( + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_cancelled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)] + ), + vec![0.0], + "shape-gated production GET must not cancel metadata responses" + ); + } + } + #[tokio::test] #[serial_test::serial(body_cache_hook)] async fn prepared_reader_rebuilds_object_info_when_precomputed_value_is_absent() { @@ -1105,7 +1438,7 @@ impl SetDisks { object: &str, opts: &ObjectOptions, ) -> Result { - let snapshot = self.get_object_fileinfo(bucket, object, opts, true, true).await?; + let snapshot = self.get_object_fileinfo(bucket, object, opts, true, false).await?; let object_info = build_get_object_info(snapshot.fi(), bucket, object, opts.versioned || opts.version_suspended); Ok(PreparedGetObjectMetadata { snapshot, @@ -3411,9 +3744,15 @@ fn collect_inline_data_shard_fileinfos_by_index<'a>( if block_index == 0 || block_index > data_shards { continue; } + if file_info.erasure.index != block_index { + continue; + } if !file_info.has_valid_erasure_geometry() { continue; } + if !core::io_primitives::metadata_early_stop_candidate_matches(file_info, fi) { + continue; + } if file_info.data.as_ref().is_none_or(|data| data.is_empty()) { continue; } @@ -9221,6 +9560,9 @@ mod tests { HashAlgorithm::HighwayHash256S }; let shards = erasure.encode_data(payload).expect("payload should encode"); + let version_id = Some(Uuid::new_v4()); + let data_dir = Some(Uuid::new_v4()); + let mod_time = Some(OffsetDateTime::now_utc()); let mut files = Vec::with_capacity(shards.len()); for shard in shards { @@ -9233,6 +9575,16 @@ mod tests { writer.shutdown().await.expect("inline writer should shutdown"); let data = writer.into_inline_data().expect("inline data should be retained"); let mut file = FileInfo::new("bucket/object", erasure.data_shards, erasure.parity_shards); + file.volume = "bucket".to_string(); + file.name = "object".to_string(); + file.size = i64::try_from(payload.len()).expect("test payload should fit i64"); + file.is_latest = true; + file.version_id = version_id; + file.data_dir = data_dir; + file.mod_time = mod_time; + file.metadata.insert("etag".to_string(), "etag-inline".to_string()); + file.add_object_part(1, "part-etag-inline".to_string(), payload.len(), file.mod_time, file.size, None, None); + file.set_inline_data(); file.erasure.index = files.len() + 1; file.data = Some(Bytes::from(data)); files.push(file); @@ -9245,16 +9597,40 @@ mod tests { inline_bitrot_files_for_payload_with_mode(payload, false).await } + fn disk_ordered_fileinfos(files: &[FileInfo]) -> Vec { + let distribution = &files + .first() + .expect("inline data shard fixture should include metadata") + .erasure + .distribution; + distribution + .iter() + .map(|block_index| { + files + .get(block_index.checked_sub(1).expect("erasure block indexes are one-based")) + .expect("inline data shard fixture should include every distributed shard") + .clone() + }) + .collect() + } + fn inline_data_shard_fileinfo( - name: &str, data_blocks: usize, parity_blocks: usize, erasure_index: usize, distribution: &[usize], data: Option<&'static [u8]>, ) -> FileInfo { - let mut fi = FileInfo::new(name, data_blocks, parity_blocks); - fi.name = name.to_string(); + let mut fi = FileInfo::new("object", data_blocks, parity_blocks); + fi.name = "object".to_string(); + fi.volume = "bucket".to_string(); + fi.size = 4; + fi.is_latest = true; + fi.data_dir = Some(Uuid::nil()); + fi.mod_time = Some(OffsetDateTime::UNIX_EPOCH); + fi.metadata.insert("etag".to_string(), "etag-inline".to_string()); + fi.add_object_part(1, "part-etag-inline".to_string(), 4, fi.mod_time, 4, None, None); + fi.set_inline_data(); fi.erasure.index = erasure_index; fi.erasure.distribution = distribution.to_vec(); fi.data = data.map(Bytes::from_static); @@ -9264,36 +9640,41 @@ mod tests { #[test] fn collect_inline_data_shards_by_index_uses_distribution_order() { let distribution = vec![3, 1, 5, 2, 4, 6]; - let mut fi = FileInfo::new("object", 4, 2); + let mut fi = inline_data_shard_fileinfo(4, 2, 1, &distribution, Some(b"x")); + fi.erasure.index = 1; fi.erasure.distribution = distribution.clone(); let files = vec![ - inline_data_shard_fileinfo("block-3", 4, 2, 3, &distribution, Some(b"c")), - inline_data_shard_fileinfo("block-1", 4, 2, 1, &distribution, Some(b"a")), - inline_data_shard_fileinfo("parity-5", 4, 2, 5, &distribution, Some(b"p")), - inline_data_shard_fileinfo("block-2", 4, 2, 2, &distribution, Some(b"b")), - inline_data_shard_fileinfo("block-4", 4, 2, 4, &distribution, Some(b"d")), - inline_data_shard_fileinfo("parity-6", 4, 2, 6, &distribution, Some(b"q")), + inline_data_shard_fileinfo(4, 2, 3, &distribution, Some(b"c")), + inline_data_shard_fileinfo(4, 2, 1, &distribution, Some(b"a")), + inline_data_shard_fileinfo(4, 2, 5, &distribution, Some(b"p")), + inline_data_shard_fileinfo(4, 2, 2, &distribution, Some(b"b")), + inline_data_shard_fileinfo(4, 2, 4, &distribution, Some(b"d")), + inline_data_shard_fileinfo(4, 2, 6, &distribution, Some(b"q")), ]; let data_files = collect_inline_data_shard_fileinfos_by_index(&files, &fi, 4, |_| true).expect("all data shards should be collected"); assert_eq!( - data_files.iter().map(|file| file.name.as_str()).collect::>(), - ["block-1", "block-2", "block-3", "block-4"] + data_files + .iter() + .map(|file| file.data.as_deref().expect("fixture carries inline bytes")) + .collect::>(), + [b"a".as_slice(), b"b".as_slice(), b"c".as_slice(), b"d".as_slice()] ); } #[test] fn collect_inline_data_shards_by_index_rejects_missing_data_shard() { let distribution = vec![1, 2, 3, 4]; - let mut fi = FileInfo::new("object", 2, 2); + let mut fi = inline_data_shard_fileinfo(2, 2, 1, &distribution, Some(b"x")); + fi.erasure.index = 1; fi.erasure.distribution = distribution.clone(); let files = vec![ - inline_data_shard_fileinfo("block-1", 2, 2, 1, &distribution, Some(b"a")), - inline_data_shard_fileinfo("block-2", 2, 2, 2, &distribution, None), - inline_data_shard_fileinfo("parity-3", 2, 2, 3, &distribution, Some(b"p")), - inline_data_shard_fileinfo("parity-4", 2, 2, 4, &distribution, Some(b"q")), + inline_data_shard_fileinfo(2, 2, 1, &distribution, Some(b"a")), + inline_data_shard_fileinfo(2, 2, 2, &distribution, None), + inline_data_shard_fileinfo(2, 2, 3, &distribution, Some(b"p")), + inline_data_shard_fileinfo(2, 2, 4, &distribution, Some(b"q")), ]; assert!(collect_inline_data_shard_fileinfos_by_index(&files, &fi, 2, |_| true).is_none()); @@ -9426,10 +9807,8 @@ mod tests { let payload = vec![b'i'; 192 * 1024]; let (erasure, files, _read_length, _checksum_algo) = inline_bitrot_files_for_payload(&payload).await; - let mut fi = FileInfo::new("bucket/object", erasure.data_shards, erasure.parity_shards); - fi.size = payload.len() as i64; - fi.data = files[0].data.clone(); - fi.add_object_part(1, String::new(), payload.len(), None, payload.len() as i64, None, None); + let fi = files[0].clone(); + let disk_files = disk_ordered_fileinfos(&files); let disks = vec![Some(disk); erasure.total_shard_count()]; let metrics_size_bucket = rustfs_io_metrics::get_object_size_bucket(fi.size); @@ -9438,7 +9817,7 @@ mod tests { "bucket", "object", &fi, - &files, + &disk_files, &disks, true, GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART, @@ -9469,10 +9848,8 @@ mod tests { let payload = vec![b'v'; 64 * 1024]; let payload_size = i64::try_from(payload.len()).expect("test payload size should fit i64"); let (erasure, files, _read_length, _checksum_algo) = inline_bitrot_files_for_payload(&payload).await; - let mut fi = FileInfo::new("bucket/object", erasure.data_shards, erasure.parity_shards); - fi.size = payload_size; - fi.data = files[0].data.clone(); - fi.add_object_part(1, String::new(), payload.len(), None, payload_size, None, None); + let fi = files[0].clone(); + let disk_files = disk_ordered_fileinfos(&files); let mut object_info = ObjectInfo { size: payload_size, @@ -9503,7 +9880,7 @@ mod tests { "bucket", "object", &fi, - &files, + &disk_files, &vec![Some(disk); erasure.total_shard_count()], true, GET_CODEC_STREAMING_OBJECT_CLASS_PLAIN_SINGLE_PART, diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 6739c8c8a..2b65f79ac 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -188,6 +188,74 @@ async fn get_object_reader_with_context( GetObjectReader::new_with_resolver(reader, range, object_info, opts, headers, ctx.object_encryption_resolver()).await } +fn data_read_metadata_early_stop_request_shape_allowed(range: &Option, opts: &ObjectOptions) -> bool { + range.is_none() + && opts.part_number.is_none() + && opts.version_id.is_none() + && !opts.incl_free_versions + && !opts.skip_free_version + && !opts.raw_data_movement_read + && !opts.data_movement + && !crate::object_api::restore_request_active(opts) +} + +#[cfg(test)] +mod data_read_metadata_early_stop_request_shape_tests { + use super::*; + + #[test] + fn data_read_metadata_early_stop_only_allows_whole_latest_plain_get_shape() { + assert!(data_read_metadata_early_stop_request_shape_allowed(&None, &ObjectOptions::default())); + + let range = Some(HTTPRangeSpec { + is_suffix_length: false, + start: 0, + end: 0, + }); + assert!(!data_read_metadata_early_stop_request_shape_allowed(&range, &ObjectOptions::default())); + + let part_opts = ObjectOptions { + part_number: Some(1), + ..Default::default() + }; + assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &part_opts)); + + let version_opts = ObjectOptions { + version_id: Some(Uuid::new_v4().to_string()), + ..Default::default() + }; + assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &version_opts)); + + let incl_free_opts = ObjectOptions { + incl_free_versions: true, + ..Default::default() + }; + assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &incl_free_opts)); + + let skip_free_opts = ObjectOptions { + skip_free_version: true, + ..Default::default() + }; + assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &skip_free_opts)); + + let data_movement_opts = ObjectOptions { + data_movement: true, + ..Default::default() + }; + assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &data_movement_opts)); + + let raw_data_movement_opts = ObjectOptions { + raw_data_movement_read: true, + ..Default::default() + }; + assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &raw_data_movement_opts)); + + let mut restore_opts = ObjectOptions::default(); + restore_opts.transition.restore_request.days = Some(1); + assert!(!data_read_metadata_early_stop_request_shape_allowed(&None, &restore_opts)); + } +} + /// Length of the full plaintext body when — and only when — this read's output /// is exactly the object's complete plaintext, so the app-layer body cache may /// serve it in place of the erasure read. @@ -431,7 +499,16 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { let (snapshot, prepared_object_info) = if let Some(prepared) = take_prepared_get_object_metadata() { (prepared.snapshot, prepared.object_info) } else { - match self.get_object_fileinfo(bucket, object, opts, true, true).await { + match self + .get_object_fileinfo( + bucket, + object, + opts, + true, + data_read_metadata_early_stop_request_shape_allowed(&range, opts), + ) + .await + { Ok(snapshot) => (snapshot, None), Err(err) => { rustfs_io_metrics::record_get_object_metadata_phase_duration(metadata_stage_start.elapsed().as_secs_f64()); @@ -6102,7 +6179,10 @@ mod inline_put_commit_path_tests { mod get_object_downstream_close_accounting_tests { use super::hermetic_set_disks_support::hermetic_set_disks; use super::*; - use crate::diagnostics::get::{GET_OBJECT_PATH_INTERNAL_META, GET_STAGE_DECODE, GET_STAGE_EMIT, GetObjectFailureReason}; + use crate::diagnostics::get::{ + GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, GET_OBJECT_PATH_INTERNAL_META, GET_STAGE_DECODE, GET_STAGE_EMIT, + GetObjectFailureReason, + }; use crate::disk::RUSTFS_META_BUCKET; use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; @@ -6221,7 +6301,22 @@ mod get_object_downstream_close_accounting_tests { let previous_gate = rustfs_io_metrics::get_stage_metrics_enabled(); rustfs_io_metrics::set_get_stage_metrics_enabled(true); - let (internal_missing, legacy_unknown, internal_fanout, legacy_fanout) = metrics::with_local_recorder(&recorder, || { + let ( + internal_missing, + legacy_unknown, + internal_fanout, + legacy_fanout, + internal_scheduled, + legacy_scheduled, + internal_completed, + legacy_completed, + internal_cancelled, + legacy_cancelled, + internal_unsafe_miss, + legacy_unsafe_miss, + internal_saved, + legacy_saved, + ) = metrics::with_local_recorder(&recorder, || { runtime.block_on(async { let (_temp_dirs, _disk_stores, set_disks) = hermetic_set_disks(4).await; let options = ObjectOptions { @@ -6265,6 +6360,54 @@ mod get_object_downstream_close_accounting_tests { "rustfs_io_get_object_metadata_fanout_error_responses", &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)], ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_scheduled", + &[("path", GET_OBJECT_PATH_INTERNAL_META)], + ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_scheduled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)], + ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_completed", + &[("path", GET_OBJECT_PATH_INTERNAL_META)], + ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_completed", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)], + ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_cancelled", + &[("path", GET_OBJECT_PATH_INTERNAL_META)], + ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_fanout_cancelled", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)], + ), + recorder.counter_value( + "rustfs_io_get_object_metadata_early_stop_total", + &[ + ("path", GET_OBJECT_PATH_INTERNAL_META), + ("decision", "miss"), + ("reason", GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST), + ], + ), + recorder.counter_value( + "rustfs_io_get_object_metadata_early_stop_total", + &[ + ("path", GET_OBJECT_PATH_LEGACY_DUPLEX), + ("decision", "miss"), + ("reason", GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST), + ], + ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_early_stop_saved_responses", + &[("path", GET_OBJECT_PATH_INTERNAL_META)], + ), + recorder.histogram_values( + "rustfs_io_get_object_metadata_early_stop_saved_responses", + &[("path", GET_OBJECT_PATH_LEGACY_DUPLEX)], + ), ) }) }); @@ -6277,6 +6420,50 @@ mod get_object_downstream_close_accounting_tests { ); assert_eq!(internal_fanout, vec![4.0], "internal metadata fanout must retain its path label"); assert!(legacy_fanout.is_empty(), "internal metadata fanout must not leak into legacy_duplex"); + assert_eq!( + internal_scheduled, + vec![4.0], + "internal metadata lifecycle scheduled count must retain its path label" + ); + assert!( + legacy_scheduled.is_empty(), + "internal metadata lifecycle scheduled count must not leak into legacy_duplex" + ); + assert_eq!( + internal_completed, + vec![4.0], + "internal metadata lifecycle completed count must retain its path label" + ); + assert!( + legacy_completed.is_empty(), + "internal metadata lifecycle completed count must not leak into legacy_duplex" + ); + assert_eq!( + internal_cancelled, + vec![0.0], + "internal metadata full-wait lifecycle must record zero cancellations" + ); + assert!( + legacy_cancelled.is_empty(), + "internal metadata lifecycle cancelled count must not leak into legacy_duplex" + ); + assert_eq!( + internal_unsafe_miss, 1, + "internal metadata unsafe early-stop miss must retain its path label" + ); + assert_eq!( + legacy_unsafe_miss, 0, + "internal metadata unsafe early-stop miss must not leak into legacy_duplex" + ); + assert_eq!( + internal_saved, + vec![0.0], + "internal metadata unsafe miss must record zero saved responses on internal_meta" + ); + assert!( + legacy_saved.is_empty(), + "internal metadata unsafe miss saved responses must not leak into legacy_duplex" + ); } } diff --git a/crates/heal/tests/endpoint_index_test.rs b/crates/heal/tests/endpoint_index_test.rs index 52684b9af..0df8ab5a6 100644 --- a/crates/heal/tests/endpoint_index_test.rs +++ b/crates/heal/tests/endpoint_index_test.rs @@ -14,6 +14,8 @@ //! test endpoint index settings +#![recursion_limit = "256"] + use std::net::SocketAddr; use tempfile::TempDir; use tokio_util::sync::CancellationToken; diff --git a/crates/heal/tests/heal_b5_versioned_regression_test.rs b/crates/heal/tests/heal_b5_versioned_regression_test.rs index 2b4a0c604..61a542955 100644 --- a/crates/heal/tests/heal_b5_versioned_regression_test.rs +++ b/crates/heal/tests/heal_b5_versioned_regression_test.rs @@ -22,6 +22,8 @@ //! bucket-metadata-sys OnceCell) — under `cargo nextest` each test runs //! in its own process so the OnceCell never collides. +#![recursion_limit = "256"] + use http::HeaderMap; use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_heal::heal::{ diff --git a/crates/heal/tests/heal_b920_subquorum_union_test.rs b/crates/heal/tests/heal_b920_subquorum_union_test.rs index f3121aaa3..6d188b3f3 100644 --- a/crates/heal/tests/heal_b920_subquorum_union_test.rs +++ b/crates/heal/tests/heal_b920_subquorum_union_test.rs @@ -21,6 +21,8 @@ //! These drive the REAL `ECStoreHealStorage` + `ECStore` against real disks. //! Every test is `#[serial]`; under `cargo nextest` each runs in its own process. +#![recursion_limit = "256"] + use http::HeaderMap; use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_heal::heal::storage::{ diff --git a/crates/heal/tests/heal_integration_test.rs b/crates/heal/tests/heal_integration_test.rs index 7094614db..21a0ab78b 100644 --- a/crates/heal/tests/heal_integration_test.rs +++ b/crates/heal/tests/heal_integration_test.rs @@ -12,6 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. +#![recursion_limit = "256"] + use http::HeaderMap; use rustfs_common::heal_channel::{HealOpts, HealScanMode}; use rustfs_heal::heal::{ diff --git a/crates/io-metrics/src/lib.rs b/crates/io-metrics/src/lib.rs index 658dee400..26723d11e 100644 --- a/crates/io-metrics/src/lib.rs +++ b/crates/io-metrics/src/lib.rs @@ -812,6 +812,17 @@ pub fn record_get_object_metadata_fanout_shape(path: &'static str, total: usize, .record(metadata_fanout_count_to_f64(non_valid)); } +/// Record task lifecycle shape for one GetObject metadata fanout. +#[inline(always)] +pub fn record_get_object_metadata_fanout_lifecycle(path: &'static str, scheduled: usize, completed: usize, cancelled: usize) { + if !get_stage_metrics_enabled() { + return; + } + histogram!("rustfs_io_get_object_metadata_fanout_scheduled", "path" => path).record(metadata_fanout_count_to_f64(scheduled)); + histogram!("rustfs_io_get_object_metadata_fanout_completed", "path" => path).record(metadata_fanout_count_to_f64(completed)); + histogram!("rustfs_io_get_object_metadata_fanout_cancelled", "path" => path).record(metadata_fanout_count_to_f64(cancelled)); +} + /// Record a guarded metadata early-stop hit for GetObject. #[inline(always)] pub fn record_get_object_metadata_early_stop_hit(path: &'static str, reason: &'static str) { @@ -2698,6 +2709,7 @@ mod tests { record_get_object_quorum_reached_latency("legacy_duplex", 0.002); record_get_object_metadata_response("legacy_duplex", "valid"); record_get_object_metadata_fanout_shape("legacy_duplex", 4, 3, 1, 1); + record_get_object_metadata_fanout_lifecycle("legacy_duplex", 4, 3, 1); record_get_object_metadata_early_stop_hit("legacy_duplex", "valid_quorum"); record_get_object_metadata_early_stop_miss("legacy_duplex", "insufficient_quorum"); record_get_object_metadata_early_stop_saved_responses("legacy_duplex", 1); @@ -2768,6 +2780,38 @@ mod tests { assert!(remote_scheduled >= remote_avoid_potential); } + #[test] + fn metadata_fanout_lifecycle_records_named_histograms() { + let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + metrics::with_local_recorder(&recorder, || { + set_get_stage_metrics_enabled(true); + record_get_object_metadata_fanout_lifecycle("legacy_duplex", 4, 3, 1); + set_get_stage_metrics_enabled(false); + }); + + let metrics = snapshotter.snapshot().into_vec(); + for (name, expected) in [ + ("rustfs_io_get_object_metadata_fanout_scheduled", 4.0), + ("rustfs_io_get_object_metadata_fanout_completed", 3.0), + ("rustfs_io_get_object_metadata_fanout_cancelled", 1.0), + ] { + let value = metrics.iter().find_map(|(composite, _, _, value)| { + let has_path = composite + .key() + .labels() + .any(|label| label.key() == "path" && label.value() == "legacy_duplex"); + (composite.kind() == MetricKind::Histogram && composite.key().name() == name && has_path).then_some(value) + }); + assert!( + matches!(value, Some(DebugValue::Histogram(values)) if values.len() == 1 && values[0].0 == expected), + "{name} must record the exact fanout lifecycle sample" + ); + } + } + #[test] fn test_record_get_object_fill_metrics() { record_get_object_fill_queued("codec_streaming", "single_inflight", 1); diff --git a/crates/protocols/tests/swift_metadata_persistence.rs b/crates/protocols/tests/swift_metadata_persistence.rs index 86f205a45..6b0b2df8d 100644 --- a/crates/protocols/tests/swift_metadata_persistence.rs +++ b/crates/protocols/tests/swift_metadata_persistence.rs @@ -23,6 +23,7 @@ //! two are tested together because a reload is the only way to tell a real //! merge from one that happened to look right in the cache. +#![recursion_limit = "256"] #![cfg(feature = "swift")] use std::collections::HashMap; diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index 6b6deb29d..4a5cc7543 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +#![recursion_limit = "256"] #![cfg_attr(docsrs, feature(doc_auto_cfg))] #![warn( // missing_docs, diff --git a/crates/scanner/tests/lifecycle_integration_test.rs b/crates/scanner/tests/lifecycle_integration_test.rs index 9fef7a372..0a43a8060 100644 --- a/crates/scanner/tests/lifecycle_integration_test.rs +++ b/crates/scanner/tests/lifecycle_integration_test.rs @@ -12,6 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. +#![recursion_limit = "256"] + use futures::FutureExt; use rustfs_config::ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT; use rustfs_scanner::scanner_folder::ScannerItem;