diff --git a/crates/ecstore/src/runtime/sources.rs b/crates/ecstore/src/runtime/sources.rs index 54150d6ae..f121cffe5 100644 --- a/crates/ecstore/src/runtime/sources.rs +++ b/crates/ecstore/src/runtime/sources.rs @@ -602,9 +602,18 @@ pub(crate) async fn initialize_local_disk_maps( } pub(crate) async fn init_tier_config_mgr(store: Arc) -> 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>, store: Arc) -> 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(()) } diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 166a90138..e7d1f5b76 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -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)]