fix(heal): discharge absent durable partial writes (#8000)

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Hauser
2026-09-18 16:07:18 +08:00
committed by GitHub
parent 0f69c073de
commit 023adb1696
9 changed files with 410 additions and 20 deletions
+2 -2
View File
@@ -1493,12 +1493,12 @@ impl Sets {
object: &str,
version_id: &str,
opts: &HealOpts,
retirement: Option<&crate::bucket::retirement::MarkerRetirementContext<'_>>,
proof: crate::set_disk::AbsenceProofRequest<'_>,
) -> 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_retirement(bucket, object, version_id, opts, &mut absence, retirement)
.heal_object_with_retirement(bucket, object, version_id, opts, &mut absence, proof)
.await?;
// A caller-owned lock does not expose its lease to this boundary.
// Keep cleanup unverified when that lease cannot be checked here.
+1
View File
@@ -874,6 +874,7 @@ mod ctx;
mod metadata;
mod ops;
pub(crate) use ops::bucket::BucketInfoQuorum;
pub(crate) use ops::heal::AbsenceProofRequest;
#[cfg(test)]
pub(crate) use ops::heal::DanglingDeleteFailure;
pub(crate) use ops::heal::HealedObjectAbsence;
+59 -11
View File
@@ -55,6 +55,13 @@ pub(crate) struct HealedObjectAbsence {
pub removed: bool,
}
/// Conditions under which a missing object may become a storage-owned proof.
#[derive(Clone, Copy)]
pub(crate) struct AbsenceProofRequest<'a> {
pub retirement: Option<&'a MarkerRetirementContext<'a>>,
pub allow_unversioned: bool,
}
fn is_unavailable_heal_rpc(error: &DiskError) -> bool {
let DiskError::Io(error) = error else {
return false;
@@ -3012,8 +3019,18 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks {
version_id: &str,
opts: &HealOpts,
) -> Result<(HealResultItem, Option<Error>)> {
self.heal_object_with_absence(bucket, object, version_id, opts, &mut None)
.await
self.heal_object_with_absence(
bucket,
object,
version_id,
opts,
&mut None,
AbsenceProofRequest {
retirement: None,
allow_unversioned: false,
},
)
.await
}
#[tracing::instrument(skip(self))]
@@ -3098,8 +3115,9 @@ impl SetDisks {
version_id: &str,
opts: &HealOpts,
absence: &mut Option<HealedObjectAbsence>,
proof: AbsenceProofRequest<'_>,
) -> Result<(HealResultItem, Option<Error>)> {
self.heal_object_with_retirement(bucket, object, version_id, opts, absence, None)
self.heal_object_with_retirement(bucket, object, version_id, opts, absence, proof)
.await
}
@@ -3110,7 +3128,7 @@ impl SetDisks {
version_id: &str,
opts: &HealOpts,
absence: &mut Option<HealedObjectAbsence>,
retirement: Option<&MarkerRetirementContext<'_>>,
proof: AbsenceProofRequest<'_>,
) -> Result<(HealResultItem, Option<Error>)> {
*absence = None;
let _write_lock_guard = if !opts.no_lock {
@@ -3176,7 +3194,7 @@ impl SetDisks {
.await;
// Check the lease after the final await before publishing the proof.
if !opts.dry_run
&& !version_id.is_empty()
&& (proof.allow_unversioned || !version_id.is_empty())
&& !disks.is_empty()
&& disks.iter().all(Option::is_some)
&& errs
@@ -3205,7 +3223,7 @@ impl SetDisks {
ExplicitVersionHeal {
opts: &inner_opts,
allow_regeneration: true,
retirement,
retirement: proof.retirement,
},
absence,
)
@@ -3225,7 +3243,7 @@ impl SetDisks {
ExplicitVersionHeal {
opts: &inner_opts,
allow_regeneration: true,
retirement,
retirement: proof.retirement,
},
absence,
)
@@ -3277,7 +3295,7 @@ impl SetDisks {
#[cfg(test)]
mod heal_result_report_tests {
use super::{
DanglingCheckPartsFailure, DanglingDeleteFailure, DanglingDeleteSafety, ReadRepairCommitFingerprint,
AbsenceProofRequest, DanglingCheckPartsFailure, DanglingDeleteFailure, DanglingDeleteSafety, ReadRepairCommitFingerprint,
ReadRepairPauseScope, SetDisks, heal_writer_error_summary,
};
use super::{HEAL_RENAME_INCOMPLETE, HealRenameFailureScope, HealWriterFailureScope};
@@ -5385,7 +5403,17 @@ mod heal_result_report_tests {
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)
.heal_object_with_absence(
bucket,
object,
&version.to_string(),
&opts,
&mut proof,
AbsenceProofRequest {
retirement: None,
allow_unversioned: false,
},
)
.await
.expect("heal should return a per-object failure");
assert!(
@@ -5401,7 +5429,17 @@ mod heal_result_report_tests {
);
drop(failure);
let (result, error) = set
.heal_object_with_absence(bucket, object, &version.to_string(), &opts, &mut proof)
.heal_object_with_absence(
bucket,
object,
&version.to_string(),
&opts,
&mut proof,
AbsenceProofRequest {
retirement: None,
allow_unversioned: false,
},
)
.await
.expect("retry should execute cleanup");
assert!(error.is_none(), "retry should complete: {error:?}");
@@ -5416,7 +5454,17 @@ mod heal_result_report_tests {
.all(|drive| drive.state == DriveState::Missing.to_string())
);
let (_, _) = set
.heal_object_with_absence(bucket, object, &version.to_string(), &opts, &mut proof)
.heal_object_with_absence(
bucket,
object,
&version.to_string(),
&opts,
&mut proof,
AbsenceProofRequest {
retirement: None,
allow_unversioned: false,
},
)
.await
.expect("already absent replay should execute");
assert!(!proof.expect("exact already-absent replay must remain provable").removed);
+72 -7
View File
@@ -296,6 +296,24 @@ impl ECStore {
.await
}
pub async fn heal_mrf_object_at_incarnation(
self: &Arc<Self>,
bucket: &str,
object: &str,
version_id: &str,
expected: uuid::Uuid,
opts: &HealOpts,
) -> Result<HealObjectStorageResult> {
let object = object.to_owned();
let version_id = version_id.to_owned();
self.run_bucket_heal_at_incarnation(bucket, expected, opts, |store, bucket, opts| async move {
store
.heal_object_with_authoritative_absence_proof(&bucket, &object, &version_id, &opts)
.await
})
.await
}
async fn acquire_heal_format_fence(
&self,
) -> Result<(
@@ -684,8 +702,18 @@ impl ECStore {
version_id: &str,
opts: &HealOpts,
) -> Result<(HealResultItem, Option<Error>)> {
self.handle_heal_object_with_absence(bucket, object, version_id, opts, &mut None, None)
.await
self.handle_heal_object_with_absence(
bucket,
object,
version_id,
opts,
&mut None,
crate::set_disk::AbsenceProofRequest {
retirement: None,
allow_unversioned: false,
},
)
.await
}
pub async fn heal_object_with_proof(
@@ -695,7 +723,34 @@ impl ECStore {
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) {
self.heal_object_with_absence_proof(bucket, object, version_id, opts, false)
.await
}
async fn heal_object_with_authoritative_absence_proof(
&self,
bucket: &str,
object: &str,
version_id: &str,
opts: &HealOpts,
) -> Result<HealObjectStorageResult> {
self.heal_object_with_absence_proof(bucket, object, version_id, opts, true)
.await
}
async fn heal_object_with_absence_proof(
&self,
bucket: &str,
object: &str,
version_id: &str,
opts: &HealOpts,
allow_unversioned_absence: bool,
) -> Result<HealObjectStorageResult> {
if opts.dry_run
|| opts.no_lock
|| (!allow_unversioned_absence && 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,
@@ -738,7 +793,17 @@ impl ECStore {
};
let mut proofs = None;
let (item, mut error) = self
.handle_heal_object_with_absence(bucket, object, version_id, opts, &mut proofs, Some(&retirement))
.handle_heal_object_with_absence(
bucket,
object,
version_id,
opts,
&mut proofs,
crate::set_disk::AbsenceProofRequest {
retirement: Some(&retirement),
allow_unversioned: allow_unversioned_absence,
},
)
.await?;
// Read the authoritative incarnation only for an absence candidate.
// The lifecycle guard has pinned it throughout the storage operation.
@@ -780,7 +845,7 @@ impl ECStore {
version_id: &str,
opts: &HealOpts,
absence: &mut Option<Vec<crate::set_disk::HealedObjectAbsence>>,
retirement: Option<&crate::bucket::retirement::MarkerRetirementContext<'_>>,
proof: crate::set_disk::AbsenceProofRequest<'_>,
) -> Result<(HealResultItem, Option<Error>)> {
trace!(
event = EVENT_HEAL_OBJECT_STARTED,
@@ -871,7 +936,7 @@ impl ECStore {
}
#[cfg(test)]
crate::core::pools::notify_decommission_external_heal_operation_started(store_id);
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts, retirement)
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts, proof)
.await
}
});
@@ -895,7 +960,7 @@ impl ECStore {
move |opts| async move {
#[cfg(test)]
crate::core::pools::notify_decommission_external_heal_operation_started(store_id);
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts, retirement)
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts, proof)
.await
},
));
+31
View File
@@ -532,6 +532,19 @@ pub trait HealStorageAPI: Send + Sync {
Err(Error::other("storage does not support incarnation-bound object healing"))
}
/// Durable MRF repair may request an authoritative unversioned absence proof.
async fn heal_mrf_object_at_incarnation(
&self,
bucket: &str,
object: &str,
version_id: Option<&str>,
expected: Uuid,
opts: &HealOpts,
) -> Result<HealStorageObjectResult> {
self.heal_object_at_incarnation(bucket, object, version_id, expected, opts)
.await
}
/// Heal object using ecstore
async fn heal_object(
&self,
@@ -967,6 +980,24 @@ impl HealStorageAPI for ECStoreHealStorage {
.await)
}
async fn heal_mrf_object_at_incarnation(
&self,
bucket: &str,
object: &str,
version_id: Option<&str>,
expected: Uuid,
opts: &HealOpts,
) -> Result<HealStorageObjectResult> {
let result = self
.ecstore
.heal_mrf_object_at_incarnation(bucket, object, version_id.unwrap_or_default(), expected, opts)
.await
.map_err(|error| incarnation_storage_error(bucket, expected, error))?;
Ok(self
.object_result_with_receipt(bucket, object, version_id, opts, result, Some(expected))
.await)
}
async fn get_object_meta(&self, bucket: &str, object: &str) -> Result<Option<HealObjectInfo>> {
debug!(
target: "rustfs::heal::storage",
+65
View File
@@ -429,6 +429,10 @@ impl HealTask {
/// Recreate missing object (for EC decode scenarios)
async fn recreate_missing_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
if self.source == HealRequestSource::Mrf {
return self.recreate_missing_mrf_object(bucket, object, version_id).await;
}
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_OBJECT_STAGE,
@@ -529,4 +533,65 @@ impl HealTask {
}
}
}
/// Durable MRF responsibilities may complete only with an exact storage proof.
async fn recreate_missing_mrf_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
let heal_opts = HealOpts {
recursive: false,
dry_run: self.options.dry_run,
remove: false,
recreate: true,
scan_mode: HealScanMode::Deep,
update_parity: true,
no_lock: self.options.no_lock,
read_repair: false,
pool: self.options.pool_index,
set: self.options.set_index,
};
let mut expected = self.outcome_identity(bucket, object, version_id, self.options.pool_index, self.options.set_index);
let bucket_incarnation_id = self
.outcome_bucket_incarnation_id(bucket, self.options.dry_run)
.await?
.ok_or_else(|| Error::TaskExecutionFailed {
message: format!("Missing bucket incarnation for durable MRF repair {bucket}/{object}"),
})?;
expected.bucket_incarnation_id = Some(bucket_incarnation_id);
let storage_result = self
.await_with_control(self.storage.heal_mrf_object_at_incarnation(
bucket,
object,
version_id,
bucket_incarnation_id,
&heal_opts,
))
.await?;
if let Some(error) = storage_result.error {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to recreate missing object {bucket}/{object}: {error}"),
});
}
let object_size = storage_result.item.object_size as u64;
let authoritatively_absent = matches!(
storage_result.receipt.as_ref().map(|receipt| &receipt.disposition),
Some(HealObjectDisposition::AuthoritativelyAbsent)
);
if !self.record_verified_storage_receipt(expected, storage_result.receipt).await {
return Err(Error::TaskExecutionFailed {
message: format!("Missing exact storage proof for durable MRF repair {bucket}/{object}"),
});
}
{
let mut progress = self.progress.write().await;
if authoritatively_absent {
progress.update_object_progress(1, 0, 0, 1, 0);
} else {
progress.update_object_progress(1, 1, 0, 0, object_size);
}
}
self.record_result_item(storage_result.item).await;
Ok(())
}
}
+61
View File
@@ -3889,6 +3889,67 @@ async fn test_heal_recreate_scanner_non_dir_not_found_fails() {
assert!(matches!(task.get_status().await, HealTaskStatus::Failed { .. }));
}
#[tokio::test]
async fn mrf_recreate_missing_object_records_exact_absence_receipt_with_scope() {
let incarnation = Uuid::new_v4();
let storage = Arc::new(MockStorage {
object_exists: Mutex::new(Some(false)),
bucket_incarnation_id: Mutex::new(Some(incarnation)),
heal_object_receipts: Mutex::new(HashMap::from([(
"deleted.bin".to_string(),
VecDeque::from([HealObjectReceipt {
identity: HealObjectIdentity {
kind: HealObjectKind::Object,
bucket: "bucket-a".to_string(),
object: "deleted.bin".to_string(),
version_id: None,
bucket_incarnation_id: Some(incarnation),
pool_index: Some(2),
set_index: Some(3),
},
disposition: HealObjectDisposition::AuthoritativelyAbsent,
}]),
)])),
..Default::default()
});
let mut request = HealRequest::new(
HealType::Object {
bucket: "bucket-a".to_string(),
object: "deleted.bin".to_string(),
version_id: None,
},
HealOptions {
recreate_missing: true,
pool_index: Some(2),
set_index: Some(3),
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Mrf;
let task = HealTask::from_request(request, storage.clone());
task.execute()
.await
.expect("a complete MRF absence receipt must complete the task");
let opts = storage.object_heal_opts.lock().unwrap()[0];
assert_eq!(opts.pool, Some(2));
assert_eq!(opts.set, Some(3));
assert!(opts.recreate);
assert_eq!(opts.scan_mode, HealScanMode::Deep);
let outcome = task.get_outcome().await;
assert_eq!(outcome.counters.unchanged, 1);
assert_eq!(outcome.counters.unknown, 0);
assert_eq!(outcome.objects.len(), 1);
assert_eq!(outcome.objects[0].disposition, HealObjectDisposition::AuthoritativelyAbsent);
assert_eq!(
(outcome.objects[0].identity.pool_index, outcome.objects[0].identity.set_index),
(Some(2), Some(3))
);
}
#[tokio::test]
async fn test_heal_scanner_missing_object_without_recreate_probes_storage() {
let storage = Arc::new(MockStorage {
@@ -717,6 +717,49 @@ mod absence_receipt_regressions {
assert_versions(&store, bucket, &old, &current).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn unversioned_missing_object_produces_fenced_absence_receipt() {
let bucket = "absence-receipt-unversioned";
let object = "deleted.bin";
let (_paths, store, storage) = heal_env_n(16).await;
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("create unversioned bucket");
let incarnation = storage
.admit_bucket_incarnation(bucket)
.await
.expect("admit bucket incarnation for the durable proof");
let result = storage
.heal_mrf_object_at_incarnation(
bucket,
object,
None,
incarnation,
&HealOpts {
pool: Some(0),
set: Some(0),
..deep_heal_opts()
},
)
.await
.expect("fenced absence probe should return a storage result");
assert!(
result.error.is_none(),
"complete unversioned absence must not remain an error: {:?}",
result.error
);
let receipt = result.receipt.expect("complete unversioned absence needs a receipt");
assert_eq!(receipt.disposition, HealObjectDisposition::AuthoritativelyAbsent);
assert_eq!(receipt.identity.bucket, bucket);
assert_eq!(receipt.identity.object, object);
assert!(receipt.identity.version_id.is_none());
assert_eq!((receipt.identity.pool_index, receipt.identity.set_index), (Some(0), Some(0)));
assert_eq!(receipt.identity.bucket_incarnation_id, Some(incarnation));
}
#[test]
#[serial]
fn historical_absence_receipt_bucket_outcome_matches_c06() {
@@ -90,6 +90,82 @@ async fn partial_write_persistence_failure_is_reported_and_retained_for_retry()
);
}
#[test]
fn unversioned_deleted_partial_write_is_discharged_by_an_absence_proof() {
const STACK_SIZE: usize = 8 * 1024 * 1024;
std::thread::Builder::new()
.name("mrf-partial-write-absence".to_owned())
.stack_size(STACK_SIZE)
.spawn(|| {
let runtime = tokio::runtime::Builder::new_current_thread()
.thread_stack_size(STACK_SIZE)
.enable_all()
.build()
.expect("partial-write absence runtime should build");
runtime.block_on(unversioned_deleted_partial_write_is_discharged_by_an_absence_proof_inner());
})
.expect("partial-write absence test thread should spawn")
.join()
.expect("partial-write absence test thread should finish");
}
async fn unversioned_deleted_partial_write_is_discharged_by_an_absence_proof_inner() {
use rustfs_common::mrf_channel::{MrfScope, persist_partial_write_intent};
temp_env::async_with_vars([("RUSTFS_HEAL_MRF_ENABLE", Some("true"))], async {
let root = tempfile::tempdir().expect("partial-write absence fixture directory");
let env = TestECStoreEnv::builder().disk_count(16).base_dir(root.path()).build().await;
env.make_bucket("partial-absence", false).await;
let mut coordinator_pool = env.endpoint_pools.as_ref()[0].clone();
let mut endpoints = coordinator_pool.endpoints.as_ref().to_vec();
for endpoint in endpoints.iter_mut().skip(4) {
endpoint.is_local = false;
}
coordinator_pool.endpoints = Endpoints::from(endpoints);
init_local_disks(EndpointServerPools::from(vec![coordinator_pool]))
.await
.expect("coordinator journal disks");
let manager = manager(&env);
mrf_queue::spawn_mrf_consumer(manager.clone());
for (object, version_id) in [
("deleted-unversioned.bin", None),
("deleted-versioned.bin", Some(uuid::Uuid::new_v4())),
] {
persist_partial_write_intent(
"partial-absence",
object,
version_id,
MrfScope {
pool_index: 0,
set_index: 0,
},
)
.await
.expect("durable partial-write responsibility must commit before scheduling");
assert!(snapshot_contains(object).await, "committed responsibility must exist before repair runs");
}
manager.start().await.expect("MRF scheduler should start");
for object in ["deleted-unversioned.bin", "deleted-versioned.bin"] {
assert!(
wait_until(|| async { !snapshot_contains(object).await }).await,
"complete absence proof must discharge the durable responsibility"
);
}
assert!(
wait_until(|| async {
let snapshot = manager.operations_snapshot().await;
snapshot.queue_length == 0 && snapshot.active_tasks == 0
})
.await,
"discharged absence repair must leave no queued work"
);
manager.stop().await.expect("absence manager should stop");
})
.await;
}
async fn wait_until<F: FnMut() -> Fut, Fut: Future<Output = bool>>(mut probe: F) -> bool {
tokio::time::timeout(Duration::from_secs(60), async {
loop {