fix(ecstore): fence free-version decommission replay

This commit is contained in:
overtrue
2026-08-23 10:34:11 +08:00
parent 5360b87459
commit dfbe165882
6 changed files with 199 additions and 71 deletions
+7 -8
View File
@@ -3911,6 +3911,7 @@ impl ECStore {
if migrated {
decommissioned += 1;
cleanup_preflight_allowed_missing.push(data_movement::source_cleanup_version_identity(version));
free_version_disposition.record_migrated();
} else {
free_version_disposition.record_retained();
@@ -6530,15 +6531,13 @@ mod pools_tests {
use super::resolve_decommission_listing_error;
use super::{
DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP, DECOMMISSION_ENTRY_QUEUE_HARD_CAP,
DECOMMISSION_FREE_VERSION_MIGRATED_REASON, DECOMMISSION_FREE_VERSION_RETAINED_REASON,
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler,
DecommissionEntryEnqueueResult, DecommissionFreeVersionDisposition, DecommissionStartPoolState,
DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus,
QueuedDecommissionEntry, apply_decommission_status_space_info, await_decommission_worker, bind_decommission_cancelers,
bind_missing_decommission_cancelers, cancel_decommission_canceler, clamp_decommission_entry_concurrency,
classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result,
decommission_entry_queue_capacity, decommission_item_size, decommission_meta_bucket_options,
decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback,
PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info,
await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers,
cancel_decommission_canceler, clamp_decommission_entry_concurrency, classify_decommission_terminal_state,
count_decommission_item, decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size,
decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
default_decommission_entry_concurrency, drain_decommission_entry_queue, enqueue_decommission_entry,
ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_generation,
ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed,
+14
View File
@@ -1950,6 +1950,20 @@ mod tests {
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
}
#[test]
fn test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
let mut free_version = cleanup_test_file_info("object.txt", Uuid::from_u128(2), "tier-cleanup");
free_version.deleted = true;
free_version.set_tier_free_version();
let mut expected = cleanup_test_versions(vec![migrated.clone()]);
expected.free_versions = vec![free_version.clone()];
let current = cleanup_test_versions(vec![migrated]);
let allowed_missing = vec![source_cleanup_version_identity(&free_version)];
assert!(source_cleanup_versions_match_with_allowed_missing(&expected, &current, &allowed_missing));
}
#[test]
fn test_decommission_cleanup_preflight_rejects_unexpected_missing_version() {
let migrated = cleanup_test_file_info("object.txt", Uuid::from_u128(1), "migrated");
+126 -23
View File
@@ -227,15 +227,15 @@ pub(super) fn restore_commit_operation_id_from_metadata(metadata: &HashMap<Strin
restore_operation_id_from_metadata(metadata)
}
async fn check_decommission_tier_free_version_target(
async fn inspect_decommission_tier_free_version_target(
disk: &DiskStore,
bucket: &str,
object: &str,
source: &FileInfo,
) -> Result<()> {
) -> Result<bool> {
let raw = match disk.read_xl(bucket, object, false).await {
Ok(raw) => raw,
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(()),
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound) => return Ok(false),
Err(err) => return Err(err.into()),
};
let meta = FileMeta::load(&raw.buf)?;
@@ -253,8 +253,11 @@ async fn check_decommission_tier_free_version_target(
all_matching_versions_equivalent = false;
}
}
if matching_count == 0 || (matching_count == 1 && all_matching_versions_equivalent) {
return Ok(());
if matching_count == 0 {
return Ok(false);
}
if matching_count == 1 && all_matching_versions_equivalent {
return Ok(true);
}
Err(StorageError::DataMovementOverwriteErr(
@@ -265,6 +268,28 @@ async fn check_decommission_tier_free_version_target(
.into())
}
fn ensure_decommission_tier_free_version_commit_fence(bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> {
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
|| opts
.bucket_lifecycle_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
Ok(())
}
impl SetDisks {
pub(super) async fn require_current_restore_operation_id(
&self,
@@ -4718,26 +4743,11 @@ impl SetDisks {
if !fi.deleted || !fi.tier_free_version() {
return Err(Error::other("decommission tier free-version write requires a free version record"));
}
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
|| opts
.bucket_lifecycle_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
self.validate_decommission_tier_free_version_target(bucket, object, fi)
.await?;
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let write_quorum = self.default_write_quorum();
@@ -4760,6 +4770,8 @@ impl SetDisks {
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object)
}
@@ -4776,13 +4788,38 @@ impl SetDisks {
let preflight = disks
.iter()
.flatten()
.map(|disk| check_decommission_tier_free_version_target(disk, bucket, object, fi));
.map(|disk| inspect_decommission_tier_free_version_target(disk, bucket, object, fi));
for result in join_all(preflight).await {
result?;
}
Ok(())
}
pub(crate) async fn has_decommission_tier_free_version_write_quorum(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<bool> {
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
let disks = self.disks.read().await.clone();
let preflight = disks.iter().map(|disk| async {
match disk {
Some(disk) => inspect_decommission_tier_free_version_target(disk, bucket, object, fi).await,
None => Ok(false),
}
});
let mut equivalent = 0;
for result in join_all(preflight).await {
if result? {
equivalent += 1;
}
}
ensure_decommission_tier_free_version_commit_fence(bucket, object, opts)?;
Ok(equivalent >= self.default_write_quorum())
}
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tiered_object(
&self,
@@ -10056,6 +10093,72 @@ mod tests {
assert_eq!(migrated.transitioned_objname, "remote/object");
}
#[tokio::test]
async fn decommission_tier_free_version_resume_requires_write_quorum() {
let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await;
let bucket = "free-version-decommission-resume";
let object = "object.txt";
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(Uuid::new_v4()),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
let opts = ObjectOptions::default();
let disks = set_disks.get_disks_internal().await;
for disk in disks.iter().take(2).flatten() {
disk.write_metadata("", bucket, object, free_version.clone())
.await
.expect("partial first attempt should leave equivalent metadata");
}
assert!(
!set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("partial target metadata should remain valid"),
"write-quorum-minus-one must not be accepted as an idempotent migration"
);
disks[2]
.as_ref()
.expect("third target disk should be online")
.write_metadata("", bucket, object, free_version.clone())
.await
.expect("third equivalent target write should complete quorum");
assert!(
set_disks
.has_decommission_tier_free_version_write_quorum(bucket, object, &free_version, &opts)
.await
.expect("write-quorum target metadata should remain valid")
);
}
#[test]
fn decommission_tier_free_version_commit_rejects_lost_fence() {
let opts = ObjectOptions {
namespace_lock_fence: Some(NamespaceLockFence::lost_for_test()),
..Default::default()
};
let err = ensure_decommission_tier_free_version_commit_fence("bucket", "object", &opts)
.expect_err("lost target lock must fail the free-version commit");
assert!(matches!(
err,
Error::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
required: 1,
achieved: 0,
..
}
));
}
#[test]
fn test_resolve_tiered_decommission_write_quorum_result_allows_successful_quorum() {
let errs = vec![None, None, Some(DiskError::DiskNotFound), None];
+43 -6
View File
@@ -3934,6 +3934,18 @@ mod tests {
.await
.expect("create free-version decommission bucket");
let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await;
let source_free = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(&bucket, object)
.await
.expect("source free-version metadata should decode")
.and_then(|versions| {
versions
.versions
.into_iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
})
.expect("source free-version identity should be present before decommission");
let mut target_reader = PutObjReader::from_vec(b"target ordinary bytes".to_vec());
let target_w = store.pools[1]
.put_object(
@@ -3988,6 +4000,15 @@ mod tests {
.iter()
.any(|version| { version.version_id == Some(free_version) && version.tier_free_version() })
);
let migrated_free = target_versions
.versions
.iter()
.find(|version| version.version_id == Some(free_version) && version.tier_free_version())
.expect("migrated free-version identity should remain readable from target disks");
assert!(
crate::store::tiered_data_movement_source_matches(&source_free, migrated_free)
.expect("migrated free-version identity should decode")
);
let retained_w = target_versions
.versions
.iter()
@@ -4010,6 +4031,22 @@ mod tests {
.is_none(),
"successful free-version migration should permit source cleanup"
);
let (heal_versions, _, _) = store
.heal_walk_versions_page(1, 0, &bucket, "", None, 2, 16, true)
.await
.expect("heal walk should decode the migrated free version");
let free_version_string = free_version.to_string();
let healed_free = heal_versions
.iter()
.find(|version| version.version_id.as_deref() == Some(free_version_string.as_str()))
.expect("heal walk should surface the migrated free version");
let healed_info = healed_free
.lifecycle_object_info
.as_ref()
.expect("heal walk should retain lifecycle identity for the migrated free version");
assert!(healed_info.transitioned_object.free_version);
assert_eq!(healed_info.transitioned_object.tier, source_free.transition_tier);
assert_eq!(healed_info.transitioned_object.name, source_free.transitioned_objname);
shutdown.cancel();
}
@@ -4132,11 +4169,11 @@ mod tests {
.filter(|version| version.header.version_id == Some(free_version))
.collect::<Vec<_>>();
assert_eq!(same_id.len(), 2, "conflict metadata should retain both same-ID records");
assert_eq!(same_id.iter().filter(|version| version.free_version()).count(), 1);
assert_eq!(same_id.iter().filter(|version| !version.free_version()).count(), 1);
assert_eq!(same_id.iter().filter(|version| version.header.free_version()).count(), 1);
assert_eq!(same_id.iter().filter(|version| !version.header.free_version()).count(), 1);
let post_conflict = same_id
.into_iter()
.find(|version| !version.free_version())
.find(|version| !version.header.free_version())
.expect("ordinary conflict version must remain addressable by the source ID");
let post_conflict_info = post_conflict
.into_fileinfo(&bucket, object, true)
@@ -4163,11 +4200,11 @@ mod tests {
.collect::<Vec<_>>();
if disk_index == 0 {
assert_eq!(same_id.len(), 2);
assert!(same_id[0].free_version());
assert!(!same_id[1].free_version());
assert!(same_id[0].header.free_version());
assert!(!same_id[1].header.free_version());
} else {
assert_eq!(same_id.len(), 1);
assert!(same_id[0].free_version());
assert!(same_id[0].header.free_version());
}
}
let sweep_err = store
+6 -34
View File
@@ -2237,36 +2237,16 @@ impl ECStore {
bucket: &str,
object: &str,
source: &rustfs_filemeta::FileInfo,
opts: &ObjectOptions,
target_pool_idx: usize,
) -> Result<bool> {
let pool = self
.pools
.get(target_pool_idx)
.ok_or_else(|| Error::other(format!("invalid tiered data movement target pool {target_pool_idx}")))?;
let logical_object = decode_dir_object(object);
let Some(versions) = pool
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, &logical_object)
.await?
else {
return Ok(false);
};
let Some(target) = versions
.versions
.iter()
.find(|version| version.version_id == source.version_id)
else {
return Ok(false);
};
if !target.tier_free_version() {
return Err(decommission_free_version_overwrite_error(bucket, object, source.version_id));
}
if tiered_data_movement_source_matches(source, target)? {
Ok(true)
} else {
Err(decommission_free_version_overwrite_error(bucket, object, source.version_id))
}
pool.get_disks_by_key(object)
.has_decommission_tier_free_version_write_quorum(bucket, object, source, opts)
.await
}
fn resolve_decommission_target_pool_idx_result(result: Result<usize>, bucket: &str, object: &str) -> Result<usize> {
@@ -2365,11 +2345,7 @@ impl ECStore {
return Err(Error::DiskFull);
}
let equivalent = if is_free_version {
self.pools[target_pool_idx]
.get_disks_by_key(&object)
.validate_decommission_tier_free_version_target(bucket, &object, &fi)
.await?;
self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, target_pool_idx)
self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, target_pool_idx)
.await?
} else {
self.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
@@ -2383,12 +2359,8 @@ impl ECStore {
}
let result = if is_free_version {
self.pools[idx]
.get_disks_by_key(&object)
.validate_decommission_tier_free_version_target(bucket, &object, &fi)
.await?;
if self
.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, idx)
.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, &opts, idx)
.await?
{
return Ok(());
@@ -235,6 +235,9 @@ report the enclosing object migration result.
Regression guard:
- `decommission_tier_free_version_preserves_remote_identity`
- `decommission_tier_free_version_resume_requires_write_quorum`
- `decommission_tier_free_version_commit_rejects_lost_fence`
- `test_decommission_cleanup_preflight_accepts_migrated_free_version_consumed_from_source`
- `decommission_entry_skips_cleanup_only_marker_when_free_version_is_present`
- `decommission_entry_rejects_subquorum_free_version_conflict_and_retains_source`