From 1f02d5a73ab26e379b8e042e95adf60313e56010 Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 2 Sep 2026 21:53:41 +0800 Subject: [PATCH] fix(ecstore): settle checkpoint capacity retries --- .../bucket/lifecycle/tier_delete_journal.rs | 2 +- crates/ecstore/src/core/pools.rs | 145 +++++++++++++----- crates/ecstore/src/store/init.rs | 44 ++++++ 3 files changed, 153 insertions(+), 38 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs index 2f519ab55..3840199c2 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs @@ -1468,7 +1468,7 @@ async fn save_decommission_manifest_checkpoint_if_match( } let write_data = next_data.to_vec(); let write = api - .run_decommission_capacity_temporary_mutation_with_capacity_lease( + .run_decommission_capacity_non_growing_replacement_with_capacity_lease( target.target_pool_index, Some(target.capacity_owner), Some(next_data.len()), diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 6ddd53bed..8e2c54bee 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -1939,6 +1939,41 @@ fn release_decommission_target_temporary_mutation( (released, released > 0 || removed) } +fn settle_decommission_target_non_growing_replacement( + meta: &mut PoolMeta, + source_pool_index: usize, + target_pool_index: usize, + mutation_id: uuid::Uuid, + now: OffsetDateTime, +) -> Result { + let pool_count = meta.pools.len(); + let pool = meta + .pools + .get_mut(source_pool_index) + .ok_or_else(|| invalid_decommission_pool_index_error(pool_count, source_pool_index))?; + let info = pool + .decommission + .as_mut() + .ok_or_else(|| decommission_metadata_not_initialized_error("settle non-growing target replacement"))?; + let reservation = info + .capacity_reservation + .as_mut() + .filter(|reservation| reservation.active()) + .ok_or_else(|| decommission_capacity_blocked_error("active reservation disappeared during target replacement"))?; + let target = reservation + .targets + .iter_mut() + .find(|target| target.pool_index == target_pool_index) + .ok_or_else(|| decommission_capacity_blocked_error("target allocation disappeared during target replacement"))?; + let (released, changed) = release_decommission_target_temporary_mutation(target, mutation_id, usize::MAX); + reservation.inflight_target_physical_bytes = reservation.inflight_target_physical_bytes.saturating_sub(released); + if changed { + renew_decommission_capacity_reservation(reservation, now, true); + pool.last_update = now; + } + Ok(changed) +} + struct DecommissionTargetConsumption { committed_data_bytes: usize, target_physical_bytes: usize, @@ -7394,6 +7429,7 @@ pub struct DecommissionCapacityReservation { enum DecommissionCapacityMutationMode { Durable, Temporary, + NonGrowingReplacement, TemporaryRelease, } @@ -9536,11 +9572,37 @@ impl ECStore { .await } + /// Run an identity-preserving replacement that the caller has already + /// proven cannot grow the target object. A failed or ambiguous write keeps + /// the ordinary temporary-mutation recovery marker, while a successful + /// write resolves the capacity intent without retaining MPU cleanup state. + pub(crate) async fn run_decommission_capacity_non_growing_replacement_with_capacity_lease( + &self, + target_pool_index: usize, + capacity_owner: Option, + expected_data_bytes: Option, + operation: F, + ) -> Result + where + F: FnOnce(Option>) -> Fut, + Fut: std::future::Future>, + { + self.run_decommission_capacity_mutation( + target_pool_index, + capacity_owner, + expected_data_bytes, + DecommissionCapacityMutationMode::NonGrowingReplacement, + |_| false, + operation, + ) + .await + } + /// Finish the capacity transaction for an identity-preserving temporary /// replacement whose target bytes are already durably present. This is - /// the crash-recovery half of `run_decommission_capacity_temporary_mutation`: - /// it never writes the target again and only resolves a pending intent - /// owned by the exact deterministic mutation id. + /// the crash-recovery half of a non-growing replacement: it never writes + /// the target again and only settles state owned by the exact deterministic + /// mutation id. pub(crate) async fn reconcile_decommission_capacity_after_equivalent_temporary_target( &self, owner: DecommissionCapacityOwner, @@ -9565,7 +9627,7 @@ impl ECStore { let mutation_id = owner .mutation_id .ok_or_else(|| decommission_capacity_blocked_error("equivalent temporary target mutation identity is missing"))?; - let (target_layout, pending_physical_bytes, pending_mutation_id, already_reconciled) = { + let (target_layout, pending_physical_bytes, pending_mutation_id, has_temporary_mutation) = { let reservation = snapshot .pools .get(source_pool_index) @@ -9588,9 +9650,8 @@ impl ECStore { .any(|mutation| mutation.mutation_id == mutation_id), ) }; - if pending_physical_bytes == 0 { - // Either the successful attempt already saved its progress, or a - // byte-neutral replacement had no inflight delta to record. + if pending_physical_bytes == 0 && !has_temporary_mutation { + // A prior successful attempt already settled both durable halves. ensure_pool_meta_write_fence(&pool_meta_guard, "equivalent temporary target reconciliation fence failed")?; if let Some(target_guard) = target_guard.as_ref() { ensure_decommission_capacity_target_fence( @@ -9601,38 +9662,34 @@ impl ECStore { } return Ok(()); } - if pending_mutation_id != Some(mutation_id) { + if pending_physical_bytes > 0 && pending_mutation_id != Some(mutation_id) { return Err(decommission_capacity_blocked_error( "equivalent temporary target pending intent belongs to another mutation", )); } - if already_reconciled { - return Err(decommission_capacity_blocked_error( - "equivalent temporary target has both pending and reconciled state", - )); - } let expected_target_physical_bytes = capacity_target_physical_bytes(expected_data_bytes.max(1), target_layout)?; - if pending_physical_bytes < expected_target_physical_bytes { + if pending_physical_bytes > 0 && pending_physical_bytes != expected_target_physical_bytes { return Err(decommission_capacity_blocked_error( - "equivalent temporary target pending capacity is smaller than the committed checkpoint", + "equivalent temporary target pending capacity does not match the committed checkpoint", )); } - resolve_decommission_target_pending( + if pending_physical_bytes > 0 { + resolve_decommission_target_pending( + &mut snapshot, + source_pool_index, + target_pool_index, + expected_target_physical_bytes, + mutation_id, + )?; + } + // Exact byte equivalence proves that this identity-preserving, + // byte-non-growing replacement committed. Settle both crash windows: + // the pending intent before finalize and the temporary marker written + // when the operation returned an error after committing its target. + settle_decommission_target_non_growing_replacement( &mut snapshot, source_pool_index, target_pool_index, - expected_target_physical_bytes, - mutation_id, - )?; - // The replacement is byte-non-growing and its exact bytes were read - // before this call, so no new physical delta is inferred on replay. - // A prior successful progress save would have taken the idempotent - // pending==0 return above. - record_decommission_target_inflight( - &mut snapshot, - source_pool_index, - target_pool_index, - 0, mutation_id, OffsetDateTime::now_utc(), )?; @@ -9773,7 +9830,11 @@ impl ECStore { Fut: std::future::Future>, P: Fn(&T) -> bool, { - let temporary = matches!(mode, DecommissionCapacityMutationMode::Temporary); + let temporary = matches!( + mode, + DecommissionCapacityMutationMode::Temporary | DecommissionCapacityMutationMode::NonGrowingReplacement + ); + let non_growing_replacement = matches!(mode, DecommissionCapacityMutationMode::NonGrowingReplacement); let temporary_release = matches!(mode, DecommissionCapacityMutationMode::TemporaryRelease); let mut operation = Some(operation); let mut save_guard = self.pool_meta_save_gate.lock().await; @@ -10051,14 +10112,24 @@ impl ECStore { mutation_id, )?; if temporary { - record_decommission_target_inflight( - &mut snapshot, - source_pool_index, - target_pool_index, - observed_physical_bytes, - mutation_id, - now, - )?; + if non_growing_replacement { + settle_decommission_target_non_growing_replacement( + &mut snapshot, + source_pool_index, + target_pool_index, + mutation_id, + now, + )?; + } else { + record_decommission_target_inflight( + &mut snapshot, + source_pool_index, + target_pool_index, + observed_physical_bytes, + mutation_id, + now, + )?; + } } else { let consumed_physical_bytes = expected_data_bytes .map(|_| expected_target_physical_bytes) diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 5d093e39b..45f910a50 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -10072,6 +10072,29 @@ mod tests { .expect("the active decommission should own the checkpoint target"); assert_eq!(targets.len(), 1); let target = targets[0].clone(); + let consumed_before_checkpoint = store.pool_meta.read().await.pools[0] + .decommission + .as_ref() + .and_then(|info| info.capacity_reservation.as_ref()) + .expect("checkpoint shutdown capacity reservation should exist") + .consumed_target_physical_bytes; + let precommit_error = store + .run_decommission_capacity_non_growing_replacement_with_capacity_lease( + target.target_pool_index, + Some(target.capacity_owner), + Some(aborting_data.len()), + |_| async { Err::<(), Error>(Error::other("injected checkpoint failure before commit")) }, + ) + .await + .expect_err("the first checkpoint attempt should fail before writing its target") + .to_string(); + assert!(precommit_error.contains("injected checkpoint failure before commit")); + assert!( + store + .has_decommission_capacity_temporary_mutation_state(target.target_pool_index, target.capacity_owner) + .await, + "a failed checkpoint attempt must retain exact retry state" + ); let barrier = crate::set_disk::PutObjectCommitBarrier::install( RUSTFS_META_BUCKET, &manifest_name, @@ -10122,6 +10145,27 @@ mod tests { .await, "the admitted target PUT must drain its capacity transaction before releasing the recovery fences" ); + { + let pool_meta = store.pool_meta.read().await; + let reservation = pool_meta.pools[0] + .decommission + .as_ref() + .and_then(|info| info.capacity_reservation.as_ref()) + .expect("checkpoint shutdown capacity reservation should remain active"); + let capacity_target = reservation + .targets + .iter() + .find(|candidate| candidate.pool_index == target.target_pool_index) + .expect("checkpoint shutdown capacity target should remain allocated"); + assert_eq!( + reservation.consumed_target_physical_bytes, consumed_before_checkpoint, + "a non-growing checkpoint replacement must not consume durable migration capacity" + ); + assert_eq!(reservation.pending_target_physical_bytes, 0); + assert_eq!(reservation.inflight_target_physical_bytes, 0); + assert_eq!(capacity_target.pending_physical_bytes, 0); + assert!(capacity_target.temporary_mutations.is_empty()); + } drop(barrier); let mut reloaded_pool_meta = PoolMeta::default(); reloaded_pool_meta