fix(rebalance): retry contended runtime metadata access (#7752)

* fix(rebalance): retry contended runtime metadata access

* chore(rebalance): shrink typed error ratchet baseline
This commit is contained in:
cxymds
2026-09-13 21:44:11 +08:00
committed by GitHub
parent 13093c5afc
commit 3ca3e26cec
5 changed files with 443 additions and 17 deletions
@@ -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<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
.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<ECStore>) {
let deployment_id = store
.ctx
+157 -4
View File
@@ -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() {
@@ -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,
+172 -5
View File
@@ -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<tokio::sync::Notify>;
}
/// 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<T, Access, AccessFuture>(
cancel: Option<&CancellationToken>,
max_attempts: usize,
mut access: Access,
) -> Result<T>
where
Access: FnMut() -> AccessFuture,
AccessFuture: std::future::Future<Output = Result<T>>,
{
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<std::sync::atomic::AtomicUsize>);
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::<Result<()>>().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 = [
+1 -1
View File
@@ -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