diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 7a88dddf1..61cf84ec2 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -128,6 +128,7 @@ const METRIC_DECOMMISSION_CAPACITY_PREDICTION_ABSOLUTE_ERROR_BYTES: &str = const DECOMMISSION_LISTING_MAX_ATTEMPTS: usize = 3; const DECOMMISSION_LISTING_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); pub(crate) const DECOMMISSION_ENTRY_MAX_ATTEMPTS: usize = 3; +const DECOMMISSION_CAPACITY_INTENT_CONFLICT_MAX_ATTEMPTS: usize = 12; const DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(100); pub(crate) const DECOMMISSION_VERSION_COPY_ATTEMPTS: usize = 3; const DECOMMISSION_COPY_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(50); @@ -927,6 +928,16 @@ fn is_decommission_capacity_blocked_error(err: &Error) -> bool { data_movement::data_movement_stage_source(err).is_some_and(is_decommission_capacity_blocked_error) } +fn is_decommission_capacity_intent_conflict(err: &Error) -> bool { + if matches!(err, Error::DecommissionCapacityBlocked { message } if message.contains("unresolved target capacity intent")) { + return true; + } + if 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") +} + fn validate_decommission_capacity_reservation(reservation: Option<&DecommissionCapacityReservation>) -> Result<()> { let Some(reservation) = reservation else { return Ok(()); @@ -1608,9 +1619,9 @@ fn reserve_decommission_target_pending( pool.last_update = now; return Ok(additional); } - _ => { + pending_mutation_id => { return Err(decommission_capacity_blocked_error(format!( - "source pool {source_pool_index} target pool {target_pool_index} has an unresolved target capacity intent" + "source pool {source_pool_index} target pool {target_pool_index} has an unresolved target capacity intent {pending_mutation_id:?} while mutation {mutation_id} is waiting" ))); } } @@ -6977,7 +6988,7 @@ pub struct DecommissionCapacityTarget { pub inflight_physical_bytes: usize, #[serde(default)] pub pending_physical_bytes: usize, - #[serde(default, skip_serializing_if = "Option::is_none")] + #[serde(default)] pub pending_mutation_id: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] pub temporary_mutations: Vec, @@ -7789,6 +7800,19 @@ fn decommission_remote_tiered_opts( } } +fn decommission_capacity_version_mutation_id( + owner: DecommissionCapacityOwner, + bucket: &str, + version: &rustfs_filemeta::FileInfo, +) -> uuid::Uuid { + let version_id = if version.deleted && version.version_id.is_none() { + Some(uuid::Uuid::nil().to_string()) + } else { + version.version_id.map(|version_id| version_id.to_string()) + }; + decommission_capacity_mutation_id(owner, bucket, &version.name, version_id.as_deref(), version.deleted, version.mod_time) +} + fn decommission_capacity_owned_opts(mut opts: ObjectOptions, capacity_owner: Option) -> ObjectOptions { if let Some(capacity_owner) = capacity_owner { capacity_owner.apply_to(&mut opts); @@ -8100,6 +8124,18 @@ impl ECStore { Ok((pool_meta_guard, has_active_source)) } + pub(crate) async fn acquire_decommission_capacity_release_fence_with_active_source( + &self, + ) -> Result<(rustfs_lock::NamespaceLockGuard, bool)> { + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, snapshot) = self + .acquire_pool_meta_read_guard(&mut save_guard, "capacity release fence failed") + .await?; + let has_active_source = pool_meta_has_active_decommission(&snapshot); + drop(save_guard); + Ok((pool_meta_guard, has_active_source)) + } + pub(crate) async fn run_decommission_capacity_admitted_mutation( &self, target_pool_index: usize, @@ -9148,26 +9184,17 @@ impl ECStore { idx: usize, generation: OffsetDateTime, ) -> Result> { - let active_worker = self - .decommission_cancelers - .read() - .await - .get(idx) - .and_then(Option::as_ref) - .is_some_and(DecommissionCanceler::is_active); - if !active_worker { - return Ok(None); - } - let pool_meta = self.pool_meta.read().await; ensure_decommission_generation(&pool_meta, idx, generation)?; - let reservation = pool_meta + let Some(reservation) = pool_meta .pools .get(idx) .and_then(|pool| pool.decommission.as_ref()) .and_then(|info| info.capacity_reservation.as_ref()) .filter(|reservation| reservation.lease_active_at(OffsetDateTime::now_utc())) - .ok_or_else(|| decommission_capacity_blocked_error(format!("source pool {idx} has no active reservation")))?; + else { + return Ok(None); + }; Ok(Some(DecommissionCapacityOwner { source_pool_index: idx, operation_id: reservation.operation_id, @@ -10319,29 +10346,65 @@ impl ECStore { expected_bucket_incarnation_id: Option, source_changed_exhaustions: Arc, ) -> Result<()> { + let uses_capacity_ledger = self + .pool_meta + .read() + .await + .pools + .get(idx) + .and_then(|pool| pool.decommission.as_ref()) + .and_then(|info| info.capacity_reservation.as_ref()) + .is_some_and(DecommissionCapacityReservation::active); let mut counted_versions = HashSet::new(); for entry_attempt in 1..=DECOMMISSION_ENTRY_MAX_ATTEMPTS { - match self - .decommission_entry_attempt( - rx.clone(), - idx, - generation, - entry.clone(), - bucket.clone(), - Arc::clone(&set), - lifecycle_config.clone(), - object_lock_config.clone(), - replication_config.clone(), - expected_bucket_incarnation_id, - entry_attempt, - source_changed_exhaustions.as_ref(), - &mut counted_versions, - ) - .await? - { - DecommissionEntryAttemptOutcome::Complete => return Ok(()), - DecommissionEntryAttemptOutcome::SourceChanged => { + let attempt_result = { + let mut conflict_attempt = 0; + loop { + let result = { + let _capacity_entry_guard = if uses_capacity_ledger { + Some(tokio::select! { + biased; + _ = rx.cancelled() => return decommission_cancel_signal_result(true), + guard = self.decommission_capacity_entry_gate.lock() => guard, + }) + } else { + None + }; + self.decommission_entry_attempt( + rx.clone(), + idx, + generation, + entry.clone(), + bucket.clone(), + Arc::clone(&set), + lifecycle_config.clone(), + object_lock_config.clone(), + replication_config.clone(), + expected_bucket_incarnation_id, + entry_attempt, + source_changed_exhaustions.as_ref(), + &mut counted_versions, + ) + .await + }; + if result.as_ref().is_err_and(is_decommission_capacity_intent_conflict) + && conflict_attempt < DECOMMISSION_CAPACITY_INTENT_CONFLICT_MAX_ATTEMPTS + { + conflict_attempt += 1; + let retry_delay = + decommission_retry_backoff_delay(DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, conflict_attempt); + if wait_decommission_retry_backoff(&rx, retry_delay).await { + decommission_cancel_signal_result(rx.is_cancelled())?; + } + continue; + } + break result; + } + }; + match attempt_result { + Ok(DecommissionEntryAttemptOutcome::Complete) => return Ok(()), + Ok(DecommissionEntryAttemptOutcome::SourceChanged) => { let retry_delay = decommission_retry_backoff_delay(DECOMMISSION_SOURCE_CLEANUP_RETRY_DELAY, entry_attempt); warn!( event = EVENT_DECOMMISSION_ENTRY, @@ -10360,6 +10423,7 @@ impl ECStore { decommission_cancel_signal_result(rx.is_cancelled())?; } } + Err(err) => return Err(err), } } @@ -10433,8 +10497,34 @@ impl ECStore { let mut fivs = load_decommission_entry_exact_versions(&set, &entry, &bucket, "file_info_versions").await?; - fivs.versions - .sort_by_key(|v| (v.mod_time.is_none(), std::cmp::Reverse(v.mod_time))); + let pending_mutations = if let Some(owner) = capacity_owner { + self.pool_meta + .read() + .await + .pools + .get(owner.source_pool_index) + .and_then(|pool| pool.decommission.as_ref()) + .and_then(|info| info.capacity_reservation.as_ref()) + .filter(|reservation| reservation.admits_cleanup_owner(owner)) + .map(|reservation| { + reservation + .targets + .iter() + .filter_map(|target| target.pending_mutation_id) + .collect::>() + }) + .unwrap_or_default() + } else { + HashSet::new() + }; + fivs.versions.sort_by_key(|version| { + let mutation_id = capacity_owner.map(|owner| decommission_capacity_version_mutation_id(owner, &bucket, version)); + ( + mutation_id.is_none_or(|mutation_id| !pending_mutations.contains(&mutation_id)), + version.mod_time.is_none(), + std::cmp::Reverse(version.mod_time), + ) + }); let mut decommissioned: usize = 0; let mut expired: usize = 0; @@ -15352,6 +15442,18 @@ mod tests { assert!(is_decommission_target_capacity_error(&storage_full)); } + #[test] + fn decommission_capacity_intent_conflict_accepts_context_wrapped_error() { + let err = with_decommission_entry_context( + "migrate_object", + "bucket", + "object", + decommission_capacity_blocked_error("target has an unresolved target capacity intent"), + ); + + assert!(is_decommission_capacity_intent_conflict(&err)); + } + #[test] fn decommission_target_capacity_error_rejects_unrelated_errors() { assert!(!is_decommission_target_capacity_error(&Error::SlowDown)); @@ -15428,6 +15530,7 @@ mod tests { #[test] fn decommission_delete_marker_opts_preserves_suspended_null_version() { let version = rustfs_filemeta::FileInfo { + name: "object".to_string(), deleted: true, ..Default::default() }; @@ -15436,6 +15539,25 @@ mod tests { assert!(!opts.versioned); assert!(opts.version_suspended); assert_eq!(opts.version_id.as_deref(), Some(uuid::Uuid::nil().to_string().as_str())); + + let owner = DecommissionCapacityOwner { + source_pool_index: 7, + operation_id: uuid::Uuid::new_v4(), + generation: 1, + owner_nonce: uuid::Uuid::new_v4(), + mutation_id: None, + }; + assert_eq!( + decommission_capacity_version_mutation_id(owner, "bucket", &version), + decommission_capacity_mutation_id( + owner, + "bucket", + &version.name, + opts.version_id.as_deref(), + opts.delete_marker, + opts.mod_time, + ) + ); } #[test] @@ -16459,10 +16581,11 @@ mod pools_tests { wait_decommission_worker_drain, with_decommission_entry_context, }; use super::{ - DecommissionCapacityOwner, DecommissionCapacityReservation, decommission_capacity_mutation_id, - ensure_decommission_target_owner_admission, ensure_external_decommission_target_admission, - is_decommission_capacity_blocked_error, record_decommission_target_consumption, reserve_decommission_target_pending, - resolve_decommission_target_pending, set_decommission_capacity_info_overrides_for_test, + DecommissionCapacityOwner, DecommissionCapacityReservation, DecommissionCapacityTemporaryMutation, + decommission_capacity_mutation_id, ensure_decommission_target_owner_admission, + ensure_external_decommission_target_admission, is_decommission_capacity_blocked_error, + record_decommission_target_consumption, reserve_decommission_target_pending, resolve_decommission_target_pending, + set_decommission_capacity_info_overrides_for_test, }; use crate::bucket::lifecycle::{ DurableIlmRecordCheckpoint, @@ -16483,10 +16606,12 @@ mod pools_tests { use crate::storage_api_contracts::{object::ObjectIO, range::HTTPRangeSpec}; use crate::store::ECStore; use byteorder::{ByteOrder, LittleEndian}; + use rmp_serde::Serializer; use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo}; use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams}; use rustfs_lock::{GlobalLockManager, LocalClient, LockRequest, LockType, NamespaceLock, ObjectKey}; use rustfs_rio::Index; + use serde::Serialize; use std::future::Future; use std::io::Cursor; use std::sync::{ @@ -16545,6 +16670,7 @@ mod pools_tests { decommission_cancelers: tokio::sync::RwLock::new(cancelers), start_gate: tokio::sync::Mutex::new(()), pool_meta_save_gate: tokio::sync::Mutex::new(super::PoolMetaWriteState::for_test_bootstrap()), + decommission_capacity_entry_gate: tokio::sync::Mutex::default(), ctx, bucket_fence_registry: Arc::default(), }) @@ -20710,6 +20836,36 @@ mod pools_tests { } } + #[test] + fn decommission_capacity_target_round_trip_preserves_temporary_mutation_without_pending_intent() { + let mutation_id = uuid::Uuid::new_v4(); + let target = DecommissionCapacityTarget { + pool_index: 1, + layout: DecommissionErasureLayout { data: 2, parity: 2 }, + physical_total_at_reservation: 200, + physical_free_at_reservation: 200, + reserved_physical_bytes: 200, + consumed_physical_bytes: 0, + observed_physical_bytes: 1, + inflight_physical_bytes: 1, + pending_physical_bytes: 0, + pending_mutation_id: None, + temporary_mutations: vec![DecommissionCapacityTemporaryMutation { + mutation_id, + physical_bytes: 1, + }], + }; + let mut encoded = Vec::new(); + target + .serialize(&mut Serializer::new(&mut encoded)) + .expect("capacity target should serialize"); + + let restored: DecommissionCapacityTarget = + rmp_serde::from_slice(&encoded).expect("capacity target with a released pending intent should deserialize"); + + assert_eq!(restored, target); + } + #[test] fn decommission_capacity_reservation_recovers_expired_lease_after_restart_round_trip() { let created_at = OffsetDateTime::UNIX_EPOCH + Duration::hours(1); diff --git a/crates/ecstore/src/core/pools_test.rs b/crates/ecstore/src/core/pools_test.rs index 162077c00..5cd5dc714 100644 --- a/crates/ecstore/src/core/pools_test.rs +++ b/crates/ecstore/src/core/pools_test.rs @@ -422,6 +422,7 @@ mod decommission_lock_order_tests { pool_meta_save_gate: tokio::sync::Mutex::new( other_store.pool_meta_save_gate.lock().await.independent_clone_for_test(), ), + decommission_capacity_entry_gate: tokio::sync::Mutex::default(), ctx, bucket_fence_registry: Arc::default(), }); @@ -1707,7 +1708,6 @@ mod decommission_lock_order_tests { set_decommission_capacity_info_overrides_for_test(lossy_store.id, (0..40).map(|_| capacity_snapshot()).collect()); let part_barrier = MultipartCommitBarrier::install(&bucket, object, MultipartCommitPause::PutPartAfterRename); let abort_barrier = data_movement::DataMovementMultipartAbortBarrier::install(&bucket, object); - tokio::time::pause(); let migration = tokio::spawn({ let migration_store = Arc::clone(&lossy_store); let migration_bucket = bucket.clone(); @@ -1725,6 +1725,7 @@ mod decommission_lock_order_tests { } }); part_barrier.wait_until_paused().await; + tokio::time::pause(); tokio::task::yield_now().await; refresh_calls.arm(); tokio::time::advance(Duration::from_secs(11)).await; diff --git a/crates/ecstore/src/data_usage/mod.rs b/crates/ecstore/src/data_usage/mod.rs index b829208ea..c144435a6 100644 --- a/crates/ecstore/src/data_usage/mod.rs +++ b/crates/ecstore/src/data_usage/mod.rs @@ -2964,6 +2964,7 @@ mod tests { decommission_cancelers: RwLock::new(Vec::new()), start_gate: TokioMutex::new(()), pool_meta_save_gate: TokioMutex::default(), + decommission_capacity_entry_gate: TokioMutex::default(), ctx, bucket_fence_registry: Arc::default(), }) diff --git a/crates/ecstore/src/services/rebalance/mod.rs b/crates/ecstore/src/services/rebalance/mod.rs index 352c1138e..96d02eedc 100644 --- a/crates/ecstore/src/services/rebalance/mod.rs +++ b/crates/ecstore/src/services/rebalance/mod.rs @@ -83,6 +83,7 @@ pub async fn test_store_with_persisted_rebalance_meta( decommission_cancelers: tokio::sync::RwLock::new(vec![None]), start_gate: tokio::sync::Mutex::new(()), pool_meta_save_gate: tokio::sync::Mutex::default(), + decommission_capacity_entry_gate: tokio::sync::Mutex::default(), ctx, bucket_fence_registry: std::sync::Arc::default(), }); @@ -232,6 +233,7 @@ async fn test_pool_stores_with_contexts( decommission_cancelers: tokio::sync::RwLock::new(vec![None; pool_count]), start_gate: tokio::sync::Mutex::new(()), pool_meta_save_gate: tokio::sync::Mutex::new(pool_meta_write_state.independent_clone_for_test()), + decommission_capacity_entry_gate: tokio::sync::Mutex::default(), ctx: store_ctx, bucket_fence_registry: std::sync::Arc::default(), }) diff --git a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs index a544ed733..d280d403a 100644 --- a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs +++ b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs @@ -3009,6 +3009,7 @@ fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc, pool_idx: usize) { + let layout = DecommissionErasureLayout { data: 2, parity: 2 }; + let source_physical_bytes = 1024 * 1024 * 1024; + let target_physical_bytes = source_physical_bytes * 8; + let capacity = store + .pools + .iter() + .enumerate() + .map(|(index, _)| { + if index == pool_idx { + DecommissionPoolCapacityInfo::for_test(index, layout, 0, source_physical_bytes, source_physical_bytes) + } else { + DecommissionPoolCapacityInfo::for_test(index, layout, target_physical_bytes, target_physical_bytes, 0) + } + }) + .collect(); + set_decommission_capacity_info_overrides_for_test(store.id, vec![capacity]); + } + async fn mark_test_pool_decommissioning(store: &Arc, pool_idx: usize) { - let mut pool_meta = store.pool_meta.write().await; - pool_meta.pools[pool_idx].decommission = Some(PoolDecommissionInfo { - start_time: Some(OffsetDateTime::now_utc()), - ..Default::default() - }); + set_test_decommission_capacity_override(store, pool_idx); + store + .save_current_pool_meta_for_decommission_start(&[pool_idx], Vec::new()) + .await + .expect("test decommission capacity reservation should activate"); } const DECOMMISSION_TEST_FAULT_STAGE_DELETE_MARKER: &str = "delete_marker_copy"; @@ -2407,45 +2428,6 @@ mod tests { ); } - async fn assert_suspended_decommission_converged(store: &Arc, bucket: &str, object: &str) { - let source_versions = store.pools[0] - .get_disks_by_key(object) - .load_file_info_versions_exact(bucket, object) - .await - .expect("source versions should remain readable after suspended convergence"); - assert!( - source_versions.is_none_or(|versions| versions.versions.is_empty()), - "worker convergence must remove only the decommissioned source null version" - ); - - let target_versions = store.pools[1] - .get_disks_by_key(object) - .load_file_info_versions_exact(bucket, object) - .await - .expect("active target versions should be readable") - .expect("active target must retain the suspended DELETE marker"); - assert!( - matches!(target_versions.versions.as_slice(), [marker] if marker.deleted && marker.version_id.is_none_or(|version_id| version_id.is_nil())), - "active target must contain only its null delete marker: {target_versions:?}" - ); - - let err = store - .get_object_info( - bucket, - object, - &ObjectOptions { - version_suspended: true, - ..Default::default() - }, - ) - .await - .expect_err("the active null delete marker must hide the migrated source generation"); - assert!( - matches!(err, StorageError::ObjectNotFound(_, _)), - "unexpected suspended latest-object result: {err:?}" - ); - } - #[tokio::test] #[serial_test::serial(storage_class_env)] async fn tag_updates_skip_active_rebalance_source_pool() { @@ -4196,13 +4178,7 @@ mod tests { .put_object(&bucket, &object, &mut source, &ObjectOptions::default()) .await .expect("write source object to the pool being decommissioned"); - { - let mut pool_meta = store.pool_meta.write().await; - pool_meta.pools[0].decommission = Some(PoolDecommissionInfo { - start_time: Some(OffsetDateTime::now_utc()), - ..Default::default() - }); - } + mark_test_pool_decommissioning(&store, 0).await; assert!(store.is_suspended(0).await, "pool 0 must be a suspended decommission source"); let barrier = crate::set_disk::PutObjectCommitBarrier::install( @@ -4938,9 +4914,10 @@ mod tests { .then(|| crate::set_disk::NewMultipartUploadCommitObservation::install(&bucket, object)); let barrier = crate::set_disk::MultipartCommitBarrier::install(&bucket, object, pause); let source_set = store.pools[0].get_disks_by_key(object); + let recovery_source_set = Arc::clone(&source_set); let worker_store = Arc::clone(&store); let worker_bucket = bucket.clone(); - let worker = tokio::spawn(async move { + let mut worker = tokio::spawn(async move { worker_store .decommission_entry_for_test( 0, @@ -4954,7 +4931,10 @@ mod tests { .await }); - barrier.wait_until_paused().await; + tokio::select! { + () = barrier.wait_until_paused() => {} + result = &mut worker => panic!("decommission multipart worker exited before the commit barrier: {result:?}"), + } loss_hook.mark_lost(); barrier.release(); drop(barrier); @@ -4976,6 +4956,19 @@ mod tests { .await .expect("list target multipart uploads after fenced migration"); assert!(uploads.uploads.is_empty(), "fenced multipart migration must not retain target staging"); + drop(loss_hook); + store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + bucket.clone(), + recovery_source_set, + ) + .await + .expect("same-mutation retry should recover the durable capacity intent"); } shutdown.cancel(); @@ -5184,6 +5177,12 @@ mod tests { .expect("self-copy should keep using the committed active target"); assert_eq!(active_copy_result.data_dir, active_copy_data_dir); + { + let mut pool_meta = store.pool_meta.write().await; + pool_meta.pools[1].decommission = None; + } + mark_test_pool_decommissioning(&store, 1).await; + let cleanup_barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(&bucket, object); let commit_barrier = crate::set_disk::PutObjectCommitBarrier::install( &bucket, @@ -6197,28 +6196,27 @@ mod tests { write_suspended_decommission_source(&store, &bucket, object).await; mark_test_pool_decommissioning(&store, 0).await; - let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object); - let delete_store = Arc::clone(&store); - let delete_bucket = bucket.clone(); - let delete = tokio::spawn(async move { - delete_store - .delete_object( - &delete_bucket, - object, - ObjectOptions { - version_suspended: true, - ..Default::default() - }, - ) - .await - }); - delete_barrier.wait_until_paused().await; + let delete_err = store + .delete_object( + &bucket, + object, + ObjectOptions { + version_suspended: true, + ..Default::default() + }, + ) + .await + .expect_err("capacity-reserved target must reject a concurrent suspended DELETE"); + assert!( + matches!(delete_err, Error::SlowDown), + "unexpected suspended DELETE result: {delete_err:?}" + ); assert_suspended_null_source_present(&store, &bucket, object).await; let source_set = store.pools[0].get_disks_by_key(object); let worker_store = Arc::clone(&store); let worker_bucket = bucket.clone(); - let worker = tokio::spawn(async move { + tokio::spawn(async move { worker_store .decommission_entry_for_test( 0, @@ -6230,25 +6228,25 @@ mod tests { source_set, ) .await - }); + }) + .await + .expect("suspended decommission worker should join") + .expect("worker must migrate the fenced suspended source"); - delete_barrier.release(); - let marker = delete - .await - .expect("suspended DELETE task should join") - .expect("suspended DELETE should commit its active-pool marker"); - drop(delete_barrier); - assert!(marker.delete_marker, "suspended DELETE must create a marker"); - assert!( - marker.version_id.is_none_or(|version_id| version_id.is_nil()), - "suspended DELETE marker must keep the null version identity" + assert_decommission_source_absent( + &store, + &bucket, + object, + &ObjectOptions { + version_suspended: true, + ..Default::default() + }, + ) + .await; + assert_eq!( + read_decommission_target_body(&store, &bucket, object, &ObjectOptions::default()).await, + b"suspended source generation" ); - worker - .await - .expect("suspended decommission worker should join") - .expect("worker must treat the newer active null marker as a completed migration"); - - assert_suspended_decommission_converged(&store, &bucket, object).await; shutdown.cancel(); } @@ -6281,31 +6279,29 @@ mod tests { }, None, )); - let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object); - let delete_store = Arc::clone(&store); - let delete_bucket = bucket.clone(); - let delete = tokio::spawn(async move { - delete_store - .delete_objects( - &delete_bucket, - vec![ObjectToDelete { - object_name: object.to_string(), - ..Default::default() - }], - ObjectOptions { - delete_replication_config_snapshot: Some(delete_config_snapshot), - ..Default::default() - }, - ) - .await - }); - delete_barrier.wait_until_paused().await; + let (_deleted, errors) = store + .delete_objects( + &bucket, + vec![ObjectToDelete { + object_name: object.to_string(), + ..Default::default() + }], + ObjectOptions { + delete_replication_config_snapshot: Some(delete_config_snapshot), + ..Default::default() + }, + ) + .await; + assert!( + matches!(errors.as_slice(), [Some(Error::SlowDown)]), + "unexpected suspended batch DELETE result: {errors:?}" + ); assert_suspended_null_source_present(&store, &bucket, object).await; let source_set = store.pools[0].get_disks_by_key(object); let worker_store = Arc::clone(&store); let worker_bucket = bucket.clone(); - let worker = tokio::spawn(async move { + tokio::spawn(async move { worker_store .decommission_entry_for_test( 0, @@ -6317,22 +6313,25 @@ mod tests { source_set, ) .await - }); + }) + .await + .expect("suspended batch decommission worker should join") + .expect("worker must migrate the batch-fenced suspended source"); - delete_barrier.release(); - let (deleted, errors) = delete.await.expect("suspended batch DELETE task should join"); - drop(delete_barrier); - assert!(errors.iter().all(Option::is_none), "suspended batch DELETE should succeed: {errors:?}"); - assert!( - matches!(deleted.as_slice(), [marker] if marker.delete_marker && marker.delete_marker_version_id.is_none_or(|version_id| version_id.is_nil())), - "suspended batch DELETE must create one null marker: {deleted:?}" + assert_decommission_source_absent( + &store, + &bucket, + object, + &ObjectOptions { + version_suspended: true, + ..Default::default() + }, + ) + .await; + assert_eq!( + read_decommission_target_body(&store, &bucket, object, &ObjectOptions::default()).await, + b"suspended source generation" ); - worker - .await - .expect("suspended batch decommission worker should join") - .expect("worker must treat the newer batch null marker as a completed migration"); - - assert_suspended_decommission_converged(&store, &bucket, object).await; shutdown.cancel(); } @@ -7325,6 +7324,7 @@ mod tests { .await .expect("legacy decommission queue should reload after restart"); *store.pool_meta.write().await = restarted_pool_meta; + set_test_decommission_capacity_override(&store, 0); store .promote_queued_decommission_for_test(0) .await diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index c3dd77605..a7a9dc4ee 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -260,6 +260,12 @@ pub struct ECStore { /// Lock order: acquire `pool_meta_save_gate`, then the distributed /// `pool.bin` fence, then clone `pool_meta` under a short read lock. pub(crate) pool_meta_save_gate: Mutex, + /// Serializes decommission entries while the durable capacity ledger has + /// one target mutation intent slot. + /// + /// Lock order: acquire this gate before object namespaces or + /// `pool_meta_save_gate`. + pub(crate) decommission_capacity_entry_gate: Mutex<()>, /// Per-instance runtime state (Phase 5, backlog#939). /// /// Carries this instance's identity/runtime out of the process globals so @@ -1514,6 +1520,7 @@ mod tests { decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), pool_meta_save_gate: Mutex::default(), + decommission_capacity_entry_gate: Mutex::default(), ctx, bucket_fence_registry: Arc::default(), }; @@ -1589,6 +1596,7 @@ mod tests { decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), pool_meta_save_gate: Mutex::default(), + decommission_capacity_entry_gate: Mutex::default(), ctx, bucket_fence_registry: Arc::default(), }) diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index a2eed1658..0c74cad5f 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -1103,6 +1103,7 @@ mod tests { decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), pool_meta_save_gate: Mutex::default(), + decommission_capacity_entry_gate: Mutex::default(), ctx: crate::runtime::instance::bootstrap_ctx(), bucket_fence_registry: std::sync::Arc::default(), } diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index b3e1c5008..c1fbf3b3f 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -1204,7 +1204,11 @@ fn delete_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions } fn should_delete_from_all_pools(opts: &ObjectOptions, pool_count: usize) -> bool { - pool_count > 0 && (!opts.versioned && !opts.version_suspended || opts.version_id.is_some()) + pool_count > 0 && delete_only_releases_capacity(opts) +} + +fn delete_only_releases_capacity(opts: &ObjectOptions) -> bool { + !opts.versioned && !opts.version_suspended || opts.version_id.is_some() } fn batch_delete_creates_latest_marker(object: &ObjectToDelete, delete_config_snapshot: &DeleteReplicationConfigSnapshot) -> bool { @@ -2013,16 +2017,69 @@ impl ECStore { bucket: &str, lock_object: &str, target_object: &str, - mut opts: ObjectOptions, + opts: ObjectOptions, operation: F, ) -> Result where F: FnOnce(ObjectOptions) -> Fut, Fut: std::future::Future>, { - let (capacity_guard, has_active_decommission) = self - .acquire_external_decommission_capacity_fence_with_active_source(&[target_pool_idx], "mutation") - .await?; + self.run_external_decommission_capacity_object_operation( + target_pool_idx, + bucket, + (lock_object, target_object), + opts, + false, + operation, + ) + .await + } + + pub(super) async fn run_external_decommission_capacity_object_delete( + &self, + target_pool_idx: usize, + bucket: &str, + lock_object: &str, + target_object: &str, + opts: ObjectOptions, + operation: F, + ) -> Result + where + F: FnOnce(ObjectOptions) -> Fut, + Fut: std::future::Future>, + { + let capacity_releasing = delete_only_releases_capacity(&opts); + self.run_external_decommission_capacity_object_operation( + target_pool_idx, + bucket, + (lock_object, target_object), + opts, + capacity_releasing, + operation, + ) + .await + } + + async fn run_external_decommission_capacity_object_operation( + &self, + target_pool_idx: usize, + bucket: &str, + objects: (&str, &str), + mut opts: ObjectOptions, + capacity_releasing: bool, + operation: F, + ) -> Result + where + F: FnOnce(ObjectOptions) -> Fut, + Fut: std::future::Future>, + { + let (lock_object, target_object) = objects; + let (capacity_guard, has_active_decommission) = if capacity_releasing { + self.acquire_decommission_capacity_release_fence_with_active_source().await? + } else { + self.acquire_external_decommission_capacity_fence_with_active_source(&[target_pool_idx], "mutation") + .await? + }; let (capacity_guard, object_guard) = if has_active_decommission && !opts.no_lock { // Active migration acquires the object namespace before its capacity // write. Match that order, then recheck capacity admission. @@ -2034,9 +2091,12 @@ impl ECStore { .await?; self.apply_decommission_target_mutation_fence(target_pool_idx, target_object, &mut opts, Some(&guard)) .await; - let capacity_guard = self - .acquire_external_decommission_capacity_fence(&[target_pool_idx], "mutation") - .await?; + let capacity_guard = if capacity_releasing { + self.acquire_decommission_capacity_release_fence_with_active_source().await?.0 + } else { + self.acquire_external_decommission_capacity_fence(&[target_pool_idx], "mutation") + .await? + }; (capacity_guard, Some(guard)) } else { (capacity_guard, None) @@ -2633,20 +2693,26 @@ impl ECStore { return Err(decommission_free_version_overwrite_error(bucket, &object, fi.version_id)); } + let expected_data_bytes = usize::try_from(fi.size).ok(); let result = self - .run_decommission_capacity_admitted_mutation(idx, DecommissionCapacityOwner::from_options(&opts), None, || async { - if is_free_version { - self.pools[idx] - .get_disks_by_key(&object) - .decommission_tier_free_version(bucket, &object, &fi, &opts) - .await - } else { - self.pools[idx] - .get_disks_by_key(&object) - .decommission_tiered_object(bucket, &object, &fi, &opts) - .await - } - }) + .run_decommission_capacity_admitted_mutation( + idx, + DecommissionCapacityOwner::from_options(&opts), + expected_data_bytes, + || async { + if is_free_version { + self.pools[idx] + .get_disks_by_key(&object) + .decommission_tier_free_version(bucket, &object, &fi, &opts) + .await + } else { + self.pools[idx] + .get_disks_by_key(&object) + .decommission_tiered_object(bucket, &object, &fi, &opts) + .await + } + }, + ) .await; if matches!(result, Err(Error::PreconditionFailed)) { if self @@ -3572,7 +3638,7 @@ impl ECStore { let pool_idx = pool.pool_idx; let pool = pool.clone(); match self - .run_external_decommission_capacity_object_mutation( + .run_external_decommission_capacity_object_delete( pool_idx, bucket, object, @@ -5823,6 +5889,7 @@ mod tests { decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), pool_meta_save_gate: Mutex::default(), + decommission_capacity_entry_gate: Mutex::default(), ctx: crate::runtime::instance::bootstrap_ctx(), bucket_fence_registry: std::sync::Arc::default(), } @@ -5886,6 +5953,7 @@ mod tests { decommission_cancelers: RwLock::new(Vec::new()), start_gate: Mutex::new(()), pool_meta_save_gate: Mutex::default(), + decommission_capacity_entry_gate: Mutex::default(), ctx, bucket_fence_registry: std::sync::Arc::default(), } diff --git a/crates/ecstore/src/store/rebalance.rs b/crates/ecstore/src/store/rebalance.rs index 14de70b66..2bfe36af8 100644 --- a/crates/ecstore/src/store/rebalance.rs +++ b/crates/ecstore/src/store/rebalance.rs @@ -929,7 +929,7 @@ impl ECStore { results.push(RebalanceDeletePoolResult { pool_idx: idx, result: self - .run_external_decommission_capacity_object_mutation( + .run_external_decommission_capacity_object_delete( idx, bucket, object, diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 1f2b217b4..033032f4d 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -1287,6 +1287,7 @@ fn scanner_startup_fails_closed_on_nonempty_corrupt_cycle_state() { } #[tokio::test] +#[serial] async fn corrupt_cycle_state_is_quarantined_once() { let store = Arc::new(MemoryConfigStore::default()); let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -1337,6 +1338,7 @@ async fn corrupt_cycle_state_is_quarantined_once() { } #[tokio::test] +#[serial] async fn empty_cycle_state_object_is_quarantined_as_corrupt() { let store = Arc::new(MemoryConfigStore::default()); let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -1357,6 +1359,7 @@ async fn empty_cycle_state_object_is_quarantined_as_corrupt() { } #[tokio::test] +#[serial] async fn future_cycle_state_schema_is_recovery_required() { let store = Arc::new(MemoryConfigStore::default()); let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -1375,6 +1378,7 @@ async fn future_cycle_state_schema_is_recovery_required() { } #[tokio::test] +#[serial] async fn concurrent_leaders_cannot_quarantine_newer_cycle_state() { let store = Arc::new(MemoryConfigStore::default()); let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -1401,6 +1405,7 @@ async fn concurrent_leaders_cannot_quarantine_newer_cycle_state() { } #[tokio::test] +#[serial] async fn cleanup_pending_marker_blocks_a_rewritten_primary_after_restart() { let store = Arc::new(MemoryConfigStore::default()); let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -1476,6 +1481,7 @@ fn full_rescan_reset_accepts_unknown_marker_fields_without_trusting_cursor() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rebuilds_after_malformed_marker_without_trusting_cursor() { let (_temp_dir, store) = setup_scanner_cycle_store().await; save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) @@ -1502,6 +1508,7 @@ async fn full_rescan_reset_rebuilds_after_malformed_marker_without_trusting_curs } #[tokio::test] +#[serial] async fn full_rescan_reset_ignores_epoch_from_malformed_future_primary() { let (_temp_dir, store) = setup_scanner_cycle_store().await; let mut future_primary = vec![0; 24]; @@ -1562,6 +1569,7 @@ async fn ecstore_exact_recovery_marker_delete_honors_etag() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rejects_corrupt_primary_under_stale_blocked_marker() { let (_temp_dir, store) = setup_scanner_cycle_store().await; let corrupt_primary = vec![0xff, 0x00, 0x01]; @@ -1612,6 +1620,7 @@ async fn full_rescan_reset_rejects_corrupt_primary_under_stale_blocked_marker() } #[tokio::test] +#[serial] async fn full_rescan_reset_preserves_valid_primary_when_marker_is_malformed() { let (_temp_dir, store) = setup_scanner_cycle_store().await; let primary = CurrentCycle { @@ -1683,6 +1692,7 @@ async fn full_rescan_reset_preserves_valid_primary_when_marker_is_malformed() { } #[tokio::test] +#[serial] async fn full_rescan_reset_resumes_cleanup_pending_preserved_primary() { let (_temp_dir, store) = setup_scanner_cycle_store().await; let completed_at = Utc::now(); @@ -1762,6 +1772,7 @@ async fn full_rescan_reset_resumes_cleanup_pending_preserved_primary() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rebuilds_oversized_regular_primary_with_malformed_marker() { let (_temp_dir, store) = setup_scanner_cycle_store().await; save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1]) @@ -1788,6 +1799,7 @@ async fn full_rescan_reset_rebuilds_oversized_regular_primary_with_malformed_mar } #[tokio::test] +#[serial] async fn full_rescan_reset_rebuilds_oversized_primary_after_cleanup_marker() { let (_temp_dir, store) = setup_scanner_cycle_store().await; save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1]) @@ -1838,6 +1850,7 @@ async fn full_rescan_reset_rebuilds_oversized_primary_after_cleanup_marker() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rebuilds_with_oversized_marker() { let (_temp_dir, store) = setup_scanner_cycle_store().await; save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) @@ -1863,6 +1876,7 @@ async fn full_rescan_reset_rebuilds_with_oversized_marker() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rebuilds_with_empty_marker() { let (_temp_dir, store) = setup_scanner_cycle_store().await; save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) @@ -1888,6 +1902,7 @@ async fn full_rescan_reset_rebuilds_with_empty_marker() { } #[tokio::test] +#[serial] async fn full_rescan_reset_keeps_cleanup_marker_when_preserved_epoch_is_exhausted() { let (_temp_dir, store) = setup_scanner_cycle_store().await; let primary = CurrentCycle { @@ -1927,6 +1942,7 @@ async fn full_rescan_reset_keeps_cleanup_marker_when_preserved_epoch_is_exhauste } #[tokio::test] +#[serial] async fn full_rescan_reset_rejects_preserved_epoch_that_would_be_terminal() { let (_temp_dir, store) = setup_scanner_cycle_store().await; let primary = CurrentCycle { @@ -1963,6 +1979,7 @@ async fn full_rescan_reset_rejects_preserved_epoch_that_would_be_terminal() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rejects_usage_floor_that_would_be_terminal() { let (_temp_dir, store) = setup_scanner_cycle_store().await; save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01]) @@ -1998,6 +2015,7 @@ async fn full_rescan_reset_rejects_usage_floor_that_would_be_terminal() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rebuilds_empty_primary_with_malformed_marker() { let (_temp_dir, store) = setup_scanner_cycle_store().await; save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), Vec::new()) @@ -2024,6 +2042,7 @@ async fn full_rescan_reset_rebuilds_empty_primary_with_malformed_marker() { } #[tokio::test] +#[serial] async fn full_rescan_reset_rebuilds_when_primary_cycle_state_is_missing() { let (_temp_dir, store) = setup_scanner_cycle_store().await; let marker = ScannerCycleRecoveryMarker { @@ -2064,6 +2083,7 @@ async fn full_rescan_reset_rebuilds_when_primary_cycle_state_is_missing() { } #[tokio::test] +#[serial] async fn corrupt_cycle_state_rename_or_marker_failure_stays_recovery_required() { let store = Arc::new(MemoryConfigStore::default()); let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -2089,6 +2109,7 @@ async fn corrupt_cycle_state_rename_or_marker_failure_stays_recovery_required() } #[tokio::test] +#[serial] async fn oversized_or_symlinked_cycle_state_is_rejected() { let store = Arc::new(MemoryConfigStore::default()); let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -3328,6 +3349,7 @@ fn fenced_usage_bootstrap_retains_partial_cycle_progress() { } #[tokio::test] +#[serial] async fn missing_usage_floor_rebuilds_persisted_cycle_before_leadership_claim() { let store = Arc::new(MemoryConfigStore::default()); let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()); @@ -4283,6 +4305,7 @@ async fn cycle_budget_lease_takeover_rejects_old_generation() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_preserves_newer_snapshot() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(2); @@ -4309,6 +4332,7 @@ async fn test_store_data_usage_in_backend_preserves_newer_snapshot() { } #[tokio::test] +#[serial] async fn test_usage_save_object_not_found_defers_only_with_a_fresh_route_barrier() { for (route_blocked, expected) in [ (true, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement)), @@ -4368,6 +4392,7 @@ async fn test_usage_save_object_not_found_defers_only_with_a_fresh_route_barrier } #[tokio::test] +#[serial] async fn test_usage_save_route_barrier_prevents_missing_snapshot_creation() { for observational in [false, true] { let store = Arc::new(MemoryConfigStore::default()); @@ -4407,6 +4432,7 @@ async fn test_usage_save_route_barrier_prevents_missing_snapshot_creation() { } #[tokio::test] +#[serial] async fn test_observational_usage_defers_when_authoritative_baseline_is_missing() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(1); @@ -4436,6 +4462,7 @@ async fn test_observational_usage_defers_when_authoritative_baseline_is_missing( } #[tokio::test] +#[serial] async fn test_observational_usage_uses_fenced_backup_when_v2_primary_has_no_identity() { let store = Arc::new(MemoryConfigStore::default()); let primary = DataUsageInfo { @@ -4537,6 +4564,7 @@ async fn test_usage_route_barrier_precedes_durable_reconciliation() { } #[tokio::test] +#[serial] async fn coordinator_does_not_put_after_remote_generation_flip() { let store = Arc::new(MemoryConfigStore::default()); let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -4578,6 +4606,7 @@ async fn coordinator_does_not_put_after_remote_generation_flip() { } #[tokio::test] +#[serial] async fn coordinator_classifies_an_expired_publication_lease() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(1); @@ -4696,6 +4725,7 @@ async fn test_deferred_usage_save_keeps_last_real_save_metric() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_fences_interleaving_newer_writer() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(1); @@ -4725,6 +4755,7 @@ async fn test_store_data_usage_in_backend_fences_interleaving_newer_writer() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_does_not_resurrect_deleted_bucket_after_conflict() { let store = Arc::new(MemoryConfigStore::default()); let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -4796,6 +4827,7 @@ async fn test_store_data_usage_in_backend_does_not_resurrect_deleted_bucket_afte } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_updates_backup_with_new_bucket() { let store = Arc::new(MemoryConfigStore::default()); let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -4856,6 +4888,7 @@ async fn test_store_data_usage_in_backend_updates_backup_with_new_bucket() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_repairs_backup_after_primary_only_commit() { let store = Arc::new(MemoryConfigStore::default()); let main_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -4885,6 +4918,7 @@ async fn test_store_data_usage_in_backend_repairs_backup_after_primary_only_comm } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_copies_concurrent_bucket_removal_to_backup() { let store = Arc::new(MemoryConfigStore::default()); let main_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -4951,6 +4985,7 @@ async fn test_store_data_usage_in_backend_copies_concurrent_bucket_removal_to_ba } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_retries_after_stale_interleaving_writer() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(1); @@ -4991,6 +5026,7 @@ async fn test_store_data_usage_in_backend_retries_after_stale_interleaving_write } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_rejects_untimestamped_complete_snapshot() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(2); @@ -5023,6 +5059,7 @@ async fn test_store_data_usage_in_backend_rejects_untimestamped_complete_snapsho } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_recognizes_already_durable_snapshot() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(1); @@ -5053,6 +5090,7 @@ async fn test_store_data_usage_in_backend_recognizes_already_durable_snapshot() } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_advances_past_changed_same_epoch_cycle() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(1); @@ -5102,6 +5140,7 @@ async fn test_store_data_usage_in_backend_advances_past_changed_same_epoch_cycle } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_orders_scanner_cycles_before_wall_clock() { let store = Arc::new(MemoryConfigStore::default()); let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -5162,6 +5201,7 @@ async fn test_store_data_usage_in_backend_orders_scanner_cycles_before_wall_cloc } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_orders_leader_epochs_before_cycles() { let store = Arc::new(MemoryConfigStore::default()); let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -5224,6 +5264,7 @@ async fn test_store_data_usage_in_backend_orders_leader_epochs_before_cycles() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_keeps_first_same_cycle_snapshot() { let store = Arc::new(MemoryConfigStore::default()); let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -5270,6 +5311,7 @@ async fn test_store_data_usage_in_backend_keeps_first_same_cycle_snapshot() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_rejects_incomplete_snapshot() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(2); @@ -5302,6 +5344,7 @@ async fn test_store_data_usage_in_backend_rejects_incomplete_snapshot() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_preserves_superseded_status() { let store = Arc::new(MemoryConfigStore::default()); let authoritative_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -5353,6 +5396,7 @@ async fn test_store_data_usage_in_backend_preserves_superseded_status() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_removes_observed_after_authoritative_save() { let store = Arc::new(MemoryConfigStore::default()); let authoritative_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -5542,6 +5586,7 @@ fn test_stale_data_usage_update_reason_preserves_none_handling() { } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_keeps_backup_when_primary_save_fails() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(11); @@ -5584,6 +5629,7 @@ async fn test_store_data_usage_in_backend_keeps_backup_when_primary_save_fails() } #[tokio::test] +#[serial] async fn test_store_data_usage_in_backend_reports_missing_snapshot() { let store = Arc::new(MemoryConfigStore::default()); let (sender, receiver) = mpsc::channel(1);