fix(ecstore): bind rebalance workers to activation id

This commit is contained in:
overtrue
2026-08-22 00:01:11 +08:00
parent a72deafc9f
commit 814e17b02a
7 changed files with 481 additions and 86 deletions
+35 -1
View File
@@ -879,6 +879,11 @@ impl PoolRebalanceActivationFence {
pub(crate) fn ensure_held(&self) -> Result<()> {
ensure_activation_locks_held(self.pool_meta_guard.is_lock_lost(), self.rebalance_meta_guard.is_lock_lost())
}
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
opts.add_namespace_lock_guard(&self.pool_meta_guard);
opts.add_namespace_lock_guard(&self.rebalance_meta_guard);
}
}
fn ensure_activation_locks_held(pool_lock_lost: bool, rebalance_lock_lost: bool) -> Result<()> {
@@ -1793,6 +1798,33 @@ impl PoolMeta {
Ok(())
}
async fn save_no_lock_with_activation_fence<S>(
&self,
pools: Vec<Arc<S>>,
activation_fence: &PoolRebalanceActivationFence,
) -> Result<()>
where
S: EcstoreObjectIO,
{
let data = self.encode_config_data()?;
if data.is_empty() {
return Ok(());
}
for pool in pools {
let mut opts = ObjectOptions {
max_parity: true,
no_lock: true,
..Default::default()
};
activation_fence.add_namespace_lock_fence(&mut opts);
activation_fence.ensure_held()?;
save_config_with_opts(pool, POOL_META_NAME, data.clone(), &opts).await?;
activation_fence.ensure_held()?;
}
Ok(())
}
pub fn decommission_cancel(&mut self, idx: usize) -> bool {
if let Some(stats) = self.pools.get_mut(idx) {
if let Some(d) = &stats.decommission {
@@ -2591,7 +2623,9 @@ impl ECStore {
}
activation_fence.ensure_held()?;
latest_pool_meta.save_no_lock(self.pools.clone()).await?;
latest_pool_meta
.save_no_lock_with_activation_fence(self.pools.clone(), &activation_fence)
.await?;
{
let mut pool_meta = self.pool_meta.write().await;
*pool_meta = latest_pool_meta;
@@ -1123,7 +1123,15 @@ impl NotificationSys {
match store.stop_rebalance_for_id(expected_rebalance_id).await {
Ok(_) => {
if let Err(err) = store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await {
let save_result = match expected_rebalance_id {
Some(expected_id) => {
store
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_id)
.await
}
None => store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await,
};
if let Err(err) = save_result {
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
+218 -24
View File
@@ -45,35 +45,73 @@ pub(super) enum RebalanceWorkerActivationFence {
NotStartedTerminal,
}
async fn merge_and_save_rebalance_meta_no_lock<S, F>(
pub(super) struct RebalanceRunGuard {
_guard: tokio::sync::OwnedRwLockReadGuard<()>,
}
async fn merge_and_save_rebalance_meta_no_lock<S>(
pool: Arc<S>,
local_snapshot: &RebalanceMeta,
stage: &str,
before_save: F,
mut opts: ObjectOptions,
activation_fence: Option<&PoolRebalanceActivationFence>,
expected_id: Option<&str>,
) -> Result<()>
where
S: EcstoreObjectIO,
F: FnOnce() -> Result<()>,
{
let opts = ObjectOptions {
no_lock: true,
..Default::default()
};
let mut merged = RebalanceMeta::new();
match merged.load_with_opts(pool.clone(), opts.clone()).await {
Ok(()) => {
if let Some(expected_id) = expected_id {
ensure_rebalance_run_id(Some(&merged), expected_id, stage)?;
}
if merge_rebalance_meta(&mut merged, local_snapshot) == RebalanceMetaMergeOutcome::RejectedActiveConflict {
return Err(Error::RebalanceAlreadyRunning);
}
}
Err(Error::ConfigNotFound) => {
if expected_id.is_some() {
return Err(rebalance_metadata_not_initialized_error(stage));
}
merged = local_snapshot.clone();
}
Err(err) => return Err(Error::other(format!("rebalance meta load before save failed during {stage}: {err}"))),
}
before_save()?;
merged.save_with_opts(pool, opts).await
if let Some(fence) = activation_fence {
fence.add_namespace_lock_fence(&mut opts);
fence.ensure_held()?;
}
merged.save_with_opts(pool, opts).await?;
if let Some(fence) = activation_fence {
fence.ensure_held()?;
}
Ok(())
}
pub(super) fn ensure_rebalance_run_id(meta: Option<&RebalanceMeta>, expected_id: &str, stage: &str) -> Result<()> {
let Some(meta) = meta else {
return Err(rebalance_metadata_not_initialized_error(stage));
};
if meta.id != expected_id {
return Err(Error::other(format!(
"stale rebalance worker rejected during {stage}: expected {expected_id}, found {}",
meta.id
)));
}
Ok(())
}
pub(super) fn ensure_rebalance_worker_active(meta: Option<&RebalanceMeta>, expected_id: &str, stage: &str) -> Result<()> {
ensure_rebalance_run_id(meta, expected_id, stage)?;
let Some(meta) = meta else {
return Err(rebalance_metadata_not_initialized_error(stage));
};
if meta.stopped_at.is_some() || !is_rebalance_conflicting_with_decommission(meta) {
return Err(Error::other(format!("inactive rebalance worker rejected during {stage}: {expected_id}")));
}
Ok(())
}
pub(super) fn validate_rebalance_disk_stats_coverage(disk_stats: &[DiskStat]) -> Result<()> {
@@ -122,22 +160,105 @@ fn clear_rebalance_status_refresh(current: &mut Option<RebalanceMeta>) {
}
impl ECStore {
// Transition order is start_gate -> activation_gate -> rebalance_meta; the probe read below is released before the gate.
async fn rebalance_activation_write_guard(
&self,
expected_id: Option<&str>,
stage: &str,
) -> Result<Option<tokio::sync::OwnedRwLockWriteGuard<()>>> {
let activation_gate = {
let meta = self.rebalance_meta.read().await;
if let Some(expected_id) = expected_id {
ensure_rebalance_run_id(meta.as_ref(), expected_id, stage)?;
}
meta.as_ref().map(|meta| Arc::clone(&meta.activation_gate))
};
Ok(match activation_gate {
Some(gate) => Some(gate.write_owned().await),
None => None,
})
}
pub(super) async fn rebalance_run_guard(&self, expected_id: &str, stage: &str) -> Result<RebalanceRunGuard> {
let activation_gate = {
let meta = self.rebalance_meta.read().await;
ensure_rebalance_worker_active(meta.as_ref(), expected_id, stage)?;
Arc::clone(
&meta
.as_ref()
.ok_or_else(|| rebalance_metadata_not_initialized_error(stage))?
.activation_gate,
)
};
let guard = Arc::clone(&activation_gate).read_owned().await;
let meta = self.rebalance_meta.read().await;
ensure_rebalance_worker_active(meta.as_ref(), expected_id, stage)?;
let current_gate = &meta
.as_ref()
.ok_or_else(|| rebalance_metadata_not_initialized_error(stage))?
.activation_gate;
if !Arc::ptr_eq(&activation_gate, current_gate) {
return Err(Error::other(format!(
"stale rebalance activation gate rejected during {stage}: {expected_id}"
)));
}
drop(meta);
Ok(RebalanceRunGuard { _guard: guard })
}
pub(super) async fn save_rebalance_meta_with_merge<S>(
&self,
pool: Arc<S>,
local_snapshot: &RebalanceMeta,
stage: &str,
) -> Result<()>
where
S: EcstoreObjectIO + StorageNamespaceLocking<Error = Error, NamespaceLock = rustfs_lock::NamespaceLockWrapper>,
{
self.save_rebalance_meta_with_merge_for_id(pool, local_snapshot, stage, None)
.await
}
pub(super) async fn save_rebalance_meta_for_id_with_merge<S>(
&self,
pool: Arc<S>,
local_snapshot: &RebalanceMeta,
stage: &str,
expected_id: &str,
) -> Result<()>
where
S: EcstoreObjectIO + StorageNamespaceLocking<Error = Error, NamespaceLock = rustfs_lock::NamespaceLockWrapper>,
{
self.save_rebalance_meta_with_merge_for_id(pool, local_snapshot, stage, Some(expected_id))
.await
}
async fn save_rebalance_meta_with_merge_for_id<S>(
&self,
pool: Arc<S>,
local_snapshot: &RebalanceMeta,
stage: &str,
expected_id: Option<&str>,
) -> Result<()>
where
S: EcstoreObjectIO + StorageNamespaceLocking<Error = Error, NamespaceLock = rustfs_lock::NamespaceLockWrapper>,
{
let ns_lock = pool.new_ns_lock(crate::disk::RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
let _guard = ns_lock
let guard = ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(rebalance_meta_lock_error)?;
let mut opts = ObjectOptions {
no_lock: true,
..Default::default()
};
opts.add_namespace_lock_guard(&guard);
merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage, || Ok(())).await
merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage, opts, None, expected_id).await?;
if guard.is_lock_lost() {
return Err(Error::other("rebalance metadata lock lost during metadata commit"));
}
Ok(())
}
async fn save_rebalance_activation_meta_with_merge<S>(
@@ -154,7 +275,18 @@ impl ECStore {
pool_meta.load_no_lock(pool.clone()).await?;
ensure_rebalance_activation_pool_meta_allowed(&pool_meta)?;
merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage, || activation_fence.ensure_held()).await
merge_and_save_rebalance_meta_no_lock(
pool,
local_snapshot,
stage,
ObjectOptions {
no_lock: true,
..Default::default()
},
Some(&activation_fence),
None,
)
.await
}
pub(super) async fn fence_rebalance_worker_activation<S>(
@@ -197,6 +329,7 @@ impl ECStore {
#[tracing::instrument(skip_all)]
pub async fn load_rebalance_meta(&self) -> Result<()> {
let _start_guard = self.start_gate.lock().await;
let mut meta = RebalanceMeta::new();
debug!(
event = EVENT_REBALANCE_STATE,
@@ -206,7 +339,9 @@ impl ECStore {
"Loading rebalance metadata"
);
let pool = clone_first_arc(&self.pools, "rebalanceMeta: no pools available")?;
if resolve_rebalance_meta_load_result(meta.load(pool).await)? {
let loaded = resolve_rebalance_meta_load_result(meta.load(pool).await)?;
let _activation_guard = self.rebalance_activation_write_guard(None, "load rebalance metadata").await?;
if loaded {
{
let mut rebalance_meta = self.rebalance_meta.write().await;
@@ -243,9 +378,14 @@ impl ECStore {
#[tracing::instrument(skip_all)]
pub async fn refresh_rebalance_status_meta(&self) -> Result<()> {
let _start_guard = self.start_gate.lock().await;
let pool = clone_first_arc(&self.pools, "refresh_rebalance_status_meta: no pools available")?;
let mut persisted = RebalanceMeta::new();
match persisted.load(pool).await {
let loaded = persisted.load(pool).await;
let _activation_guard = self
.rebalance_activation_write_guard(None, "refresh rebalance metadata")
.await?;
match loaded {
Ok(()) => {
let mut rebalance_meta = self.rebalance_meta.write().await;
merge_rebalance_status_refresh(&mut rebalance_meta, persisted);
@@ -311,7 +451,7 @@ impl ECStore {
if let Some(meta) = rebalance_meta.as_ref() {
let pool = clone_first_arc(&self.pools, "update_rebalance_stats: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, meta, "update_rebalance_stats")
self.save_rebalance_meta_for_id_with_merge(pool, meta, "update_rebalance_stats", meta.id.as_str())
.await,
"update_rebalance_stats",
)?;
@@ -436,11 +576,14 @@ impl ECStore {
pub async fn init_rebalance_start(self: &Arc<Self>, bucktes: Vec<String>) -> Result<String> {
let _start_guard = self.start_gate.lock().await;
let decommission_running = self.is_decommission_running().await;
{
let decommission_running = self.is_decommission_running().await;
let rebalance_meta = self.rebalance_meta.read().await;
validate_init_rebalance_state(decommission_running, rebalance_meta.as_ref())?;
}
let _activation_guard = self
.rebalance_activation_write_guard(None, "initialize replacement rebalance")
.await?;
self.init_rebalance_meta(bucktes).await
}
@@ -480,11 +623,35 @@ impl ECStore {
#[tracing::instrument(skip(self, versions))]
pub async fn update_pool_stats_batch(&self, pool_index: usize, bucket: String, versions: &[&FileInfo]) -> Result<()> {
self.update_pool_stats_batch_for_id(pool_index, bucket, versions, None).await
}
pub(super) async fn update_pool_stats_batch_for_rebalance(
&self,
pool_index: usize,
bucket: String,
versions: &[&FileInfo],
expected_id: &str,
) -> Result<()> {
self.update_pool_stats_batch_for_id(pool_index, bucket, versions, Some(expected_id))
.await
}
async fn update_pool_stats_batch_for_id(
&self,
pool_index: usize,
bucket: String,
versions: &[&FileInfo],
expected_id: Option<&str>,
) -> Result<()> {
if versions.is_empty() {
return Ok(());
}
let mut rebalance_meta = self.rebalance_meta.write().await;
if let Some(expected_id) = expected_id {
ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "update rebalance pool stats")?;
}
if let Some(meta) = rebalance_meta.as_mut() {
if !should_accept_rebalance_stats_update(meta, pool_index) {
return Ok(());
@@ -499,8 +666,9 @@ impl ECStore {
}
#[tracing::instrument(skip(self))]
pub async fn next_rebal_bucket(&self, pool_index: usize) -> Result<Option<String>> {
pub async fn next_rebal_bucket(&self, pool_index: usize, expected_id: &str) -> Result<Option<String>> {
let rebalance_meta = self.rebalance_meta.read().await;
ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "next rebalance bucket")?;
debug!(
event = EVENT_REBALANCE_BUCKET,
component = LOG_COMPONENT_ECSTORE,
@@ -514,8 +682,9 @@ impl ECStore {
}
#[tracing::instrument(skip(self))]
pub async fn bucket_rebalance_done(&self, pool_index: usize, bucket: String) -> Result<()> {
pub async fn bucket_rebalance_done(&self, pool_index: usize, bucket: String, expected_id: &str) -> Result<()> {
let mut rebalance_meta = self.rebalance_meta.write().await;
ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "mark rebalance bucket done")?;
mark_rebalance_bucket_done(rebalance_meta.as_mut(), pool_index, &bucket)
}
@@ -525,8 +694,10 @@ impl ECStore {
bucket: &str,
object: &str,
message: String,
expected_id: &str,
) -> Result<()> {
let mut rebalance_meta = self.rebalance_meta.write().await;
ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "record rebalance cleanup warning")?;
record_rebalance_cleanup_warning_in_meta(
rebalance_meta.as_mut(),
pool_index,
@@ -537,8 +708,15 @@ impl ECStore {
)
}
pub(super) async fn defer_rebalance_bucket(&self, pool_index: usize, bucket: String, last_error: String) -> Result<()> {
pub(super) async fn defer_rebalance_bucket(
&self,
pool_index: usize,
bucket: String,
last_error: String,
expected_id: &str,
) -> Result<()> {
let mut rebalance_meta = self.rebalance_meta.write().await;
ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "defer rebalance bucket")?;
let Some(meta) = rebalance_meta.as_mut() else {
return Err(rebalance_metadata_not_initialized_error("defer rebalance bucket"));
};
@@ -634,6 +812,8 @@ impl ECStore {
#[tracing::instrument(skip(self))]
pub async fn stop_rebalance_for_id(self: &Arc<Self>, expected_id: Option<&str>) -> Result<()> {
let _start_guard = self.start_gate.lock().await;
let _activation_guard = self.rebalance_activation_write_guard(expected_id, "stop rebalance").await?;
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
stop_rebalance_meta_snapshot_for_id(rebalance_meta.as_mut(), OffsetDateTime::now_utc(), expected_id)
@@ -642,7 +822,7 @@ impl ECStore {
if let Some(meta_to_save) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "stop_rebalance: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, "stop_rebalance")
self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, "stop_rebalance", meta_to_save.id.as_str())
.await,
"stop_rebalance",
)?;
@@ -656,6 +836,10 @@ impl ECStore {
expected_id: Option<&str>,
start_error: String,
) -> Result<()> {
let _start_guard = self.start_gate.lock().await;
let _activation_guard = self
.rebalance_activation_write_guard(expected_id, "rollback rebalance start")
.await?;
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_rebalance_start_meta_snapshot_for_id(
@@ -669,8 +853,13 @@ impl ECStore {
if let Some(meta_to_save) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "rollback_rebalance_start: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, "rollback_rebalance_start")
.await,
self.save_rebalance_meta_for_id_with_merge(
pool,
&meta_to_save,
"rollback_rebalance_start",
meta_to_save.id.as_str(),
)
.await,
"rollback_rebalance_start",
)?;
}
@@ -692,8 +881,13 @@ impl ECStore {
if let Some(meta_to_save) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "record_rebalance_stop_propagation: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, "record_rebalance_stop_propagation")
.await,
self.save_rebalance_meta_for_id_with_merge(
pool,
&meta_to_save,
"record_rebalance_stop_propagation",
meta_to_save.id.as_str(),
)
.await,
"record_rebalance_stop_propagation",
)?;
}
+63 -19
View File
@@ -50,16 +50,20 @@ impl ECStore {
bucket: &str,
object: &str,
stats_updates: &[&FileInfo],
expected_id: &str,
cleanup: impl std::future::Future<Output = std::result::Result<ObjectInfo, data_movement::SourceCleanupError>>,
) -> Result<RebalanceEntryCleanupResult> {
// Persisted stats can complete a pool on restart, so source cleanup must resolve first.
let cleanup_result = resolve_rebalance_entry_cleanup_delete_result(cleanup.await, bucket, object);
let run_guard = self.rebalance_run_guard(expected_id, "rebalance source cleanup").await?;
let cleanup_result = cleanup.await;
drop(run_guard);
let cleanup_result = resolve_rebalance_entry_cleanup_delete_result(cleanup_result, bucket, object);
let RebalanceEntryCleanupResult::Completed { warning } = cleanup_result else {
return Ok(cleanup_result);
};
if let Some(message) = warning.as_ref()
&& let Err(err) = self
.record_rebalance_cleanup_warning(pool_index, bucket, object, message.clone())
.record_rebalance_cleanup_warning(pool_index, bucket, object, message.clone(), expected_id)
.await
{
error!(
@@ -76,7 +80,7 @@ impl ECStore {
}
resolve_rebalance_stats_update_result(
self.update_pool_stats_batch(pool_index, bucket.to_string(), stats_updates)
self.update_pool_stats_batch_for_rebalance(pool_index, bucket.to_string(), stats_updates, expected_id)
.await,
pool_index,
bucket,
@@ -95,6 +99,7 @@ impl ECStore {
entry: MetaCacheEntry,
set: Arc<SetDisks>,
bucket_configs: Arc<RebalanceBucketConfigs>,
rebalance_id: Arc<str>,
// wk: Arc<Workers>,
) -> Result<RebalanceEntryOutcome> {
debug!(
@@ -129,7 +134,7 @@ impl ECStore {
return Ok(RebalanceEntryOutcome::Completed);
}
if self.check_if_rebalance_done(pool_index).await {
if self.check_if_rebalance_done(pool_index, rebalance_id.as_ref()).await? {
debug!(
event = EVENT_REBALANCE_ENTRY,
component = LOG_COMPONENT_ECSTORE,
@@ -160,7 +165,10 @@ impl ECStore {
let mut cleanup_preflight_allowed_missing = Vec::new();
let mut stats_updates = Vec::with_capacity(fivs.versions.len());
for version in fivs.versions.iter() {
if crate::core::pools::should_skip_lifecycle_for_data_movement(
let run_guard = self
.rebalance_run_guard(rebalance_id.as_ref(), "rebalance lifecycle mutation")
.await?;
let expired_by_lifecycle = crate::core::pools::should_skip_lifecycle_for_data_movement(
self.clone(),
&bucket,
version,
@@ -169,8 +177,9 @@ impl ECStore {
true,
&crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc::Rebal,
)
.await?
{
.await?;
drop(run_guard);
if expired_by_lifecycle {
expired += 1;
// The lifecycle expiry above physically deleted this version from the source set.
// Record its identity so the source-cleanup preflight tolerates its absence,
@@ -211,17 +220,32 @@ impl ECStore {
let expected_bucket_incarnation_id = bucket_configs.bucket_incarnation_id;
let mut transfer = |src_pool_idx: usize, bucket: String, rd: GetObjectReader| {
let store = self.clone();
let rebalance_id = Arc::clone(&rebalance_id);
async move {
store
let run_guard = store
.rebalance_run_guard(rebalance_id.as_ref(), "rebalance object migration")
.await?;
let result = store
.clone()
.rebalance_object(src_pool_idx, bucket, rd, expected_bucket_incarnation_id)
.await
.await;
drop(run_guard);
result
}
};
// Route delete-marker migration through the store layer so it lands on the
// cross-pool target (excluding the source pool), not back onto the source set.
let mut delete_marker = |bucket: String, object: String, opts: ObjectOptions| {
let store = self.clone();
async move { store.delete_object(&bucket, &object, opts).await }
let rebalance_id = Arc::clone(&rebalance_id);
async move {
let run_guard = store
.rebalance_run_guard(rebalance_id.as_ref(), "rebalance delete-marker migration")
.await?;
let result = store.delete_object(&bucket, &object, opts).await;
drop(run_guard);
result
}
};
let result = migrate_entry_version(
&RebalanceMigrationBackend::new(set.as_ref(), self.as_ref()),
@@ -280,7 +304,10 @@ impl ECStore {
error = %err,
"Deferred rebalance entry after transient migration failure"
);
if let Err(stats_err) = self.update_rebalance_last_error(pool_index, deferred_error.clone()).await {
if let Err(stats_err) = self
.update_rebalance_last_error(pool_index, deferred_error.clone(), rebalance_id.as_ref())
.await
{
error!(
"rebalance_entry {} failed to record deferred transient failure for {}: {}",
&bucket, &entry.name, stats_err
@@ -295,7 +322,12 @@ impl ECStore {
if !stats_updates.is_empty()
&& let Err(stats_err) = self
.update_pool_stats_batch(pool_index, bucket.clone(), stats_updates.as_slice())
.update_pool_stats_batch_for_rebalance(
pool_index,
bucket.clone(),
stats_updates.as_slice(),
rebalance_id.as_ref(),
)
.await
{
error!(
@@ -323,6 +355,7 @@ impl ECStore {
bucket.as_str(),
entry.name.as_str(),
stats_updates.as_slice(),
rebalance_id.as_ref(),
data_movement::cleanup_source_entry_if_unchanged(
set.clone(),
bucket.as_str(),
@@ -397,8 +430,13 @@ impl ECStore {
);
resolve_rebalance_stats_update_result(
self.update_pool_stats_batch(pool_index, bucket.clone(), stats_updates.as_slice())
.await,
self.update_pool_stats_batch_for_rebalance(
pool_index,
bucket.clone(),
stats_updates.as_slice(),
rebalance_id.as_ref(),
)
.await,
pool_index,
bucket.as_str(),
entry.name.as_str(),
@@ -419,8 +457,9 @@ impl ECStore {
data_movement::migrate_object(self, pool_idx, bucket, rd, expected_bucket_incarnation_id, "rebalance_object").await
}
async fn update_rebalance_last_error(&self, pool_idx: usize, message: String) -> Result<()> {
async fn update_rebalance_last_error(&self, pool_idx: usize, message: String, expected_id: &str) -> Result<()> {
let mut rebalance_meta = self.rebalance_meta.write().await;
super::control::ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "record rebalance last error")?;
let Some(meta) = rebalance_meta.as_mut() else {
return Err(rebalance_metadata_not_initialized_error("record rebalance last error"));
};
@@ -441,6 +480,7 @@ impl ECStore {
rx: CancellationToken,
bucket: String,
pool_index: usize,
rebalance_id: Arc<str>,
) -> Result<RebalanceBucketOutcome> {
ensure_valid_rebalance_pool_index(self.pools.len(), pool_index)?;
@@ -473,6 +513,7 @@ impl ECStore {
let bucket_configs = bucket_configs.clone();
let entry_tasks = entry_tasks.clone();
let entry_workers = entry_workers.clone();
let rebalance_id = Arc::clone(&rebalance_id);
move |entry: MetaCacheEntry| {
let this = this.clone();
let bucket = bucket.clone();
@@ -482,6 +523,7 @@ impl ECStore {
let bucket_configs = bucket_configs.clone();
let entry_tasks = entry_tasks.clone();
let entry_workers = entry_workers.clone();
let rebalance_id = Arc::clone(&rebalance_id);
Box::pin(async move {
if callback_rx.is_cancelled() {
return;
@@ -534,7 +576,9 @@ impl ECStore {
state = "task_started",
"Started rebalance entry task"
);
let result = this.rebalance_entry(bucket, pool_index, entry, set, bucket_configs).await;
let result = this
.rebalance_entry(bucket, pool_index, entry, set, bucket_configs, rebalance_id)
.await;
if let Err(err) = &result {
error!("rebalance_entry: rebalance entry failed: {err}");
let mut first_err = entry_error.lock().await;
@@ -679,7 +723,7 @@ mod tests {
let finish_store = Arc::clone(&store);
let finish = tokio::spawn(async move {
finish_store
.finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&version], async move {
.finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&version], "", async move {
cleanup_released.await.expect("cleanup release sender should remain alive");
Ok(ObjectInfo::default())
})
@@ -726,7 +770,7 @@ mod tests {
meta.as_mut().expect("rebalance metadata should exist").pool_stats[0].bytes = 0;
}
let warning_result = store
.finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], async {
.finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], "", async {
Err(Error::SlowDown.into())
})
.await
@@ -743,7 +787,7 @@ mod tests {
meta.as_mut().expect("rebalance metadata should exist").pool_stats[0].bytes = 0;
}
let deferred = store
.finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], async {
.finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], "", async {
Err(data_movement::SourceCleanupError::SourceChanged)
})
.await
@@ -66,10 +66,12 @@ use rustfs_filemeta::{FileInfo, MetaCacheEntry};
use rustfs_rio::Index;
use s3s::dto::ReplicationConfiguration;
use serde::Serialize;
use std::future::Future;
use std::io::Cursor;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::task::{Context, Poll};
use time::OffsetDateTime;
use tokio::sync::mpsc;
use tokio::time::Duration;
@@ -2711,9 +2713,83 @@ async fn test_start_rebalance_for_id_rejects_stopped_metadata() {
assert!(err.to_string().contains("was stopped before start"));
}
#[test]
fn test_stopped_activation_state_prevents_worker_token_commit() {
let mut meta = RebalanceMeta {
id: "rebalance-a".to_string(),
stopped_at: Some(OffsetDateTime::now_utc()),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let outcome = commit_local_rebalance_worker_activation(&mut meta, "rebalance-a", tokio_util::sync::CancellationToken::new())
.expect("stopped metadata should produce a non-start outcome");
assert_eq!(outcome, RebalanceLocalActivationOutcome::NotStartedTerminal);
assert!(meta.cancel.is_none(), "stopped rebalance must not receive a worker token");
}
#[tokio::test]
async fn test_stop_at_activation_barrier_prevents_worker_token_commit() {
let meta = Arc::new(tokio::sync::RwLock::new(RebalanceMeta {
async fn test_old_worker_cannot_mutate_replacement_rebalance_state() {
let meta = RebalanceMeta {
id: "rebalance-b".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
buckets: vec!["bucket-a".to_string()],
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let store = test_store_with_rebalance_meta(meta);
let mut fi = FileInfo::default();
fi.size = 128;
for err in [
store
.next_rebal_bucket(0, "rebalance-a")
.await
.expect_err("old worker must not read replacement work"),
store
.bucket_rebalance_done(0, "bucket-a".to_string(), "rebalance-a")
.await
.expect_err("old worker must not complete replacement bucket"),
store
.update_pool_stats_batch_for_rebalance(0, "bucket-a".to_string(), &[&fi], "rebalance-a")
.await
.expect_err("old worker must not update replacement stats"),
store
.check_if_rebalance_done(0, "rebalance-a")
.await
.expect_err("old worker must not complete replacement pool"),
store
.save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, "rebalance-a")
.await
.expect_err("old save task must not persist replacement metadata"),
] {
assert!(err.to_string().contains("stale rebalance worker rejected"));
}
let meta = store.rebalance_meta.read().await;
let meta = meta.as_ref().expect("replacement metadata should remain present");
assert_eq!(meta.id, "rebalance-b");
assert!(meta.pool_stats[0].rebalanced_buckets.is_empty());
assert_eq!(meta.pool_stats[0].bytes, 0);
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
}
#[tokio::test]
async fn test_stop_waits_for_active_rebalance_side_effect_guard() {
let meta = RebalanceMeta {
id: "rebalance-a".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
@@ -2724,28 +2800,38 @@ async fn test_stop_at_activation_barrier_prevents_worker_token_commit() {
..Default::default()
}],
..Default::default()
}));
let fence_reached = Arc::new(tokio::sync::Barrier::new(2));
let stop_committed = Arc::new(tokio::sync::Barrier::new(2));
};
let store = test_store_with_rebalance_meta(meta);
let run_guard = store
.rebalance_run_guard("rebalance-a", "test side effect")
.await
.expect("active run should admit the side effect");
let mut stop = Box::pin(store.stop_rebalance_for_id(Some("rebalance-a")));
let mut context = Context::from_waker(futures::task::noop_waker_ref());
let stop_meta = Arc::clone(&meta);
let stop_fence_reached = Arc::clone(&fence_reached);
let stop_committed_signal = Arc::clone(&stop_committed);
let stop = tokio::spawn(async move {
stop_fence_reached.wait().await;
stop_meta.write().await.stopped_at = Some(OffsetDateTime::now_utc());
stop_committed_signal.wait().await;
});
assert!(matches!(stop.as_mut().poll(&mut context), Poll::Pending));
assert!(
store
.rebalance_meta
.read()
.await
.as_ref()
.is_some_and(|meta| meta.stopped_at.is_none()),
"stop must not change run state while a fenced side effect is active"
);
fence_reached.wait().await;
stop_committed.wait().await;
stop.await.expect("stop barrier task should finish");
let mut meta = meta.write().await;
let outcome = commit_local_rebalance_worker_activation(&mut meta, "rebalance-a", tokio_util::sync::CancellationToken::new())
.expect("stopped metadata should produce a non-start outcome");
assert_eq!(outcome, RebalanceLocalActivationOutcome::NotStartedTerminal);
assert!(meta.cancel.is_none(), "stopped rebalance must not receive a worker token");
drop(run_guard);
stop.await
.expect_err("empty test store should fail only after committing the local stop state");
assert!(
store
.rebalance_meta
.read()
.await
.as_ref()
.is_some_and(|meta| meta.stopped_at.is_some()),
"stop should commit after the production side-effect fence is released"
);
}
fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc<crate::store::ECStore> {
@@ -79,12 +79,12 @@ impl ECStore {
state = "starting",
"Starting rebalance"
);
let expected_id = {
let expected_id: Arc<str> = {
let rebalance_meta = self.rebalance_meta.read().await;
rebalance_meta.as_ref().ok_or(Error::ConfigNotFound)?.id.clone()
Arc::from(rebalance_meta.as_ref().ok_or(Error::ConfigNotFound)?.id.as_str())
};
let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?;
let activation_fence = match self.fence_rebalance_worker_activation(pool, &expected_id).await? {
let activation_fence = match self.fence_rebalance_worker_activation(pool, expected_id.as_ref()).await? {
RebalanceWorkerActivationFence::Ready(fence) => fence,
RebalanceWorkerActivationFence::NotStartedTerminal => return Ok(()),
};
@@ -122,7 +122,7 @@ impl ECStore {
meta_to_save = Some(meta.clone());
}
activation_fence.ensure_held()?;
activation_outcome = commit_local_rebalance_worker_activation(meta, &expected_id, cancel_tx)?;
activation_outcome = commit_local_rebalance_worker_activation(meta, expected_id.as_ref(), cancel_tx)?;
drop(rebalance_meta);
}
@@ -198,9 +198,10 @@ impl ECStore {
let pool_idx = idx;
let store = self.clone();
let rx_clone = rx.clone();
let worker_id = Arc::clone(&expected_id);
workers_started += 1;
tokio::spawn(async move {
if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx).await {
if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx, worker_id).await {
error!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -247,13 +248,14 @@ impl ECStore {
}
#[tracing::instrument(skip(self, rx))]
async fn rebalance_buckets(self: &Arc<Self>, rx: CancellationToken, pool_index: usize) -> Result<()> {
async fn rebalance_buckets(self: &Arc<Self>, rx: CancellationToken, pool_index: usize, rebalance_id: Arc<str>) -> Result<()> {
ensure_valid_rebalance_pool_index(self.pools.len(), pool_index)?;
let (done_tx, mut done_rx) = tokio::sync::mpsc::channel::<Result<()>>(1);
// Save rebalance metadata periodically
let store = self.clone();
let save_rebalance_id = Arc::clone(&rebalance_id);
let save_task = tokio::spawn(async move {
let mut timer = tokio::time::interval_at(Instant::now() + Duration::from_secs(30), Duration::from_secs(10));
let mut msg: String;
@@ -267,6 +269,11 @@ impl ECStore {
let terminal_event = classify_rebalance_terminal_event(result, now);
msg = terminal_event.message().to_string();
let mut rebalance_meta = store.rebalance_meta.write().await;
super::control::ensure_rebalance_run_id(
rebalance_meta.as_ref(),
save_rebalance_id.as_ref(),
"apply rebalance terminal event",
)?;
if let Some(meta) = rebalance_meta.as_mut() {
let meta_stopped = meta.stopped_at.is_some();
if let Some(pool_stat) = meta.pool_stats.get_mut(pool_index) {
@@ -315,7 +322,10 @@ impl ECStore {
}
}
if let Err(err) = store.save_rebalance_stats(pool_index, RebalSaveOpt::Stats).await {
if let Err(err) = store
.save_rebalance_stats_for_id(pool_index, RebalSaveOpt::Stats, save_rebalance_id.as_ref())
.await
{
let wrapped = Error::other(format!("rebalance save_task stats save failed for pool {pool_index}: {err}"));
error!("{} err: {:?}", msg, wrapped);
if quit {
@@ -381,7 +391,7 @@ impl ECStore {
break;
}
let next_bucket = match self.next_rebal_bucket(pool_index).await {
let next_bucket = match self.next_rebal_bucket(pool_index, rebalance_id.as_ref()).await {
Ok(bucket) => bucket,
Err(err) => {
error!(
@@ -413,7 +423,8 @@ impl ECStore {
);
let outcome = match resolve_rebalance_bucket_result(
self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index).await,
self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index, Arc::clone(&rebalance_id))
.await,
pool_index,
&bucket,
) {
@@ -476,7 +487,7 @@ impl ECStore {
"Deferred rebalance bucket after transient object failures"
);
if let Err(err) = self
.defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone())
.defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone(), rebalance_id.as_ref())
.await
{
error!(
@@ -540,7 +551,7 @@ impl ECStore {
"Completed rebalance bucket"
);
source_cleanup_deferred_attempts.remove(&bucket);
if let Err(err) = self.bucket_rebalance_done(pool_index, bucket).await {
if let Err(err) = self.bucket_rebalance_done(pool_index, bucket, rebalance_id.as_ref()).await {
error!(
event = EVENT_REBALANCE_BUCKET,
component = LOG_COMPONENT_ECSTORE,
@@ -601,8 +612,9 @@ impl ECStore {
final_result
}
pub(super) async fn check_if_rebalance_done(&self, pool_index: usize) -> bool {
pub(super) async fn check_if_rebalance_done(&self, pool_index: usize, expected_id: &str) -> Result<bool> {
let mut rebalance_meta = self.rebalance_meta.write().await;
super::control::ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "check rebalance completion")?;
if let Some(meta) = rebalance_meta.as_mut()
&& let Some(pool_stat) = meta.pool_stats.get_mut(pool_index)
@@ -617,7 +629,7 @@ impl ECStore {
state = "already_completed",
"Rebalance pool is already completed"
);
return true;
return Ok(true);
}
// Mark pool rebalance as done only after it reaches the PercentFreeGoal.
@@ -647,19 +659,30 @@ impl ECStore {
percent_free = pfi,
"Marked rebalance pool completed"
);
return true;
return Ok(true);
}
}
false
Ok(false)
}
}
impl ECStore {
#[tracing::instrument(skip(self))]
pub async fn save_rebalance_stats(&self, pool_idx: usize, opt: RebalSaveOpt) -> Result<()> {
self.save_rebalance_stats_inner(pool_idx, opt, None).await
}
pub(crate) async fn save_rebalance_stats_for_id(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: &str) -> Result<()> {
self.save_rebalance_stats_inner(pool_idx, opt, Some(expected_id)).await
}
async fn save_rebalance_stats_inner(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: Option<&str>) -> Result<()> {
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
if let Some(expected_id) = expected_id {
super::control::ensure_rebalance_run_id(rebalance_meta.as_ref(), expected_id, "save rebalance stats")?;
}
let Some(meta) = rebalance_meta.as_mut() else {
return Ok(());
};
@@ -681,10 +704,14 @@ impl ECStore {
"Rebalance metadata save requested"
);
let stage = format!("save_rebalance_stats for pool {pool_idx} opt {opt:?}");
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await,
stage.as_str(),
)?;
let save_result = match expected_id {
Some(expected_id) => {
self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, stage.as_str(), expected_id)
.await
}
None => self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await,
};
resolve_rebalance_meta_save_result(save_result, stage.as_str())?;
Ok(())
}
@@ -144,6 +144,8 @@ pub struct RebalanceMeta {
#[serde(skip)]
pub cancel: Option<CancellationToken>, // To be invoked on rebalance-stop
#[serde(skip)]
pub activation_gate: std::sync::Arc<tokio::sync::RwLock<()>>,
#[serde(skip)]
pub last_refreshed_at: Option<OffsetDateTime>,
#[serde(rename = "stopTs")]
pub stopped_at: Option<OffsetDateTime>, // Time when rebalance-stop was issued