From d087416ff9ce90c1cd105c5898f5855a4ee0bc4b Mon Sep 17 00:00:00 2001 From: Chris Date: Sun, 4 Oct 2026 01:04:37 +0800 Subject: [PATCH] refactor(storage): separate metadata scheduling from quorum decisions (#8327) --- .../src/set_disk/core/io_primitives.rs | 1316 +-------------- .../src/set_disk/core/metadata_quorum.rs | 213 +++ .../src/set_disk/core/metadata_read.rs | 1443 +++++++++++++++++ crates/ecstore/src/set_disk/core/mod.rs | 3 +- crates/ecstore/src/set_disk/mod.rs | 23 +- crates/ecstore/src/set_disk/ops/heal.rs | 46 +- crates/ecstore/src/set_disk/read.rs | 39 +- docs/architecture/compat-cleanup-register.md | 2 + .../architecture/ecstore-module-split-plan.md | 36 + 9 files changed, 1827 insertions(+), 1294 deletions(-) create mode 100644 crates/ecstore/src/set_disk/core/metadata_read.rs diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 0d025b45d..3757b59c6 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -16,11 +16,9 @@ //! operation families (P5 of the God-Object split, tracking backlog#815, //! issue backlog#820). //! -//! These are the metadata-fanout quorum accumulator, bitrot reader -//! scheduling/creation, shard-cost classification, and read-repair heal -//! dedup helpers relocated verbatim from `set_disk/read.rs`. Bodies are -//! byte-identical to the pre-move sources; only the module header -//! (`use super::*;` -> `use super::super::*;`) and item visibility change. +//! Bitrot reader setup, shard-cost classification, read repair and object +//! mutation primitives. Metadata scheduling and decisions have separate owners +//! in `metadata_read` and `metadata_quorum`; legacy callers retain adapters here. #[cfg(test)] use super::super::ENV_RUSTFS_GET_METADATA_DATA_READ_EARLY_STOP_ENABLE; @@ -46,48 +44,61 @@ use super::super::{ GetCodecStreamingFallbackReason, GetObjectMetadataCacheEntry, HTTPPreconditions, HashAlgorithm, HealAdmissionResult, HealChannelPriority, HealRequestSource, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_SET_DISK, MultipartWriteQuorumContext, OBJECT_OP_IGNORED_ERRS, ObjectOptions, ObjectPartInfo, OffsetDateTime, RUSTFS_META_BUCKET, RUSTFS_META_MULTIPART_BUCKET, - RawFileInfo, ReadMultipleReq, ReadMultipleResp, ReadOptions, Result, SLASH_SEPARATOR, STORAGE_FORMAT_FILE, SetDisks, - SnapshotLeaseToken, StorageError, UpdateMetadataOpts, Uuid, build_inline_bitrot_readers_from_refs, - can_try_inline_data_shards_direct, capacity_scope_from_disks, codec_streaming_rollout_applies, coding, - collect_inline_data_shard_fileinfos_by_index_or_reason, current_dirty_generation, debug, disk, - file_info_is_valid_for_metadata, get_metadata_slowtail_fault_request, info, inline_erasure_shard_file_offset, - inline_erasure_shard_size, is_get_metadata_data_read_early_stop_enabled, is_get_metadata_early_stop_bounded_fanout_enabled, - is_get_metadata_early_stop_enabled, is_get_metadata_non_inline_data_read_early_stop_enabled, is_object_dangling, - is_version_early_stop_enabled, issue3031_diag_enabled, join_all, join_errs, log_multipart_write_quorum_failure, - merge_file_meta_versions, object_fits_single_block, path_join_buf, record_global_dirty_scope, reduce_read_quorum_errs, - reduce_write_quorum_errs, send_heal_request_with_admission, should_prevent_write, to_object_err, - try_read_inline_data_shards_direct, warn, + RawFileInfo, ReadMultipleReq, ReadMultipleResp, Result, SLASH_SEPARATOR, STORAGE_FORMAT_FILE, SetDisks, SnapshotLeaseToken, + StorageError, UpdateMetadataOpts, Uuid, capacity_scope_from_disks, current_dirty_generation, debug, disk, + file_info_is_valid_for_metadata, info, is_object_dangling, issue3031_diag_enabled, join_all, join_errs, + log_multipart_write_quorum_failure, merge_file_meta_versions, path_join_buf, record_global_dirty_scope, + reduce_read_quorum_errs, reduce_write_quorum_errs, send_heal_request_with_admission, should_prevent_write, to_object_err, + warn, }; #[cfg(test)] +use super::super::{ReadOptions, coding}; +#[cfg(test)] pub(in crate::set_disk) use super::metadata_quorum::MetadataEarlyStopDecision; -pub(in crate::set_disk) use super::metadata_quorum::{ - MetadataQuorumAccumulator, is_metadata_fanout_ignored_error, metadata_early_stop_candidate_matches, +#[cfg(test)] +use super::metadata_quorum::MetadataQuorumAccumulator; +pub(in crate::set_disk) use super::metadata_quorum::metadata_early_stop_candidate_matches; +#[cfg(test)] +use super::metadata_read::MetadataReadResult; +#[cfg(test)] +use super::metadata_read::{ + ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE, ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, + ExpectedBatchReadVersionItem, data_read_early_stop_inline_body_miss_reason, map_batch_read_version_responses, + metadata_distribution_key, +}; +#[cfg(test)] +pub(in crate::set_disk) use super::metadata_read::{ + MetadataFanoutDiagnostics, MetadataFanoutObservation, bounded_metadata_fanout_order, non_inline_data_read_early_stop_allowed, }; #[cfg(test)] use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE; #[cfg(test)] -use crate::diagnostics::get::GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH; -#[cfg(test)] use crate::diagnostics::get::GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_PAYLOAD; +#[cfg(test)] use crate::diagnostics::get::{ GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE, - GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED, GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM, - GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, 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_DIRECT_MEMORY, - GET_OBJECT_PATH_INTERNAL_META, GET_OBJECT_PATH_LEGACY_DUPLEX, 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, get_stage_timer_if_enabled, - record_get_stage_duration_if_enabled, + GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED, GET_METADATA_RESPONSE_CORRUPT, + GET_METADATA_RESPONSE_DISK_NOT_FOUND, GET_METADATA_RESPONSE_ERROR, GET_METADATA_RESPONSE_IGNORED, + GET_METADATA_RESPONSE_VALID, +}; +#[cfg(test)] +use crate::diagnostics::get::{ + GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH, GET_OBJECT_PATH_INTERNAL_META, + GET_OBJECT_PATH_LEGACY_DUPLEX, }; #[cfg(test)] use crate::diagnostics::get::{ GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER, GET_METADATA_EARLY_STOP_REASON_ERROR, GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM, }; +use crate::diagnostics::get::{ + GET_OBJECT_PATH_DIRECT_MEMORY, 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, get_stage_timer_if_enabled, record_get_stage_duration_if_enabled, +}; #[cfg(test)] use crate::disk::CHECK_PART_FILE_NOT_FOUND; use crate::disk::DiskAPI; @@ -95,7 +106,6 @@ use crate::disk::DiskAPI; use crate::disk::DiskOption; #[cfg(test)] use crate::disk::RUSTFS_META_TMP_BUCKET; -use crate::disk::disk_store::get_drive_metadata_timeout; #[cfg(test)] use crate::disk::endpoint::Endpoint; #[cfg(test)] @@ -103,10 +113,11 @@ use crate::disk::format::FormatV3; use crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX; #[cfg(test)] use crate::disk::new_disk; +#[cfg(test)] +use crate::disk::{BatchReadVersionItem, BatchReadVersionResp}; use crate::disk::{ - BATCH_READ_VERSION_MAX_ITEMS, BatchReadVersionItem, BatchReadVersionReq, BatchReadVersionResp, DataDirDeleteStatus, Disk, - OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK, PartTransactionAction, - STORAGE_FORMAT_FILE_BACKUP, part_transaction_path, + DataDirDeleteStatus, OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK, + PartTransactionAction, STORAGE_FORMAT_FILE_BACKUP, part_transaction_path, }; use crate::erasure::coding::BitrotReader; use crate::io_support::bitrot::ShardReader; @@ -122,9 +133,6 @@ use crate::storage_api_contracts::object::ObjectOperations; use futures::FutureExt as _; use futures::stream::{FuturesUnordered, StreamExt}; use metrics::counter; -use rustfs_io_metrics::internode_metrics::{ - INTERNODE_STAGE_BATCH_READ_VERSION_COALESCER_WAIT, INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_MAP, -}; use std::{ collections::{HashMap, HashSet, VecDeque}, future::Future, @@ -136,353 +144,12 @@ 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("/") -} - -fn read_version_coalescing_enabled() -> bool { - let enabled = || { - rustfs_utils::get_env_opt_str(ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE) - .is_some_and(|value| value.eq_ignore_ascii_case("auto") || value.eq_ignore_ascii_case("on")) - }; - - #[cfg(test)] - { - enabled() - } - - #[cfg(not(test))] - { - static ENABLED: OnceLock = OnceLock::new(); - *ENABLED.get_or_init(enabled) - } -} - -fn read_version_coalescing_delay() -> Duration { - #[cfg(test)] - { - let micros = rustfs_utils::get_env_u64( - ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, - DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, - ); - Duration::from_micros(micros) - } - - #[cfg(not(test))] - { - static DELAY: OnceLock = OnceLock::new(); - *DELAY.get_or_init(|| { - Duration::from_micros(rustfs_utils::get_env_u64( - ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, - DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, - )) - }) - } -} - -struct CoalescedReadVersionRequest { - item: BatchReadVersionItem, - tx: oneshot::Sender>, -} - -#[derive(Clone, Debug, Eq, PartialEq)] -struct ExpectedBatchReadVersionItem { - path: String, - version_id: String, -} - -impl From<&BatchReadVersionItem> for ExpectedBatchReadVersionItem { - fn from(item: &BatchReadVersionItem) -> Self { - Self { - path: item.path.clone(), - version_id: item.version_id.clone(), - } - } -} - -#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] -struct ReadVersionCoalescerKey { - disk: usize, - incl_free_versions: bool, - read_data: bool, - healing: bool, -} - -impl ReadVersionCoalescerKey { - fn new(disk: &DiskStore, opts: &ReadOptions) -> Self { - Self { - disk: Arc::as_ptr(disk) as usize, - incl_free_versions: opts.incl_free_versions, - read_data: opts.read_data, - healing: opts.healing, - } - } -} - -#[derive(Default)] -struct ReadVersionCoalescer { - lanes: HashMap>, -} - -fn read_version_coalescer() -> &'static Mutex { - static COALESCER: OnceLock> = OnceLock::new(); - COALESCER.get_or_init(|| Mutex::new(ReadVersionCoalescer::default())) -} - -fn record_read_version_coalescer_event(event: &'static str, item_count: usize) { - counter!( - METRIC_GET_METADATA_READ_VERSION_COALESCER_TOTAL, - "event" => event, - "item_count" => item_count.to_string() - ) - .increment(1); -} - -fn batch_read_version_stage_timer() -> Option { - rustfs_io_metrics::get_stage_metrics_enabled().then(Instant::now) -} - -fn record_batch_read_version_stage(stage: &'static str, started_at: Option) { - if let Some(started_at) = started_at { - crate::cluster::rpc::runtime_sources::record_remote_disk_grpc_batch_read_version_stage(stage, started_at.elapsed()); - } -} - -async fn read_version_via_coalescer( - disk: DiskStore, - org_bucket: &str, - bucket: &str, - object: &str, - version_id: &str, - opts: &ReadOptions, - allow_coalescing: bool, -) -> disk::error::Result { - if !allow_coalescing || !read_version_coalescing_enabled() { - return disk.read_version(org_bucket, bucket, object, version_id, opts).await; - } - if !matches!(disk.as_ref(), Disk::Remote(_)) { - record_read_version_coalescer_event("bypass_non_remote", 1); - return disk.read_version(org_bucket, bucket, object, version_id, opts).await; - } - - let (tx, rx) = oneshot::channel(); - let item = BatchReadVersionItem { - org_volume: org_bucket.to_string(), - volume: bucket.to_string(), - path: object.to_string(), - version_id: version_id.to_string(), - }; - let lane_key = ReadVersionCoalescerKey::new(&disk, opts); - let pending = { - let mut coalescer = read_version_coalescer().lock().await; - let lane = coalescer.lanes.entry(lane_key).or_default(); - let schedule_delayed_flush = lane.is_empty(); - lane.push(CoalescedReadVersionRequest { item, tx }); - if lane.len() >= BATCH_READ_VERSION_MAX_ITEMS { - coalescer.lanes.remove(&lane_key) - } else if schedule_delayed_flush { - let disk = disk.clone(); - let task_opts = *opts; - tokio::spawn(async move { - tokio::time::sleep(read_version_coalescing_delay()).await; - flush_read_version_coalescer_lane(lane_key, disk, task_opts).await; - }); - None - } else { - None - } - }; - - if let Some(pending) = pending { - flush_read_version_coalescer_pending(lane_key, disk, *opts, pending).await; - } - - let wait_started = batch_read_version_stage_timer(); - let response = rx - .await - .unwrap_or_else(|_| Err(DiskError::other("coalesced read_version response channel closed"))); - record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_COALESCER_WAIT, wait_started); - response -} - -async fn flush_read_version_coalescer_lane(lane_key: ReadVersionCoalescerKey, disk: DiskStore, opts: ReadOptions) { - let pending = { - let mut coalescer = read_version_coalescer().lock().await; - coalescer.lanes.remove(&lane_key).unwrap_or_default() - }; - flush_read_version_coalescer_pending(lane_key, disk, opts, pending).await; -} - -async fn flush_read_version_coalescer_pending( - lane_key: ReadVersionCoalescerKey, - disk: DiskStore, - opts: ReadOptions, - pending: Vec, -) { - if pending.is_empty() { - return; - } - - // Only the #[cfg(test)] counter-recording block below reads this. - #[cfg(not(test))] - let _ = lane_key; - #[cfg(test)] - { - let mut observed_paths = HashSet::new(); - for request in &pending { - if observed_paths.insert(request.item.path.as_str()) { - disk_call_counters::record(&request.item.path, disk_call_counters::KIND_BATCH_READ_VERSION, lane_key.disk); - } - } - } - - let mut senders = Vec::with_capacity(pending.len()); - let mut items = Vec::with_capacity(pending.len()); - for request in pending { - senders.push(request.tx); - items.push(request.item); - } - - let expected_items = items.iter().map(ExpectedBatchReadVersionItem::from).collect::>(); - record_read_version_coalescer_event("attempted_batch", items.len()); - let result = - match tokio::time::timeout(get_drive_metadata_timeout(), disk.batch_read_version(BatchReadVersionReq { items, opts })) - .await - { - Ok(result) => result, - Err(_) => Err(DiskError::Timeout), - }; - match result { - Ok(responses) => { - let map_started = batch_read_version_stage_timer(); - let results = map_batch_read_version_responses(&expected_items, responses); - record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_MAP, map_started); - for (tx, result) in senders.into_iter().zip(results) { - let _ = tx.send(result); - } - } - Err(err) => { - let message = err.to_string(); - for tx in senders { - let _ = tx.send(Err(DiskError::other(message.clone()))); - } - } - } -} - -fn map_batch_read_version_responses( - expected_items: &[ExpectedBatchReadVersionItem], - responses: Vec, -) -> Vec> { - let mut results = (0..expected_items.len()) - .map(|_| Err(DiskError::other("coalesced read_version response missing"))) - .collect::>(); - let mut seen = vec![false; expected_items.len()]; - for response in responses { - let Some(expected) = expected_items.get(response.index) else { - continue; - }; - let Some(slot) = results.get_mut(response.index) else { - continue; - }; - if seen[response.index] { - *slot = Err(DiskError::other("coalesced read_version response duplicate index")); - continue; - } - seen[response.index] = true; - if response.path != expected.path || response.version_id != expected.version_id { - *slot = Err(DiskError::other("coalesced read_version response identity mismatch")); - } else { - *slot = if response.success { - Ok(response.file_info) - } else { - Err(batch_read_version_response_error(response.error_code, response.error)) - }; - } - } - results -} - -fn batch_read_version_response_error(error_code: u32, error: String) -> DiskError { - match DiskError::from_u32(error_code) { - Some(DiskError::Io(_)) | None => DiskError::other(error), - Some(error) => error, - } -} - -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::{Mutex, RwLock, oneshot}; +use tokio::sync::RwLock; use tokio::task::JoinSet; -struct AbortOnDropJoinHandle(tokio::task::JoinHandle); - -impl Future for AbortOnDropJoinHandle { - type Output = std::result::Result; - - fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { - Pin::new(&mut self.0).poll(cx) - } -} - -impl Drop for AbortOnDropJoinHandle { - fn drop(&mut self) { - self.0.abort(); - } -} - pub(in crate::set_disk) const EVENT_SET_DISK_READ: &str = "set_disk_read"; pub(in crate::set_disk) const ENV_RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP: &str = "RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP"; -const ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE: &str = "RUSTFS_GET_METADATA_READ_VERSION_COALESCE"; -const ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS: &str = "RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS"; -const DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS: u64 = 200; -const METRIC_GET_METADATA_READ_VERSION_COALESCER_TOTAL: &str = "rustfs_get_metadata_read_version_coalescer_total"; pub(crate) const ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE: &str = "RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE"; /// Default reader-setup strategy for the GET read path (rustfs/backlog#1215, /// #1159, #923). @@ -570,134 +237,6 @@ pub(in crate::set_disk) fn codec_streaming_reader_setup_fallback_reason( (missing_shards > 0).then_some(GetCodecStreamingFallbackReason::ReadQuorumNotSafe) } -#[derive(Clone, Copy, Debug)] -pub(in crate::set_disk) struct MetadataFanoutObservation { - pub(in crate::set_disk) outcome: &'static str, - pub(in crate::set_disk) elapsed: Duration, - pub(in crate::set_disk) valid: bool, - pub(in crate::set_disk) ignored: bool, -} - -impl MetadataFanoutObservation { - pub(in crate::set_disk) fn from_file_info(file_info: &FileInfo, elapsed: Duration) -> Self { - if file_info_is_valid_for_metadata(file_info) { - Self { - outcome: GET_METADATA_RESPONSE_VALID, - elapsed, - valid: true, - ignored: false, - } - } else { - Self { - outcome: GET_METADATA_RESPONSE_ERROR, - elapsed, - valid: false, - ignored: false, - } - } - } - - pub(in crate::set_disk) fn from_error(err: &DiskError, elapsed: Duration) -> Self { - Self { - outcome: classify_metadata_response_error(err), - elapsed, - valid: false, - ignored: is_metadata_fanout_ignored_error(err), - } - } -} - -#[derive(Clone, Debug, Default)] -pub(in crate::set_disk) struct MetadataFanoutDiagnostics { - pub(in crate::set_disk) fanout_duration: Duration, - pub(in crate::set_disk) observations: Vec, -} - -impl MetadataFanoutDiagnostics { - pub(in crate::set_disk) fn new(fanout_duration: Duration, observations: Vec) -> Self { - Self { - fanout_duration, - observations, - } - } - - pub(in crate::set_disk) fn total_responses(&self) -> usize { - self.observations.len() - } - - pub(in crate::set_disk) fn valid_responses(&self) -> usize { - self.observations.iter().filter(|observation| observation.valid).count() - } - - pub(in crate::set_disk) fn ignored_responses(&self) -> usize { - self.observations.iter().filter(|observation| observation.ignored).count() - } - - pub(in crate::set_disk) fn non_valid_responses(&self) -> usize { - self.total_responses().saturating_sub(self.valid_responses()) - } - - pub(in crate::set_disk) fn first_response_latency(&self) -> Option { - self.observations.iter().map(|observation| observation.elapsed).min() - } - - pub(in crate::set_disk) fn first_valid_response_latency(&self) -> Option { - self.observations - .iter() - .filter(|observation| observation.valid) - .map(|observation| observation.elapsed) - .min() - } - - pub(in crate::set_disk) fn slowest_response_latency(&self) -> Option { - self.observations.iter().map(|observation| observation.elapsed).max() - } - - pub(in crate::set_disk) fn quorum_candidate_latency(&self, read_quorum: usize) -> Option { - if read_quorum == 0 { - return Some(Duration::ZERO); - } - - let mut valid_latencies = self - .observations - .iter() - .filter(|observation| observation.valid) - .map(|observation| observation.elapsed) - .collect::>(); - valid_latencies.sort_unstable(); - valid_latencies.get(read_quorum.saturating_sub(1)).copied() - } - - pub(in crate::set_disk) fn record(&self, path: &'static str) { - rustfs_io_metrics::record_get_object_metadata_fanout_duration(path, self.fanout_duration.as_secs_f64()); - if let Some(latency) = self.first_response_latency() { - rustfs_io_metrics::record_get_object_first_metadata_response_latency(path, latency.as_secs_f64()); - } - if let Some(latency) = self.first_valid_response_latency() { - rustfs_io_metrics::record_get_object_first_valid_metadata_response_latency(path, latency.as_secs_f64()); - } - if let Some(latency) = self.slowest_response_latency() { - rustfs_io_metrics::record_get_object_slowest_metadata_response_latency(path, latency.as_secs_f64()); - } - rustfs_io_metrics::record_get_object_metadata_fanout_shape( - path, - self.total_responses(), - self.valid_responses(), - self.ignored_responses(), - self.non_valid_responses(), - ); - for observation in &self.observations { - rustfs_io_metrics::record_get_object_metadata_response(path, observation.outcome); - } - } - - pub(in crate::set_disk) fn record_quorum_candidate_latency(&self, path: &'static str, read_quorum: usize) { - if let Some(latency) = self.quorum_candidate_latency(read_quorum) { - rustfs_io_metrics::record_get_object_quorum_reached_latency(path, latency.as_secs_f64()); - } - } -} - #[derive(Clone, Debug)] pub(in crate::set_disk) enum MetadataCacheLookup { Hit(Arc), @@ -705,205 +244,6 @@ pub(in crate::set_disk) enum MetadataCacheLookup { RejectedInsufficientQuorum, } -pub(in crate::set_disk) async fn data_read_early_stop_inline_body_miss_reason( - bucket: &str, - object: &str, - candidate: &FileInfo, - parts_metadata: &[FileInfo], - disks: &[Option], -) -> Option<&'static str> { - if let Some(reason) = data_read_early_stop_inline_candidate_miss_reason(candidate) { - return Some(reason); - } - - 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 Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY); - }; - let data_files = - match collect_inline_data_shard_fileinfos_by_index_or_reason(parts_metadata, candidate, erasure.data_shards, |index| { - disks.get(index).is_some_and(Option::is_some) - }) { - Ok(data_files) => data_files, - Err(reason) => return Some(reason), - }; - - let Some(part) = candidate.parts.first() else { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE); - }; - let Ok(object_size) = usize::try_from(candidate.size) else { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); - }; - 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 Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY); - }; - - match try_read_inline_data_shards_direct(&mut readers, erasure.data_shards, read_length, object_size).await { - Some(body) if body.len() == object_size => None, - _ => Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY), - } -} - -fn data_read_early_stop_inline_candidate_miss_reason(candidate: &FileInfo) -> Option<&'static str> { - // `inline_data` excludes remote objects; this diagnostic reports them separately. - if !rustfs_utils::http::contains_key_str(&candidate.metadata, rustfs_utils::http::SUFFIX_INLINE_DATA) { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE); - } - if candidate.is_compressed() - || candidate - .metadata - .keys() - .any(|key| rustfs_utils::http::is_object_encryption_marker(key)) - { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED); - } - if candidate.is_remote() { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE); - } - if candidate.deleted { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED); - } - if candidate.size <= 0 { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); - } - if candidate.parts.len() != 1 { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE); - } - if !candidate.has_valid_erasure_geometry() { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY); - } - - let Ok(object_size) = usize::try_from(candidate.size) else { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); - }; - if candidate.parts.first().is_none_or(|part| part.size != object_size) { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE); - } - if !can_try_inline_data_shards_direct(object_size, candidate.erasure.block_size) { - return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); - } - None -} - -pub(in crate::set_disk) fn non_inline_data_read_candidate_is_safe(candidate: &FileInfo) -> bool { - if candidate.inline_data() - || candidate.is_compressed() - || candidate.is_remote() - || candidate - .metadata - .keys() - .any(|key| rustfs_utils::http::is_object_encryption_marker(key)) - || candidate.parts.len() != 1 - { - return false; - } - candidate.has_valid_erasure_geometry() -} - -pub(in crate::set_disk) fn late_materialization_candidate_is_safe(candidate: &FileInfo) -> bool { - non_inline_data_read_candidate_is_safe(candidate) - && candidate.size > 512 * 1024 - && object_fits_single_block(candidate.size, candidate.erasure.block_size) -} - -pub(in crate::set_disk) fn non_inline_data_read_early_stop_allowed(read_data: bool, bucket: &str, object: &str) -> bool { - read_data && is_get_metadata_non_inline_data_read_early_stop_enabled() && !codec_streaming_rollout_applies(bucket, object) -} - -const NON_INLINE_SINGLE_PENDING_HEDGE_DELAY: Duration = Duration::from_millis(100); - -fn data_read_inline_missing_shards_are_pending( - candidate: &FileInfo, - parts_metadata: &[FileInfo], - errors: &[Option], - disks: &[Option], - fanout_order: &[usize], - scheduled_fanout_len: usize, -) -> bool { - 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 distribution = &candidate.erasure.distribution; - let mut data_shards_seen_or_pending = vec![false; erasure.data_shards]; - let mut missing_pending_data_shards = 0usize; - - for (disk_index, file_info) in parts_metadata.iter().enumerate() { - let Some(&block_index) = distribution.get(disk_index) else { - return false; - }; - if block_index == 0 || block_index > erasure.data_shards { - continue; - } - if !disks.get(disk_index).is_some_and(Option::is_some) { - return false; - } - - let data_slot = block_index - 1; - if file_info.name.is_empty() { - let scheduled_and_not_failed = fanout_order - .get(..scheduled_fanout_len) - .is_some_and(|scheduled_disks| scheduled_disks.contains(&disk_index)) - && errors.get(disk_index).is_some_and(Option::is_none); - if scheduled_and_not_failed { - data_shards_seen_or_pending[data_slot] = true; - missing_pending_data_shards = missing_pending_data_shards.saturating_add(1); - continue; - } - return false; - } - if file_info.erasure.index != block_index - || !file_info.has_valid_erasure_geometry() - || !metadata_early_stop_candidate_matches(file_info, candidate) - || file_info.data.as_ref().is_none_or(|data| data.is_empty()) - { - return false; - } - data_shards_seen_or_pending[data_slot] = true; - } - - missing_pending_data_shards > 0 && data_shards_seen_or_pending.into_iter().all(|seen_or_pending| seen_or_pending) -} - -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, - DiskError::FileVersionNotFound => GET_METADATA_RESPONSE_VERSION_NOT_FOUND, - DiskError::DiskNotFound => GET_METADATA_RESPONSE_DISK_NOT_FOUND, - DiskError::FileCorrupt | DiskError::CorruptedFormat | DiskError::CorruptedBackend | DiskError::OutdatedXLMeta => { - GET_METADATA_RESPONSE_CORRUPT - } - DiskError::Timeout => GET_METADATA_RESPONSE_TIMEOUT, - DiskError::FaultyDisk | DiskError::FaultyRemoteDisk => GET_METADATA_RESPONSE_IGNORED, - _ => GET_METADATA_RESPONSE_ERROR, - } -} - pub(in crate::set_disk) fn is_confirmed_missing_part_error(err: Option<&str>) -> bool { let Some(err) = err else { return false; @@ -2451,37 +1791,6 @@ where // metadata cache) stays in read.rs and reaches these through the SetDisks core. // =========================================================================== -pub(in crate::set_disk) fn should_allow_metadata_early_stop( - read_data: bool, - version_id: &str, - healing: bool, - incl_free_versions: bool, -) -> bool { - if read_data && !is_get_metadata_data_read_early_stop_enabled() { - return false; - } - - (is_get_metadata_early_stop_enabled() && version_id.is_empty() && !healing && !incl_free_versions) - || (is_version_early_stop_enabled() && !version_id.is_empty() && !healing && !incl_free_versions) -} - -/// Final gate for the metadata early-stop fast path. -/// -/// `caller_allows_early_stop=false` unconditionally forces the full quorum -/// fanout so read-before-write callers (object tagging) get the complete -/// online-disk set as their write target; the early-stop subset would only -/// carry read quorum and fail write quorum (backlog#872 regression). -pub(in crate::set_disk) fn metadata_early_stop_permitted( - caller_allows_early_stop: bool, - observe: bool, - read_data: bool, - version_id: &str, - healing: bool, - incl_free_versions: bool, -) -> bool { - caller_allows_early_stop && observe && should_allow_metadata_early_stop(read_data, version_id, healing, incl_free_versions) -} - impl SetDisks { pub(in crate::set_disk) async fn read_parts( disks: &[Option], @@ -2548,7 +1857,7 @@ impl SetDisks { healing: bool, incl_free_versions: bool, ) -> disk::error::Result<(Vec, Vec>)> { - let (ress, errors, _) = Self::read_all_fileinfo_inner( + let result = Self::read_metadata_inner( disks, org_bucket, bucket, @@ -2563,9 +1872,11 @@ impl SetDisks { false, ) .await?; - Ok((ress, errors)) + let (metadata, errors, _) = result.into_legacy(); + Ok((metadata, errors)) } + #[cfg(test)] #[allow(clippy::too_many_arguments)] pub(in crate::set_disk) async fn read_all_fileinfo_observed( disks: &[Option], @@ -2579,7 +1890,7 @@ impl SetDisks { caller_allows_early_stop: bool, default_parity_count: usize, ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { - Self::read_all_fileinfo_inner( + Self::read_metadata_inner( disks, org_bucket, bucket, @@ -2594,8 +1905,10 @@ impl SetDisks { false, ) .await + .map(MetadataReadResult::into_legacy) } + #[cfg(test)] #[allow(clippy::too_many_arguments)] pub(in crate::set_disk) async fn read_all_fileinfo_observed_for_get_object( disks: &[Option], @@ -2608,7 +1921,7 @@ impl SetDisks { caller_allows_early_stop: bool, default_parity_count: usize, ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { - Self::read_all_fileinfo_inner( + Self::read_metadata_inner( disks, org_bucket, bucket, @@ -2623,497 +1936,7 @@ impl SetDisks { true, ) .await - } - - #[allow(clippy::too_many_arguments)] - async fn read_all_fileinfo_inner( - disks: &[Option], - org_bucket: &str, - bucket: &str, - object: &str, - version_id: &str, - read_data: bool, - healing: bool, - incl_free_versions: bool, - observe: bool, - // When false, the caller opts out of the early-stop fast path even for - // otherwise-eligible reads. Read-before-write callers (e.g. object - // tagging) must set this so the returned online-disk set reflects the - // full quorum fanout rather than the early-stop subset — writing to the - // subset would fail write quorum (backlog#872 regression). - caller_allows_early_stop: bool, - default_parity_count: usize, - allow_coalescing: bool, - ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { - let early_stop_enabled = - caller_allows_early_stop && observe && (is_get_metadata_early_stop_enabled() || is_version_early_stop_enabled()); - let allow_early_stop = - metadata_early_stop_permitted(caller_allows_early_stop, observe, read_data, version_id, healing, incl_free_versions); - if allow_early_stop { - return Self::read_all_fileinfo_early_stop( - disks, - org_bucket, - bucket, - object, - version_id, - read_data, - healing, - incl_free_versions, - non_inline_data_read_early_stop_allowed(read_data, bucket, object), - default_parity_count, - allow_coalescing, - ) - .await; - } - if early_stop_enabled { - let metrics_path = metadata_metrics_path(bucket); - rustfs_io_metrics::record_get_object_metadata_early_stop_miss( - metrics_path, - GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, - ); - rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(metrics_path, 0); - } - - Self::read_all_fileinfo_full_wait( - disks, - org_bucket, - bucket, - object, - version_id, - read_data, - healing, - incl_free_versions, - observe, - allow_coalescing, - ) - .await - } - - #[allow(clippy::too_many_arguments)] - async fn read_all_fileinfo_full_wait( - disks: &[Option], - org_bucket: &str, - bucket: &str, - object: &str, - version_id: &str, - read_data: bool, - healing: bool, - incl_free_versions: bool, - observe: bool, - allow_coalescing: bool, - ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { - let fanout_start = observe.then(Instant::now); - 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, - healing, - }; - let org_bucket: Arc = Arc::from(org_bucket); - let bucket: Arc = Arc::from(bucket); - let object: Arc = Arc::from(object); - let version_id: Arc = Arc::from(version_id); - let slowtail_fault = get_metadata_slowtail_fault_request(bucket.as_ref(), object.as_ref(), read_data); - let futures = disks.iter().enumerate().map(|(disk_index, disk)| { - let disk = disk.clone(); - let task_opts = opts; - let org_bucket = org_bucket.clone(); - let bucket = bucket.clone(); - let object = object.clone(); - let version_id = version_id.clone(); - let slowtail_fault = slowtail_fault.clone(); - AbortOnDropJoinHandle(tokio::spawn(async move { - let response_start = observe.then(Instant::now); - let result = if let Some(disk) = disk { - Self::record_read_version_call(&object, disk_index); - if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(disk_index)) { - Self::record_metadata_slowtail_fault(&object, disk_index); - tokio::time::sleep(delay).await; - } - read_version_via_coalescer(disk, &org_bucket, &bucket, &object, &version_id, &task_opts, allow_coalescing) - .await - } else { - Err(DiskError::DiskNotFound) - }; - let elapsed = response_start.map(|start| start.elapsed()); - (result, elapsed) - })) - }); - - // Wait for all futures to complete - let results = join_all(futures).await; - - for join_result in results { - match join_result { - Ok((res, elapsed)) => match res { - Ok(file_info) => { - if let (Some(observations), Some(elapsed)) = (&mut observations, elapsed) { - observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed)); - } - ress.push(file_info); - errors.push(None); - } - Err(e) => { - if let (Some(observations), Some(elapsed)) = (&mut observations, elapsed) { - observations.push(MetadataFanoutObservation::from_error(&e, elapsed)); - } - ress.push(FileInfo::default()); - errors.push(Some(e)); - } - }, - Err(_join_err) => { - // A spawned task panicked — treat as unexpected disk error - if let Some(observations) = &mut observations { - observations.push(MetadataFanoutObservation::from_error(&DiskError::Unexpected, Duration::ZERO)); - } - ress.push(FileInfo::default()); - errors.push(Some(DiskError::Unexpected)); - } - } - } - let diagnostics = match (fanout_start, observations) { - (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)) - } - - #[allow(clippy::too_many_arguments)] - async fn read_all_fileinfo_early_stop( - disks: &[Option], - org_bucket: &str, - bucket: &str, - object: &str, - version_id: &str, - read_data: bool, - healing: bool, - incl_free_versions: bool, - allow_non_inline_data_read_early_stop: bool, - default_parity_count: usize, - allow_coalescing: bool, - ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { - let fanout_start = Instant::now(); - let mut ress = vec![FileInfo::default(); disks.len()]; - let mut errors = vec![None; disks.len()]; - let mut observations = Vec::with_capacity(disks.len()); - let mut accumulator = - MetadataQuorumAccumulator::new(disks.len(), default_parity_count, true).with_requested_version_id(version_id); - let opts = ReadOptions { - incl_free_versions, - read_data, - healing, - }; - let org_bucket: Arc = Arc::from(org_bucket); - 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 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 mut final_miss_reason_override = None; - let mut non_inline_candidate_eligible = None; - let mut single_pending_hedge_deadline = None; - let slowtail_fault = get_metadata_slowtail_fault_request(bucket.as_ref(), object.as_ref(), read_data); - let spawn_read_version = - |join_set: &mut JoinSet<(usize, disk::error::Result, Duration)>, index: usize, disk: Option| { - let task_opts = opts; - let org_bucket = org_bucket.clone(); - let bucket = bucket.clone(); - let object = object.clone(); - let version_id = version_id.clone(); - let slowtail_fault = slowtail_fault.clone(); - join_set.spawn(async move { - let response_start = Instant::now(); - let result = if let Some(disk) = disk { - #[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; - if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(index)) { - Self::record_metadata_slowtail_fault(&object, index); - tokio::time::sleep(delay).await; - } - read_version_via_coalescer(disk, &org_bucket, &bucket, &object, &version_id, &task_opts, allow_coalescing) - .await - } else { - Err(DiskError::DiskNotFound) - }; - (index, result, response_start.elapsed()) - }); - }; - - if bounded_fanout { - let initial_target = accumulator.default_write_quorum().min(disks.len()); - 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_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); - } - } - - loop { - let mut defer_pending_inline_data_shard = false; - let result = if let Some(deadline) = single_pending_hedge_deadline.take() { - tokio::select! { - result = join_set.join_next() => result, - _ = async { - tokio::time::sleep_until(deadline).await; - #[cfg(test)] - rename_fanout_barrier::checkpoint( - object.as_ref(), - 0, - rename_fanout_barrier::PHASE_NON_INLINE_HEDGE_TIMER, - ) - .await; - } => { - if bounded_fanout - && !force_full_wait - && join_set.len() == 1 - && non_inline_candidate_eligible == Some(true) - && !accumulator.candidate_has_read_reserve() - && next_fanout_index < disks.len() - { - while next_fanout_index < disks.len() { - let disk_index = fanout_order[next_fanout_index]; - next_fanout_index = next_fanout_index.saturating_add(1); - 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); - break; - } - } - } - continue; - } - } - } else { - join_set.join_next().await - }; - let Some(result) = result else { break }; - match result { - Ok((index, res, elapsed)) => match res { - Ok(file_info) => { - observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed)); - if allow_non_inline_data_read_early_stop { - accumulator.observe_file_info_at(index, &file_info); - } else { - accumulator.observe_file_info(&file_info); - } - if allow_non_inline_data_read_early_stop && non_inline_candidate_eligible.is_none() { - non_inline_candidate_eligible = - accumulator.candidate.as_ref().map(non_inline_data_read_candidate_is_safe); - } - if bounded_fanout - && read_data - && !force_full_wait - && let Some(reason) = data_read_early_stop_inline_candidate_miss_reason(&file_info) - && !(non_inline_candidate_eligible == Some(true) - && reason == GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE) - { - force_full_wait = true; - final_miss_reason_override.get_or_insert(reason); - } - if let Some(slot) = ress.get_mut(index) { - *slot = file_info; - } - } - Err(err) => { - observations.push(MetadataFanoutObservation::from_error(&err, elapsed)); - accumulator.observe_error(&err); - if let Some(slot) = errors.get_mut(index) { - *slot = Some(err); - } - } - }, - Err(_) => { - let err = DiskError::Unexpected; - observations.push(MetadataFanoutObservation::from_error(&err, fanout_start.elapsed())); - accumulator.observe_error(&err); - } - } - - if !force_full_wait - && let Some(decision) = accumulator - .early_stop_decision() - .or_else(|| accumulator.version_early_stop_decision()) - { - let should_return_early = if read_data { - match accumulator.candidate.as_ref() { - Some(_candidate) if non_inline_candidate_eligible == Some(true) => { - accumulator.candidate_has_read_reserve() - } - Some(candidate) => match data_read_early_stop_inline_body_miss_reason( - bucket.as_ref(), - object.as_ref(), - candidate, - &ress, - disks, - ) - .await - { - None => true, - Some(reason) => { - final_miss_reason_override = Some(reason); - if bounded_fanout - && reason == GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD - && data_read_inline_missing_shards_are_pending( - candidate, - &ress, - &errors, - disks, - &fanout_order, - next_fanout_index, - ) - { - defer_pending_inline_data_shard = true; - } else { - force_full_wait = true; - } - false - } - }, - None => { - force_full_wait = true; - final_miss_reason_override = Some(GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM); - false - } - } - } else { - true - }; - - 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(); - // Inline verification can still depend on a missing data shard; - // issue one immediate spare when only that shard remains. The - // non-inline path keeps its delayed hedge below to avoid healthy - // reads paying speculative I/O before the candidate is classified. - let should_hedge_single_pending_inline_read = read_data - && !force_full_wait - && !defer_pending_inline_data_shard - && pending_responses == 1 - && non_inline_candidate_eligible != Some(true) - && accumulator.can_still_reach_early_stop_with_pending(pending_responses); - // A non-inline plan must retain one extra matching shard as a - // reconstruction reserve. Schedule that reserve only after the - // candidate is known to be eligible, so inline GETs do not pay an - // extra fanout and the healthy path remains allocation-free. - let needs_non_inline_read_reserve = non_inline_candidate_eligible == Some(true) - && !accumulator.candidate_has_read_reserve() - && accumulator - .candidate_read_reserve_target() - .is_some_and(|reserve_target| scheduled_count < reserve_target || pending_responses == 0); - if bounded_fanout - && !force_full_wait - && (needs_non_inline_read_reserve || should_hedge_single_pending_inline_read) - && 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 && 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 - && !defer_pending_inline_data_shard - && next_fanout_index < disks.len() - && !accumulator.can_still_reach_early_stop_with_pending(pending_responses) - { - 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); - } - if bounded_fanout - && !force_full_wait - && !defer_pending_inline_data_shard - && join_set.len() == 1 - && non_inline_candidate_eligible == Some(true) - && !accumulator.candidate_has_read_reserve() - && accumulator.can_still_reach_early_stop_with_pending(join_set.len()) - && next_fanout_index < disks.len() - { - single_pending_hedge_deadline = Some(tokio::time::Instant::now() + NON_INLINE_SINGLE_PENDING_HEDGE_DELAY); - } - } - - let accumulator_miss_reason = accumulator.final_miss_reason(); - let final_miss_reason = match (final_miss_reason_override, accumulator_miss_reason) { - (Some(reason), GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM) => reason, - _ => accumulator_miss_reason, - }; - rustfs_io_metrics::record_get_object_metadata_early_stop_miss(metrics_path, 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)) + .map(MetadataReadResult::into_legacy) } pub(in crate::set_disk) async fn read_all_xl( @@ -5330,37 +4153,6 @@ impl SetDisks { None } - /// Test-only seam that records one per-disk `read_version` metadata RPC for - /// the call-counter registry (backlog#1325). In production this is inlined to - /// nothing and adds no behavior; only the `#[cfg(test)]` variant touches the - /// registry. Placed inside the fan-out spawn tasks so counts are observed - /// even though the increments run on arbitrary runtime worker threads. - #[cfg(test)] - #[inline] - fn record_read_version_call(object: &str, disk_index: usize) { - disk_call_counters::record(object, disk_call_counters::KIND_READ_VERSION, disk_index); - } - - #[cfg(not(test))] - #[inline(always)] - fn record_read_version_call(_object: &str, _disk_index: usize) {} - - #[cfg(test)] - #[inline] - fn record_metadata_slowtail_fault(object: &str, disk_index: usize) { - disk_call_counters::record(object, disk_call_counters::KIND_METADATA_SLOWTAIL_FAULT, disk_index); - } - - #[cfg(not(test))] - #[inline(always)] - fn record_metadata_slowtail_fault(_object: &str, _disk_index: usize) {} - - #[cfg(test)] - #[inline] - async fn read_version_fanout_barrier(object: &str, disk_index: usize) { - rename_fanout_barrier::checkpoint(object, disk_index, rename_fanout_barrier::PHASE_READ_VERSION).await; - } - /// Test-only awaitable pause point for the rename/commit fan-out (backlog#1325, /// serving the barrier-style acceptances of #1312 / #1319 / #1313). `phase` is /// [`rename_fanout_barrier::PHASE_RENAME`] or `PHASE_CLEANUP`. When a test has @@ -7263,7 +6055,7 @@ pub(crate) mod disk_call_counters { /// Record one call of `kind` against `object`'s `disk_index`. A no-op unless /// an observing scope for `object` is currently installed. - pub(super) fn record(object: &str, kind: &str, disk_index: usize) { + pub(in crate::set_disk) fn record(object: &str, kind: &str, disk_index: usize) { let mut reg = registry().lock().expect("disk call-counter registry poisoned"); if !reg.observed.contains(object) { return; diff --git a/crates/ecstore/src/set_disk/core/metadata_quorum.rs b/crates/ecstore/src/set_disk/core/metadata_quorum.rs index 2922b6fc6..1b97adec8 100644 --- a/crates/ecstore/src/set_disk/core/metadata_quorum.rs +++ b/crates/ecstore/src/set_disk/core/metadata_quorum.rs @@ -29,6 +29,27 @@ use crate::disk::error_reduce::OBJECT_OP_IGNORED_ERRS; use crate::set_disk::file_info_is_valid_for_metadata; use rustfs_filemeta::FileInfo; +/// One disk's observation. Pending is neither a successful vote nor an offline +/// disk: the scheduler may not have issued this slot or may still be waiting. +// Option> keeps FileInfo inline without a per-response Box. +// None is pending; Some(Ok(_)) is success; Some(Err(_)) retains a disk error. +pub(in crate::set_disk) type MetadataDiskResult = Option>; + +#[derive(Debug)] +pub(in crate::set_disk) struct MetadataDiskObservation { + pub(in crate::set_disk) disk_index: usize, + pub(in crate::set_disk) result: MetadataDiskResult, +} + +impl MetadataDiskObservation { + pub(in crate::set_disk) fn file_info(&self) -> Option<&FileInfo> { + match &self.result { + Some(Ok(metadata)) => Some(metadata), + None | Some(Err(_)) => None, + } + } +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub(in crate::set_disk) struct MetadataEarlyStopDecision { pub(in crate::set_disk) reason: &'static str, @@ -59,6 +80,14 @@ pub(in crate::set_disk) struct MetadataQuorumAccumulator { } impl MetadataQuorumAccumulator { + pub(in crate::set_disk) fn observe(&mut self, observation: &MetadataDiskObservation) { + match &observation.result { + None => {} + Some(Ok(metadata)) => self.observe_file_info_at(observation.disk_index, metadata), + Some(Err(error)) => self.observe_error(error), + } + } + pub(in crate::set_disk) fn new(total_disks: usize, default_parity_count: usize, allow_early_stop: bool) -> Self { Self { total_disks, @@ -86,6 +115,7 @@ impl MetadataQuorumAccumulator { self } + #[cfg(test)] pub(in crate::set_disk) fn observe_file_info(&mut self, file_info: &FileInfo) { self.observe_file_info_with_index(None, file_info); } @@ -383,3 +413,186 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo pub(in crate::set_disk) fn is_metadata_fanout_ignored_error(err: &DiskError) -> bool { OBJECT_OP_IGNORED_ERRS.iter().any(|ignored| ignored == err) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::set_disk::SetDisks; + use time::OffsetDateTime; + use uuid::Uuid; + + fn payload(disk_index: usize) -> FileInfo { + let mut metadata = FileInfo::new("object", 2, 2); + metadata.volume = "bucket".to_string(); + metadata.name = "object".to_string(); + metadata.size = 1; + metadata.mod_time = Some(OffsetDateTime::UNIX_EPOCH); + metadata.data_dir = Some(Uuid::from_u128(1)); + metadata.erasure.index = metadata.erasure.distribution[disk_index]; + metadata.metadata.insert("etag".to_string(), "etag".to_string()); + metadata.add_object_part(1, "etag".to_string(), 1, None, 1, None, None); + assert!(file_info_is_valid_for_metadata(&metadata), "the observation corpus needs a valid payload"); + metadata + } + + fn observation(disk_index: usize, state: usize) -> MetadataDiskObservation { + MetadataDiskObservation { + disk_index, + result: match state { + 0 => Some(Ok(payload(disk_index))), + 1 => Some(Err(DiskError::FileNotFound)), + 2 => Some(Err(DiskError::FileCorrupt)), + 3 => Some(Err(DiskError::DiskNotFound)), + 4 => None, + _ => panic!("unexpected test state"), + }, + } + } + + fn assert_same_reduction(typed: &MetadataQuorumAccumulator, legacy: &MetadataQuorumAccumulator) { + assert_eq!(typed.early_stop_decision(), legacy.early_stop_decision()); + assert_eq!(typed.version_early_stop_decision(), legacy.version_early_stop_decision()); + assert_eq!(typed.final_miss_reason(), legacy.final_miss_reason()); + assert_eq!(typed.valid_responses, legacy.valid_responses); + assert_eq!(typed.not_found_responses, legacy.not_found_responses); + assert_eq!(typed.version_not_found_responses, legacy.version_not_found_responses); + assert_eq!(typed.ignored_errors, legacy.ignored_errors); + assert_eq!(typed.hard_errors, legacy.hard_errors); + assert_eq!(typed.candidate_votes, legacy.candidate_votes); + assert_eq!(typed.candidate_shard_mask, legacy.candidate_shard_mask); + assert_eq!(typed.matching_version_votes, legacy.matching_version_votes); + assert_eq!(typed.delete_marker_votes, legacy.delete_marker_votes); + assert_eq!(typed.conflicting_metadata, legacy.conflicting_metadata); + assert_eq!( + typed.candidate.as_ref().map(SetDisks::file_info_quorum_hash), + legacy.candidate.as_ref().map(SetDisks::file_info_quorum_hash) + ); + } + + fn observe_legacy(accumulator: &mut MetadataQuorumAccumulator, observation: &MetadataDiskObservation) { + match &observation.result { + None => {} + Some(Ok(metadata)) => accumulator.observe_file_info_at(observation.disk_index, metadata), + Some(Err(error)) => accumulator.observe_error(error), + } + } + + #[test] + fn metadata_observation_all_four_slot_states_and_arrival_orders_preserve_reduction() { + let permutations = (0..4) + .flat_map(|a| (0..4).filter(move |b| *b != a).map(move |b| (a, b))) + .flat_map(|(a, b)| (0..4).filter(move |c| *c != a && *c != b).map(move |c| (a, b, c))) + .map(|(a, b, c)| [a, b, c, 6 - a - b - c]) + .collect::>(); + assert_eq!(permutations.len(), 24); + // 5^4 states and all 4! arrivals, both with early-stop enabled and disabled. + for encoded in 0..625usize { + let states = [encoded % 5, encoded / 5 % 5, encoded / 25 % 5, encoded / 125 % 5]; + let inputs = std::array::from_fn::<_, 4, _>(|index| observation(index, states[index])); + for enabled in [false, true] { + for order in &permutations { + let mut typed = MetadataQuorumAccumulator::new(4, 2, enabled); + let mut legacy = MetadataQuorumAccumulator::new(4, 2, enabled); + for &index in order { + typed.observe(&inputs[index]); + observe_legacy(&mut legacy, &inputs[index]); + assert_same_reduction(&typed, &legacy); + } + assert_eq!(typed.valid_responses, states.iter().filter(|&&state| state == 0).count()); + assert_eq!(typed.not_found_responses, states.iter().filter(|&&state| state == 1).count()); + assert_eq!(typed.hard_errors, states.iter().filter(|&&state| state == 2).count()); + assert_eq!(typed.ignored_errors, states.iter().filter(|&&state| state == 3).count()); + let expected_early_stop = + enabled && typed.valid_responses >= 3 && !states.iter().any(|&state| state == 1 || state == 2); + assert_eq!( + typed.early_stop_decision().is_some(), + expected_early_stop, + "states={states:?}, order={order:?}" + ); + } + } + } + } + + #[test] + fn metadata_observation_pending_newer_version_and_same_time_different_directory_force_full_wait() { + let mut accumulator = MetadataQuorumAccumulator::new(4, 2, true); + for index in 0..2 { + accumulator.observe(&observation(index, 0)); + } + accumulator.observe(&observation(2, 4)); + assert_eq!(accumulator.candidate_votes, 2, "an outstanding response cannot supply the deciding vote"); + assert_eq!(accumulator.early_stop_decision(), None); + + let mut newer = payload(2); + newer.mod_time = Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(1)); + accumulator.observe(&MetadataDiskObservation { + disk_index: 2, + result: Some(Ok(newer)), + }); + assert!(accumulator.conflicting_metadata); + assert_eq!(accumulator.early_stop_decision(), None); + + let mut same_time_different_directory = payload(3); + same_time_different_directory.data_dir = Some(Uuid::from_u128(2)); + accumulator.observe(&MetadataDiskObservation { + disk_index: 3, + result: Some(Ok(same_time_different_directory)), + }); + assert_eq!(accumulator.early_stop_decision(), None); + } + + #[test] + fn metadata_observation_duplicate_shard_cannot_supply_a_read_reserve() { + let mut accumulator = MetadataQuorumAccumulator::new(4, 2, true); + let first = payload(0); + for disk_index in 0..3 { + accumulator.observe(&MetadataDiskObservation { + disk_index, + result: Some(Ok(first.clone())), + }); + } + assert_eq!(accumulator.candidate_votes, 3); + assert!( + !accumulator.candidate_has_read_reserve(), + "copied erasure indexes cannot supply independent data shards" + ); + } + + #[test] + fn metadata_observation_null_versions_markers_and_invalid_success_keep_their_meaning() { + let mut null_version = MetadataQuorumAccumulator::new(4, 2, true); + for index in 0..3 { + null_version.observe(&observation(index, 0)); + } + assert!(null_version.early_stop_decision().is_some()); + assert_eq!(null_version.matching_version_votes, 0); + + let mut marker = MetadataQuorumAccumulator::new(4, 2, true); + for disk_index in 0..3 { + let metadata = FileInfo { + volume: "bucket".to_string(), + name: "object".to_string(), + deleted: true, + mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(1)), + ..Default::default() + }; + assert!(metadata.is_canonical_delete_marker(), "the null marker fixture must reach marker voting"); + marker.observe(&MetadataDiskObservation { + disk_index, + result: Some(Ok(metadata)), + }); + } + assert_eq!(marker.delete_marker_votes, 3); + assert!(marker.early_stop_decision().is_some()); + + let mut invalid = MetadataQuorumAccumulator::new(4, 2, true); + invalid.observe(&MetadataDiskObservation { + disk_index: 0, + result: Some(Ok(FileInfo::default())), + }); + assert_eq!(invalid.valid_responses, 0); + assert_eq!(invalid.hard_errors, 1); + assert_eq!(invalid.early_stop_decision(), None); + } +} diff --git a/crates/ecstore/src/set_disk/core/metadata_read.rs b/crates/ecstore/src/set_disk/core/metadata_read.rs new file mode 100644 index 000000000..82dd3b8a4 --- /dev/null +++ b/crates/ecstore/src/set_disk/core/metadata_read.rs @@ -0,0 +1,1443 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Metadata read scheduling, cancellation and coalescing owned by `SetDisks`. +//! Quorum decisions consume slot observations in `metadata_quorum`. + +#[cfg(test)] +use super::io_primitives::{disk_call_counters, rename_fanout_barrier}; +use super::metadata_quorum::{ + MetadataDiskObservation, MetadataQuorumAccumulator, is_metadata_fanout_ignored_error, metadata_early_stop_candidate_matches, +}; +use crate::diagnostics::get::{ + GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED, + GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD, + GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE, + GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE, GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE, + GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED, GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM, + GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, 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_INTERNAL_META, + GET_OBJECT_PATH_LEGACY_DUPLEX, +}; +use crate::disk::disk_store::get_drive_metadata_timeout; +use crate::disk::{ + self, BATCH_READ_VERSION_MAX_ITEMS, BatchReadVersionItem, BatchReadVersionReq, BatchReadVersionResp, Disk, DiskAPI, +}; +use crate::set_disk::{ + DiskError, DiskStore, FileInfo, HashAlgorithm, ReadOptions, SetDisks, build_inline_bitrot_readers_from_refs, + can_try_inline_data_shards_direct, codec_streaming_rollout_applies, coding, + collect_inline_data_shard_fileinfos_from_observations, file_info_is_valid_for_metadata, get_metadata_slowtail_fault_request, + inline_erasure_shard_file_offset, inline_erasure_shard_size, is_get_metadata_data_read_early_stop_enabled, + is_get_metadata_early_stop_bounded_fanout_enabled, is_get_metadata_early_stop_enabled, + is_get_metadata_non_inline_data_read_early_stop_enabled, is_version_early_stop_enabled, object_fits_single_block, + try_read_inline_data_shards_direct, +}; +use futures::future::join_all; +use metrics::counter; +use rustfs_io_metrics::internode_metrics::{ + INTERNODE_STAGE_BATCH_READ_VERSION_COALESCER_WAIT, INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_MAP, +}; +#[cfg(test)] +use std::collections::HashSet; +use std::{ + collections::HashMap, + future::Future, + pin::Pin, + sync::{Arc, OnceLock}, + task::{Context, Poll}, + time::{Duration, Instant}, +}; +use tokio::sync::{Mutex, oneshot}; +use tokio::task::JoinSet; + +#[derive(Debug)] +pub(in crate::set_disk) struct MetadataReadResult { + pub(in crate::set_disk) slots: Vec, + pub(in crate::set_disk) diagnostics: MetadataFanoutDiagnostics, +} + +impl MetadataReadResult { + fn pending(total_disks: usize) -> Vec { + (0..total_disks) + .map(|disk_index| MetadataDiskObservation { + disk_index, + result: None, + }) + .collect() + } + + pub(in crate::set_disk) fn is_complete(&self) -> bool { + self.slots.iter().all(|slot| slot.result.is_some()) + } + + /// Legacy quorum/layout helpers still require aligned slices. Only this + /// boundary produces their empty placeholders; the scheduler and reducer + /// never infer a response from a default `FileInfo` or a missing error. + // RUSTFS_COMPAT_TODO(backlog-2253-metadata-observation-adapter): Existing layout, shard and exact-version writer/delete callers still consume aligned slices. Remove after those consumers use typed observations in a separate cleanup. + pub(in crate::set_disk) fn into_legacy(self) -> (Vec, Vec>, MetadataFanoutDiagnostics) { + let mut metadata = Vec::with_capacity(self.slots.len()); + let mut errors = Vec::with_capacity(self.slots.len()); + for slot in self.slots { + match slot.result { + None => { + metadata.push(FileInfo::default()); + errors.push(None); + } + Some(Ok(file_info)) => { + metadata.push(file_info); + errors.push(None); + } + Some(Err(error)) => { + metadata.push(FileInfo::default()); + errors.push(Some(error)); + } + } + } + (metadata, errors, self.diagnostics) + } +} + +pub(super) const ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE: &str = "RUSTFS_GET_METADATA_READ_VERSION_COALESCE"; +pub(super) const ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS: &str = + "RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS"; +const DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS: u64 = 200; +const METRIC_GET_METADATA_READ_VERSION_COALESCER_TOTAL: &str = "rustfs_get_metadata_read_version_coalescer_total"; + +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 + } +} + +pub(super) fn metadata_distribution_key(bucket: &str, object: &str) -> String { + [bucket, object].join("/") +} + +fn read_version_coalescing_enabled() -> bool { + let enabled = || { + rustfs_utils::get_env_opt_str(ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE) + .is_some_and(|value| value.eq_ignore_ascii_case("auto") || value.eq_ignore_ascii_case("on")) + }; + + #[cfg(test)] + { + enabled() + } + + #[cfg(not(test))] + { + static ENABLED: OnceLock = OnceLock::new(); + *ENABLED.get_or_init(enabled) + } +} + +fn read_version_coalescing_delay() -> Duration { + #[cfg(test)] + { + let micros = rustfs_utils::get_env_u64( + ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, + DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, + ); + Duration::from_micros(micros) + } + + #[cfg(not(test))] + { + static DELAY: OnceLock = OnceLock::new(); + *DELAY.get_or_init(|| { + Duration::from_micros(rustfs_utils::get_env_u64( + ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, + DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, + )) + }) + } +} + +struct CoalescedReadVersionRequest { + item: BatchReadVersionItem, + tx: oneshot::Sender>, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub(super) struct ExpectedBatchReadVersionItem { + pub(super) path: String, + pub(super) version_id: String, +} + +impl From<&BatchReadVersionItem> for ExpectedBatchReadVersionItem { + fn from(item: &BatchReadVersionItem) -> Self { + Self { + path: item.path.clone(), + version_id: item.version_id.clone(), + } + } +} + +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +struct ReadVersionCoalescerKey { + disk: usize, + incl_free_versions: bool, + read_data: bool, + healing: bool, +} + +impl ReadVersionCoalescerKey { + fn new(disk: &DiskStore, opts: &ReadOptions) -> Self { + Self { + disk: Arc::as_ptr(disk) as usize, + incl_free_versions: opts.incl_free_versions, + read_data: opts.read_data, + healing: opts.healing, + } + } +} + +#[derive(Default)] +struct ReadVersionCoalescer { + lanes: HashMap>, +} + +fn read_version_coalescer() -> &'static Mutex { + static COALESCER: OnceLock> = OnceLock::new(); + COALESCER.get_or_init(|| Mutex::new(ReadVersionCoalescer::default())) +} + +fn record_read_version_coalescer_event(event: &'static str, item_count: usize) { + counter!( + METRIC_GET_METADATA_READ_VERSION_COALESCER_TOTAL, + "event" => event, + "item_count" => item_count.to_string() + ) + .increment(1); +} + +fn batch_read_version_stage_timer() -> Option { + rustfs_io_metrics::get_stage_metrics_enabled().then(Instant::now) +} + +fn record_batch_read_version_stage(stage: &'static str, started_at: Option) { + if let Some(started_at) = started_at { + crate::cluster::rpc::runtime_sources::record_remote_disk_grpc_batch_read_version_stage(stage, started_at.elapsed()); + } +} + +async fn read_version_via_coalescer( + disk: DiskStore, + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + opts: &ReadOptions, + allow_coalescing: bool, +) -> disk::error::Result { + if !allow_coalescing || !read_version_coalescing_enabled() { + return disk.read_version(org_bucket, bucket, object, version_id, opts).await; + } + if !matches!(disk.as_ref(), Disk::Remote(_)) { + record_read_version_coalescer_event("bypass_non_remote", 1); + return disk.read_version(org_bucket, bucket, object, version_id, opts).await; + } + + let (tx, rx) = oneshot::channel(); + let item = BatchReadVersionItem { + org_volume: org_bucket.to_string(), + volume: bucket.to_string(), + path: object.to_string(), + version_id: version_id.to_string(), + }; + let lane_key = ReadVersionCoalescerKey::new(&disk, opts); + let pending = { + let mut coalescer = read_version_coalescer().lock().await; + let lane = coalescer.lanes.entry(lane_key).or_default(); + let schedule_delayed_flush = lane.is_empty(); + lane.push(CoalescedReadVersionRequest { item, tx }); + if lane.len() >= BATCH_READ_VERSION_MAX_ITEMS { + coalescer.lanes.remove(&lane_key) + } else if schedule_delayed_flush { + let disk = disk.clone(); + let task_opts = *opts; + tokio::spawn(async move { + tokio::time::sleep(read_version_coalescing_delay()).await; + flush_read_version_coalescer_lane(lane_key, disk, task_opts).await; + }); + None + } else { + None + } + }; + + if let Some(pending) = pending { + flush_read_version_coalescer_pending(lane_key, disk, *opts, pending).await; + } + + let wait_started = batch_read_version_stage_timer(); + let response = rx + .await + .unwrap_or_else(|_| Err(DiskError::other("coalesced read_version response channel closed"))); + record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_COALESCER_WAIT, wait_started); + response +} + +async fn flush_read_version_coalescer_lane(lane_key: ReadVersionCoalescerKey, disk: DiskStore, opts: ReadOptions) { + let pending = { + let mut coalescer = read_version_coalescer().lock().await; + coalescer.lanes.remove(&lane_key).unwrap_or_default() + }; + flush_read_version_coalescer_pending(lane_key, disk, opts, pending).await; +} + +async fn flush_read_version_coalescer_pending( + lane_key: ReadVersionCoalescerKey, + disk: DiskStore, + opts: ReadOptions, + pending: Vec, +) { + if pending.is_empty() { + return; + } + + // Only the #[cfg(test)] counter-recording block below reads this. + #[cfg(not(test))] + let _ = lane_key; + #[cfg(test)] + { + let mut observed_paths = HashSet::new(); + for request in &pending { + if observed_paths.insert(request.item.path.as_str()) { + disk_call_counters::record(&request.item.path, disk_call_counters::KIND_BATCH_READ_VERSION, lane_key.disk); + } + } + } + + let mut senders = Vec::with_capacity(pending.len()); + let mut items = Vec::with_capacity(pending.len()); + for request in pending { + senders.push(request.tx); + items.push(request.item); + } + + let expected_items = items.iter().map(ExpectedBatchReadVersionItem::from).collect::>(); + record_read_version_coalescer_event("attempted_batch", items.len()); + let result = + match tokio::time::timeout(get_drive_metadata_timeout(), disk.batch_read_version(BatchReadVersionReq { items, opts })) + .await + { + Ok(result) => result, + Err(_) => Err(DiskError::Timeout), + }; + match result { + Ok(responses) => { + let map_started = batch_read_version_stage_timer(); + let results = map_batch_read_version_responses(&expected_items, responses); + record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_MAP, map_started); + for (tx, result) in senders.into_iter().zip(results) { + let _ = tx.send(result); + } + } + Err(err) => { + let message = err.to_string(); + for tx in senders { + let _ = tx.send(Err(DiskError::other(message.clone()))); + } + } + } +} + +pub(super) fn map_batch_read_version_responses( + expected_items: &[ExpectedBatchReadVersionItem], + responses: Vec, +) -> Vec> { + let mut results = (0..expected_items.len()) + .map(|_| Err(DiskError::other("coalesced read_version response missing"))) + .collect::>(); + let mut seen = vec![false; expected_items.len()]; + for response in responses { + let Some(expected) = expected_items.get(response.index) else { + continue; + }; + let Some(slot) = results.get_mut(response.index) else { + continue; + }; + if seen[response.index] { + *slot = Err(DiskError::other("coalesced read_version response duplicate index")); + continue; + } + seen[response.index] = true; + if response.path != expected.path || response.version_id != expected.version_id { + *slot = Err(DiskError::other("coalesced read_version response identity mismatch")); + } else { + *slot = if response.success { + Ok(response.file_info) + } else { + Err(batch_read_version_response_error(response.error_code, response.error)) + }; + } + } + results +} + +fn batch_read_version_response_error(error_code: u32, error: String) -> DiskError { + match DiskError::from_u32(error_code) { + Some(DiskError::Io(_)) | None => DiskError::other(error), + Some(error) => error, + } +} + +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 +} + +struct AbortOnDropJoinHandle(tokio::task::JoinHandle); + +impl Future for AbortOnDropJoinHandle { + type Output = std::result::Result; + + fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + Pin::new(&mut self.0).poll(cx) + } +} + +impl Drop for AbortOnDropJoinHandle { + fn drop(&mut self) { + self.0.abort(); + } +} + +#[derive(Clone, Copy, Debug)] +pub(in crate::set_disk) struct MetadataFanoutObservation { + pub(in crate::set_disk) outcome: &'static str, + pub(in crate::set_disk) elapsed: Duration, + pub(in crate::set_disk) valid: bool, + pub(in crate::set_disk) ignored: bool, +} + +impl MetadataFanoutObservation { + pub(in crate::set_disk) fn from_file_info(file_info: &FileInfo, elapsed: Duration) -> Self { + if file_info_is_valid_for_metadata(file_info) { + Self { + outcome: GET_METADATA_RESPONSE_VALID, + elapsed, + valid: true, + ignored: false, + } + } else { + Self { + outcome: GET_METADATA_RESPONSE_ERROR, + elapsed, + valid: false, + ignored: false, + } + } + } + + pub(in crate::set_disk) fn from_error(err: &DiskError, elapsed: Duration) -> Self { + Self { + outcome: classify_metadata_response_error(err), + elapsed, + valid: false, + ignored: is_metadata_fanout_ignored_error(err), + } + } +} + +#[derive(Clone, Debug, Default)] +pub(in crate::set_disk) struct MetadataFanoutDiagnostics { + pub(in crate::set_disk) fanout_duration: Duration, + pub(in crate::set_disk) observations: Vec, +} + +impl MetadataFanoutDiagnostics { + pub(in crate::set_disk) fn new(fanout_duration: Duration, observations: Vec) -> Self { + Self { + fanout_duration, + observations, + } + } + + pub(in crate::set_disk) fn total_responses(&self) -> usize { + self.observations.len() + } + + pub(in crate::set_disk) fn valid_responses(&self) -> usize { + self.observations.iter().filter(|observation| observation.valid).count() + } + + pub(in crate::set_disk) fn ignored_responses(&self) -> usize { + self.observations.iter().filter(|observation| observation.ignored).count() + } + + pub(in crate::set_disk) fn non_valid_responses(&self) -> usize { + self.total_responses().saturating_sub(self.valid_responses()) + } + + pub(in crate::set_disk) fn first_response_latency(&self) -> Option { + self.observations.iter().map(|observation| observation.elapsed).min() + } + + pub(in crate::set_disk) fn first_valid_response_latency(&self) -> Option { + self.observations + .iter() + .filter(|observation| observation.valid) + .map(|observation| observation.elapsed) + .min() + } + + pub(in crate::set_disk) fn slowest_response_latency(&self) -> Option { + self.observations.iter().map(|observation| observation.elapsed).max() + } + + pub(in crate::set_disk) fn quorum_candidate_latency(&self, read_quorum: usize) -> Option { + if read_quorum == 0 { + return Some(Duration::ZERO); + } + + let mut valid_latencies = self + .observations + .iter() + .filter(|observation| observation.valid) + .map(|observation| observation.elapsed) + .collect::>(); + valid_latencies.sort_unstable(); + valid_latencies.get(read_quorum.saturating_sub(1)).copied() + } + + pub(in crate::set_disk) fn record(&self, path: &'static str) { + rustfs_io_metrics::record_get_object_metadata_fanout_duration(path, self.fanout_duration.as_secs_f64()); + if let Some(latency) = self.first_response_latency() { + rustfs_io_metrics::record_get_object_first_metadata_response_latency(path, latency.as_secs_f64()); + } + if let Some(latency) = self.first_valid_response_latency() { + rustfs_io_metrics::record_get_object_first_valid_metadata_response_latency(path, latency.as_secs_f64()); + } + if let Some(latency) = self.slowest_response_latency() { + rustfs_io_metrics::record_get_object_slowest_metadata_response_latency(path, latency.as_secs_f64()); + } + rustfs_io_metrics::record_get_object_metadata_fanout_shape( + path, + self.total_responses(), + self.valid_responses(), + self.ignored_responses(), + self.non_valid_responses(), + ); + for observation in &self.observations { + rustfs_io_metrics::record_get_object_metadata_response(path, observation.outcome); + } + } + + pub(in crate::set_disk) fn record_quorum_candidate_latency(&self, path: &'static str, read_quorum: usize) { + if let Some(latency) = self.quorum_candidate_latency(read_quorum) { + rustfs_io_metrics::record_get_object_quorum_reached_latency(path, latency.as_secs_f64()); + } + } +} + +#[cfg(test)] +pub(in crate::set_disk) async fn data_read_early_stop_inline_body_miss_reason( + bucket: &str, + object: &str, + candidate: &FileInfo, + parts_metadata: &[FileInfo], + disks: &[Option], +) -> Option<&'static str> { + inline_body_miss_reason_from_observations( + bucket, + object, + candidate, + parts_metadata + .iter() + .enumerate() + .map(|(index, metadata)| (index, Some(metadata))), + disks, + ) + .await +} + +async fn inline_body_miss_reason_from_observations<'a>( + bucket: &str, + object: &str, + candidate: &FileInfo, + observations: impl IntoIterator)>, + disks: &[Option], +) -> Option<&'static str> { + if let Some(reason) = data_read_early_stop_inline_candidate_miss_reason(candidate) { + return Some(reason); + } + + 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 Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY); + }; + let data_files = + match collect_inline_data_shard_fileinfos_from_observations(observations, candidate, erasure.data_shards, |index| { + disks.get(index).is_some_and(Option::is_some) + }) { + Ok(data_files) => data_files, + Err(reason) => return Some(reason), + }; + + let Some(part) = candidate.parts.first() else { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE); + }; + let Ok(object_size) = usize::try_from(candidate.size) else { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); + }; + 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 Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY); + }; + + match try_read_inline_data_shards_direct(&mut readers, erasure.data_shards, read_length, object_size).await { + Some(body) if body.len() == object_size => None, + _ => Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_BODY_VERIFY), + } +} + +fn data_read_early_stop_inline_candidate_miss_reason(candidate: &FileInfo) -> Option<&'static str> { + // `inline_data` excludes remote objects; this diagnostic reports them separately. + if !rustfs_utils::http::contains_key_str(&candidate.metadata, rustfs_utils::http::SUFFIX_INLINE_DATA) { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE); + } + if candidate.is_compressed() + || candidate + .metadata + .keys() + .any(|key| rustfs_utils::http::is_object_encryption_marker(key)) + { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_TRANSFORMED); + } + if candidate.is_remote() { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_REMOTE); + } + if candidate.deleted { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_DELETED); + } + if candidate.size <= 0 { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); + } + if candidate.parts.len() != 1 { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE); + } + if !candidate.has_valid_erasure_geometry() { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_GEOMETRY); + } + + let Ok(object_size) = usize::try_from(candidate.size) else { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); + }; + if candidate.parts.first().is_none_or(|part| part.size != object_size) { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_PART_SHAPE); + } + if !can_try_inline_data_shards_direct(object_size, candidate.erasure.block_size) { + return Some(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_SIZE); + } + None +} + +pub(in crate::set_disk) fn non_inline_data_read_candidate_is_safe(candidate: &FileInfo) -> bool { + if candidate.inline_data() + || candidate.is_compressed() + || candidate.is_remote() + || candidate + .metadata + .keys() + .any(|key| rustfs_utils::http::is_object_encryption_marker(key)) + || candidate.parts.len() != 1 + { + return false; + } + candidate.has_valid_erasure_geometry() +} + +pub(in crate::set_disk) fn late_materialization_candidate_is_safe(candidate: &FileInfo) -> bool { + non_inline_data_read_candidate_is_safe(candidate) + && candidate.size > 512 * 1024 + && object_fits_single_block(candidate.size, candidate.erasure.block_size) +} + +pub(in crate::set_disk) fn non_inline_data_read_early_stop_allowed(read_data: bool, bucket: &str, object: &str) -> bool { + read_data && is_get_metadata_non_inline_data_read_early_stop_enabled() && !codec_streaming_rollout_applies(bucket, object) +} + +const NON_INLINE_SINGLE_PENDING_HEDGE_DELAY: Duration = Duration::from_millis(100); + +fn data_read_inline_missing_shards_are_pending( + candidate: &FileInfo, + slots: &[MetadataDiskObservation], + disks: &[Option], + fanout_order: &[usize], + scheduled_fanout_len: usize, +) -> bool { + 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 distribution = &candidate.erasure.distribution; + let mut data_shards_seen_or_pending = vec![false; erasure.data_shards]; + let mut missing_pending_data_shards = 0usize; + + for slot in slots { + let disk_index = slot.disk_index; + let Some(&block_index) = distribution.get(disk_index) else { + return false; + }; + if block_index == 0 || block_index > erasure.data_shards { + continue; + } + if !disks.get(disk_index).is_some_and(Option::is_some) { + return false; + } + + let data_slot = block_index - 1; + if slot.file_info().is_none_or(|metadata| metadata.name.is_empty()) { + let scheduled_and_not_failed = fanout_order + .get(..scheduled_fanout_len) + .is_some_and(|scheduled_disks| scheduled_disks.contains(&disk_index)) + && slot.result.is_none(); + if scheduled_and_not_failed { + data_shards_seen_or_pending[data_slot] = true; + missing_pending_data_shards = missing_pending_data_shards.saturating_add(1); + continue; + } + return false; + } + let Some(file_info) = slot.file_info() else { return false }; + if file_info.erasure.index != block_index + || !file_info.has_valid_erasure_geometry() + || !metadata_early_stop_candidate_matches(file_info, candidate) + || file_info.data.as_ref().is_none_or(|data| data.is_empty()) + { + return false; + } + data_shards_seen_or_pending[data_slot] = true; + } + + missing_pending_data_shards > 0 && data_shards_seen_or_pending.into_iter().all(|seen_or_pending| seen_or_pending) +} + +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, + DiskError::FileVersionNotFound => GET_METADATA_RESPONSE_VERSION_NOT_FOUND, + DiskError::DiskNotFound => GET_METADATA_RESPONSE_DISK_NOT_FOUND, + DiskError::FileCorrupt | DiskError::CorruptedFormat | DiskError::CorruptedBackend | DiskError::OutdatedXLMeta => { + GET_METADATA_RESPONSE_CORRUPT + } + DiskError::Timeout => GET_METADATA_RESPONSE_TIMEOUT, + DiskError::FaultyDisk | DiskError::FaultyRemoteDisk => GET_METADATA_RESPONSE_IGNORED, + _ => GET_METADATA_RESPONSE_ERROR, + } +} + +pub(in crate::set_disk) fn should_allow_metadata_early_stop( + read_data: bool, + version_id: &str, + healing: bool, + incl_free_versions: bool, +) -> bool { + if read_data && !is_get_metadata_data_read_early_stop_enabled() { + return false; + } + + (is_get_metadata_early_stop_enabled() && version_id.is_empty() && !healing && !incl_free_versions) + || (is_version_early_stop_enabled() && !version_id.is_empty() && !healing && !incl_free_versions) +} + +/// Final gate for the metadata early-stop fast path. +/// +/// `caller_allows_early_stop=false` unconditionally forces the full quorum +/// fanout so read-before-write callers (object tagging) get the complete +/// online-disk set as their write target; the early-stop subset would only +/// carry read quorum and fail write quorum (backlog#872 regression). +pub(in crate::set_disk) fn metadata_early_stop_permitted( + caller_allows_early_stop: bool, + observe: bool, + read_data: bool, + version_id: &str, + healing: bool, + incl_free_versions: bool, +) -> bool { + caller_allows_early_stop && observe && should_allow_metadata_early_stop(read_data, version_id, healing, incl_free_versions) +} + +impl SetDisks { + #[allow(clippy::too_many_arguments)] + #[tracing::instrument(level = "debug", skip(disks))] + pub(in crate::set_disk) async fn read_metadata( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + healing: bool, + incl_free_versions: bool, + ) -> disk::error::Result { + Self::read_metadata_inner( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + healing, + incl_free_versions, + false, + true, + 0, + false, + ) + .await + } + + #[allow(clippy::too_many_arguments)] + pub(in crate::set_disk) async fn read_metadata_observed( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + healing: bool, + incl_free_versions: bool, + caller_allows_early_stop: bool, + default_parity_count: usize, + ) -> disk::error::Result { + Self::read_metadata_inner( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + healing, + incl_free_versions, + true, + caller_allows_early_stop, + default_parity_count, + false, + ) + .await + } + + #[allow(clippy::too_many_arguments)] + pub(in crate::set_disk) async fn read_metadata_for_get_object( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + incl_free_versions: bool, + caller_allows_early_stop: bool, + default_parity_count: usize, + ) -> disk::error::Result { + Self::read_metadata_inner( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + false, + incl_free_versions, + true, + caller_allows_early_stop, + default_parity_count, + true, + ) + .await + } + + #[allow(clippy::too_many_arguments)] + pub(in crate::set_disk) async fn read_metadata_inner( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + healing: bool, + incl_free_versions: bool, + observe: bool, + // When false, the caller opts out of the early-stop fast path even for + // otherwise-eligible reads. Read-before-write callers (e.g. object + // tagging) must set this so the returned online-disk set reflects the + // full quorum fanout rather than the early-stop subset — writing to the + // subset would fail write quorum (backlog#872 regression). + caller_allows_early_stop: bool, + default_parity_count: usize, + allow_coalescing: bool, + ) -> disk::error::Result { + let early_stop_enabled = + caller_allows_early_stop && observe && (is_get_metadata_early_stop_enabled() || is_version_early_stop_enabled()); + let allow_early_stop = + metadata_early_stop_permitted(caller_allows_early_stop, observe, read_data, version_id, healing, incl_free_versions); + if allow_early_stop { + return Self::read_metadata_early_stop( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + healing, + incl_free_versions, + non_inline_data_read_early_stop_allowed(read_data, bucket, object), + default_parity_count, + allow_coalescing, + ) + .await; + } + if early_stop_enabled { + let metrics_path = metadata_metrics_path(bucket); + rustfs_io_metrics::record_get_object_metadata_early_stop_miss( + metrics_path, + GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, + ); + rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(metrics_path, 0); + } + + Self::read_metadata_full_wait( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + healing, + incl_free_versions, + observe, + allow_coalescing, + ) + .await + } + + #[allow(clippy::too_many_arguments)] + async fn read_metadata_full_wait( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + healing: bool, + incl_free_versions: bool, + observe: bool, + allow_coalescing: bool, + ) -> disk::error::Result { + let fanout_start = observe.then(Instant::now); + let mut slots = MetadataReadResult::pending(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, + healing, + }; + let org_bucket: Arc = Arc::from(org_bucket); + let bucket: Arc = Arc::from(bucket); + let object: Arc = Arc::from(object); + let version_id: Arc = Arc::from(version_id); + let slowtail_fault = get_metadata_slowtail_fault_request(bucket.as_ref(), object.as_ref(), read_data); + let futures = disks.iter().enumerate().map(|(disk_index, disk)| { + let disk = disk.clone(); + let task_opts = opts; + let org_bucket = org_bucket.clone(); + let bucket = bucket.clone(); + let object = object.clone(); + let version_id = version_id.clone(); + let slowtail_fault = slowtail_fault.clone(); + AbortOnDropJoinHandle(tokio::spawn(async move { + let response_start = observe.then(Instant::now); + let result = if let Some(disk) = disk { + Self::record_read_version_call(&object, disk_index); + if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(disk_index)) { + Self::record_metadata_slowtail_fault(&object, disk_index); + tokio::time::sleep(delay).await; + } + read_version_via_coalescer(disk, &org_bucket, &bucket, &object, &version_id, &task_opts, allow_coalescing) + .await + } else { + Err(DiskError::DiskNotFound) + }; + let elapsed = response_start.map(|start| start.elapsed()); + (result, elapsed) + })) + }); + + // Wait for all futures to complete + let results = join_all(futures).await; + + for (index, join_result) in results.into_iter().enumerate() { + match join_result { + Ok((res, elapsed)) => match res { + Ok(file_info) => { + if let (Some(observations), Some(elapsed)) = (&mut observations, elapsed) { + observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed)); + } + slots[index].result = Some(Ok(file_info)); + } + Err(e) => { + if let (Some(observations), Some(elapsed)) = (&mut observations, elapsed) { + observations.push(MetadataFanoutObservation::from_error(&e, elapsed)); + } + slots[index].result = Some(Err(e)); + } + }, + Err(_join_err) => { + // A spawned task panicked — treat as unexpected disk error + if let Some(observations) = &mut observations { + observations.push(MetadataFanoutObservation::from_error(&DiskError::Unexpected, Duration::ZERO)); + } + slots[index].result = Some(Err(DiskError::Unexpected)); + } + } + } + let diagnostics = match (fanout_start, observations) { + (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(MetadataReadResult { slots, diagnostics }) + } + + #[allow(clippy::too_many_arguments)] + async fn read_metadata_early_stop( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + healing: bool, + incl_free_versions: bool, + allow_non_inline_data_read_early_stop: bool, + default_parity_count: usize, + allow_coalescing: bool, + ) -> disk::error::Result { + let fanout_start = Instant::now(); + let mut slots = MetadataReadResult::pending(disks.len()); + let mut task_ids = vec![None; disks.len()]; + let mut observations = Vec::with_capacity(disks.len()); + let mut accumulator = + MetadataQuorumAccumulator::new(disks.len(), default_parity_count, true).with_requested_version_id(version_id); + let opts = ReadOptions { + incl_free_versions, + read_data, + healing, + }; + let org_bucket: Arc = Arc::from(org_bucket); + 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 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 mut final_miss_reason_override = None; + let mut non_inline_candidate_eligible = None; + let mut single_pending_hedge_deadline = None; + let slowtail_fault = get_metadata_slowtail_fault_request(bucket.as_ref(), object.as_ref(), read_data); + let spawn_read_version = |join_set: &mut JoinSet<(usize, disk::error::Result, Duration)>, + task_ids: &mut [Option], + index: usize, + disk: Option| { + let task_opts = opts; + let org_bucket = org_bucket.clone(); + let bucket = bucket.clone(); + let object = object.clone(); + let version_id = version_id.clone(); + let slowtail_fault = slowtail_fault.clone(); + let task = join_set.spawn(async move { + let response_start = Instant::now(); + let result = if let Some(disk) = disk { + #[cfg(test)] + let _fanout_task_guard = rename_fanout_barrier::task_guard(&object); + Self::record_read_version_call(&object, index); + #[cfg(test)] + Self::read_version_fanout_barrier(&object, index).await; + if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(index)) { + Self::record_metadata_slowtail_fault(&object, index); + tokio::time::sleep(delay).await; + } + read_version_via_coalescer(disk, &org_bucket, &bucket, &object, &version_id, &task_opts, allow_coalescing) + .await + } else { + Err(DiskError::DiskNotFound) + }; + (index, result, response_start.elapsed()) + }); + task_ids[index] = Some(task.id()); + }; + + if bounded_fanout { + let initial_target = accumulator.default_write_quorum().min(disks.len()); + 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, &mut task_ids, disk_index, disk); + scheduled_count = scheduled_count.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, &mut task_ids, index, disk); + scheduled_count = scheduled_count.saturating_add(1); + } + } + + loop { + let mut defer_pending_inline_data_shard = false; + let result = if let Some(deadline) = single_pending_hedge_deadline.take() { + tokio::select! { + result = join_set.join_next_with_id() => result, + _ = async { + tokio::time::sleep_until(deadline).await; + #[cfg(test)] + rename_fanout_barrier::checkpoint( + object.as_ref(), + 0, + rename_fanout_barrier::PHASE_NON_INLINE_HEDGE_TIMER, + ) + .await; + } => { + if bounded_fanout + && !force_full_wait + && join_set.len() == 1 + && non_inline_candidate_eligible == Some(true) + && !accumulator.candidate_has_read_reserve() + && next_fanout_index < disks.len() + { + while next_fanout_index < disks.len() { + let disk_index = fanout_order[next_fanout_index]; + next_fanout_index = next_fanout_index.saturating_add(1); + if let Some(disk) = disks.get(disk_index).cloned() { + spawn_read_version(&mut join_set, &mut task_ids, disk_index, disk); + scheduled_count = scheduled_count.saturating_add(1); + break; + } + } + } + continue; + } + } + } else { + join_set.join_next_with_id().await + }; + let Some(result) = result else { break }; + match result { + Ok((_task_id, (index, res, elapsed))) => match res { + Ok(file_info) => { + observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed)); + let inline_candidate_miss_reason = if bounded_fanout && read_data && !force_full_wait { + data_read_early_stop_inline_candidate_miss_reason(&file_info) + } else { + None + }; + slots[index].result = Some(Ok(file_info)); + accumulator.observe(&slots[index]); + if allow_non_inline_data_read_early_stop && non_inline_candidate_eligible.is_none() { + non_inline_candidate_eligible = + accumulator.candidate.as_ref().map(non_inline_data_read_candidate_is_safe); + } + if bounded_fanout + && read_data + && !force_full_wait + && let Some(reason) = inline_candidate_miss_reason + && !(non_inline_candidate_eligible == Some(true) + && reason == GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_NOT_INLINE) + { + force_full_wait = true; + final_miss_reason_override.get_or_insert(reason); + } + } + Err(err) => { + observations.push(MetadataFanoutObservation::from_error(&err, elapsed)); + slots[index].result = Some(Err(err)); + accumulator.observe(&slots[index]); + } + }, + Err(join_error) => { + let err = DiskError::Unexpected; + observations.push(MetadataFanoutObservation::from_error(&err, fanout_start.elapsed())); + if let Some(index) = task_ids.iter().position(|id| *id == Some(join_error.id())) { + slots[index].result = Some(Err(err)); + accumulator.observe(&slots[index]); + } + } + } + + if !force_full_wait + && let Some(decision) = accumulator + .early_stop_decision() + .or_else(|| accumulator.version_early_stop_decision()) + { + let should_return_early = if read_data { + match accumulator.candidate.as_ref() { + Some(_candidate) if non_inline_candidate_eligible == Some(true) => { + accumulator.candidate_has_read_reserve() + } + Some(candidate) => match inline_body_miss_reason_from_observations( + bucket.as_ref(), + object.as_ref(), + candidate, + slots.iter().map(|slot| (slot.disk_index, slot.file_info())), + disks, + ) + .await + { + None => true, + Some(reason) => { + final_miss_reason_override = Some(reason); + if bounded_fanout + && reason == GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD + && data_read_inline_missing_shards_are_pending( + candidate, + &slots, + disks, + &fanout_order, + next_fanout_index, + ) + { + defer_pending_inline_data_shard = true; + } else { + force_full_wait = true; + } + false + } + }, + None => { + force_full_wait = true; + final_miss_reason_override = Some(GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM); + false + } + } + } else { + true + }; + + 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_with_id().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(MetadataReadResult { slots, diagnostics }); + } + } + + let pending_responses = join_set.len(); + // Inline verification can still depend on a missing data shard; + // issue one immediate spare when only that shard remains. The + // non-inline path keeps its delayed hedge below to avoid healthy + // reads paying speculative I/O before the candidate is classified. + let should_hedge_single_pending_inline_read = read_data + && !force_full_wait + && !defer_pending_inline_data_shard + && pending_responses == 1 + && non_inline_candidate_eligible != Some(true) + && accumulator.can_still_reach_early_stop_with_pending(pending_responses); + // A non-inline plan must retain one extra matching shard as a + // reconstruction reserve. Schedule that reserve only after the + // candidate is known to be eligible, so inline GETs do not pay an + // extra fanout and the healthy path remains allocation-free. + let needs_non_inline_read_reserve = non_inline_candidate_eligible == Some(true) + && !accumulator.candidate_has_read_reserve() + && accumulator + .candidate_read_reserve_target() + .is_some_and(|reserve_target| scheduled_count < reserve_target || pending_responses == 0); + if bounded_fanout + && !force_full_wait + && (needs_non_inline_read_reserve || should_hedge_single_pending_inline_read) + && 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, &mut task_ids, disk_index, disk); + scheduled_count = scheduled_count.saturating_add(1); + } + next_fanout_index = next_fanout_index.saturating_add(1); + } else 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, &mut task_ids, disk_index, disk); + scheduled_count = scheduled_count.saturating_add(1); + } + next_fanout_index = next_fanout_index.saturating_add(1); + } + } else if bounded_fanout + && !defer_pending_inline_data_shard + && next_fanout_index < disks.len() + && !accumulator.can_still_reach_early_stop_with_pending(pending_responses) + { + let disk_index = fanout_order[next_fanout_index]; + if let Some(disk) = disks.get(disk_index).cloned() { + spawn_read_version(&mut join_set, &mut task_ids, disk_index, disk); + scheduled_count = scheduled_count.saturating_add(1); + } + next_fanout_index = next_fanout_index.saturating_add(1); + } + if bounded_fanout + && !force_full_wait + && !defer_pending_inline_data_shard + && join_set.len() == 1 + && non_inline_candidate_eligible == Some(true) + && !accumulator.candidate_has_read_reserve() + && accumulator.can_still_reach_early_stop_with_pending(join_set.len()) + && next_fanout_index < disks.len() + { + single_pending_hedge_deadline = Some(tokio::time::Instant::now() + NON_INLINE_SINGLE_PENDING_HEDGE_DELAY); + } + } + + let accumulator_miss_reason = accumulator.final_miss_reason(); + let final_miss_reason = match (final_miss_reason_override, accumulator_miss_reason) { + (Some(reason), GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM) => reason, + _ => accumulator_miss_reason, + }; + rustfs_io_metrics::record_get_object_metadata_early_stop_miss(metrics_path, 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(MetadataReadResult { slots, diagnostics }) + } + + /// Test-only seam that records one per-disk `read_version` metadata RPC for + /// the call-counter registry (backlog#1325). In production this is inlined to + /// nothing and adds no behavior; only the `#[cfg(test)]` variant touches the + /// registry. Placed inside the fan-out spawn tasks so counts are observed + /// even though the increments run on arbitrary runtime worker threads. + #[cfg(test)] + #[inline] + fn record_read_version_call(object: &str, disk_index: usize) { + disk_call_counters::record(object, disk_call_counters::KIND_READ_VERSION, disk_index); + } + + #[cfg(not(test))] + #[inline(always)] + fn record_read_version_call(_object: &str, _disk_index: usize) {} + + #[cfg(test)] + #[inline] + fn record_metadata_slowtail_fault(object: &str, disk_index: usize) { + disk_call_counters::record(object, disk_call_counters::KIND_METADATA_SLOWTAIL_FAULT, disk_index); + } + + #[cfg(not(test))] + #[inline(always)] + fn record_metadata_slowtail_fault(_object: &str, _disk_index: usize) {} + + #[cfg(test)] + #[inline] + async fn read_version_fanout_barrier(object: &str, disk_index: usize) { + rename_fanout_barrier::checkpoint(object, disk_index, rename_fanout_barrier::PHASE_READ_VERSION).await; + } +} diff --git a/crates/ecstore/src/set_disk/core/mod.rs b/crates/ecstore/src/set_disk/core/mod.rs index b8507603e..8ecce941c 100644 --- a/crates/ecstore/src/set_disk/core/mod.rs +++ b/crates/ecstore/src/set_disk/core/mod.rs @@ -18,4 +18,5 @@ //! duplicating read/write/erasure logic. pub(crate) mod io_primitives; -mod metadata_quorum; +pub(in crate::set_disk) mod metadata_quorum; +pub(in crate::set_disk) mod metadata_read; diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 78edf3df5..e12e42feb 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -5527,12 +5527,29 @@ fn collect_inline_data_shard_fileinfos_by_index_or_reason<'a>( parts_metadata: &'a [FileInfo], fi: &FileInfo, data_shards: usize, + disk_is_online: impl FnMut(usize) -> bool, +) -> std::result::Result, &'static str> { + collect_inline_data_shard_fileinfos_from_observations( + parts_metadata + .iter() + .enumerate() + .map(|(index, metadata)| (index, Some(metadata))), + fi, + data_shards, + disk_is_online, + ) +} + +fn collect_inline_data_shard_fileinfos_from_observations<'a>( + observations: impl IntoIterator)>, + fi: &FileInfo, + data_shards: usize, mut disk_is_online: impl FnMut(usize) -> bool, ) -> std::result::Result, &'static str> { let distribution = &fi.erasure.distribution; let mut data_files = vec![None; data_shards]; - for (disk_index, file_info) in parts_metadata.iter().enumerate() { + for (disk_index, file_info) in observations { if !disk_is_online(disk_index) { continue; } @@ -5542,9 +5559,9 @@ fn collect_inline_data_shard_fileinfos_by_index_or_reason<'a>( if block_index == 0 || block_index > data_shards { continue; } - if file_info.name.is_empty() { + let Some(file_info) = file_info.filter(|metadata| !metadata.name.is_empty()) else { return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_MISSING_SHARD); - } + }; if file_info.erasure.index != block_index { return Err(GET_METADATA_EARLY_STOP_REASON_DATA_READ_INLINE_IDENTITY_MISMATCH); } diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 617c49d56..39924f1a4 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -598,17 +598,18 @@ impl SetDisks { } scope.check()?; - let (_, errors) = - Self::read_all_fileinfo(disks.as_slice(), "", bucket, object, &version.to_string(), false, false, false).await?; + let metadata_read = + Self::read_metadata(disks.as_slice(), "", bucket, object, &version.to_string(), false, false, false).await?; + let all_absent = metadata_read + .slots + .iter() + .all(|slot| matches!(&slot.result, Some(Err(DiskError::FileNotFound | DiskError::FileVersionNotFound)))); + let selected = self.get_disks_internal().await; let same_targets = selected.len() == disks.len() && selected.iter().zip(&disks).all(|(current, original)| { matches!((current, original), (Some(current), Some(original)) if std::sync::Arc::ptr_eq(current, original)) }); - if !same_targets - || !errors - .iter() - .all(|error| matches!(error, Some(DiskError::FileNotFound | DiskError::FileVersionNotFound))) - { + if !same_targets || !all_absent { return Err(StorageError::SlowDown); } scope.check()?; @@ -734,7 +735,7 @@ impl SetDisks { object: &str, version_id: &str, ) -> disk::error::Result { - let (parts_metadata, errs, _) = Self::read_all_fileinfo_observed( + let metadata_read = Self::read_metadata_observed( disks, "", bucket, @@ -747,6 +748,7 @@ impl SetDisks { self.default_parity_count, ) .await?; + let (parts_metadata, errs, _) = metadata_read.into_legacy(); let (read_quorum, _) = Self::object_quorum_from_meta(&parts_metadata, &errs, self.default_parity_count)?; let read_quorum = usize::try_from(read_quorum).map_err(|_| DiskError::ErasureReadQuorum)?; let (_, quorum_mod_time, quorum_etag) = Self::list_online_disks(disks, &parts_metadata, &errs, read_quorum); @@ -931,8 +933,9 @@ impl SetDisks { } }; - let (mut parts_metadata, errs) = - Self::read_all_fileinfo(&disks, "", bucket, object, version_id, true, true, false).await?; + let (mut parts_metadata, errs, _) = Self::read_metadata(&disks, "", bucket, object, version_id, true, true, false) + .await? + .into_legacy(); trace!( event = EVENT_SET_DISK_HEAL, @@ -1284,8 +1287,10 @@ impl SetDisks { // in `parts_metadata` to defaults. Re-read only before // destructive cleanup so the guard sees every original // identity. - let (delete_guard_metadata, delete_guard_errs) = - Self::read_all_fileinfo(&disks, "", bucket, object, version_id, true, true, false).await?; + let (delete_guard_metadata, delete_guard_errs, _) = + Self::read_metadata(&disks, "", bucket, object, version_id, true, true, false) + .await? + .into_legacy(); if self .dangling_delete_safety(bucket, object, &delete_guard_metadata, &delete_guard_errs, &disks) .await? @@ -2316,16 +2321,16 @@ impl SetDisks { return Err(DiskError::retired_marker_deferred(format!("conditional marker deletion failed: {error}"))); } } - let (_, errors) = Self::read_all_fileinfo(disks, "", bucket, object, &version.to_string(), false, false, false).await?; + let metadata_read = Self::read_metadata(disks, "", bucket, object, &version.to_string(), false, false, false).await?; + let all_absent = metadata_read + .slots + .iter() + .all(|slot| matches!(&slot.result, Some(Err(DiskError::FileNotFound | DiskError::FileVersionNotFound)))); let selected = self.get_disks_internal().await; let same_targets = selected.len() == disks.len() && selected.iter().zip(disks).all(|(current, original)| { matches!((current, original), (Some(current), Some(original)) if std::sync::Arc::ptr_eq(current, original)) }); - if !same_targets - || !errors - .iter() - .all(|error| matches!(error, Some(DiskError::FileNotFound | DiskError::FileVersionNotFound))) - { + if !same_targets || !all_absent { return Err(DiskError::retired_marker_deferred( "complete marker absence could not be verified under the original target set", )); @@ -3383,9 +3388,10 @@ impl SetDisks { // The inner heal and missing-object report read the registry again; // release this snapshot guard before a topology writer can queue between reads. let disks = self.get_disks_internal().await; - let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false, false) + let (_, errs, _) = Self::read_metadata(&disks, "", bucket, object, version_id, false, false, false) .await - .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; + .map_err(|e| to_object_err(e.into(), vec![bucket, object]))? + .into_legacy(); if DiskError::is_all_not_found(&errs) { debug!( event = EVENT_SET_DISK_HEAL, diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 9e4061ebf..3f5c4e2df 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -113,7 +113,28 @@ use super::GetCodecStreamingObjectClass; use super::GetCodecStreamingRollout; #[cfg(test)] use super::classify_get_codec_streaming_object_class; -use super::core::io_primitives::*; +use super::core::io_primitives::{ + BitrotReaderSetup, BitrotReaderSetupAttribution, BitrotReaderSetupMode, DeferredReaderReopener, EVENT_SET_DISK_READ, + GetCodecStreamingReaderBuildOutcome, MetadataCacheLookup, ObjectBitrotReader, ReadRepairAdmissionSubmitter, + ReadRepairHealSubmission, SLOW_OBJECT_READ_LOG_THRESHOLD, codec_streaming_reader_setup_fallback_reason, + create_bitrot_readers_until_quorum_with_preference, create_data_block_bitrot_readers, resolved_read_repair_version_id, + send_read_repair_heal_request, shard_read_costs_for_disks, submit_read_repair_heal, submit_read_repair_heal_with_submitter, +}; +#[cfg(test)] +use super::core::io_primitives::{ + ENV_RUSTFS_GET_CODEC_STREAMING_DATA_BLOCKS_FIRST_READER_SETUP, ENV_RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP, + MultipartCodecStreamingReader, ReadRepairAdmissionFuture, ReadRepairAdmissionOutcome, collect_read_multiple_results, + collect_read_parts_results, create_bitrot_readers_until_quorum, create_bitrot_readers_until_quorum_all_shards, + release_read_repair_heal_reservation, reserve_read_repair_heal, resolve_read_part_from_responses, shard_read_cost_for_disk, + shard_read_cost_for_endpoint, +}; +#[cfg(test)] +use super::core::metadata_quorum::{MetadataEarlyStopDecision, MetadataQuorumAccumulator}; +#[cfg(test)] +use super::core::metadata_read::{ + MetadataFanoutDiagnostics, MetadataFanoutObservation, metadata_early_stop_permitted, should_allow_metadata_early_stop, +}; +use super::core::metadata_read::{late_materialization_candidate_is_safe, non_inline_data_read_early_stop_allowed}; #[cfg(test)] use super::get_codec_streaming_config_cached_core; #[cfg(test)] @@ -607,11 +628,11 @@ impl SetDisks { }; // Early-stop for safe metadata reads is handled inside - // read_all_fileinfo_observed (see read_all_fileinfo_early_stop in - // core/io_primitives.rs); unsafe requests and callers that opt out + // read_metadata_observed (see read_all_fileinfo_early_stop in + // core/metadata_read.rs); unsafe requests and callers that opt out // (allow_early_stop=false) fall back to full-wait. - let (mut parts_metadata, errs, metadata_fanout_diagnostics) = if allow_read_version_coalescing { - Self::read_all_fileinfo_observed_for_get_object( + let metadata_read = if allow_read_version_coalescing { + Self::read_metadata_for_get_object( &disks, "", bucket, @@ -624,7 +645,7 @@ impl SetDisks { ) .await? } else { - Self::read_all_fileinfo_observed( + Self::read_metadata_observed( &disks, "", bucket, @@ -638,13 +659,14 @@ impl SetDisks { ) .await? }; + let metadata_fanout_complete = metadata_read.is_complete(); + let (mut parts_metadata, errs, metadata_fanout_diagnostics) = metadata_read.into_legacy(); let metadata_metrics_path = if crate::bucket::utils::is_meta_bucketname(bucket) { GET_OBJECT_PATH_INTERNAL_META } else { GET_OBJECT_PATH_LEGACY_DUPLEX }; metadata_fanout_diagnostics.record(metadata_metrics_path); - let metadata_fanout_complete = metadata_fanout_diagnostics.total_responses() >= disks.len(); // warn!("get_object_fileinfo parts_metadata {:?}", &parts_metadata); // warn!("get_object_fileinfo {}/{} errs {:?}", bucket, object, &errs); @@ -2117,7 +2139,7 @@ impl SetDisks { expected: &LateMetadataIdentity, metrics_path: &'static str, ) -> Result<(FileInfo, Vec, Vec>)> { - let (mut parts_metadata, errs, diagnostics) = SetDisks::read_all_fileinfo_observed( + let metadata_read = SetDisks::read_metadata_observed( fallback_disks, "", bucket, @@ -2130,6 +2152,7 @@ impl SetDisks { expected.parity_blocks, ) .await?; + let (mut parts_metadata, errs, diagnostics) = metadata_read.into_legacy(); diagnostics.record(metrics_path); let (read_quorum, write_quorum) = SetDisks::object_quorum_from_meta(&parts_metadata, &errs, expected.parity_blocks) diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index e1d0cb5f5..423ca0bd5 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -11,6 +11,8 @@ ## Open Items +- `backlog-2253-metadata-observation-adapter` metadata observations: preserve aligned metadata/error slices for existing layout, shard and exact-version PUT, MPU and delete consumers while read/heal scheduling uses typed disk observations. Pending slots acquire empty placeholders only in this adapter. Remove after those remaining consumers use typed observations; delete the adapter and old facade in a separate cleanup. + - `connect-894` pending Connect heartbeats: replay the exact preceding producer-capability list, with or without jobs, when upgrading to the memory-service capability. Preserve request ID, sequence, and persisted body rather than inserting the new capability into a retry. Remove after upgrades from the pre-service-memory capability set are unsupported. - `backlog-2539` administrator erasure-set scope: decode historical pending intents whose empty bucket list omitted the all-buckets marker. Normalize that representation to the explicit scope before replay or checkpoint comparison. Remove after all pre-marker administrator ErasureSet intents are retired; deployment verification must confirm that no such pending records remain on coordinator disks. diff --git a/docs/architecture/ecstore-module-split-plan.md b/docs/architecture/ecstore-module-split-plan.md index eb26b73f4..dd9cd8d84 100644 --- a/docs/architecture/ecstore-module-split-plan.md +++ b/docs/architecture/ecstore-module-split-plan.md @@ -83,3 +83,39 @@ A candidate is ready for code movement only when all of these hold: - old facade names have compatibility tests or explicit deprecation coverage; - focused tests cover the changed owner path before any full gate is attempted; - rollback preserves object IO, quorum, lifecycle/replication queues, scanner repair, notification/audit events, and metadata compatibility. + +## Metadata read ownership and identity + +`set_disk/core/metadata_read.rs` owns disk/RPC scheduling, coalescing, timeouts, +hedging and cancellation. Its result retains each disk slot as pending, a +successful `FileInfo`, or a typed `DiskError`. Pending includes both an +unscheduled disk and a response not yet observed; it contributes no vote and +does not establish absence. Task IDs retain the disk slot when a task panics. +The existing full-wait and early-stop gates, coalescer admission, cancellation +drain and fallback policy remain at this scheduling boundary. + +`set_disk/core/metadata_quorum.rs` consumes those observations without IO. +`SetDisks` remains the sole owner of topology, locks and caches. Read and heal +callers use the typed result; `MetadataReadResult::into_legacy` is the explicit +adapter for existing quorum/layout and shard-reader helpers. It moves successful +metadata without cloning and creates empty placeholders only at that adapter. +Old full-wait facade callers retain the same aligned metadata/error slices. + +These identities have different purposes and must not become one generic hash: + +| Decision | Equal fields | Deliberately different or separately checked | +| --- | --- | --- | +| Quorum grouping (`file_info_quorum_hash`) | Size, deletion/restore/transition state, timestamps, version and data-directory UUIDs, checksums, mode/writer version, normalized metadata, target delete-marker versions and part identities. Payloads also include data/parity counts and distribution. | Object name and volume come from the request; per-disk erasure index and shard bytes are not an object identity. Canonical delete markers and zero-size objects do not add payload geometry. | +| Early-stop strict match (`metadata_early_stop_candidate_matches`) | Request name/volume, version/latest/deletion state, transition fields, size/time, mode/writer version, complete metadata/replication state/parts/checksum, version-history fields, data directory, algorithm/block size/data/parity/distribution. | Per-disk erasure index and data bytes differ. A non-inline read reserve additionally proves distinct indexes mapped to their disk slots; inline reads verify shard identity and bitrot/body content before stopping. | +| Late materialization (`LateMetadataIdentity`) | Request name/volume, algorithm, block size, legacy-checksum mode and quorum hash. | Each late shard must match the original distribution at its disk slot. A changed identity or insufficient matching shards fails read quorum; it cannot mix generations into a body. | + +`metadata_observation_all_four_slot_states_and_arrival_orders_preserve_reduction` +checks all 625 four-disk success/missing/corrupt/offline/pending combinations and +24 arrival orders with early-stop on and off. The observation regressions also +cover a newer pending response, same-time different directories, duplicated +shard indexes, null versions, canonical delete markers and invalid success +payloads. Existing bounded-fanout, legacy inline bitrot, cache, late-metadata and +GET/Range regressions continue to exercise the real scheduling and reader paths; +the pure matrix does not replace those checks. Their call counters bound actual +disk reads, including canceled tails, without adding a scheduling policy or a +per-request metadata clone.