mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 04:21:35 +00:00
fix(tier): keep recovery worker after startup reload failure (#8260)
This commit is contained in:
@@ -602,9 +602,18 @@ pub(crate) async fn initialize_local_disk_maps(
|
||||
}
|
||||
|
||||
pub(crate) async fn init_tier_config_mgr(store: Arc<ECStore>) -> Result<()> {
|
||||
let handle = get_global_tier_config_mgr();
|
||||
TierConfigMgr::reload_handle(&handle, store.clone()).await?;
|
||||
tokio::spawn(TierConfigMgr::refresh_tier_config_handle(handle, store));
|
||||
init_tier_config_mgr_handle(get_global_tier_config_mgr(), store).await
|
||||
}
|
||||
|
||||
pub(crate) async fn init_tier_config_mgr_handle(handle: Arc<RwLock<TierConfigMgr>>, store: Arc<ECStore>) -> Result<()> {
|
||||
let initial_reload = TierConfigMgr::reload_handle(&handle, store.clone()).await;
|
||||
// Startup may continue after a transient recovery failure. Its retained
|
||||
// mutation fences still need a worker to consume later commit notifications.
|
||||
tokio::spawn(TierConfigMgr::refresh_tier_config_handle(handle.clone(), store));
|
||||
if initial_reload.is_err() {
|
||||
TierConfigMgr::request_committed_mutation_refresh(&handle).await;
|
||||
}
|
||||
initial_reload?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -20362,6 +20362,89 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn tier_config_initial_reload_lock_failure_keeps_recovery_worker() {
|
||||
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
|
||||
|
||||
let temp_dir = tempfile::tempdir().expect("create contended tier startup store dir");
|
||||
let (ctx, store, shutdown) =
|
||||
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-startup-contention", &[4])).await;
|
||||
let manager = ctx.tier_config_mgr();
|
||||
register_mock_tier(&manager, "COLD-A").await;
|
||||
let candidate = TierConfigMgr::new();
|
||||
let candidate = candidate.read().await;
|
||||
let candidate_digest = tier_config_candidate_digest(&candidate).expect("empty candidate digest should resolve");
|
||||
candidate
|
||||
.save_tiering_config(store.clone())
|
||||
.await
|
||||
.expect("committed removal config should persist");
|
||||
let config_info = store
|
||||
.get_object_info(
|
||||
RUSTFS_META_BUCKET,
|
||||
&format!("{}/{TIER_CONFIG_FILE}", com::CONFIG_PREFIX),
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await
|
||||
.expect("committed removal config metadata should load");
|
||||
let mutation_id = uuid::Uuid::new_v4();
|
||||
let mut intent = tier_mutation_peer_test_intent(mutation_id, "COLD-A", candidate_digest);
|
||||
intent.kind = TierMutationIntentKind::Remove;
|
||||
intent.affected_targets[0].new_backend_identity = None;
|
||||
intent
|
||||
.advance(
|
||||
TierMutationIntentState::Committed,
|
||||
Some(config_info.etag.expect("committed config should have an ETag")),
|
||||
)
|
||||
.expect("removal intent should commit");
|
||||
save_tier_mutation_intent_record(store.clone(), &intent)
|
||||
.await
|
||||
.expect("committed removal intent should persist");
|
||||
|
||||
let config_lock = format!("{}/{TIER_CONFIG_FILE}.lock", com::CONFIG_PREFIX);
|
||||
let namespace = store
|
||||
.new_ns_lock(RUSTFS_META_BUCKET, &config_lock)
|
||||
.await
|
||||
.expect("tier config namespace should resolve");
|
||||
let owner = namespace
|
||||
.get_write_lock(crate::set_disk::get_lock_acquire_timeout())
|
||||
.await
|
||||
.expect("competing recovery should hold the config lock");
|
||||
let initial_error = runtime_sources::init_tier_config_mgr_handle(manager.clone(), store.clone())
|
||||
.await
|
||||
.expect_err("startup must still report the initial recovery lock failure");
|
||||
assert!(initial_error.to_string().contains(&config_lock));
|
||||
match TierConfigMgr::acquire_operation_lease(&manager, "COLD-A").await {
|
||||
Err(err) => assert_eq!(err.message, "Remote tier configuration is being replaced"),
|
||||
Ok(_) => panic!("failed initial recovery must retain its mutation fence"),
|
||||
}
|
||||
drop(owner);
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(30), async {
|
||||
loop {
|
||||
if manager.read().await.tiers.is_empty()
|
||||
&& matches!(
|
||||
TierConfigMgr::acquire_operation_lease(&manager, "COLD-A").await,
|
||||
Err(err) if err.code == crate::services::tier::tier_handlers::ERR_TIER_NOT_FOUND.code
|
||||
)
|
||||
{
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("recovery worker must publish and clear its fence without another admin mutation");
|
||||
assert_eq!(
|
||||
load_tier_mutation_intent_record(store.clone(), mutation_id)
|
||||
.await
|
||||
.expect("unexpired committed tombstone must remain durable"),
|
||||
intent
|
||||
);
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
|
||||
Reference in New Issue
Block a user