mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 20:59:05 +00:00
fix(ecstore): close rebalance activation races
This commit is contained in:
@@ -2639,6 +2639,11 @@ impl ECStore {
|
|||||||
ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)
|
ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn ensure_decommission_rebalance_idle_after_refresh_under_start_gate(&self) -> Result<()> {
|
||||||
|
self.load_rebalance_meta_under_start_gate().await?;
|
||||||
|
ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn status(&self, idx: usize) -> Result<PoolStatus> {
|
pub async fn status(&self, idx: usize) -> Result<PoolStatus> {
|
||||||
let space_info = self.get_decommission_pool_space_info(idx).await?;
|
let space_info = self.get_decommission_pool_space_info(idx).await?;
|
||||||
|
|
||||||
@@ -4148,7 +4153,8 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let _start_guard = self.start_gate.lock().await;
|
let _start_guard = self.start_gate.lock().await;
|
||||||
self.ensure_decommission_rebalance_idle_after_refresh().await?;
|
self.ensure_decommission_rebalance_idle_after_refresh_under_start_gate()
|
||||||
|
.await?;
|
||||||
|
|
||||||
let all_space_infos = self.get_decommission_all_pool_space_infos().await?;
|
let all_space_infos = self.get_decommission_all_pool_space_infos().await?;
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -1121,15 +1121,19 @@ impl NotificationSys {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
match store.stop_rebalance_for_id(expected_rebalance_id).await {
|
let local_rebalance_id = match expected_rebalance_id {
|
||||||
|
Some(expected_id) => Some(expected_id.to_owned()),
|
||||||
|
None => store.current_rebalance_id().await,
|
||||||
|
};
|
||||||
|
match store.stop_rebalance_for_id(local_rebalance_id.as_deref()).await {
|
||||||
Ok(_) => {
|
Ok(_) => {
|
||||||
let save_result = match expected_rebalance_id {
|
let save_result = match local_rebalance_id.as_deref() {
|
||||||
Some(expected_id) => {
|
Some(expected_id) => {
|
||||||
store
|
store
|
||||||
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_id)
|
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_id)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
None => store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await,
|
None => Ok(()),
|
||||||
};
|
};
|
||||||
if let Err(err) = save_result {
|
if let Err(err) = save_result {
|
||||||
error!(
|
error!(
|
||||||
|
|||||||
@@ -330,6 +330,10 @@ impl ECStore {
|
|||||||
#[tracing::instrument(skip_all)]
|
#[tracing::instrument(skip_all)]
|
||||||
pub async fn load_rebalance_meta(&self) -> Result<()> {
|
pub async fn load_rebalance_meta(&self) -> Result<()> {
|
||||||
let _start_guard = self.start_gate.lock().await;
|
let _start_guard = self.start_gate.lock().await;
|
||||||
|
self.load_rebalance_meta_under_start_gate().await
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn load_rebalance_meta_under_start_gate(&self) -> Result<()> {
|
||||||
let mut meta = RebalanceMeta::new();
|
let mut meta = RebalanceMeta::new();
|
||||||
debug!(
|
debug!(
|
||||||
event = EVENT_REBALANCE_STATE,
|
event = EVENT_REBALANCE_STATE,
|
||||||
@@ -867,27 +871,31 @@ impl ECStore {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn record_rebalance_stop_propagation(self: &Arc<Self>, record: RebalanceStopPropagationRecord) -> Result<()> {
|
pub async fn record_rebalance_stop_propagation(
|
||||||
|
self: &Arc<Self>,
|
||||||
|
expected_id: &str,
|
||||||
|
record: RebalanceStopPropagationRecord,
|
||||||
|
) -> Result<()> {
|
||||||
if !record.has_failures() {
|
if !record.has_failures() {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let _start_guard = self.start_gate.lock().await;
|
||||||
|
let _activation_guard = self
|
||||||
|
.rebalance_activation_write_guard(Some(expected_id), "record rebalance stop propagation")
|
||||||
|
.await?;
|
||||||
let encoded_error = encode_rebalance_stop_propagation_record(&record);
|
let encoded_error = encode_rebalance_stop_propagation_record(&record);
|
||||||
let meta_to_save = {
|
let meta_to_save = {
|
||||||
let mut rebalance_meta = self.rebalance_meta.write().await;
|
let mut rebalance_meta = self.rebalance_meta.write().await;
|
||||||
|
ensure_rebalance_run_id(rebalance_meta.as_ref(), expected_id, "record rebalance stop propagation")?;
|
||||||
record_rebalance_stop_propagation_snapshot(rebalance_meta.as_mut(), encoded_error, OffsetDateTime::now_utc())
|
record_rebalance_stop_propagation_snapshot(rebalance_meta.as_mut(), encoded_error, OffsetDateTime::now_utc())
|
||||||
};
|
};
|
||||||
|
|
||||||
if let Some(meta_to_save) = meta_to_save {
|
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")?;
|
let pool = clone_first_arc(self.pools.as_slice(), "record_rebalance_stop_propagation: no pools available")?;
|
||||||
resolve_rebalance_meta_save_result(
|
resolve_rebalance_meta_save_result(
|
||||||
self.save_rebalance_meta_for_id_with_merge(
|
self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, "record_rebalance_stop_propagation", expected_id)
|
||||||
pool,
|
.await,
|
||||||
&meta_to_save,
|
|
||||||
"record_rebalance_stop_propagation",
|
|
||||||
meta_to_save.id.as_str(),
|
|
||||||
)
|
|
||||||
.await,
|
|
||||||
"record_rebalance_stop_propagation",
|
"record_rebalance_stop_propagation",
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -220,33 +220,22 @@ impl ECStore {
|
|||||||
let expected_bucket_incarnation_id = bucket_configs.bucket_incarnation_id;
|
let expected_bucket_incarnation_id = bucket_configs.bucket_incarnation_id;
|
||||||
let mut transfer = |src_pool_idx: usize, bucket: String, rd: GetObjectReader| {
|
let mut transfer = |src_pool_idx: usize, bucket: String, rd: GetObjectReader| {
|
||||||
let store = self.clone();
|
let store = self.clone();
|
||||||
let rebalance_id = Arc::clone(&rebalance_id);
|
|
||||||
async move {
|
async move {
|
||||||
let run_guard = store
|
store
|
||||||
.rebalance_run_guard(rebalance_id.as_ref(), "rebalance object migration")
|
|
||||||
.await?;
|
|
||||||
let result = store
|
|
||||||
.clone()
|
.clone()
|
||||||
.rebalance_object(src_pool_idx, bucket, rd, expected_bucket_incarnation_id)
|
.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
|
// 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.
|
// cross-pool target (excluding the source pool), not back onto the source set.
|
||||||
let mut delete_marker = |bucket: String, object: String, opts: ObjectOptions| {
|
let mut delete_marker = |bucket: String, object: String, opts: ObjectOptions| {
|
||||||
let store = self.clone();
|
let store = self.clone();
|
||||||
let rebalance_id = Arc::clone(&rebalance_id);
|
async move { store.delete_object(&bucket, &object, opts).await }
|
||||||
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 run_guard = self
|
||||||
|
.rebalance_run_guard(rebalance_id.as_ref(), "rebalance version migration")
|
||||||
|
.await?;
|
||||||
let result = migrate_entry_version(
|
let result = migrate_entry_version(
|
||||||
&RebalanceMigrationBackend::new(set.as_ref(), self.as_ref()),
|
&RebalanceMigrationBackend::new(set.as_ref(), self.as_ref()),
|
||||||
bucket.clone(),
|
bucket.clone(),
|
||||||
@@ -260,6 +249,7 @@ impl ECStore {
|
|||||||
&mut delete_marker,
|
&mut delete_marker,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
|
drop(run_guard);
|
||||||
|
|
||||||
if result.ignored {
|
if result.ignored {
|
||||||
if should_count_rebalance_version_complete(&result) {
|
if should_count_rebalance_version_complete(&result) {
|
||||||
|
|||||||
@@ -2775,6 +2775,10 @@ async fn test_old_worker_cannot_mutate_replacement_rebalance_state() {
|
|||||||
.save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, "rebalance-a")
|
.save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, "rebalance-a")
|
||||||
.await
|
.await
|
||||||
.expect_err("old save task must not persist replacement metadata"),
|
.expect_err("old save task must not persist replacement metadata"),
|
||||||
|
store
|
||||||
|
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, "rebalance-a")
|
||||||
|
.await
|
||||||
|
.expect_err("old stop path must not persist replacement metadata"),
|
||||||
] {
|
] {
|
||||||
assert!(err.to_string().contains("stale rebalance worker rejected"));
|
assert!(err.to_string().contains("stale rebalance worker rejected"));
|
||||||
}
|
}
|
||||||
@@ -2785,10 +2789,11 @@ async fn test_old_worker_cannot_mutate_replacement_rebalance_state() {
|
|||||||
assert!(meta.pool_stats[0].rebalanced_buckets.is_empty());
|
assert!(meta.pool_stats[0].rebalanced_buckets.is_empty());
|
||||||
assert_eq!(meta.pool_stats[0].bytes, 0);
|
assert_eq!(meta.pool_stats[0].bytes, 0);
|
||||||
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
|
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
|
||||||
|
assert!(meta.stopped_at.is_none());
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_stop_waits_for_active_rebalance_side_effect_guard() {
|
async fn test_stop_waits_for_active_rebalance_migration_guard() {
|
||||||
let meta = RebalanceMeta {
|
let meta = RebalanceMeta {
|
||||||
id: "rebalance-a".to_string(),
|
id: "rebalance-a".to_string(),
|
||||||
pool_stats: vec![RebalanceStats {
|
pool_stats: vec![RebalanceStats {
|
||||||
@@ -2803,9 +2808,9 @@ async fn test_stop_waits_for_active_rebalance_side_effect_guard() {
|
|||||||
};
|
};
|
||||||
let store = test_store_with_rebalance_meta(meta);
|
let store = test_store_with_rebalance_meta(meta);
|
||||||
let run_guard = store
|
let run_guard = store
|
||||||
.rebalance_run_guard("rebalance-a", "test side effect")
|
.rebalance_run_guard("rebalance-a", "rebalance remote-tier migration")
|
||||||
.await
|
.await
|
||||||
.expect("active run should admit the side effect");
|
.expect("active run should admit the remote-tier migration");
|
||||||
let mut stop = Box::pin(store.stop_rebalance_for_id(Some("rebalance-a")));
|
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 mut context = Context::from_waker(futures::task::noop_waker_ref());
|
||||||
|
|
||||||
@@ -2834,6 +2839,52 @@ async fn test_stop_waits_for_active_rebalance_side_effect_guard() {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_rebalance_metadata_reload_under_start_gate_does_not_reacquire_gate() {
|
||||||
|
let store = test_store_with_rebalance_meta(RebalanceMeta::default());
|
||||||
|
let _start_guard = store.start_gate.lock().await;
|
||||||
|
|
||||||
|
let err = store
|
||||||
|
.load_rebalance_meta_under_start_gate()
|
||||||
|
.await
|
||||||
|
.expect_err("empty test store should reach the metadata load without waiting on start_gate again");
|
||||||
|
|
||||||
|
assert!(err.to_string().contains("no pools available"));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_stale_stop_propagation_cannot_mutate_replacement_rebalance() {
|
||||||
|
let meta = RebalanceMeta {
|
||||||
|
id: "rebalance-b".to_string(),
|
||||||
|
pool_stats: vec![RebalanceStats {
|
||||||
|
participating: true,
|
||||||
|
info: RebalanceInfo {
|
||||||
|
status: RebalStatus::Started,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
..Default::default()
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let store = test_store_with_rebalance_meta(meta);
|
||||||
|
let record = RebalanceStopPropagationRecord {
|
||||||
|
stop_failures: vec!["old rebalance stop failed".to_string()],
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let err = store
|
||||||
|
.record_rebalance_stop_propagation("rebalance-a", record)
|
||||||
|
.await
|
||||||
|
.expect_err("old propagation failure must not mutate 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.last_refreshed_at.is_none());
|
||||||
|
assert!(meta.pool_stats[0].info.last_error.is_none());
|
||||||
|
}
|
||||||
|
|
||||||
fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc<crate::store::ECStore> {
|
fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc<crate::store::ECStore> {
|
||||||
let endpoint_pools: crate::layout::endpoints::EndpointServerPools = Vec::new().into();
|
let endpoint_pools: crate::layout::endpoints::EndpointServerPools = Vec::new().into();
|
||||||
Arc::new(crate::store::ECStore {
|
Arc::new(crate::store::ECStore {
|
||||||
|
|||||||
@@ -673,7 +673,7 @@ impl ECStore {
|
|||||||
self.save_rebalance_stats_inner(pool_idx, opt, None).await
|
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<()> {
|
pub 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
|
self.save_rebalance_stats_inner(pool_idx, opt, Some(expected_id)).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -177,12 +177,15 @@ async fn rollback_cluster_rebalance_start(
|
|||||||
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
|
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
|
||||||
terminal_reload_failures: terminal_reload_failures.clone(),
|
terminal_reload_failures: terminal_reload_failures.clone(),
|
||||||
};
|
};
|
||||||
store.record_rebalance_stop_propagation(record).await.map_err(|err| {
|
store
|
||||||
format!(
|
.record_rebalance_stop_propagation(rebalance_id, record)
|
||||||
"cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}",
|
.await
|
||||||
rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures)
|
.map_err(|err| {
|
||||||
)
|
format!(
|
||||||
})?;
|
"cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}",
|
||||||
|
rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures)
|
||||||
|
)
|
||||||
|
})?;
|
||||||
return Err(rebalance_rollback_failure_message(
|
return Err(rebalance_rollback_failure_message(
|
||||||
rebalance_id,
|
rebalance_id,
|
||||||
&stop_failures,
|
&stop_failures,
|
||||||
@@ -197,7 +200,7 @@ async fn rollback_cluster_rebalance_start(
|
|||||||
.await
|
.await
|
||||||
.map_err(|err| format!("local stop_rebalance rollback for {rebalance_id} failed: {err}"))?;
|
.map_err(|err| format!("local stop_rebalance rollback for {rebalance_id} failed: {err}"))?;
|
||||||
store
|
store
|
||||||
.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt)
|
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, rebalance_id)
|
||||||
.await
|
.await
|
||||||
.map_err(|err| format!("local rollback stop metadata save for {rebalance_id} failed: {err}"))?;
|
.map_err(|err| format!("local rollback stop metadata save for {rebalance_id} failed: {err}"))?;
|
||||||
Ok(())
|
Ok(())
|
||||||
@@ -679,7 +682,7 @@ impl Operation for RebalanceStart {
|
|||||||
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
|
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
|
||||||
terminal_reload_failures: terminal_reload_failures.clone(),
|
terminal_reload_failures: terminal_reload_failures.clone(),
|
||||||
};
|
};
|
||||||
store.record_rebalance_stop_propagation(record).await.map_err(|err| {
|
store.record_rebalance_stop_propagation(&id, record).await.map_err(|err| {
|
||||||
rebalance_internal_error(format!(
|
rebalance_internal_error(format!(
|
||||||
"failed to persist rebalance local-start rollback propagation metadata: {err}"
|
"failed to persist rebalance local-start rollback propagation metadata: {err}"
|
||||||
))
|
))
|
||||||
@@ -926,23 +929,26 @@ impl Operation for RebalanceStop {
|
|||||||
log_rebalance_request_rejected("stop", "rebalance_not_started", &request_id, &actor, &remote_addr);
|
log_rebalance_request_rejected("stop", "rebalance_not_started", &request_id, &actor, &remote_addr);
|
||||||
return Err(s3_error!(NoSuchResource, "pool rebalance is not started"));
|
return Err(s3_error!(NoSuchResource, "pool rebalance is not started"));
|
||||||
}
|
}
|
||||||
|
let Some(expected_rebalance_id) = expected_rebalance_id else {
|
||||||
|
return Err(s3_error!(InternalError, "active rebalance metadata has no activation id"));
|
||||||
|
};
|
||||||
|
|
||||||
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 mut stop_failures = Vec::new();
|
||||||
if let Some(notification_sys) = notification_sys.as_ref() {
|
if let Some(notification_sys) = notification_sys.as_ref() {
|
||||||
stop_failures = notification_sys
|
stop_failures = notification_sys
|
||||||
.stop_rebalance_failures(expected_rebalance_id.as_deref())
|
.stop_rebalance_failures(Some(expected_rebalance_id.as_str()))
|
||||||
.await
|
.await
|
||||||
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e))?;
|
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e))?;
|
||||||
} else {
|
} else {
|
||||||
store
|
store
|
||||||
.stop_rebalance_for_id(expected_rebalance_id.as_deref())
|
.stop_rebalance_for_id(Some(expected_rebalance_id.as_str()))
|
||||||
.await
|
.await
|
||||||
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?;
|
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?;
|
||||||
|
|
||||||
store
|
store
|
||||||
.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt)
|
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_rebalance_id.as_str())
|
||||||
.await
|
.await
|
||||||
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?;
|
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?;
|
||||||
}
|
}
|
||||||
@@ -1007,7 +1013,7 @@ impl Operation for RebalanceStop {
|
|||||||
terminal_reload_failures: terminal_reload_failures.clone(),
|
terminal_reload_failures: terminal_reload_failures.clone(),
|
||||||
};
|
};
|
||||||
store
|
store
|
||||||
.record_rebalance_stop_propagation(record)
|
.record_rebalance_stop_propagation(expected_rebalance_id.as_str(), record)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop propagation metadata: {}", e))?;
|
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop propagation metadata: {}", e))?;
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user