From 3cbe3d6b9486023fe15c9c24576b5c65c2608e7d Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 25 Jul 2026 23:20:54 +0800 Subject: [PATCH] fix(tier): reconcile committed mutation replay (#5230) Co-authored-by: heihutu --- .../src/inline_fast_path_cluster_test.rs | 45 +++++++ .../src/cluster/rpc/peer_rest_client.rs | 35 +++--- crates/ecstore/src/services/tier/tier.rs | 110 +++++++++++++++++- 3 files changed, 171 insertions(+), 19 deletions(-) diff --git a/crates/e2e_test/src/inline_fast_path_cluster_test.rs b/crates/e2e_test/src/inline_fast_path_cluster_test.rs index 93874d77e..30132a5e0 100644 --- a/crates/e2e_test/src/inline_fast_path_cluster_test.rs +++ b/crates/e2e_test/src/inline_fast_path_cluster_test.rs @@ -1030,6 +1030,32 @@ async fn wait_for_tier_verifiable(hot: &RustFSTestClusterEnvironment, tier_name: .into()) } +async fn wait_for_tier_converged(hot: &RustFSTestClusterEnvironment, tier_name: &str, add_tier_responses: &str) -> TestResult { + let deadline = Instant::now() + Duration::from_secs(60); + let final_error = loop { + let snapshot = tier_readiness_snapshot(hot, tier_name).await?; + if snapshot + .iter() + .all(|node| node.list_status.is_success() && node.list_has_tier && node.verify_status.is_success()) + { + return Ok(()); + } + if snapshot.iter().any(|node| { + !node.list_status.is_success() || (!node.verify_status.is_success() && !is_retryable_tier_error(&node.verify_body)) + }) { + break format_tier_readiness_snapshot(&snapshot); + } + if Instant::now() >= deadline { + break format_tier_readiness_snapshot(&snapshot); + } + sleep(Duration::from_millis(500)).await; + }; + Err(format!( + "tier {tier_name} did not converge on every hot node within 60s after AddTier({add_tier_responses}): {final_error}" + ) + .into()) +} + struct TierNodeReadiness { node_index: usize, node_url: String, @@ -1521,6 +1547,25 @@ async fn four_node_mixed_msgpack_compat_mode_preserves_fallback_controls() -> Te Ok(()) } +#[tokio::test] +#[serial] +async fn four_node_add_tier_committed_replay_converges() -> TestResult { + init_logging(); + + let mut cold = RustFSTestEnvironment::new().await?; + cold.access_key = "inlineconcurrentcoldadmin".to_string(); + cold.secret_key = "inlineconcurrentcoldsecret".to_string(); + cold.start_rustfs_server_without_cleanup(vec![]).await?; + cold.create_s3_client().create_bucket().bucket(TIER_BUCKET).send().await?; + + let mut hot = RustFSTestClusterEnvironment::new(4).await?; + hot.start().await?; + + let tier_name = unique_tier_name(); + add_rustfs_tier(&hot, &cold, &tier_name).await?; + wait_for_tier_converged(&hot, &tier_name, "committed AddTier replay").await +} + #[tokio::test] #[serial] async fn four_node_mixed_msgpack_compat_mode_preserves_fallback_controls_during_transition() -> TestResult { diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index 8617ac416..7fcc1ae47 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -1631,24 +1631,29 @@ impl PeerRestClient { } pub async fn load_transition_tier_config(&self) -> Result<()> { - self.finalize_result( - async { - let mut client = self.get_client().await?; - let request = Request::new(LoadTransitionTierConfigRequest {}); + let result = self.load_transition_tier_config_inner().await; + if let Err(err) = &result + && Self::is_network_like_error(err) + { + self.prepare_retry().await; + return self.finalize_result(self.load_transition_tier_config_inner().await).await; + } + self.finalize_result(result).await + } - let response = client.load_transition_tier_config(request).await?.into_inner(); - if !response.success { - if let Some(msg) = response.error_info { - return Err(Error::other(msg)); - } - return Err(Error::other("")); - } + async fn load_transition_tier_config_inner(&self) -> Result<()> { + let mut client = self.get_client().await?; + let request = Request::new(LoadTransitionTierConfigRequest {}); - Ok(()) + let response = client.load_transition_tier_config(request).await?.into_inner(); + if !response.success { + if let Some(msg) = response.error_info { + return Err(Error::other(msg)); } - .await, - ) - .await + return Err(Error::other("")); + } + + Ok(()) } } diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 15651efac..9ac5ed661 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -930,6 +930,7 @@ async fn prepare_tier_mutation_peers( let result = peer.prepare_tier_mutation(mutation_id, payload.clone()).await; match result { Ok(PeerTierMutationState::Prepared) => prepared.push(peer), + Ok(PeerTierMutationState::Committed) => {} Ok(state) => { return Err(TierMutationPrepareFailure { error: tier_mutation_fanout_admin_error( @@ -5980,6 +5981,41 @@ mod tests { ); } + #[tokio::test] + async fn prepare_tier_mutation_peers_accepts_replayed_committed_peer() { + let mutation_id = uuid::Uuid::from_u128(36); + let intent = prepared_remove_intent("COLD-A", mutation_id); + let calls = Arc::new(Mutex::new(Vec::new())); + let peers = vec![ + FakeTierMutationPeer::boxed_with_prepare_commit( + "peer-a", + calls.clone(), + Ok(PeerTierMutationState::Committed), + Ok(PeerTierMutationState::Committed), + ), + ConcurrencyTrackingTierMutationPeer::boxed( + "peer-b", + calls.clone(), + Arc::new(AtomicUsize::new(0)), + Arc::new(AtomicUsize::new(0)), + ), + ]; + + let prepared = match prepare_tier_mutation_peers(mutation_id, peers, &intent).await { + Ok(prepared) => prepared, + Err(_) => panic!("replayed committed peer should not fail the prepare fanout"), + }; + + assert_eq!(prepared.len(), 1); + let calls = lock_unpoisoned(&calls); + assert_eq!(calls.len(), 2); + assert!( + calls[0].starts_with("peer-a:prepare:"), + "replayed peer must be prepared before the remaining peer" + ); + assert_eq!(calls[1], "peer-b:prepare"); + } + #[tokio::test] async fn commit_tier_mutation_peers_serializes_shared_record_writes() { let mutation_id = uuid::Uuid::from_u128(35); @@ -6178,6 +6214,72 @@ mod tests { assert_eq!(observed.committed_config_etag.as_deref(), Some("etag-new")); } + #[tokio::test] + async fn intent_advance_reconciles_matching_commit_after_final_cas_race() { + let store = Arc::new(CasConfigStore::default()); + let mutation_id = uuid::Uuid::from_u128(37); + let prepared = prepared_remove_intent("COLD-A", mutation_id); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &prepared) + .await + .expect("prepared intent fixture should persist"); + + let committed = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + let object = crate::services::tier::tier_mutation_intent::tier_mutation_intent_record_object_name(mutation_id) + .expect("record object should build"); + store + .rewrite_on_next_if_match(object.clone(), prepared.encode().expect("prepared intent fixture should encode")) + .await; + store + .rewrite_on_next_if_match(object.clone(), prepared.encode().expect("prepared intent fixture should encode")) + .await; + store + .rewrite_on_next_if_match(object, committed.encode().expect("committed intent fixture should encode")) + .await; + + let (observed, applied) = crate::services::tier::tier_mutation_intent::advance_tier_mutation_intent_record_idempotent( + store, + mutation_id, + TierMutationIntentState::Committed, + Some("etag-new".to_string()), + ) + .await + .expect("matching commit after the final CAS race should be idempotent"); + + assert!(!applied); + assert_eq!(observed.state, TierMutationIntentState::Committed); + assert_eq!(observed.committed_config_etag.as_deref(), Some("etag-new")); + } + + #[tokio::test] + async fn intent_advance_rejects_unresolved_final_cas_race() { + let store = Arc::new(CasConfigStore::default()); + let mutation_id = uuid::Uuid::from_u128(38); + let prepared = prepared_remove_intent("COLD-A", mutation_id); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &prepared) + .await + .expect("prepared intent fixture should persist"); + + let object = crate::services::tier::tier_mutation_intent::tier_mutation_intent_record_object_name(mutation_id) + .expect("record object should build"); + const CAS_ATTEMPTS_UNDER_TEST: usize = 3; + for _ in 0..CAS_ATTEMPTS_UNDER_TEST { + store + .rewrite_on_next_if_match(object.clone(), prepared.encode().expect("prepared intent fixture should encode")) + .await; + } + + let err = crate::services::tier::tier_mutation_intent::advance_tier_mutation_intent_record_idempotent( + store, + mutation_id, + TierMutationIntentState::Committed, + Some("etag-new".to_string()), + ) + .await + .expect_err("a still-prepared intent after the final CAS race must fail closed"); + + assert!(matches!(err, Error::PreconditionFailed)); + } + #[tokio::test] async fn committed_mutation_recovery_requires_every_peer_to_commit() { let mutation_id = uuid::Uuid::from_u128(17); @@ -8270,7 +8372,7 @@ mod tests { struct CasConfigStore { objects: tokio::sync::Mutex, String)>>, legacy_state: tokio::sync::Mutex>>, - if_match_race_rewrite: tokio::sync::Mutex)>>, + if_match_race_rewrite: tokio::sync::Mutex)>>, next_etag: AtomicUsize, fail_put: AtomicBool, truncate_reference_page_without_marker: AtomicBool, @@ -8284,7 +8386,7 @@ mod tests { Self { objects: tokio::sync::Mutex::new(HashMap::new()), legacy_state: tokio::sync::Mutex::new(None), - if_match_race_rewrite: tokio::sync::Mutex::new(None), + if_match_race_rewrite: tokio::sync::Mutex::new(std::collections::VecDeque::new()), next_etag: AtomicUsize::new(0), fail_put: AtomicBool::new(false), truncate_reference_page_without_marker: AtomicBool::new(false), @@ -8311,7 +8413,7 @@ mod tests { } async fn rewrite_on_next_if_match(&self, object: String, data: Vec) { - *self.if_match_race_rewrite.lock().await = Some((object, data)); + self.if_match_race_rewrite.lock().await.push_back((object, data)); } fn omit_truncated_reference_marker(&self) { @@ -8380,7 +8482,7 @@ mod tests { .and_then(HTTPPreconditions::if_match_value) .is_some() { - self.if_match_race_rewrite.lock().await.take() + self.if_match_race_rewrite.lock().await.pop_front() } else { None };