From c9acf01fd03653cac6806b71b5485b19109bf1b2 Mon Sep 17 00:00:00 2001 From: Chris Date: Tue, 29 Sep 2026 07:16:01 +0800 Subject: [PATCH] fix(tier): reread mutation intents after lock contention (#8209) --- .../src/services/tier/tier_mutation_peer.rs | 6 +- crates/ecstore/src/store/init.rs | 170 ++++++++++++++++++ 2 files changed, 175 insertions(+), 1 deletion(-) diff --git a/crates/ecstore/src/services/tier/tier_mutation_peer.rs b/crates/ecstore/src/services/tier/tier_mutation_peer.rs index c81c1ed5f..6379a4c0f 100644 --- a/crates/ecstore/src/services/tier/tier_mutation_peer.rs +++ b/crates/ecstore/src/services/tier/tier_mutation_peer.rs @@ -107,7 +107,9 @@ async fn handle_prepare( } let tier_config_mgr = api.tier_config_mgr(); - for _ in 0..3 { + // Peers can race to create the same durable intent. Re-read after bounded + // lock contention so the winning publication follows the normal identity checks. + for attempt in 0..3 { let (stored, applied) = match load_tier_mutation_intent_record(api.clone(), mutation_id).await { Ok(existing) => { if !existing.same_identity_as(&intent) { @@ -123,9 +125,11 @@ async fn handle_prepare( match save_tier_mutation_intent_record_if_absent(api.clone(), &intent).await { Ok(()) => (intent.clone(), true), Err(Error::PreconditionFailed) => continue, + Err(Error::Lock(rustfs_lock::LockError::Timeout { .. })) if attempt < 2 => continue, Err(err) => return Err(err.into()), } } + Err(Error::Lock(rustfs_lock::LockError::Timeout { .. })) if attempt < 2 => continue, Err(err) => return Err(err.into()), }; diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index e57d54e54..166a90138 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -20362,6 +20362,176 @@ mod tests { } } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn tier_mutation_peer_prepare_reloads_after_save_and_load_lock_timeouts() { + use crate::storage_api_contracts::namespace::NamespaceLocking as _; + use std::sync::atomic::Ordering; + + let temp_dir = tempfile::tempdir().expect("create contended tier prepare store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-prepare-contention", &[4])).await; + let mutation_id = uuid::Uuid::new_v4(); + let intent = tier_mutation_peer_test_intent(mutation_id, "COLD-A", [9; 32]); + let object = crate::services::tier::tier_mutation_intent::tier_mutation_intent_record_object_name(mutation_id) + .expect("intent record object should resolve"); + let barrier = crate::set_disk::PutObjectCommitBarrier::install( + RUSTFS_META_BUCKET, + &object, + crate::set_disk::PutObjectCommitPause::BeforeNamespace, + ); + let prepare_store = store.clone(); + let payload = intent.encode().expect("prepare intent should encode"); + let mut prepare = tokio::spawn(async move { + handle_tier_mutation_peer_request( + prepare_store, + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + &payload, + ) + .await + }); + barrier.wait_until_paused().await; + let namespace = store.pools[0].disk_set[0] + .new_ns_lock(RUSTFS_META_BUCKET, &object) + .await + .expect("intent record namespace should resolve"); + let owner = namespace + .get_write_lock(crate::set_disk::get_lock_acquire_timeout()) + .await + .expect("competing publisher should hold the intent namespace"); + // The competing publisher already owns the namespace. Bypass only its + // recursive acquisition and drain the durable tail before releasing it. + com::save_config_with_opts( + store.clone(), + &object, + intent.encode().expect("winning intent should encode"), + &ObjectOptions { + no_lock: true, + max_parity: true, + write_completion: crate::object_api::WriteCompletion::TailDrained, + ..Default::default() + }, + ) + .await + .expect("competing publisher should persist the matching intent under its lock"); + let manager = ctx + .lock_manager() + .as_fast_lock_manager() + .expect("fixture should use real locks"); + let key = rustfs_lock::fast_lock::types::ObjectKey { + bucket: Arc::from(RUSTFS_META_BUCKET), + object: Arc::from(object.as_str()), + version: None, + }; + let timeouts = &manager.get_shard(&key).metrics().timeouts; + let before = timeouts.load(Ordering::Relaxed); + barrier.release_and_wait_until_namespace_pending().await; + drop(barrier); + tokio::select! { + result = &mut prepare => panic!("Prepare must retry the contended save and load before completing: {result:?}"), + result = tokio::time::timeout(Duration::from_secs(30), async { + while timeouts.load(Ordering::Relaxed) - before < 2 { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) => result.expect("the save and subsequent load should each exhaust their lock wait"), + } + assert_eq!(timeouts.load(Ordering::Relaxed) - before, 2); + drop(owner); + let outcome = tokio::time::timeout(Duration::from_secs(30), prepare) + .await + .expect("third prepare round should complete after the competing publisher releases its lock") + .expect("prepare task should join") + .expect("third prepare round should reload the matching durable intent"); + assert_eq!(outcome.state, TierMutationPeerState::Prepared); + assert!(!outcome.applied, "the retry must reuse the winning record rather than recreate it"); + assert_eq!( + load_tier_mutation_intent_record(store.clone(), mutation_id) + .await + .expect("winning record should remain readable"), + intent + ); + assert_eq!(timeouts.load(Ordering::Relaxed) - before, 2); + shutdown.cancel(); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn tier_mutation_peer_prepare_preserves_lock_timeout_after_three_rounds() { + use crate::storage_api_contracts::namespace::NamespaceLocking as _; + use std::sync::atomic::Ordering; + + let temp_dir = tempfile::tempdir().expect("create exhausted tier prepare store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-prepare-exhaustion", &[4])).await; + let mutation_id = uuid::Uuid::new_v4(); + let intent = tier_mutation_peer_test_intent(mutation_id, "COLD-A", [9; 32]); + let object = crate::services::tier::tier_mutation_intent::tier_mutation_intent_record_object_name(mutation_id) + .expect("intent record object should resolve"); + let barrier = crate::set_disk::PutObjectCommitBarrier::install( + RUSTFS_META_BUCKET, + &object, + crate::set_disk::PutObjectCommitPause::BeforeNamespace, + ); + let prepare_store = store.clone(); + let payload = intent.encode().expect("prepare intent should encode"); + let prepare = tokio::spawn(async move { + handle_tier_mutation_peer_request( + prepare_store, + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + &payload, + ) + .await + }); + barrier.wait_until_paused().await; + let namespace = store.pools[0].disk_set[0] + .new_ns_lock(RUSTFS_META_BUCKET, &object) + .await + .expect("intent record namespace should resolve"); + let lock_timeout = crate::set_disk::get_lock_acquire_timeout(); + let owner = namespace + .get_write_lock(lock_timeout) + .await + .expect("competing publisher should hold the intent namespace"); + let manager = ctx + .lock_manager() + .as_fast_lock_manager() + .expect("fixture should use real locks"); + let key = rustfs_lock::fast_lock::types::ObjectKey { + bucket: Arc::from(RUSTFS_META_BUCKET), + object: Arc::from(object.as_str()), + version: None, + }; + let timeouts = &manager.get_shard(&key).metrics().timeouts; + let before = timeouts.load(Ordering::Relaxed); + barrier.release_and_wait_until_namespace_pending().await; + drop(barrier); + let error = tokio::time::timeout(Duration::from_secs(30), prepare) + .await + .expect("three contended prepare rounds should exhaust the bounded budget") + .expect("prepare task should join") + .expect_err("Prepare must fail while the competing publisher keeps the namespace"); + match error { + TierMutationPeerError::Store(Error::Lock(rustfs_lock::LockError::Timeout { resource, timeout })) => { + assert_eq!(resource, key.to_string()); + assert_eq!(timeout, lock_timeout); + } + other => panic!("Prepare must preserve the final typed lock timeout: {other:?}"), + } + assert_eq!(timeouts.load(Ordering::Relaxed) - before, 3, "save and load must share three rounds"); + drop(owner); + assert!(matches!( + load_tier_mutation_intent_record(store, mutation_id).await, + Err(Error::ConfigNotFound) + )); + shutdown.cancel(); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)]