mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 20:43:04 +00:00
fix(heal): report verified normal shard repairs (#7890)
This commit is contained in:
@@ -1732,6 +1732,17 @@ impl SetDisks {
|
||||
return Ok((result, Some(error)));
|
||||
}
|
||||
|
||||
result.repair_verified = protected
|
||||
&& !latest_meta.deleted
|
||||
&& !latest_meta.is_remote()
|
||||
&& !read_repair_uses_shared_lock
|
||||
&& result.drives_healed().is_some_and(|healed| healed > 0)
|
||||
&& result
|
||||
.after
|
||||
.drives
|
||||
.iter()
|
||||
.all(|drive| drive.state == DriveState::Ok.to_string());
|
||||
|
||||
// The object is healthy here; sweep any data dirs left behind
|
||||
// by pre-#3510 unversioned overwrites, which the dangling paths
|
||||
// above never touch (issues #3231, #3191). Best effort — a
|
||||
|
||||
@@ -105,7 +105,7 @@ fn verified_object_receipt(
|
||||
item: &HealResultItem,
|
||||
bucket_incarnation_id: Uuid,
|
||||
) -> Option<HealObjectReceipt> {
|
||||
if opts.dry_run || !item.integrity_verified {
|
||||
if opts.dry_run || (!item.integrity_verified && !item.repair_verified) {
|
||||
return None;
|
||||
}
|
||||
let resolved_version = Uuid::from_bytes(item.resolved_version_id?);
|
||||
@@ -116,6 +116,9 @@ fn verified_object_receipt(
|
||||
}
|
||||
item.drives_reported()?;
|
||||
let drives_healed = item.drives_healed()?;
|
||||
if item.repair_verified && drives_healed == 0 {
|
||||
return None;
|
||||
}
|
||||
let ok_drive_state = DriveState::Ok.to_string();
|
||||
if !item.after.drives.iter().all(|drive| drive.state == ok_drive_state) {
|
||||
return None;
|
||||
@@ -132,8 +135,10 @@ fn verified_object_receipt(
|
||||
},
|
||||
disposition: if drives_healed > 0 {
|
||||
HealObjectDisposition::Repaired
|
||||
} else {
|
||||
} else if item.integrity_verified {
|
||||
HealObjectDisposition::VerifiedHealthy
|
||||
} else {
|
||||
return None;
|
||||
},
|
||||
})
|
||||
}
|
||||
@@ -682,7 +687,7 @@ impl ECStoreHealStorage {
|
||||
} else {
|
||||
None
|
||||
}
|
||||
} else if error.is_none() && !opts.dry_run && item.integrity_verified {
|
||||
} else if error.is_none() && !opts.dry_run && (item.integrity_verified || item.repair_verified) {
|
||||
let bucket_incarnation_id = match expected {
|
||||
Some(expected) => Some(expected),
|
||||
None => self.ecstore.bucket_incarnation_id(bucket).await.ok(),
|
||||
@@ -1848,6 +1853,20 @@ mod tests {
|
||||
"the exact version cannot certify unverified shard integrity"
|
||||
);
|
||||
assert!(verified_object_receipt("bucket", "object", None, &options, &item, incarnation).is_none());
|
||||
item.repair_verified = true;
|
||||
item.resolved_version_id = Some([0; 16]);
|
||||
assert_eq!(
|
||||
verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation)
|
||||
.expect("a protected repaired shard should carry a repair receipt")
|
||||
.disposition,
|
||||
HealObjectDisposition::Repaired
|
||||
);
|
||||
item.before.drives[0].state = "ok".to_string();
|
||||
assert!(
|
||||
verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation).is_none(),
|
||||
"repair proof without a repaired drive must remain unresolved"
|
||||
);
|
||||
item.repair_verified = false;
|
||||
item.integrity_verified = true;
|
||||
item.resolved_version_id = None;
|
||||
assert!(
|
||||
|
||||
@@ -379,3 +379,162 @@ fn admin_pool_set_repairs_five_shards_and_retains_exact_terminal_after_restart()
|
||||
.join()
|
||||
.expect("heal fixture result");
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial]
|
||||
fn admin_normal_protected_repairs_report_repaired_for_data_and_parity() {
|
||||
std::thread::Builder::new()
|
||||
.stack_size(8 * 1024 * 1024)
|
||||
.spawn(|| {
|
||||
tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
.expect("heal fixture runtime")
|
||||
.block_on(async {
|
||||
use rustfs_heal::heal::{
|
||||
manager::{HealConfig, HealManager},
|
||||
task::{HealOptions, HealPriority, HealRequest, HealTaskStatus, HealType},
|
||||
};
|
||||
use rustfs_heal_contracts::heal_channel::HealRequestSource;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use storage_api::integration::WriteCompletion;
|
||||
|
||||
let root = tempfile::tempdir().expect("normal repair fixture");
|
||||
let env = TestECStoreEnv::builder()
|
||||
.base_dir(root.path())
|
||||
.prefix("admin_normal_repair_receipts")
|
||||
.build()
|
||||
.await;
|
||||
let bucket = "admin-normal-repair";
|
||||
env.make_bucket(bucket, false).await;
|
||||
let set = env.ecstore.pools[0].get_disks(0);
|
||||
let disks = set.disks.read().await.iter().flatten().cloned().collect::<Vec<_>>();
|
||||
let mut faults = Vec::new();
|
||||
for (object, want_data, byte) in [("missing-data", true, 0x31_u8), ("missing-parity", false, 0x72_u8)] {
|
||||
let body = vec![byte; 1024 * 1024 + 37];
|
||||
env.ecstore
|
||||
.put_object(
|
||||
bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(body.clone()),
|
||||
&ObjectOptions {
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Protected),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("commit protected fixture");
|
||||
let mut metadata = Vec::new();
|
||||
for disk in &disks {
|
||||
metadata.push(
|
||||
disk.read_version("", bucket, object, "", &ReadOptions::default())
|
||||
.await
|
||||
.expect("fixture metadata"),
|
||||
);
|
||||
}
|
||||
let target_slot = metadata
|
||||
.iter()
|
||||
.position(|part| (part.erasure.index <= part.erasure.data_blocks) == want_data)
|
||||
.expect("fixture must expose both data and parity shards");
|
||||
let target = &metadata[target_slot];
|
||||
let path = env.disk_paths[target_slot]
|
||||
.join(bucket)
|
||||
.join(object)
|
||||
.join(target.data_dir.expect("external shard").to_string())
|
||||
.join("part.1");
|
||||
let original = tokio::fs::read(&path).await.expect("original shard bytes");
|
||||
tokio::fs::remove_file(&path).await.expect("inject missing shard");
|
||||
faults.push((object.to_owned(), path, original, body));
|
||||
}
|
||||
|
||||
let storage = Arc::new(ECStoreHealStorage::new(env.ecstore.clone()));
|
||||
let manager = HealManager::new(
|
||||
storage.clone(),
|
||||
Some(HealConfig {
|
||||
enable_auto_heal: false,
|
||||
..Default::default()
|
||||
}),
|
||||
);
|
||||
manager.start().await.expect("start heal manager");
|
||||
let mut request = HealRequest::new(
|
||||
HealType::ErasureSet {
|
||||
buckets: Vec::new(),
|
||||
set_disk_id: "pool_0_set_0".to_owned(),
|
||||
},
|
||||
HealOptions {
|
||||
recursive: true,
|
||||
scan_mode: HealScanMode::Normal,
|
||||
pool_index: Some(0),
|
||||
set_index: Some(0),
|
||||
timeout: Some(Duration::from_secs(60)),
|
||||
..Default::default()
|
||||
},
|
||||
HealPriority::High,
|
||||
);
|
||||
request.source = HealRequestSource::Admin;
|
||||
let token = request.id.clone();
|
||||
manager
|
||||
.submit_heal_request(request)
|
||||
.await
|
||||
.expect("admit normal all-buckets heal");
|
||||
let terminal = tokio::time::timeout(Duration::from_secs(60), async {
|
||||
loop {
|
||||
let report = manager.get_task_report(&token).await.expect("same-token report");
|
||||
if matches!(
|
||||
report.status,
|
||||
HealTaskStatus::Completed | HealTaskStatus::Failed { .. } | HealTaskStatus::Timeout
|
||||
) {
|
||||
break report;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(20)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("terminal deadline");
|
||||
assert_eq!(terminal.status, HealTaskStatus::Completed);
|
||||
let outcome = terminal.outcome.as_deref().expect("canonical terminal");
|
||||
assert_eq!(
|
||||
(
|
||||
outcome.counters.processed,
|
||||
outcome.counters.healed,
|
||||
outcome.counters.unchanged,
|
||||
outcome.counters.skipped,
|
||||
outcome.counters.failed,
|
||||
outcome.counters.unknown,
|
||||
),
|
||||
(2, 2, 0, 0, 0, 0),
|
||||
"normal protected repairs must be canonical: {outcome:?}"
|
||||
);
|
||||
assert_eq!(outcome.objects.len(), 2);
|
||||
for (object, path, original, body) in &faults {
|
||||
assert_eq!(&tokio::fs::read(path).await.expect("rebuilt physical shard"), original);
|
||||
let receipt = outcome
|
||||
.objects
|
||||
.iter()
|
||||
.find(|receipt| receipt.identity.bucket == bucket && receipt.identity.object == *object)
|
||||
.expect("one receipt per repaired object");
|
||||
assert_eq!(receipt.disposition, HealObjectDisposition::Repaired);
|
||||
assert_eq!((receipt.identity.pool_index, receipt.identity.set_index), (Some(0), Some(0)));
|
||||
assert!(receipt.identity.version_id.is_some());
|
||||
assert_eq!(
|
||||
receipt.identity.bucket_incarnation_id,
|
||||
Some(storage.admit_bucket_incarnation(bucket).await.expect("original bucket"))
|
||||
);
|
||||
let mut reader = env
|
||||
.ecstore
|
||||
.get_object_reader(bucket, object, None, Default::default(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("read repaired object");
|
||||
let mut actual = Vec::new();
|
||||
reader.stream.read_to_end(&mut actual).await.expect("read repaired body");
|
||||
assert_eq!(&actual, body);
|
||||
}
|
||||
manager.stop().await.expect("stop heal manager");
|
||||
});
|
||||
})
|
||||
.expect("heal fixture thread")
|
||||
.join()
|
||||
.expect("heal fixture result");
|
||||
}
|
||||
|
||||
@@ -40,6 +40,11 @@ pub struct HealResultItem {
|
||||
/// or old-peer results default to unproven and cannot advance heal receipts.
|
||||
#[serde(skip)]
|
||||
pub integrity_verified: bool,
|
||||
/// Storage-owner proof that a protected missing shard was reconstructed and
|
||||
/// committed. This is intentionally in-process only and does not certify a
|
||||
/// healthy or unchanged object.
|
||||
#[serde(skip)]
|
||||
pub repair_verified: bool,
|
||||
#[serde(rename = "resultId")]
|
||||
pub result_index: usize,
|
||||
#[serde(rename = "type")]
|
||||
@@ -113,6 +118,8 @@ mod tests {
|
||||
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());
|
||||
assert!(wire.get("repair_verified").is_none());
|
||||
assert!(wire.get("repairVerified").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");
|
||||
|
||||
Reference in New Issue
Block a user