diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index abf69d5be..2a35d732b 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -8,8 +8,8 @@ use super::meta::{ validate_init_rebalance_state, }; use super::worker::{ - rebalance_meta_lock_error, resolve_load_rebalance_stats_update_result, resolve_rebalance_meta_load_result, - resolve_rebalance_meta_save_result, + rebalance_max_attempts, rebalance_meta_lock_error, resolve_load_rebalance_stats_update_result, + resolve_rebalance_meta_load_result, resolve_rebalance_meta_save_result, retry_rebalance_metadata_access, }; use super::{ DiskStat, EVENT_REBALANCE_BUCKET, EVENT_REBALANCE_STATE, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REBALANCE, REBAL_META_NAME, @@ -513,10 +513,13 @@ impl ECStore { S: EcstoreObjectIO + StorageNamespaceLocking, { let ns_lock = pool.new_ns_lock(crate::disk::RUSTFS_META_BUCKET, REBAL_META_NAME).await?; - let guard = ns_lock - .get_write_lock(get_lock_acquire_timeout()) - .await - .map_err(|err| rebalance_meta_lock_error(err, "write"))?; + let guard = retry_rebalance_metadata_access(None, rebalance_max_attempts(), || async { + ns_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|err| rebalance_meta_lock_error(err, "write")) + }) + .await?; let mut opts = ObjectOptions { no_lock: true, ..Default::default() @@ -1348,6 +1351,105 @@ mod tests { use crate::object_api::NamespaceLockFence; use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause, hermetic_set_disks_isolated}; + #[tokio::test] + async fn rebalance_status_read_recovers_after_metadata_lock_timeout() { + let id = "rebalance-status-lock-retry"; + let (_temp_dirs, store) = crate::services::rebalance::test_store_with_persisted_rebalance_meta(RebalanceMeta { + id: id.to_string(), + pool_stats: vec![RebalanceStats::default()], + ..Default::default() + }) + .await; + let ns_lock = store.pools[0] + .new_ns_lock(crate::disk::RUSTFS_META_BUCKET, REBAL_META_NAME) + .await + .expect("metadata lock should be available"); + let writer = ns_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .expect("hold metadata writer"); + let retried = Arc::new(tokio::sync::Notify::new()); + let refresh = super::super::worker::REBALANCE_METADATA_RETRY_PROBE + .scope(Arc::clone(&retried), store.refresh_rebalance_status_meta()); + tokio::pin!(refresh); + tokio::time::timeout(std::time::Duration::from_secs(30), async { + tokio::select! { + _ = retried.notified() => {}, + result = &mut refresh => panic!("status read must retry its contended metadata lock: {result:?}"), + } + }) + .await + .expect("status read should encounter a real lock timeout"); + drop(writer); + tokio::time::timeout(std::time::Duration::from_secs(30), refresh) + .await + .expect("status refresh should finish after the writer releases") + .expect("transient metadata contention must not fail status refresh"); + assert_eq!( + store + .rebalance_meta + .read() + .await + .as_ref() + .expect("metadata must remain present") + .id, + id + ); + } + + #[tokio::test] + async fn rebalance_stats_save_recovers_after_metadata_lock_timeout() { + let id = "rebalance-stats-lock-retry"; + let (_temp_dirs, store) = crate::services::rebalance::test_store_with_persisted_rebalance_meta(RebalanceMeta { + id: id.to_string(), + pool_stats: vec![RebalanceStats::default()], + ..Default::default() + }) + .await; + store + .rebalance_meta + .write() + .await + .as_mut() + .expect("local metadata must exist") + .pool_stats[0] + .num_objects = 7; + let ns_lock = store.pools[0] + .new_ns_lock(crate::disk::RUSTFS_META_BUCKET, REBAL_META_NAME) + .await + .expect("metadata lock should be available"); + let reader = ns_lock + .get_read_lock(get_lock_acquire_timeout()) + .await + .expect("hold a migration read fence"); + let retried = Arc::new(tokio::sync::Notify::new()); + let save = super::super::worker::REBALANCE_METADATA_RETRY_PROBE.scope( + Arc::clone(&retried), + store.save_rebalance_stats_for_id(0, super::super::RebalSaveOpt::Stats, id), + ); + tokio::pin!(save); + tokio::time::timeout(std::time::Duration::from_secs(30), async { + tokio::select! { + _ = retried.notified() => {}, + result = &mut save => panic!("stats save must retry its contended metadata lock: {result:?}"), + } + }) + .await + .expect("stats save should encounter a real lock timeout"); + drop(reader); + tokio::time::timeout(std::time::Duration::from_secs(30), save) + .await + .expect("stats save should finish after the migration read fence releases") + .expect("transient metadata contention must not lose stats"); + let mut persisted = RebalanceMeta::new(); + persisted + .load(Arc::clone(&store.pools[0])) + .await + .expect("persisted stats must be readable"); + assert_eq!(persisted.id, id); + assert_eq!(persisted.pool_stats[0].num_objects, 7); + } + async fn persist_initialized_identity_then_remove_pool_meta(store: &Arc) { let deployment_id = store .ctx diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 2a64c7c61..a58d76edc 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -21,9 +21,9 @@ use super::worker::{ RebalanceEntryCleanupResult, RebalanceEntryTask, load_rebalance_bucket_configs, rebalance_max_attempts, record_rebalance_error, resolve_rebalance_bucket_error, resolve_rebalance_entry_cleanup_delete_result, resolve_rebalance_file_info_versions_result, resolve_rebalance_migrate_result_error, resolve_rebalance_stats_update_result, - resolve_rebalance_worker_result, run_rebalance_listing_with_retry, should_cleanup_rebalance_source_entry, - should_count_rebalance_version_complete, should_defer_rebalance_entry_failure, should_skip_rebalance_delete_marker, - wait_rebalance_entry_tasks, with_rebalance_entry_context, + resolve_rebalance_worker_result, retry_rebalance_metadata_access, run_rebalance_listing_with_retry, + should_cleanup_rebalance_source_entry, should_count_rebalance_version_complete, should_defer_rebalance_entry_failure, + should_skip_rebalance_delete_marker, wait_rebalance_entry_tasks, with_rebalance_entry_context, }; use super::{ EVENT_REBALANCE_BUCKET, EVENT_REBALANCE_ENTRY, EVENT_REBALANCE_STATE, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REBALANCE, @@ -264,7 +264,10 @@ impl ECStore { // Target capacity admission can then acquire pool.bin under the run fence. // Stop waits for in-flight entries through cleanup, but not for entries admitted later. ensure_rebalance_entry_active(&cancel)?; - let run_guard = self.rebalance_run_guard(rebalance_id.as_ref(), "rebalance entry").await?; + let run_guard = retry_rebalance_metadata_access(Some(&cancel), rebalance_max_attempts(), || { + self.rebalance_run_guard(rebalance_id.as_ref(), "rebalance entry") + }) + .await?; #[cfg(test)] if let Ok((arrived, release)) = REBALANCE_ENTRY_RUN_FENCE_BARRIER.try_with(Clone::clone) { arrived.notify_one(); @@ -981,6 +984,7 @@ mod tests { }; use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; use crate::storage_api_contracts::multipart::{CompletePart, MultipartOperations as _}; + use crate::storage_api_contracts::namespace::NamespaceLocking as _; use crate::storage_api_contracts::object::ObjectIO as _; use http::HeaderMap; use rustfs_filemeta::{FileInfo, FileMeta, ObjectPartInfo, TransitionVersionState}; @@ -1248,6 +1252,155 @@ mod tests { assert_eq!(pool_stats.cleanup_warnings.count, 1, "deferred cleanup must not add a permanent warning"); } + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_entry_recovers_after_run_fence_lock_timeout() { + assert_real_rebalance_entry_metadata_retry(EntryRetryAction::Complete).await; + } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_entry_cancels_after_run_fence_lock_timeout() { + assert_real_rebalance_entry_metadata_retry(EntryRetryAction::Cancel).await; + } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_entry_rechecks_persisted_run_after_lock_timeout() { + assert_real_rebalance_entry_metadata_retry(EntryRetryAction::ReplaceRun).await; + } + + enum EntryRetryAction { + Complete, + Cancel, + ReplaceRun, + } + + async fn assert_real_rebalance_entry_metadata_retry(action: EntryRetryAction) { + const REBALANCE_ID: &str = "rebalance-entry-lock-retry"; + let (_temp_dirs, store, _peer) = crate::services::rebalance::test_two_pool_stores_with_isolated_node_contexts(Some( + active_rebalance_meta(REBALANCE_ID), + )) + .await; + let bucket = crate::disk::RUSTFS_META_BUCKET; + let object = "rebalance-entry-lock-retry-object"; + let payload = b"metadata contention must not fail a rebalance entry".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 opts = ObjectOptions { + versioned: true, + version_id: Some(uuid::Uuid::new_v4().to_string()), + ..Default::default() + }; + let source_before = source_set + .put_object(bucket, object, &mut PutObjReader::from_vec(payload.clone()), &opts) + .await + .expect("source version should be written"); + let entry = metacache_entry_from_source(&source_set, bucket, object).await; + let ns_lock = store.pools[0] + .new_ns_lock(bucket, super::super::REBAL_META_NAME) + .await + .expect("metadata lock should be available"); + let writer = ns_lock + .get_write_lock(crate::set_disk::get_lock_acquire_timeout()) + .await + .expect("hold the metadata write lock"); + let retried = Arc::new(tokio::sync::Notify::new()); + let cancel = CancellationToken::new(); + let transfer = super::super::worker::REBALANCE_METADATA_RETRY_PROBE.scope( + Arc::clone(&retried), + Arc::clone(&store).rebalance_entry( + RebalanceEntryTarget { + bucket: bucket.to_string(), + pool_index: 0, + }, + entry, + source_set, + Arc::new(RebalanceBucketConfigs::default()), + Arc::from(REBALANCE_ID), + cancel.clone(), + ), + ); + tokio::pin!(transfer); + tokio::time::timeout(StdDuration::from_secs(30), async { + tokio::select! { + _ = retried.notified() => {}, + result = &mut transfer => panic!("entry admission must retry its contended run fence: {result:?}"), + } + }) + .await + .expect("entry should encounter a real metadata lock timeout"); + match action { + EntryRetryAction::Complete => {} + EntryRetryAction::Cancel => cancel.cancel(), + EntryRetryAction::ReplaceRun => { + let mut save_opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + save_opts.add_namespace_lock_guard(&writer); + active_rebalance_meta("replacement-rebalance-run") + .save_with_opts(Arc::clone(&store.pools[0]), save_opts) + .await + .expect("a peer may replace the persisted run while holding the metadata writer"); + } + } + drop(writer); + let result = tokio::time::timeout(StdDuration::from_secs(30), transfer) + .await + .expect("entry should finish after the metadata writer releases"); + if !matches!(action, EntryRetryAction::Complete) { + let err = result.expect_err("a cancelled or replaced run must not resume migration"); + match action { + EntryRetryAction::Cancel => assert!(matches!(err, Error::OperationCanceled)), + EntryRetryAction::ReplaceRun => assert!(err.to_string().contains("stale rebalance run rejected")), + EntryRetryAction::Complete => unreachable!(), + } + let mut reader = store.pools[0] + .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + .await + .expect("rejected entry must retain its source"); + let mut actual = Vec::new(); + reader + .stream + .read_to_end(&mut actual) + .await + .expect("retained source must remain readable"); + assert_eq!(actual, payload); + let target_err = target_set + .get_object_info(bucket, object, &opts) + .await + .expect_err("rejected entry must not publish a target version"); + assert!(crate::error::is_err_object_not_found(&target_err) || crate::error::is_err_version_not_found(&target_err)); + return; + } + assert!(matches!( + result.expect("transient metadata contention must not fail the entry"), + RebalanceEntryOutcome::Completed + )); + let mut reader = target_set + .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) + .await + .expect("target version must be readable"); + let mut actual = Vec::new(); + reader + .stream + .read_to_end(&mut actual) + .await + .expect("target body should drain completely"); + assert_eq!(actual, payload); + assert_eq!(reader.object_info.version_id, source_before.version_id); + assert_eq!(reader.object_info.etag, source_before.etag); + let source_error = store.pools[0] + .get_object_info(bucket, object, &opts) + .await + .expect_err("completed migration should remove the source version"); + assert!(crate::error::is_err_object_not_found(&source_error) || crate::error::is_err_version_not_found(&source_error)); + let meta = store.rebalance_meta.read().await; + let stats = &meta.as_ref().expect("run must remain installed").pool_stats[0]; + assert_eq!((stats.num_objects, stats.num_versions, stats.cleanup_warnings.count), (1, 1, 0)); + } + #[tokio::test] #[serial_test::serial] async fn real_rebalance_entry_progresses_while_peer_activation_waits_for_run_fence() { diff --git a/crates/ecstore/src/services/rebalance/meta.rs b/crates/ecstore/src/services/rebalance/meta.rs index 8e383065d..8983f40cd 100644 --- a/crates/ecstore/src/services/rebalance/meta.rs +++ b/crates/ecstore/src/services/rebalance/meta.rs @@ -1,3 +1,4 @@ +use super::worker::{rebalance_max_attempts, retry_rebalance_metadata_access}; use super::{ EVENT_REBALANCE_BUCKET, EVENT_REBALANCE_STATE, Error, GetObjectReader, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REBALANCE, ObjectInfo, ObjectOptions, PutObjReader, REBAL_META_FMT, REBAL_META_NAME, REBAL_META_VER, @@ -100,7 +101,10 @@ impl RebalanceMeta { PutObjectReader = PutObjReader, >, { - let (data, _) = read_config_with_metadata(store, REBAL_META_NAME, &opts).await?; + let (data, _) = retry_rebalance_metadata_access(None, rebalance_max_attempts(), || { + read_config_with_metadata(Arc::clone(&store), REBAL_META_NAME, &opts) + }) + .await?; if data.is_empty() { debug!( event = EVENT_REBALANCE_STATE, diff --git a/crates/ecstore/src/services/rebalance/worker.rs b/crates/ecstore/src/services/rebalance/worker.rs index c0c93a7f4..c61a78339 100644 --- a/crates/ecstore/src/services/rebalance/worker.rs +++ b/crates/ecstore/src/services/rebalance/worker.rs @@ -21,6 +21,11 @@ use tokio::time::Duration; use tokio_util::sync::CancellationToken; use tracing::{debug, error, info}; +#[cfg(test)] +tokio::task_local! { + pub(super) static REBALANCE_METADATA_RETRY_PROBE: Arc; +} + /// Background walks skip the total timeout, so the per-read stall budget is what /// catches a drive that stops answering. Keep it generous: rebalance is not /// latency-sensitive, and one slow read is not a dead drive. @@ -116,11 +121,14 @@ pub(super) fn rebalance_meta_lock_error(err: rustfs_lock::LockError, mode: &'sta required, achieved, }, - other => Error::other(format!( - "failed to acquire rebalance metadata {mode} lock on {}/{}: {other}", - crate::disk::RUSTFS_META_BUCKET, - REBAL_META_NAME - )), + other => crate::data_movement::data_movement_context_error( + format!( + "failed to acquire rebalance metadata {mode} lock on {}/{}: {other}", + crate::disk::RUSTFS_META_BUCKET, + REBAL_META_NAME + ), + Error::Lock(other), + ), } } @@ -313,6 +321,49 @@ pub(super) fn rebalance_max_attempts() -> usize { parse_rebalance_max_attempts(std::env::var(REBALANCE_MAX_ATTEMPTS_ENV).ok().as_deref()) } +/// Retry lock admission or read-only metadata access, never a data mutation. +/// Each failed attempt must release its guards before the backoff so queued +/// writers and stop/replacement activation can make progress. +pub(super) async fn retry_rebalance_metadata_access( + cancel: Option<&CancellationToken>, + max_attempts: usize, + mut access: Access, +) -> Result +where + Access: FnMut() -> AccessFuture, + AccessFuture: std::future::Future>, +{ + let mut attempt = 0usize; + loop { + let result = match cancel { + Some(cancel) => tokio::select! { + biased; + _ = cancel.cancelled() => return Err(Error::OperationCanceled), + result = access() => result, + }, + None => access().await, + }; + match result { + Ok(value) => return Ok(value), + Err(err) => { + if attempt.saturating_add(1) >= max_attempts.max(1) + || !matches!(rebalance_error_source(&err), Error::Lock(lock_err) if is_rebalance_transient_lock_error(lock_err)) + { + return Err(err); + } + #[cfg(test)] + let _ = REBALANCE_METADATA_RETRY_PROBE.try_with(|probe| probe.notify_one()); + let delay = rebalance_migration_retry_delay(attempt, &err); + match cancel { + Some(cancel) => wait_rebalance_listing_retry(cancel, delay).await?, + None => tokio::time::sleep(delay).await, + } + attempt += 1; + } + } + } +} + pub(super) fn rebalance_listing_retry_delay(attempt: usize) -> Duration { let multiplier = u32::try_from(attempt.saturating_add(1)).unwrap_or(u32::MAX); REBALANCE_LISTING_RETRY_BASE_DELAY.saturating_mul(multiplier) @@ -601,6 +652,122 @@ impl SetDisks { mod error_source_tests { use super::*; + #[tokio::test] + async fn rebalance_metadata_retry_is_bounded_and_retains_the_timeout() { + for max_attempts in [0, 1, 3] { + let mut attempts = 0; + let result = retry_rebalance_metadata_access(None, max_attempts, || { + attempts += 1; + std::future::ready(Err::<(), _>(rebalance_meta_lock_error( + rustfs_lock::LockError::timeout(".rustfs.sys/rebalance.bin@latest", Duration::from_secs(5)), + "read", + ))) + }) + .await; + let err = result.expect_err("persistent lock contention must not become success"); + assert_eq!(attempts, max_attempts.max(1)); + assert!(matches!( + rebalance_error_source(&err), + Error::Lock(rustfs_lock::LockError::Timeout { .. }) + )); + } + } + + #[tokio::test] + async fn rebalance_metadata_retry_does_not_retry_permanent_or_untyped_errors() { + for err in [ + Error::FileAccessDenied, + Error::DiskFull, + Error::NamespaceLockQuorumUnavailable { + mode: "read", + bucket: crate::disk::RUSTFS_META_BUCKET.to_string(), + object: REBAL_META_NAME.to_string(), + required: 3, + achieved: 2, + }, + Error::other("stale rebalance run rejected: lock acquisition timed out"), + Error::other("rebalance distributed run fence lost"), + Error::Io(std::io::Error::from(std::io::ErrorKind::TimedOut)), + ] { + let expected = err.to_string(); + let mut error = Some(crate::data_movement::data_movement_context_error(expected.clone(), err)); + let mut attempts = 0; + let result = retry_rebalance_metadata_access(None, 3, || { + attempts += 1; + std::future::ready(Err::<(), _>(error.take().expect("permanent metadata failures must not be retried"))) + }) + .await; + assert_eq!(attempts, 1); + assert_eq!( + rebalance_error_source(&result.expect_err("failure must remain visible")).to_string(), + expected + ); + } + } + + #[tokio::test] + async fn rebalance_metadata_retry_cancels_a_pending_attempt() { + struct DropProbe(Arc); + impl Drop for DropProbe { + fn drop(&mut self) { + self.0.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + } + } + let cancel = CancellationToken::new(); + let started = tokio::sync::Notify::new(); + let dropped = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let access = retry_rebalance_metadata_access(Some(&cancel), 3, || async { + let _probe = DropProbe(Arc::clone(&dropped)); + started.notify_one(); + std::future::pending::>().await + }); + let (result, ()) = tokio::time::timeout(Duration::from_secs(5), async { + tokio::join!(access, async { + started.notified().await; + cancel.cancel(); + }) + }) + .await + .expect("cancellation must interrupt a pending lock attempt"); + assert!(matches!(result, Err(Error::OperationCanceled))); + assert_eq!(dropped.load(std::sync::atomic::Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn rebalance_metadata_retry_cancels_before_another_attempt() { + let cancel = CancellationToken::new(); + let mut attempts = 0; + let result = retry_rebalance_metadata_access(Some(&cancel), 3, || { + attempts += 1; + let cancel = &cancel; + async move { + cancel.cancel(); + Err::<(), _>(Error::Lock(rustfs_lock::LockError::timeout(REBAL_META_NAME, Duration::from_secs(5)))) + } + }) + .await; + assert!(matches!(result, Err(Error::OperationCanceled))); + assert_eq!(attempts, 1); + } + + #[test] + fn rebalance_metadata_lock_timeout_preserves_retryable_source() { + let resource = ".rustfs.sys/rebalance.bin@latest"; + for mode in ["read", "write"] { + let error = rebalance_meta_lock_error(rustfs_lock::LockError::timeout(resource, Duration::from_secs(5)), mode); + assert!( + is_transient_rebalance_error(&error), + "metadata lock contention must remain retryable: {error}" + ); + assert!(is_rebalance_lock_or_rpc_timeout(&error)); + assert!(matches!( + rebalance_error_source(&error), + Error::Lock(rustfs_lock::LockError::Timeout { resource: actual, .. }) if actual == resource + )); + assert!(error.to_string().contains(&format!("rebalance metadata {mode} lock"))); + } + } + #[test] fn stage_wrapped_errors_select_the_source_backoff_policy() { let cases = [ diff --git a/scripts/error-other-format-baseline.txt b/scripts/error-other-format-baseline.txt index 7b7abe72a..5939a7027 100644 --- a/scripts/error-other-format-baseline.txt +++ b/scripts/error-other-format-baseline.txt @@ -50,7 +50,7 @@ 1|crates/ecstore/src/services/rebalance/entry.rs 8|crates/ecstore/src/services/rebalance/meta.rs 8|crates/ecstore/src/services/rebalance/runtime.rs -19|crates/ecstore/src/services/rebalance/worker.rs +18|crates/ecstore/src/services/rebalance/worker.rs 33|crates/ecstore/src/services/tier/tier.rs 1|crates/ecstore/src/services/tier/tier_config.rs 1|crates/ecstore/src/services/tier/warm_backend.rs