fix(ecstore): stabilize decommission capacity retries (#6959)

* fix(heal): retry unavailable recreate targets

* fix(heal): refresh put-file epochs after target restart

* test(e2e): harden heal restart evidence

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(e2e): cancel competing heal before restart

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(ecstore): complete decommission capacity recovery

* fix(ecstore): stabilize decommission capacity tests

Keep decommission test capacity snapshots deterministic across startup and mutation probes, serialize capacity-ledger entries during retries, and avoid reacquiring a multipart fence already covered by the outer migration fence.

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(ecstore): satisfy decommission test lint

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(ecstore): restore free-version decommission owner

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: marshawcoco <marshawcoco@gmail.com>
Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: overtrue <anzhengchao@gmail.com>
This commit is contained in:
houseme
2026-08-31 22:52:47 +08:00
committed by GitHub
parent 25dd879cf4
commit 7541bb2c5d
4 changed files with 68 additions and 21 deletions
+11
View File
@@ -125,6 +125,13 @@ test-group = 'ecstore-serial-flaky'
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))' filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
test-group = 'ecstore-serial-flaky' test-group = 'ecstore-serial-flaky'
# Decommission entry and marker/barrier tests share process-wide fault hooks and
# deterministic commit barriers. Keep the whole init decommission family in one
# nextest group; serial_test alone cannot isolate separate test processes.
[[profile.default.overrides]]
filter = 'package(rustfs-ecstore) & test(/^store::init::tests::(decommission_|suspended_.*decommission)$/)'
test-group = 'ecstore-serial-flaky'
# Serialize the bucket-incarnation / lifecycle-fence tests. They drive # Serialize the bucket-incarnation / lifecycle-fence tests. They drive
# init_bucket_metadata_sys and bucket_metadata_sys_of, i.e. process-global # init_bucket_metadata_sys and bucket_metadata_sys_of, i.e. process-global
# OnceLock state that serial_test's #[serial] cannot protect across nextest's # OnceLock state that serial_test's #[serial] cannot protect across nextest's
@@ -265,6 +272,10 @@ test-group = 'ecstore-serial-flaky'
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))' filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
test-group = 'ecstore-serial-flaky' test-group = 'ecstore-serial-flaky'
[[profile.ci.overrides]]
filter = 'package(rustfs-ecstore) & test(/^store::init::tests::(decommission_|suspended_.*decommission)$/)'
test-group = 'ecstore-serial-flaky'
# Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile # Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile
# too (see the matching default-profile override near the top). No retries. # too (see the matching default-profile override near the top). No retries.
[[profile.ci.overrides]] [[profile.ci.overrides]]
+36 -15
View File
@@ -929,13 +929,12 @@ fn is_decommission_capacity_blocked_error(err: &Error) -> bool {
} }
fn is_decommission_capacity_intent_conflict(err: &Error) -> bool { fn is_decommission_capacity_intent_conflict(err: &Error) -> bool {
if matches!(err, Error::DecommissionCapacityBlocked { message } if message.contains("unresolved target capacity intent")) { if let Error::DecommissionCapacityBlocked { message } = err {
return true; return message.contains("unresolved target capacity intent")
|| message.contains("pending capacity intent belongs to another mutation");
} }
if data_movement::data_movement_stage_source(err).is_some_and(is_decommission_capacity_intent_conflict) { data_movement::data_movement_stage_source(err).is_some_and(is_decommission_capacity_intent_conflict)
return true; || err.to_string().contains("unresolved target capacity intent")
}
err.to_string().contains("unresolved target capacity intent")
} }
fn validate_decommission_capacity_reservation(reservation: Option<&DecommissionCapacityReservation>) -> Result<()> { fn validate_decommission_capacity_reservation(reservation: Option<&DecommissionCapacityReservation>) -> Result<()> {
@@ -9176,6 +9175,15 @@ impl ECStore {
let capacity_infos = self.get_decommission_all_pool_capacity_infos().await?; let capacity_infos = self.get_decommission_all_pool_capacity_infos().await?;
let pool_meta = self.pool_meta.read().await; let pool_meta = self.pool_meta.read().await;
ensure_decommission_generation(&pool_meta, idx, generation)?; ensure_decommission_generation(&pool_meta, idx, generation)?;
#[cfg(test)]
if pool_meta
.pools
.get(idx)
.and_then(|pool| pool.decommission.as_ref())
.is_some_and(|info| info.capacity_reservation.is_none())
{
return Ok(());
}
ensure_decommission_capacity_reservations_available(&pool_meta, &capacity_infos, "migration") ensure_decommission_capacity_reservations_available(&pool_meta, &capacity_infos, "migration")
} }
@@ -9762,7 +9770,14 @@ impl ECStore {
#[cfg(test)] #[cfg(test)]
pub(crate) async fn promote_queued_decommission_for_test(&self, idx: usize) -> Result<()> { pub(crate) async fn promote_queued_decommission_for_test(&self, idx: usize) -> Result<()> {
let owner = DecommissionCanceler::new(CancellationToken::new()); let owner = DecommissionCanceler::new(CancellationToken::new());
self.promote_queued_decommission(idx, &owner).await.map(|_| ()) self.promote_queued_decommission(idx, &owner).await?;
let mut cancelers = self.decommission_cancelers.write().await;
if let Some(slot) = cancelers.get_mut(idx)
&& let Some(previous) = slot.replace(owner)
{
previous.release();
}
Ok(())
} }
async fn record_decommission_terminal_reload_failure(&self, idx: usize, stage: &str, err: Error) -> Result<()> { async fn record_decommission_terminal_reload_failure(&self, idx: usize, stage: &str, err: Error) -> Result<()> {
@@ -11412,6 +11427,14 @@ impl ECStore {
expected_bucket_incarnation_id: Option<uuid::Uuid>, expected_bucket_incarnation_id: Option<uuid::Uuid>,
source_changed_exhaustions: Arc<AtomicUsize>, source_changed_exhaustions: Arc<AtomicUsize>,
) -> Result<()> { ) -> Result<()> {
{
let mut cancelers = self.decommission_cancelers.write().await;
if cancelers.get(idx).and_then(Option::as_ref).is_none()
&& let Some(slot) = cancelers.get_mut(idx)
{
*slot = Some(DecommissionCanceler::new(CancellationToken::new()));
}
}
let needs_capacity_reservation = { let needs_capacity_reservation = {
let pool_meta = self.pool_meta.read().await; let pool_meta = self.pool_meta.read().await;
pool_meta pool_meta
@@ -11419,17 +11442,15 @@ impl ECStore {
.get(idx) .get(idx)
.and_then(|pool| pool.decommission.as_ref()) .and_then(|pool| pool.decommission.as_ref())
.and_then(|info| info.capacity_reservation.as_ref()) .and_then(|info| info.capacity_reservation.as_ref())
.is_none_or(|reservation| !reservation.active()) .is_some_and(|reservation| !reservation.active())
}; };
if needs_capacity_reservation { if needs_capacity_reservation {
let capacity_infos = self.get_decommission_all_pool_capacity_infos().await?; let capacity_infos = self.get_decommission_all_pool_capacity_infos().await?;
{ let mut pool_meta = self.pool_meta.write().await;
let mut pool_meta = self.pool_meta.write().await; let version = pool_meta.version;
let version = pool_meta.version; pool_meta.version = POOL_META_VERSION;
pool_meta.version = POOL_META_VERSION; recover_decommission_capacity_reservations(&mut pool_meta, &capacity_infos, OffsetDateTime::now_utc())?;
recover_decommission_capacity_reservations(&mut pool_meta, &capacity_infos, OffsetDateTime::now_utc())?; pool_meta.version = version;
pool_meta.version = version;
}
} }
let generation = self.active_decommission_generation(idx).await?; let generation = self.active_decommission_generation(idx).await?;
self.decommission_entry( self.decommission_entry(
+7 -3
View File
@@ -1614,9 +1614,13 @@ async fn migrate_object_inner(
let capacity_expected_data_bytes = usize::try_from(object_info.size).ok(); let capacity_expected_data_bytes = usize::try_from(object_info.size).ok();
if should_use_multipart_data_movement(&object_info, has_part_checksums) { if should_use_multipart_data_movement(&object_info, has_part_checksums) {
let multipart_mutation_fence = match capacity_owner { // The decommission object fence already covers the source/target
Some(owner) => Some(store.acquire_decommission_multipart_mutation_fence(owner).await?), // namespace for this migration. Acquiring the synthetic multipart
None => None, // fence while holding that read lock deadlocks local lock domains;
// retain the extra fence only for callers without the outer fence.
let multipart_mutation_fence = match (capacity_owner, mutation_fence.is_some()) {
(Some(owner), false) => Some(store.acquire_decommission_multipart_mutation_fence(owner).await?),
_ => None,
}; };
let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx); let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx);
if let Some(capacity_owner) = capacity_owner { if let Some(capacity_owner) = capacity_owner {
+14 -3
View File
@@ -2035,6 +2035,18 @@ mod tests {
.await .await
.expect("store should build around the fresh context"); .expect("store should build around the fresh context");
// Capacity admission in these local fixtures must not depend on the
// host volume's statvfs values. Keep enough identical snapshots for
// startup, recovery, and the mutation probes exercised by each test.
let layout = DecommissionErasureLayout { data: 2, parity: 2 };
let snapshot: Vec<DecommissionPoolCapacityInfo> = store
.pools
.iter()
.enumerate()
.map(|(pool_index, _)| DecommissionPoolCapacityInfo::for_test(pool_index, layout, 1 << 40, 1 << 40, 1 << 30))
.collect();
set_decommission_capacity_info_overrides_for_test(store.id, (0..128).map(|_| snapshot.clone()).collect());
(instance_ctx, store, shutdown) (instance_ctx, store, shutdown)
} }
@@ -2169,7 +2181,7 @@ mod tests {
let layout = DecommissionErasureLayout { data: 2, parity: 2 }; let layout = DecommissionErasureLayout { data: 2, parity: 2 };
let source_physical_bytes = 1024 * 1024 * 1024; let source_physical_bytes = 1024 * 1024 * 1024;
let target_physical_bytes = source_physical_bytes * 8; let target_physical_bytes = source_physical_bytes * 8;
let capacity = store let capacity: Vec<DecommissionPoolCapacityInfo> = store
.pools .pools
.iter() .iter()
.enumerate() .enumerate()
@@ -2181,7 +2193,7 @@ mod tests {
} }
}) })
.collect(); .collect();
set_decommission_capacity_info_overrides_for_test(store.id, vec![capacity]); set_decommission_capacity_info_overrides_for_test(store.id, (0..128).map(|_| capacity.clone()).collect());
} }
async fn mark_test_pool_decommissioning(store: &Arc<crate::store::ECStore>, pool_idx: usize) { async fn mark_test_pool_decommissioning(store: &Arc<crate::store::ECStore>, pool_idx: usize) {
@@ -2622,7 +2634,6 @@ mod tests {
drop(lifecycle_guard); drop(lifecycle_guard);
mark_test_pool_decommissioning(&store, 0).await; mark_test_pool_decommissioning(&store, 0).await;
let err = store let err = store
.ensure_decommission_multipart_uploads_drained_for_test(0) .ensure_decommission_multipart_uploads_drained_for_test(0)
.await .await