fix(tier): reread mutation intents after lock contention (#8209)

This commit is contained in:
Chris
2026-09-29 07:16:01 +08:00
committed by GitHub
parent 87e9a84a4e
commit c9acf01fd0
2 changed files with 175 additions and 1 deletions
@@ -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()),
};
+170
View File
@@ -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)]