mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-23 19:06:30 +00:00
fix(heal): close tail-truncation and stale current-marker gaps (#8019)
* fix(heal): verify complete encoded shard tails * fix(heal): prove stale current marker deletion
This commit is contained in:
@@ -587,7 +587,11 @@ pub(crate) async fn verify_deep_parts(
|
||||
return Err(DiskError::FileCorrupt);
|
||||
}
|
||||
let part_status = statuses.get_mut(&part_index).ok_or(DiskError::FileCorrupt)?;
|
||||
let length = erasure.shard_file_offset(0, part.size, part.size);
|
||||
// Deep verification reads the complete encoded shard. `shard_file_offset`
|
||||
// is a range-read helper and can stop before the final encoded byte for
|
||||
// a partial data block, which would let a one-byte tail truncation pass
|
||||
// as healthy. Use the physical shard length for the integrity proof.
|
||||
let length = usize::try_from(erasure.shard_file_size(part.size as i64)).map_err(|_| DiskError::FileCorrupt)?;
|
||||
let mut readers = Vec::with_capacity(disks.len());
|
||||
for (index, disk) in disks.iter().enumerate() {
|
||||
if part_status[index] != CHECK_PART_UNKNOWN {
|
||||
|
||||
@@ -3915,6 +3915,72 @@ pub struct SetDisks {
|
||||
>,
|
||||
}
|
||||
|
||||
fn marker_purge_receipt_path(
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version: Uuid,
|
||||
purge: &rustfs_common::mrf_channel::MrfDeleteMarkerPurge,
|
||||
) -> String {
|
||||
let object_hash = rustfs_utils::crypto::hex(Sha256::digest(object.as_bytes()));
|
||||
let marker_hash = rustfs_utils::crypto::hex(purge.marker_identity);
|
||||
format!(
|
||||
"marker-purge-receipts/{}/{}/{}-{}-{}.json",
|
||||
bucket, purge.bucket_incarnation_id, version, object_hash, marker_hash
|
||||
)
|
||||
}
|
||||
|
||||
impl SetDisks {
|
||||
/// Persist a quorum DELETE receipt outside the object marker itself. The
|
||||
/// receipt is written to the metadata bucket so it survives a node restart
|
||||
/// while an inaccessible member still holds the stale marker.
|
||||
pub(crate) async fn persist_marker_purge_receipt(
|
||||
&self,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version: Uuid,
|
||||
purge: &rustfs_common::mrf_channel::MrfDeleteMarkerPurge,
|
||||
) -> bool {
|
||||
let path = marker_purge_receipt_path(bucket, object, version, purge);
|
||||
let api = Arc::new(self.clone());
|
||||
crate::config::com::save_config_with_opts(
|
||||
api,
|
||||
&path,
|
||||
b"rustfs-marker-purge-receipt-v1".to_vec(),
|
||||
&ObjectOptions {
|
||||
max_parity: true,
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.is_ok()
|
||||
}
|
||||
|
||||
pub(crate) async fn has_marker_purge_receipt(
|
||||
&self,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version: Uuid,
|
||||
purge: &rustfs_common::mrf_channel::MrfDeleteMarkerPurge,
|
||||
) -> bool {
|
||||
let path = marker_purge_receipt_path(bucket, object, version, purge);
|
||||
crate::config::com::read_config_limited_preserve_empty(Arc::new(self.clone()), &path, 128)
|
||||
.await
|
||||
.is_ok()
|
||||
}
|
||||
|
||||
pub(crate) async fn consume_marker_purge_receipt(
|
||||
&self,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version: Uuid,
|
||||
purge: &rustfs_common::mrf_channel::MrfDeleteMarkerPurge,
|
||||
) {
|
||||
let path = marker_purge_receipt_path(bucket, object, version, purge);
|
||||
let _ = crate::config::com::delete_config_no_lock(Arc::new(self.clone()), &path).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Read every physical copy before selecting a version quorum. A minority
|
||||
/// legacy record is still evidence and must not disappear behind a majority
|
||||
/// not-found result. Only an explicit file/volume absence produces `None`;
|
||||
|
||||
@@ -38,6 +38,7 @@ use crate::disk::DiskAPI;
|
||||
use crate::disk::local::{DELETE_DATA_DIR_MARKER_PREFIX, metadata_less_part_file};
|
||||
use crate::io_support::bitrot::object_mmap_read_enabled;
|
||||
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
|
||||
use rustfs_common::mrf_channel::MrfDeleteMarkerPurge;
|
||||
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit};
|
||||
use tracing::trace;
|
||||
|
||||
@@ -1469,7 +1470,12 @@ impl SetDisks {
|
||||
})
|
||||
.transpose()
|
||||
.map_err(DiskError::from)?;
|
||||
let till_offset = erasure.shard_file_offset(0, part.size, part.size);
|
||||
// This reader covers the whole encoded shard. A range
|
||||
// offset can exclude the final encoded byte when the
|
||||
// logical part ends inside a data block, allowing a
|
||||
// one-byte tail truncation to evade deep Heal.
|
||||
let till_offset = usize::try_from(erasure.shard_file_size(part.size as i64))
|
||||
.map_err(|_| DiskError::FileCorrupt)?;
|
||||
let use_mmap_read = object_mmap_read_enabled();
|
||||
|
||||
let mut readers = Vec::with_capacity(latest_disks.len());
|
||||
@@ -1907,13 +1913,26 @@ impl SetDisks {
|
||||
.zip(&errs)
|
||||
.any(|(info, error)| error.is_none() && info.is_canonical_delete_marker())
|
||||
{
|
||||
let cleanup = match Self::retired_marker_candidate(&parts_metadata, &errs) {
|
||||
Ok(marker) => {
|
||||
self.remove_retired_marker(retirement, bucket, object, marker, opts, &disks)
|
||||
.await
|
||||
}
|
||||
Err(error) => Err(error),
|
||||
};
|
||||
// A current-generation delete marker can be left behind on
|
||||
// an inaccessible member after an acknowledged S3 DELETE.
|
||||
// That case has no bucket-retirement record, so evaluate its
|
||||
// stricter same-incarnation proof before the retired-marker
|
||||
// path. The proof requires an exact marker identity, at
|
||||
// least one replica already absent, and every target online.
|
||||
let cleanup =
|
||||
match Self::same_incarnation_marker_candidate(&parts_metadata, &errs, retirement.current_incarnation) {
|
||||
Ok(marker) => {
|
||||
self.remove_current_marker(retirement, bucket, object, marker, opts, &disks)
|
||||
.await
|
||||
}
|
||||
Err(_) => match Self::retired_marker_candidate(&parts_metadata, &errs) {
|
||||
Ok(marker) => {
|
||||
self.remove_retired_marker(retirement, bucket, object, marker, opts, &disks)
|
||||
.await
|
||||
}
|
||||
Err(error) => Err(error),
|
||||
},
|
||||
};
|
||||
let absent = cleanup.is_ok();
|
||||
let item = self
|
||||
.dangling_heal_result(FileInfo::default(), &errs, bucket, object, version_id, absent)
|
||||
@@ -2027,6 +2046,41 @@ impl SetDisks {
|
||||
candidate.ok_or_else(|| DiskError::retired_marker_deferred("no exact marker candidate"))
|
||||
}
|
||||
|
||||
fn same_incarnation_marker_candidate(
|
||||
metadata: &[FileInfo],
|
||||
errors: &[Option<DiskError>],
|
||||
current_incarnation: Option<Uuid>,
|
||||
) -> disk::error::Result<rustfs_filemeta::MetaDeleteMarker> {
|
||||
let defer = DiskError::retired_marker_deferred;
|
||||
let Some(current) = current_incarnation.filter(|id| !id.is_nil()) else {
|
||||
return Err(defer("current bucket incarnation is unavailable"));
|
||||
};
|
||||
let mut candidate = None;
|
||||
let mut missing_replica = false;
|
||||
for (info, error) in metadata.iter().zip(errors) {
|
||||
match error {
|
||||
Some(DiskError::FileNotFound | DiskError::FileVersionNotFound) => {
|
||||
missing_replica = true;
|
||||
continue;
|
||||
}
|
||||
Some(_) => return Err(defer("not every replica is readable")),
|
||||
None => {}
|
||||
}
|
||||
if !info.is_canonical_delete_marker() || info.delete_marker_incarnation() != Some(current) {
|
||||
return Err(defer("marker is not a canonical current-generation marker"));
|
||||
}
|
||||
let marker = rustfs_filemeta::MetaDeleteMarker::from(info.clone());
|
||||
if candidate.as_ref().is_some_and(|previous| previous != &marker) {
|
||||
return Err(defer("surviving marker identities conflict"));
|
||||
}
|
||||
candidate = Some(marker);
|
||||
}
|
||||
if !missing_replica {
|
||||
return Err(defer("marker is present on every replica"));
|
||||
}
|
||||
candidate.ok_or_else(|| defer("no exact current-generation marker candidate"))
|
||||
}
|
||||
|
||||
async fn remove_retired_marker(
|
||||
&self,
|
||||
context: &MarkerRetirementContext<'_>,
|
||||
@@ -2065,6 +2119,72 @@ impl SetDisks {
|
||||
if context.lifecycle_guard.is_lock_lost() {
|
||||
return Err(defer("bucket lifecycle fence was lost"));
|
||||
}
|
||||
let result = self
|
||||
.remove_marker_exact_under_fence(bucket, object, marker, version, disks)
|
||||
.await;
|
||||
if context.lifecycle_guard.is_lock_lost() {
|
||||
return Err(defer("bucket lifecycle fence was lost"));
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
async fn remove_current_marker(
|
||||
&self,
|
||||
context: &MarkerRetirementContext<'_>,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
marker: rustfs_filemeta::MetaDeleteMarker,
|
||||
opts: &HealOpts,
|
||||
disks: &[Option<DiskStore>],
|
||||
) -> disk::error::Result<()> {
|
||||
let defer = DiskError::retired_marker_deferred;
|
||||
if opts.dry_run || !opts.remove {
|
||||
return Err(defer("cleanup requires remove=true and dry_run=false"));
|
||||
}
|
||||
if disks.len() != self.set_drive_count || disks.iter().any(Option::is_none) {
|
||||
return Err(defer("every target disk must be online"));
|
||||
}
|
||||
let Some(current) = context.current_incarnation.filter(|id| !id.is_nil()) else {
|
||||
return Err(defer("current bucket incarnation is unavailable"));
|
||||
};
|
||||
if context.lifecycle_guard.is_lock_lost() {
|
||||
return Err(defer("bucket lifecycle fence was lost"));
|
||||
}
|
||||
let incarnation = marker.into_fileinfo(bucket, object, false)?.delete_marker_incarnation();
|
||||
if incarnation != Some(current) {
|
||||
return Err(defer("marker is not from the current bucket incarnation"));
|
||||
}
|
||||
let Some(version) = marker.version_id.filter(|id| !id.is_nil()) else {
|
||||
return Err(defer("cleanup requires an explicit non-null version"));
|
||||
};
|
||||
let marker_identity = marker.stable_identity();
|
||||
let marker_bytes = marker.marshal_msg().map_err(DiskError::other)?;
|
||||
let Some(purge) = MrfDeleteMarkerPurge::new(current, current, marker_identity, marker_bytes) else {
|
||||
return Err(defer("marker purge proof payload is invalid"));
|
||||
};
|
||||
if !self.has_marker_purge_receipt(bucket, object, version, &purge).await {
|
||||
return Err(defer("same-incarnation marker lacks an acknowledged delete receipt"));
|
||||
}
|
||||
let result = self
|
||||
.remove_marker_exact_under_fence(bucket, object, marker, version, disks)
|
||||
.await;
|
||||
if context.lifecycle_guard.is_lock_lost() {
|
||||
return Err(defer("bucket lifecycle fence was lost"));
|
||||
}
|
||||
if result.is_ok() {
|
||||
self.consume_marker_purge_receipt(bucket, object, version, &purge).await;
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
async fn remove_marker_exact_under_fence(
|
||||
&self,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
marker: rustfs_filemeta::MetaDeleteMarker,
|
||||
version: Uuid,
|
||||
disks: &[Option<DiskStore>],
|
||||
) -> disk::error::Result<()> {
|
||||
let request = FileInfo {
|
||||
volume: bucket.to_owned(),
|
||||
name: object.to_owned(),
|
||||
@@ -2105,13 +2225,14 @@ impl SetDisks {
|
||||
let same_targets = selected.len() == disks.len() && selected.iter().zip(disks).all(|(current, original)| {
|
||||
matches!((current, original), (Some(current), Some(original)) if std::sync::Arc::ptr_eq(current, original))
|
||||
});
|
||||
if context.lifecycle_guard.is_lock_lost()
|
||||
|| !same_targets
|
||||
if !same_targets
|
||||
|| !errors
|
||||
.iter()
|
||||
.all(|error| matches!(error, Some(DiskError::FileNotFound | DiskError::FileVersionNotFound)))
|
||||
{
|
||||
return Err(defer("complete marker absence could not be verified under the original fence"));
|
||||
return Err(DiskError::retired_marker_deferred(
|
||||
"complete marker absence could not be verified under the original target set",
|
||||
));
|
||||
}
|
||||
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||
Ok(())
|
||||
@@ -3319,7 +3440,9 @@ mod heal_result_report_tests {
|
||||
use crate::disk::{DiskAPI as _, DiskOption, DiskStore, RUSTFS_META_TMP_BUCKET, ReadOptions, STORAGE_FORMAT_FILE, new_disk};
|
||||
use crate::error::Error;
|
||||
use crate::object_api::{ObjectOptions, PutObjReader};
|
||||
use crate::set_disk::ops::object::hermetic_set_disks_support::hermetic_set_disks_isolated;
|
||||
use crate::set_disk::ops::object::hermetic_set_disks_support::{
|
||||
hermetic_set_disks_for_pool_with_default_parity_isolated, hermetic_set_disks_isolated,
|
||||
};
|
||||
use crate::storage_api_contracts::bucket::{BucketOperations as _, DeleteBucketOptions, MakeBucketOptions};
|
||||
use crate::storage_api_contracts::heal::HealOperations as _;
|
||||
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _};
|
||||
@@ -4183,6 +4306,89 @@ mod heal_result_report_tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn deep_heal_rebuilds_h2_near_tail_truncated_part() {
|
||||
let (temp_dirs, disks, set) = hermetic_set_disks_for_pool_with_default_parity_isolated(16, 0, 4).await;
|
||||
let bucket = "deep-heal-h2-near-tail-truncation";
|
||||
let object = "object.bin";
|
||||
for disk in &disks {
|
||||
disk.make_volume(bucket).await.expect("bucket volume should be created");
|
||||
}
|
||||
|
||||
let expected_payload = vec![0x6d; 5 * 1024 * 1024 + 123];
|
||||
set.put_object(
|
||||
bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(expected_payload.clone()),
|
||||
&ObjectOptions {
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("source object should be written before shard truncation");
|
||||
let source = disks[2]
|
||||
.read_version("", bucket, object, "", &ReadOptions::default())
|
||||
.await
|
||||
.expect("source metadata should be readable");
|
||||
let data_dir = source.data_dir.expect("non-inline source should have a data directory");
|
||||
let truncated_part = temp_dirs[1]
|
||||
.path()
|
||||
.join(bucket)
|
||||
.join(object)
|
||||
.join(data_dir.to_string())
|
||||
.join("part.1");
|
||||
let original_len = tokio::fs::metadata(&truncated_part)
|
||||
.await
|
||||
.expect("target shard should exist")
|
||||
.len();
|
||||
assert!(original_len > 1, "test shard must be large enough to truncate");
|
||||
let file = tokio::fs::OpenOptions::new()
|
||||
.write(true)
|
||||
.open(&truncated_part)
|
||||
.await
|
||||
.expect("target shard should be writable");
|
||||
file.set_len(original_len - 1)
|
||||
.await
|
||||
.expect("target shard should be truncated");
|
||||
|
||||
let mut reader = set
|
||||
.get_object_reader(bucket, object, None, Default::default(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("GET should remain readable after a one-byte shard truncation");
|
||||
let mut read_back = Vec::new();
|
||||
tokio::io::copy(&mut reader, &mut read_back)
|
||||
.await
|
||||
.expect("GET should reconstruct the truncated shard through EC");
|
||||
assert_eq!(read_back, expected_payload, "EC GET must preserve the object bytes");
|
||||
|
||||
let (result, error) = set
|
||||
.heal_object(
|
||||
bucket,
|
||||
object,
|
||||
"",
|
||||
&HealOpts {
|
||||
no_lock: true,
|
||||
scan_mode: HealScanMode::Deep,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("deep heal should finish after a one-byte H2 tail truncation");
|
||||
|
||||
assert!(error.is_none(), "deep heal should recover the H2 truncated shard: {error:?}");
|
||||
assert_eq!(result.drives_healed(), Some(1));
|
||||
assert_eq!(result.before.drives[1].state, DriveState::Corrupt.to_string());
|
||||
assert_eq!(
|
||||
tokio::fs::metadata(&truncated_part)
|
||||
.await
|
||||
.expect("repaired H2 shard should exist")
|
||||
.len(),
|
||||
original_len,
|
||||
"deep heal must restore the complete H2 shard length"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn deep_heal_inline_bitrot_rebuilds_physical_data_and_parity_shards() {
|
||||
use crate::storage_api_contracts::range::HTTPRangeSpec;
|
||||
|
||||
@@ -7703,10 +7703,11 @@ impl SetDisks {
|
||||
drop(namespace_owner);
|
||||
if quorum_result.is_ok()
|
||||
&& errs.iter().any(Option::is_some)
|
||||
&& let Some(purge) = delete_marker_purge
|
||||
&& let Some(purge) = delete_marker_purge.as_ref()
|
||||
&& let Some(version) = fi.version_id.filter(|version| !version.is_nil())
|
||||
{
|
||||
let _ = self.persist_delete_marker_purge(bucket, object, version, purge).await;
|
||||
let _ = self.persist_marker_purge_receipt(bucket, object, version, purge).await;
|
||||
let _ = self.persist_delete_marker_purge(bucket, object, version, purge.clone()).await;
|
||||
}
|
||||
// An explicit purge can carry deleted=true for the existing marker.
|
||||
// It must not create a repair intent that could reintroduce that marker.
|
||||
|
||||
@@ -9124,6 +9124,100 @@ mod tests {
|
||||
assert_eq!(dirs.len(), 16);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn current_marker_delete_ack_removes_stale_rejoined_copy() {
|
||||
let bucket = "current-marker-stale-rejoin";
|
||||
let object = "marker.bin";
|
||||
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
|
||||
let (_dirs, original) = make_local_set_disks_with_ctx(16, 4, ctx.clone()).await;
|
||||
let store = Arc::new(new_prepared_reader_test_store_with_ctx(&[original], ctx).await);
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
store
|
||||
.handle_make_bucket(
|
||||
bucket,
|
||||
&MakeBucketOptions {
|
||||
versioning_enabled: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("create versioned bucket");
|
||||
let marker = store
|
||||
.delete_object(
|
||||
bucket,
|
||||
object,
|
||||
ObjectOptions {
|
||||
versioned: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("create current marker")
|
||||
.version_id
|
||||
.expect("marker has explicit version");
|
||||
let current = store
|
||||
.bucket_incarnation_id_from_disk(bucket)
|
||||
.await
|
||||
.expect("bucket incarnation");
|
||||
let set = &store.pools[0].disk_set[0];
|
||||
let before = set.disks.read().await.clone();
|
||||
for slot in 12..16 {
|
||||
set.disks.write().await[slot] = None;
|
||||
}
|
||||
store
|
||||
.delete_object(
|
||||
bucket,
|
||||
object,
|
||||
ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(marker.to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("acknowledged delete should remove the marker from the online quorum");
|
||||
*set.disks.write().await = before.clone();
|
||||
|
||||
for disk in before.iter().take(12).flatten() {
|
||||
assert!(matches!(
|
||||
disk.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
|
||||
.await,
|
||||
Err(crate::disk::error::DiskError::FileVersionNotFound | crate::disk::error::DiskError::FileNotFound)
|
||||
));
|
||||
}
|
||||
for disk in before.iter().skip(12).flatten() {
|
||||
assert!(
|
||||
disk.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
|
||||
.await
|
||||
.is_ok(),
|
||||
"offline member must retain the stale current marker"
|
||||
);
|
||||
}
|
||||
let healed = store
|
||||
.heal_object_at_incarnation(
|
||||
bucket,
|
||||
object,
|
||||
&marker.to_string(),
|
||||
current,
|
||||
&rustfs_heal_contracts::heal_channel::HealOpts {
|
||||
remove: true,
|
||||
scan_mode: rustfs_heal_contracts::heal_channel::HealScanMode::Deep,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("stale current marker heal should complete");
|
||||
assert!(healed.error.is_none(), "stale current marker heal: {:?}", healed.error);
|
||||
assert!(healed.absence.as_ref().is_some_and(|proof| proof.removed));
|
||||
for disk in before.iter().flatten() {
|
||||
assert!(matches!(
|
||||
disk.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
|
||||
.await,
|
||||
Err(crate::disk::error::DiskError::FileVersionNotFound | crate::disk::error::DiskError::FileNotFound)
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn multipool_delete_marker_stays_with_existing_versions() {
|
||||
let bucket = "multipool-marker-routing";
|
||||
|
||||
Reference in New Issue
Block a user