diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 17fffac3d..c7dbf5686 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -440,6 +440,11 @@ pub mod rebalance { RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record, encode_rebalance_stop_propagation_record, }; + + #[cfg(feature = "test-util")] + pub mod test_util { + pub use crate::services::rebalance::entry::test_util::PausedRebalanceEntryTestFixture; + } } pub mod rio { diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 854a3854b..20b4b1ebf 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -2500,6 +2500,85 @@ fn lifecycle_action_skips_heal_version(action: IlmAction) -> bool { action.delete() } +#[cfg(test)] +struct LifecycleDataMovementMutationBarrierState { + bucket: String, + object: String, + arrived: tokio::sync::Notify, + release: tokio::sync::Notify, +} + +#[cfg(test)] +pub(crate) struct LifecycleDataMovementMutationBarrier { + state: Arc, +} + +#[cfg(test)] +static LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +impl LifecycleDataMovementMutationBarrier { + pub(crate) fn install(bucket: &str, object: &str) -> Self { + let state = Arc::new(LifecycleDataMovementMutationBarrierState { + bucket: bucket.to_string(), + object: object.to_string(), + arrived: tokio::sync::Notify::new(), + release: tokio::sync::Notify::new(), + }); + let mut slot = LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("lifecycle data movement mutation barrier should not be poisoned"); + assert!(slot.is_none(), "lifecycle data movement mutation barrier must be unique"); + *slot = Some(Arc::clone(&state)); + Self { state } + } + + pub(crate) async fn wait_until_paused(&self) { + tokio::time::timeout(std::time::Duration::from_secs(30), self.state.arrived.notified()) + .await + .expect("lifecycle data movement should reach its mutation boundary"); + } + + pub(crate) fn release(&self) { + self.state.release.notify_one(); + } +} + +#[cfg(test)] +impl Drop for LifecycleDataMovementMutationBarrier { + fn drop(&mut self) { + self.state.release.notify_one(); + let mut slot = LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("lifecycle data movement mutation barrier should not be poisoned"); + if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *slot = None; + } + } +} + +#[cfg(test)] +async fn pause_lifecycle_data_movement_mutation(bucket: &str, object: &str, has_run_fence_signal: bool) { + if !has_run_fence_signal { + return; + } + let barrier = LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("lifecycle data movement mutation barrier should not be poisoned") + .as_ref() + .filter(|barrier| barrier.bucket == bucket && barrier.object == object) + .cloned(); + if let Some(barrier) = barrier { + barrier.arrived.notify_one(); + barrier.release.notified().await; + } +} + fn resolve_data_movement_lifecycle_expiry_result(action: IlmAction, apply_actions: bool, applied: bool) -> Result { if !apply_actions || applied { return Ok(true); @@ -2548,6 +2627,8 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement( Ok(false) } action if lifecycle_action_removes_data_movement_version(action) => { + #[cfg(test)] + pause_lifecycle_data_movement_mutation(bucket, &version.name, lock_lost_signal.is_some()).await; if lifecycle_delete_all_versions_blocked_by_replication(store.clone(), bucket, &object_info.name, action).await? { return Ok(false); } diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 31682e2e2..70c20b30d 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -1226,12 +1226,12 @@ pub(crate) async fn make_local_two_set_sets() -> (Vec, Arc) -> (Vec, Arc) { make_local_two_set_sets_for_pool_with_ctx(ctx, 0).await } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] pub(crate) async fn make_local_two_set_sets_for_pool_with_ctx( ctx: Arc, pool_idx: usize, diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index caccd1e94..85213a579 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -1062,7 +1062,7 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged( ensure_source_cleanup_versions_match(expected, ¤t, allowed_missing) } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] struct SourceCleanupDeleteBarrierState { bucket: String, object: String, @@ -1070,7 +1070,7 @@ struct SourceCleanupDeleteBarrierState { release: tokio::sync::Notify, } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] #[allow( dead_code, reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)" @@ -1079,11 +1079,11 @@ pub(crate) struct SourceCleanupDeleteBarrier { state: Arc, } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock>>> = std::sync::OnceLock::new(); -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] #[allow( dead_code, reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)" @@ -1116,7 +1116,7 @@ impl SourceCleanupDeleteBarrier { } } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] impl Drop for SourceCleanupDeleteBarrier { fn drop(&mut self) { self.state.release.notify_one(); @@ -1130,7 +1130,7 @@ impl Drop for SourceCleanupDeleteBarrier { } } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) { let barrier = SOURCE_CLEANUP_DELETE_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) @@ -1172,7 +1172,7 @@ pub(crate) async fn cleanup_source_entry_if_unchanged( ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?; - #[cfg(test)] + #[cfg(any(test, feature = "test-util"))] pause_source_cleanup_before_delete(bucket, object).await; let mut opts = ObjectOptions { diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index 530aac3f7..d788ef4dd 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -82,6 +82,75 @@ pub(super) struct RebalanceRunGuard { persisted_guard: rustfs_lock::NamespaceLockGuard, } +#[cfg(any(test, feature = "test-util"))] +struct RebalanceStopWaitProbeState { + expected_id: String, + attempted: std::sync::atomic::AtomicBool, + notify: tokio::sync::Notify, +} + +#[cfg(any(test, feature = "test-util"))] +static REBALANCE_STOP_WAIT_PROBES: std::sync::OnceLock>>> = + std::sync::OnceLock::new(); + +#[cfg(any(test, feature = "test-util"))] +pub(super) struct RebalanceStopWaitProbe { + state: Arc, +} + +#[cfg(any(test, feature = "test-util"))] +impl RebalanceStopWaitProbe { + pub(super) fn install(expected_id: &str) -> Self { + let state = Arc::new(RebalanceStopWaitProbeState { + expected_id: expected_id.to_string(), + attempted: std::sync::atomic::AtomicBool::new(false), + notify: tokio::sync::Notify::new(), + }); + REBALANCE_STOP_WAIT_PROBES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("rebalance stop wait probe should not be poisoned") + .push(Arc::clone(&state)); + Self { state } + } + + pub(super) async fn wait_until_attempted(&self) { + tokio::time::timeout(std::time::Duration::from_secs(30), async { + while !self.state.attempted.load(std::sync::atomic::Ordering::Acquire) { + self.state.notify.notified().await; + } + }) + .await + .expect("rebalance stop should reach the deterministic activation wait probe"); + } +} + +#[cfg(any(test, feature = "test-util"))] +impl Drop for RebalanceStopWaitProbe { + fn drop(&mut self) { + let mut probes = REBALANCE_STOP_WAIT_PROBES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("rebalance stop wait probe should not be poisoned"); + probes.retain(|state| !Arc::ptr_eq(state, &self.state)); + } +} + +#[cfg(any(test, feature = "test-util"))] +fn observe_rebalance_stop_wait_attempt(expected_id: Option<&str>) { + let probes = REBALANCE_STOP_WAIT_PROBES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("rebalance stop wait probe should not be poisoned") + .clone(); + for state in probes { + if expected_id == Some(state.expected_id.as_str()) { + state.attempted.store(true, std::sync::atomic::Ordering::Release); + state.notify.notify_one(); + } + } +} + impl RebalanceRunGuard { pub(super) fn ensure_held(&self, stage: &str) -> Result<()> { #[cfg(test)] @@ -945,6 +1014,24 @@ impl ECStore { .and_then(|meta| (!meta.id.is_empty()).then(|| meta.id.clone())) } + pub async fn cancel_rebalance_admission_for_id(&self, expected_id: &str) -> Result<()> { + let _start_guard = self.start_gate.lock().await; + let mut rebalance_meta = self.rebalance_meta.write().await; + ensure_rebalance_run_id(rebalance_meta.as_ref(), expected_id, "cancel rebalance admission")?; + let meta = rebalance_meta + .as_mut() + .ok_or_else(|| rebalance_metadata_not_initialized_error("cancel rebalance admission"))?; + if meta.stopped_at.is_some() || !is_rebalance_conflicting_with_decommission(meta) { + return Err(Error::other(format!( + "inactive rebalance rejected while cancelling admission: {expected_id}" + ))); + } + meta.cancel + .get_or_insert_with(tokio_util::sync::CancellationToken::new) + .cancel(); + Ok(()) + } + #[tracing::instrument(skip(self))] pub async fn stop_rebalance(self: &Arc) -> Result<()> { self.stop_rebalance_for_id(None).await @@ -965,7 +1052,11 @@ impl ECStore { }) }; let _activation_guard = match activation_gate { - Some(gate) => Some(gate.write_owned().await), + Some(gate) => { + #[cfg(any(test, feature = "test-util"))] + observe_rebalance_stop_wait_attempt(expected_id); + Some(gate.write_owned().await) + } None => None, }; let meta_to_save = { @@ -1061,6 +1152,46 @@ mod tests { use crate::object_api::NamespaceLockFence; use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause, hermetic_set_disks_isolated}; + #[tokio::test] + async fn cancel_rebalance_admission_is_id_checked_and_idempotent() { + let rebalance_id = "rebalance-admission-current"; + let cancel = tokio_util::sync::CancellationToken::new(); + let (_temp_dirs, store) = crate::services::rebalance::test_store_with_persisted_rebalance_meta(RebalanceMeta { + id: rebalance_id.to_string(), + percent_free_goal: 1.0, + cancel: Some(cancel.clone()), + pool_stats: vec![RebalanceStats { + participating: true, + init_capacity: 100, + info: RebalanceInfo { + start_time: Some(OffsetDateTime::now_utc()), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }) + .await; + + let error = store + .cancel_rebalance_admission_for_id("rebalance-admission-replacement") + .await + .expect_err("a stale admin stop must not cancel the current run"); + assert!(error.to_string().contains("expected rebalance-admission-replacement")); + assert!(!cancel.is_cancelled()); + + store + .cancel_rebalance_admission_for_id(rebalance_id) + .await + .expect("the current run admission should close"); + store + .cancel_rebalance_admission_for_id(rebalance_id) + .await + .expect("retrying admission cancellation should be idempotent"); + assert!(cancel.is_cancelled()); + } + async fn assert_real_activation_start_race(paused_kind: PoolActivationStartKind) { let (_temp_dirs, rebalance_store, decommission_store) = crate::services::rebalance::test_two_pool_stores(None).await; set_rebalance_disk_stats_override_for_test( diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 271da9133..2dd0d3f99 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -786,20 +786,267 @@ impl ECStore { } } +#[cfg(any(test, feature = "test-util"))] +pub mod test_util { + use super::super::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; + use super::*; + use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; + use crate::storage_api_contracts::object::ObjectOperations as _; + use rustfs_filemeta::FileMeta; + + pub struct PausedRebalanceEntryTestFixture { + _temp_dirs: Vec, + store: Arc, + cancel: CancellationToken, + barrier: crate::data_movement::SourceCleanupDeleteBarrier, + stop_probe: super::super::control::RebalanceStopWaitProbe, + entry_task: Option>>, + } + + impl PausedRebalanceEntryTestFixture { + pub async fn new(rebalance_id: &'static str) -> Self { + let cancel = CancellationToken::new(); + let (_temp_dirs, store) = super::super::test_store_with_persisted_rebalance_meta(RebalanceMeta { + id: rebalance_id.to_string(), + percent_free_goal: 1.0, + cancel: Some(cancel.clone()), + pool_stats: vec![RebalanceStats { + participating: true, + init_capacity: 100, + buckets: vec!["bucket".to_string()], + info: RebalanceInfo { + start_time: Some(OffsetDateTime::now_utc()), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }) + .await; + let bucket = "bucket"; + let object = "delete-marker"; + let set = store.pools[0].get_disks_by_key(object); + set.make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("source bucket should be created"); + set.delete_object( + bucket, + object, + ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("source delete marker should be created"); + let source_versions = set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata should be readable") + .expect("source delete marker should exist"); + assert_eq!(source_versions.versions.len(), 1); + assert!(source_versions.versions[0].deleted); + let mut file_meta = FileMeta::new(); + for version in source_versions.versions { + file_meta + .add_version(version) + .expect("source version should encode into a metacache entry"); + } + let entry = MetaCacheEntry { + name: object.to_string(), + metadata: file_meta.marshal_msg().expect("source metadata should marshal"), + cached: Some(file_meta), + reusable: false, + }; + let barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(bucket, object); + let entry_store = Arc::clone(&store); + let entry_set = Arc::clone(&set); + let entry_cancel = cancel.clone(); + let entry_task = tokio::spawn(async move { + entry_store + .rebalance_entry( + bucket.to_string(), + 0, + entry, + entry_set, + Arc::new(RebalanceBucketConfigs::default()), + Arc::from(rebalance_id), + entry_cancel, + ) + .await + }); + + Self { + _temp_dirs, + store, + cancel, + barrier, + stop_probe: super::super::control::RebalanceStopWaitProbe::install(rebalance_id), + entry_task: Some(entry_task), + } + } + + pub fn store(&self) -> Arc { + Arc::clone(&self.store) + } + + pub async fn wait_until_entry_paused(&self) { + self.barrier.wait_until_paused().await; + } + + pub async fn wait_until_admission_cancelled(&self) { + tokio::time::timeout(std::time::Duration::from_secs(5), self.cancel.cancelled()) + .await + .expect("admin stop should close admission before waiting for the entry guard"); + } + + pub async fn wait_until_stop_waiting_for_entry(&self) { + self.stop_probe.wait_until_attempted().await; + } + + pub fn release_entry(&self) { + self.barrier.release(); + } + + pub async fn assert_entry_cancelled(&mut self) { + let error = tokio::time::timeout( + std::time::Duration::from_secs(5), + self.entry_task.take().expect("entry task should still be available"), + ) + .await + .expect("the cancelled entry should finish") + .expect("entry task should not panic") + .expect_err("the entry must observe stop cancellation"); + assert!(matches!(error, Error::OperationCanceled)); + } + } + + impl Drop for PausedRebalanceEntryTestFixture { + fn drop(&mut self) { + self.barrier.release(); + if let Some(task) = self.entry_task.take() { + task.abort(); + } + } + } +} + #[cfg(test)] mod tests { use super::*; + use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE; + use crate::client::api_s3_datatypes::CompletePart; use crate::object_api::PutObjReader; use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; - use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause}; - use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; + use crate::set_disk::{ + DeleteObjectCommitBarrier, MultipartCommitBarrier, MultipartCommitPause, PutObjectCommitBarrier, PutObjectCommitPause, + TieredMetadataCommitBarrier, + }; + use crate::storage_api_contracts::multipart::MultipartOperations as _; use crate::storage_api_contracts::object::ObjectOperations as _; use http::HeaderMap; - use rustfs_filemeta::{FileInfo, FileMeta}; + use rustfs_filemeta::{FileInfo, FileMeta, ObjectPartInfo, TransitionVersionState}; + use s3s::dto::{BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule}; use std::time::Duration as StdDuration; use time::OffsetDateTime; use tokio::io::AsyncReadExt; + fn active_rebalance_meta(rebalance_id: &str) -> RebalanceMeta { + RebalanceMeta { + id: rebalance_id.to_string(), + percent_free_goal: 1.0, + cancel: Some(CancellationToken::new()), + pool_stats: vec![ + RebalanceStats { + participating: true, + init_capacity: 100, + buckets: vec![crate::disk::RUSTFS_META_BUCKET.to_string()], + info: RebalanceInfo { + start_time: Some(OffsetDateTime::now_utc()), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + RebalanceStats::default(), + ], + ..Default::default() + } + } + + fn expired_delete_marker_lifecycle() -> BucketLifecycleConfiguration { + BucketLifecycleConfiguration { + expiry_updated_at: None, + rules: vec![LifecycleRule { + status: ExpirationStatus::from_static(ExpirationStatus::ENABLED), + expiration: Some(LifecycleExpiration { + expired_object_delete_marker: Some(true), + ..Default::default() + }), + abort_incomplete_multipart_upload: None, + del_marker_expiration: None, + filter: None, + id: Some("expired-marker".to_string()), + noncurrent_version_expiration: None, + noncurrent_version_transitions: None, + prefix: None, + transitions: None, + }], + } + } + + async fn metacache_entry_from_source(set: &SetDisks, bucket: &str, object: &str) -> MetaCacheEntry { + let source_versions = set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata should be readable") + .expect("source version should exist"); + let mut file_meta = FileMeta::new(); + for version in source_versions.versions { + file_meta + .add_version(version) + .expect("source version should encode into a metacache entry"); + } + MetaCacheEntry { + name: object.to_string(), + metadata: file_meta.marshal_msg().expect("source metadata should marshal"), + cached: Some(file_meta), + reusable: false, + } + } + + fn spawn_real_rebalance_entry( + store: Arc, + set: Arc, + entry: MetaCacheEntry, + rebalance_id: &'static str, + bucket_configs: Arc, + ) -> tokio::task::JoinHandle> { + tokio::spawn(async move { + store + .rebalance_entry( + crate::disk::RUSTFS_META_BUCKET.to_string(), + 0, + entry, + set, + bucket_configs, + Arc::from(rebalance_id), + CancellationToken::new(), + ) + .await + }) + } + + async fn assert_real_entry_rejected_after_run_fence_loss(task: tokio::task::JoinHandle>) { + let error = tokio::time::timeout(StdDuration::from_secs(5), task) + .await + .expect("fenced real rebalance entry should finish") + .expect("rebalance entry task should not panic") + .expect_err("lost run fence must reject the entry commit"); + assert!(error.to_string().contains("run fence lost"), "unexpected fence error: {error}"); + } + #[tokio::test] async fn rebalance_stats_wait_for_source_cleanup_result() { let rebalance_id = "rebalance-test"; @@ -940,118 +1187,6 @@ mod tests { assert_eq!(pool_stats.cleanup_warnings.count, 1, "deferred cleanup must not add a permanent warning"); } - #[tokio::test] - #[serial_test::serial] - async fn stop_cancels_real_rebalance_entry_before_waiting_for_cleanup_guard() { - let rebalance_id = "rebalance-stop-entry"; - let cancel = CancellationToken::new(); - let (_temp_dirs, store) = crate::services::rebalance::test_store_with_persisted_rebalance_meta(RebalanceMeta { - id: rebalance_id.to_string(), - percent_free_goal: 1.0, - cancel: Some(cancel.clone()), - pool_stats: vec![RebalanceStats { - participating: true, - init_capacity: 100, - buckets: vec!["bucket".to_string()], - info: RebalanceInfo { - start_time: Some(OffsetDateTime::now_utc()), - status: RebalStatus::Started, - ..Default::default() - }, - ..Default::default() - }], - ..Default::default() - }) - .await; - let bucket = "bucket"; - let object = "delete-marker"; - let set = store.pools[0].get_disks_by_key(object); - set.make_bucket(bucket, &MakeBucketOptions::default()) - .await - .expect("source bucket should be created"); - set.delete_object( - bucket, - object, - ObjectOptions { - versioned: true, - ..Default::default() - }, - ) - .await - .expect("source delete marker should be created"); - let source_versions = set - .load_file_info_versions_exact(bucket, object) - .await - .expect("source metadata should be readable") - .expect("source delete marker should exist"); - assert_eq!(source_versions.versions.len(), 1); - assert!(source_versions.versions[0].deleted); - let mut file_meta = FileMeta::new(); - for version in &source_versions.versions { - file_meta - .add_version(version.clone()) - .expect("source version should encode into a metacache entry"); - } - let entry = MetaCacheEntry { - name: object.to_string(), - metadata: file_meta.marshal_msg().expect("source metadata should marshal"), - cached: Some(file_meta), - reusable: false, - }; - let barrier = data_movement::SourceCleanupDeleteBarrier::install(bucket, object); - let entry_store = Arc::clone(&store); - let entry_set = Arc::clone(&set); - let entry_cancel = cancel.clone(); - let mut entry_task = tokio::spawn(async move { - entry_store - .rebalance_entry( - bucket.to_string(), - 0, - entry, - entry_set, - Arc::new(RebalanceBucketConfigs::default()), - Arc::from(rebalance_id), - entry_cancel, - ) - .await - }); - barrier.wait_until_paused().await; - - let stop_store = Arc::clone(&store); - let mut stop_task = tokio::spawn(async move { stop_store.stop_rebalance_for_id(Some(rebalance_id)).await }); - tokio::time::timeout(StdDuration::from_secs(1), cancel.cancelled()) - .await - .expect("stop must cancel the in-flight entry before waiting for its guard"); - assert!( - tokio::time::timeout(StdDuration::from_millis(50), &mut stop_task) - .await - .is_err(), - "stop must wait for the real entry cleanup guard to drain" - ); - - barrier.release(); - let entry_error = tokio::time::timeout(StdDuration::from_secs(5), &mut entry_task) - .await - .expect("the cancelled entry should finish in bounded time") - .expect("entry task should not panic") - .expect_err("the entry must observe stop cancellation"); - assert!(matches!(entry_error, Error::OperationCanceled)); - tokio::time::timeout(StdDuration::from_secs(5), &mut stop_task) - .await - .expect("stop should finish after the entry guard drains") - .expect("stop task should not panic") - .expect("stop should persist the terminal state"); - assert!( - store - .rebalance_meta - .read() - .await - .as_ref() - .is_some_and(|meta| meta.stopped_at.is_some()), - "stop should publish the terminal state after entry cleanup drains" - ); - } - #[tokio::test] #[serial_test::serial] async fn real_rebalance_run_fence_loss_before_target_commit_preserves_target_and_source() { @@ -1185,4 +1320,323 @@ mod tests { .expect("source body should drain"); assert_eq!(source_body, payload, "source version must remain byte-identical"); } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_run_fence_loss_blocks_multipart_completion() { + const REBALANCE_ID: &str = "rebalance-multipart-commit-fence"; + let bucket = crate::disk::RUSTFS_META_BUCKET; + let object = "rebalance-multipart-commit-fence-object"; + let (_temp_dirs, store, _unused_store) = + crate::services::rebalance::test_two_pool_stores(Some(active_rebalance_meta(REBALANCE_ID))).await; + let source_set = store.pools[0].get_disks_by_key(object); + let target_set = store.pools[1].get_disks_by_key(object); + let upload = source_set + .new_multipart_upload(bucket, object, &ObjectOptions::default()) + .await + .expect("source multipart upload should be created"); + let mut reader = PutObjReader::from_vec(b"multipart source payload".repeat(1024)); + let part = source_set + .put_object_part(bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default()) + .await + .expect("source multipart part should be written"); + source_set + .clone() + .complete_multipart_upload( + bucket, + object, + &upload.upload_id, + vec![CompletePart { + part_num: part.part_num, + etag: part.etag, + ..Default::default() + }], + &ObjectOptions::default(), + ) + .await + .expect("source multipart object should commit"); + + let entry = metacache_entry_from_source(source_set.as_ref(), bucket, object).await; + let run_signal_fence = RebalanceRunSignalTestFence::install(REBALANCE_ID); + let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::BeforeQuotaRename); + let task = spawn_real_rebalance_entry( + Arc::clone(&store), + Arc::clone(&source_set), + entry, + REBALANCE_ID, + Arc::new(RebalanceBucketConfigs::default()), + ); + barrier.wait_until_paused().await; + run_signal_fence.mark_lost(); + barrier.release(); + drop(barrier); + assert_real_entry_rejected_after_run_fence_loss(task).await; + + assert!( + target_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("target metadata lookup should succeed") + .is_none(), + "lost run fence must not publish the target multipart object" + ); + assert!( + source_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata lookup should succeed") + .is_some(), + "lost run fence must preserve the source multipart object" + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_run_fence_loss_blocks_remote_tier_metadata_commit() { + const REBALANCE_ID: &str = "rebalance-tiered-commit-fence"; + let bucket = crate::disk::RUSTFS_META_BUCKET; + let object = "rebalance-tiered-commit-fence-object"; + let (_temp_dirs, store, _unused_store) = + crate::services::rebalance::test_two_pool_stores(Some(active_rebalance_meta(REBALANCE_ID))).await; + let source_set = store.pools[0].get_disks_by_key(object); + let target_set = store.pools[1].get_disks_by_key(object); + let version_id = uuid::Uuid::new_v4(); + let source = FileInfo { + volume: bucket.to_string(), + name: object.to_string(), + version_id: Some(version_id), + mod_time: Some(OffsetDateTime::now_utc()), + size: 32, + parts: vec![ObjectPartInfo { + number: 1, + size: 32, + actual_size: 32, + etag: "tiered-part".to_string(), + ..Default::default() + }], + transition_status: TRANSITION_COMPLETE.to_string(), + transition_tier: "WARM".to_string(), + transitioned_objname: "remote/rebalance-tiered-commit-fence-object".to_string(), + transition_version: Some("remote-version".to_string()), + transition_version_state: TransitionVersionState::Exact, + fresh: true, + ..Default::default() + }; + source_set + .decommission_tiered_object( + bucket, + object, + &source, + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: source.mod_time, + ..Default::default() + }, + ) + .await + .expect("source tiered metadata should be written"); + + let entry = metacache_entry_from_source(source_set.as_ref(), bucket, object).await; + let run_signal_fence = RebalanceRunSignalTestFence::install(REBALANCE_ID); + let barrier = TieredMetadataCommitBarrier::install(bucket, object); + let task = spawn_real_rebalance_entry( + Arc::clone(&store), + Arc::clone(&source_set), + entry, + REBALANCE_ID, + Arc::new(RebalanceBucketConfigs::default()), + ); + barrier.wait_until_paused().await; + run_signal_fence.mark_lost(); + barrier.release(); + drop(barrier); + assert_real_entry_rejected_after_run_fence_loss(task).await; + + assert!( + target_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("target metadata lookup should succeed") + .is_none(), + "lost run fence must not publish target tier metadata" + ); + assert!( + source_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata lookup should succeed") + .is_some(), + "lost run fence must preserve source tier metadata" + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_run_fence_loss_blocks_delete_marker_commit() { + const REBALANCE_ID: &str = "rebalance-delete-marker-commit-fence"; + let bucket = crate::disk::RUSTFS_META_BUCKET; + let object = "rebalance-delete-marker-commit-fence-object"; + let (_temp_dirs, store, _unused_store) = + crate::services::rebalance::test_two_pool_stores(Some(active_rebalance_meta(REBALANCE_ID))).await; + let source_set = store.pools[0].get_disks_by_key(object); + let target_set = store.pools[1].get_disks_by_key(object); + let mut reader = PutObjReader::from_vec(b"source version beneath delete marker".to_vec()); + source_set + .put_object( + bucket, + object, + &mut reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("source version should be written"); + source_set + .delete_object( + bucket, + object, + ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("source delete marker should be written"); + + let entry = metacache_entry_from_source(source_set.as_ref(), bucket, object).await; + let run_signal_fence = RebalanceRunSignalTestFence::install(REBALANCE_ID); + let barrier = DeleteObjectCommitBarrier::install(bucket, object); + let task = spawn_real_rebalance_entry( + Arc::clone(&store), + Arc::clone(&source_set), + entry, + REBALANCE_ID, + Arc::new(RebalanceBucketConfigs::default()), + ); + barrier.wait_until_paused().await; + run_signal_fence.mark_lost(); + barrier.release(); + drop(barrier); + assert_real_entry_rejected_after_run_fence_loss(task).await; + + assert!( + target_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("target metadata lookup should succeed") + .is_none(), + "lost run fence must not publish the target delete marker" + ); + assert_eq!( + source_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata lookup should succeed") + .expect("source versions should remain") + .versions + .len(), + 2, + "lost run fence must preserve both source versions" + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_run_fence_loss_blocks_lifecycle_mutation_dispatch() { + const REBALANCE_ID: &str = "rebalance-lifecycle-mutation-fence"; + let bucket = crate::disk::RUSTFS_META_BUCKET; + let object = "rebalance-lifecycle-mutation-fence-object"; + let (_temp_dirs, store, _unused_store) = + crate::services::rebalance::test_two_pool_stores(Some(active_rebalance_meta(REBALANCE_ID))).await; + let source_set = store.pools[0].get_disks_by_key(object); + source_set + .delete_object( + bucket, + object, + ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("source delete marker should be written"); + + let entry = metacache_entry_from_source(source_set.as_ref(), bucket, object).await; + let run_signal_fence = RebalanceRunSignalTestFence::install(REBALANCE_ID); + let barrier = crate::core::pools::LifecycleDataMovementMutationBarrier::install(bucket, object); + let task = spawn_real_rebalance_entry( + Arc::clone(&store), + Arc::clone(&source_set), + entry, + REBALANCE_ID, + Arc::new(RebalanceBucketConfigs { + lifecycle_config: Some(expired_delete_marker_lifecycle()), + ..Default::default() + }), + ); + barrier.wait_until_paused().await; + run_signal_fence.mark_lost(); + barrier.release(); + drop(barrier); + assert_real_entry_rejected_after_run_fence_loss(task).await; + + assert!( + source_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata lookup should succeed") + .is_some(), + "lost run fence must preserve the source before lifecycle mutation" + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_run_fence_loss_blocks_source_cleanup_delete() { + const REBALANCE_ID: &str = "rebalance-source-cleanup-fence"; + let bucket = crate::disk::RUSTFS_META_BUCKET; + let object = "rebalance-source-cleanup-fence-object"; + let (_temp_dirs, store, _unused_store) = + crate::services::rebalance::test_two_pool_stores(Some(active_rebalance_meta(REBALANCE_ID))).await; + let source_set = store.pools[0].get_disks_by_key(object); + source_set + .delete_object( + bucket, + object, + ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("source delete marker should be written"); + + let entry = metacache_entry_from_source(source_set.as_ref(), bucket, object).await; + let run_signal_fence = RebalanceRunSignalTestFence::install(REBALANCE_ID); + let barrier = data_movement::SourceCleanupDeleteBarrier::install(bucket, object); + let task = spawn_real_rebalance_entry( + Arc::clone(&store), + Arc::clone(&source_set), + entry, + REBALANCE_ID, + Arc::new(RebalanceBucketConfigs::default()), + ); + barrier.wait_until_paused().await; + run_signal_fence.mark_lost(); + barrier.release(); + drop(barrier); + assert_real_entry_rejected_after_run_fence_loss(task).await; + + assert!( + source_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata lookup should succeed") + .is_some(), + "lost run fence must preserve the source during cleanup" + ); + } } diff --git a/crates/ecstore/src/services/rebalance/mod.rs b/crates/ecstore/src/services/rebalance/mod.rs index 4e44727e3..bc3fdeb3b 100644 --- a/crates/ecstore/src/services/rebalance/mod.rs +++ b/crates/ecstore/src/services/rebalance/mod.rs @@ -53,7 +53,7 @@ pub use types::{ }; use types::{RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome}; -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] pub(crate) async fn test_store_with_persisted_rebalance_meta( meta: RebalanceMeta, ) -> (Vec, std::sync::Arc) { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 4d1d56578..078c4a5d5 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -740,6 +740,8 @@ mod ops; pub(crate) use ops::hermetic_set_disks_isolated; #[cfg(any(test, feature = "test-util"))] pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause}; +#[cfg(test)] +pub(crate) use ops::object::DeleteObjectCommitBarrier; #[cfg(feature = "test-util")] pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier; pub(crate) use ops::object::body_cache_plaintext_len; @@ -1448,6 +1450,21 @@ mod prepared_get_object_metadata_tests { } impl SetDisks { + #[cfg(test)] + async fn pause_tiered_metadata_commit(bucket: &str, object: &str) { + let barrier = TIERED_METADATA_COMMIT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("tiered metadata commit barrier should not be poisoned") + .as_ref() + .filter(|barrier| barrier.bucket == bucket && barrier.object == object) + .cloned(); + if let Some(barrier) = barrier { + barrier.arrived.notify_one(); + barrier.release.notified().await; + } + } + pub(crate) async fn prepare_get_object_metadata( &self, bucket: &str, @@ -4697,6 +4714,8 @@ impl SetDisks { )?; let fi = build_tiered_decommission_file_info(bucket, object, fi, layout); let write_quorum = layout.write_quorum; + #[cfg(test)] + Self::pause_tiered_metadata_commit(bucket, object).await; if _lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) || opts .namespace_lock_fence @@ -4744,6 +4763,66 @@ impl SetDisks { } } +#[cfg(test)] +struct TieredMetadataCommitBarrierState { + bucket: String, + object: String, + arrived: tokio::sync::Notify, + release: tokio::sync::Notify, +} + +#[cfg(test)] +pub(crate) struct TieredMetadataCommitBarrier { + state: Arc, +} + +#[cfg(test)] +static TIERED_METADATA_COMMIT_BARRIER: std::sync::OnceLock>>> = + std::sync::OnceLock::new(); + +#[cfg(test)] +impl TieredMetadataCommitBarrier { + pub(crate) fn install(bucket: &str, object: &str) -> Self { + let state = Arc::new(TieredMetadataCommitBarrierState { + bucket: bucket.to_string(), + object: object.to_string(), + arrived: tokio::sync::Notify::new(), + release: tokio::sync::Notify::new(), + }); + let mut slot = TIERED_METADATA_COMMIT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("tiered metadata commit barrier should not be poisoned"); + assert!(slot.is_none(), "tiered metadata commit barrier must be unique"); + *slot = Some(Arc::clone(&state)); + Self { state } + } + + pub(crate) async fn wait_until_paused(&self) { + tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified()) + .await + .expect("tiered metadata write should reach its deterministic commit barrier"); + } + + pub(crate) fn release(&self) { + self.state.release.notify_one(); + } +} + +#[cfg(test)] +impl Drop for TieredMetadataCommitBarrier { + fn drop(&mut self) { + self.state.release.notify_one(); + let mut slot = TIERED_METADATA_COMMIT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("tiered metadata commit barrier should not be poisoned"); + if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *slot = None; + } + } +} + #[derive(Debug, PartialEq, Eq)] struct ObjProps { successor_mod_time: Option, diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 570a405a1..4a0eaee59 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -4750,7 +4750,7 @@ struct DeleteObjectCommitBarrierState { } #[cfg(test)] -struct DeleteObjectCommitBarrier { +pub(crate) struct DeleteObjectCommitBarrier { state: Arc, } @@ -4760,7 +4760,7 @@ static DELETE_OBJECT_COMMIT_BARRIER: std::sync::OnceLock Self { + pub(crate) fn install(bucket: &str, object: &str) -> Self { let state = Arc::new(DeleteObjectCommitBarrierState { bucket: bucket.to_string(), object: object.to_string(), @@ -4776,13 +4776,13 @@ impl DeleteObjectCommitBarrier { Self { state } } - async fn wait_until_paused(&self) { + pub(crate) async fn wait_until_paused(&self) { tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified()) .await .expect("delete object should reach the deterministic commit barrier"); } - fn release(&self) { + pub(crate) fn release(&self) { self.state.release.notify_one(); } } @@ -6453,6 +6453,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }; let find_vid = Uuid::new_v4(); + #[cfg(test)] + pause_delete_object_commit(bucket, object).await; if mark_delete && (opts.versioned || opts.version_suspended) { if !delete_marker { @@ -6521,8 +6523,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { dfi.set_skip_tier_free_version(); } - #[cfg(test)] - pause_delete_object_commit(bucket, object).await; ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?; self.delete_object_version(bucket, object, &dfi, opts.delete_marker) .await diff --git a/rustfs/src/admin/handlers/rebalance.rs b/rustfs/src/admin/handlers/rebalance.rs index 775f099cd..f0edb8fec 100644 --- a/rustfs/src/admin/handlers/rebalance.rs +++ b/rustfs/src/admin/handlers/rebalance.rs @@ -872,6 +872,34 @@ impl Operation for RebalanceStatus { } } +async fn stop_rebalance_admission_first( + store: &Arc, + notification_sys: Option<&NotificationSys>, + expected_rebalance_id: &str, +) -> S3Result> { + store + .cancel_rebalance_admission_for_id(expected_rebalance_id) + .await + .map_err(|e| s3_error!(InternalError, "failed to close rebalance admission before stop: {}", e))?; + + if let Some(notification_sys) = notification_sys { + return notification_sys + .stop_rebalance_failures(Some(expected_rebalance_id)) + .await + .map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e)); + } + + store + .stop_rebalance_for_id(Some(expected_rebalance_id)) + .await + .map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?; + store + .save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_rebalance_id) + .await + .map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?; + Ok(Vec::new()) +} + // RebalanceStop pub struct RebalanceStop {} @@ -919,10 +947,6 @@ impl Operation for RebalanceStop { return Err(s3_error!(InternalError, "object layer is not initialized")); }; - store - .load_rebalance_meta() - .await - .map_err(|e| s3_error!(InternalError, "failed to load rebalance metadata before stop: {}", e))?; let expected_rebalance_id = store.current_rebalance_id().await; if !store.is_rebalance_conflicting_with_decommission().await { @@ -935,23 +959,8 @@ impl Operation for RebalanceStop { let notification_sys = current_notification_system(); let stop_attempt_at = OffsetDateTime::now_utc(); - let mut stop_failures = Vec::new(); - if let Some(notification_sys) = notification_sys.as_ref() { - stop_failures = notification_sys - .stop_rebalance_failures(Some(expected_rebalance_id.as_str())) - .await - .map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e))?; - } else { - store - .stop_rebalance_for_id(Some(expected_rebalance_id.as_str())) - .await - .map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?; - - store - .save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_rebalance_id.as_str()) - .await - .map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?; - } + let stop_failures = + stop_rebalance_admission_first(&store, notification_sys.as_deref(), expected_rebalance_id.as_str()).await?; info!( event = EVENT_ADMIN_REBALANCE_STATE, @@ -1088,7 +1097,7 @@ mod rebalance_handler_tests { build_rebalance_admin_status, build_rebalance_pool_statuses, build_rebalance_stop_propagation_status, rebalance_pool_used, rebalance_query_present, rebalance_remaining_buckets, rebalance_rollback_failure_message, rebalance_rollback_stop_failure_message, rebalance_start_rollback_error, rebalance_start_steps, rebalance_used_pct, - rollback_result_label, + rollback_result_label, stop_rebalance_admission_first, }; use crate::admin::storage_api::rebalance::{ DiskStat, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta, @@ -1096,6 +1105,30 @@ mod rebalance_handler_tests { }; use time::OffsetDateTime; + #[tokio::test] + #[serial_test::serial] + async fn real_admin_stop_cancels_paused_entry_before_waiting_for_activation_gate() { + const REBALANCE_ID: &str = "admin-stop-paused-entry"; + let mut fixture = rustfs_ecstore::api::rebalance::test_util::PausedRebalanceEntryTestFixture::new(REBALANCE_ID).await; + fixture.wait_until_entry_paused().await; + + let stop_store = fixture.store(); + let mut stop_task = tokio::spawn(async move { stop_rebalance_admission_first(&stop_store, None, REBALANCE_ID).await }); + fixture.wait_until_admission_cancelled().await; + fixture.wait_until_stop_waiting_for_entry().await; + assert!(!stop_task.is_finished(), "admin stop must wait for the paused entry guard to drain"); + + fixture.release_entry(); + fixture.assert_entry_cancelled().await; + let stop_failures = tokio::time::timeout(std::time::Duration::from_secs(5), &mut stop_task) + .await + .expect("admin stop should finish after the entry guard drains") + .expect("admin stop task should not panic") + .expect("admin stop should persist the terminal state"); + assert!(stop_failures.is_empty()); + assert!(!fixture.store().is_rebalance_conflicting_with_decommission().await); + } + #[test] fn test_calculate_rebalance_progress_running() { let start = OffsetDateTime::from_unix_timestamp(1_000).unwrap();