From 41770983d6ee51dae5b7c7276f31a400a0a18d7f Mon Sep 17 00:00:00 2001 From: cxymds Date: Sun, 13 Sep 2026 21:43:43 +0800 Subject: [PATCH] fix(heal): preserve historical null version identity (#7749) --- Cargo.lock | 1 + crates/ecstore/src/set_disk/ops/heal.rs | 1 + crates/ecstore/src/set_disk/ops/heal_walk.rs | 61 +++- crates/heal/Cargo.toml | 1 + crates/heal/src/heal/resume.rs | 25 +- crates/heal/src/heal/resume/checkpoint.rs | 8 +- crates/heal/src/heal/resume/tests.rs | 146 ++++++++ crates/heal/src/heal/storage.rs | 149 ++++++-- .../heal_b5_versioned_regression_test.rs | 322 +++++++++++++++++- crates/madmin/src/heal_commands.rs | 18 + 10 files changed, 662 insertions(+), 70 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index e36abcc99..b090c62c6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9935,6 +9935,7 @@ dependencies = [ "rustfs-concurrency", "rustfs-config", "rustfs-ecstore", + "rustfs-filemeta", "rustfs-heal-contracts", "rustfs-lock", "rustfs-madmin", diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 12e347d91..12605ec7f 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -905,6 +905,7 @@ impl SetDisks { result.object_size = ObjectInfo::from_file_info(&latest_meta, bucket, object, true).get_actual_size()? as usize; + result.resolved_version_id = Some(*latest_meta.version_id.unwrap_or_default().as_bytes()); // Loop to find number of disks with valid data, per-drive // data state and a list of outdated disks on which data needs // to be healed. diff --git a/crates/ecstore/src/set_disk/ops/heal_walk.rs b/crates/ecstore/src/set_disk/ops/heal_walk.rs index aa627e415..2dbf2fdd5 100644 --- a/crates/ecstore/src/set_disk/ops/heal_walk.rs +++ b/crates/ecstore/src/set_disk/ops/heal_walk.rs @@ -45,12 +45,12 @@ const BACKGROUND_WALKDIR_STALL_TIMEOUT: Duration = Duration::from_secs(60); /// `is_delete_marker` is OBSERVABILITY-ONLY (metrics / logging / e2e assertions); /// it must not gate healing logic — the delete-marker vs data path is chosen /// inside `ops/heal.rs` from the resolved latest metadata. `version_id` is -/// normalized (nil/absent UUID => `None`). +/// exact: an enumerated nil/absent UUID selects the null slot, never latest. #[derive(Debug, Clone)] pub struct HealWalkVersion { /// object key pub name: String, - /// normalized version id (`None` when the version is nil/absent) + /// Exact version id, including the nil UUID for the null slot. pub version_id: Option, /// version modification time as Unix nanoseconds pub mod_time_unix_nanos: Option, @@ -136,8 +136,7 @@ impl HealWalkCollector { }; versions.push(HealWalkVersion { name: entry.name.clone(), - // Normalize: nil/absent version id => None. - version_id: version_uuid.map(|u| u.to_string()), + version_id: Some(fi.version_id.unwrap_or_default().to_string()), mod_time_unix_nanos: fi.mod_time.map(|mod_time| mod_time.unix_timestamp_nanos()), lifecycle_object_info, is_delete_marker: fi.deleted, @@ -194,7 +193,7 @@ impl HealWalkCollector { }; for fi in fiv.versions.iter().chain(fiv.free_versions.iter()) { let version_uuid = fi.version_id.filter(|version_id| !version_id.is_nil()); - let vid = version_uuid.map(|u| u.to_string()); + let vid = Some(fi.version_id.unwrap_or_default().to_string()); if seen.insert(vid.clone()) { let lifecycle_object_info = if self.include_lifecycle_object_info { Some(ObjectInfo::from_file_info_with_version_id(fi, &self.bucket, &entry.name, version_uuid)) @@ -386,6 +385,58 @@ mod tests { }) } + #[test] + fn collectors_preserve_historical_null_identity() { + use rustfs_filemeta::FileInfo; + let latest = Uuid::from_u128(7); + for null_marker in [false, true] { + let mut metadata = FileMeta::new(); + for (version, seconds, deleted) in [(Uuid::nil(), 1, null_marker), (latest, 2, false)] { + // Delete markers have no erasure payload geometry. + let mut info = if deleted { + FileInfo::default() + } else { + FileInfo::new("object", 4, 2) + }; + info.name = "object".to_string(); + info.volume = "bucket".to_string(); + info.version_id = Some(version); + info.versioned = true; + info.deleted = deleted; + info.size = if deleted { 0 } else { 100 }; + info.mod_time = Some(OffsetDateTime::from_unix_timestamp(seconds).expect("fixture timestamp")); + metadata.add_version(info).expect("fixture version should be valid"); + } + let entry = MetaCacheEntry { + name: "object".to_string(), + metadata: metadata.marshal_msg().expect("fixture metadata should serialize"), + ..Default::default() + }; + for merged in [false, true] { + let collector = test_collector(); + if merged { + collector.ingest_merged(&MetaCacheEntries(vec![Some(entry.clone()), Some(entry.clone())])); + } else { + collector.ingest(entry.clone()); + } + let objects = collector.lock_objects().expect("collector should remain readable"); + assert_eq!(objects.len(), 1); + let versions = &objects[0].versions; + assert_eq!(versions.len(), 2, "one exact unit per version, including merged duplicates"); + let null = versions + .iter() + .find(|item| item.version_id.as_deref() == Some(Uuid::nil().to_string().as_str())) + .expect("historical null must remain an explicit selector"); + assert_eq!(null.is_delete_marker, null_marker); + assert!( + versions + .iter() + .any(|item| item.version_id.as_deref() == Some(latest.to_string().as_str())) + ); + } + } + } + fn crc_valid_semantically_corrupt_entry(name: &str) -> MetaCacheEntry { let mut metadata = FileMeta::new(); metadata diff --git a/crates/heal/Cargo.toml b/crates/heal/Cargo.toml index aabd39a6f..0792afa6a 100644 --- a/crates/heal/Cargo.toml +++ b/crates/heal/Cargo.toml @@ -97,6 +97,7 @@ crc-fast = { workspace = true } sha2 = { workspace = true } [dev-dependencies] +rustfs-filemeta = { workspace = true } serde_json = { workspace = true, features = ["raw_value"] } rustfs-test-utils = { workspace = true } serial_test = { workspace = true } diff --git a/crates/heal/src/heal/resume.rs b/crates/heal/src/heal/resume.rs index 19d9522ae..7d873f39c 100644 --- a/crates/heal/src/heal/resume.rs +++ b/crates/heal/src/heal/resume.rs @@ -62,10 +62,9 @@ const REPLACEMENT_RECOVERY_CONFLICT_PREFIX: &str = "replacement recovery conflic const REPLACEMENT_RECOVERY_CORRUPTION_PREFIX: &str = "replacement recovery corruption:"; /// Current on-disk schema version for `ResumeState`. Snapshots written by an -/// older schema (which tracked latest-only object names and a positional -/// cursor) are incompatible with the per-version resume cursor, so they are -/// discarded on load and the scan restarts from the beginning. -const CURRENT_RESUME_SCHEMA: u32 = 5; +/// older schema could mark historical null versions covered after reading +/// latest. Discard their cursor and progress and scan from the beginning. +const CURRENT_RESUME_SCHEMA: u32 = 6; /// Persistence throttle for per-object bookkeeping: flush after this many /// buffered mutations or once the interval elapses, whichever comes first. @@ -812,9 +811,7 @@ impl ResumeManager { }); } - // A snapshot written by an older schema tracked a latest-only positional - // cursor that is meaningless under per-version resume. Discard the stale - // progress so the scan restarts cleanly, then stamp the current schema. + // Older progress cannot prove that the exact null slot was inspected. if state.schema_version > CURRENT_RESUME_SCHEMA { return Err(Error::TaskExecutionFailed { message: format!( @@ -824,6 +821,14 @@ impl ResumeManager { }); } if state.schema_version < CURRENT_RESUME_SCHEMA { + // Replacement intents may already have a separate completion proof. + // Resetting only their cursor could revive that stale proof or reopen + // a target for formatting. Preserve ownership for explicit recovery. + if state.replacement_generation.is_some() || !state.replacement_targets.is_empty() { + return Err(replacement_recovery_conflict(format!( + "Replacement intent {task_id} has legacy version coverage; explicit recovery is required before resuming" + ))); + } warn!( target: "rustfs::heal::resume", event = EVENT_HEAL_RESUME_STATE, @@ -849,7 +854,11 @@ impl ResumeManager { state.baseline_known = false; state.counter_unknown = false; state.completed = false; - state.completed_buckets.clear(); + state.pending_buckets.append(&mut state.completed_buckets); + state.pending_buckets.sort(); + state.pending_buckets.dedup(); + state.current_bucket = None; + state.current_object = None; state.schema_version = CURRENT_RESUME_SCHEMA; } diff --git a/crates/heal/src/heal/resume/checkpoint.rs b/crates/heal/src/heal/resume/checkpoint.rs index 505f10a4d..d6dfbb4b9 100644 --- a/crates/heal/src/heal/resume/checkpoint.rs +++ b/crates/heal/src/heal/resume/checkpoint.rs @@ -31,12 +31,12 @@ use super::{ const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state"; const RESUME_CHECKPOINT_DIGEST_FILE: &str = "ahm_checkpoint.sha256"; -const CHECKPOINT_PER_VERSION_SCHEMA: u32 = 5; +const CHECKPOINT_PER_VERSION_SCHEMA: u32 = 7; /// Current on-disk schema version for `ResumeCheckpoint`. Schema 5 could -/// persist dedup identities without the aggregate counters needed to restore -/// them safely, so stale checkpoints are discarded and replayed. -pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6; +/// persist incomplete counters; schema 6 could acknowledge null as latest. +/// Discard those positions and dedup identities and replay the scan. +pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 7; #[derive(Debug, Clone, Copy)] pub enum CheckpointObjectOutcome { diff --git a/crates/heal/src/heal/resume/tests.rs b/crates/heal/src/heal/resume/tests.rs index 73d379b04..7d4d4b87f 100644 --- a/crates/heal/src/heal/resume/tests.rs +++ b/crates/heal/src/heal/resume/tests.rs @@ -1602,6 +1602,152 @@ async fn test_resumestate_schema_v0_discarded_on_load() { temp_dir.close().expect("remove schema test directory"); } +#[tokio::test] +async fn historical_null_progress_is_replayed_after_upgrade() { + use sha2::{Digest, Sha256}; + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let mut state = ResumeState::new( + task_id.clone(), + "erasure_set".to_string(), + "pool_0_set_0".to_string(), + vec!["pending-bucket".to_string()], + ); + state.schema_version = 5; + state.resume_cursor = Some("dw1:old-page".to_string()); + state.completed_buckets = vec!["completed-bucket".to_string()]; + state.processed_objects = 5; + state.successful_objects = 5; + state.completed = true; + let state_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_STATE_FILE}"); + disk.write_all( + RUSTFS_META_BUCKET, + &state_path, + serde_json::to_vec(&state).expect("serialize schema 5").into(), + ) + .await + .expect("persist old resume state"); + + let mut checkpoint = ResumeCheckpoint::new(task_id.clone()); + checkpoint.schema_version = 6; + checkpoint.current_bucket_index = 1; + checkpoint.current_object_index = 5; + checkpoint.processed_objects.insert(compose_key("versions/object.bin", None)); + checkpoint.successful_objects = 5; + // Schema 6 used a digest over the canonical JSON with a null digest field. + let mut wire = serde_json::to_value(&checkpoint).expect("serialize schema 6"); + let unsigned = serde_json::to_vec(&wire).expect("serialize unsigned schema 6"); + wire["integrity_digest"] = serde_json::json!(base64_simd::STANDARD.encode_to_string(Sha256::digest(&unsigned))); + let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}"); + disk.write_all( + RUSTFS_META_BUCKET, + &checkpoint_path, + serde_json::to_vec(&wire).expect("serialize signed schema 6").into(), + ) + .await + .expect("persist old checkpoint"); + + let restored = ResumeManager::load_from_disk(disk.clone(), &task_id) + .await + .expect("load old resume state") + .get_state() + .await; + assert_eq!(restored.schema_version, CURRENT_RESUME_SCHEMA); + assert_eq!(restored.resume_cursor, None); + assert_eq!(restored.processed_objects, 0); + assert_eq!(restored.successful_objects, 0); + assert!(!restored.completed); + assert!(restored.completed_buckets.is_empty()); + assert_eq!( + restored.pending_buckets, + ["completed-bucket", "pending-bucket"], + "formerly completed buckets must be rescanned too" + ); + let restored = CheckpointManager::load_from_disk(disk, &task_id) + .await + .expect("load old signed checkpoint") + .get_checkpoint() + .await; + assert_eq!(restored.schema_version, CURRENT_CHECKPOINT_SCHEMA); + assert!( + restored.processed_objects.is_empty(), + "ambiguous null coverage cannot survive the upgrade" + ); + assert_eq!(restored.current_bucket_index, 0); + assert_eq!(restored.current_object_index, 0); + assert_eq!(restored.successful_objects, 0); + temp_dir.close().expect("remove upgrade fixture"); +} + +#[tokio::test] +async fn legacy_replacement_null_coverage_cannot_reuse_a_completion_proof() { + let (temp_dir, disk) = schema_test_disk().await; + let task_id = ResumeUtils::generate_task_id(); + let manager = ResumeManager::new_replacement_intent( + disk.clone(), + task_id.clone(), + "pool_0_set_0".to_string(), + vec!["bucket".to_string()], + vec!["replacement-a".to_string()], + vec![ReplacementTargetIdentity { + endpoint: "replacement-a".to_string(), + canonical_path: "/mnt/replacement-a".to_string(), + physical_device_ids: vec!["device-a".to_string()], + filesystem_identity: "1:2:3".to_string(), + }], + ) + .await + .expect("persist replacement intent"); + manager + .mark_replacement_completed_and_verified() + .await + .expect("persist the old completion proof"); + let proof_path = replacement_completion_proof_path(&task_id); + let proof_path = proof_path.to_str().expect("proof path must be UTF-8"); + let proof_before = disk + .read_all(RUSTFS_META_BUCKET, proof_path) + .await + .expect("read completion proof"); + let intent_path = ResumeStateFile::ReplacementIntent.path(&task_id); + let intent_path = intent_path.to_str().expect("intent path must be UTF-8"); + for phase in [ + ReplacementPhase::Intent, + ReplacementPhase::Rebuilding, + ReplacementPhase::Verified, + ReplacementPhase::CleanupPending, + ] { + let mut legacy = manager.get_state().await; + legacy.schema_version = 5; + legacy.replacement_phase = phase; + legacy.completed = matches!(phase, ReplacementPhase::Verified | ReplacementPhase::CleanupPending); + let bytes = serde_json::to_vec(&legacy).expect("serialize old replacement state"); + disk.write_all(RUSTFS_META_BUCKET, intent_path, bytes.clone().into()) + .await + .expect("write old replacement state"); + let error = ResumeManager::load_replacement_intent(disk.clone(), &task_id) + .await + .err() + .expect("legacy replacement coverage must require explicit recovery"); + assert!( + error.to_string().contains("legacy version coverage"), + "unexpected recovery error: {error}" + ); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, intent_path) + .await + .expect("retained intent") + .as_ref(), + bytes + ); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, proof_path).await.expect("retained proof"), + proof_before, + "an upgrade must not discard the ownership/completion fence" + ); + } + temp_dir.close().expect("remove replacement upgrade fixture"); +} + #[tokio::test] async fn test_checkpoint_schema_v5_discarded_on_load() { let (temp_dir, disk) = schema_test_disk().await; diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index e02471798..962055742 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -86,6 +86,47 @@ impl From<(HealResultItem, Option)> for HealStorageObjectResult { } } +fn verified_object_receipt( + bucket: &str, + object: &str, + version_id: Option<&str>, + opts: &HealOpts, + item: &HealResultItem, + bucket_incarnation_id: Uuid, +) -> Option { + if opts.dry_run || !item.integrity_verified { + return None; + } + let resolved_version = Uuid::from_bytes(item.resolved_version_id?); + if let Some(requested) = version_id.filter(|version| !version.is_empty()) + && Uuid::parse_str(requested).ok()? != resolved_version + { + return None; + } + item.drives_reported()?; + let drives_healed = item.drives_healed()?; + let ok_drive_state = DriveState::Ok.to_string(); + if !item.after.drives.iter().all(|drive| drive.state == ok_drive_state) { + return None; + } + Some(HealObjectReceipt { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: bucket.to_string(), + object: object.to_string(), + version_id: version_id.map(ToOwned::to_owned), + bucket_incarnation_id: Some(bucket_incarnation_id), + pool_index: opts.pool, + set_index: opts.set, + }, + disposition: if drives_healed > 0 { + HealObjectDisposition::Repaired + } else { + HealObjectDisposition::VerifiedHealthy + }, + }) +} + const LOG_COMPONENT_HEAL: &str = "heal"; const LOG_SUBSYSTEM_STORAGE: &str = "storage"; const EVENT_HEAL_STORAGE_OBJECT_IO: &str = "heal_storage_object_io"; @@ -320,13 +361,13 @@ pub(crate) fn decode_disk_walk_token(token: &str) -> Option { /// `is_delete_marker` is OBSERVABILITY-ONLY (metrics / logging / e2e /// assertions); it MUST NOT gate healing logic. Whether the delete-marker path /// or the data path is taken is decided internally in `ops/heal.rs` from -/// `latest_meta.deleted`. `version_id` is normalized (nil/absent UUID => `None`) -/// at the single construction point in `list_objects_for_heal_page`. +/// `latest_meta.deleted`. Every enumerated version has an exact selector: +/// nil/absent metadata UUIDs select the nil UUID, never an unspecified latest. #[derive(Debug, Clone)] pub struct HealListItem { /// object key pub name: String, - /// normalized version id (`None` when the version is nil/absent) + /// Exact version id, including the nil UUID for the null slot. pub version_id: Option, /// version modification time as Unix nanoseconds pub mod_time_unix_nanos: Option, @@ -1170,35 +1211,11 @@ impl HealStorageAPI for ECStoreHealStorage { None } } else if error.is_none() && !opts.dry_run && item.integrity_verified { - 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 ( - self.ecstore.bucket_incarnation_id(bucket).await, - item.drives_reported(), - item.drives_healed(), - all_after_drives_ok, - ) { - (Ok(bucket_incarnation_id), Some(_), Some(drives_healed), true) => { - let disposition = if drives_healed > 0 { - HealObjectDisposition::Repaired - } else { - HealObjectDisposition::VerifiedHealthy - }; - Some(HealObjectReceipt { - identity: HealObjectIdentity { - kind: HealObjectKind::Object, - bucket: bucket.to_string(), - object: object.to_string(), - version_id: version_id.map(ToOwned::to_owned), - bucket_incarnation_id: Some(bucket_incarnation_id), - pool_index: opts.pool, - set_index: opts.set, - }, - disposition, - }) - } - _ => None, - } + self.ecstore + .bucket_incarnation_id(bucket) + .await + .ok() + .and_then(|incarnation| verified_object_receipt(bucket, object, version_id, opts, &item, incarnation)) } else { None }; @@ -1381,14 +1398,14 @@ impl HealStorageAPI for ECStoreHealStorage { } }; - // Collect versions from this page. version_id is normalized to Option - // here at the single construction point: nil/absent UUID => None. + // Listing has already selected a concrete version. Preserve the null + // slot's identity even when a newer UUID has become latest. let page_objects: Vec = list_info .objects .into_iter() .map(|mut obj| { + let version_id = Some(obj.version_id.unwrap_or_default().to_string()); obj.version_id = obj.version_id.filter(|u| !u.is_nil()); - let version_id = obj.version_id.map(|u| u.to_string()); let mod_time_unix_nanos = obj.mod_time.map(|mod_time| mod_time.unix_timestamp_nanos()); let is_delete_marker = obj.delete_marker; if include_lifecycle_object_info { @@ -1648,6 +1665,68 @@ mod tests { is_transient_object_exists_message, next_heal_listing_token, }; + #[test] + fn object_receipt_requires_integrity_and_resolved_version_evidence() { + use super::{HealObjectDisposition, HealOpts, HealResultItem, Uuid, verified_object_receipt}; + use rustfs_madmin::heal_commands::HealDriveInfo; + let incarnation = Uuid::new_v4(); + let null = Uuid::nil().to_string(); + let latest = Uuid::new_v4(); + let options = HealOpts::default(); + let mut item = HealResultItem { + integrity_verified: true, + version_id: null.clone(), + resolved_version_id: Some(*latest.as_bytes()), + ..Default::default() + }; + let healthy = HealDriveInfo { + state: "ok".to_string(), + ..Default::default() + }; + item.before.drives.push(healthy.clone()); + item.after.drives.push(healthy); + assert!( + verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation).is_none(), + "echoing null cannot certify a different resolved version" + ); + assert!( + verified_object_receipt("bucket", "object", None, &options, &item, incarnation).is_some(), + "an omitted selector still means latest" + ); + assert!(verified_object_receipt("bucket", "object", Some(""), &options, &item, incarnation).is_some()); + assert!( + verified_object_receipt("bucket", "object", Some("null"), &options, &item, incarnation).is_none(), + "the internal boundary requires a UUID, not an S3 spelling" + ); + assert!(verified_object_receipt("bucket", "object", Some(&latest.to_string()), &options, &item, incarnation).is_some()); + + item.resolved_version_id = Some([0; 16]); + let receipt = verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation) + .expect("the exact healthy null version should be certifiable"); + assert_eq!(receipt.identity.version_id.as_deref(), Some(null.as_str())); + assert_eq!(receipt.disposition, HealObjectDisposition::VerifiedHealthy); + item.before.drives[0].state = "missing".to_string(); + assert_eq!( + verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation) + .expect("the restored null version should be certifiable") + .disposition, + HealObjectDisposition::Repaired + ); + item.integrity_verified = false; + assert!( + verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation).is_none(), + "the exact version cannot certify unverified shard integrity" + ); + assert!(verified_object_receipt("bucket", "object", None, &options, &item, incarnation).is_none()); + item.integrity_verified = true; + item.resolved_version_id = None; + assert!( + verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation).is_none(), + "legacy results cannot prove the selected version" + ); + assert!(verified_object_receipt("bucket", "object", None, &options, &item, incarnation).is_none()); + } + #[test] fn next_heal_listing_token_returns_none_for_complete_page() { assert_eq!( diff --git a/crates/heal/tests/heal_b5_versioned_regression_test.rs b/crates/heal/tests/heal_b5_versioned_regression_test.rs index 7b291f3db..2e6210350 100644 --- a/crates/heal/tests/heal_b5_versioned_regression_test.rs +++ b/crates/heal/tests/heal_b5_versioned_regression_test.rs @@ -25,6 +25,7 @@ #![recursion_limit = "256"] use http::HeaderMap; +use rustfs_filemeta::{FileInfo, FileMeta}; use rustfs_heal::heal::{ manager::{HealConfig, HealManager}, storage::{ @@ -46,6 +47,7 @@ mod storage_api; use storage_api::integration::{ BucketOperations, ECStore, MakeBucketOptions, NamespaceLocking as _, ObjectIO as _, ObjectOperations as _, + ShardIntegrityWriteMode, }; /// 256 KiB + change: large enough to be stored as non-inline erasure shards @@ -102,6 +104,7 @@ async fn put_versioned(ecstore: &Arc, bucket: &str, object: &str, data: let mut reader = PutObjReader::from_vec(data.to_vec()); let opts = ObjectOptions { versioned: true, + shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Protected), ..Default::default() }; let info = (**ecstore) @@ -117,7 +120,15 @@ async fn put_versioned(ecstore: &Arc, bucket: &str, object: &str, data: async fn put_unversioned(ecstore: &Arc, bucket: &str, object: &str, data: &[u8]) { let mut reader = PutObjReader::from_vec(data.to_vec()); (**ecstore) - .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .put_object( + bucket, + object, + &mut reader, + &ObjectOptions { + shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Protected), + ..Default::default() + }, + ) .await .expect("unversioned put_object failed"); wait_for_put_tail(ecstore, bucket, object).await; @@ -145,8 +156,8 @@ fn object_dir(disk: &Path, bucket: &str, object: &str) -> PathBuf { disk.join(bucket).join(object) } -/// Count `part.*` data-shard files two levels below the object dir -/// (`//part.N`). One data dir per non-delete-marker version. +/// Count `part.N` data-shard files two levels below the object dir, excluding +/// integrity indexes. One data dir per non-delete-marker version. fn count_part_files(obj_dir: &Path) -> usize { if !obj_dir.exists() { return 0; @@ -156,7 +167,14 @@ fn count_part_files(obj_dir: &Path) -> usize { .max_depth(2) .into_iter() .filter_map(Result::ok) - .filter(|e| e.file_type().is_file() && e.file_name().to_str().map(|n| n.starts_with("part.")).unwrap_or(false)) + .filter(|entry| { + entry.file_type().is_file() + && entry + .file_name() + .to_str() + .and_then(|name| name.strip_prefix("part.")) + .is_some_and(|number| number.parse::().is_ok()) + }) .count() } @@ -252,6 +270,278 @@ async fn read_version(ecstore: &Arc, bucket: &str, object: &str, versio mod serial_tests { use super::*; + fn physical_version(disk: &Path, bucket: &str, object: &str, version: &str) -> FileInfo { + let bytes = std::fs::read(xl_meta_path(&object_dir(disk, bucket, object))).expect("physical xl.meta must exist"); + FileMeta::load_or_convert(&bytes) + .expect("physical metadata must decode") + .into_fileinfo(bucket, object, version, true, false, true) + .expect("the exact physical version must exist") + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[serial] + async fn test_exact_null_heal_with_current_and_historical_markers() { + let (disks, ecstore, storage) = heal_env().await; + let null = uuid::Uuid::nil().to_string(); + let object = "object.bin"; + for (null_marker, newer, latest_marker) in [ + (false, false, false), + (true, false, false), + (true, true, false), + (false, true, true), + (true, true, true), + ] { + let bucket = format!("null-marker-{null_marker}-{newer}-{latest_marker}"); + create_versioned_bucket(&ecstore, &bucket).await; + let old = put_versioned(&ecstore, &bucket, object, &versioned_test_data(10)).await; + put_unversioned(&ecstore, &bucket, object, &versioned_test_data(11)).await; + if null_marker { + let info = ecstore + .delete_object( + &bucket, + object, + ObjectOptions { + version_suspended: true, + ..Default::default() + }, + ) + .await + .expect("suspended delete must create a null marker"); + assert!(info.delete_marker); + assert_eq!(info.version_id.unwrap_or_default(), uuid::Uuid::nil()); + } + let latest = if !newer { + null.clone() + } else if latest_marker { + put_delete_marker(&ecstore, &bucket, object).await + } else { + put_versioned(&ecstore, &bucket, object, &versioned_test_data(12)).await + }; + let mut expected = vec![old, null.clone()]; + if newer { + expected.push(latest.clone()); + } + let originals: Vec<_> = expected + .iter() + .map(|version| physical_version(&disks[0], &bucket, object, version)) + .collect(); + assert_eq!(originals[1].deleted, null_marker); + assert_eq!(originals[1].is_latest, !newer); + let listed = enumerate_all_versions(&storage, &bucket).await; + let (walked, _, truncated) = storage + .list_versions_for_heal_page_disk_walk(SET_DISK_ID, &bucket, "", None, false) + .await + .expect("disk walk must enumerate null and UUID versions"); + assert!(!truncated); + for items in [&listed, &walked] { + assert_eq!(items.len(), expected.len()); + for version in &expected { + assert_eq!( + items + .iter() + .filter(|item| item.version_id.as_deref() == Some(version.as_str())) + .count(), + 1 + ); + } + let item = items + .iter() + .find(|item| item.version_id.as_deref() == Some(null.as_str())) + .expect("exact null entry"); + assert_eq!(item.is_delete_marker, null_marker); + } + std::fs::remove_file(xl_meta_path(&object_dir(&disks[0], &bucket, object))).expect("remove one member's metadata"); + let options = HealOpts { + scan_mode: HealScanMode::Deep, + pool: Some(0), + set: Some(0), + ..Default::default() + }; + for item in listed { + let healed = storage + .heal_object_with_receipt(&bucket, object, item.version_id.as_deref(), &options) + .await + .expect("heal exact version"); + assert!(healed.error.is_none(), "exact version repair failed: {:?}", healed.error); + if item.is_delete_marker { + assert!(!healed.item.integrity_verified); + assert!(healed.receipt.is_none(), "delete markers carry no shard-integrity proof"); + } else { + assert!(healed.item.integrity_verified); + let receipt = healed + .receipt + .expect("verified exact data version repair must produce a receipt"); + assert_eq!(receipt.identity.version_id, item.version_id); + } + assert_eq!( + healed.item.resolved_version_id, + Some( + *uuid::Uuid::parse_str(item.version_id.as_deref().expect("exact selector")) + .expect("UUID selector") + .as_bytes() + ) + ); + } + for (version, original) in expected.iter().zip(originals) { + let repaired = physical_version(&disks[0], &bucket, object, version); + assert_eq!(repaired.version_id, original.version_id); + assert_eq!(repaired.deleted, original.deleted); + assert_eq!(repaired.data_dir, original.data_dir); + assert_eq!(repaired.size, original.size); + } + let latest_result = storage + .heal_object_with_receipt(&bucket, object, None, &options) + .await + .expect("heal latest"); + assert!(latest_result.error.is_none()); + assert_eq!( + latest_result.item.resolved_version_id, + Some(*uuid::Uuid::parse_str(&latest).expect("latest UUID").as_bytes()) + ); + } + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[serial] + async fn test_historical_null_recursive_heal_restores_physical_versions() { + use rustfs_heal::heal::outcome::HealObjectDisposition; + let (disk_paths, ecstore, storage) = heal_env().await; + let manager = HealManager::new( + storage.clone(), + Some(HealConfig { + heal_interval: Duration::from_millis(1), + ..Default::default() + }), + ); + manager.start().await.expect("heal manager should start"); + let null = uuid::Uuid::nil().to_string(); + let object = "versions/object.bin"; + + for (parity, missing_metadata) in [(false, false), (true, false), (false, true), (true, true)] { + let bucket = format!("historical-null-{parity}-{missing_metadata}"); + create_versioned_bucket(&ecstore, &bucket).await; + let mut versions = Vec::new(); + for seed in 1..=3 { + let data = versioned_test_data(seed); + let version = put_versioned(&ecstore, &bucket, object, &data).await; + versions.push((version, data)); + } + // Exercise the suspended null slot, including an overwrite, before + // a new versioned write makes that slot historical. + put_unversioned(&ecstore, &bucket, object, &versioned_test_data(4)).await; + let null_data = versioned_test_data(5); + put_unversioned(&ecstore, &bucket, object, &null_data).await; + versions.push((null.clone(), null_data.clone())); + let mut latest_data = versioned_test_data(6); + latest_data.extend_from_slice(b"new-uuid"); + let latest = put_versioned(&ecstore, &bucket, object, &latest_data).await; + versions.push((latest, latest_data.clone())); + + let target = disk_paths + .iter() + .find(|disk| { + let info = physical_version(disk, &bucket, object, &null); + assert!(!info.is_latest, "the damaged null must be historical"); + info.erasure.index == if parity { info.erasure.data_blocks + 1 } else { 1 } + }) + .expect("the requested data/parity member must exist"); + let originals: Vec<_> = versions + .iter() + .map(|(version, _)| { + let info = physical_version(target, &bucket, object, version); + let part = object_dir(target, &bucket, object) + .join(info.data_dir.expect("data version must have a data directory").to_string()) + .join("part.1"); + let bytes = std::fs::read(&part).expect("original physical shard must exist"); + (info, part, bytes) + }) + .collect(); + if missing_metadata { + std::fs::remove_file(xl_meta_path(&object_dir(target, &bucket, object))).expect("remove one xl.meta"); + } else { + std::fs::remove_file(&originals[3].1).expect("remove only the historical null shard"); + } + + let request = HealRequest::new( + HealType::Bucket { bucket: bucket.clone() }, + HealOptions { + recursive: true, + scan_mode: HealScanMode::Deep, + pool_index: Some(0), + set_index: Some(0), + ..Default::default() + }, + HealPriority::Normal, + ); + let task_id = request.id.clone(); + assert!( + manager + .submit_heal_request(request) + .await + .expect("submit recursive heal") + .is_admitted() + ); + wait_for_task(&manager, &task_id, Duration::from_secs(60)).await; + + // Inspect physical repair before any GET can trigger read repair. + for ((version, _), (original, part, bytes)) in versions.iter().zip(&originals) { + assert_eq!(std::fs::read(part).expect("heal must restore every physical shard"), *bytes); + let repaired = physical_version(target, &bucket, object, version); + assert_eq!(repaired.version_id, original.version_id); + assert_eq!(repaired.data_dir, original.data_dir); + assert_eq!(repaired.size, original.size); + } + let report = manager.get_task_report(&task_id).await.expect("completed task report"); + let outcome = report.outcome.expect("recursive task must have an outcome"); + assert_eq!(outcome.counters.processed, 5); + assert_eq!(outcome.counters.failed, 0); + assert_eq!(outcome.counters.unknown, 0); + assert_eq!(outcome.counters.healed, if missing_metadata { 5 } else { 1 }); + let null_outcome = outcome + .objects + .iter() + .find(|item| item.identity.version_id.as_deref() == Some(null.as_str())) + .expect("historical null must have its own exact receipt"); + assert_eq!(null_outcome.disposition, HealObjectDisposition::Repaired); + let null_item = report + .result_items + .iter() + .find(|item| item.version_id == null) + .expect("historical null must have its own result item"); + assert_eq!(null_item.object_size, null_data.len(), "null must not report the latest UUID's size"); + + let listed = enumerate_all_versions(&storage, &bucket).await; + assert_eq!(listed.len(), 5); + assert!(listed.iter().any(|item| item.version_id.as_deref() == Some(null.as_str()))); + let latest_result = storage + .heal_object_with_receipt(&bucket, object, None, &HealOpts::default()) + .await + .expect("an omitted selector must still heal latest"); + assert!(latest_result.error.is_none()); + assert_eq!(latest_result.item.object_size, latest_data.len()); + assert!(latest_result.receipt.is_none(), "a normal scan cannot certify payload integrity"); + let verified_latest = storage + .heal_object_with_receipt( + &bucket, + object, + None, + &HealOpts { + scan_mode: HealScanMode::Deep, + ..Default::default() + }, + ) + .await + .expect("deep heal of latest"); + assert!(verified_latest.error.is_none()); + assert_eq!(verified_latest.item.object_size, latest_data.len()); + assert!(verified_latest.receipt.is_some(), "a deep scan can certify the exact latest version"); + for (version, data) in &versions { + assert_eq!(&read_version(&ecstore, &bucket, object, version).await, data); + } + } + manager.stop().await.expect("heal manager should stop"); + } + /// Directly exercises `ECStoreHealStorage::list_objects_for_heal_page` on a /// real versioned fixture: two data versions + a delete-marker-latest. Proves /// enumeration returns one `HealListItem` per version, flags the delete marker @@ -416,8 +706,7 @@ mod serial_tests { ); } - /// Unversioned objects normalize to `version_id == None` and enumerate once - /// each (never a phantom `Some("null")`). + /// Enumerated unversioned objects select the exact null slot once each. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] async fn test_version_id_normalization_null_and_unversioned_real_fixture() { @@ -434,21 +723,14 @@ mod serial_tests { let items = enumerate_all_versions(&heal_storage, bucket).await; assert_eq!(items.len(), 2, "unversioned bucket => exactly one heal unit per object, got {items:?}"); assert!( - items.iter().all(|it| it.version_id.is_none()), - "unversioned objects must normalize to version_id == None, never Some(\"null\"): {items:?}" + items + .iter() + .all(|it| it.version_id.as_deref() == Some(uuid::Uuid::nil().to_string().as_str())), + "enumerated null versions must retain an exact selector: {items:?}" ); assert!(items.iter().all(|it| !it.is_delete_marker), "no delete markers expected"); let names: std::collections::HashSet<&str> = items.iter().map(|it| it.name.as_str()).collect(); assert!(names.contains("a.bin") && names.contains("b.bin"), "both objects enumerated once"); - - // NOTE: A literal MinIO-interop "null" *string* version id is only minted - // by the S3 request layer (versioning-suspended writes); it cannot be - // constructed through the internal ECStore put path used here (suspended/ - // unversioned writes store no version id at all -> None). The nil-UUID -> - // None normalization itself is unit-covered in B5-1 - // (crates/heal storage list_objects_for_heal_page maps - // `version_id.filter(|u| !u.is_nil())`). This test covers the real - // unversioned-object => None path end-to-end. } /// Unversioned bucket, full heal path via `HealManager` (recursive bucket @@ -479,7 +761,11 @@ mod serial_tests { // Enumeration returns exactly one unit per object (no duplicates). let items = enumerate_all_versions(&heal_storage, bucket).await; assert_eq!(items.len(), objects.len(), "one heal unit per object"); - assert!(items.iter().all(|it| it.version_id.is_none())); + assert!( + items + .iter() + .all(|it| it.version_id.as_deref() == Some(uuid::Uuid::nil().to_string().as_str())) + ); // Drive the real recursive bucket heal through the HealManager task loop. let cfg = HealConfig { diff --git a/crates/madmin/src/heal_commands.rs b/crates/madmin/src/heal_commands.rs index ac4340bec..b2ac27930 100644 --- a/crates/madmin/src/heal_commands.rs +++ b/crates/madmin/src/heal_commands.rs @@ -50,6 +50,10 @@ pub struct HealResultItem { pub object: String, #[serde(rename = "versionId")] pub version_id: String, + /// Exact version selected by the storage owner, including the nil UUID. + /// Wire-decoded and legacy results cannot supply this in-process proof. + #[serde(skip)] + pub resolved_version_id: Option<[u8; 16]>, #[serde(rename = "detail")] pub detail: String, #[serde(rename = "parityBlocks")] @@ -100,6 +104,20 @@ impl HealResultItem { mod tests { use super::*; + #[test] + fn resolved_heal_version_is_not_wire_proof() { + let item = HealResultItem { + resolved_version_id: Some([0; 16]), + ..Default::default() + }; + let mut wire = serde_json::to_value(&item).expect("heal result should serialize"); + assert!(wire.get("resolved_version_id").is_none()); + assert!(wire.get("resolvedVersionId").is_none()); + wire["resolved_version_id"] = serde_json::json!(vec![0; 16]); + let decoded: HealResultItem = serde_json::from_value(wire).expect("legacy wire shape should remain readable"); + assert_eq!(decoded.resolved_version_id, None, "wire input cannot supply owner proof"); + } + fn drive(state: &str) -> HealDriveInfo { HealDriveInfo { uuid: String::new(),