From bb1b5dea1672de6570fa963f8d2425b251385381 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 9 Sep 2026 16:32:17 +0800 Subject: [PATCH] fix(ecstore): recover expansion buckets and resume retirement (#7574) * fix(ecstore): recover partially created expansion bucket volumes * fix(ecstore): resume retirement after startup safety checks --- crates/ecstore/src/bucket/metadata_sys.rs | 10 +- crates/ecstore/src/core/pools.rs | 3 + crates/ecstore/src/store/bucket.rs | 117 ++++++++++++++++++++++ crates/ecstore/src/store/init.rs | 93 ++++++++++++----- 4 files changed, 195 insertions(+), 28 deletions(-) diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index c20825612..1cdea7e5c 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -1681,9 +1681,13 @@ impl BucketMetadataSys { expected: Option<&Arc>, namespace_guard: &rustfs_lock::NamespaceLockGuard, ) -> Result<()> { - if !self - .bucket_exists(bucket, namespace_guard, "bucket metadata existence check") - .await? + if !await_bucket_namespace_operation( + Some(namespace_guard), + bucket, + "bucket metadata heal existence check", + self.api.bucket_exists_for_heal(bucket), + ) + .await? { if matches!(mode, MetadataLoadMode::Refresh) { let _publish_guard = self diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 39bdfff70..3fe028e48 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -12721,6 +12721,9 @@ impl ECStore { self: &Arc, rx: CancellationToken, ) -> Result<()> { + #[cfg(test)] + let endpoints = self.instance_endpoints().unwrap_or_else(|| self.endpoints()); + #[cfg(not(test))] let endpoints = self.endpoints(); let index_cancelers = self.reserve_missing_local_decommission_routines(&rx, &endpoints).await?; if index_cancelers.is_empty() { diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index c648255ad..132cb8b67 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -770,6 +770,31 @@ impl ECStore { Ok(()) } + /// Prove a live bucket generation before repairing missing expansion volumes. + /// Unlike request validation, repair only needs one erasure set to confirm + /// existence; an incomplete expansion set is precisely what repair fixes. + /// Callers must hold the bucket namespace lock through the subsequent heal. + pub(crate) async fn bucket_exists_for_heal(&self, bucket: &str) -> Result { + let results = futures::future::join_all( + self.bucket_sets() + .map(|(_, _, set)| async move { set.get_bucket_info(bucket, &BucketOptions::default()).await }), + ) + .await; + let mut first_error = None; + for result in results { + match result { + Ok(_) => return Ok(true), + Err(err) if is_err_strict_volume_not_found(&err) => {} + Err(err) if first_error.is_none() => first_error = Some(err), + Err(_) => {} + } + } + match first_error { + Some(err) => Err(err), + None => Ok(false), + } + } + #[instrument(skip(self))] pub(crate) async fn get_bucket_info_from_sets(&self, bucket: &str, opts: &BucketOptions) -> Result { // One host may participate in several pools after expansion. Resolve the @@ -2050,6 +2075,98 @@ mod tests { .expect("metadata initialization should recreate the bucket volume in the new pool"); } + #[tokio::test] + #[serial] + async fn bucket_metadata_init_repairs_half_created_expansion_pool() { + // Check both pool orders: an incomplete set must not hide a later + // complete set, and a complete set must not weaken request validation. + for complete_pool in 0..2 { + let (temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("partial-expansion-{}", Uuid::new_v4().simple()); + for pool_index in 0..2 { + let present_disks = if pool_index == complete_pool { 4 } else { 2 }; + for disk_index in 0..present_disks { + tokio::fs::create_dir( + temp_dir + .path() + .join(format!("pool{pool_index}-disk{disk_index}")) + .join(&bucket), + ) + .await + .expect("fixture bucket volume should be created"); + } + } + assert_eq!( + ecstore + .get_bucket_info_from_sets(&bucket, &BucketOptions::default()) + .await + .expect_err("request validation must reject a half-created expansion set"), + StorageError::ErasureWriteQuorum + ); + + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), vec![bucket.clone()]).await; + + for pool_index in 0..2 { + for disk_index in 0..4 { + assert!( + temp_dir + .path() + .join(format!("pool{pool_index}-disk{disk_index}")) + .join(&bucket) + .is_dir(), + "metadata initialization must heal every missing expansion volume" + ); + } + } + ecstore + .get_bucket_info_from_sets(&bucket, &BucketOptions::default()) + .await + .expect("strict request validation should succeed after volume repair"); + } + } + + #[tokio::test] + #[serial] + async fn bucket_metadata_init_does_not_combine_partial_set_evidence() { + let (temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; + let bucket = format!("no-quorum-expansion-{}", Uuid::new_v4().simple()); + for pool_index in 0..2 { + for disk_index in 0..2 { + tokio::fs::create_dir( + temp_dir + .path() + .join(format!("pool{pool_index}-disk{disk_index}")) + .join(&bucket), + ) + .await + .expect("fixture bucket volume should be created"); + } + } + assert_eq!( + ecstore + .bucket_exists_for_heal(&bucket) + .await + .expect_err("repair must require a complete quorum within one set"), + StorageError::ErasureWriteQuorum + ); + + metadata_sys::init_bucket_metadata_sys(ecstore.clone(), vec![bucket.clone()]).await; + + for pool_index in 0..2 { + for disk_index in 0..4 { + assert_eq!( + temp_dir + .path() + .join(format!("pool{pool_index}-disk{disk_index}")) + .join(&bucket) + .is_dir(), + disk_index < 2, + "unproven bucket generations must not recreate missing volumes" + ); + } + } + } + #[tokio::test] #[serial] async fn bucket_metadata_init_does_not_recreate_stale_bucket_name() { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index d64083a39..c8d44beb0 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -95,7 +95,6 @@ fn preflight_startup_rpc_secret_with( } } -const LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY: Duration = Duration::from_secs(60 * 3); const LOCAL_DECOMMISSION_RESUME_RETRY_DELAY: Duration = Duration::from_secs(30); const LOCAL_DECOMMISSION_WATCHDOG_INTERVAL: Duration = Duration::from_secs(30); const LOCAL_DECOMMISSION_WATCHDOG_MAX_RETRY_DELAY: Duration = Duration::from_secs(60 * 5); @@ -280,20 +279,26 @@ where } } +async fn reconcile_local_decommission_after_init(store: &Arc, rx: CancellationToken) -> Result<()> { + store + .ensure_pool_meta_side_effects_safe("decommission worker recovery blocked while pool metadata requires recovery") + .await?; + if store.has_active_local_decommission_worker().await { + return Ok(()); + } + store.refresh_pool_status_meta().await?; + let resume_required = pool_meta_has_active_decommission(&*store.pool_meta.read().await); + if resume_required { + crate::core::pools::acquire_pool_activation_fleet_proof(&store.ctx).await?; + } + store.spawn_missing_local_decommission_routines_with_token(rx).await +} + async fn supervise_local_decommission_after_init(store: Arc, rx: CancellationToken) { run_local_decommission_watchdog(rx.clone(), || { let store = store.clone(); let worker_rx = rx.clone(); - async move { - store - .ensure_pool_meta_side_effects_safe("decommission worker recovery blocked while pool metadata requires recovery") - .await?; - if store.has_active_local_decommission_worker().await { - return Ok(()); - } - store.refresh_pool_status_meta().await?; - store.spawn_missing_local_decommission_routines_with_token(worker_rx).await - } + async move { reconcile_local_decommission_after_init(&store, worker_rx).await } }) .await; } @@ -784,14 +789,9 @@ impl ECStore { ); } if has_local_decommission_leadership { - let store = self.clone(); - let decommission_rx = rx.clone(); - tokio::spawn(async move { - if !wait_for_local_decommission_resume_delay(&decommission_rx, LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY).await { - return; - } - supervise_local_decommission_after_init(store, decommission_rx).await; - }); + // The watchdog checks recovery safety and retries transient failures. + // Resume persisted work without an unconditional cold-start delay. + tokio::spawn(supervise_local_decommission_after_init(self.clone(), rx.clone())); } let recovery_store = self.clone(); @@ -2104,6 +2104,44 @@ mod tests { ); } + #[tokio::test(start_paused = true)] + async fn test_local_decommission_watchdog_cancelled_start_does_not_reconcile() { + let rx = CancellationToken::new(); + rx.cancel(); + run_local_decommission_watchdog(rx, || async { + panic!("cancelled startup must not schedule persisted work"); + }) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn test_local_decommission_recovery_waits_for_live_fleet_proof_before_reserving_worker() { + let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await; + mark_test_pool_decommissioning(&store, 0).await; + assert!(store.ctx.is_dist_erasure().await); + let worker_rx = CancellationToken::new(); + + { + let _proof_guard = crate::services::notification_sys::without_cross_pool_fence_fleet_proof_for_test(); + let err = super::reconcile_local_decommission_after_init(&store, worker_rx.clone()) + .await + .expect_err("cold distributed recovery must wait for live fleet proof"); + assert!( + crate::core::pools::is_pool_activation_fleet_proof_error(&err), + "recovery must reach the live fleet proof gate: {err:?}" + ); + assert!(store.decommission_cancelers.read().await.iter().all(Option::is_none)); + assert!(pool_meta_has_active_decommission(&*store.pool_meta.read().await)); + } + + super::reconcile_local_decommission_after_init(&store, worker_rx.clone()) + .await + .expect("restored fleet proof should admit the persisted worker"); + assert!(store.has_active_local_decommission_worker().await); + worker_rx.cancel(); + } + #[tokio::test(start_paused = true)] async fn test_local_decommission_watchdog_retries_general_failures_until_cancelled() { let rx = CancellationToken::new(); @@ -2115,11 +2153,13 @@ mod tests { let attempts = attempts.clone(); let rx = rx.clone(); async move { - if attempts.fetch_add(1, Ordering::SeqCst) == 0 { - Err(StorageError::SlowDown) - } else { - rx.cancel(); - Ok(()) + match attempts.fetch_add(1, Ordering::SeqCst) { + 0 => Err(StorageError::other("pool activation requires a live fleet capability proof")), + 1 => Err(StorageError::SlowDown), + _ => { + rx.cancel(); + Ok(()) + } } } } @@ -2128,8 +2168,11 @@ mod tests { tokio::task::yield_now().await; assert_eq!(attempts.load(Ordering::SeqCst), 1); tokio::time::advance(LOCAL_DECOMMISSION_RESUME_RETRY_DELAY).await; - task.await.expect("watchdog task should exit after cancellation"); + tokio::task::yield_now().await; assert_eq!(attempts.load(Ordering::SeqCst), 2); + tokio::time::advance(local_decommission_watchdog_retry_delay(2)).await; + task.await.expect("watchdog task should exit after cancellation"); + assert_eq!(attempts.load(Ordering::SeqCst), 3); } #[tokio::test(start_paused = true)]