From 166679a7236f9f3584359a22a3455fcbfb55f730 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sun, 5 Jul 2026 22:35:05 +0800 Subject: [PATCH] refactor(ecstore): sink write + shared-read primitives into set_disk::core::io_primitives (backlog#820) (#4288) * refactor(ecstore): sink write/rename/delete primitives into set_disk::core::io_primitives (backlog#820) P5 step 2 of the SetDisks God-Object split (tracking backlog#815, issue backlog#820; follows step 1 #4285). Relocate the entire set_disk/write.rs primitive family into the core/io_primitives.rs module home established by step 1: - Module items: dangling_delete_grace() + its env consts, and the OrphanDirScan scan-result enum. - impl SetDisks write/rename/delete primitives shared across mod.rs and the ops/ operation families: rename_data, commit_rename_data_dir, cleanup_multipart_path, rename_part, eval_disks, write_unique_file_info, update_object_meta(_with_opts), delete_if_dangling, delete_prefix, scan_orphan_dir, purge_orphan_dir_object, check_write_precondition, default_read_quorum, default_write_quorum. - The dangling_delete_grace unit tests move with their subject. set_disk/write.rs is removed and 'mod write;' dropped from mod.rs. Pure move + visibility adjustment, zero logic change. Method bodies are moved verbatim; pub(super) items are widened to pub(in crate::set_disk) so mod.rs / ops still reach them (write.rs's super was set_disk; the deeper module needs the explicit path to preserve identical reach). Callers use inherent self./Self:: calls, so no call sites change. Verification: - cargo check / clippy -D warnings -p rustfs-ecstore --all-targets: clean - cargo test -p rustfs-ecstore --lib: 1841 passed, 0 failed (moved dangling_delete_grace tests run as set_disk::core::io_primitives::tests::*) - all five arch guard scripts: pass - token-stream diff of moved block vs original write.rs: identical modulo visibility tokens (2713 tokens each) * refactor(ecstore): sink shared read primitives into set_disk::core::io_primitives (backlog#820) P5 step 3 (final) of the SetDisks God-Object split (tracking backlog#815, issue backlog#820; follows steps 1 #4285 and 2). Relocate the shared, low-level metadata/erasure READ PRIMITIVE methods out of set_disk/read.rs into the core/io_primitives.rs module home, leaving the object-read operation itself in read.rs. Moved into a new impl SetDisks block in io_primitives.rs (verbatim bodies): - read_parts, read_all_fileinfo (+ _observed / _inner / _full_wait / _early_stop variants), read_all_xl, load_file_info_versions_exact, read_all_raw_file_info, pick_latest_quorum_files_info, read_multiple_files. - The should_allow_metadata_early_stop free helper (called by the moved read_all_fileinfo_observed). Kept in read.rs: the object-read operation and its private helpers (read_version_optimized, get_object_fileinfo, get_object_info_and_quorum, try_get_object_direct_data_shards_with_fileinfo, get_object_with_fileinfo, get_object_decode_reader_with_fileinfo, build_codec_streaming_part_reader) and the metadata-cache helpers/tests. read.rs reaches the moved primitives through the SetDisks core (inherent self./Self:: calls; the sole cross-boundary edge, get_object_fileinfo -> read_all_fileinfo_observed, is handled by widening that one method to pub(in crate::set_disk)). Pure move + visibility adjustment, zero logic change. pub(super) primitives are widened to pub(in crate::set_disk); the three internal-only fanout variants (_inner/_full_wait/_early_stop) stay private; load_file_info_versions_exact stays pub(crate). Method bodies are moved verbatim. Verification: - cargo check / clippy -D warnings -p rustfs-ecstore --all-targets: clean - cargo test -p rustfs-ecstore --lib: 1841 passed, 0 failed - all five arch guard scripts: pass - token-stream diff of the three moved regions vs original read.rs: identical modulo visibility tokens, the new impl-block scaffolding, and rustfmt signature re-wrapping (read.rs diff is pure deletion, zero added lines) --- .../src/set_disk/core/io_primitives.rs | 1571 +++++++++++++++++ crates/ecstore/src/set_disk/mod.rs | 1 - crates/ecstore/src/set_disk/read.rs | 626 ------- crates/ecstore/src/set_disk/write.rs | 941 ---------- 4 files changed, 1571 insertions(+), 1568 deletions(-) delete mode 100644 crates/ecstore/src/set_disk/write.rs diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 7cca91923..940462a6c 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -1699,3 +1699,1574 @@ where Ok((responses, errors)) } + +// =========================================================================== +// Shared metadata/erasure read primitives (relocated verbatim from +// set_disk/read.rs, P5 step 3, tracking backlog#815, issue backlog#820). +// The object-read operation itself (get_object_*, read_version_optimized, the +// 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 { + 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) +} + +impl SetDisks { + pub(in crate::set_disk) async fn read_parts( + disks: &[Option], + bucket: &str, + part_meta_paths: &[String], + part_numbers: &[usize], + read_quorum: usize, + ) -> disk::error::Result> { + let bucket = bucket.to_string(); + let part_meta_paths = part_meta_paths.to_vec(); + + let tasks: Vec<_> = disks + .iter() + .map(|disk| { + let disk = disk.clone(); + let bucket = bucket.clone(); + let part_meta_paths = part_meta_paths.clone(); + + async move { + if let Some(disk) = disk { + disk.read_parts(&bucket, &part_meta_paths).await + } else { + Err(DiskError::DiskNotFound) + } + } + }) + .collect(); + + let (responses, collected_errors) = match collect_read_parts_results(tasks, read_quorum).await { + Ok(collected) => collected, + Err(()) => return Err(DiskError::ErasureReadQuorum), + }; + + if let Some(err) = reduce_read_quorum_errs(&collected_errors, OBJECT_OP_IGNORED_ERRS, read_quorum) { + return Err(err); + } + + let mut ret = vec![ObjectPartInfo::default(); part_meta_paths.len()]; + + for (part_idx, part_info) in part_meta_paths.iter().enumerate() { + ret[part_idx] = resolve_read_part_from_responses( + &bucket, + part_info, + part_numbers[part_idx], + part_idx, + part_meta_paths.len(), + &responses, + read_quorum, + )?; + } + + Ok(ret) + } + + #[allow(clippy::too_many_arguments)] + #[tracing::instrument(level = "debug", skip(disks))] + pub(in crate::set_disk) async fn read_all_fileinfo( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + healing: bool, + incl_free_versions: bool, + ) -> disk::error::Result<(Vec, Vec>)> { + let (ress, errors, _) = Self::read_all_fileinfo_inner( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + healing, + incl_free_versions, + false, + 0, + ) + .await?; + Ok((ress, errors)) + } + + #[allow(clippy::too_many_arguments)] + pub(in crate::set_disk) async fn read_all_fileinfo_observed( + disks: &[Option], + org_bucket: &str, + bucket: &str, + object: &str, + version_id: &str, + read_data: bool, + healing: bool, + incl_free_versions: bool, + default_parity_count: usize, + ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { + Self::read_all_fileinfo_inner( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + healing, + incl_free_versions, + true, + default_parity_count, + ) + .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, + default_parity_count: usize, + ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { + let early_stop_enabled = observe && (is_get_metadata_early_stop_enabled() || is_version_early_stop_enabled()); + let allow_early_stop = observe && should_allow_metadata_early_stop(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, + default_parity_count, + ) + .await; + } + if early_stop_enabled { + rustfs_io_metrics::record_get_object_metadata_early_stop_miss( + GET_OBJECT_PATH_LEGACY_DUPLEX, + GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, + ); + rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0); + } + + Self::read_all_fileinfo_full_wait( + disks, + org_bucket, + bucket, + object, + version_id, + read_data, + healing, + incl_free_versions, + observe, + ) + .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, + ) -> 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 opts = Arc::new(ReadOptions { + incl_free_versions, + read_data, + healing, + }); + let org_bucket = Arc::new(org_bucket.to_string()); + let bucket = Arc::new(bucket.to_string()); + let object = Arc::new(object.to_string()); + let version_id = Arc::new(version_id.to_string()); + let futures = disks.iter().map(|disk| { + let disk = disk.clone(); + let opts = opts.clone(); + let org_bucket = org_bucket.clone(); + let bucket = bucket.clone(); + let object = object.clone(); + let version_id = version_id.clone(); + tokio::spawn(async move { + let response_start = observe.then(Instant::now); + let result = if let Some(disk) = disk { + disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).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(), + }; + 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, + default_parity_count: usize, + ) -> 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 = Arc::new(ReadOptions { + incl_free_versions, + read_data, + healing, + }); + let org_bucket = Arc::new(org_bucket.to_string()); + let bucket = Arc::new(bucket.to_string()); + let object = Arc::new(object.to_string()); + let version_id = Arc::new(version_id.to_string()); + let mut join_set = JoinSet::new(); + + for (index, disk) in disks.iter().cloned().enumerate() { + let opts = opts.clone(); + let org_bucket = org_bucket.clone(); + let bucket = bucket.clone(); + let object = object.clone(); + let version_id = version_id.clone(); + join_set.spawn(async move { + let response_start = Instant::now(); + let result = if let Some(disk) = disk { + disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await + } else { + Err(DiskError::DiskNotFound) + }; + (index, result, response_start.elapsed()) + }); + } + + while let Some(result) = join_set.join_next().await { + match result { + Ok((index, res, elapsed)) => match res { + Ok(file_info) => { + observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed)); + accumulator.observe_file_info(&file_info); + 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 let Some(decision) = accumulator + .early_stop_decision() + .or_else(|| accumulator.version_early_stop_decision()) + { + let saved_responses = join_set.len(); + join_set.abort_all(); + rustfs_io_metrics::record_get_object_metadata_early_stop_hit(GET_OBJECT_PATH_LEGACY_DUPLEX, decision.reason); + rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses( + GET_OBJECT_PATH_LEGACY_DUPLEX, + saved_responses, + ); + while join_set.join_next().await.is_some() {} + let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); + return Ok((ress, errors, diagnostics)); + } + } + + rustfs_io_metrics::record_get_object_metadata_early_stop_miss( + GET_OBJECT_PATH_LEGACY_DUPLEX, + accumulator.final_miss_reason(), + ); + rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0); + let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); + Ok((ress, errors, diagnostics)) + } + + pub(in crate::set_disk) async fn read_all_xl( + disks: &[Option], + bucket: &str, + object: &str, + read_data: bool, + incl_free_vers: bool, + ) -> (Vec, Vec>) { + let (fileinfos, errs) = Self::read_all_raw_file_info(disks, bucket, object, read_data).await; + + Self::pick_latest_quorum_files_info(fileinfos, errs, bucket, object, read_data, incl_free_vers).await + } + + pub(crate) async fn load_file_info_versions_exact( + &self, + bucket: &str, + object: &str, + ) -> Result> { + let disks = self.get_disks_internal().await; + if disks.is_empty() { + return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object])); + } + + let read_quorum = disks.len().div_ceil(2).max(1); + let (raw_fileinfos, errs) = Self::read_all_raw_file_info(&disks, bucket, object, false).await; + + if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) { + let object_err = to_object_err(err.into(), vec![bucket, object]); + if is_err_object_not_found(&object_err) || is_err_version_not_found(&object_err) { + return Ok(None); + } + return Err(object_err); + } + + let mut shallow_versions = Vec::with_capacity(raw_fileinfos.len()); + for raw_fileinfo in raw_fileinfos.into_iter().flatten() { + let meta = FileMeta::load(&raw_fileinfo.buf) + .map_err(|err| Error::other(format!("exact object metadata decode failed for {bucket}/{object}: {err}")))?; + shallow_versions.push(meta.versions); + } + + if shallow_versions.len() < read_quorum { + return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object])); + } + + let versions = merge_file_meta_versions(read_quorum, true, 0, &shallow_versions); + if versions.is_empty() { + return Err(Error::other(format!( + "exact object metadata read returned no quorum versions for {bucket}/{object}" + ))); + } + + FileMeta { + versions, + ..Default::default() + } + .get_all_file_info_versions(bucket, object, true) + .map(Some) + .map_err(|err| Error::other(format!("exact object versions decode failed for {bucket}/{object}: {err}"))) + } + + pub(in crate::set_disk) async fn read_all_raw_file_info( + disks: &[Option], + bucket: &str, + object: &str, + read_data: bool, + ) -> (Vec>, Vec>) { + let mut ress = Vec::with_capacity(disks.len()); + let mut errors = Vec::with_capacity(disks.len()); + + let mut futures = Vec::with_capacity(disks.len()); + + for disk in disks.iter() { + futures.push(async move { + if let Some(disk) = disk { + disk.read_xl(bucket, object, read_data).await + } else { + Err(DiskError::DiskNotFound) + } + }); + } + + let results = join_all(futures).await; + for result in results { + match result { + Ok(res) => { + ress.push(Some(res)); + errors.push(None); + } + Err(e) => { + ress.push(None); + errors.push(Some(e)); + } + } + } + + (ress, errors) + } + + pub(in crate::set_disk) async fn pick_latest_quorum_files_info( + fileinfos: Vec>, + errs: Vec>, + bucket: &str, + object: &str, + read_data: bool, + incl_free_vers: bool, + ) -> (Vec, Vec>) { + let mut metadata_array = vec![None; fileinfos.len()]; + let mut meta_file_infos = vec![FileInfo::default(); fileinfos.len()]; + let mut metadata_shallow_versions = vec![None; fileinfos.len()]; + + let mut v2_bufs = { + if !read_data { + vec![Vec::new(); fileinfos.len()] + } else { + Vec::new() + } + }; + + let mut errs = errs; + + for (idx, info_op) in fileinfos.iter().enumerate() { + if let Some(info) = info_op { + if !read_data { + v2_bufs[idx] = info.buf.clone(); + } + + let xlmeta = match FileMeta::load(&info.buf) { + Ok(res) => res, + Err(err) => { + errs[idx] = Some(err.into()); + continue; + } + }; + + metadata_array[idx] = Some(xlmeta); + meta_file_infos[idx] = FileInfo::default(); + } + } + + for (idx, info_op) in metadata_array.iter().enumerate() { + if let Some(info) = info_op { + metadata_shallow_versions[idx] = Some(info.versions.clone()); + } + } + + let shallow_versions: Vec> = metadata_shallow_versions.iter().flatten().cloned().collect(); + + let read_quorum = fileinfos.len().div_ceil(2); + let versions = merge_file_meta_versions(read_quorum, false, 1, &shallow_versions); + let meta = FileMeta { + versions, + ..Default::default() + }; + + let finfo = match meta.into_fileinfo(bucket, object, "", true, incl_free_vers, true) { + Ok(res) => res, + Err(err) => { + for item in errs.iter_mut() { + if item.is_none() { + *item = Some(err.clone().into()); + } + } + + return (meta_file_infos, errs); + } + }; + + if !finfo.is_valid() { + for item in errs.iter_mut() { + if item.is_none() { + *item = Some(DiskError::FileCorrupt); + } + } + + return (meta_file_infos, errs); + } + + let vid = finfo.version_id.unwrap_or(Uuid::nil()); + + for (idx, meta_op) in metadata_array.iter().enumerate() { + if let Some(meta) = meta_op { + match meta.into_fileinfo(bucket, object, vid.to_string().as_str(), read_data, incl_free_vers, true) { + Ok(res) => meta_file_infos[idx] = res, + Err(err) => errs[idx] = Some(err.into()), + } + } + } + + (meta_file_infos, errs) + } + + pub(in crate::set_disk) async fn read_multiple_files( + disks: &[Option], + req: ReadMultipleReq, + read_quorum: usize, + ) -> Vec { + let mut futures = Vec::with_capacity(disks.len()); + let empty_quorum_result = || { + req.files + .iter() + .map(|want| ReadMultipleResp { + bucket: req.bucket.clone(), + prefix: req.prefix.clone(), + file: want.clone(), + exists: false, + error: Error::ErasureReadQuorum.to_string(), + data: Vec::new(), + mod_time: None, + }) + .collect::>() + }; + + for disk in disks.iter() { + let disk = disk.clone(); + let req = req.clone(); + futures.push(async move { + if let Some(disk) = disk { + disk.read_multiple(req).await + } else { + Err(DiskError::DiskNotFound) + } + }); + } + + let (ress, errors) = match collect_read_multiple_results(futures, read_quorum).await { + Ok(collected) => collected, + Err(()) => return empty_quorum_result(), + }; + + // debug!("ReadMultipleResp ress {:?}", ress); + // debug!("ReadMultipleResp errors {:?}", errors); + + let mut ret = Vec::with_capacity(req.files.len()); + + for want in req.files.iter() { + let mut quorum = 0; + + let mut get_res = ReadMultipleResp::default(); + + for res in ress.iter() { + if res.is_none() { + continue; + } + + let disk_res = res.as_ref().unwrap(); + + for resp in disk_res.iter() { + if !resp.error.is_empty() || !resp.exists { + continue; + } + + if &resp.file != want || resp.bucket != req.bucket || resp.prefix != req.prefix { + continue; + } + quorum += 1; + + if get_res.mod_time > resp.mod_time || get_res.data.len() > resp.data.len() { + continue; + } + + get_res = resp.clone(); + } + } + + if quorum < read_quorum { + // debug!("quorum < read_quorum: {} < {}", quorum, read_quorum); + get_res.exists = false; + get_res.error = Error::ErasureReadQuorum.to_string(); + get_res.data = Vec::new(); + } + + ret.push(get_res); + } + + // log err + + ret + } +} + +// =========================================================================== +// Write / rename / delete primitives (relocated verbatim from set_disk/write.rs, +// P5 step 2, tracking backlog#815, issue backlog#820). +// =========================================================================== + +/// Grace window during which a recently modified object is never deleted as +/// dangling. 0 disables the grace window. +const ENV_HEAL_DANGLING_DELETE_GRACE_SECS: &str = "RUSTFS_HEAL_DANGLING_DELETE_GRACE_SECS"; +const DEFAULT_HEAL_DANGLING_DELETE_GRACE_SECS: u64 = 3600; + +fn dangling_delete_grace() -> time::Duration { + let secs = rustfs_utils::get_env_u64(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, DEFAULT_HEAL_DANGLING_DELETE_GRACE_SECS); + time::Duration::seconds(i64::try_from(secs).unwrap_or(i64::MAX)) +} + +/// Result of scanning one disk's copy of a directory prefix while deciding +/// whether an orphan (metadata-less) directory tree can be safely purged. +enum OrphanDirScan { + /// The subtree holds at least one regular file (object metadata or data), so + /// it is a real object and must not be purged. + HasData, + /// The prefix exists on this disk and contains only nested empty directories. + /// Carries every directory path in pre-order (parents before children). + Empty(Vec), + /// The prefix does not exist on this disk. + Missing, +} + +impl SetDisks { + pub(in crate::set_disk) fn default_read_quorum(&self) -> usize { + self.set_drive_count - self.default_parity_count + } + + pub(in crate::set_disk) fn default_write_quorum(&self) -> usize { + let mut data_count = self.set_drive_count - self.default_parity_count; + if data_count == self.default_parity_count { + data_count += 1 + } + + data_count + } + + #[tracing::instrument(level = "debug", skip(disks, file_infos))] + #[allow(clippy::type_complexity)] + pub(in crate::set_disk) async fn rename_data( + disks: &[Option], + src_bucket: &str, + src_object: &str, + file_infos: &[FileInfo], + dst_bucket: &str, + dst_object: &str, + write_quorum: usize, + ) -> disk::error::Result<(Vec>, Option>, Option, Vec>)> { + let mut futures = Vec::with_capacity(disks.len()); + + let mut errs = Vec::with_capacity(disks.len()); + + let src_bucket = Arc::new(src_bucket.to_string()); + let src_object = Arc::new(src_object.to_string()); + let dst_bucket = Arc::new(dst_bucket.to_string()); + let dst_object = Arc::new(dst_object.to_string()); + + for (i, (disk, file_info)) in disks.iter().zip(file_infos.iter()).enumerate() { + let mut file_info = file_info.clone(); + let disk = disk.clone(); + let src_bucket = src_bucket.clone(); + let src_object = src_object.clone(); + let dst_object = dst_object.clone(); + let dst_bucket = dst_bucket.clone(); + + futures.push(tokio::spawn(async move { + if file_info.erasure.index == 0 { + file_info.erasure.index = i + 1; + } + + if !file_info.is_valid() { + return Err(DiskError::FileCorrupt); + } + + if let Some(disk) = disk { + disk.rename_data(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object) + .await + } else { + Err(DiskError::DiskNotFound) + } + })); + } + + let mut disk_versions = vec![None; disks.len()]; + let mut data_dirs = vec![None; disks.len()]; + + let results = join_all(futures).await; + + for (idx, result) in results.iter().enumerate() { + match result.as_ref().map_err(|_| DiskError::Unexpected)? { + Ok(res) => { + data_dirs[idx] = res.old_data_dir; + disk_versions[idx].clone_from(&res.sign); + errs.push(None); + } + Err(e) => { + errs.push(Some(e.clone())); + } + } + } + + if issue3031_diag_enabled() { + let success_count = errs.iter().filter(|err| err.is_none()).count(); + let failure_count = errs.len().saturating_sub(success_count); + let ignored_failure_count = errs + .iter() + .filter(|err| err.as_ref().is_some_and(|err| OBJECT_OP_IGNORED_ERRS.contains(err))) + .count(); + let data_dir_vote_count = data_dirs.iter().filter(|data_dir| data_dir.is_some()).count(); + let reduced_data_dir = Self::reduce_common_data_dir(&data_dirs, write_quorum); + warn!( + target: "rustfs_ecstore::set_disk", + src_bucket = %src_bucket, + src_object = %src_object, + dst_bucket = %dst_bucket, + dst_object = %dst_object, + write_quorum, + disk_count = errs.len(), + success_count, + failure_count, + ignored_failure_count, + data_dir_vote_count, + reduced_data_dir = ?reduced_data_dir, + errs = ?errs, + data_dirs = ?data_dirs, + "issue3031_rename_data_quorum_context" + ); + } + + let mut futures = Vec::with_capacity(disks.len()); + if let Some(ret_err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + // TODO: add concurrency + for (i, err) in errs.iter().enumerate() { + if err.is_some() { + continue; + } + + if let Some(disk) = disks[i].as_ref() { + let fi = file_infos[i].clone(); + let old_data_dir = data_dirs[i]; + let disk = disk.clone(); + let dst_bucket = dst_bucket.clone(); + let dst_object = dst_object.clone(); + futures.push(tokio::spawn(async move { + disk.delete_version( + &dst_bucket, + &dst_object, + fi, + false, + DeleteOptions { + undo_write: true, + old_data_dir, + ..Default::default() + }, + ) + .await + })); + } + } + + if issue3031_diag_enabled() { + warn!( + target: "rustfs_ecstore::set_disk", + src_bucket = %src_bucket, + src_object = %src_object, + dst_bucket = %dst_bucket, + dst_object = %dst_object, + write_quorum, + ret_err = %ret_err, + errs = ?errs, + data_dirs = ?data_dirs, + "issue3031_rename_data_quorum_failed" + ); + } + + let undo_results = join_all(futures).await; + let undo_error_count = undo_results + .iter() + .filter(|result| match result { + Err(_) | Ok(Err(_)) => true, + Ok(Ok(_)) => false, + }) + .count(); + if undo_error_count > 0 { + warn!( + target: "rustfs_ecstore::set_disk", + dst_bucket = %dst_bucket, + dst_object = %dst_object, + undo_error_count, + "rename_data quorum rollback reported errors" + ); + } + return Err(ret_err); + } + + let versions = None; + // TODO: reduceCommonVersions + + let data_dir = Self::reduce_common_data_dir(&data_dirs, write_quorum); + let online_disks = Self::eval_disks(disks, &errs); + let cleanup_disks = if let Some(data_dir) = data_dir { + disks + .iter() + .zip(errs.iter()) + .zip(data_dirs.iter()) + .map(|((disk, err), old_data_dir)| { + if err.is_none() && *old_data_dir == Some(data_dir) { + disk.clone() + } else { + None + } + }) + .collect() + } else { + vec![None; disks.len()] + }; + + // // TODO: reduce_common_data_dir + // if let Some(old_dir) = rename_ress + // .iter() + // .filter_map(|v| if v.is_some() { v.as_ref().unwrap().old_data_dir } else { None }) + // .map(|v| v.to_string()) + // .next() + // { + // let cm_errs = self.commit_rename_data_dir(&shuffle_disks, &bucket, &object, &old_dir).await; + // warn!("put_object commit_rename_data_dir:{:?}", &cm_errs); + // } + + // self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_dir).await?; + + Ok((online_disks, versions, data_dir, cleanup_disks)) + } + + #[allow(dead_code)] + #[tracing::instrument(level = "debug", skip(self, disks))] + pub(in crate::set_disk) async fn commit_rename_data_dir( + &self, + disks: &[Option], + bucket: &str, + object: &str, + data_dir: &str, + write_quorum: usize, + ) -> disk::error::Result<()> { + let file_path = Arc::new(format!("{object}/{data_dir}")); + let bucket = Arc::new(bucket.to_string()); + let futures = disks.iter().map(|disk| { + let file_path = file_path.clone(); + let bucket = bucket.clone(); + let disk = disk.clone(); + tokio::spawn(async move { + if let Some(disk) = disk { + (disk + .delete( + &bucket, + &file_path, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await) + .err() + } else { + Some(DiskError::DiskNotFound) + } + }) + }); + let errs: Vec> = join_all(futures) + .await + .into_iter() + .map(|e| e.unwrap_or(Some(DiskError::Unexpected))) + .collect(); + + if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + return Err(err); + } + + Ok(()) + } + + #[tracing::instrument(skip(self))] + pub(in crate::set_disk) async fn cleanup_multipart_path(&self, paths: &[String]) { + if paths.is_empty() { + return; + } + let disks = self.get_disks_internal().await; + + let mut errs = Vec::with_capacity(disks.len()); + + // Use improved simple batch processor instead of join_all for better performance + let processor = runtime_sources::batch_processors().write_processor(); + + let tasks: Vec<_> = disks + .iter() + .map(|disk| { + let disk = disk.clone(); + let paths = paths.to_vec(); + + async move { + if let Some(disk) = disk { + disk.delete_paths(RUSTFS_META_MULTIPART_BUCKET, &paths).await + } else { + Err(DiskError::DiskNotFound) + } + } + }) + .collect(); + + let results = processor.execute_batch(tasks).await; + for result in results { + match result { + Ok(_) => { + errs.push(None); + } + Err(e) => { + errs.push(Some(e)); + } + } + } + + if errs.iter().any(|e| e.is_some()) { + warn!("cleanup_multipart_path errs {:?}", &errs); + } + } + + #[tracing::instrument(skip(disks, meta))] + #[allow(clippy::too_many_arguments)] + pub(in crate::set_disk) async fn rename_part( + &self, + disks: &[Option], + src_bucket: &str, + src_object: &str, + dst_bucket: &str, + dst_object: &str, + meta: Bytes, + write_quorum: usize, + quorum_context: Option>, + ) -> disk::error::Result>> { + let src_bucket = Arc::new(src_bucket.to_string()); + let src_object = Arc::new(src_object.to_string()); + let dst_bucket = Arc::new(dst_bucket.to_string()); + let dst_object = Arc::new(dst_object.to_string()); + + // Match MinIO's multipart overwrite semantics: clear any stale destination + // part payload and metadata before the new per-disk rename fan-out begins. + self.cleanup_multipart_path(&[dst_object.to_string(), format!("{dst_object}.meta")]) + .await; + + let mut errs = Vec::with_capacity(disks.len()); + + let futures = disks.iter().map(|disk| { + let disk = disk.clone(); + let meta = meta.clone(); + let src_bucket = src_bucket.clone(); + let src_object = src_object.clone(); + let dst_bucket = dst_bucket.clone(); + let dst_object = dst_object.clone(); + tokio::spawn(async move { + if let Some(disk) = disk { + disk.rename_part(&src_bucket, &src_object, &dst_bucket, &dst_object, meta) + .await + } else { + Err(DiskError::DiskNotFound) + } + }) + }); + + let results = join_all(futures).await; + for result in results { + match result? { + Ok(_) => { + errs.push(None); + } + Err(e) => { + errs.push(Some(e)); + } + } + } + + if issue3031_diag_enabled() { + let success_count = errs.iter().filter(|err| err.is_none()).count(); + let error_count = errs.len().saturating_sub(success_count); + let disk_not_found_count = errs.iter().filter(|err| matches!(err, Some(DiskError::DiskNotFound))).count(); + let file_not_found_count = errs.iter().filter(|err| matches!(err, Some(DiskError::FileNotFound))).count(); + warn!( + target: "rustfs_ecstore::set_disk", + src_bucket = %src_bucket, + src_object = %src_object, + dst_bucket = %dst_bucket, + dst_object = %dst_object, + write_quorum = write_quorum, + disk_count = errs.len(), + success_count = success_count, + error_count = error_count, + disk_not_found_count = disk_not_found_count, + file_not_found_count = file_not_found_count, + "issue3031_rename_part_context" + ); + } + + if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + if let Some(context) = quorum_context { + log_multipart_write_quorum_failure(context, &errs, write_quorum, &err); + } else { + warn!("rename_part errs {:?}", &errs); + } + self.cleanup_multipart_path(&[dst_object.to_string(), format!("{dst_object}.meta")]) + .await; + return Err(err); + } + + let disks = Self::eval_disks(disks, &errs); + Ok(disks) + } + + pub(in crate::set_disk) fn eval_disks(disks: &[Option], errs: &[Option]) -> Vec> { + if disks.len() != errs.len() { + return Vec::new(); + } + + let mut online_disks = vec![None; disks.len()]; + + for (i, err_op) in errs.iter().enumerate() { + if err_op.is_none() { + online_disks[i].clone_from(&disks[i]); + } + } + + online_disks + } + + #[tracing::instrument(skip(disks, files))] + pub(in crate::set_disk) async fn write_unique_file_info( + disks: &[Option], + org_bucket: &str, + bucket: &str, + prefix: &str, + files: &[FileInfo], + write_quorum: usize, + ) -> disk::error::Result<()> { + let mut futures = Vec::with_capacity(disks.len()); + let mut errs = Vec::with_capacity(disks.len()); + + for (i, disk) in disks.iter().enumerate() { + let mut file_info = files[i].clone(); + file_info.erasure.index = i + 1; + futures.push(async move { + if let Some(disk) = disk { + disk.write_metadata(org_bucket, bucket, prefix, file_info).await + } else { + Err(DiskError::DiskNotFound) + } + }); + } + + let results = join_all(futures).await; + for result in results { + match result { + Ok(_) => { + errs.push(None); + } + Err(e) => { + errs.push(Some(e)); + } + } + } + + if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + let mut revert_futures = Vec::with_capacity(disks.len()); + for (i, err) in errs.iter().enumerate() { + if err.is_some() { + continue; + } + + if let Some(disk) = disks[i].as_ref() { + let disk = disk.clone(); + let bucket = bucket.to_string(); + let path = path_join_buf(&[prefix, STORAGE_FORMAT_FILE]); + revert_futures.push(async move { + if let Err(err) = disk + .delete( + &bucket, + &path, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await + { + warn!("write meta revert err {:?}", err); + } + }); + } + } + + join_all(revert_futures).await; + return Err(err); + } + Ok(()) + } + + pub(in crate::set_disk) async fn update_object_meta( + &self, + bucket: &str, + object: &str, + fi: FileInfo, + disks: &[Option], + ) -> disk::error::Result<()> { + self.update_object_meta_with_opts(bucket, object, fi, disks, &UpdateMetadataOpts::default()) + .await + } + + pub(in crate::set_disk) async fn update_object_meta_with_opts( + &self, + bucket: &str, + object: &str, + fi: FileInfo, + disks: &[Option], + opts: &UpdateMetadataOpts, + ) -> disk::error::Result<()> { + if fi.metadata.is_empty() && !opts.replace_user_metadata { + return Ok(()); + } + + self.invalidate_get_object_metadata_cache(bucket, object).await; + + let mut futures = Vec::with_capacity(disks.len()); + + let mut errs = Vec::with_capacity(disks.len()); + + for disk in disks.iter() { + let fi = fi.clone(); + futures.push(async move { + if let Some(disk) = disk { + disk.update_metadata(bucket, object, fi, opts).await + } else { + Err(DiskError::DiskNotFound) + } + }) + } + + let results = join_all(futures).await; + for result in results { + match result { + Ok(_) => { + errs.push(None); + } + Err(e) => { + errs.push(Some(e)); + } + } + } + + if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, fi.write_quorum(self.default_write_quorum())) { + return Err(err); + } + + self.invalidate_get_object_metadata_cache(bucket, object).await; + + Ok(()) + } + + pub(in crate::set_disk) async fn delete_if_dangling( + &self, + bucket: &str, + object: &str, + meta_arr: &[FileInfo], + errs: &[Option], + data_errs_by_part: &HashMap>, + opts: ObjectOptions, + ) -> disk::error::Result { + let (m, can_heal) = is_object_dangling(meta_arr, errs, data_errs_by_part); + + if !can_heal { + return Err(DiskError::ErasureReadQuorum); + } + + // Recently written objects get a grace window before dangling cleanup: after + // an unclean shutdown some disks may still be catching up (or carry writes + // that were never made durable), and deleting the surviving shards right away + // turns a partial loss into a total one. Skip deletion and leave the object + // for a later heal/scanner pass to re-evaluate. + if m.is_valid() + && let Some(mod_time) = m.mod_time + { + let grace = dangling_delete_grace(); + if !grace.is_zero() && OffsetDateTime::now_utc() - mod_time < grace { + info!( + bucket = bucket, + object = object, + mod_time = %mod_time, + grace_secs = grace.whole_seconds(), + "skipping dangling-object deletion within grace window" + ); + return Err(DiskError::ErasureReadQuorum); + } + } + + let mut tags: HashMap = HashMap::new(); + tags.insert("set".to_string(), self.set_index.to_string()); + tags.insert("pool".to_string(), self.pool_index.to_string()); + tags.insert("merrs".to_string(), join_errs(errs)); + tags.insert("derrs".to_string(), format!("{data_errs_by_part:?}")); + if m.is_valid() { + tags.insert("sz".to_string(), m.size.to_string()); + tags.insert( + "mt".to_string(), + m.mod_time + .as_ref() + .map_or(String::new(), |mod_time| mod_time.unix_timestamp().to_string()), + ); + tags.insert("d:p".to_string(), format!("{}:{}", m.erasure.data_blocks, m.erasure.parity_blocks)); + } else { + tags.insert("invalid".to_string(), "1".to_string()); + tags.insert( + "d:p".to_string(), + format!("{}:{}", self.set_drive_count - self.default_parity_count, self.default_parity_count), + ); + } + let mut offline = 0; + for (i, err) in errs.iter().enumerate() { + let mut found = false; + if let Some(err) = err + && err == &DiskError::DiskNotFound + { + found = true; + } + for p in data_errs_by_part { + if let Some(v) = p.1.get(i) + && *v == CHECK_PART_DISK_NOT_FOUND + { + found = true; + break; + } + } + + if found { + offline += 1; + } + } + + if offline > 0 { + tags.insert("offline".to_string(), offline.to_string()); + } + + let mut fi = FileInfo::default(); + if let Some(ref version_id) = opts.version_id { + fi.version_id = Uuid::parse_str(version_id).ok(); + } + + fi.set_tier_free_version_id(&Uuid::new_v4().to_string()); + + let disks = self.get_disks_internal().await; + + let mut futures = Vec::with_capacity(disks.len()); + for disk_op in disks.iter() { + let bucket = bucket.to_string(); + let object = object.to_string(); + let fi = fi.clone(); + futures.push(async move { + if let Some(disk) = disk_op { + disk.delete_version(&bucket, &object, fi, false, DeleteOptions::default()) + .await + } else { + Err(DiskError::DiskNotFound) + } + }); + } + + let results = join_all(futures).await; + let mut delete_errs = Vec::with_capacity(results.len()); + for (index, result) in results.into_iter().enumerate() { + let key = format!("ddisk-{index}"); + let already_absent = matches!( + errs.get(index).and_then(Option::as_ref), + Some(DiskError::FileNotFound | DiskError::FileVersionNotFound) + ); + match result { + Ok(_) => { + tags.insert(key, "".to_string()); + delete_errs.push(None); + } + Err(e) => { + tags.insert(key, e.to_string()); + if already_absent || matches!(&e, DiskError::FileNotFound | DiskError::FileVersionNotFound) { + delete_errs.push(None); + } else { + delete_errs.push(Some(e)); + } + } + } + } + + let write_quorum = if m.is_valid() { + m.write_quorum(self.default_write_quorum()) + } else { + self.default_write_quorum() + }; + if let Some(err) = reduce_write_quorum_errs(&delete_errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + return Err(err); + } + + Ok(m) + } + + pub(in crate::set_disk) async fn delete_prefix(&self, bucket: &str, prefix: &str) -> disk::error::Result<()> { + let disks = self.get_disks_internal().await; + let write_quorum = disks.len() / 2 + 1; + + let mut futures = Vec::with_capacity(disks.len()); + + for disk_op in disks.iter() { + let bucket = bucket.to_string(); + let prefix = prefix.to_string(); + futures.push(async move { + if let Some(disk) = disk_op { + disk.delete( + &bucket, + &prefix, + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) + .await + } else { + Ok(()) + } + }); + } + + let errs = join_all(futures).await.into_iter().map(|v| v.err()).collect::>(); + + if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { + return Err(err); + } + + Ok(()) + } + + /// Scan a single disk's copy of `prefix` and decide whether it is an orphan + /// (metadata-less) directory subtree. Walks the tree iteratively and returns + /// [`OrphanDirScan::HasData`] as soon as any regular file is found. + async fn scan_orphan_dir(disk: &DiskStore, bucket: &str, prefix: &str) -> OrphanDirScan { + let root = prefix.trim_end_matches(SLASH_SEPARATOR).to_string(); + let mut stack = vec![root.clone()]; + // Pre-order list of directories (a parent always precedes its descendants), + // so reversing it yields a safe children-first removal order. + let mut dirs: Vec = Vec::new(); + let mut existed = false; + + while let Some(dir) = stack.pop() { + let entries = match disk.list_dir("", bucket, &dir, 0).await { + Ok(entries) => entries, + Err(_) => { + // The root missing (or never existing) means there is nothing to + // purge on this disk. A nested directory vanishing mid-scan is a + // benign race, so skip it and keep walking. + if dir == root { + return OrphanDirScan::Missing; + } + continue; + } + }; + + existed = true; + dirs.push(dir.clone()); + + for entry in entries { + match entry.strip_suffix(SLASH_SEPARATOR) { + // `read_dir` marks directories with a trailing slash; anything else + // is a regular file, which means real object data lives here. + Some(child) => stack.push(format!("{dir}{SLASH_SEPARATOR}{child}")), + None => return OrphanDirScan::HasData, + } + } + } + + if existed { + OrphanDirScan::Empty(dirs) + } else { + OrphanDirScan::Missing + } + } + + /// Purge an orphan directory prefix — a trailing-slash key that exists on disk + /// as an empty directory tree with no object metadata on any disk of this set. + /// Such prefixes are listable (see `scan_dir`) yet are not real objects, so the + /// normal delete path returns NotFound and leaves them stranded (issue #4189). + /// + /// Callers pass the *decoded* directory name (`prefix/`), not the `__XLDIR__` + /// encoded object key — the orphan tree on disk uses the plain path. + /// + /// Returns `Ok(true)` when the prefix was an orphan tree on this set and was + /// removed, `Ok(false)` when it holds real data or does not exist on any disk + /// of this set (the caller should surface the original NotFound), and `Err` on + /// a hard disk failure. + pub(crate) async fn purge_orphan_dir_object(&self, bucket: &str, object: &str) -> disk::error::Result { + let disks = self.get_disks_internal().await; + + // Phase 1: classify every online disk. Refuse to purge if ANY disk holds + // object data under the prefix, so a degraded/healable object is never + // destroyed. + let mut per_disk_dirs: Vec<(usize, Vec)> = Vec::new(); + let mut existed = false; + for (i, disk) in disks.iter().enumerate() { + let Some(disk) = disk else { continue }; + match Self::scan_orphan_dir(disk, bucket, object).await { + OrphanDirScan::HasData => return Ok(false), + OrphanDirScan::Empty(dirs) => { + existed = true; + per_disk_dirs.push((i, dirs)); + } + OrphanDirScan::Missing => {} + } + } + + if !existed { + return Ok(false); + } + + // Phase 2: remove the empty directories children-first on each disk. A + // non-recursive delete performs an empty-only `rmdir`, so a directory that + // concurrently gained an object fails with DirectoryNotEmpty and is skipped — + // a racing PutObject is never clobbered. + for (i, mut dirs) in per_disk_dirs { + let Some(disk) = disks[i].as_ref() else { continue }; + dirs.reverse(); + for dir in dirs { + if let Err(err) = disk + .delete( + bucket, + &dir, + DeleteOptions { + recursive: false, + immediate: true, + ..Default::default() + }, + ) + .await + { + // Best effort: a sibling removal may have already cleared a shared + // parent, or a concurrent writer repopulated the directory. Neither + // is fatal to purging the orphan tree. + debug!(bucket, object, dir, error = ?err, "purge_orphan_dir_object: skipped non-empty/absent directory"); + } + } + } + + Ok(true) + } + + pub(in crate::set_disk) async fn check_write_precondition( + &self, + bucket: &str, + object: &str, + opts: &ObjectOptions, + ) -> Option { + let mut opts = opts.clone(); + + let http_preconditions = opts.http_preconditions?; + opts.http_preconditions = None; + + // Never claim a lock here, to avoid deadlock + // - If no_lock is false, we must have obtained the lock out side of this function + // - If no_lock is true, we should not obtain locks + opts.no_lock = true; + let oi = self.get_object_info(bucket, object, &opts).await; + + match oi { + Ok(oi) => { + // If top level is a delete marker proceed to upload. + if oi.delete_marker { + return None; + } + let if_none_match = http_preconditions.if_none_match_value().map(str::to_owned); + let if_match = http_preconditions.if_match_value().map(str::to_owned); + if should_prevent_write(&oi, if_none_match, if_match) { + return Some(StorageError::PreconditionFailed); + } + } + + Err(StorageError::VersionNotFound(_, _, _)) + | Err(StorageError::ObjectNotFound(_, _)) + | Err(StorageError::ErasureReadQuorum) => { + // When the object is not found, + // - if If-Match is set, we should return 404 NotFound + // - if If-None-Match is set, we should be able to proceed with the request + if http_preconditions.if_match_value().is_some() { + return Some(StorageError::ObjectNotFound(bucket.to_string(), object.to_string())); + } + } + + Err(e) => { + return Some(e); + } + } + + None + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn dangling_delete_grace_defaults_to_one_hour() { + temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, None::<&str>, || { + assert_eq!(dangling_delete_grace(), time::Duration::seconds(3600)); + }); + } + + #[test] + fn dangling_delete_grace_env_override_and_disable() { + temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, Some("120"), || { + assert_eq!(dangling_delete_grace(), time::Duration::seconds(120)); + }); + temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, Some("0"), || { + assert!(dangling_delete_grace().is_zero()); + }); + } +} diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 3c162243c..839f17f1c 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -451,7 +451,6 @@ mod ops; mod read; mod replication; pub(crate) mod shard_source; -mod write; /// Get lock acquire timeout from environment variable RUSTFS_LOCK_ACQUIRE_TIMEOUT (in seconds) /// Defaults to 30 seconds if not set or invalid diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 2fece0c14..107e1addc 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -120,344 +120,6 @@ impl SetDisks { .await; } - pub(super) async fn read_parts( - disks: &[Option], - bucket: &str, - part_meta_paths: &[String], - part_numbers: &[usize], - read_quorum: usize, - ) -> disk::error::Result> { - let bucket = bucket.to_string(); - let part_meta_paths = part_meta_paths.to_vec(); - - let tasks: Vec<_> = disks - .iter() - .map(|disk| { - let disk = disk.clone(); - let bucket = bucket.clone(); - let part_meta_paths = part_meta_paths.clone(); - - async move { - if let Some(disk) = disk { - disk.read_parts(&bucket, &part_meta_paths).await - } else { - Err(DiskError::DiskNotFound) - } - } - }) - .collect(); - - let (responses, collected_errors) = match collect_read_parts_results(tasks, read_quorum).await { - Ok(collected) => collected, - Err(()) => return Err(DiskError::ErasureReadQuorum), - }; - - if let Some(err) = reduce_read_quorum_errs(&collected_errors, OBJECT_OP_IGNORED_ERRS, read_quorum) { - return Err(err); - } - - let mut ret = vec![ObjectPartInfo::default(); part_meta_paths.len()]; - - for (part_idx, part_info) in part_meta_paths.iter().enumerate() { - ret[part_idx] = resolve_read_part_from_responses( - &bucket, - part_info, - part_numbers[part_idx], - part_idx, - part_meta_paths.len(), - &responses, - read_quorum, - )?; - } - - Ok(ret) - } - - #[allow(clippy::too_many_arguments)] - #[tracing::instrument(level = "debug", skip(disks))] - pub(super) async fn read_all_fileinfo( - disks: &[Option], - org_bucket: &str, - bucket: &str, - object: &str, - version_id: &str, - read_data: bool, - healing: bool, - incl_free_versions: bool, - ) -> disk::error::Result<(Vec, Vec>)> { - let (ress, errors, _) = Self::read_all_fileinfo_inner( - disks, - org_bucket, - bucket, - object, - version_id, - read_data, - healing, - incl_free_versions, - false, - 0, - ) - .await?; - Ok((ress, errors)) - } - - #[allow(clippy::too_many_arguments)] - async fn read_all_fileinfo_observed( - disks: &[Option], - org_bucket: &str, - bucket: &str, - object: &str, - version_id: &str, - read_data: bool, - healing: bool, - incl_free_versions: bool, - default_parity_count: usize, - ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { - Self::read_all_fileinfo_inner( - disks, - org_bucket, - bucket, - object, - version_id, - read_data, - healing, - incl_free_versions, - true, - default_parity_count, - ) - .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, - default_parity_count: usize, - ) -> disk::error::Result<(Vec, Vec>, MetadataFanoutDiagnostics)> { - let early_stop_enabled = observe && (is_get_metadata_early_stop_enabled() || is_version_early_stop_enabled()); - let allow_early_stop = observe && should_allow_metadata_early_stop(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, - default_parity_count, - ) - .await; - } - if early_stop_enabled { - rustfs_io_metrics::record_get_object_metadata_early_stop_miss( - GET_OBJECT_PATH_LEGACY_DUPLEX, - GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST, - ); - rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0); - } - - Self::read_all_fileinfo_full_wait( - disks, - org_bucket, - bucket, - object, - version_id, - read_data, - healing, - incl_free_versions, - observe, - ) - .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, - ) -> 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 opts = Arc::new(ReadOptions { - incl_free_versions, - read_data, - healing, - }); - let org_bucket = Arc::new(org_bucket.to_string()); - let bucket = Arc::new(bucket.to_string()); - let object = Arc::new(object.to_string()); - let version_id = Arc::new(version_id.to_string()); - let futures = disks.iter().map(|disk| { - let disk = disk.clone(); - let opts = opts.clone(); - let org_bucket = org_bucket.clone(); - let bucket = bucket.clone(); - let object = object.clone(); - let version_id = version_id.clone(); - tokio::spawn(async move { - let response_start = observe.then(Instant::now); - let result = if let Some(disk) = disk { - disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).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(), - }; - 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, - default_parity_count: usize, - ) -> 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 = Arc::new(ReadOptions { - incl_free_versions, - read_data, - healing, - }); - let org_bucket = Arc::new(org_bucket.to_string()); - let bucket = Arc::new(bucket.to_string()); - let object = Arc::new(object.to_string()); - let version_id = Arc::new(version_id.to_string()); - let mut join_set = JoinSet::new(); - - for (index, disk) in disks.iter().cloned().enumerate() { - let opts = opts.clone(); - let org_bucket = org_bucket.clone(); - let bucket = bucket.clone(); - let object = object.clone(); - let version_id = version_id.clone(); - join_set.spawn(async move { - let response_start = Instant::now(); - let result = if let Some(disk) = disk { - disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await - } else { - Err(DiskError::DiskNotFound) - }; - (index, result, response_start.elapsed()) - }); - } - - while let Some(result) = join_set.join_next().await { - match result { - Ok((index, res, elapsed)) => match res { - Ok(file_info) => { - observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed)); - accumulator.observe_file_info(&file_info); - 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 let Some(decision) = accumulator - .early_stop_decision() - .or_else(|| accumulator.version_early_stop_decision()) - { - let saved_responses = join_set.len(); - join_set.abort_all(); - rustfs_io_metrics::record_get_object_metadata_early_stop_hit(GET_OBJECT_PATH_LEGACY_DUPLEX, decision.reason); - rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses( - GET_OBJECT_PATH_LEGACY_DUPLEX, - saved_responses, - ); - while join_set.join_next().await.is_some() {} - let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); - return Ok((ress, errors, diagnostics)); - } - } - - rustfs_io_metrics::record_get_object_metadata_early_stop_miss( - GET_OBJECT_PATH_LEGACY_DUPLEX, - accumulator.final_miss_reason(), - ); - rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0); - let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations); - Ok((ress, errors, diagnostics)) - } - pub async fn read_version_optimized( &self, bucket: &str, @@ -498,285 +160,6 @@ impl SetDisks { } } - pub(super) async fn read_all_xl( - disks: &[Option], - bucket: &str, - object: &str, - read_data: bool, - incl_free_vers: bool, - ) -> (Vec, Vec>) { - let (fileinfos, errs) = Self::read_all_raw_file_info(disks, bucket, object, read_data).await; - - Self::pick_latest_quorum_files_info(fileinfos, errs, bucket, object, read_data, incl_free_vers).await - } - - pub(crate) async fn load_file_info_versions_exact( - &self, - bucket: &str, - object: &str, - ) -> Result> { - let disks = self.get_disks_internal().await; - if disks.is_empty() { - return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object])); - } - - let read_quorum = disks.len().div_ceil(2).max(1); - let (raw_fileinfos, errs) = Self::read_all_raw_file_info(&disks, bucket, object, false).await; - - if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) { - let object_err = to_object_err(err.into(), vec![bucket, object]); - if is_err_object_not_found(&object_err) || is_err_version_not_found(&object_err) { - return Ok(None); - } - return Err(object_err); - } - - let mut shallow_versions = Vec::with_capacity(raw_fileinfos.len()); - for raw_fileinfo in raw_fileinfos.into_iter().flatten() { - let meta = FileMeta::load(&raw_fileinfo.buf) - .map_err(|err| Error::other(format!("exact object metadata decode failed for {bucket}/{object}: {err}")))?; - shallow_versions.push(meta.versions); - } - - if shallow_versions.len() < read_quorum { - return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object])); - } - - let versions = merge_file_meta_versions(read_quorum, true, 0, &shallow_versions); - if versions.is_empty() { - return Err(Error::other(format!( - "exact object metadata read returned no quorum versions for {bucket}/{object}" - ))); - } - - FileMeta { - versions, - ..Default::default() - } - .get_all_file_info_versions(bucket, object, true) - .map(Some) - .map_err(|err| Error::other(format!("exact object versions decode failed for {bucket}/{object}: {err}"))) - } - - pub(super) async fn read_all_raw_file_info( - disks: &[Option], - bucket: &str, - object: &str, - read_data: bool, - ) -> (Vec>, Vec>) { - let mut ress = Vec::with_capacity(disks.len()); - let mut errors = Vec::with_capacity(disks.len()); - - let mut futures = Vec::with_capacity(disks.len()); - - for disk in disks.iter() { - futures.push(async move { - if let Some(disk) = disk { - disk.read_xl(bucket, object, read_data).await - } else { - Err(DiskError::DiskNotFound) - } - }); - } - - let results = join_all(futures).await; - for result in results { - match result { - Ok(res) => { - ress.push(Some(res)); - errors.push(None); - } - Err(e) => { - ress.push(None); - errors.push(Some(e)); - } - } - } - - (ress, errors) - } - - pub(super) async fn pick_latest_quorum_files_info( - fileinfos: Vec>, - errs: Vec>, - bucket: &str, - object: &str, - read_data: bool, - incl_free_vers: bool, - ) -> (Vec, Vec>) { - let mut metadata_array = vec![None; fileinfos.len()]; - let mut meta_file_infos = vec![FileInfo::default(); fileinfos.len()]; - let mut metadata_shallow_versions = vec![None; fileinfos.len()]; - - let mut v2_bufs = { - if !read_data { - vec![Vec::new(); fileinfos.len()] - } else { - Vec::new() - } - }; - - let mut errs = errs; - - for (idx, info_op) in fileinfos.iter().enumerate() { - if let Some(info) = info_op { - if !read_data { - v2_bufs[idx] = info.buf.clone(); - } - - let xlmeta = match FileMeta::load(&info.buf) { - Ok(res) => res, - Err(err) => { - errs[idx] = Some(err.into()); - continue; - } - }; - - metadata_array[idx] = Some(xlmeta); - meta_file_infos[idx] = FileInfo::default(); - } - } - - for (idx, info_op) in metadata_array.iter().enumerate() { - if let Some(info) = info_op { - metadata_shallow_versions[idx] = Some(info.versions.clone()); - } - } - - let shallow_versions: Vec> = metadata_shallow_versions.iter().flatten().cloned().collect(); - - let read_quorum = fileinfos.len().div_ceil(2); - let versions = merge_file_meta_versions(read_quorum, false, 1, &shallow_versions); - let meta = FileMeta { - versions, - ..Default::default() - }; - - let finfo = match meta.into_fileinfo(bucket, object, "", true, incl_free_vers, true) { - Ok(res) => res, - Err(err) => { - for item in errs.iter_mut() { - if item.is_none() { - *item = Some(err.clone().into()); - } - } - - return (meta_file_infos, errs); - } - }; - - if !finfo.is_valid() { - for item in errs.iter_mut() { - if item.is_none() { - *item = Some(DiskError::FileCorrupt); - } - } - - return (meta_file_infos, errs); - } - - let vid = finfo.version_id.unwrap_or(Uuid::nil()); - - for (idx, meta_op) in metadata_array.iter().enumerate() { - if let Some(meta) = meta_op { - match meta.into_fileinfo(bucket, object, vid.to_string().as_str(), read_data, incl_free_vers, true) { - Ok(res) => meta_file_infos[idx] = res, - Err(err) => errs[idx] = Some(err.into()), - } - } - } - - (meta_file_infos, errs) - } - - pub(super) async fn read_multiple_files( - disks: &[Option], - req: ReadMultipleReq, - read_quorum: usize, - ) -> Vec { - let mut futures = Vec::with_capacity(disks.len()); - let empty_quorum_result = || { - req.files - .iter() - .map(|want| ReadMultipleResp { - bucket: req.bucket.clone(), - prefix: req.prefix.clone(), - file: want.clone(), - exists: false, - error: Error::ErasureReadQuorum.to_string(), - data: Vec::new(), - mod_time: None, - }) - .collect::>() - }; - - for disk in disks.iter() { - let disk = disk.clone(); - let req = req.clone(); - futures.push(async move { - if let Some(disk) = disk { - disk.read_multiple(req).await - } else { - Err(DiskError::DiskNotFound) - } - }); - } - - let (ress, errors) = match collect_read_multiple_results(futures, read_quorum).await { - Ok(collected) => collected, - Err(()) => return empty_quorum_result(), - }; - - // debug!("ReadMultipleResp ress {:?}", ress); - // debug!("ReadMultipleResp errors {:?}", errors); - - let mut ret = Vec::with_capacity(req.files.len()); - - for want in req.files.iter() { - let mut quorum = 0; - - let mut get_res = ReadMultipleResp::default(); - - for res in ress.iter() { - if res.is_none() { - continue; - } - - let disk_res = res.as_ref().unwrap(); - - for resp in disk_res.iter() { - if !resp.error.is_empty() || !resp.exists { - continue; - } - - if &resp.file != want || resp.bucket != req.bucket || resp.prefix != req.prefix { - continue; - } - quorum += 1; - - if get_res.mod_time > resp.mod_time || get_res.data.len() > resp.data.len() { - continue; - } - - get_res = resp.clone(); - } - } - - if quorum < read_quorum { - // debug!("quorum < read_quorum: {} < {}", quorum, read_quorum); - get_res.exists = false; - get_res.error = Error::ErasureReadQuorum.to_string(); - get_res.data = Vec::new(); - } - - ret.push(get_res); - } - - // log err - - ret - } - #[tracing::instrument(level = "debug", skip(self))] pub(super) async fn get_object_fileinfo( &self, @@ -1772,15 +1155,6 @@ impl SetDisks { } } -fn should_allow_metadata_early_stop(read_data: bool, version_id: &str, healing: bool, incl_free_versions: bool) -> bool { - if read_data { - 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) -} - fn get_object_metadata_cache_request_bypass_reason(bucket: &str, opts: &ObjectOptions, read_data: bool) -> Option<&'static str> { if !read_data { return Some(GET_METADATA_CACHE_REASON_NOT_READ_DATA); diff --git a/crates/ecstore/src/set_disk/write.rs b/crates/ecstore/src/set_disk/write.rs deleted file mode 100644 index 365c199b1..000000000 --- a/crates/ecstore/src/set_disk/write.rs +++ /dev/null @@ -1,941 +0,0 @@ -// 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. - -use super::*; - -/// Grace window during which a recently modified object is never deleted as -/// dangling. 0 disables the grace window. -const ENV_HEAL_DANGLING_DELETE_GRACE_SECS: &str = "RUSTFS_HEAL_DANGLING_DELETE_GRACE_SECS"; -const DEFAULT_HEAL_DANGLING_DELETE_GRACE_SECS: u64 = 3600; - -fn dangling_delete_grace() -> time::Duration { - let secs = rustfs_utils::get_env_u64(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, DEFAULT_HEAL_DANGLING_DELETE_GRACE_SECS); - time::Duration::seconds(i64::try_from(secs).unwrap_or(i64::MAX)) -} - -/// Result of scanning one disk's copy of a directory prefix while deciding -/// whether an orphan (metadata-less) directory tree can be safely purged. -enum OrphanDirScan { - /// The subtree holds at least one regular file (object metadata or data), so - /// it is a real object and must not be purged. - HasData, - /// The prefix exists on this disk and contains only nested empty directories. - /// Carries every directory path in pre-order (parents before children). - Empty(Vec), - /// The prefix does not exist on this disk. - Missing, -} - -impl SetDisks { - pub(super) fn default_read_quorum(&self) -> usize { - self.set_drive_count - self.default_parity_count - } - - pub(super) fn default_write_quorum(&self) -> usize { - let mut data_count = self.set_drive_count - self.default_parity_count; - if data_count == self.default_parity_count { - data_count += 1 - } - - data_count - } - - #[tracing::instrument(level = "debug", skip(disks, file_infos))] - #[allow(clippy::type_complexity)] - pub(super) async fn rename_data( - disks: &[Option], - src_bucket: &str, - src_object: &str, - file_infos: &[FileInfo], - dst_bucket: &str, - dst_object: &str, - write_quorum: usize, - ) -> disk::error::Result<(Vec>, Option>, Option, Vec>)> { - let mut futures = Vec::with_capacity(disks.len()); - - let mut errs = Vec::with_capacity(disks.len()); - - let src_bucket = Arc::new(src_bucket.to_string()); - let src_object = Arc::new(src_object.to_string()); - let dst_bucket = Arc::new(dst_bucket.to_string()); - let dst_object = Arc::new(dst_object.to_string()); - - for (i, (disk, file_info)) in disks.iter().zip(file_infos.iter()).enumerate() { - let mut file_info = file_info.clone(); - let disk = disk.clone(); - let src_bucket = src_bucket.clone(); - let src_object = src_object.clone(); - let dst_object = dst_object.clone(); - let dst_bucket = dst_bucket.clone(); - - futures.push(tokio::spawn(async move { - if file_info.erasure.index == 0 { - file_info.erasure.index = i + 1; - } - - if !file_info.is_valid() { - return Err(DiskError::FileCorrupt); - } - - if let Some(disk) = disk { - disk.rename_data(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object) - .await - } else { - Err(DiskError::DiskNotFound) - } - })); - } - - let mut disk_versions = vec![None; disks.len()]; - let mut data_dirs = vec![None; disks.len()]; - - let results = join_all(futures).await; - - for (idx, result) in results.iter().enumerate() { - match result.as_ref().map_err(|_| DiskError::Unexpected)? { - Ok(res) => { - data_dirs[idx] = res.old_data_dir; - disk_versions[idx].clone_from(&res.sign); - errs.push(None); - } - Err(e) => { - errs.push(Some(e.clone())); - } - } - } - - if issue3031_diag_enabled() { - let success_count = errs.iter().filter(|err| err.is_none()).count(); - let failure_count = errs.len().saturating_sub(success_count); - let ignored_failure_count = errs - .iter() - .filter(|err| err.as_ref().is_some_and(|err| OBJECT_OP_IGNORED_ERRS.contains(err))) - .count(); - let data_dir_vote_count = data_dirs.iter().filter(|data_dir| data_dir.is_some()).count(); - let reduced_data_dir = Self::reduce_common_data_dir(&data_dirs, write_quorum); - warn!( - target: "rustfs_ecstore::set_disk", - src_bucket = %src_bucket, - src_object = %src_object, - dst_bucket = %dst_bucket, - dst_object = %dst_object, - write_quorum, - disk_count = errs.len(), - success_count, - failure_count, - ignored_failure_count, - data_dir_vote_count, - reduced_data_dir = ?reduced_data_dir, - errs = ?errs, - data_dirs = ?data_dirs, - "issue3031_rename_data_quorum_context" - ); - } - - let mut futures = Vec::with_capacity(disks.len()); - if let Some(ret_err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - // TODO: add concurrency - for (i, err) in errs.iter().enumerate() { - if err.is_some() { - continue; - } - - if let Some(disk) = disks[i].as_ref() { - let fi = file_infos[i].clone(); - let old_data_dir = data_dirs[i]; - let disk = disk.clone(); - let dst_bucket = dst_bucket.clone(); - let dst_object = dst_object.clone(); - futures.push(tokio::spawn(async move { - disk.delete_version( - &dst_bucket, - &dst_object, - fi, - false, - DeleteOptions { - undo_write: true, - old_data_dir, - ..Default::default() - }, - ) - .await - })); - } - } - - if issue3031_diag_enabled() { - warn!( - target: "rustfs_ecstore::set_disk", - src_bucket = %src_bucket, - src_object = %src_object, - dst_bucket = %dst_bucket, - dst_object = %dst_object, - write_quorum, - ret_err = %ret_err, - errs = ?errs, - data_dirs = ?data_dirs, - "issue3031_rename_data_quorum_failed" - ); - } - - let undo_results = join_all(futures).await; - let undo_error_count = undo_results - .iter() - .filter(|result| match result { - Err(_) | Ok(Err(_)) => true, - Ok(Ok(_)) => false, - }) - .count(); - if undo_error_count > 0 { - warn!( - target: "rustfs_ecstore::set_disk", - dst_bucket = %dst_bucket, - dst_object = %dst_object, - undo_error_count, - "rename_data quorum rollback reported errors" - ); - } - return Err(ret_err); - } - - let versions = None; - // TODO: reduceCommonVersions - - let data_dir = Self::reduce_common_data_dir(&data_dirs, write_quorum); - let online_disks = Self::eval_disks(disks, &errs); - let cleanup_disks = if let Some(data_dir) = data_dir { - disks - .iter() - .zip(errs.iter()) - .zip(data_dirs.iter()) - .map(|((disk, err), old_data_dir)| { - if err.is_none() && *old_data_dir == Some(data_dir) { - disk.clone() - } else { - None - } - }) - .collect() - } else { - vec![None; disks.len()] - }; - - // // TODO: reduce_common_data_dir - // if let Some(old_dir) = rename_ress - // .iter() - // .filter_map(|v| if v.is_some() { v.as_ref().unwrap().old_data_dir } else { None }) - // .map(|v| v.to_string()) - // .next() - // { - // let cm_errs = self.commit_rename_data_dir(&shuffle_disks, &bucket, &object, &old_dir).await; - // warn!("put_object commit_rename_data_dir:{:?}", &cm_errs); - // } - - // self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_dir).await?; - - Ok((online_disks, versions, data_dir, cleanup_disks)) - } - - #[allow(dead_code)] - #[tracing::instrument(level = "debug", skip(self, disks))] - pub(super) async fn commit_rename_data_dir( - &self, - disks: &[Option], - bucket: &str, - object: &str, - data_dir: &str, - write_quorum: usize, - ) -> disk::error::Result<()> { - let file_path = Arc::new(format!("{object}/{data_dir}")); - let bucket = Arc::new(bucket.to_string()); - let futures = disks.iter().map(|disk| { - let file_path = file_path.clone(); - let bucket = bucket.clone(); - let disk = disk.clone(); - tokio::spawn(async move { - if let Some(disk) = disk { - (disk - .delete( - &bucket, - &file_path, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await) - .err() - } else { - Some(DiskError::DiskNotFound) - } - }) - }); - let errs: Vec> = join_all(futures) - .await - .into_iter() - .map(|e| e.unwrap_or(Some(DiskError::Unexpected))) - .collect(); - - if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - return Err(err); - } - - Ok(()) - } - - #[tracing::instrument(skip(self))] - pub(super) async fn cleanup_multipart_path(&self, paths: &[String]) { - if paths.is_empty() { - return; - } - let disks = self.get_disks_internal().await; - - let mut errs = Vec::with_capacity(disks.len()); - - // Use improved simple batch processor instead of join_all for better performance - let processor = runtime_sources::batch_processors().write_processor(); - - let tasks: Vec<_> = disks - .iter() - .map(|disk| { - let disk = disk.clone(); - let paths = paths.to_vec(); - - async move { - if let Some(disk) = disk { - disk.delete_paths(RUSTFS_META_MULTIPART_BUCKET, &paths).await - } else { - Err(DiskError::DiskNotFound) - } - } - }) - .collect(); - - let results = processor.execute_batch(tasks).await; - for result in results { - match result { - Ok(_) => { - errs.push(None); - } - Err(e) => { - errs.push(Some(e)); - } - } - } - - if errs.iter().any(|e| e.is_some()) { - warn!("cleanup_multipart_path errs {:?}", &errs); - } - } - - #[tracing::instrument(skip(disks, meta))] - #[allow(clippy::too_many_arguments)] - pub(super) async fn rename_part( - &self, - disks: &[Option], - src_bucket: &str, - src_object: &str, - dst_bucket: &str, - dst_object: &str, - meta: Bytes, - write_quorum: usize, - quorum_context: Option>, - ) -> disk::error::Result>> { - let src_bucket = Arc::new(src_bucket.to_string()); - let src_object = Arc::new(src_object.to_string()); - let dst_bucket = Arc::new(dst_bucket.to_string()); - let dst_object = Arc::new(dst_object.to_string()); - - // Match MinIO's multipart overwrite semantics: clear any stale destination - // part payload and metadata before the new per-disk rename fan-out begins. - self.cleanup_multipart_path(&[dst_object.to_string(), format!("{dst_object}.meta")]) - .await; - - let mut errs = Vec::with_capacity(disks.len()); - - let futures = disks.iter().map(|disk| { - let disk = disk.clone(); - let meta = meta.clone(); - let src_bucket = src_bucket.clone(); - let src_object = src_object.clone(); - let dst_bucket = dst_bucket.clone(); - let dst_object = dst_object.clone(); - tokio::spawn(async move { - if let Some(disk) = disk { - disk.rename_part(&src_bucket, &src_object, &dst_bucket, &dst_object, meta) - .await - } else { - Err(DiskError::DiskNotFound) - } - }) - }); - - let results = join_all(futures).await; - for result in results { - match result? { - Ok(_) => { - errs.push(None); - } - Err(e) => { - errs.push(Some(e)); - } - } - } - - if issue3031_diag_enabled() { - let success_count = errs.iter().filter(|err| err.is_none()).count(); - let error_count = errs.len().saturating_sub(success_count); - let disk_not_found_count = errs.iter().filter(|err| matches!(err, Some(DiskError::DiskNotFound))).count(); - let file_not_found_count = errs.iter().filter(|err| matches!(err, Some(DiskError::FileNotFound))).count(); - warn!( - target: "rustfs_ecstore::set_disk", - src_bucket = %src_bucket, - src_object = %src_object, - dst_bucket = %dst_bucket, - dst_object = %dst_object, - write_quorum = write_quorum, - disk_count = errs.len(), - success_count = success_count, - error_count = error_count, - disk_not_found_count = disk_not_found_count, - file_not_found_count = file_not_found_count, - "issue3031_rename_part_context" - ); - } - - if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - if let Some(context) = quorum_context { - log_multipart_write_quorum_failure(context, &errs, write_quorum, &err); - } else { - warn!("rename_part errs {:?}", &errs); - } - self.cleanup_multipart_path(&[dst_object.to_string(), format!("{dst_object}.meta")]) - .await; - return Err(err); - } - - let disks = Self::eval_disks(disks, &errs); - Ok(disks) - } - - pub(super) fn eval_disks(disks: &[Option], errs: &[Option]) -> Vec> { - if disks.len() != errs.len() { - return Vec::new(); - } - - let mut online_disks = vec![None; disks.len()]; - - for (i, err_op) in errs.iter().enumerate() { - if err_op.is_none() { - online_disks[i].clone_from(&disks[i]); - } - } - - online_disks - } - - #[tracing::instrument(skip(disks, files))] - pub(super) async fn write_unique_file_info( - disks: &[Option], - org_bucket: &str, - bucket: &str, - prefix: &str, - files: &[FileInfo], - write_quorum: usize, - ) -> disk::error::Result<()> { - let mut futures = Vec::with_capacity(disks.len()); - let mut errs = Vec::with_capacity(disks.len()); - - for (i, disk) in disks.iter().enumerate() { - let mut file_info = files[i].clone(); - file_info.erasure.index = i + 1; - futures.push(async move { - if let Some(disk) = disk { - disk.write_metadata(org_bucket, bucket, prefix, file_info).await - } else { - Err(DiskError::DiskNotFound) - } - }); - } - - let results = join_all(futures).await; - for result in results { - match result { - Ok(_) => { - errs.push(None); - } - Err(e) => { - errs.push(Some(e)); - } - } - } - - if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - let mut revert_futures = Vec::with_capacity(disks.len()); - for (i, err) in errs.iter().enumerate() { - if err.is_some() { - continue; - } - - if let Some(disk) = disks[i].as_ref() { - let disk = disk.clone(); - let bucket = bucket.to_string(); - let path = path_join_buf(&[prefix, STORAGE_FORMAT_FILE]); - revert_futures.push(async move { - if let Err(err) = disk - .delete( - &bucket, - &path, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - { - warn!("write meta revert err {:?}", err); - } - }); - } - } - - join_all(revert_futures).await; - return Err(err); - } - Ok(()) - } - - pub(super) async fn update_object_meta( - &self, - bucket: &str, - object: &str, - fi: FileInfo, - disks: &[Option], - ) -> disk::error::Result<()> { - self.update_object_meta_with_opts(bucket, object, fi, disks, &UpdateMetadataOpts::default()) - .await - } - - pub(super) async fn update_object_meta_with_opts( - &self, - bucket: &str, - object: &str, - fi: FileInfo, - disks: &[Option], - opts: &UpdateMetadataOpts, - ) -> disk::error::Result<()> { - if fi.metadata.is_empty() && !opts.replace_user_metadata { - return Ok(()); - } - - self.invalidate_get_object_metadata_cache(bucket, object).await; - - let mut futures = Vec::with_capacity(disks.len()); - - let mut errs = Vec::with_capacity(disks.len()); - - for disk in disks.iter() { - let fi = fi.clone(); - futures.push(async move { - if let Some(disk) = disk { - disk.update_metadata(bucket, object, fi, opts).await - } else { - Err(DiskError::DiskNotFound) - } - }) - } - - let results = join_all(futures).await; - for result in results { - match result { - Ok(_) => { - errs.push(None); - } - Err(e) => { - errs.push(Some(e)); - } - } - } - - if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, fi.write_quorum(self.default_write_quorum())) { - return Err(err); - } - - self.invalidate_get_object_metadata_cache(bucket, object).await; - - Ok(()) - } - - pub(super) async fn delete_if_dangling( - &self, - bucket: &str, - object: &str, - meta_arr: &[FileInfo], - errs: &[Option], - data_errs_by_part: &HashMap>, - opts: ObjectOptions, - ) -> disk::error::Result { - let (m, can_heal) = is_object_dangling(meta_arr, errs, data_errs_by_part); - - if !can_heal { - return Err(DiskError::ErasureReadQuorum); - } - - // Recently written objects get a grace window before dangling cleanup: after - // an unclean shutdown some disks may still be catching up (or carry writes - // that were never made durable), and deleting the surviving shards right away - // turns a partial loss into a total one. Skip deletion and leave the object - // for a later heal/scanner pass to re-evaluate. - if m.is_valid() - && let Some(mod_time) = m.mod_time - { - let grace = dangling_delete_grace(); - if !grace.is_zero() && OffsetDateTime::now_utc() - mod_time < grace { - info!( - bucket = bucket, - object = object, - mod_time = %mod_time, - grace_secs = grace.whole_seconds(), - "skipping dangling-object deletion within grace window" - ); - return Err(DiskError::ErasureReadQuorum); - } - } - - let mut tags: HashMap = HashMap::new(); - tags.insert("set".to_string(), self.set_index.to_string()); - tags.insert("pool".to_string(), self.pool_index.to_string()); - tags.insert("merrs".to_string(), join_errs(errs)); - tags.insert("derrs".to_string(), format!("{data_errs_by_part:?}")); - if m.is_valid() { - tags.insert("sz".to_string(), m.size.to_string()); - tags.insert( - "mt".to_string(), - m.mod_time - .as_ref() - .map_or(String::new(), |mod_time| mod_time.unix_timestamp().to_string()), - ); - tags.insert("d:p".to_string(), format!("{}:{}", m.erasure.data_blocks, m.erasure.parity_blocks)); - } else { - tags.insert("invalid".to_string(), "1".to_string()); - tags.insert( - "d:p".to_string(), - format!("{}:{}", self.set_drive_count - self.default_parity_count, self.default_parity_count), - ); - } - let mut offline = 0; - for (i, err) in errs.iter().enumerate() { - let mut found = false; - if let Some(err) = err - && err == &DiskError::DiskNotFound - { - found = true; - } - for p in data_errs_by_part { - if let Some(v) = p.1.get(i) - && *v == CHECK_PART_DISK_NOT_FOUND - { - found = true; - break; - } - } - - if found { - offline += 1; - } - } - - if offline > 0 { - tags.insert("offline".to_string(), offline.to_string()); - } - - let mut fi = FileInfo::default(); - if let Some(ref version_id) = opts.version_id { - fi.version_id = Uuid::parse_str(version_id).ok(); - } - - fi.set_tier_free_version_id(&Uuid::new_v4().to_string()); - - let disks = self.get_disks_internal().await; - - let mut futures = Vec::with_capacity(disks.len()); - for disk_op in disks.iter() { - let bucket = bucket.to_string(); - let object = object.to_string(); - let fi = fi.clone(); - futures.push(async move { - if let Some(disk) = disk_op { - disk.delete_version(&bucket, &object, fi, false, DeleteOptions::default()) - .await - } else { - Err(DiskError::DiskNotFound) - } - }); - } - - let results = join_all(futures).await; - let mut delete_errs = Vec::with_capacity(results.len()); - for (index, result) in results.into_iter().enumerate() { - let key = format!("ddisk-{index}"); - let already_absent = matches!( - errs.get(index).and_then(Option::as_ref), - Some(DiskError::FileNotFound | DiskError::FileVersionNotFound) - ); - match result { - Ok(_) => { - tags.insert(key, "".to_string()); - delete_errs.push(None); - } - Err(e) => { - tags.insert(key, e.to_string()); - if already_absent || matches!(&e, DiskError::FileNotFound | DiskError::FileVersionNotFound) { - delete_errs.push(None); - } else { - delete_errs.push(Some(e)); - } - } - } - } - - let write_quorum = if m.is_valid() { - m.write_quorum(self.default_write_quorum()) - } else { - self.default_write_quorum() - }; - if let Some(err) = reduce_write_quorum_errs(&delete_errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - return Err(err); - } - - Ok(m) - } - - pub(super) async fn delete_prefix(&self, bucket: &str, prefix: &str) -> disk::error::Result<()> { - let disks = self.get_disks_internal().await; - let write_quorum = disks.len() / 2 + 1; - - let mut futures = Vec::with_capacity(disks.len()); - - for disk_op in disks.iter() { - let bucket = bucket.to_string(); - let prefix = prefix.to_string(); - futures.push(async move { - if let Some(disk) = disk_op { - disk.delete( - &bucket, - &prefix, - DeleteOptions { - recursive: true, - immediate: true, - ..Default::default() - }, - ) - .await - } else { - Ok(()) - } - }); - } - - let errs = join_all(futures).await.into_iter().map(|v| v.err()).collect::>(); - - if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - return Err(err); - } - - Ok(()) - } - - /// Scan a single disk's copy of `prefix` and decide whether it is an orphan - /// (metadata-less) directory subtree. Walks the tree iteratively and returns - /// [`OrphanDirScan::HasData`] as soon as any regular file is found. - async fn scan_orphan_dir(disk: &DiskStore, bucket: &str, prefix: &str) -> OrphanDirScan { - let root = prefix.trim_end_matches(SLASH_SEPARATOR).to_string(); - let mut stack = vec![root.clone()]; - // Pre-order list of directories (a parent always precedes its descendants), - // so reversing it yields a safe children-first removal order. - let mut dirs: Vec = Vec::new(); - let mut existed = false; - - while let Some(dir) = stack.pop() { - let entries = match disk.list_dir("", bucket, &dir, 0).await { - Ok(entries) => entries, - Err(_) => { - // The root missing (or never existing) means there is nothing to - // purge on this disk. A nested directory vanishing mid-scan is a - // benign race, so skip it and keep walking. - if dir == root { - return OrphanDirScan::Missing; - } - continue; - } - }; - - existed = true; - dirs.push(dir.clone()); - - for entry in entries { - match entry.strip_suffix(SLASH_SEPARATOR) { - // `read_dir` marks directories with a trailing slash; anything else - // is a regular file, which means real object data lives here. - Some(child) => stack.push(format!("{dir}{SLASH_SEPARATOR}{child}")), - None => return OrphanDirScan::HasData, - } - } - } - - if existed { - OrphanDirScan::Empty(dirs) - } else { - OrphanDirScan::Missing - } - } - - /// Purge an orphan directory prefix — a trailing-slash key that exists on disk - /// as an empty directory tree with no object metadata on any disk of this set. - /// Such prefixes are listable (see `scan_dir`) yet are not real objects, so the - /// normal delete path returns NotFound and leaves them stranded (issue #4189). - /// - /// Callers pass the *decoded* directory name (`prefix/`), not the `__XLDIR__` - /// encoded object key — the orphan tree on disk uses the plain path. - /// - /// Returns `Ok(true)` when the prefix was an orphan tree on this set and was - /// removed, `Ok(false)` when it holds real data or does not exist on any disk - /// of this set (the caller should surface the original NotFound), and `Err` on - /// a hard disk failure. - pub(crate) async fn purge_orphan_dir_object(&self, bucket: &str, object: &str) -> disk::error::Result { - let disks = self.get_disks_internal().await; - - // Phase 1: classify every online disk. Refuse to purge if ANY disk holds - // object data under the prefix, so a degraded/healable object is never - // destroyed. - let mut per_disk_dirs: Vec<(usize, Vec)> = Vec::new(); - let mut existed = false; - for (i, disk) in disks.iter().enumerate() { - let Some(disk) = disk else { continue }; - match Self::scan_orphan_dir(disk, bucket, object).await { - OrphanDirScan::HasData => return Ok(false), - OrphanDirScan::Empty(dirs) => { - existed = true; - per_disk_dirs.push((i, dirs)); - } - OrphanDirScan::Missing => {} - } - } - - if !existed { - return Ok(false); - } - - // Phase 2: remove the empty directories children-first on each disk. A - // non-recursive delete performs an empty-only `rmdir`, so a directory that - // concurrently gained an object fails with DirectoryNotEmpty and is skipped — - // a racing PutObject is never clobbered. - for (i, mut dirs) in per_disk_dirs { - let Some(disk) = disks[i].as_ref() else { continue }; - dirs.reverse(); - for dir in dirs { - if let Err(err) = disk - .delete( - bucket, - &dir, - DeleteOptions { - recursive: false, - immediate: true, - ..Default::default() - }, - ) - .await - { - // Best effort: a sibling removal may have already cleared a shared - // parent, or a concurrent writer repopulated the directory. Neither - // is fatal to purging the orphan tree. - debug!(bucket, object, dir, error = ?err, "purge_orphan_dir_object: skipped non-empty/absent directory"); - } - } - } - - Ok(true) - } - - pub(super) async fn check_write_precondition( - &self, - bucket: &str, - object: &str, - opts: &ObjectOptions, - ) -> Option { - let mut opts = opts.clone(); - - let http_preconditions = opts.http_preconditions?; - opts.http_preconditions = None; - - // Never claim a lock here, to avoid deadlock - // - If no_lock is false, we must have obtained the lock out side of this function - // - If no_lock is true, we should not obtain locks - opts.no_lock = true; - let oi = self.get_object_info(bucket, object, &opts).await; - - match oi { - Ok(oi) => { - // If top level is a delete marker proceed to upload. - if oi.delete_marker { - return None; - } - let if_none_match = http_preconditions.if_none_match_value().map(str::to_owned); - let if_match = http_preconditions.if_match_value().map(str::to_owned); - if should_prevent_write(&oi, if_none_match, if_match) { - return Some(StorageError::PreconditionFailed); - } - } - - Err(StorageError::VersionNotFound(_, _, _)) - | Err(StorageError::ObjectNotFound(_, _)) - | Err(StorageError::ErasureReadQuorum) => { - // When the object is not found, - // - if If-Match is set, we should return 404 NotFound - // - if If-None-Match is set, we should be able to proceed with the request - if http_preconditions.if_match_value().is_some() { - return Some(StorageError::ObjectNotFound(bucket.to_string(), object.to_string())); - } - } - - Err(e) => { - return Some(e); - } - } - - None - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn dangling_delete_grace_defaults_to_one_hour() { - temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, None::<&str>, || { - assert_eq!(dangling_delete_grace(), time::Duration::seconds(3600)); - }); - } - - #[test] - fn dangling_delete_grace_env_override_and_disable() { - temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, Some("120"), || { - assert_eq!(dangling_delete_grace(), time::Duration::seconds(120)); - }); - temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, Some("0"), || { - assert!(dangling_delete_grace().is_zero()); - }); - } -}