fix(heal): prove completed historical version cleanup (#7736)

* fix(heal): prove completed historical version cleanup

* test(heal): use debug runtime stack for C06 regression

* test(ecstore): acknowledge partial PUT heal admission concurrently

* test(mrf): fix journal fault and ownership fixtures
This commit is contained in:
cxymds
2026-09-13 19:59:47 +08:00
committed by GitHub
parent f21b06dfd2
commit 39ccd3abb0
15 changed files with 801 additions and 114 deletions
+1
View File
@@ -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 {
+22
View File
@@ -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<Error>, Option<crate::set_disk::HealedObjectAbsence>)> {
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::*;
+35 -2
View File
@@ -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<io::Error> {
let grace = error.get_ref()?.downcast_ref::<DanglingDeleteGraceError>()?;
Some(io::Error::new(error.kind(), grace.clone()))
}
pub fn dangling_delete_retry_after(&self) -> Option<std::time::Duration> {
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<std::time::Duration> {
let grace = error.get_ref()?.downcast_ref::<DanglingDeleteGraceError>()?;
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(
+10
View File
@@ -396,6 +396,13 @@ impl StorageError {
)
}
pub fn dangling_delete_retry_after(&self) -> Option<std::time::Duration> {
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(),
@@ -5979,6 +5979,20 @@ impl SetDisks {
data_errs_by_part: &HashMap<usize, Vec<usize>>,
opts: ObjectOptions,
) -> disk::error::Result<FileInfo> {
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<DiskError>],
data_errs_by_part: &HashMap<usize, Vec<usize>>,
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<disk::error::Result<()>>, write_quorum: usize) -> disk::error::Result<()> {
+1
View File
@@ -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;
+251 -99
View File
@@ -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<DiskError> {
pub(in crate::set_disk) fn injected_dangling_delete_error(bucket: &str, object: &str, disk_index: usize) -> Option<DiskError> {
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<DiskError>)> {
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<HealedObjectAbsence>,
) -> disk::error::Result<(HealResultItem, Option<DiskError>)> {
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<Error>)> {
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<HealedObjectAbsence>,
) -> Result<(HealResultItem, Option<Error>)> {
*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<DiskError>],
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() {
+139 -3
View File
@@ -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<usize>,
pub set_index: Option<usize>,
pub locations: Vec<(usize, usize)>,
pub removed: bool,
}
#[derive(Debug)]
pub struct HealObjectStorageResult {
pub item: HealResultItem,
pub error: Option<Error>,
pub absence: Option<HealObjectAbsenceProof>,
}
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<Error>)> {
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<HealObjectStorageResult> {
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<Vec<crate::set_disk::HealedObjectAbsence>>,
) -> Result<(HealResultItem, Option<Error>)> {
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()
+1
View File
@@ -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;
+10
View File
@@ -145,6 +145,16 @@ impl Error {
}
}
pub(crate) fn dangling_delete_retry_not_before(&self) -> Option<std::time::SystemTime> {
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(),
+43 -2
View File
@@ -1127,8 +1127,49 @@ impl HealStorageAPI for ECStoreHealStorage {
version_id: Option<&str>,
opts: &HealOpts,
) -> Result<HealStorageObjectResult> {
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 (
+8 -5
View File
@@ -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<SystemTime>) {
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",
+1 -1
View File
@@ -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!(
@@ -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<PathBuf>, Arc<ECStore>, Arc<ECStoreHealStorage>, 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<ECStore>, 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, &current).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, &current).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, &current).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);
}
}
+1
View File
@@ -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};
}