diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 65995b95a..7b295eec7 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -589,6 +589,7 @@ pub mod storage { all_local_disk_path, find_local_disk_by_ref, init_local_disks, init_local_disks_with_instance_ctx, init_lock_clients, prewarm_local_disk_id_map, prewarm_local_disk_id_map_with_instance_ctx, }; + pub use crate::store::{HealObjectAbsenceProof, HealObjectStorageResult}; } pub mod tier { diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index bfe88aae1..2a09bf0c7 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -1475,6 +1475,28 @@ pub(crate) async fn make_local_two_set_sets_for_pool_with_drive_count_and_ctx( (temp_dirs, sets) } +impl Sets { + pub(crate) async fn heal_object_with_absence( + &self, + bucket: &str, + object: &str, + version_id: &str, + opts: &HealOpts, + ) -> Result<(HealResultItem, Option, Option)> { + let mut absence = None; + let (item, error) = self + .get_disks_for_heal_object(object, opts)? + .heal_object_with_absence(bucket, object, version_id, opts, &mut absence) + .await?; + // A caller-owned lock does not expose its lease to this boundary. + // Keep cleanup unverified when that lease cannot be checked here. + if opts.no_lock { + absence = None; + } + Ok((item, error, absence)) + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/ecstore/src/disk/error.rs b/crates/ecstore/src/disk/error.rs index 8c00ab93e..37c959608 100644 --- a/crates/ecstore/src/disk/error.rs +++ b/crates/ecstore/src/disk/error.rs @@ -35,7 +35,7 @@ pub(crate) struct TerminalReadError { source: DiskError, } -#[derive(Debug)] +#[derive(Debug, Clone)] struct DanglingDeleteGraceError { retry_after_secs: i64, grace_secs: i64, @@ -329,6 +329,23 @@ impl DiskError { ) } + pub(crate) fn clone_dangling_delete_grace(error: &io::Error) -> Option { + let grace = error.get_ref()?.downcast_ref::()?; + Some(io::Error::new(error.kind(), grace.clone())) + } + + pub fn dangling_delete_retry_after(&self) -> Option { + match self { + Self::Io(error) => Self::io_error_dangling_delete_retry_after(error), + _ => None, + } + } + + pub fn io_error_dangling_delete_retry_after(error: &io::Error) -> Option { + let grace = error.get_ref()?.downcast_ref::()?; + u64::try_from(grace.retry_after_secs).ok().map(std::time::Duration::from_secs) + } + pub fn is_dangling_delete_grace(&self) -> bool { matches!(self, DiskError::Io(io_error) if Self::io_error_is_dangling_delete_grace(io_error)) } @@ -667,7 +684,8 @@ impl Clone for DiskError { DiskError::conditional_file_not_committed(io::Error::new(io_error.kind(), io_error.to_string())), ), DiskError::Io(io_error) => DiskError::Io( - rustfs_rio::clone_internode_http_io_error(io_error) + Self::clone_dangling_delete_grace(io_error) + .or_else(|| rustfs_rio::clone_internode_http_io_error(io_error)) .and_then(std::io::Error::into_inner) // The helper derives a kind from the source; Clone must retain the original outer kind. .map(|source| std::io::Error::new(io_error.kind(), source)) @@ -859,6 +877,21 @@ mod tests { use super::*; use std::collections::HashMap; + #[test] + fn dangling_grace_retry_timing_survives_disk_and_storage_clones() { + let original = super::DiskError::dangling_delete_grace(21, 3600); + let disk = original.clone(); + assert_eq!(disk.dangling_delete_retry_after(), Some(std::time::Duration::from_secs(21))); + let storage: crate::error::StorageError = disk.into(); + let cloned = storage.clone(); + assert!(cloned.is_dangling_delete_grace()); + assert_eq!(cloned.dangling_delete_retry_after(), Some(std::time::Duration::from_secs(21))); + assert_eq!(original.dangling_delete_retry_after(), Some(std::time::Duration::from_secs(21))); + assert_eq!(storage.dangling_delete_retry_after(), Some(std::time::Duration::from_secs(21))); + assert_eq!(super::DiskError::dangling_delete_grace(-1, 3600).dangling_delete_retry_after(), None); + assert_eq!(super::DiskError::FaultyDisk.dangling_delete_retry_after(), None); + } + #[test] fn conditional_file_not_committed_marker_is_explicit_and_clone_safe() { let marked = DiskError::from(DiskError::conditional_file_not_committed(io::Error::new( diff --git a/crates/ecstore/src/error/mod.rs b/crates/ecstore/src/error/mod.rs index 281c70a7e..5916e2163 100644 --- a/crates/ecstore/src/error/mod.rs +++ b/crates/ecstore/src/error/mod.rs @@ -396,6 +396,13 @@ impl StorageError { ) } + pub fn dangling_delete_retry_after(&self) -> Option { + match self { + Self::Io(error) => DiskError::io_error_dangling_delete_retry_after(error), + _ => None, + } + } + pub fn is_dangling_delete_grace(&self) -> bool { matches!(self, StorageError::Io(io_error) if DiskError::io_error_is_dangling_delete_grace(io_error)) } @@ -608,6 +615,9 @@ impl Clone for StorageError { fn clone(&self) -> Self { match self { StorageError::Io(e) => { + if let Some(error) = DiskError::clone_dangling_delete_grace(e) { + return StorageError::Io(error); + } if let Some(context) = self.pool_metadata_failure() { Self::Io(std::io::Error::new( e.kind(), diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index b74a397e4..9f12d53d8 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -5979,6 +5979,20 @@ impl SetDisks { data_errs_by_part: &HashMap>, opts: ObjectOptions, ) -> disk::error::Result { + self.delete_if_dangling_with_proof(bucket, object, meta_arr, errs, data_errs_by_part, opts) + .await + .map(|(metadata, _)| metadata) + } + + pub(in crate::set_disk) async fn delete_if_dangling_with_proof( + &self, + bucket: &str, + object: &str, + meta_arr: &[FileInfo], + errs: &[Option], + data_errs_by_part: &HashMap>, + opts: ObjectOptions, + ) -> disk::error::Result<(FileInfo, bool)> { let (m, can_heal) = is_object_dangling(meta_arr, errs, data_errs_by_part); if !can_heal { @@ -6065,11 +6079,17 @@ impl SetDisks { let disks = self.get_disks_internal().await; let mut futures = Vec::with_capacity(disks.len()); - for disk_op in disks.iter() { + for (disk_index, disk_op) in disks.iter().enumerate() { + #[cfg(not(test))] + let _ = disk_index; let bucket = bucket.to_string(); let object = object.to_string(); let fi = fi.clone(); futures.push(async move { + #[cfg(test)] + if let Some(error) = crate::set_disk::ops::heal::injected_dangling_delete_error(&bucket, &object, disk_index) { + return Err(error); + } if let Some(disk) = disk_op { disk.delete_version(&bucket, &object, fi, false, DeleteOptions::default()) .await @@ -6080,6 +6100,7 @@ impl SetDisks { } let results = join_all(futures).await; + let mut all_deleted = !results.is_empty(); let mut delete_errs = Vec::with_capacity(results.len()); for (index, result) in results.into_iter().enumerate() { let key = format!("ddisk-{index}"); @@ -6093,6 +6114,7 @@ impl SetDisks { delete_errs.push(None); } Err(e) => { + all_deleted &= matches!(&e, DiskError::FileNotFound | DiskError::FileVersionNotFound); tags.insert(key, e.to_string()); if already_absent || matches!(&e, DiskError::FileNotFound | DiskError::FileVersionNotFound) { delete_errs.push(None); @@ -6112,7 +6134,30 @@ impl SetDisks { return Err(err); } - Ok(m) + // Quorum success alone may leave the only stale replica behind. The + // proof uses the same disk snapshot as deletion and exact-version reads. + let absent = if all_deleted { + match Self::read_all_fileinfo( + &disks, + "", + bucket, + object, + opts.version_id.as_deref().unwrap_or(""), + false, + false, + false, + ) + .await + { + Ok((_, after)) => after + .iter() + .all(|err| matches!(err, Some(DiskError::FileNotFound | DiskError::FileVersionNotFound))), + Err(_) => false, + } + } else { + false + }; + Ok((m, absent)) } fn reduce_delete_prefix_results(results: Vec>, write_quorum: usize) -> disk::error::Result<()> { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 363f261d5..fa46f6110 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -866,6 +866,7 @@ mod ctx; mod metadata; mod ops; pub(crate) use ops::bucket::BucketInfoQuorum; +pub(crate) use ops::heal::HealedObjectAbsence; #[cfg(test)] pub(crate) use ops::hermetic_set_disks_isolated; diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 73e18e4aa..65ae1bd36 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -36,6 +36,14 @@ const EVENT_HEAL_OBJECT_RENAME: &str = "heal_object_rename"; const HEAL_RENAME_INCOMPLETE: &str = "heal rename incomplete"; const READ_REPAIR_DATA_PHASE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60 * 60); +/// Exact set-local absence, established while the object mutation lock is held. +#[derive(Debug)] +pub(crate) struct HealedObjectAbsence { + pub pool_index: usize, + pub set_index: usize, + pub removed: bool, +} + fn heal_drive_state_for_error(error: &DiskError) -> DriveState { match error { DiskError::DiskNotFound | DiskError::RemoteClientUnavailable(_) => DriveState::Offline, @@ -391,7 +399,7 @@ impl Drop for DanglingDeleteFailure { } #[cfg(test)] -fn injected_dangling_delete_error(bucket: &str, object: &str, disk_index: usize) -> Option { +pub(in crate::set_disk) fn injected_dangling_delete_error(bucket: &str, object: &str, disk_index: usize) -> Option { dangling_delete_failures() .lock() .expect("dangling delete failure registry should not poison") @@ -541,6 +549,7 @@ impl SetDisks { .all(|committed| committed)) } + #[cfg(test)] #[tracing::instrument(level = "trace", skip(self, opts), fields(bucket = %bucket, object = %object, version_id = %version_id))] pub(in crate::set_disk) async fn heal_object( &self, @@ -549,7 +558,7 @@ impl SetDisks { 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 + Box::pin(self.heal_object_with_explicit_version_regen(bucket, object, version_id, opts, true, &mut None)).await } async fn read_repair_commit_fingerprint( @@ -661,6 +670,7 @@ impl SetDisks { version_id: &str, opts: &HealOpts, allow_explicit_version_regen: bool, + absence: &mut Option, ) -> disk::error::Result<(HealResultItem, Option)> { trace!( event = EVENT_SET_DISK_HEAL, @@ -1025,7 +1035,7 @@ impl SetDisks { // 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( + .delete_if_dangling_with_proof( bucket, object, &parts_metadata, @@ -1038,12 +1048,18 @@ impl SetDisks { ) .await { - Ok(m) => { - let mut t_errs = Vec::with_capacity(errs.len()); - for _ in 0..errs.len() { - t_errs.push(None); + Ok((m, absent)) => { + if absent { + *absence = Some(HealedObjectAbsence { + pool_index: self.pool_index, + set_index: self.set_index, + removed: true, + }); } - Ok((self.default_heal_result(m, &t_errs, bucket, object, version_id).await, None)) + Ok(( + self.dangling_heal_result(m, &errs, bucket, object, version_id, absent).await, + (!absent).then_some(DiskError::ErasureWriteQuorum), + )) } Err(err) => { error!( @@ -1562,7 +1578,10 @@ impl SetDisks { .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; + return Box::pin( + self.heal_object_with_explicit_version_regen(bucket, object, version_id, opts, false, absence), + ) + .await; } if opts.dry_run { @@ -1587,7 +1606,7 @@ impl SetDisks { let data_errs_by_part = HashMap::new(); match self - .delete_if_dangling( + .delete_if_dangling_with_proof( bucket, object, &parts_metadata, @@ -1600,7 +1619,19 @@ impl SetDisks { ) .await { - Ok(m) => Ok((self.default_heal_result(m, &errs, bucket, object, version_id).await, None)), + Ok((m, absent)) => { + if absent { + *absence = Some(HealedObjectAbsence { + pool_index: self.pool_index, + set_index: self.set_index, + removed: true, + }); + } + Ok(( + self.dangling_heal_result(m, &errs, bucket, object, version_id, absent).await, + (!absent).then_some(DiskError::ErasureWriteQuorum), + )) + } Err(cleanup_err) => Ok(( self.default_heal_result(FileInfo::default(), &errs, bucket, object, version_id) .await, @@ -2465,95 +2496,8 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks { 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 - .map_err(|e| e.narrow_to_disk().unwrap_or_else(DiskError::other))?; - Some(ns_lock.get_write_lock(get_lock_acquire_timeout()).await.map_err(|e| { - self.map_namespace_lock_error(bucket, object, "write", e) - .narrow_to_disk() - .unwrap_or_else(DiskError::other) - })?) - } 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()))); - } - - // The inner heal and missing-object report read the registry again; - // release this snapshot guard before a topology writer can queue between reads. - let disks = self.get_disks_internal().await; - let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false, false) + self.heal_object_with_absence(bucket, object, version_id, opts, &mut None) .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 - }; - if version_id.is_empty() - && (opts.remove || opts.dry_run) - && let Some(cleanup) = self - .cleanup_metadata_less_data_dirs(bucket, object, &disks, opts.dry_run) - .await - .map_err(|e| to_object_err(e.into(), vec![bucket, object]))? - { - let mut result = self - .metadata_less_data_dir_heal_result(bucket, object, &cleanup, opts.dry_run) - .await; - result.detail = format!( - "metadata-less data directories matched={}, removed={}, dry_run={}", - cleanup.matched, cleanup.removed, opts.dry_run - ); - let err = cleanup.first_error.map(Error::from).or(Some(err)); - return Ok((result, err)); - } - 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))] @@ -2630,6 +2574,149 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks { } } +impl SetDisks { + pub(crate) async fn heal_object_with_absence( + &self, + bucket: &str, + object: &str, + version_id: &str, + opts: &HealOpts, + absence: &mut Option, + ) -> Result<(HealResultItem, Option)> { + *absence = None; + let _write_lock_guard = if !opts.no_lock { + let ns_lock = self + .new_ns_lock(bucket, object) + .await + .map_err(|e| e.narrow_to_disk().unwrap_or_else(DiskError::other))?; + Some(ns_lock.get_write_lock(get_lock_acquire_timeout()).await.map_err(|e| { + self.map_namespace_lock_error(bucket, object, "write", e) + .narrow_to_disk() + .unwrap_or_else(DiskError::other) + })?) + } 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()))); + } + + // The inner heal and missing-object report read the registry again; + // release this snapshot guard before a topology writer can queue between reads. + let disks = self.get_disks_internal().await; + let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false, false) + .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 + }; + if version_id.is_empty() + && (opts.remove || opts.dry_run) + && let Some(cleanup) = self + .cleanup_metadata_less_data_dirs(bucket, object, &disks, opts.dry_run) + .await + .map_err(|e| to_object_err(e.into(), vec![bucket, object]))? + { + let mut result = self + .metadata_less_data_dir_heal_result(bucket, object, &cleanup, opts.dry_run) + .await; + result.detail = format!( + "metadata-less data directories matched={}, removed={}, dry_run={}", + cleanup.matched, cleanup.removed, opts.dry_run + ); + let err = cleanup.first_error.map(Error::from).or(Some(err)); + return Ok((result, err)); + } + let result = self + .default_heal_result(FileInfo::default(), &errs, bucket, object, version_id) + .await; + // Check the lease after the final await before publishing the proof. + if !opts.dry_run + && !version_id.is_empty() + && !disks.is_empty() + && disks.iter().all(Option::is_some) + && errs + .iter() + .all(|err| matches!(err, Some(DiskError::FileNotFound | DiskError::FileVersionNotFound))) + && !_write_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) + { + *absence = Some(HealedObjectAbsence { + pool_index: self.pool_index, + set_index: self.set_index, + removed: false, + }); + } + return Ok((result, 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_with_explicit_version_regen(bucket, object, version_id, &inner_opts, true, absence) + .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_with_explicit_version_regen(bucket, object, version_id, &inner_opts, true, absence) + .await + .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; + if _write_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) { + *absence = None; + } + return Ok((result, err.map(|e| e.into()))); + } + _ => {} + } + } + if _write_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) { + *absence = None; + } + Ok((result, err.map(|e| e.into()))) + } +} + +impl SetDisks { + async fn dangling_heal_result( + &self, + metadata: FileInfo, + errs: &[Option], + bucket: &str, + object: &str, + version_id: &str, + absent: bool, + ) -> HealResultItem { + let mut item = self.default_heal_result(metadata, errs, bucket, object, version_id).await; + if absent { + for drive in &mut item.after.drives { + drive.state = DriveState::Missing.to_string(); + } + } + item + } +} + #[cfg(test)] mod heal_result_report_tests { use super::{ @@ -4629,6 +4716,71 @@ mod heal_result_report_tests { ); } + #[tokio::test] + #[serial_test::serial] + async fn dangling_absence_proof_rejects_failed_stale_replica_and_accepts_retry() { + temp_env::async_with_vars([("RUSTFS_HEAL_DANGLING_DELETE_GRACE_SECS", Some("0"))], async { + let bucket = "dangling-absence-partial-delete"; + let object = "history.txt"; + let (_temp_dirs, set, disks) = + dangling_inline_test_fixture(bucket, object, OffsetDateTime::now_utc() - time::Duration::hours(2)).await; + let version = Uuid::new_v4(); + let disk = disks[0].as_ref().expect("stale disk must be online"); + let mut metadata = disk + .read_version("", bucket, object, "", &ReadOptions::default()) + .await + .expect("load stale inline metadata"); + metadata.version_id = Some(version); + disk.write_metadata("", bucket, object, metadata) + .await + .expect("seed exact stale historical version"); + let opts = HealOpts { + no_lock: true, + scan_mode: HealScanMode::Deep, + ..Default::default() + }; + let failure = DanglingDeleteFailure::install(bucket, object, 0, DiskError::FaultyDisk); + let mut proof = None; + let (_, error) = set + .heal_object_with_absence(bucket, object, &version.to_string(), &opts, &mut proof) + .await + .expect("heal should return a per-object failure"); + assert!( + error.is_some(), + "three absent slots meeting write quorum cannot hide the failed stale slot" + ); + assert!(proof.is_none(), "partial cleanup must not produce an absence proof"); + assert!( + disk.read_version("", bucket, object, &version.to_string(), &ReadOptions::default()) + .await + .is_ok(), + "the failed historical version must remain for retry" + ); + drop(failure); + let (result, error) = set + .heal_object_with_absence(bucket, object, &version.to_string(), &opts, &mut proof) + .await + .expect("retry should execute cleanup"); + assert!(error.is_none(), "retry should complete: {error:?}"); + let receipt = proof.take().expect("successful exact cleanup must produce proof"); + assert!(receipt.removed); + assert_eq!((receipt.pool_index, receipt.set_index), (set.pool_index, set.set_index)); + assert!( + result + .after + .drives + .iter() + .all(|drive| drive.state == DriveState::Missing.to_string()) + ); + let (_, _) = set + .heal_object_with_absence(bucket, object, &version.to_string(), &opts, &mut proof) + .await + .expect("already absent replay should execute"); + assert!(!proof.expect("exact already-absent replay must remain provable").removed); + }) + .await; + } + #[tokio::test] #[serial_test::serial] async fn heal_reports_success_after_dangling_inline_cleanup() { diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index 1b4866c74..818d70bd1 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -28,6 +28,27 @@ const EVENT_HEAL_ABANDONED_PARTS: &str = "heal_abandoned_parts"; const EVENT_HEAL_FORMAT_COMPLETED: &str = "heal_format_completed"; const EVENT_HEAL_OBJECT_STARTED: &str = "heal_object_started"; +/// Storage-owned proof for the exact version and every selected erasure location. +/// This is an in-process result, never reconstructed from admin drive telemetry. +#[derive(Debug)] +pub struct HealObjectAbsenceProof { + pub bucket: String, + pub object: String, + pub version_id: String, + pub bucket_incarnation_id: Uuid, + pub pool_index: Option, + pub set_index: Option, + pub locations: Vec<(usize, usize)>, + pub removed: bool, +} + +#[derive(Debug)] +pub struct HealObjectStorageResult { + pub item: HealResultItem, + pub error: Option, + pub absence: Option, +} + fn invalid_heal_pool_index(pool_idx: usize, pool_count: usize) -> Error { StorageError::InvalidArgument( "heal".to_string(), @@ -462,6 +483,70 @@ impl ECStore { object: &str, version_id: &str, opts: &HealOpts, + ) -> Result<(HealResultItem, Option)> { + self.handle_heal_object_with_absence(bucket, object, version_id, opts, &mut None) + .await + } + + pub async fn heal_object_with_proof( + &self, + bucket: &str, + object: &str, + version_id: &str, + opts: &HealOpts, + ) -> Result { + if opts.dry_run || opts.no_lock || version_id.is_empty() || super::utils::is_reserved_or_invalid_bucket(bucket, false) { + let (item, error) = self.handle_heal_object(bucket, object, version_id, opts).await?; + return Ok(HealObjectStorageResult { + item, + error, + absence: None, + }); + } + + // Match object publication: bucket lifecycle before capacity and object + // namespace locks. Keep the incarnation pinned through proof delivery. + let guard = self.acquire_bucket_lifecycle_read_lock(bucket).await?; + let mut proofs = None; + let (item, mut error) = self + .handle_heal_object_with_absence(bucket, object, version_id, opts, &mut proofs) + .await?; + // Read the authoritative incarnation only for an absence candidate. + // The lifecycle guard has pinned it throughout the storage operation. + let incarnation = if proofs.is_some() && !guard.is_lock_lost() { + self.bucket_incarnation_id_from_disk(bucket) + .await + .ok() + .filter(|id| !id.is_nil()) + } else { + None + }; + let absence = match (incarnation, proofs) { + (Some(incarnation), Some(proofs)) if !guard.is_lock_lost() => Some(HealObjectAbsenceProof { + bucket: bucket.to_owned(), + object: object.to_owned(), + version_id: version_id.to_owned(), + bucket_incarnation_id: incarnation, + pool_index: opts.pool, + set_index: opts.set, + removed: proofs.iter().any(|proof| proof.removed), + locations: proofs.into_iter().map(|proof| (proof.pool_index, proof.set_index)).collect(), + }), + _ => None, + }; + if absence.is_some() { + error = None; + } + Ok(HealObjectStorageResult { item, error, absence }) + } + + async fn handle_heal_object_with_absence( + &self, + bucket: &str, + object: &str, + version_id: &str, + opts: &HealOpts, + absence: &mut Option>, ) -> Result<(HealResultItem, Option)> { trace!( event = EVENT_HEAL_OBJECT_STARTED, @@ -476,7 +561,9 @@ impl ECStore { ); let object = encode_dir_object(object); + *absence = None; let pools = self.get_pools_for_heal_object(opts)?; + let requested_pool_count = pools.len(); if let Some(set_idx) = opts.set { for pool in &pools { if set_idx >= pool.disk_set.len() { @@ -550,7 +637,7 @@ impl ECStore { } #[cfg(test)] crate::core::pools::notify_decommission_external_heal_operation_started(store_id); - pool.heal_object(bucket, &pool_object, version_id, &opts).await + pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts).await } }); let results = join_all(futures).await; @@ -573,19 +660,28 @@ impl ECStore { move |opts| async move { #[cfg(test)] crate::core::pools::notify_decommission_external_heal_operation_started(store_id); - pool.heal_object(bucket, &pool_object, version_id, &opts).await + pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts).await }, )); } join_all(futures).await }; + let mut proofs = Vec::with_capacity(requested_pool_count); let mut errs = Vec::with_capacity(self.pools.len()); let mut ress = Vec::with_capacity(self.pools.len()); for res in results.into_iter() { match res { - Ok((result, err)) => { + Ok((result, err, proof)) => { + if let Some(proof) = proof + && (err.is_none() + || err + .as_ref() + .is_some_and(|err| is_err_object_not_found(err) || is_err_version_not_found(err))) + { + proofs.push(proof); + } let mut result = result; result.object = decode_dir_object(&result.object); ress.push(result); @@ -598,6 +694,12 @@ impl ECStore { } } + // Absence in one pool cannot discharge a responsibility covering other + // pools, including skipped decommission sources or failed lookups. + if requested_pool_count > 0 && proofs.len() == requested_pool_count { + *absence = Some(proofs); + } + for (idx, err) in errs.iter().enumerate() { if err.is_none() { return Ok((ress.remove(idx), None)); @@ -1090,6 +1192,40 @@ mod tests { (temp_dir, store, shutdown) } + #[tokio::test] + #[serial_test::serial] + async fn absence_proof_requires_every_selected_pool() { + let (_temp_dir, store, shutdown) = multi_pool_heal_store().await; + let bucket = format!("absence-scope-{}", Uuid::new_v4().simple()); + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create bucket in both pools"); + let version = Uuid::new_v4().to_string(); + let result = store + .heal_object_with_proof(&bucket, "history.txt", &version, &HealOpts::default()) + .await + .expect("exact absence lookup should complete across both pools"); + assert!(result.error.is_none()); + let proof = result.absence.expect("all selected pools proved the exact version absent"); + assert_eq!(proof.locations, vec![(0, 0), (1, 0)]); + assert_eq!((proof.pool_index, proof.set_index), (None, None)); + assert_eq!(proof.bucket, bucket); + assert_eq!(proof.version_id, version); + assert!(!proof.removed, "already absent versions do not count as another cleanup"); + + store.pool_meta.write().await.pools[1].decommission = Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }); + let partial = store + .heal_object_with_proof(&bucket, "history.txt", &version, &HealOpts::default()) + .await + .expect("unscoped heal should retain its legacy suspended-pool behavior"); + assert!(partial.absence.is_none(), "an uninspected suspended pool prevents global absence proof"); + shutdown.cancel(); + } + fn heal_test_format_path(temp_dir: &tempfile::TempDir, pool_index: usize, disk_index: usize) -> std::path::PathBuf { temp_dir .path() diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index ba1f70c5e..6dfc2b482 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -419,6 +419,7 @@ mod bucket_fence; pub(crate) use bucket::await_bucket_namespace_operation; pub use bucket_fence::BucketIncarnationFenceGuard; mod heal; +pub use heal::{HealObjectAbsenceProof, HealObjectStorageResult}; mod heal_walk; pub use heal_walk::HealWalkVersion; mod init; diff --git a/crates/heal/src/error.rs b/crates/heal/src/error.rs index 7fa8d7717..effd74d9a 100644 --- a/crates/heal/src/error.rs +++ b/crates/heal/src/error.rs @@ -145,6 +145,16 @@ impl Error { } } + pub(crate) fn dangling_delete_retry_not_before(&self) -> Option { + let after = match self { + Self::Storage(error) => error.dangling_delete_retry_after(), + Self::Disk(error) => error.dangling_delete_retry_after(), + Self::Io(error) => DiskError::io_error_dangling_delete_retry_after(error), + _ => None, + }?; + std::time::SystemTime::now().checked_add(after) + } + pub(crate) fn is_dangling_delete_grace(&self) -> bool { match self { Error::Storage(err) => err.is_dangling_delete_grace(), diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index 1dc6fdddc..559ed2e68 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -1127,8 +1127,49 @@ impl HealStorageAPI for ECStoreHealStorage { version_id: Option<&str>, opts: &HealOpts, ) -> Result { - let (item, error) = self.heal_object(bucket, object, version_id, opts).await?; - let receipt = if error.is_none() && !opts.dry_run { + let result = self + .ecstore + .heal_object_with_proof(bucket, object, version_id.unwrap_or(""), opts) + .await + .map_err(Error::Storage)?; + let item = result.item; + let error = result.error.map(Error::Storage); + let receipt = if let Some(proof) = result.absence { + if error.is_none() + && !opts.dry_run + && proof.bucket == bucket + && proof.object == object + && proof.version_id == version_id.unwrap_or("") + && proof.pool_index == opts.pool + && proof.set_index == opts.set + && !proof.bucket_incarnation_id.is_nil() + && !proof.locations.is_empty() + && proof.locations.iter().all(|(pool, set)| { + opts.pool.is_none_or(|expected| expected == *pool) && opts.set.is_none_or(|expected| expected == *set) + }) + { + Some(HealObjectReceipt { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: proof.bucket, + object: proof.object, + version_id: version_id.map(ToOwned::to_owned), + bucket_incarnation_id: Some(proof.bucket_incarnation_id), + pool_index: proof.pool_index, + set_index: proof.set_index, + }, + // A committed cleanup repaired the stale replica. A replay + // observing an already absent version made no new repair. + disposition: if proof.removed { + HealObjectDisposition::Repaired + } else { + HealObjectDisposition::AuthoritativelyAbsent + }, + }) + } else { + None + } + } else if error.is_none() && !opts.dry_run { let ok_drive_state = DriveState::Ok.to_string(); let all_after_drives_ok = item.after.drives.iter().all(|drive| drive.state == ok_drive_state); match ( diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index 1615b9077..76d477c45 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -632,7 +632,7 @@ impl HealTask { true } - async fn record_deferred_object(&self, reason: HealDeferredReason) { + async fn record_deferred_object(&self, reason: HealDeferredReason, retry_not_before: Option) { if let Some(identity) = self.single_object_identity() { let mut outcome = self.outcome.write().await; outcome.attempt_failed(); @@ -640,7 +640,7 @@ impl HealTask { identity, disposition: HealObjectDisposition::Deferred { reason, - retry_not_before: None, + retry_not_before, }, detail: None, }); @@ -782,7 +782,8 @@ impl HealTask { } async fn skip_due_to_transient_object_exists(&self, bucket: &str, object: &str, err: &Error) -> Result<()> { - self.record_deferred_object(HealDeferredReason::TransientExistenceCheck).await; + self.record_deferred_object(HealDeferredReason::TransientExistenceCheck, None) + .await; warn!( target: "rustfs::heal::task", event = EVENT_HEAL_OBJECT_RESULT, @@ -882,7 +883,8 @@ impl HealTask { return false; } - self.record_deferred_object(HealDeferredReason::TransientUsageCache).await; + self.record_deferred_object(HealDeferredReason::TransientUsageCache, None) + .await; warn!( target: "rustfs::heal::task", @@ -906,7 +908,8 @@ impl HealTask { return false; } - self.record_deferred_object(HealDeferredReason::DanglingDeleteGrace).await; + self.record_deferred_object(HealDeferredReason::DanglingDeleteGrace, err.dangling_delete_retry_not_before()) + .await; warn!( target: "rustfs::heal::task", diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index cec17886a..1ba504686 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -679,7 +679,7 @@ impl HealTask { if Self::is_dangling_delete_grace_error(&err) { disposition = HealObjectDisposition::Deferred { reason: HealDeferredReason::DanglingDeleteGrace, - retry_not_before: None, + retry_not_before: err.dangling_delete_retry_not_before(), }; telemetry_unknown |= !increment_counter(&mut skipped); warn!( diff --git a/crates/heal/tests/heal_b920_subquorum_union_test.rs b/crates/heal/tests/heal_b920_subquorum_union_test.rs index 6be2d0595..c3cfe3905 100644 --- a/crates/heal/tests/heal_b920_subquorum_union_test.rs +++ b/crates/heal/tests/heal_b920_subquorum_union_test.rs @@ -577,3 +577,234 @@ mod serial_tests { ); } } + +mod absence_receipt_regressions { + use super::*; + use rustfs_heal::heal::outcome::{HealDeferredReason, HealObjectDisposition}; + use rustfs_heal::heal::{HealOptions, HealPriority, HealRequest, HealTask, HealType}; + use storage_api::integration::{DiskAPI as _, DiskSetSelector, ObjectOperations as _, ReadOptions, StorageAdminApi as _}; + + const OBJECT: &str = "history.txt"; + const CURRENT: &[u8] = b"retained-current-version"; + + async fn stale_history(bucket: &str) -> (Vec, Arc, Arc, String, String) { + let (paths, store, storage) = heal_env_n(16).await; + create_versioned_bucket(&store, bucket).await; + let old = put_versioned(&store, bucket, OBJECT, b"stale-historical-version").await; + let current = put_versioned(&store, bucket, OBJECT, CURRENT).await; + let target = xl_meta_path(&object_dir(&paths[0], bucket, OBJECT)); + let stale = std::fs::read(&target).expect("capture both versions before the historical delete"); + store + .delete_object( + bucket, + OBJECT, + ObjectOptions { + version_id: Some(old.clone()), + versioned: true, + ..Default::default() + }, + ) + .await + .expect("delete the exact historical version on every disk"); + std::fs::write(&target, stale).expect("rejoin one disk retaining the deleted version"); + let inventory = store + .disk_set_inventory(DiskSetSelector::new(0, 0)) + .await + .expect("inspect actual disk inventory"); + for (index, disk) in inventory.iter().enumerate() { + let old_meta = disk + .as_ref() + .expect("fixture disks must be online") + .read_version("", bucket, OBJECT, &old, &ReadOptions::default()) + .await; + assert_eq!(old_meta.is_ok(), index == 0, "only the selected stale disk must retain the old version"); + } + (paths, store, storage, old, current) + } + + fn request(bucket: &str, old: &str) -> HealRequest { + HealRequest::new( + HealType::Object { + bucket: bucket.to_owned(), + object: OBJECT.to_owned(), + version_id: Some(old.to_owned()), + }, + HealOptions { + scan_mode: HealScanMode::Deep, + pool_index: Some(0), + set_index: Some(0), + ..Default::default() + }, + HealPriority::Normal, + ) + } + + async fn assert_versions(store: &Arc, bucket: &str, old: &str, current: &str) { + let inventory = store + .disk_set_inventory(DiskSetSelector::new(0, 0)) + .await + .expect("inspect post-heal disks"); + for disk in inventory.iter().flatten() { + assert!( + disk.read_version("", bucket, OBJECT, old, &ReadOptions::default()) + .await + .is_err_and(|error| matches!( + error, + storage_api::integration::DiskError::FileNotFound + | storage_api::integration::DiskError::FileVersionNotFound + )), + "the exact old version must be absent on each physical disk" + ); + let retained = disk + .read_version("", bucket, OBJECT, current, &ReadOptions::default()) + .await + .expect("cleanup must retain the current version on every physical disk"); + assert_eq!(retained.version_id.map(|id| id.to_string()).as_deref(), Some(current)); + assert_eq!( + (retained.erasure.data_blocks, retained.erasure.parity_blocks), + (12, 4), + "C06 fixture must exercise the production EC12+4 geometry" + ); + } + assert_eq!(read_version(store, bucket, OBJECT, current).await, CURRENT); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[serial] + async fn historical_absence_receipt_repairs_and_replays() { + let bucket = "absence-receipt-replay"; + let (_paths, store, storage, old, current) = stale_history(bucket).await; + let opts = HealOpts { + pool: Some(0), + set: Some(0), + ..deep_heal_opts() + }; + let result = with_dangling_grace_disabled(storage.heal_object_with_receipt(bucket, OBJECT, Some(&old), &opts)) + .await + .expect("historical cleanup should complete"); + assert!(result.error.is_none(), "cleanup failed: {:?}", result.error); + let receipt = result + .receipt + .expect("committed cleanup must produce a receipt without all drives being OK"); + assert_eq!(receipt.disposition, HealObjectDisposition::Repaired); + assert_eq!(receipt.identity.bucket, bucket); + assert_eq!(receipt.identity.object, OBJECT); + assert_eq!(receipt.identity.version_id.as_deref(), Some(old.as_str())); + assert_eq!((receipt.identity.pool_index, receipt.identity.set_index), (Some(0), Some(0))); + assert_eq!( + receipt.identity.bucket_incarnation_id, + Some(store.bucket_incarnation_id_from_disk(bucket).await.expect("bucket identity")) + ); + assert!(result.item.after.drives.iter().all(|drive| drive.state == "missing")); + assert_versions(&store, bucket, &old, ¤t).await; + + // Reconstruct the task and storage facade as after a lost response or + // task restart. The on-disk absence supplies a new exact no-op proof. + let restarted = Arc::new(ECStoreHealStorage::new(store.clone())); + let task = HealTask::from_request(request(bucket, &old), restarted); + task.execute().await.expect("replayed exact-version heal should complete"); + let outcome = task.get_outcome().await; + assert_eq!(outcome.counters.processed, 1); + assert_eq!(outcome.counters.unchanged, 1); + assert_eq!(outcome.counters.healed, 0); + assert_eq!(outcome.counters.unknown, 0); + assert_eq!(outcome.counters.failed, 0); + assert_eq!(outcome.objects.len(), 1); + assert_eq!(outcome.objects[0].disposition, HealObjectDisposition::AuthoritativelyAbsent); + assert_versions(&store, bucket, &old, ¤t).await; + } + + #[test] + #[serial] + fn historical_absence_receipt_bucket_outcome_matches_c06() { + // The real ECStore initialization and bucket traversal need the debug + // server's stack budget, which exceeds libtest's default on Linux. + const STACK_SIZE: usize = 8 * 1024 * 1024; + std::thread::Builder::new() + .name("absence-receipt-c06".to_owned()) + .stack_size(STACK_SIZE) + .spawn(|| { + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(4) + .thread_stack_size(STACK_SIZE) + .enable_all() + .build() + .expect("C06 test runtime should build"); + + runtime.block_on(historical_absence_receipt_bucket_outcome_matches_c06_inner()); + }) + .expect("C06 test thread should spawn") + .join() + .expect("C06 test thread should finish"); + } + + async fn historical_absence_receipt_bucket_outcome_matches_c06_inner() { + let bucket = "absence-receipt-c06"; + let (_paths, store, storage, old, current) = stale_history(bucket).await; + put_versioned(&store, bucket, "healthy.txt", b"already healthy").await; + let task = HealTask::from_request( + HealRequest::new( + HealType::Bucket { + bucket: bucket.to_owned(), + }, + HealOptions { + recursive: true, + scan_mode: HealScanMode::Deep, + ..Default::default() + }, + HealPriority::Normal, + ), + storage, + ); + with_dangling_grace_disabled(task.execute()) + .await + .expect("C06 bucket traversal should complete"); + let outcome = task.get_outcome().await; + assert_eq!(outcome.counters.processed, 3); + assert_eq!(outcome.counters.healed, 1, "completed cleanup must be repaired: {outcome:?}"); + assert_eq!(outcome.counters.unchanged, 2); + assert_eq!(outcome.counters.unknown, 0); + assert_eq!(outcome.counters.failed, 0); + assert_eq!(outcome.counters.skipped, 0); + assert_versions(&store, bucket, &old, ¤t).await; + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[serial] + async fn historical_absence_receipt_grace_and_dry_run_preserve_version() { + let bucket = "absence-receipt-grace"; + let (paths, _store, storage, old, _current) = stale_history(bucket).await; + let target = xl_meta_path(&object_dir(&paths[0], bucket, OBJECT)); + let before = std::fs::read(&target).expect("read pre-heal stale metadata"); + let result = with_dangling_grace_disabled(storage.heal_object_with_receipt( + bucket, + OBJECT, + Some(&old), + &HealOpts { + dry_run: true, + ..deep_heal_opts() + }, + )) + .await + .expect("dry run should return a result"); + assert!(result.receipt.is_none(), "dry run cannot issue a cleanup receipt"); + assert_eq!(std::fs::read(&target).expect("read dry-run metadata"), before); + + let task = HealTask::from_request(request(bucket, &old), storage); + temp_env::async_with_vars([(GRACE_ENV, Some("3600"))], task.execute()) + .await + .expect("grace should defer without failing execution"); + let outcome = task.get_outcome().await; + assert_eq!(outcome.counters.processed, 1); + assert_eq!(outcome.counters.healed, 0); + assert_eq!(outcome.counters.unknown, 0); + assert!( + matches!(outcome.objects[0].disposition, HealObjectDisposition::Deferred { + reason: HealDeferredReason::DanglingDeleteGrace, retry_not_before: Some(due), + } if due > std::time::SystemTime::now()), + "grace must expose its retry deadline: {:?}", + outcome.objects + ); + assert_eq!(std::fs::read(&target).expect("read grace-protected metadata"), before); + } +} diff --git a/crates/heal/tests/storage_api.rs b/crates/heal/tests/storage_api.rs index 9c0fdf956..ae9f3449c 100644 --- a/crates/heal/tests/storage_api.rs +++ b/crates/heal/tests/storage_api.rs @@ -30,4 +30,5 @@ pub(crate) mod integration { pub(crate) use rustfs_storage_api::NamespaceLocking; pub(crate) use rustfs_storage_api::ObjectIO; pub(crate) use rustfs_storage_api::ObjectOperations; + pub(crate) use rustfs_storage_api::{DiskSetSelector, StorageAdminApi}; }