fix(heal): preserve historical null version identity (#7749)

This commit is contained in:
cxymds
2026-09-13 21:43:43 +08:00
committed by GitHub
parent 244e7dfb99
commit 41770983d6
10 changed files with 662 additions and 70 deletions
Generated
+1
View File
@@ -9935,6 +9935,7 @@ dependencies = [
"rustfs-concurrency",
"rustfs-config",
"rustfs-ecstore",
"rustfs-filemeta",
"rustfs-heal-contracts",
"rustfs-lock",
"rustfs-madmin",
+1
View File
@@ -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.
+56 -5
View File
@@ -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<String>,
/// version modification time as Unix nanoseconds
pub mod_time_unix_nanos: Option<i128>,
@@ -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
+1
View File
@@ -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 }
+17 -8
View File
@@ -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;
}
+4 -4
View File
@@ -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 {
+146
View File
@@ -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;
+114 -35
View File
@@ -86,6 +86,47 @@ impl From<(HealResultItem, Option<Error>)> for HealStorageObjectResult {
}
}
fn verified_object_receipt(
bucket: &str,
object: &str,
version_id: Option<&str>,
opts: &HealOpts,
item: &HealResultItem,
bucket_incarnation_id: Uuid,
) -> Option<HealObjectReceipt> {
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<String> {
/// `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<String>,
/// version modification time as Unix nanoseconds
pub mod_time_unix_nanos: Option<i128>,
@@ -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<String>
// 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<HealListItem> = 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!(
@@ -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<ECStore>, 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<ECStore>, bucket: &str, object: &str, data:
async fn put_unversioned(ecstore: &Arc<ECStore>, 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
/// (`<object>/<data-uuid>/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::<usize>().is_ok())
})
.count()
}
@@ -252,6 +270,278 @@ async fn read_version(ecstore: &Arc<ECStore>, 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 {
+18
View File
@@ -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(),