mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-25 05:26:50 +00:00
fix(rebalance): cancel admin stop before activation wait
This commit is contained in:
@@ -440,6 +440,11 @@ pub mod rebalance {
|
|||||||
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
|
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
|
||||||
encode_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 {
|
pub mod rio {
|
||||||
|
|||||||
@@ -2500,6 +2500,85 @@ fn lifecycle_action_skips_heal_version(action: IlmAction) -> bool {
|
|||||||
action.delete()
|
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<LifecycleDataMovementMutationBarrierState>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
static LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER: std::sync::OnceLock<
|
||||||
|
std::sync::Mutex<Option<Arc<LifecycleDataMovementMutationBarrierState>>>,
|
||||||
|
> = 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<bool> {
|
fn resolve_data_movement_lifecycle_expiry_result(action: IlmAction, apply_actions: bool, applied: bool) -> Result<bool> {
|
||||||
if !apply_actions || applied {
|
if !apply_actions || applied {
|
||||||
return Ok(true);
|
return Ok(true);
|
||||||
@@ -2548,6 +2627,8 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement(
|
|||||||
Ok(false)
|
Ok(false)
|
||||||
}
|
}
|
||||||
action if lifecycle_action_removes_data_movement_version(action) => {
|
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? {
|
if lifecycle_delete_all_versions_blocked_by_replication(store.clone(), bucket, &object_info.name, action).await? {
|
||||||
return Ok(false);
|
return Ok(false);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1226,12 +1226,12 @@ pub(crate) async fn make_local_two_set_sets() -> (Vec<tempfile::TempDir>, Arc<Se
|
|||||||
make_local_two_set_sets_with_ctx(bootstrap_ctx()).await
|
make_local_two_set_sets_with_ctx(bootstrap_ctx()).await
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>) -> (Vec<tempfile::TempDir>, Arc<Sets>) {
|
pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>) -> (Vec<tempfile::TempDir>, Arc<Sets>) {
|
||||||
make_local_two_set_sets_for_pool_with_ctx(ctx, 0).await
|
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(
|
pub(crate) async fn make_local_two_set_sets_for_pool_with_ctx(
|
||||||
ctx: Arc<InstanceContext>,
|
ctx: Arc<InstanceContext>,
|
||||||
pool_idx: usize,
|
pool_idx: usize,
|
||||||
|
|||||||
@@ -1062,7 +1062,7 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged(
|
|||||||
ensure_source_cleanup_versions_match(expected, ¤t, allowed_missing)
|
ensure_source_cleanup_versions_match(expected, ¤t, allowed_missing)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
struct SourceCleanupDeleteBarrierState {
|
struct SourceCleanupDeleteBarrierState {
|
||||||
bucket: String,
|
bucket: String,
|
||||||
object: String,
|
object: String,
|
||||||
@@ -1070,7 +1070,7 @@ struct SourceCleanupDeleteBarrierState {
|
|||||||
release: tokio::sync::Notify,
|
release: tokio::sync::Notify,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
#[allow(
|
#[allow(
|
||||||
dead_code,
|
dead_code,
|
||||||
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
|
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
|
||||||
@@ -1079,11 +1079,11 @@ pub(crate) struct SourceCleanupDeleteBarrier {
|
|||||||
state: Arc<SourceCleanupDeleteBarrierState>,
|
state: Arc<SourceCleanupDeleteBarrierState>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<SourceCleanupDeleteBarrierState>>>> =
|
static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<SourceCleanupDeleteBarrierState>>>> =
|
||||||
std::sync::OnceLock::new();
|
std::sync::OnceLock::new();
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
#[allow(
|
#[allow(
|
||||||
dead_code,
|
dead_code,
|
||||||
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
|
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 {
|
impl Drop for SourceCleanupDeleteBarrier {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
self.state.release.notify_one();
|
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) {
|
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
|
||||||
let barrier = SOURCE_CLEANUP_DELETE_BARRIER
|
let barrier = SOURCE_CLEANUP_DELETE_BARRIER
|
||||||
.get_or_init(|| std::sync::Mutex::new(None))
|
.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?;
|
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;
|
pause_source_cleanup_before_delete(bucket, object).await;
|
||||||
|
|
||||||
let mut opts = ObjectOptions {
|
let mut opts = ObjectOptions {
|
||||||
|
|||||||
@@ -82,6 +82,75 @@ pub(super) struct RebalanceRunGuard {
|
|||||||
persisted_guard: rustfs_lock::NamespaceLockGuard,
|
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::Mutex<Vec<Arc<RebalanceStopWaitProbeState>>>> =
|
||||||
|
std::sync::OnceLock::new();
|
||||||
|
|
||||||
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
|
pub(super) struct RebalanceStopWaitProbe {
|
||||||
|
state: Arc<RebalanceStopWaitProbeState>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[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 {
|
impl RebalanceRunGuard {
|
||||||
pub(super) fn ensure_held(&self, stage: &str) -> Result<()> {
|
pub(super) fn ensure_held(&self, stage: &str) -> Result<()> {
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -945,6 +1014,24 @@ impl ECStore {
|
|||||||
.and_then(|meta| (!meta.id.is_empty()).then(|| meta.id.clone()))
|
.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))]
|
#[tracing::instrument(skip(self))]
|
||||||
pub async fn stop_rebalance(self: &Arc<Self>) -> Result<()> {
|
pub async fn stop_rebalance(self: &Arc<Self>) -> Result<()> {
|
||||||
self.stop_rebalance_for_id(None).await
|
self.stop_rebalance_for_id(None).await
|
||||||
@@ -965,7 +1052,11 @@ impl ECStore {
|
|||||||
})
|
})
|
||||||
};
|
};
|
||||||
let _activation_guard = match activation_gate {
|
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,
|
None => None,
|
||||||
};
|
};
|
||||||
let meta_to_save = {
|
let meta_to_save = {
|
||||||
@@ -1061,6 +1152,46 @@ mod tests {
|
|||||||
use crate::object_api::NamespaceLockFence;
|
use crate::object_api::NamespaceLockFence;
|
||||||
use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause, hermetic_set_disks_isolated};
|
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) {
|
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;
|
let (_temp_dirs, rebalance_store, decommission_store) = crate::services::rebalance::test_two_pool_stores(None).await;
|
||||||
set_rebalance_disk_stats_override_for_test(
|
set_rebalance_disk_stats_override_for_test(
|
||||||
|
|||||||
@@ -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<tempfile::TempDir>,
|
||||||
|
store: Arc<ECStore>,
|
||||||
|
cancel: CancellationToken,
|
||||||
|
barrier: crate::data_movement::SourceCleanupDeleteBarrier,
|
||||||
|
stop_probe: super::super::control::RebalanceStopWaitProbe,
|
||||||
|
entry_task: Option<tokio::task::JoinHandle<Result<RebalanceEntryOutcome>>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<ECStore> {
|
||||||
|
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)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
||||||
|
use crate::client::api_s3_datatypes::CompletePart;
|
||||||
use crate::object_api::PutObjReader;
|
use crate::object_api::PutObjReader;
|
||||||
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats};
|
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats};
|
||||||
use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause};
|
use crate::set_disk::{
|
||||||
use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions};
|
DeleteObjectCommitBarrier, MultipartCommitBarrier, MultipartCommitPause, PutObjectCommitBarrier, PutObjectCommitPause,
|
||||||
|
TieredMetadataCommitBarrier,
|
||||||
|
};
|
||||||
|
use crate::storage_api_contracts::multipart::MultipartOperations as _;
|
||||||
use crate::storage_api_contracts::object::ObjectOperations as _;
|
use crate::storage_api_contracts::object::ObjectOperations as _;
|
||||||
use http::HeaderMap;
|
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 std::time::Duration as StdDuration;
|
||||||
use time::OffsetDateTime;
|
use time::OffsetDateTime;
|
||||||
use tokio::io::AsyncReadExt;
|
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<ECStore>,
|
||||||
|
set: Arc<SetDisks>,
|
||||||
|
entry: MetaCacheEntry,
|
||||||
|
rebalance_id: &'static str,
|
||||||
|
bucket_configs: Arc<RebalanceBucketConfigs>,
|
||||||
|
) -> tokio::task::JoinHandle<Result<RebalanceEntryOutcome>> {
|
||||||
|
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<Result<RebalanceEntryOutcome>>) {
|
||||||
|
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]
|
#[tokio::test]
|
||||||
async fn rebalance_stats_wait_for_source_cleanup_result() {
|
async fn rebalance_stats_wait_for_source_cleanup_result() {
|
||||||
let rebalance_id = "rebalance-test";
|
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");
|
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]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[serial_test::serial]
|
||||||
async fn real_rebalance_run_fence_loss_before_target_commit_preserves_target_and_source() {
|
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");
|
.expect("source body should drain");
|
||||||
assert_eq!(source_body, payload, "source version must remain byte-identical");
|
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"
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -53,7 +53,7 @@ pub use types::{
|
|||||||
};
|
};
|
||||||
use types::{RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome};
|
use types::{RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome};
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
pub(crate) async fn test_store_with_persisted_rebalance_meta(
|
pub(crate) async fn test_store_with_persisted_rebalance_meta(
|
||||||
meta: RebalanceMeta,
|
meta: RebalanceMeta,
|
||||||
) -> (Vec<tempfile::TempDir>, std::sync::Arc<crate::store::ECStore>) {
|
) -> (Vec<tempfile::TempDir>, std::sync::Arc<crate::store::ECStore>) {
|
||||||
|
|||||||
@@ -740,6 +740,8 @@ mod ops;
|
|||||||
pub(crate) use ops::hermetic_set_disks_isolated;
|
pub(crate) use ops::hermetic_set_disks_isolated;
|
||||||
#[cfg(any(test, feature = "test-util"))]
|
#[cfg(any(test, feature = "test-util"))]
|
||||||
pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
|
pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) use ops::object::DeleteObjectCommitBarrier;
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
||||||
pub(crate) use ops::object::body_cache_plaintext_len;
|
pub(crate) use ops::object::body_cache_plaintext_len;
|
||||||
@@ -1448,6 +1450,21 @@ mod prepared_get_object_metadata_tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl SetDisks {
|
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(
|
pub(crate) async fn prepare_get_object_metadata(
|
||||||
&self,
|
&self,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
@@ -4697,6 +4714,8 @@ impl SetDisks {
|
|||||||
)?;
|
)?;
|
||||||
let fi = build_tiered_decommission_file_info(bucket, object, fi, layout);
|
let fi = build_tiered_decommission_file_info(bucket, object, fi, layout);
|
||||||
let write_quorum = layout.write_quorum;
|
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())
|
if _lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
|
||||||
|| opts
|
|| opts
|
||||||
.namespace_lock_fence
|
.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<TieredMetadataCommitBarrierState>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
static TIERED_METADATA_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<TieredMetadataCommitBarrierState>>>> =
|
||||||
|
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)]
|
#[derive(Debug, PartialEq, Eq)]
|
||||||
struct ObjProps {
|
struct ObjProps {
|
||||||
successor_mod_time: Option<OffsetDateTime>,
|
successor_mod_time: Option<OffsetDateTime>,
|
||||||
|
|||||||
@@ -4750,7 +4750,7 @@ struct DeleteObjectCommitBarrierState {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
struct DeleteObjectCommitBarrier {
|
pub(crate) struct DeleteObjectCommitBarrier {
|
||||||
state: Arc<DeleteObjectCommitBarrierState>,
|
state: Arc<DeleteObjectCommitBarrierState>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -4760,7 +4760,7 @@ static DELETE_OBJECT_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
impl DeleteObjectCommitBarrier {
|
impl DeleteObjectCommitBarrier {
|
||||||
fn install(bucket: &str, object: &str) -> Self {
|
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
||||||
let state = Arc::new(DeleteObjectCommitBarrierState {
|
let state = Arc::new(DeleteObjectCommitBarrierState {
|
||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
object: object.to_string(),
|
object: object.to_string(),
|
||||||
@@ -4776,13 +4776,13 @@ impl DeleteObjectCommitBarrier {
|
|||||||
Self { state }
|
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())
|
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
|
||||||
.await
|
.await
|
||||||
.expect("delete object should reach the deterministic commit barrier");
|
.expect("delete object should reach the deterministic commit barrier");
|
||||||
}
|
}
|
||||||
|
|
||||||
fn release(&self) {
|
pub(crate) fn release(&self) {
|
||||||
self.state.release.notify_one();
|
self.state.release.notify_one();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -6453,6 +6453,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
};
|
};
|
||||||
|
|
||||||
let find_vid = Uuid::new_v4();
|
let find_vid = Uuid::new_v4();
|
||||||
|
#[cfg(test)]
|
||||||
|
pause_delete_object_commit(bucket, object).await;
|
||||||
|
|
||||||
if mark_delete && (opts.versioned || opts.version_suspended) {
|
if mark_delete && (opts.versioned || opts.version_suspended) {
|
||||||
if !delete_marker {
|
if !delete_marker {
|
||||||
@@ -6521,8 +6523,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
|||||||
dfi.set_skip_tier_free_version();
|
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)?;
|
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
|
||||||
self.delete_object_version(bucket, object, &dfi, opts.delete_marker)
|
self.delete_object_version(bucket, object, &dfi, opts.delete_marker)
|
||||||
.await
|
.await
|
||||||
|
|||||||
@@ -872,6 +872,34 @@ impl Operation for RebalanceStatus {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn stop_rebalance_admission_first(
|
||||||
|
store: &Arc<ECStore>,
|
||||||
|
notification_sys: Option<&NotificationSys>,
|
||||||
|
expected_rebalance_id: &str,
|
||||||
|
) -> S3Result<Vec<String>> {
|
||||||
|
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
|
// RebalanceStop
|
||||||
pub struct RebalanceStop {}
|
pub struct RebalanceStop {}
|
||||||
|
|
||||||
@@ -919,10 +947,6 @@ impl Operation for RebalanceStop {
|
|||||||
return Err(s3_error!(InternalError, "object layer is not initialized"));
|
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;
|
let expected_rebalance_id = store.current_rebalance_id().await;
|
||||||
|
|
||||||
if !store.is_rebalance_conflicting_with_decommission().await {
|
if !store.is_rebalance_conflicting_with_decommission().await {
|
||||||
@@ -935,23 +959,8 @@ impl Operation for RebalanceStop {
|
|||||||
|
|
||||||
let notification_sys = current_notification_system();
|
let notification_sys = current_notification_system();
|
||||||
let stop_attempt_at = OffsetDateTime::now_utc();
|
let stop_attempt_at = OffsetDateTime::now_utc();
|
||||||
let mut stop_failures = Vec::new();
|
let stop_failures =
|
||||||
if let Some(notification_sys) = notification_sys.as_ref() {
|
stop_rebalance_admission_first(&store, notification_sys.as_deref(), expected_rebalance_id.as_str()).await?;
|
||||||
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))?;
|
|
||||||
}
|
|
||||||
|
|
||||||
info!(
|
info!(
|
||||||
event = EVENT_ADMIN_REBALANCE_STATE,
|
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,
|
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_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,
|
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::{
|
use crate::admin::storage_api::rebalance::{
|
||||||
DiskStat, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta,
|
DiskStat, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta,
|
||||||
@@ -1096,6 +1105,30 @@ mod rebalance_handler_tests {
|
|||||||
};
|
};
|
||||||
use time::OffsetDateTime;
|
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]
|
#[test]
|
||||||
fn test_calculate_rebalance_progress_running() {
|
fn test_calculate_rebalance_progress_running() {
|
||||||
let start = OffsetDateTime::from_unix_timestamp(1_000).unwrap();
|
let start = OffsetDateTime::from_unix_timestamp(1_000).unwrap();
|
||||||
|
|||||||
Reference in New Issue
Block a user