diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 01d2e0864..c4b09de6a 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -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, Option)> { 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. diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index f75a0a177..623485066 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -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; diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 921315701..82cc66610 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -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)> { - 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, + proof: AbsenceProofRequest<'_>, ) -> Result<(HealResultItem, Option)> { - 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, - retirement: Option<&MarkerRetirementContext<'_>>, + proof: AbsenceProofRequest<'_>, ) -> Result<(HealResultItem, Option)> { *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); diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index b3708cda5..0e7f86168 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -296,6 +296,24 @@ impl ECStore { .await } + pub async fn heal_mrf_object_at_incarnation( + self: &Arc, + bucket: &str, + object: &str, + version_id: &str, + expected: uuid::Uuid, + opts: &HealOpts, + ) -> Result { + 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)> { - 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 { - 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 { + 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 { + 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>, - retirement: Option<&crate::bucket::retirement::MarkerRetirementContext<'_>>, + proof: crate::set_disk::AbsenceProofRequest<'_>, ) -> Result<(HealResultItem, Option)> { 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 }, )); diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index 956130e54..a82a6452c 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -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 { + 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 { + 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> { debug!( target: "rustfs::heal::storage", diff --git a/crates/heal/src/heal/task/heal_object.rs b/crates/heal/src/heal/task/heal_object.rs index 50e5e1fc8..dbde2fad4 100644 --- a/crates/heal/src/heal/task/heal_object.rs +++ b/crates/heal/src/heal/task/heal_object.rs @@ -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(()) + } } diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index cfb9deb2d..7104ae6fe 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -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 { diff --git a/crates/heal/tests/heal_b920_subquorum_union_test.rs b/crates/heal/tests/heal_b920_subquorum_union_test.rs index f7546bfb7..76bafa06e 100644 --- a/crates/heal/tests/heal_b920_subquorum_union_test.rs +++ b/crates/heal/tests/heal_b920_subquorum_union_test.rs @@ -717,6 +717,49 @@ mod absence_receipt_regressions { assert_versions(&store, bucket, &old, ¤t).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() { diff --git a/crates/heal/tests/mrf_partial_write_test.rs b/crates/heal/tests/mrf_partial_write_test.rs index 26035987d..c8340e786 100644 --- a/crates/heal/tests/mrf_partial_write_test.rs +++ b/crates/heal/tests/mrf_partial_write_test.rs @@ -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 Fut, Fut: Future>(mut probe: F) -> bool { tokio::time::timeout(Duration::from_secs(60), async { loop {