mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
test(ecstore): exercise real rebalance fences
This commit is contained in:
@@ -19,8 +19,8 @@ use crate::bucket::{
|
||||
LifecycleExpiryConfigs,
|
||||
bucket_lifecycle_audit::LcEventSrc,
|
||||
bucket_lifecycle_ops::{
|
||||
LifecycleOps, apply_expiry_rule_for_data_movement, apply_expiry_rule_in, eval_action_from_lifecycle,
|
||||
lifecycle_delete_all_versions_blocked_by_replication,
|
||||
LifecycleOps, apply_expiry_on_transitioned_object, apply_expiry_rule_for_data_movement, apply_expiry_rule_in,
|
||||
eval_action_from_lifecycle, lifecycle_delete_all_versions_blocked_by_replication,
|
||||
},
|
||||
get_expiry_configs,
|
||||
lifecycle::IlmAction,
|
||||
@@ -914,6 +914,79 @@ where
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub(crate) enum PoolActivationStartKind {
|
||||
Rebalance,
|
||||
Decommission,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
struct PoolActivationStartProbeState {
|
||||
kind: PoolActivationStartKind,
|
||||
attempted: std::sync::atomic::AtomicBool,
|
||||
notify: tokio::sync::Notify,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
static POOL_ACTIVATION_START_PROBES: std::sync::OnceLock<std::sync::Mutex<Vec<Arc<PoolActivationStartProbeState>>>> =
|
||||
std::sync::OnceLock::new();
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) struct PoolActivationStartProbe {
|
||||
state: Arc<PoolActivationStartProbeState>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl PoolActivationStartProbe {
|
||||
pub(crate) fn install(kind: PoolActivationStartKind) -> Self {
|
||||
let state = Arc::new(PoolActivationStartProbeState {
|
||||
kind,
|
||||
attempted: std::sync::atomic::AtomicBool::new(false),
|
||||
notify: tokio::sync::Notify::new(),
|
||||
});
|
||||
POOL_ACTIVATION_START_PROBES
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("pool activation start probe should not be poisoned")
|
||||
.push(Arc::clone(&state));
|
||||
Self { state }
|
||||
}
|
||||
|
||||
pub(crate) async fn wait_until_attempted(&self) {
|
||||
while !self.state.attempted.load(Ordering::Acquire) {
|
||||
self.state.notify.notified().await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Drop for PoolActivationStartProbe {
|
||||
fn drop(&mut self) {
|
||||
let mut probes = POOL_ACTIVATION_START_PROBES
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("pool activation start probe should not be poisoned");
|
||||
probes.retain(|state| !Arc::ptr_eq(state, &self.state));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn observe_pool_activation_start_attempt(kind: PoolActivationStartKind) {
|
||||
let probes = POOL_ACTIVATION_START_PROBES
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("pool activation start probe should not be poisoned")
|
||||
.iter()
|
||||
.filter(|state| state.kind == kind)
|
||||
.cloned()
|
||||
.collect::<Vec<_>>();
|
||||
for state in probes {
|
||||
state.attempted.store(true, Ordering::Release);
|
||||
state.notify.notify_one();
|
||||
}
|
||||
}
|
||||
|
||||
fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) {
|
||||
*pool_meta = previous_pool_meta;
|
||||
}
|
||||
@@ -2462,7 +2535,15 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement(
|
||||
let Ok(bucket_incarnation_id) = store.bucket_incarnation_id_from_disk(bucket).await else {
|
||||
return Ok(false);
|
||||
};
|
||||
let _ = apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, lock_lost_signal).await;
|
||||
let _ = match lock_lost_signal {
|
||||
Some(signal) => {
|
||||
apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, Some(signal)).await
|
||||
}
|
||||
None => {
|
||||
apply_expiry_on_transitioned_object(store, &object_info, &event, event_source, bucket_incarnation_id)
|
||||
.await
|
||||
}
|
||||
};
|
||||
}
|
||||
Ok(false)
|
||||
}
|
||||
@@ -2471,7 +2552,12 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement(
|
||||
return Ok(false);
|
||||
}
|
||||
let applied = !apply_actions
|
||||
|| apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, lock_lost_signal).await;
|
||||
|| match lock_lost_signal {
|
||||
Some(signal) => {
|
||||
apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, Some(signal)).await
|
||||
}
|
||||
None => apply_expiry_rule_in(store, &event, event_source, &object_info).await,
|
||||
};
|
||||
resolve_data_movement_lifecycle_expiry_result(action, apply_actions, applied)
|
||||
}
|
||||
_ => Ok(false),
|
||||
@@ -2573,6 +2659,8 @@ impl ECStore {
|
||||
.first()
|
||||
.cloned()
|
||||
.ok_or_else(|| Error::other("decommission start rebalance metadata load failed: no storage pools available"))?;
|
||||
#[cfg(test)]
|
||||
observe_pool_activation_start_attempt(PoolActivationStartKind::Decommission);
|
||||
let activation_fence = acquire_pool_rebalance_activation_locks(rebalance_pool.clone()).await?;
|
||||
|
||||
let mut rebalance_meta = RebalanceMeta::new();
|
||||
@@ -4116,7 +4204,11 @@ impl ECStore {
|
||||
validate_start_decommission_request(&indices, self.single_pool())?;
|
||||
|
||||
self.ensure_decommission_rebalance_idle_after_refresh().await?;
|
||||
ensure_decommission_start_local_leader(&self.endpoints(), &indices)?;
|
||||
#[cfg(test)]
|
||||
let endpoints = self.instance_endpoints().unwrap_or_else(|| self.endpoints());
|
||||
#[cfg(not(test))]
|
||||
let endpoints = self.endpoints();
|
||||
ensure_decommission_start_local_leader(&endpoints, &indices)?;
|
||||
|
||||
for idx in indices.iter().copied() {
|
||||
ensure_valid_decommission_pool_index(self.pools.len(), idx)?;
|
||||
|
||||
@@ -1228,6 +1228,14 @@ pub(crate) async fn make_local_two_set_sets() -> (Vec<tempfile::TempDir>, Arc<Se
|
||||
|
||||
#[cfg(test)]
|
||||
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
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn make_local_two_set_sets_for_pool_with_ctx(
|
||||
ctx: Arc<InstanceContext>,
|
||||
pool_idx: usize,
|
||||
) -> (Vec<tempfile::TempDir>, Arc<Sets>) {
|
||||
use crate::layout::endpoint::Endpoint;
|
||||
use rustfs_lock::client::local::LocalClient;
|
||||
|
||||
@@ -1243,7 +1251,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
|
||||
let temp_dir = tempfile::tempdir().expect("tempdir should be created");
|
||||
let mut endpoint = Endpoint::try_from(temp_dir.path().to_str().expect("tempdir path should be utf8"))
|
||||
.expect("endpoint should parse");
|
||||
endpoint.set_pool_index(0);
|
||||
endpoint.set_pool_index(pool_idx);
|
||||
endpoint.set_set_index(set_index);
|
||||
endpoint.set_disk_index(disk_index);
|
||||
let disk = new_disk(
|
||||
@@ -1279,7 +1287,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
|
||||
2,
|
||||
1,
|
||||
set_index,
|
||||
0,
|
||||
pool_idx,
|
||||
endpoints,
|
||||
format.clone(),
|
||||
lockers,
|
||||
@@ -1292,7 +1300,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
|
||||
let sets = Arc::new(Sets {
|
||||
id: format.id,
|
||||
disk_set: disk_sets,
|
||||
pool_idx: 0,
|
||||
pool_idx,
|
||||
endpoints: PoolEndpoints {
|
||||
legacy: false,
|
||||
set_count: 2,
|
||||
|
||||
@@ -24,7 +24,7 @@ use crate::storage_api_contracts::{
|
||||
pub struct NamespaceLockFence {
|
||||
signals: Arc<Vec<Arc<rustfs_lock::distributed_lock::LockLostSignal>>>,
|
||||
#[cfg(test)]
|
||||
forced_lost: Arc<std::sync::atomic::AtomicBool>,
|
||||
forced_lost: Arc<Vec<Arc<std::sync::atomic::AtomicBool>>>,
|
||||
}
|
||||
|
||||
impl Debug for NamespaceLockFence {
|
||||
@@ -40,13 +40,17 @@ impl NamespaceLockFence {
|
||||
Self {
|
||||
signals: Arc::default(),
|
||||
#[cfg(test)]
|
||||
forced_lost: Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||
forced_lost: Arc::new(vec![Arc::new(std::sync::atomic::AtomicBool::new(false))]),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn is_lock_lost(&self) -> bool {
|
||||
#[cfg(test)]
|
||||
if self.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
|
||||
if self
|
||||
.forced_lost
|
||||
.iter()
|
||||
.any(|lost| lost.load(std::sync::atomic::Ordering::Acquire))
|
||||
{
|
||||
return true;
|
||||
}
|
||||
self.signals.iter().any(|signal| signal.is_lost())
|
||||
@@ -62,25 +66,80 @@ impl NamespaceLockFence {
|
||||
}
|
||||
Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned());
|
||||
#[cfg(test)]
|
||||
if other.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
|
||||
self.forced_lost.store(true, std::sync::atomic::Ordering::Release);
|
||||
if !Arc::ptr_eq(&self.forced_lost, &other.forced_lost) {
|
||||
Arc::make_mut(&mut self.forced_lost).extend(other.forced_lost.iter().cloned());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn lost_for_test() -> Self {
|
||||
let fence = Self::new();
|
||||
fence.forced_lost.store(true, std::sync::atomic::Ordering::Release);
|
||||
fence.forced_lost[0].store(true, std::sync::atomic::Ordering::Release);
|
||||
fence
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn loss_handle_for_test() -> (Self, Arc<std::sync::atomic::AtomicBool>) {
|
||||
let fence = Self::new();
|
||||
(fence.clone(), Arc::clone(&fence.forced_lost))
|
||||
(fence.clone(), Arc::clone(&fence.forced_lost[0]))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
static NAMESPACE_LOCK_SIGNAL_TEST_FENCES: std::sync::OnceLock<std::sync::Mutex<Vec<(usize, NamespaceLockFence)>>> =
|
||||
std::sync::OnceLock::new();
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) struct NamespaceLockSignalTestFence {
|
||||
signal_key: usize,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl NamespaceLockSignalTestFence {
|
||||
pub(crate) fn install_with_loss_handle(
|
||||
signal: &Arc<rustfs_lock::distributed_lock::LockLostSignal>,
|
||||
loss_handle: Arc<std::sync::atomic::AtomicBool>,
|
||||
) -> Self {
|
||||
let fence = NamespaceLockFence {
|
||||
signals: Arc::default(),
|
||||
forced_lost: Arc::new(vec![loss_handle]),
|
||||
};
|
||||
let signal_key = Arc::as_ptr(signal) as usize;
|
||||
let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("namespace lock signal test fence should not be poisoned");
|
||||
assert!(
|
||||
!fences.iter().any(|(key, _)| *key == signal_key),
|
||||
"namespace lock signal test fence must be unique"
|
||||
);
|
||||
fences.push((signal_key, fence));
|
||||
Self { signal_key }
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Drop for NamespaceLockSignalTestFence {
|
||||
fn drop(&mut self) {
|
||||
let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("namespace lock signal test fence should not be poisoned");
|
||||
fences.retain(|(key, _)| *key != self.signal_key);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn namespace_lock_signal_test_fence_is_lost(signal: &Arc<rustfs_lock::distributed_lock::LockLostSignal>) -> bool {
|
||||
NAMESPACE_LOCK_SIGNAL_TEST_FENCES
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("namespace lock signal test fence should not be poisoned")
|
||||
.iter()
|
||||
.find(|(key, _)| *key == Arc::as_ptr(signal) as usize)
|
||||
.is_some_and(|(_, fence)| fence.is_lock_lost())
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct ObjectLockConfigSnapshot {
|
||||
store_id: Option<Uuid>,
|
||||
@@ -402,9 +461,23 @@ impl ObjectOptions {
|
||||
}
|
||||
|
||||
pub(crate) fn add_namespace_lock_lost_signal(&mut self, signal: Arc<rustfs_lock::distributed_lock::LockLostSignal>) {
|
||||
#[cfg(test)]
|
||||
let test_fence = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("namespace lock signal test fence should not be poisoned")
|
||||
.iter()
|
||||
.find(|(key, _)| *key == Arc::as_ptr(&signal) as usize)
|
||||
.map(|(_, fence)| fence.clone());
|
||||
self.namespace_lock_fence
|
||||
.get_or_insert_with(NamespaceLockFence::new)
|
||||
.add_signal(signal);
|
||||
#[cfg(test)]
|
||||
if let Some(test_fence) = test_fence {
|
||||
self.namespace_lock_fence
|
||||
.get_or_insert_with(NamespaceLockFence::new)
|
||||
.extend(&test_fence);
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn ensure_namespace_lock_fence(&mut self) {
|
||||
|
||||
@@ -35,6 +35,10 @@ use uuid::Uuid;
|
||||
#[cfg(test)]
|
||||
static FAIL_NEXT_REBALANCE_ACTIVATION_SAVE: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
|
||||
|
||||
#[cfg(test)]
|
||||
static REBALANCE_DISK_STATS_OVERRIDES: std::sync::OnceLock<std::sync::Mutex<std::collections::HashMap<Uuid, Vec<DiskStat>>>> =
|
||||
std::sync::OnceLock::new();
|
||||
|
||||
#[cfg(test)]
|
||||
pub(super) fn fail_next_rebalance_activation_save_for_test(rebalance_id: &str) {
|
||||
*FAIL_NEXT_REBALANCE_ACTIVATION_SAVE
|
||||
@@ -42,6 +46,24 @@ pub(super) fn fail_next_rebalance_activation_save_for_test(rebalance_id: &str) {
|
||||
.expect("rebalance activation save failure hook should not be poisoned") = Some(rebalance_id.to_string());
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn set_rebalance_disk_stats_override_for_test(store_id: Uuid, disk_stats: Vec<DiskStat>) {
|
||||
REBALANCE_DISK_STATS_OVERRIDES
|
||||
.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
|
||||
.lock()
|
||||
.expect("rebalance disk stats override should not be poisoned")
|
||||
.insert(store_id, disk_stats);
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn take_rebalance_disk_stats_override_for_test(store_id: Uuid) -> Option<Vec<DiskStat>> {
|
||||
REBALANCE_DISK_STATS_OVERRIDES
|
||||
.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
|
||||
.lock()
|
||||
.expect("rebalance disk stats override should not be poisoned")
|
||||
.remove(&store_id)
|
||||
}
|
||||
|
||||
fn ensure_rebalance_activation_pool_meta_allowed(meta: &PoolMeta) -> Result<()> {
|
||||
if pool_meta_has_active_decommission(meta) {
|
||||
return Err(Error::DecommissionAlreadyRunning);
|
||||
@@ -62,7 +84,14 @@ pub(super) struct RebalanceRunGuard {
|
||||
|
||||
impl RebalanceRunGuard {
|
||||
pub(super) fn ensure_held(&self, stage: &str) -> Result<()> {
|
||||
if self.persisted_guard.is_lock_lost() {
|
||||
#[cfg(test)]
|
||||
let forced_lost = self
|
||||
.persisted_guard
|
||||
.lock_lost_signal()
|
||||
.is_some_and(|signal| crate::object_api::namespace_lock_signal_test_fence_is_lost(&signal));
|
||||
#[cfg(not(test))]
|
||||
let forced_lost = false;
|
||||
if self.persisted_guard.is_lock_lost() || forced_lost {
|
||||
return Err(Error::other(format!("rebalance distributed run fence lost during {stage}")));
|
||||
}
|
||||
Ok(())
|
||||
@@ -368,6 +397,8 @@ impl ECStore {
|
||||
where
|
||||
S: EcstoreObjectIO + StorageNamespaceLocking<Error = Error, NamespaceLock = rustfs_lock::NamespaceLockWrapper>,
|
||||
{
|
||||
#[cfg(test)]
|
||||
crate::core::pools::observe_pool_activation_start_attempt(crate::core::pools::PoolActivationStartKind::Rebalance);
|
||||
let activation_fence = acquire_pool_rebalance_activation_locks(pool.clone()).await?;
|
||||
let mut pool_meta = PoolMeta::default();
|
||||
pool_meta.load_no_lock(pool.clone()).await?;
|
||||
@@ -591,6 +622,13 @@ impl ECStore {
|
||||
disk_stats[disk.pool_index as usize].available_space += disk.available_space;
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
if let Some(overridden) = take_rebalance_disk_stats_override_for_test(self.id) {
|
||||
disk_stats = overridden;
|
||||
total_cap = disk_stats.iter().map(|stat| stat.total_space).sum();
|
||||
total_free = disk_stats.iter().map(|stat| stat.available_space).sum();
|
||||
}
|
||||
|
||||
let percent_free_goal = percent_free_ratio(total_free, total_cap);
|
||||
validate_rebalance_disk_stats_coverage(&disk_stats)?;
|
||||
|
||||
@@ -1019,8 +1057,153 @@ impl ECStore {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::core::pools::{PoolActivationStartKind, PoolActivationStartProbe};
|
||||
use crate::object_api::NamespaceLockFence;
|
||||
use crate::set_disk::hermetic_set_disks_isolated;
|
||||
use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause, hermetic_set_disks_isolated};
|
||||
|
||||
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(
|
||||
rebalance_store.id,
|
||||
vec![
|
||||
DiskStat {
|
||||
total_space: 100,
|
||||
available_space: 0,
|
||||
},
|
||||
DiskStat {
|
||||
total_space: 100,
|
||||
available_space: 100,
|
||||
},
|
||||
],
|
||||
);
|
||||
let (first_object, competing_object, competing_kind) = match paused_kind {
|
||||
PoolActivationStartKind::Rebalance => {
|
||||
(REBAL_META_NAME, crate::core::pools::POOL_META_NAME, PoolActivationStartKind::Decommission)
|
||||
}
|
||||
PoolActivationStartKind::Decommission => {
|
||||
(crate::core::pools::POOL_META_NAME, REBAL_META_NAME, PoolActivationStartKind::Rebalance)
|
||||
}
|
||||
};
|
||||
let first_barrier = PutObjectCommitBarrier::install(
|
||||
crate::disk::RUSTFS_META_BUCKET,
|
||||
first_object,
|
||||
PutObjectCommitPause::BeforeQuotaRename,
|
||||
);
|
||||
let competing_barrier = PutObjectCommitBarrier::install(
|
||||
crate::disk::RUSTFS_META_BUCKET,
|
||||
competing_object,
|
||||
PutObjectCommitPause::BeforeQuotaRename,
|
||||
);
|
||||
|
||||
let mut rebalance_task;
|
||||
let mut decommission_task;
|
||||
if paused_kind == PoolActivationStartKind::Rebalance {
|
||||
let store = Arc::clone(&rebalance_store);
|
||||
rebalance_task =
|
||||
tokio::spawn(async move { store.init_rebalance_start(vec!["bucket".to_string()]).await.map(|_| ()) });
|
||||
first_barrier.wait_until_paused().await;
|
||||
let probe = PoolActivationStartProbe::install(competing_kind);
|
||||
let store = Arc::clone(&decommission_store);
|
||||
decommission_task = tokio::spawn(async move { store.start_decommission(vec![0]).await });
|
||||
tokio::time::timeout(std::time::Duration::from_secs(15), probe.wait_until_attempted())
|
||||
.await
|
||||
.expect("real decommission start should reach activation lock acquisition");
|
||||
} else {
|
||||
let store = Arc::clone(&decommission_store);
|
||||
decommission_task = tokio::spawn(async move { store.start_decommission(vec![0]).await });
|
||||
first_barrier.wait_until_paused().await;
|
||||
let probe = PoolActivationStartProbe::install(competing_kind);
|
||||
let store = Arc::clone(&rebalance_store);
|
||||
rebalance_task =
|
||||
tokio::spawn(async move { store.init_rebalance_start(vec!["bucket".to_string()]).await.map(|_| ()) });
|
||||
tokio::time::timeout(std::time::Duration::from_secs(15), probe.wait_until_attempted())
|
||||
.await
|
||||
.expect("real rebalance start should reach activation lock acquisition");
|
||||
}
|
||||
|
||||
let mut observed_rebalance_result = None;
|
||||
let mut observed_decommission_result = None;
|
||||
let competing_reached_commit = tokio::time::timeout(std::time::Duration::from_secs(15), async {
|
||||
if paused_kind == PoolActivationStartKind::Rebalance {
|
||||
tokio::select! {
|
||||
_ = competing_barrier.wait_until_paused() => true,
|
||||
result = &mut decommission_task => {
|
||||
observed_decommission_result = Some(result.expect("decommission activation task should not panic"));
|
||||
false
|
||||
}
|
||||
}
|
||||
} else {
|
||||
tokio::select! {
|
||||
_ = competing_barrier.wait_until_paused() => true,
|
||||
result = &mut rebalance_task => {
|
||||
observed_rebalance_result = Some(result.expect("rebalance activation task should not panic"));
|
||||
false
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("competing activation should either block to timeout or reach its commit");
|
||||
if competing_reached_commit {
|
||||
competing_barrier.release();
|
||||
}
|
||||
drop(competing_barrier);
|
||||
first_barrier.release();
|
||||
drop(first_barrier);
|
||||
|
||||
let rebalance_result = match observed_rebalance_result {
|
||||
Some(result) => result,
|
||||
None => tokio::time::timeout(std::time::Duration::from_secs(15), &mut rebalance_task)
|
||||
.await
|
||||
.expect("real rebalance activation should finish")
|
||||
.expect("rebalance activation task should not panic"),
|
||||
};
|
||||
let decommission_result = match observed_decommission_result {
|
||||
Some(result) => result,
|
||||
None => tokio::time::timeout(std::time::Duration::from_secs(15), &mut decommission_task)
|
||||
.await
|
||||
.expect("real decommission activation should finish")
|
||||
.expect("decommission activation task should not panic"),
|
||||
};
|
||||
assert!(
|
||||
!competing_reached_commit,
|
||||
"the competing real start reached its commit while the other activation lock was held"
|
||||
);
|
||||
assert_ne!(
|
||||
rebalance_result.is_ok(),
|
||||
decommission_result.is_ok(),
|
||||
"exactly one real activation entry may commit"
|
||||
);
|
||||
|
||||
let mut persisted_rebalance = RebalanceMeta::new();
|
||||
let rebalance_committed = persisted_rebalance
|
||||
.load(rebalance_store.pools[0].clone())
|
||||
.await
|
||||
.is_ok_and(|_| is_rebalance_conflicting_with_decommission(&persisted_rebalance));
|
||||
let mut persisted_pool = PoolMeta::default();
|
||||
persisted_pool
|
||||
.load_no_lock(rebalance_store.pools[0].clone())
|
||||
.await
|
||||
.expect("persisted pool metadata should remain readable");
|
||||
let decommission_committed = pool_meta_has_active_decommission(&persisted_pool);
|
||||
assert_eq!(
|
||||
usize::from(rebalance_committed) + usize::from(decommission_committed),
|
||||
1,
|
||||
"at most one persisted activation state may be active"
|
||||
);
|
||||
assert_eq!(rebalance_committed, paused_kind == PoolActivationStartKind::Rebalance);
|
||||
assert_eq!(decommission_committed, paused_kind == PoolActivationStartKind::Decommission);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn real_rebalance_and_decommission_starts_commit_at_most_one_side() {
|
||||
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async {
|
||||
assert_real_activation_start_race(PoolActivationStartKind::Rebalance).await;
|
||||
assert_real_activation_start_race(PoolActivationStartKind::Decommission).await;
|
||||
})
|
||||
.await;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rebalance_activation_rejects_persisted_decommission_despite_idle_local_snapshot() {
|
||||
|
||||
@@ -50,6 +50,66 @@ fn ensure_rebalance_entry_active(cancel: &CancellationToken) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
static REBALANCE_RUN_SIGNAL_TEST_FENCES: std::sync::OnceLock<
|
||||
std::sync::Mutex<std::collections::HashMap<String, Arc<std::sync::atomic::AtomicBool>>>,
|
||||
> = std::sync::OnceLock::new();
|
||||
|
||||
#[cfg(test)]
|
||||
struct RebalanceRunSignalTestFence {
|
||||
rebalance_id: String,
|
||||
loss_handle: Arc<std::sync::atomic::AtomicBool>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl RebalanceRunSignalTestFence {
|
||||
fn install(rebalance_id: &str) -> Self {
|
||||
let loss_handle = Arc::new(std::sync::atomic::AtomicBool::new(false));
|
||||
let previous = REBALANCE_RUN_SIGNAL_TEST_FENCES
|
||||
.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
|
||||
.lock()
|
||||
.expect("rebalance run signal test fence should not be poisoned")
|
||||
.insert(rebalance_id.to_string(), Arc::clone(&loss_handle));
|
||||
assert!(previous.is_none(), "rebalance run signal test fence must be unique");
|
||||
Self {
|
||||
rebalance_id: rebalance_id.to_string(),
|
||||
loss_handle,
|
||||
}
|
||||
}
|
||||
|
||||
fn mark_lost(&self) {
|
||||
self.loss_handle.store(true, std::sync::atomic::Ordering::Release);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Drop for RebalanceRunSignalTestFence {
|
||||
fn drop(&mut self) {
|
||||
REBALANCE_RUN_SIGNAL_TEST_FENCES
|
||||
.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
|
||||
.lock()
|
||||
.expect("rebalance run signal test fence should not be poisoned")
|
||||
.remove(&self.rebalance_id);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn attach_rebalance_run_signal_test_fence(
|
||||
rebalance_id: &str,
|
||||
signal: &Arc<rustfs_lock::distributed_lock::LockLostSignal>,
|
||||
) -> Option<crate::object_api::NamespaceLockSignalTestFence> {
|
||||
let loss_handle = REBALANCE_RUN_SIGNAL_TEST_FENCES
|
||||
.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
|
||||
.lock()
|
||||
.expect("rebalance run signal test fence should not be poisoned")
|
||||
.get(rebalance_id)
|
||||
.cloned()?;
|
||||
Some(crate::object_api::NamespaceLockSignalTestFence::install_with_loss_handle(
|
||||
signal,
|
||||
loss_handle,
|
||||
))
|
||||
}
|
||||
|
||||
impl ECStore {
|
||||
async fn finish_rebalance_entry_after_cleanup(
|
||||
&self,
|
||||
@@ -179,6 +239,10 @@ impl ECStore {
|
||||
ensure_rebalance_entry_active(&cancel)?;
|
||||
let run_guard = self.rebalance_run_guard(rebalance_id.as_ref(), "rebalance entry").await?;
|
||||
let lock_lost_signal = run_guard.lock_lost_signal();
|
||||
#[cfg(test)]
|
||||
let _run_signal_test_fence = lock_lost_signal
|
||||
.as_ref()
|
||||
.and_then(|signal| attach_rebalance_run_signal_test_fence(rebalance_id.as_ref(), signal));
|
||||
|
||||
let mut rebalanced: usize = 0;
|
||||
let mut expired: usize = 0;
|
||||
@@ -725,12 +789,16 @@ impl ECStore {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
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::storage_api_contracts::object::ObjectOperations as _;
|
||||
use http::HeaderMap;
|
||||
use rustfs_filemeta::{FileInfo, FileMeta};
|
||||
use std::time::Duration as StdDuration;
|
||||
use time::OffsetDateTime;
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
#[tokio::test]
|
||||
async fn rebalance_stats_wait_for_source_cleanup_result() {
|
||||
@@ -983,4 +1051,138 @@ mod tests {
|
||||
"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() {
|
||||
let rebalance_id = "rebalance-target-commit-fence";
|
||||
let (_temp_dirs, store, _unused_store) = crate::services::rebalance::test_two_pool_stores(Some(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()
|
||||
}))
|
||||
.await;
|
||||
let bucket = crate::disk::RUSTFS_META_BUCKET;
|
||||
let object = "rebalance-commit-fence-object";
|
||||
let version_id = uuid::Uuid::new_v4();
|
||||
let payload = b"rebalance target commit must not survive its run fence".repeat(1024);
|
||||
let source_set = store.pools[0].get_disks_by_key(object);
|
||||
let target_set = store.pools[1].get_disks_by_key(object);
|
||||
let mut writer = PutObjReader::from_vec(payload.clone());
|
||||
let source_before = source_set
|
||||
.put_object(
|
||||
bucket,
|
||||
object,
|
||||
&mut writer,
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(version_id.to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("source version should be written");
|
||||
let source_versions = source_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.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 run_signal_fence = RebalanceRunSignalTestFence::install(rebalance_id);
|
||||
let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::BeforeQuotaRename);
|
||||
let entry_store = Arc::clone(&store);
|
||||
let entry_set = Arc::clone(&source_set);
|
||||
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),
|
||||
CancellationToken::new(),
|
||||
)
|
||||
.await
|
||||
});
|
||||
barrier.wait_until_paused().await;
|
||||
run_signal_fence.mark_lost();
|
||||
barrier.release();
|
||||
drop(barrier);
|
||||
let entry_error = tokio::time::timeout(StdDuration::from_secs(5), &mut entry_task)
|
||||
.await
|
||||
.expect("fenced real rebalance entry should finish")
|
||||
.expect("rebalance entry task should not panic")
|
||||
.expect_err("lost run fence must reject the target commit");
|
||||
assert!(entry_error.to_string().contains("run fence lost"));
|
||||
let target_error = target_set
|
||||
.get_object_info(
|
||||
bucket,
|
||||
object,
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(version_id.to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect_err("target version must remain absent after the lost commit fence");
|
||||
assert!(
|
||||
crate::error::is_err_object_not_found(&target_error) || crate::error::is_err_version_not_found(&target_error),
|
||||
"target version must remain absent after the lost commit fence"
|
||||
);
|
||||
let mut source_body = Vec::new();
|
||||
let mut source_after = source_set
|
||||
.get_object_reader(
|
||||
bucket,
|
||||
object,
|
||||
None,
|
||||
HeaderMap::new(),
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(version_id.to_string()),
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("source version must remain readable");
|
||||
assert_eq!(source_after.object_info.version_id, source_before.version_id);
|
||||
assert_eq!(source_after.object_info.data_dir, source_before.data_dir);
|
||||
assert_eq!(source_after.object_info.mod_time, source_before.mod_time);
|
||||
assert_eq!(source_after.object_info.size, source_before.size);
|
||||
assert_eq!(source_after.object_info.etag, source_before.etag);
|
||||
source_after
|
||||
.stream
|
||||
.read_to_end(&mut source_body)
|
||||
.await
|
||||
.expect("source body should drain");
|
||||
assert_eq!(source_body, payload, "source version must remain byte-identical");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -79,5 +79,65 @@ pub(crate) async fn test_store_with_persisted_rebalance_meta(
|
||||
(temp_dirs, store)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn test_two_pool_stores(
|
||||
rebalance_meta: Option<RebalanceMeta>,
|
||||
) -> (
|
||||
Vec<tempfile::TempDir>,
|
||||
std::sync::Arc<crate::store::ECStore>,
|
||||
std::sync::Arc<crate::store::ECStore>,
|
||||
) {
|
||||
use crate::core::pools::PoolMeta;
|
||||
use crate::layout::endpoints::{EndpointServerPools, SetupType};
|
||||
|
||||
let ctx = std::sync::Arc::new(crate::runtime::instance::InstanceContext::new());
|
||||
ctx.update_erasure_type(SetupType::DistErasure).await;
|
||||
let (mut temp_dirs, first_pool) =
|
||||
crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 0).await;
|
||||
let (second_temp_dirs, second_pool) =
|
||||
crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 1).await;
|
||||
temp_dirs.extend(second_temp_dirs);
|
||||
let pools = vec![first_pool, second_pool];
|
||||
{
|
||||
let local_disk_map = ctx.local_disk_map();
|
||||
let mut local_disk_map = local_disk_map.write().await;
|
||||
for pool in &pools {
|
||||
for set in &pool.disk_set {
|
||||
for disk in set.disks.read().await.iter().flatten() {
|
||||
local_disk_map.insert(disk.endpoint().to_string(), Some(disk.clone()));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
let pool_meta = PoolMeta::new(&pools, &PoolMeta::default());
|
||||
pool_meta
|
||||
.save(pools.clone())
|
||||
.await
|
||||
.expect("baseline pool metadata should be persisted");
|
||||
if let Some(meta) = rebalance_meta.as_ref() {
|
||||
meta.save(pools[0].clone())
|
||||
.await
|
||||
.expect("active rebalance metadata should be persisted");
|
||||
}
|
||||
let endpoint_pools: EndpointServerPools = pools.iter().map(|pool| pool.endpoints.clone()).collect::<Vec<_>>().into();
|
||||
ctx.set_endpoints(endpoint_pools.clone());
|
||||
let make_store = || {
|
||||
std::sync::Arc::new(crate::store::ECStore {
|
||||
id: uuid::Uuid::new_v4(),
|
||||
disk_map: std::collections::HashMap::new(),
|
||||
pools: pools.clone(),
|
||||
peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, std::sync::Arc::clone(&ctx)),
|
||||
pool_meta: tokio::sync::RwLock::new(pool_meta.clone()),
|
||||
rebalance_meta: tokio::sync::RwLock::new(rebalance_meta.clone()),
|
||||
decommission_cancelers: tokio::sync::RwLock::new(vec![None, None]),
|
||||
start_gate: tokio::sync::Mutex::new(()),
|
||||
pool_meta_save_gate: tokio::sync::Mutex::new(()),
|
||||
ctx: std::sync::Arc::clone(&ctx),
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
})
|
||||
};
|
||||
(temp_dirs, make_store(), make_store())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod rebalance_unit_tests;
|
||||
|
||||
Reference in New Issue
Block a user