// 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::super::*; use crate::io_support::bitrot::object_mmap_read_enabled; use crate::storage_api_contracts::namespace::NamespaceLocking as _; const LOG_COMPONENT_ECSTORE: &str = "ecstore"; const LOG_SUBSYSTEM_HEAL: &str = "heal"; const EVENT_HEAL_OBJECT_RENAME: &str = "heal_object_rename"; const HEAL_RENAME_INCOMPLETE: &str = "heal rename incomplete"; #[cfg(test)] static HEAL_RENAME_FAILURES: std::sync::Mutex> = std::sync::Mutex::new(Vec::new()); #[cfg(test)] struct HealRenameFailureScope { bucket: String, object: String, } #[cfg(test)] impl HealRenameFailureScope { fn install(bucket: &str, object: &str, disk_indexes: &[usize]) -> Self { let mut failures = HEAL_RENAME_FAILURES .lock() .expect("heal rename failure registry should not poison"); assert!( !failures .iter() .any(|(registered_bucket, registered_object, _)| { registered_bucket == bucket && registered_object == object }), "heal rename failures must be installed once per object" ); failures.extend( disk_indexes .iter() .map(|index| (bucket.to_string(), object.to_string(), *index)), ); Self { bucket: bucket.to_string(), object: object.to_string(), } } } #[cfg(test)] impl Drop for HealRenameFailureScope { fn drop(&mut self) { HEAL_RENAME_FAILURES .lock() .expect("heal rename failure registry should not poison") .retain(|(bucket, object, _)| bucket != &self.bucket || object != &self.object); } } #[cfg(test)] fn should_fail_heal_rename(bucket: &str, object: &str, disk_index: usize) -> bool { let mut failures = HEAL_RENAME_FAILURES .lock() .expect("heal rename failure registry should not poison"); if let Some(position) = failures .iter() .position(|entry| entry == &(bucket.to_string(), object.to_string(), disk_index)) { failures.swap_remove(position); true } else { false } } #[cfg(not(test))] fn should_fail_heal_rename(_bucket: &str, _object: &str, _disk_index: usize) -> bool { false } #[derive(Clone, Copy, Debug, PartialEq, Eq)] struct PartFailureSummary { part_number: usize, failed_shards: usize, bitrot_failure: bool, } #[derive(Clone)] struct RecoverableMetaCandidate { identity: [u8; 32], file_info: FileInfo, data_count: usize, local_payload: bool, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum DanglingDeleteSafety { UnsafeToDelete, NoRecoverableCandidate, } #[cfg(test)] struct DanglingCheckPartsFailure { key: DanglingCheckPartsFailureKey, } #[cfg(test)] type DanglingCheckPartsFailureKey = (String, String, usize); #[cfg(test)] type DanglingCheckPartsFailures = HashMap; #[cfg(test)] fn dangling_check_parts_failures() -> &'static std::sync::Mutex { static FAILURES: std::sync::OnceLock> = std::sync::OnceLock::new(); FAILURES.get_or_init(|| std::sync::Mutex::new(HashMap::new())) } #[cfg(test)] impl DanglingCheckPartsFailure { fn install(bucket: &str, object: &str, disk_index: usize, error: DiskError) -> Self { let key = (bucket.to_string(), object.to_string(), disk_index); let previous = dangling_check_parts_failures() .lock() .expect("dangling check-parts failure registry should not poison") .insert(key.clone(), error); assert!(previous.is_none(), "dangling check-parts failure already installed"); Self { key } } } #[cfg(test)] impl Drop for DanglingCheckPartsFailure { fn drop(&mut self) { dangling_check_parts_failures() .lock() .expect("dangling check-parts failure registry should not poison") .remove(&self.key); } } #[cfg(test)] fn injected_dangling_check_parts_error(bucket: &str, object: &str, disk_index: usize) -> Option { dangling_check_parts_failures() .lock() .expect("dangling check-parts failure registry should not poison") .get(&(bucket.to_string(), object.to_string(), disk_index)) .cloned() } fn first_unhealthy_part_summary( data_errs_by_part: &HashMap>, parts: &[ObjectPartInfo], ) -> Option { data_errs_by_part .iter() .filter_map(|(part_index, part_errs)| { let failed_shards = count_part_not_success(part_errs); if failed_shards == 0 { return None; } Some(( *part_index, PartFailureSummary { part_number: parts.get(*part_index).map(|part| part.number).unwrap_or(part_index + 1), failed_shards, bitrot_failure: part_errs.contains(&CHECK_PART_FILE_CORRUPT), }, )) }) .min_by_key(|(part_index, _)| *part_index) .map(|(_, summary)| summary) } impl SetDisks { #[tracing::instrument(skip(self, opts), fields(bucket = %bucket, object = %object, version_id = %version_id))] pub(in crate::set_disk) async fn heal_object( &self, bucket: &str, object: &str, version_id: &str, opts: &HealOpts, ) -> disk::error::Result<(HealResultItem, Option)> { Box::pin(self.heal_object_with_explicit_version_regen(bucket, object, version_id, opts, true)).await } #[allow(clippy::too_many_lines)] async fn heal_object_with_explicit_version_regen( &self, bucket: &str, object: &str, version_id: &str, opts: &HealOpts, allow_explicit_version_regen: bool, ) -> disk::error::Result<(HealResultItem, Option)> { info!(?opts, "Starting heal_object"); let disks = self.get_disks_internal().await; let mut result = HealResultItem { heal_item_type: HealItemType::Object.to_string(), bucket: bucket.to_string(), object: object.to_string(), version_id: version_id.to_string(), disk_count: disks.len(), ..Default::default() }; let write_lock_guard = if !opts.no_lock { let ns_lock = self.new_ns_lock(bucket, object).await?; Some( ns_lock .get_write_lock(get_lock_acquire_timeout()) .await .map_err(|e| self.map_namespace_lock_error(bucket, object, "write", e))?, ) } else { None }; let version_id_op = { if version_id.is_empty() { None } else { Some(version_id.to_string()) } }; let (mut parts_metadata, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, true, true, false).await?; info!( parts_count = parts_metadata.len(), bucket = bucket, object = object, version_id = version_id, ?errs, "File info read complete" ); if DiskError::is_all_not_found(&errs) { debug!(bucket, object, version_id, "heal_object skipped missing object"); let err = if !version_id.is_empty() { DiskError::FileVersionNotFound } else { DiskError::FileNotFound }; // Nothing to do, file is already gone. return Ok(( self.default_heal_result(FileInfo::default(), &errs, bucket, object, version_id) .await, Some(err), )); } info!(parts_count = parts_metadata.len(), "heal_object Initiating quorum check"); match Self::object_quorum_from_meta(&parts_metadata, &errs, self.default_parity_count) { Ok((read_quorum, _)) => { result.parity_blocks = result.disk_count - read_quorum as usize; result.data_blocks = read_quorum as usize; let ((mut online_disks, quorum_mod_time, quorum_etag), disk_len) = { let disks = self.disks.read().await; let disk_len = disks.len(); (Self::list_online_disks(&disks, &parts_metadata, &errs, read_quorum as usize), disk_len) }; info!(?parts_metadata, ?errs, ?read_quorum, ?disk_len, "heal_object List disks metadata"); info!(?online_disks, ?quorum_mod_time, ?quorum_etag, "heal_object List online disks"); let filter_by_etag = quorum_etag.is_some(); match Self::pick_valid_fileinfo(&parts_metadata, quorum_mod_time, quorum_etag.clone(), read_quorum as usize) { Ok(latest_meta) => { info!("heal_object latest_meta: {:?}", latest_meta); let (data_errs_by_disk, data_errs_by_part) = disks_with_all_parts( &mut online_disks, &mut parts_metadata, &errs, &latest_meta, filter_by_etag, bucket, object, opts.scan_mode, ) .await?; info!( "disks_with_all_parts heal_object results: available_disks count={}, total_disks={}", online_disks.iter().filter(|d| d.is_some()).count(), online_disks.len() ); let erasure = if !latest_meta.deleted && !latest_meta.is_remote() { // Initialize erasure coding; use legacy mode for old-version files coding::Erasure::try_new_with_options( latest_meta.erasure.data_blocks, latest_meta.erasure.parity_blocks, latest_meta.erasure.block_size, latest_meta.uses_legacy_checksum, ) .map_err(DiskError::from)? } else { coding::Erasure::default() }; result.object_size = ObjectInfo::from_file_info(&latest_meta, bucket, object, true).get_actual_size()? as usize; // Loop to find number of disks with valid data, per-drive // data state and a list of outdated disks on which data needs // to be healed. let mut out_dated_disks = vec![None; disk_len]; let mut disks_to_heal_count = 0; let mut meta_to_heal_count = 0; for index in 0..online_disks.len() { let (yes, is_meta, reason) = should_heal_object_on_disk( &errs[index], &data_errs_by_disk[&index], &parts_metadata[index], &latest_meta, ); if yes { out_dated_disks[index] = disks[index].clone(); disks_to_heal_count += 1; if is_meta { meta_to_heal_count += 1; } debug!("heal_object Disk {} marked for healing (endpoint={})", index, self.set_endpoints[index]); } let drive_state = match reason { Some(err) => match err { DiskError::DiskNotFound => DriveState::Offline.to_string(), DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound | DiskError::PartMissingOrCorrupt | DiskError::OutdatedXLMeta => DriveState::Missing.to_string(), DiskError::FileCorrupt => DriveState::Corrupt.to_string(), _ => DriveState::Unknown(err.to_string()).to_string(), }, None => DriveState::Ok.to_string(), }; result.before.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint: self.set_endpoints[index].to_string(), state: drive_state.to_string(), }); result.after.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint: self.set_endpoints[index].to_string(), state: drive_state.to_string(), }); } if disks_to_heal_count == 0 { // The object is already healthy: no disk needs healing. // This is the common case for the very objects PR #4356 // targets — a valid `xl.meta` plus a leaked pre-#3510 // data dir needs no shard healing, so it would otherwise // return here and never reach the post-heal reclaim tail // below. Sweep the strays on this path too (issues #3231, // #3191). Skipped on dry-run, like every mutating step. if !opts.dry_run { self.reclaim_orphan_data_dirs_best_effort(bucket, object).await; } return Ok((result, None)); } if opts.dry_run { return Ok((result, None)); } let mut cannot_heal = !latest_meta.deleted && meta_to_heal_count > latest_meta.erasure.parity_blocks; if cannot_heal && quorum_etag.is_some() { cannot_heal = false; } if !latest_meta.deleted && !latest_meta.is_remote() { for part_errs in data_errs_by_part.values() { if count_part_not_success(part_errs) > latest_meta.erasure.parity_blocks { cannot_heal = true; break; } } } if cannot_heal { let total_disks = parts_metadata.len(); let healthy_count = total_disks.saturating_sub(disks_to_heal_count); let required_data = total_disks.saturating_sub(latest_meta.erasure.parity_blocks); let no_parity_failure = (!latest_meta.deleted && !latest_meta.is_remote() && latest_meta.erasure.parity_blocks == 0) .then(|| first_unhealthy_part_summary(&data_errs_by_part, &latest_meta.parts)) .flatten(); let cannot_heal_err = if no_parity_failure.is_some_and(|failure| failure.bitrot_failure) { DiskError::FileCorrupt } else { DiskError::ErasureReadQuorum }; if let Some(failure) = no_parity_failure { result.detail = format!( "no-parity object is unrecoverable: part {} has {} missing or corrupt data shard(s), bitrot_failure={}, data_blocks={}, parity_blocks=0", failure.part_number, failure.failed_shards, failure.bitrot_failure, latest_meta.erasure.data_blocks ); error!( bucket, object, version_id, data_shards = latest_meta.erasure.data_blocks, parity_shards = latest_meta.erasure.parity_blocks, required_data_shards = required_data, healthy_shards = healthy_count, missing_or_corrupt_shards = disks_to_heal_count, part_number = failure.part_number, bitrot_failure = failure.bitrot_failure, "No-parity object failed integrity or availability validation and cannot be reconstructed" ); } else { result.detail = format!( "object cannot be reconstructed with available shards: required_data_shards={required_data}, healthy_shards={healthy_count}, missing_or_corrupt_shards={disks_to_heal_count}, parity_shards={}", latest_meta.erasure.parity_blocks ); error!( bucket, object, version_id, required_data_shards = required_data, healthy_shards = healthy_count, missing_or_corrupt_shards = disks_to_heal_count, parity_shards = latest_meta.erasure.parity_blocks, "Heal object cannot reconstruct with available shards" ); } // `disks_with_all_parts` normalizes conflicting entries // 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?; if self .dangling_delete_safety(bucket, object, &delete_guard_metadata, &delete_guard_errs, &disks) .await? == DanglingDeleteSafety::UnsafeToDelete { return Ok((result, Some(cannot_heal_err))); } // Allow for dangling deletes, on versions that have DataDir missing etc. // this would end up restoring the correct readable versions. return match self .delete_if_dangling( bucket, object, &parts_metadata, &errs, &data_errs_by_part, ObjectOptions { version_id: version_id_op.clone(), ..Default::default() }, ) .await { Ok(m) => { let derr = if !version_id.is_empty() { DiskError::FileVersionNotFound } else { DiskError::FileNotFound }; let mut t_errs = Vec::with_capacity(errs.len()); for _ in 0..errs.len() { t_errs.push(None); } Ok((self.default_heal_result(m, &t_errs, bucket, object, version_id).await, Some(derr))) } Err(err) => { error!( bucket, object, version_id, error = %err, returned_error = %cannot_heal_err, "Heal object dangling cleanup could not prove object deletion" ); Ok((result, Some(cannot_heal_err))) } }; } if !latest_meta.deleted && latest_meta.erasure.distribution.len() != online_disks.len() { let err_str = format!( "unexpected file distribution ({:?}) from available disks ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})", latest_meta.erasure.distribution, online_disks, bucket, object, version_id ); warn!(err_str); let err = DiskError::other(err_str); return Ok(( self.default_heal_result(latest_meta, &errs, bucket, object, version_id).await, Some(err), )); } let latest_disks = Self::shuffle_disks(&online_disks, &latest_meta.erasure.distribution); if !latest_meta.deleted && latest_meta.erasure.distribution.len() != out_dated_disks.len() { let err_str = format!( "unexpected file distribution ({:?}) from outdated disks ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})", latest_meta.erasure.distribution, out_dated_disks, bucket, object, version_id ); warn!(err_str); let err = DiskError::other(err_str); return Ok(( self.default_heal_result(latest_meta, &errs, bucket, object, version_id).await, Some(err), )); } if !latest_meta.deleted && latest_meta.erasure.distribution.len() != parts_metadata.len() { let err_str = format!( "unexpected file distribution ({:?}) from metadata entries ({:?}), looks like backend disks have been manually modified refusing to heal {}/{}({})", latest_meta.erasure.distribution, parts_metadata.len(), bucket, object, version_id ); warn!(err_str); let err = DiskError::other(err_str); return Ok(( self.default_heal_result(latest_meta, &errs, bucket, object, version_id).await, Some(err), )); } out_dated_disks = Self::shuffle_disks(&out_dated_disks, &latest_meta.erasure.distribution); let mut parts_metadata = Self::shuffle_parts_metadata(&parts_metadata, &latest_meta.erasure.distribution); let mut copy_parts_metadata = vec![None; parts_metadata.len()]; for (index, disk) in latest_disks.iter().enumerate() { if disk.is_some() { copy_parts_metadata[index] = Some(parts_metadata[index].clone()); } } let clean_file_info = |fi: &FileInfo| -> FileInfo { let mut nfi = fi.clone(); if !nfi.is_remote() { nfi.data = None; nfi.erasure.index = 0; nfi.erasure.checksums = Vec::new(); } nfi }; for (index, disk) in out_dated_disks.iter().enumerate() { if disk.is_some() { // Make sure to write the FileInfo information // that is expected to be in quorum. parts_metadata[index] = clean_file_info(&latest_meta); } } // We write at temporary location and then rename to final location. let tmp_id = Uuid::new_v4().to_string(); // Delete markers and remote (transitioned) objects carry no data_dir and // skip the data-heal block below, so a nil placeholder is safe for them. // For a regular object a missing data_dir means the latest metadata is // corrupt; fail this object's heal with a clear error instead of building // part paths under a nil UUID directory. let data_dir = match latest_meta.data_dir { Some(data_dir) => data_dir, None => { if !latest_meta.deleted && !latest_meta.is_remote() { error!( "heal: latest metadata for {}/{} has no data_dir, cannot heal object data", bucket, object ); return Err(DiskError::FileCorrupt); } Uuid::nil() } }; let src_data_dir = data_dir.to_string(); let dst_data_dir = data_dir; if !latest_meta.deleted && !latest_meta.is_remote() { let erasure_info = latest_meta.erasure.clone(); for (part_index, part) in latest_meta.parts.iter().enumerate() { let till_offset = erasure.shard_file_offset(0, part.size, part.size); let use_mmap_read = object_mmap_read_enabled(); let mut readers = Vec::with_capacity(latest_disks.len()); let mut writers = Vec::with_capacity(out_dated_disks.len()); // let mut errors = Vec::with_capacity(out_dated_disks.len()); let mut prefer = vec![false; latest_disks.len()]; for (index, disk) in latest_disks.iter().enumerate() { let this_part_errs = Self::shuffle_check_parts(&data_errs_by_part[&part_index], &erasure_info.distribution); if this_part_errs[index] != CHECK_PART_SUCCESS { info!( "reading part {}: index={}, part_errs={:?}, skipping", part.number, index, this_part_errs[index] ); readers.push(None); continue; } if let (Some(disk), Some(metadata)) = (disk, ©_parts_metadata[index]) { let checksum_info = metadata.erasure.get_checksum_info(part.number); let checksum_algo = if metadata.uses_legacy_checksum && checksum_info.algorithm == HashAlgorithm::HighwayHash256S { HashAlgorithm::HighwayHash256SLegacy } else { checksum_info.algorithm }; match create_bitrot_reader( metadata.data.as_deref(), Some(disk), bucket, &path_join_buf(&[object, &src_data_dir, &format!("part.{}", part.number)]), 0, till_offset, erasure.shard_size(), checksum_algo.clone(), false, use_mmap_read, ) .await { Ok(Some(reader)) => { readers.push(Some(reader)); } Ok(None) => { readers.push(None); continue; } Err(e) => { readers.push(None); continue; } } prefer[index] = disk.host_name().is_empty(); } else { readers.push(None); // errors.push(Some(DiskError::DiskNotFound)); } } // Preserve the committed layout: recomputing inline-ness here // (with a hardcoded unversioned threshold) makes healed replicas // diverge from healthy ones in quorum identity, so heal would // flag them forever. let is_inline_buffer = latest_meta.inline_data(); // create writers for all disk positions, but only for outdated disks for (index, disk_op) in out_dated_disks.iter().enumerate() { if let Some(outdated_disk) = disk_op { let writer = match create_bitrot_writer( is_inline_buffer, Some(outdated_disk), RUSTFS_META_TMP_BUCKET, &path_join_buf(&[ &tmp_id.to_string(), &dst_data_dir.to_string(), &format!("part.{}", part.number), ]), erasure.shard_file_size(part.size as i64), erasure.shard_size(), HashAlgorithm::HighwayHash256S, ) .await { Ok(writer) => writer, Err(err) => { info!( "create_bitrot_writer disk {}, err {:?}, skipping operation", outdated_disk.to_string(), err ); writers.push(None); continue; } }; writers.push(Some(writer)); } else { writers.push(None); } } // Heal each part. erasure.Heal() will write the healed // part to .rustfs/tmp/uuid/ which needs to be renamed // later to the final location. if let Err(e) = erasure.heal(&mut writers, readers, part.size, &prefer).await { // Don't leak the partially-written healed shards in // .rustfs/tmp when heal fails midway (backlog#799 B20). let _ = self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_id).await; return Err(e); } // close_bitrot_writers(&mut writers).await?; for (index, disk_op) in out_dated_disks.iter_mut().enumerate() { if disk_op.is_none() { continue; } if writers[index].is_none() { *disk_op = None; disks_to_heal_count -= 1; continue; } parts_metadata[index].data_dir = Some(dst_data_dir); parts_metadata[index].add_object_part( part.number, part.etag.clone(), part.size, part.mod_time, part.actual_size, part.index.clone(), part.checksums.clone(), ); if is_inline_buffer { if let Some(writer) = writers[index].take() { // if let Some(w) = writer.as_any().downcast_ref::() { // parts_metadata[index].data = Some(w.inline_data().to_vec()); // } parts_metadata[index].data = Some(writer.into_inline_data().map(Bytes::from).unwrap_or_default()); } parts_metadata[index].set_inline_data(); } else { parts_metadata[index].data = None; } } if disks_to_heal_count == 0 { // Clean up healed shards written to .rustfs/tmp before bailing (B20). let _ = self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_id).await; return Ok(( result, Some(DiskError::other(format!( "all drives had write errors, unable to heal {bucket}/{object}" ))), )); } } } // Rename from tmp location to the actual location. // MinIO stops on the first RenameData error. RustFS intentionally // continues per target, but reports any residue after all attempts // so successful repairs survive and failed targets remain retryable. let mut rename_attempts = 0usize; let mut rename_successes = 0usize; let mut healed_disks = vec![None; out_dated_disks.len()]; for (index, outdated_disk) in out_dated_disks.iter().enumerate() { if let Some(disk) = outdated_disk { rename_attempts += 1; // record the index of the updated disks parts_metadata[index].erasure.index = index + 1; // Attempt a rename now from healed data to final location. parts_metadata[index].set_healing(); let rename_result = if should_fail_heal_rename(bucket, object, index) { Err(DiskError::Unexpected) } else { disk.rename_data( RUSTFS_META_TMP_BUCKET, &tmp_id, parts_metadata[index].clone(), bucket, object, ) .await }; if let Err(err) = &rename_result { warn!( event = EVENT_HEAL_OBJECT_RENAME, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_HEAL, bucket, object, version_id, disk_index = index, endpoint = %disk.endpoint(), tmp_id, result = "failed", error = %err, "Heal object rename failed" ); } else { rename_successes += 1; healed_disks[index] = Some(disk.clone()); if parts_metadata[index].is_remote() { let rm_data_dir = parts_metadata[index].data_dir.expect("operation should succeed").to_string(); let d_path = Path::new(&encode_dir_object(object)).join(rm_data_dir); if let Err(e) = disk .delete( bucket, d_path.to_str().expect("operation should succeed"), DeleteOptions { immediate: true, recursive: true, ..Default::default() }, ) .await { // The healed shard has already been renamed into place; a // failure cleaning up the old remote data dir must not abort // the heal and leak the tmp shards (backlog#799 B20). warn!( component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_HEAL, bucket, object, error = %e, "Heal remote data-dir cleanup failed" ); } } for (i, v) in result.before.drives.iter().enumerate() { if v.endpoint == disk.endpoint().to_string() { result.after.drives[i].state = DriveState::Ok.to_string(); } } } } } self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_id) .await .map_err(DiskError::other)?; self.record_healed_capacity_scope(&healed_disks); if rename_successes < rename_attempts { return Ok(( result, Some(DiskError::other(format!( "{HEAL_RENAME_INCOMPLETE}: {rename_successes} of {rename_attempts} targets committed for \ {bucket}/{object}" ))), )); } // The object is healthy here; sweep any data dirs left behind // by pre-#3510 unversioned overwrites, which the dangling paths // above never touch (issues #3231, #3191). Best effort — a // failure must not fail the heal. self.reclaim_orphan_data_dirs_best_effort(bucket, object).await; Ok((result, None)) } Err(err) => Ok((result, Some(err))), } } Err(err) => { if allow_explicit_version_regen && !version_id.is_empty() && self .try_regenerate_explicit_version_meta(bucket, object, version_id, &parts_metadata, &errs, &disks) .await? { return Box::pin(self.heal_object_with_explicit_version_regen(bucket, object, version_id, opts, false)).await; } if self .dangling_delete_safety(bucket, object, &parts_metadata, &errs, &disks) .await? == DanglingDeleteSafety::UnsafeToDelete { return Ok(( self.default_heal_result(FileInfo::default(), &errs, bucket, object, version_id) .await, Some(err), )); } let data_errs_by_part = HashMap::new(); match self .delete_if_dangling( bucket, object, &parts_metadata, &errs, &data_errs_by_part, ObjectOptions { version_id: version_id_op.clone(), ..Default::default() }, ) .await { Ok(m) => { let err = if !version_id.is_empty() { DiskError::FileVersionNotFound } else { DiskError::FileNotFound }; Ok((self.default_heal_result(m, &errs, bucket, object, version_id).await, Some(err))) } Err(_) => Ok(( self.default_heal_result(FileInfo::default(), &errs, bucket, object, version_id) .await, Some(err), )), } } } } async fn try_regenerate_explicit_version_meta( &self, bucket: &str, object: &str, version_id: &str, parts_metadata: &[FileInfo], errs: &[Option], disks: &[Option], ) -> disk::error::Result { let Ok(version_id) = Uuid::parse_str(version_id) else { return Ok(false); }; let candidates = parts_metadata .iter() .zip(errs.iter()) .filter_map(|(file_info, err)| { (err.is_none() && file_info_is_valid_for_metadata(file_info) && file_info.version_id == Some(version_id) && file_info.has_valid_erasure_geometry() && !file_info.deleted && !file_info.is_remote() && file_info.data_dir.is_some() && !file_info.parts.is_empty() && file_info.erasure.data_blocks > 0 && file_info .erasure .data_blocks .checked_add(file_info.erasure.parity_blocks) .is_some_and(|shards| shards == disks.len())) .then_some(file_info) }) .collect::>(); let Some(candidate) = candidates.first().copied() else { return Ok(false); }; let identity = Self::file_info_quorum_hash(candidate); if candidates .iter() .any(|file_info| Self::file_info_quorum_hash(file_info) != identity) { return Ok(false); } let mut available = 0usize; for disk in disks { let Some(disk) = disk else { return Ok(false); }; match disk.check_parts(bucket, object, candidate).await { Ok(response) if !response.results.is_empty() && response.results.iter().all(|result| *result == CHECK_PART_SUCCESS) => { available += 1; } Ok(_) | Err( DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::PathNotFound | DiskError::VolumeNotFound, ) => {} Err(_) => return Ok(false), } } if available < candidate.erasure.data_blocks { return Ok(false); } let mut wrote = 0usize; for (index, disk) in disks.iter().enumerate() { let Some(disk) = disk else { return Ok(false); }; let metadata_absent = matches!( errs.get(index).and_then(Option::as_ref), Some(DiskError::FileNotFound | DiskError::FileVersionNotFound) ); if !metadata_absent { continue; } let Some(&shard_index) = candidate.erasure.distribution.get(index) else { return Ok(false); }; let mut regenerated = candidate.clone(); regenerated.fresh = false; regenerated.erasure.index = shard_index; match disk.write_metadata("", bucket, object, regenerated).await { Ok(()) => wrote += 1, Err(error) => { warn!( bucket, object, disk_index = index, error = %error, "failed to regenerate recoverable xl.meta" ); } } } Ok(wrote > 0) } /// Best-effort orphan-data-dir reclaim for an object that is healthy on this /// set. Wraps [`Self::reclaim_orphan_data_dirs`] with the shared logging so /// both `heal_object` exits — the already-healthy early return and the /// post-heal tail — reclaim identically. Never fails the heal: delete errors /// are logged and swallowed. Callers must gate this on `!opts.dry_run`. async fn reclaim_orphan_data_dirs_best_effort(&self, bucket: &str, object: &str) { match self.reclaim_orphan_data_dirs(bucket, object).await { Ok(removed) if removed > 0 => { info!(bucket, object, removed, "heal_object: reclaimed orphaned data directories"); } Ok(_) => {} Err(e) => { warn!(bucket, object, error = %e, "heal_object: orphan data-dir reclaim failed"); } } } /// Prevent dangling cleanup when surviving state cannot prove that deletion /// is safe. Part presence proves only recoverability, never commit: the write /// path can durably rename data before xl.meta is committed. async fn dangling_delete_safety( &self, bucket: &str, object: &str, parts_metadata: &[FileInfo], errs: &[Option], disks: &[Option], ) -> disk::error::Result { if disks.iter().any(Option::is_none) || errs.iter().flatten().any(|err| { !matches!( err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::PathNotFound | DiskError::VolumeNotFound ) }) { return Ok(DanglingDeleteSafety::UnsafeToDelete); } let mut candidates = Vec::::with_capacity(parts_metadata.len()); for (fi, err) in parts_metadata.iter().zip(errs.iter()) { if err.is_some() || !file_info_is_valid_for_metadata(fi) { continue; } let identity = Self::file_info_quorum_hash(fi); if !candidates.iter().any(|candidate| candidate.identity == identity) { let local_payload = fi.has_valid_erasure_geometry() && !fi.deleted && !fi.is_remote() && fi.data_dir.is_some() && !fi.parts.is_empty() && fi.erasure.data_blocks > 0 && fi .erasure .data_blocks .checked_add(fi.erasure.parity_blocks) .is_some_and(|shards| shards == disks.len()); candidates.push(RecoverableMetaCandidate { identity, file_info: fi.clone(), data_count: 0, local_payload, }); } } if candidates .iter() .any(|candidate| candidate.file_info.deleted || candidate.file_info.is_remote()) || candidates.len() > 1 { return Ok(DanglingDeleteSafety::UnsafeToDelete); } for candidate in candidates.iter_mut().filter(|candidate| candidate.local_payload) { for (disk_index, disk) in disks.iter().enumerate() { let Some(disk) = disk else { return Ok(DanglingDeleteSafety::UnsafeToDelete); }; #[cfg(test)] let check_result = match injected_dangling_check_parts_error(bucket, object, disk_index) { Some(error) => Err(error), None => disk.check_parts(bucket, object, &candidate.file_info).await, }; #[cfg(not(test))] let check_result = disk.check_parts(bucket, object, &candidate.file_info).await; match check_result { Ok(resp) if !resp.results.is_empty() && resp.results.iter().all(|result| *result == CHECK_PART_SUCCESS) => { candidate.data_count += 1; } Ok(_) => {} Err( DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::PathNotFound | DiskError::VolumeNotFound, ) => {} Err(_) => return Ok(DanglingDeleteSafety::UnsafeToDelete), } } } Ok( if candidates .iter() .any(|candidate| candidate.local_payload && candidate.data_count >= candidate.file_info.erasure.data_blocks) { DanglingDeleteSafety::UnsafeToDelete } else { DanglingDeleteSafety::NoRecoverableCandidate }, ) } pub(in crate::set_disk) async fn heal_object_dir_locked( &self, bucket: &str, object: &str, dry_run: bool, remove: bool, ) -> Result<(HealResultItem, Option)> { let disks = { let disks = self.disks.read().await; disks.clone() }; let mut result = HealResultItem { heal_item_type: HealItemType::Object.to_string(), bucket: bucket.to_string(), object: object.to_string(), disk_count: self.disks.read().await.len(), parity_blocks: self.default_parity_count, data_blocks: disks.len() - self.default_parity_count, object_size: 0, ..Default::default() }; // Filled below by pushing one entry per disk while zipping the (index-aligned) `errs`. // Pre-filling here would double the reported drive list once the push loop runs. result.before.drives = Vec::with_capacity(disks.len()); result.after.drives = Vec::with_capacity(disks.len()); let errs = stat_all_dirs(&disks, bucket, object).await; let dangling_object = is_object_dir_dangling(&errs); if dangling_object && !dry_run && remove { let mut futures = Vec::with_capacity(disks.len()); for disk in disks.iter().flatten() { let disk = disk.clone(); let bucket = bucket.to_string(); let object = object.to_string(); futures.push(tokio::spawn(async move { let _ = disk .delete( &bucket, &object, DeleteOptions { recursive: false, immediate: false, ..Default::default() }, ) .await; })); } // ignore errors let _ = join_all(futures).await; } for (err, drive) in errs.iter().zip(self.set_endpoints.iter()) { let endpoint = drive.to_string(); let drive_state = match err { Some(err) => match err { DiskError::DiskNotFound => DriveState::Offline.to_string(), DiskError::FileNotFound | DiskError::VolumeNotFound => DriveState::Missing.to_string(), _ => DriveState::Corrupt.to_string(), }, None => DriveState::Ok.to_string(), }; result.before.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint: endpoint.clone(), state: drive_state.to_string(), }); result.after.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint, state: drive_state.to_string(), }); } if dangling_object || DiskError::is_all_not_found(&errs) { return Ok((result, Some(DiskError::FileNotFound))); } if dry_run { // Quit without try to heal the object dir return Ok((result, None)); } for (index, (err, disk)) in errs.iter().zip(disks.iter()).enumerate() { if let (Some(DiskError::VolumeNotFound | DiskError::FileNotFound), Some(disk)) = (err, disk) { let vol_path = Path::new(bucket).join(object); let drive_state = match disk.make_volume(vol_path.to_str().expect("operation should succeed")).await { Ok(_) => DriveState::Ok.to_string(), Err(merr) => match merr { DiskError::VolumeExists => DriveState::Ok.to_string(), DiskError::DiskNotFound => DriveState::Offline.to_string(), _ => DriveState::Corrupt.to_string(), }, }; result.after.drives[index].state = drive_state.to_string(); } } Ok((result, None)) } #[tracing::instrument(skip(self))] pub(in crate::set_disk) async fn heal_object_dir( &self, bucket: &str, object: &str, dry_run: bool, remove: bool, ) -> Result<(HealResultItem, Option)> { let _write_lock_guard = self .new_ns_lock(bucket, object) .await? .get_write_lock(get_lock_acquire_timeout()) .await .map_err(|e| DiskError::other(self.map_namespace_lock_error(bucket, object, "write", e).to_string()))?; self.heal_object_dir_locked(bucket, object, dry_run, remove).await } pub(in crate::set_disk) async fn default_heal_result( &self, lfi: FileInfo, errs: &[Option], bucket: &str, object: &str, version_id: &str, ) -> HealResultItem { // Take a single snapshot of the disk vector and drive both `disk_len` and // the per-drive loop below from it, so the reported `disk_count` and the // pushed drive records always agree (previously two independent // `self.disks.read()` calls could observe different lengths). let disks = self.disks.read().await; let disk_len = disks.len(); let mut result = HealResultItem { heal_item_type: HealItemType::Object.to_string(), bucket: bucket.to_string(), object: object.to_string(), object_size: lfi.size as usize, version_id: version_id.to_string(), disk_count: disk_len, ..Default::default() }; // Report the object's own parity only when it actually carries erasure // geometry; delete markers and geometry-less versions fall back to the // pool default. Uses `has_valid_erasure_geometry()` (not `is_valid()`) // to stay in step with the rest of the metadata-predicate migration — // `is_valid()` now requires full payload validation and returns `false` // for delete markers, which would misreport their parity here. if lfi.has_valid_erasure_geometry() { result.parity_blocks = lfi.erasure.parity_blocks; } else { result.parity_blocks = self.default_parity_count; } result.data_blocks = disk_len - result.parity_blocks; // `errs` is index-aligned with the disk vector; only the online path below // indexes into it (the offline branch `continue`s before touching it). debug_assert_eq!(errs.len(), disk_len, "errs length must match the disk count"); for (index, disk) in disks.iter().enumerate() { if disk.is_none() { result.before.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint: self.set_endpoints[index].to_string(), state: DriveState::Offline.to_string(), }); result.after.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint: self.set_endpoints[index].to_string(), state: DriveState::Offline.to_string(), }); // Offline disks contribute exactly one record; without this the // control flow fell through and pushed a second (Corrupt) record // for the same disk, doubling the list and breaking index alignment. continue; } let mut drive_state = DriveState::Corrupt; if let Some(err) = &errs[index] { if err == &DiskError::FileNotFound || err == &DiskError::VolumeNotFound { drive_state = DriveState::Missing; } } else { drive_state = DriveState::Ok; } result.before.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint: self.set_endpoints[index].to_string(), state: drive_state.to_string(), }); result.after.drives.push(HealDriveInfo { uuid: "".to_string(), endpoint: self.set_endpoints[index].to_string(), state: drive_state.to_string(), }); } result } } // Heal operation family: the storage-api `HealOperations` contract stays // implemented `for SetDisks` (contract bounds unchanged) but now lives beside // its inherent helpers in the `set_disk::ops::heal` module. Bodies are moved // unchanged; `get_pool_and_set` reads the core through `SetDisksCtx` to keep // the Heal family aligned with the borrow pattern from #816. #[async_trait::async_trait] impl crate::storage_api_contracts::heal::HealOperations for SetDisks { type Error = Error; type HealResultItem = HealResultItem; type HealOptions = HealOpts; #[tracing::instrument(skip(self))] async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { let disks = self.disks.read().await.clone(); let (formats, errs) = load_format_erasure_all(&disks, true).await; let ref_format = match get_format_erasure_in_quorum(&formats) { Ok(format) => format, Err(err) => { let can_use_cached_layout = count_errs(&errs, &DiskError::UnformattedDisk) > 0 && formats.iter().flatten().all(|format| self.format.check_other(format).is_ok()) && errs .iter() .all(|err| err.is_none() || matches!(err, Some(DiskError::UnformattedDisk))); if can_use_cached_layout { self.format.clone() } else { return Ok((HealResultItem::default(), Some(err))); } } }; let endpoints = crate::layout::endpoints::Endpoints::from(self.set_endpoints.clone()); let before_drives = crate::layout::set_heal::formats_to_drives_info(&endpoints, &formats, &errs); let mut result = HealResultItem { heal_item_type: HealItemType::Metadata.to_string(), detail: "disk-format".to_string(), disk_count: self.set_drive_count, set_count: 1, before: Infos { drives: before_drives.clone(), }, after: Infos { drives: before_drives }, ..Default::default() }; if count_errs(&errs, &DiskError::UnformattedDisk) == 0 { info!("set disk formats success, NoHealRequired, errs: {:?}", errs); return Ok((result, Some(StorageError::NoHealRequired))); } if !dry_run { for (disk_idx, err) in errs.iter().enumerate() { if !matches!(err, Some(DiskError::UnformattedDisk)) { continue; } let mut new_format = ref_format.clone(); new_format.erasure.this = ref_format.erasure.sets[self.set_index][disk_idx]; if save_format_file(&disks[disk_idx], &Some(new_format.clone())).await.is_ok() { result.after.drives[disk_idx].uuid = new_format.erasure.this.to_string(); result.after.drives[disk_idx].state = DriveState::Ok.to_string(); } } } Ok((result, None)) } #[tracing::instrument(skip(self))] async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result { let mut result = heal_bucket_local_on_disks(bucket, opts, self.disk_inventory().await).await?; result.set_count = 1; Ok(result) } #[tracing::instrument(skip(self))] async fn heal_object( &self, bucket: &str, object: &str, version_id: &str, opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { let _write_lock_guard = if !opts.no_lock { let ns_lock = self.new_ns_lock(bucket, object).await?; Some( ns_lock .get_write_lock(get_lock_acquire_timeout()) .await .map_err(|e| self.map_namespace_lock_error(bucket, object, "write", e))?, ) } else { None }; if has_suffix(object, SLASH_SEPARATOR) { let (result, err) = self.heal_object_dir_locked(bucket, object, opts.dry_run, opts.remove).await?; return Ok((result, err.map(|e| e.into()))); } let disks = self.disks.read().await; let disks = disks.clone(); let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false, false) .await .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; if DiskError::is_all_not_found(&errs) { debug!( event = EVENT_SET_DISK_HEAL, component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_SET_DISK, bucket, object, version_id, state = "missing_object_skipped", "Set disk heal skipped missing object" ); let err = if !version_id.is_empty() { Error::FileVersionNotFound } else { Error::FileNotFound }; return Ok(( self.default_heal_result(FileInfo::default(), &errs, bucket, object, version_id) .await, Some(err), )); } // Heal the object. // Pass no_lock=true since we already obtained write lock (or are already called with no_lock=true) let mut inner_opts = *opts; inner_opts.no_lock = true; let (result, err) = self .heal_object(bucket, object, version_id, &inner_opts) .await .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; if let Some(err) = err.as_ref() { match err { &DiskError::FileCorrupt if opts.scan_mode != HealScanMode::Deep => { // Instead of returning an error when a bitrot error is detected // during a normal heal scan, heal again with bitrot flag enabled. inner_opts.scan_mode = HealScanMode::Deep; let (result, err) = self .heal_object(bucket, object, version_id, &inner_opts) .await .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; return Ok((result, err.map(|e| e.into()))); } _ => {} } } Ok((result, err.map(|e| e.into()))) } #[tracing::instrument(skip(self))] async fn get_pool_and_set(&self, id: &str) -> Result<(Option, Option, Option)> { let ctx = self.ctx(); for (set_idx, set) in ctx.format().erasure.sets.iter().enumerate() { for (disk_idx, disk_id) in set.iter().enumerate() { if disk_id.to_string() == id { return Ok((Some(ctx.pool_index()), Some(set_idx), Some(disk_idx))); } } } Err(Error::DiskNotFound) } #[tracing::instrument(skip(self))] async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> { // Multipart orphan reconciliation is intentionally retained above the set layer // until there is a concrete caller and a stable lower-level contract to implement. Err(StorageError::NotImplemented) } } #[cfg(test)] mod heal_result_report_tests { use super::{DanglingCheckPartsFailure, DanglingDeleteSafety, SetDisks}; use super::{HEAL_RENAME_INCOMPLETE, HealRenameFailureScope}; use crate::disk::endpoint::Endpoint; use crate::disk::error::DiskError; use crate::disk::format::FormatV3; use crate::disk::{DiskAPI as _, DiskOption, DiskStore, RUSTFS_META_TMP_BUCKET, ReadOptions, new_disk}; use crate::object_api::{ObjectOptions, PutObjReader}; use crate::set_disk::ops::object::hermetic_set_disks_support::hermetic_set_disks_isolated; use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; use crate::{config::storageclass, store::init_format::save_format_file}; use rustfs_common::heal_channel::{DriveState, HealOpts, HealScanMode}; use rustfs_filemeta::{BLOCK_SIZE_V2, FileInfo, ObjectPartInfo, TRANSITION_COMPLETE}; use std::sync::Arc; use tempfile::TempDir; use time::OffsetDateTime; use tokio::sync::RwLock; use uuid::Uuid; async fn real_disk() -> (TempDir, Endpoint, DiskStore) { let dir = tempfile::tempdir().expect("tempdir should be created"); let endpoint = Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); let disk = new_disk( &endpoint, &DiskOption { cleanup: false, health_check: false, }, ) .await .expect("disk should be created"); (dir, endpoint, disk) } async fn set_disks_with( disks: Vec>, endpoints: Vec, default_parity_count: usize, ) -> Arc { let set_drive_count = disks.len(); SetDisks::new( "test-owner".to_string(), Arc::new(RwLock::new(disks)), set_drive_count, default_parity_count, 0, 0, endpoints, FormatV3::new(1, set_drive_count), vec![], ) .await } fn meta_regen_test_fileinfo(object: &str, data_dir: Uuid, mod_time: i64, disk_index: usize) -> FileInfo { let mut fi = FileInfo::new(object, 2, 2); fi.data_dir = Some(data_dir); fi.mod_time = Some(OffsetDateTime::from_unix_timestamp(mod_time).expect("test timestamp should parse")); fi.size = 1; fi.parts = vec![ObjectPartInfo { number: 1, size: 1, actual_size: 1, ..Default::default() }]; fi.erasure.index = fi.erasure.distribution[disk_index]; fi } async fn meta_regen_test_set( bucket: &str, object: &str, data_dirs: &[(Uuid, usize)], ) -> (Vec, Arc, Vec>) { let mut temp_dirs = Vec::new(); let mut endpoints = Vec::new(); let mut disks = Vec::new(); for disk_index in 0..4 { let (temp_dir, endpoint, disk) = real_disk().await; disk.make_volume(bucket).await.expect("test bucket should be created"); for (data_dir, shard_count) in data_dirs { if disk_index >= *shard_count { continue; } let part_dir = temp_dir.path().join(bucket).join(object).join(data_dir.to_string()); tokio::fs::create_dir_all(&part_dir) .await .expect("test data directory should be created"); tokio::fs::write(part_dir.join("part.1"), [1u8; 2]) .await .expect("test data shard should be written"); } temp_dirs.push(temp_dir); endpoints.push(endpoint); disks.push(Some(disk)); } let set = set_disks_with(disks.clone(), endpoints, 2).await; (temp_dirs, set, disks) } async fn seed_meta_regen_test_metadata( disks: &[Option], disk_index: usize, bucket: &str, object: &str, file_info: &FileInfo, ) { disks[disk_index] .as_ref() .expect("metadata test disk should be online") .write_metadata("", bucket, object, file_info.clone()) .await .expect("test metadata should be written"); } async fn formatted_single_disk_no_parity_set() -> (TempDir, Arc) { let format = FormatV3::new(1, 1); let dir = tempfile::tempdir().expect("tempdir should be created"); let mut endpoint = Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse"); endpoint.set_pool_index(0); endpoint.set_set_index(0); endpoint.set_disk_index(0); let disk = new_disk( &endpoint, &DiskOption { cleanup: false, health_check: false, }, ) .await .expect("disk should be created"); let mut disk_format = format.clone(); disk_format.erasure.this = format.erasure.sets[0][0]; save_format_file(&Some(disk.clone()), &Some(disk_format)) .await .expect("format should be saved"); let set = SetDisks::new( "test-owner".to_string(), Arc::new(RwLock::new(vec![Some(disk)])), 1, 0, 0, 0, vec![endpoint], format, vec![], ) .await; set.set_test_storage_class_config( storageclass::lookup_config_for_pools_without_env(&rustfs_config::server_config::KVS::new(), &[1]) .expect("test storage class should resolve for one local drive"), ); (dir, set) } async fn non_trash_tmp_entries(temp_dirs: &[TempDir]) -> Vec { let mut entries = Vec::new(); for dir in temp_dirs { let tmp = dir.path().join(RUSTFS_META_TMP_BUCKET); let mut read_dir = match tokio::fs::read_dir(&tmp).await { Ok(read_dir) => read_dir, Err(err) if err.kind() == std::io::ErrorKind::NotFound => continue, Err(err) => panic!("tmp directory should be readable: {err}"), }; while let Some(entry) = read_dir.next_entry().await.expect("tmp entry should be readable") { let name = entry.file_name().to_string_lossy().into_owned(); if name != ".trash" { entries.push(name); } } } entries } #[tokio::test] #[serial_test::serial] async fn heal_rename_outcome_matrix_reports_partial_and_retries_failed_targets() { for (case, failed_attempts, expect_error) in [ ("ok-ok", Vec::new(), false), ("ok-err", vec![1], true), ("err-ok", vec![0], true), ("err-err", vec![0, 1], true), ] { let (temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await; let bucket = format!("heal-rename-{case}"); let object = "object.bin"; for disk in &disks { disk.make_volume(&bucket).await.expect("bucket volume should be created"); } let payload = vec![0x5a; 1024 * 1024]; let mut reader = PutObjReader::from_vec(payload); set.put_object(&bucket, object, &mut reader, &ObjectOptions::default()) .await .expect("source object should be written"); let source = disks[2] .read_version("", &bucket, object, "", &ReadOptions::default()) .await .expect("source metadata should be readable"); let data_dir = source.data_dir.expect("non-inline source should have a data directory"); let tmp_entries_before_heal = non_trash_tmp_entries(&temp_dirs).await; let target_slots = { let mut slots = [source.erasure.distribution[0] - 1, source.erasure.distribution[1] - 1]; slots.sort_unstable(); slots }; let failed_slots = failed_attempts .iter() .map(|attempt| target_slots[*attempt]) .collect::>(); let failed_physical_indexes = [0, 1] .into_iter() .filter(|index| failed_slots.contains(&(source.erasure.distribution[*index] - 1))) .collect::>(); for index in [0, 1] { tokio::fs::remove_file( temp_dirs[index] .path() .join(&bucket) .join(object) .join(data_dir.to_string()) .join("part.1"), ) .await .expect("target shard should be removed before heal"); } let failure_scope = HealRenameFailureScope::install(&bucket, object, &failed_slots); let (first_result, first_error) = set .heal_object( &bucket, object, "", &HealOpts { no_lock: true, scan_mode: HealScanMode::Deep, ..Default::default() }, ) .await .expect("heal should report its per-target rename outcome"); drop(failure_scope); assert_eq!(first_error.is_some(), expect_error, "{case}: aggregate status must match target outcomes"); if let Some(error) = first_error { let error = error.to_string(); assert!( error.contains(HEAL_RENAME_INCOMPLETE), "{case}: partial/all failure must have an explicit retryable status: {error}" ); assert!( error.contains(&format!("{} of 2 targets committed", 2 - failed_slots.len())), "{case}: aggregate status must distinguish partial from all-target failure: {error}" ); } for index in [0, 1] { let expected = if failed_physical_indexes.contains(&index) { DriveState::Missing } else { DriveState::Ok }; assert_eq!( first_result.after.drives[index].state, expected.to_string(), "{case}: after.drives must reflect the actual rename outcome at index {index}" ); assert_eq!( temp_dirs[index] .path() .join(&bucket) .join(object) .join(data_dir.to_string()) .join("part.1") .exists(), !failed_physical_indexes.contains(&index), "{case}: tmp cleanup must neither delete committed shards nor expose failed targets" ); } let tmp_entries_after_heal = non_trash_tmp_entries(&temp_dirs).await; assert!( tmp_entries_after_heal .iter() .all(|entry| tmp_entries_before_heal.contains(entry)), "{case}: first heal must not leave a new temporary shard: {tmp_entries_after_heal:?}" ); if !failed_slots.is_empty() { let (retry_result, retry_error) = set .heal_object( &bucket, object, "", &HealOpts { no_lock: true, scan_mode: HealScanMode::Deep, ..Default::default() }, ) .await .expect("second heal should retry failed targets"); assert!(retry_error.is_none(), "{case}: second heal should complete remaining targets"); for index in [0, 1] { assert_eq!( retry_result.after.drives[index].state, DriveState::Ok.to_string(), "{case}: second heal must converge target {index}" ); } let tmp_entries_after_retry = non_trash_tmp_entries(&temp_dirs).await; assert!( tmp_entries_after_retry .iter() .all(|entry| tmp_entries_before_heal.contains(entry)), "{case}: retry must not leave a new temporary shard: {tmp_entries_after_retry:?}" ); } } } // Regression for #955: an offline disk must contribute exactly one drive // record. Before the fix the offline branch fell through and pushed a second // (Corrupt) record for the same disk, so `before/after.drives` grew to // `disk_count + offline_count` and every entry after the first offline slot // was misaligned relative to its disk index. #[tokio::test] async fn default_heal_result_reports_one_record_per_disk_and_stays_aligned() { // index 0: online, no error -> Ok // index 1: offline (None) -> Offline (single record) // index 2: online, FileNotFound -> Missing // index 3: online, DiskAccessDenied-> Corrupt let (_d0, ep0, disk0) = real_disk().await; let (_d2, ep2, disk2) = real_disk().await; let (_d3, ep3, disk3) = real_disk().await; let ep1 = Endpoint::try_from("http://127.0.0.1:9001/data").expect("endpoint should parse"); let disks = vec![Some(disk0), None, Some(disk2), Some(disk3)]; let endpoints = vec![ep0, ep1, ep2, ep3]; let set = set_disks_with(disks, endpoints, 1).await; let errs = vec![ None, Some(DiskError::DiskNotFound), Some(DiskError::FileNotFound), Some(DiskError::DiskAccessDenied), ]; let result = set .default_heal_result(FileInfo::default(), &errs, "bucket", "object", "") .await; // Exactly one record per disk (not disk_count + offline_count). assert_eq!(result.disk_count, 4); assert_eq!(result.before.drives.len(), 4, "one before record per disk"); assert_eq!(result.after.drives.len(), 4, "one after record per disk"); // Records stay index-aligned with the disk vector and set_endpoints. let expected_states = [ DriveState::Ok.to_string(), DriveState::Offline.to_string(), DriveState::Missing.to_string(), DriveState::Corrupt.to_string(), ]; for (i, expected) in expected_states.iter().enumerate() { assert_eq!(&result.before.drives[i].state, expected, "before state at {i}"); assert_eq!(&result.after.drives[i].state, expected, "after state at {i}"); assert_eq!( result.before.drives[i].endpoint, set.set_endpoints[i].to_string(), "before endpoint aligned at {i}" ); assert_eq!( result.after.drives[i].endpoint, set.set_endpoints[i].to_string(), "after endpoint aligned at {i}" ); } // The offline endpoint appears exactly once, never as a second Corrupt row. let offline_ep = set.set_endpoints[1].to_string(); assert_eq!( result.before.drives.iter().filter(|d| d.endpoint == offline_ep).count(), 1, "offline disk must not produce a duplicate record" ); } // Two interleaved offline disks: assert every record still maps to its own // set_endpoints[index] (no cumulative drift after the first offline slot). #[tokio::test] async fn default_heal_result_alignment_with_multiple_offline_disks() { let (_d1, ep1, disk1) = real_disk().await; let (_d3, ep3, disk3) = real_disk().await; let ep0 = Endpoint::try_from("http://127.0.0.1:9000/data").expect("endpoint should parse"); let ep2 = Endpoint::try_from("http://127.0.0.1:9002/data").expect("endpoint should parse"); // index 0 offline, 1 online, 2 offline, 3 online. let disks = vec![None, Some(disk1), None, Some(disk3)]; let endpoints = vec![ep0, ep1, ep2, ep3]; let set = set_disks_with(disks, endpoints, 1).await; let errs = vec![Some(DiskError::DiskNotFound), None, Some(DiskError::DiskNotFound), None]; let result = set .default_heal_result(FileInfo::default(), &errs, "bucket", "object", "") .await; assert_eq!(result.before.drives.len(), 4); assert_eq!(result.after.drives.len(), 4); for i in 0..4 { assert_eq!(result.before.drives[i].endpoint, set.set_endpoints[i].to_string(), "aligned at {i}"); } assert_eq!(result.before.drives[0].state, DriveState::Offline.to_string()); assert_eq!(result.before.drives[1].state, DriveState::Ok.to_string()); assert_eq!(result.before.drives[2].state, DriveState::Offline.to_string()); assert_eq!(result.before.drives[3].state, DriveState::Ok.to_string()); } #[tokio::test] async fn dangling_delete_guard_preserves_conflicting_identities_without_writing_metadata() { let bucket = "bucket-delete-guard-conflict"; let object = "object.bin"; let old_data_dir = Uuid::parse_str("99999999-9999-9999-9999-999999999999").expect("old data dir should parse"); let new_data_dir = Uuid::parse_str("aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa").expect("new data dir should parse"); let (_temp_dirs, set, disks) = meta_regen_test_set(bucket, object, &[(old_data_dir, 4), (new_data_dir, 2)]).await; let version_id = Uuid::parse_str("bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb").expect("version id should parse"); let mut metadata = vec![ meta_regen_test_fileinfo(object, old_data_dir, 9, 0), meta_regen_test_fileinfo(object, new_data_dir, 10, 1), FileInfo::default(), FileInfo::default(), ]; metadata[0].version_id = Some(version_id); metadata[1].version_id = Some(version_id); assert_eq!( metadata[0].version_id, metadata[1].version_id, "the conflicting candidates must share one version id" ); seed_meta_regen_test_metadata(&disks, 0, bucket, object, &metadata[0]).await; seed_meta_regen_test_metadata(&disks, 1, bucket, object, &metadata[1]).await; let errs = vec![None, None, Some(DiskError::FileNotFound), Some(DiskError::FileNotFound)]; assert!( set.dangling_delete_safety(bucket, object, &metadata, &errs, &disks) .await .expect("conflicting identities should be classified") == DanglingDeleteSafety::UnsafeToDelete ); let reversed = vec![ metadata[1].clone(), metadata[0].clone(), FileInfo::default(), FileInfo::default(), ]; assert!( set.dangling_delete_safety(bucket, object, &reversed, &errs, &disks) .await .expect("reversed identities should be classified") == DanglingDeleteSafety::UnsafeToDelete ); let version_id = version_id.to_string(); assert!( !set.try_regenerate_explicit_version_meta(bucket, object, &version_id, &metadata, &errs, &disks) .await .expect("conflicting explicit-version candidates should be rejected"), "an explicit version must not select between conflicting metadata identities" ); for disk_index in [2, 3] { assert!( matches!( disks[disk_index] .as_ref() .expect("test disk should be online") .read_version("", bucket, object, "", &ReadOptions::default()) .await, Err(DiskError::FileNotFound) ), "the delete guard must not manufacture metadata on missing disks" ); } let old = disks[0] .as_ref() .expect("first test disk should be online") .read_version("", bucket, object, "", &ReadOptions::default()) .await .expect("old metadata should remain readable"); let new = disks[1] .as_ref() .expect("second test disk should be online") .read_version("", bucket, object, "", &ReadOptions::default()) .await .expect("new metadata should remain readable"); assert_eq!(old.data_dir, Some(old_data_dir)); assert_eq!(new.data_dir, Some(new_data_dir)); } #[tokio::test] async fn heal_meta_quorum_failure_preserves_reconstructable_uncommitted_candidate() { let bucket = "bucket-delete-guard-reconstructable"; let object = "object.bin"; let data_dir = Uuid::parse_str("33333333-3333-3333-3333-333333333333").expect("data dir should parse"); let (_temp_dirs, set, disks) = meta_regen_test_set(bucket, object, &[(data_dir, 2)]).await; let metadata = [ meta_regen_test_fileinfo(object, data_dir, 3, 0), FileInfo::default(), FileInfo::default(), FileInfo::default(), ]; seed_meta_regen_test_metadata(&disks, 0, bucket, object, &metadata[0]).await; let (observed_metadata, observed_errs) = SetDisks::read_all_fileinfo(&disks, "", bucket, object, "", true, true, false) .await .expect("test metadata should be readable across the set"); assert_eq!( set.dangling_delete_safety(bucket, object, &observed_metadata, &observed_errs, &disks) .await .expect("observed reconstructable candidate should be classified"), DanglingDeleteSafety::UnsafeToDelete ); let (_, err) = set .heal_object( bucket, object, "", &HealOpts { no_lock: true, ..Default::default() }, ) .await .expect("unsafe dangling state should be reported without deletion"); assert_eq!(err, Some(DiskError::FileNotFound)); let surviving = disks[0] .as_ref() .expect("first test disk should be online") .read_version("", bucket, object, "", &ReadOptions::default()) .await .expect("the only metadata copy must be preserved"); assert_eq!(surviving.data_dir, Some(data_dir)); assert!( matches!( disks[1] .as_ref() .expect("second test disk should be online") .read_version("", bucket, object, "", &ReadOptions::default()) .await, Err(DiskError::FileNotFound) ), "the delete guard must not propagate metadata" ); } #[tokio::test] async fn heal_meta_quorum_failure_preserves_candidate_when_required_shard_disk_is_offline() { let bucket = "bucket-delete-guard-offline"; let object = "object.bin"; let data_dir = Uuid::parse_str("44444444-4444-4444-4444-444444444444").expect("data dir should parse"); let (temp_dirs, set, disks) = meta_regen_test_set(bucket, object, &[(data_dir, 2)]).await; let metadata = meta_regen_test_fileinfo(object, data_dir, 4, 0); seed_meta_regen_test_metadata(&disks, 0, bucket, object, &metadata).await; set.disks.write().await[1] = None; let (_, err) = set .heal_object( bucket, object, "", &HealOpts { no_lock: true, ..Default::default() }, ) .await .expect("offline shard state should be reported without deletion"); assert_eq!(err, Some(DiskError::FileNotFound)); let surviving = disks[0] .as_ref() .expect("first test disk should be online") .read_version("", bucket, object, "", &ReadOptions::default()) .await .expect("offline uncertainty must preserve the surviving metadata"); assert_eq!(surviving.data_dir, Some(data_dir)); assert!( temp_dirs[0] .path() .join(bucket) .join(object) .join(data_dir.to_string()) .join("part.1") .is_file(), "offline uncertainty must preserve the last online shard" ); } #[tokio::test] async fn heal_meta_quorum_failure_preserves_candidate_when_part_probe_times_out() { let bucket = "bucket-delete-guard-timeout"; let object = "object.bin"; let data_dir = Uuid::parse_str("55555555-5555-5555-5555-555555555555").expect("data dir should parse"); let (temp_dirs, set, disks) = meta_regen_test_set(bucket, object, &[(data_dir, 2)]).await; let metadata = meta_regen_test_fileinfo(object, data_dir, 5, 0); seed_meta_regen_test_metadata(&disks, 0, bucket, object, &metadata).await; let _failure = DanglingCheckPartsFailure::install(bucket, object, 1, DiskError::Timeout); let (_, err) = set .heal_object( bucket, object, "", &HealOpts { no_lock: true, ..Default::default() }, ) .await .expect("part probe timeout should be reported without deletion"); assert_eq!(err, Some(DiskError::FileNotFound)); let surviving = disks[0] .as_ref() .expect("first test disk should be online") .read_version("", bucket, object, "", &ReadOptions::default()) .await .expect("probe uncertainty must preserve the surviving metadata"); assert_eq!(surviving.data_dir, Some(data_dir)); assert!( temp_dirs[0] .path() .join(bucket) .join(object) .join(data_dir.to_string()) .join("part.1") .is_file(), "probe uncertainty must preserve the last confirmed shard" ); } #[tokio::test] async fn dangling_delete_guard_ignores_set_incompatible_geometry() { let bucket = "bucket-delete-guard-short-geometry"; let object = "object.bin"; let data_dir = Uuid::parse_str("abababab-abab-abab-abab-abababababab").expect("data dir should parse"); let (_temp_dirs, set, disks) = meta_regen_test_set(bucket, object, &[(data_dir, 1)]).await; let mut candidate = FileInfo::new(object, 1, 0); candidate.data_dir = Some(data_dir); candidate.mod_time = Some(OffsetDateTime::from_unix_timestamp(18).expect("timestamp should parse")); candidate.size = 1; candidate.parts = vec![ObjectPartInfo { number: 1, size: 1, actual_size: 1, ..Default::default() }]; candidate.erasure.index = candidate.erasure.distribution[0]; seed_meta_regen_test_metadata(&disks, 0, bucket, object, &candidate).await; let metadata = vec![candidate, FileInfo::default(), FileInfo::default(), FileInfo::default()]; let errs = vec![ None, Some(DiskError::FileNotFound), Some(DiskError::FileNotFound), Some(DiskError::FileNotFound), ]; assert!( set.dangling_delete_safety(bucket, object, &metadata, &errs, &disks) .await .expect("set-incompatible geometry should be classified") == DanglingDeleteSafety::NoRecoverableCandidate ); } #[tokio::test] async fn dangling_delete_guard_preserves_delete_marker_and_remote_metadata() { let bucket = "bucket-delete-guard-nonlocal"; let object = "object.bin"; let (_temp_dirs, set, disks) = meta_regen_test_set(bucket, object, &[]).await; let marker = FileInfo { name: object.to_string(), version_id: Some(Uuid::parse_str("eeeeeeee-eeee-eeee-eeee-eeeeeeeeeeee").expect("version id should parse")), deleted: true, mod_time: Some(OffsetDateTime::from_unix_timestamp(14).expect("marker timestamp should parse")), ..Default::default() }; let remote_dir = Uuid::parse_str("89898989-8989-8989-8989-898989898989").expect("remote data dir should parse"); let mut remote = meta_regen_test_fileinfo(object, remote_dir, 15, 1); remote.transition_status = TRANSITION_COMPLETE.to_string(); remote.transition_tier = "WARM".to_string(); remote.transitioned_objname = "remote/object.bin".to_string(); for metadata in [marker, remote] { let candidates = vec![metadata, FileInfo::default(), FileInfo::default(), FileInfo::default()]; let errs = vec![ None, Some(DiskError::FileNotFound), Some(DiskError::FileNotFound), Some(DiskError::FileNotFound), ]; assert_eq!( set.dangling_delete_safety(bucket, object, &candidates, &errs, &disks) .await .expect("non-local metadata should be classified"), DanglingDeleteSafety::UnsafeToDelete ); } } #[tokio::test] async fn dangling_delete_guard_preserves_metadata_read_uncertainty() { let bucket = "bucket-delete-guard-read-error"; let object = "object.bin"; let (_temp_dirs, set, disks) = meta_regen_test_set(bucket, object, &[]).await; let metadata = vec![FileInfo::default(); disks.len()]; for read_error in [DiskError::Timeout, DiskError::DiskAccessDenied, DiskError::DiskNotFound] { let mut errs = vec![Some(DiskError::FileNotFound); disks.len()]; errs[0] = Some(read_error); assert_eq!( set.dangling_delete_safety(bucket, object, &metadata, &errs, &disks) .await .expect("metadata read uncertainty should be classified"), DanglingDeleteSafety::UnsafeToDelete ); } } #[tokio::test] async fn heal_no_parity_bitrot_reports_unrecoverable_integrity_failure() { let (dir, set) = formatted_single_disk_no_parity_set().await; let bucket = "bucket-no-parity-bitrot"; let object = "bad-object.bin"; let payload = (0..(BLOCK_SIZE_V2 + 17)).map(|idx| (idx % 251) as u8).collect::>(); let opts = ObjectOptions { no_lock: true, ..Default::default() }; set.make_bucket(bucket, &MakeBucketOptions::default()) .await .expect("bucket should be created"); let mut reader = PutObjReader::from_vec(payload); set.put_object(bucket, object, &mut reader, &opts) .await .expect("object should be written"); let (fi, _, _) = set .get_object_fileinfo(bucket, object, &opts, true, false) .await .expect("object metadata should resolve"); assert_eq!(fi.erasure.parity_blocks, 0); let data_dir = fi.data_dir.expect("non-inline object should have a data directory"); let part_path = dir.path().join(bucket).join(object).join(data_dir.to_string()).join("part.1"); let mut part = tokio::fs::read(&part_path).await.expect("part should be readable"); part[0] ^= 0xff; tokio::fs::write(&part_path, part) .await .expect("part corruption should be written"); let (result, err) = set .heal_object( bucket, object, "", &HealOpts { no_lock: true, scan_mode: HealScanMode::Deep, ..Default::default() }, ) .await .expect("heal should report the unrecoverable object without panicking"); assert_eq!(err, Some(DiskError::FileCorrupt)); assert_eq!(result.bucket, bucket); assert_eq!(result.object, object); assert_eq!(result.data_blocks, 1); assert_eq!(result.parity_blocks, 0); assert_eq!(result.before.drives[0].state, DriveState::Corrupt.to_string()); assert!(result.detail.contains("no-parity object is unrecoverable")); assert!(result.detail.contains("part 1")); assert!(result.detail.contains("bitrot_failure=true")); } }