From 3f60cc743e26031a2b794398e6c95d93d28dbc94 Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 22 Jul 2026 11:06:55 +0800 Subject: [PATCH] fix(tier): replay committed mutation intents (#5107) Keep recovered tier mutation blocks installed until the local publish transition has atomically established draining for the affected tiers, so old-generation leases cannot slip in after peer commit replay and before local publish. Delay committed intent cleanup until local publish succeeds. If peer replay succeeds but local publish fails, the durable committed intent remains available for retry and the runtime block stays fail-closed. Co-authored-by: heihutu --- crates/ecstore/src/services/tier/tier.rs | 999 ++++++++++++++++-- .../src/services/tier/tier_mutation_intent.rs | 67 +- .../src/services/tier/tier_mutation_peer.rs | 24 +- crates/ecstore/src/store/init.rs | 43 +- 4 files changed, 1045 insertions(+), 88 deletions(-) diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 5e8f8a3b0..9b13be9f8 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -20,7 +20,7 @@ use byteorder::{ByteOrder, LittleEndian}; use bytes::Bytes; -use futures::FutureExt; +use futures::{FutureExt, future::join_all}; use http::HeaderMap; use http::status::StatusCode; use lazy_static::lazy_static; @@ -65,6 +65,7 @@ use crate::storage_api_contracts::{ range::HTTPRangeSpec, }; use crate::{ + cluster::rpc::peer_rest_client::{PeerRestClient, PeerTierMutationState}, config::com::{CONFIG_PREFIX, read_config, read_config_with_metadata}, disk::{MIGRATING_META_BUCKET, RUSTFS_META_BUCKET}, object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader}, @@ -81,7 +82,7 @@ use super::{ tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_PERM_ERR}, tier_mutation_intent::{ TierMutationIntent, TierMutationIntentKind, TierMutationIntentState, TierMutationIntentTarget, - list_tier_mutation_intent_records, + delete_tier_mutation_intent_record, list_tier_mutation_intent_records, }, warm_backend::WarmBackendImpl, }; @@ -384,6 +385,19 @@ fn runtime_blocks_tier(runtime: &TierDriverRuntime, tier_name: &str) -> bool { runtime.draining.contains_key(tier_name) || runtime.prepared_mutation_blocks.contains_key(tier_name) } +fn runtime_blocks_tier_transition( + runtime: &TierDriverRuntime, + tier_name: &str, + allowed_mutation_blocks: Option<&HashSet>, +) -> bool { + if runtime.draining.contains_key(tier_name) { + return true; + } + runtime.prepared_mutation_blocks.get(tier_name).is_some_and(|mutation_id| { + !allowed_mutation_blocks.is_some_and(|allowed_mutation_blocks| allowed_mutation_blocks.contains(mutation_id)) + }) +} + fn runtime_has_mutation_block(runtime: &TierDriverRuntime) -> bool { !runtime.draining.is_empty() || !runtime.prepared_mutation_blocks.is_empty() } @@ -704,6 +718,89 @@ fn tier_reference_proof_admin_error(err: impl std::fmt::Display) -> AdminError { admin_err } +fn tier_mutation_replay_error(err: impl std::fmt::Display) -> io::Error { + io::Error::other(format!("Remote tier mutation committed replay failed: {err}")) +} + +fn ensure_complete_tier_mutation_commit_peer_set(peer_count: usize, remote_host_count: usize) -> io::Result<()> { + if peer_count != remote_host_count { + return Err(tier_mutation_replay_error( + "cluster endpoint topology has remote hosts without peer commit clients", + )); + } + Ok(()) +} + +#[async_trait::async_trait] +trait TierMutationCommitPeer: Send + Sync { + fn peer_label(&self) -> String; + async fn commit_tier_mutation(&self, mutation_id: uuid::Uuid, canonical_payload: Bytes) -> Result; +} + +#[async_trait::async_trait] +impl TierMutationCommitPeer for PeerRestClient { + fn peer_label(&self) -> String { + self.grid_host.clone() + } + + async fn commit_tier_mutation(&self, mutation_id: uuid::Uuid, canonical_payload: Bytes) -> Result { + Ok(PeerRestClient::commit_tier_mutation(self, mutation_id, canonical_payload) + .await? + .state) + } +} + +#[cfg(test)] +tokio::task_local! { + static TIER_MUTATION_TEST_COMMIT_PEERS: Vec>; +} + +async fn remote_tier_mutation_commit_peers() -> io::Result>> { + #[cfg(test)] + if let Ok(peers) = TIER_MUTATION_TEST_COMMIT_PEERS.try_with(|peers| peers.clone()) { + return Ok(peers); + } + let Some(endpoints) = runtime_sources::endpoint_pools() else { + return Err(tier_mutation_replay_error("cluster endpoint topology is not initialized")); + }; + let remote_host_count = endpoints.hosts_sorted().iter().flatten().count(); + let (peers, _) = PeerRestClient::new_clients(endpoints).await; + let peers = peers + .into_iter() + .flatten() + .map(|peer| Arc::new(peer) as Arc) + .collect::>(); + ensure_complete_tier_mutation_commit_peer_set(peers.len(), remote_host_count)?; + Ok(peers) +} + +async fn commit_tier_mutation_peers( + mutation_id: uuid::Uuid, + peers: Vec>, + committed_config_etag: &str, +) -> io::Result<()> { + let payload = Bytes::copy_from_slice(committed_config_etag.as_bytes()); + let results = join_all(peers.into_iter().map(|peer| { + let payload = payload.clone(); + async move { + let label = peer.peer_label(); + let result = peer.commit_tier_mutation(mutation_id, payload).await; + (label, result) + } + })) + .await; + for (label, result) in results { + match result { + Ok(PeerTierMutationState::Committed) => {} + Ok(state) => { + return Err(tier_mutation_replay_error(format!("peer {label} returned unexpected state {state:?}"))); + } + Err(err) => return Err(tier_mutation_replay_error(format!("peer {label}: {err}"))), + } + } + Ok(()) +} + async fn apply_tier_candidate_mutation( mutation: TierCandidateMutation, candidate: &mut TierConfigMgr, @@ -2290,10 +2387,19 @@ impl TierConfigMgr { handle: &Arc>, manager: &mut Self, candidate: &Self, + ) -> std::result::Result { + Self::begin_publish_transition_with_allowed_mutation_blocks(handle, manager, candidate, None) + } + + fn begin_publish_transition_with_allowed_mutation_blocks( + handle: &Arc>, + manager: &mut Self, + candidate: &Self, + allowed_mutation_blocks: Option<&HashSet>, ) -> std::result::Result { let changed = changed_tier_names(manager, candidate); let replaced_destinations = replaced_tier_destinations(manager, candidate)?; - Self::begin_tier_transition_with_destinations(handle, manager, changed, replaced_destinations) + Self::begin_tier_transition_with_destinations(handle, manager, changed, replaced_destinations, allowed_mutation_blocks) } fn begin_tier_transition( @@ -2301,7 +2407,7 @@ impl TierConfigMgr { manager: &mut Self, changed: HashSet, ) -> std::result::Result { - Self::begin_tier_transition_with_destinations(handle, manager, changed, HashMap::new()) + Self::begin_tier_transition_with_destinations(handle, manager, changed, HashMap::new(), None) } fn begin_tier_transition_with_destinations( @@ -2309,6 +2415,7 @@ impl TierConfigMgr { manager: &mut Self, changed: HashSet, replaced_destinations: HashMap, + allowed_mutation_blocks: Option<&HashSet>, ) -> std::result::Result { let runtime = tier_driver_runtime(handle, manager); for tier_name in replaced_destinations.keys() { @@ -2318,13 +2425,14 @@ impl TierConfigMgr { manager.replace_driver(tier_name, driver)?; } } - Self::begin_tier_transition_in_runtime(runtime, changed, replaced_destinations) + Self::begin_tier_transition_in_runtime(runtime, changed, replaced_destinations, allowed_mutation_blocks) } fn begin_tier_transition_in_runtime( runtime: Arc>, changed: HashSet, replaced_destinations: HashMap, + allowed_mutation_blocks: Option<&HashSet>, ) -> std::result::Result { let count = u64::try_from(changed.len()).map_err(|_| { let mut err = ERR_TIER_INVALID_CONFIG.clone(); @@ -2332,7 +2440,10 @@ impl TierConfigMgr { err })?; let mut runtime_guard = lock_unpoisoned(&runtime); - if changed.iter().any(|tier_name| runtime_blocks_tier(&runtime_guard, tier_name)) { + if changed + .iter() + .any(|tier_name| runtime_blocks_tier_transition(&runtime_guard, tier_name, allowed_mutation_blocks)) + { let mut err = ERR_TIER_INVALID_CONFIG.clone(); err.message = "Remote tier configuration is already being replaced".to_string(); return Err(err); @@ -2369,10 +2480,24 @@ impl TierConfigMgr { handle: &Arc>, candidate: Self, driver_tier: Option<&str>, + ) -> std::result::Result<(), AdminError> { + Self::publish_candidate_inner_with_allowed_mutation_blocks(handle, candidate, driver_tier, None).await + } + + async fn publish_candidate_inner_with_allowed_mutation_blocks( + handle: &Arc>, + candidate: Self, + driver_tier: Option<&str>, + allowed_mutation_blocks: Option<&HashSet>, ) -> std::result::Result<(), AdminError> { let transition = { let mut manager = handle.write().await; - Self::begin_publish_transition(handle, &mut manager, &candidate)? + Self::begin_publish_transition_with_allowed_mutation_blocks( + handle, + &mut manager, + &candidate, + allowed_mutation_blocks, + )? }; transition.wait_for_active_leases().await?; @@ -2470,12 +2595,29 @@ impl TierConfigMgr { candidate: Self, driver_tier: Option, update: tokio::sync::OwnedMutexGuard<()>, + ) -> std::result::Result<(), AdminError> { + Self::publish_candidate_owned_with_allowed_mutation_blocks(handle, candidate, driver_tier, update, HashSet::new()).await + } + + async fn publish_candidate_owned_with_allowed_mutation_blocks( + handle: &Arc>, + candidate: Self, + driver_tier: Option, + update: tokio::sync::OwnedMutexGuard<()>, + allowed_mutation_blocks: HashSet, ) -> std::result::Result<(), AdminError> { let handle = handle.clone(); tokio::spawn(async move { match AssertUnwindSafe(async move { let _update = update; - Self::publish_candidate_inner(&handle, candidate, driver_tier.as_deref()).await + let allowed_mutation_blocks = (!allowed_mutation_blocks.is_empty()).then_some(&allowed_mutation_blocks); + Self::publish_candidate_inner_with_allowed_mutation_blocks( + &handle, + candidate, + driver_tier.as_deref(), + allowed_mutation_blocks, + ) + .await }) .catch_unwind() .await @@ -2580,20 +2722,60 @@ impl TierConfigMgr { } pub async fn reload_handle(handle: &Arc>, api: Arc) -> io::Result<()> { - let prepared_intents = Self::load_prepared_mutation_intents(api.clone()).await?; - let update = Self::admin_update_lock(handle).await; - let candidate = load_tier_config(api).await?; - Self::publish_candidate_owned(handle, candidate, None, update) - .await - .map_err(io::Error::other)?; - Self::reconcile_prepared_mutation_intents(handle, &prepared_intents) - .await - .map_err(io::Error::other) + Self::reload_handle_with(handle, api).await } - async fn load_prepared_mutation_intents(api: Arc) -> io::Result> { + async fn reload_handle_with(handle: &Arc>, api: Arc) -> io::Result<()> + where + S: EcstoreObjectIO + + EcstoreObjectOperations + + NamespaceLocking + + ListOperations>, + { + let mut mutation_intents = Self::load_tier_mutation_intents(api.clone()).await?; + // Lock order matches update_candidate_with_config_lock: namespace tier-config lock before admin_updates. + let config_lock = if mutation_intents + .iter() + .any(|intent| intent.state == TierMutationIntentState::Committed) + { + let guard = Self::acquire_tier_config_write_lock(api.clone()) + .await + .map_err(|err| match err { + TierConfigUpdateError::Load(err) => err, + other => io::Error::other(format!("{other:?}")), + })?; + mutation_intents = Self::load_tier_mutation_intents(api.clone()).await?; + Some(guard) + } else { + None + }; + Self::reconcile_committed_mutation_intent_blocks(handle, &mutation_intents).await?; + Self::replay_committed_mutation_intents(&mutation_intents).await?; + let allowed_mutation_blocks = mutation_intents + .iter() + .filter(|intent| matches!(intent.state, TierMutationIntentState::Prepared | TierMutationIntentState::Committed)) + .map(|intent| intent.mutation_id) + .collect::>(); + let update = Self::admin_update_lock(handle).await; + let candidate = load_tier_config(api.clone()).await?; + Self::publish_candidate_owned_with_allowed_mutation_blocks(handle, candidate, None, update, allowed_mutation_blocks) + .await + .map_err(io::Error::other)?; + Self::delete_replayed_committed_mutation_intents(api, &mutation_intents).await?; + mutation_intents.retain(|intent| intent.state == TierMutationIntentState::Prepared); + Self::reconcile_prepared_mutation_intents(handle, &mutation_intents) + .await + .map_err(io::Error::other)?; + drop(config_lock); + Ok(()) + } + + async fn load_tier_mutation_intents(api: Arc) -> io::Result> + where + S: EcstoreObjectIO + ListOperations>, + { let mut marker = None; - let mut prepared = Vec::new(); + let mut intents = Vec::new(); loop { let scan = list_tier_mutation_intent_records(api.clone(), TIER_MUTATION_INTENT_RECOVERY_SCAN_LIMIT, marker) .await @@ -2604,13 +2786,9 @@ impl TierConfigMgr { scan.failed ))); } - prepared.extend( - scan.intents - .into_iter() - .filter(|intent| intent.state == TierMutationIntentState::Prepared), - ); + intents.extend(scan.intents); if !scan.truncated { - return Ok(prepared); + return Ok(intents); } let next_marker = scan.next_marker.ok_or_else(|| { io::Error::other("tier mutation intent recovery scan was truncated without a continuation marker") @@ -2619,13 +2797,49 @@ impl TierConfigMgr { } } + async fn replay_committed_mutation_intents(intents: &[TierMutationIntent]) -> io::Result<()> { + if !intents + .iter() + .any(|intent| intent.state == TierMutationIntentState::Committed) + { + return Ok(()); + } + let peers = remote_tier_mutation_commit_peers().await?; + for intent in intents + .iter() + .filter(|intent| intent.state == TierMutationIntentState::Committed) + { + let committed_config_etag = intent + .committed_config_etag + .as_deref() + .ok_or_else(|| io::Error::other("committed tier mutation intent is missing committed config etag"))?; + commit_tier_mutation_peers(intent.mutation_id, peers.clone(), committed_config_etag).await?; + } + Ok(()) + } + + async fn delete_replayed_committed_mutation_intents(api: Arc, intents: &[TierMutationIntent]) -> io::Result<()> + where + S: EcstoreObjectOperations, + { + for intent in intents + .iter() + .filter(|intent| intent.state == TierMutationIntentState::Committed) + { + delete_tier_mutation_intent_record(api.clone(), intent.mutation_id) + .await + .map_err(tier_mutation_replay_error)?; + } + Ok(()) + } + async fn reconcile_prepared_mutation_intents( handle: &Arc>, intents: &[TierMutationIntent], ) -> std::result::Result<(), AdminError> { let mut prepared_mutation_blocks = HashMap::new(); for intent in intents { - Self::collect_prepared_mutation_intent_block(&mut prepared_mutation_blocks, intent)?; + Self::collect_mutation_intent_block(&mut prepared_mutation_blocks, intent, TierMutationIntentState::Prepared)?; } let manager = handle.read().await; let runtime = tier_driver_runtime(handle, &manager); @@ -2634,6 +2848,22 @@ impl TierConfigMgr { Ok(()) } + async fn reconcile_committed_mutation_intent_blocks( + handle: &Arc>, + intents: &[TierMutationIntent], + ) -> io::Result<()> { + let mut mutation_blocks = HashMap::new(); + for intent in intents { + Self::collect_mutation_intent_block(&mut mutation_blocks, intent, TierMutationIntentState::Committed) + .map_err(io::Error::other)?; + } + let manager = handle.read().await; + let runtime = tier_driver_runtime(handle, &manager); + let mut runtime = lock_unpoisoned(&runtime); + runtime.prepared_mutation_blocks = mutation_blocks; + Ok(()) + } + pub(crate) async fn apply_prepared_mutation_intent_block( handle: &Arc>, intent: &TierMutationIntent, @@ -2642,7 +2872,7 @@ impl TierConfigMgr { let runtime = tier_driver_runtime(handle, &manager); let mut runtime = lock_unpoisoned(&runtime); let mut prepared_mutation_blocks = runtime.prepared_mutation_blocks.clone(); - Self::collect_prepared_mutation_intent_block(&mut prepared_mutation_blocks, intent)?; + Self::collect_mutation_intent_block(&mut prepared_mutation_blocks, intent, TierMutationIntentState::Prepared)?; runtime.prepared_mutation_blocks = prepared_mutation_blocks; Ok(()) } @@ -2657,11 +2887,12 @@ impl TierConfigMgr { .retain(|_, blocked_mutation_id| *blocked_mutation_id != mutation_id); } - fn collect_prepared_mutation_intent_block( + fn collect_mutation_intent_block( prepared_mutation_blocks: &mut HashMap, intent: &TierMutationIntent, + state: TierMutationIntentState, ) -> std::result::Result<(), AdminError> { - if intent.state != TierMutationIntentState::Prepared { + if intent.state != state { return Ok(()); } for target in &intent.affected_targets { @@ -2926,8 +3157,8 @@ impl TierConfigMgr { self.replace_driver(tier_name, driver).map_err(io::Error::other)?; } } - let mut transition = - Self::begin_tier_transition_in_runtime(runtime, changed, replaced_destinations).map_err(io::Error::other)?; + let mut transition = Self::begin_tier_transition_in_runtime(runtime, changed, replaced_destinations, None) + .map_err(io::Error::other)?; transition.wait_for_active_leases().await.map_err(io::Error::other)?; transition .ensure_replaced_destinations_are_empty() @@ -3209,7 +3440,10 @@ where } #[tracing::instrument(level = "debug", name = "load_tier_config", skip(api))] -async fn load_tier_config(api: Arc) -> std::result::Result { +async fn load_tier_config(api: Arc) -> std::result::Result +where + S: EcstoreObjectIO, +{ let config_file = tier_config_path(TIER_CONFIG_FILE); match read_config(api.clone(), config_file.as_str()).await { Ok(data) => decode_tiering_config_blob(&data), @@ -3257,15 +3491,7 @@ async fn load_tier_config(api: Arc) -> std::result::Result(api: Arc) -> std::result::Result<(TierConfigMgr, Option), std::io::Error> where - S: ObjectIO< - Error = Error, - RangeSpec = HTTPRangeSpec, - HeaderMap = HeaderMap, - ObjectOptions = ObjectOptions, - ObjectInfo = ObjectInfo, - GetObjectReader = GetObjectReader, - PutObjectReader = PutObjReader, - >, + S: EcstoreObjectIO, { let config_file = tier_config_path(TIER_CONFIG_FILE); match read_config_with_metadata(api.clone(), &config_file, &ObjectOptions::default()).await { @@ -3295,6 +3521,14 @@ where } } +pub(crate) async fn tier_config_etag_matches(api: Arc, expected: &str) -> io::Result +where + S: EcstoreObjectIO, +{ + let (_, etag) = load_tier_config_for_update(api).await?; + Ok(etag.as_deref() == Some(expected)) +} + async fn read_tier_config_from_bucket( api: Arc, bucket: &str, @@ -4983,6 +5217,497 @@ mod tests { } } + fn committed_remove_intent(tier_name: &str, mutation_id: uuid::Uuid, etag: &str) -> TierMutationIntent { + let mut intent = prepared_remove_intent(tier_name, mutation_id); + intent + .advance(TierMutationIntentState::Committed, Some(etag.to_string())) + .expect("fixture intent should commit"); + intent + } + + struct FakeTierMutationCommitPeer { + label: &'static str, + calls: Arc>>, + commit: std::result::Result, + } + + impl FakeTierMutationCommitPeer { + fn boxed( + label: &'static str, + calls: Arc>>, + commit: std::result::Result, + ) -> Arc { + Arc::new(Self { label, calls, commit }) + } + + fn record(&self, call: String) { + lock_unpoisoned(&self.calls).push(call); + } + } + + #[async_trait::async_trait] + impl TierMutationCommitPeer for FakeTierMutationCommitPeer { + fn peer_label(&self) -> String { + self.label.to_string() + } + + async fn commit_tier_mutation(&self, mutation_id: uuid::Uuid, canonical_payload: Bytes) -> Result { + self.record(format!( + "{}:commit:{}:{}", + self.label, + mutation_id, + String::from_utf8_lossy(canonical_payload.as_ref()) + )); + self.commit.map_err(Error::other) + } + } + + #[tokio::test] + async fn committed_mutation_recovery_replays_peer_commit() { + let mutation_id = uuid::Uuid::from_u128(13); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + let calls = Arc::new(Mutex::new(Vec::new())); + + TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Ok(PeerTierMutationState::Committed), + )], + async { + TierConfigMgr::replay_committed_mutation_intents(&[intent]) + .await + .expect("committed recovery should replay peer commit"); + }, + ) + .await; + + assert_eq!(lock_unpoisoned(&calls).as_slice(), &[format!("peer-a:commit:{mutation_id}:etag-new")]); + } + + #[tokio::test] + async fn committed_mutation_recovery_fails_closed_on_peer_commit_error() { + let mutation_id = uuid::Uuid::from_u128(14); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + let calls = Arc::new(Mutex::new(Vec::new())); + + let err = TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Err("network down"), + )], + async { TierConfigMgr::replay_committed_mutation_intents(&[intent]).await }, + ) + .await + .expect_err("committed recovery must fail closed when peer commit cannot replay"); + + assert!(err.to_string().contains("network down"), "{err}"); + assert_eq!(lock_unpoisoned(&calls).as_slice(), &[format!("peer-a:commit:{mutation_id}:etag-new")]); + } + + #[tokio::test] + async fn committed_mutation_recovery_requires_every_peer_to_commit() { + let mutation_id = uuid::Uuid::from_u128(17); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + let calls = Arc::new(Mutex::new(Vec::new())); + + let err = TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![ + FakeTierMutationCommitPeer::boxed("peer-a", calls.clone(), Ok(PeerTierMutationState::Committed)), + FakeTierMutationCommitPeer::boxed("peer-b", calls.clone(), Ok(PeerTierMutationState::Prepared)), + ], + async { TierConfigMgr::replay_committed_mutation_intents(&[intent]).await }, + ) + .await + .expect_err("committed recovery must fail unless every peer reports committed"); + + assert!(err.to_string().contains("unexpected state"), "{err}"); + assert_eq!( + lock_unpoisoned(&calls).as_slice(), + &[ + format!("peer-a:commit:{mutation_id}:etag-new"), + format!("peer-b:commit:{mutation_id}:etag-new") + ] + ); + } + + #[test] + fn committed_replay_rejects_partial_peer_discovery() { + ensure_complete_tier_mutation_commit_peer_set(0, 0).expect("local-only topology has no remote peers"); + ensure_complete_tier_mutation_commit_peer_set(2, 2).expect("two remote hosts should require two peers"); + let err = ensure_complete_tier_mutation_commit_peer_set(1, 2).expect_err("missing one expected peer must fail closed"); + assert!(err.to_string().contains("without peer commit clients"), "{err}"); + } + + #[tokio::test] + async fn reload_handle_replays_committed_mutation_and_removes_intent_record() { + let store = Arc::new(CasConfigStore::default()); + let mut persisted = empty_mgr(); + persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + persisted + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("tier config fixture should persist"); + + let mutation_id = uuid::Uuid::from_u128(15); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &intent) + .await + .expect("committed intent fixture should persist"); + let calls = Arc::new(Mutex::new(Vec::new())); + let handle = TierConfigMgr::new(); + + TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Ok(PeerTierMutationState::Committed), + )], + async { + TierConfigMgr::reload_handle_with(&handle, store.clone()) + .await + .expect("reload should replay committed intent before publishing"); + }, + ) + .await; + + assert_eq!(lock_unpoisoned(&calls).as_slice(), &[format!("peer-a:commit:{mutation_id}:etag-new")]); + assert!(handle.read().await.tiers.contains_key("COLD-A")); + let reloaded = TierConfigMgr::load_tier_mutation_intents(store) + .await + .expect("intent scan should succeed after cleanup"); + assert!( + reloaded.iter().all(|intent| intent.mutation_id != mutation_id), + "committed intent should be removed after successful replay" + ); + } + + #[tokio::test] + async fn reload_handle_replays_committed_and_restores_prepared_blocks() { + let store = Arc::new(CasConfigStore::default()); + let mut persisted = empty_mgr(); + persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + persisted.tiers.insert("COLD-B".to_string(), build_rustfs_tier("COLD-B")); + persisted + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("tier config fixture should persist"); + + let committed_id = uuid::Uuid::from_u128(18); + let committed = committed_remove_intent("COLD-A", committed_id, "etag-new"); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &committed) + .await + .expect("committed intent fixture should persist"); + let prepared_id = uuid::Uuid::from_u128(19); + let prepared = prepared_remove_intent("COLD-B", prepared_id); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &prepared) + .await + .expect("prepared intent fixture should persist"); + let calls = Arc::new(Mutex::new(Vec::new())); + let handle = TierConfigMgr::new(); + + TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Ok(PeerTierMutationState::Committed), + )], + async { + TierConfigMgr::reload_handle_with(&handle, store.clone()) + .await + .expect("reload should replay committed intent and restore prepared blocks"); + }, + ) + .await; + + assert_eq!(lock_unpoisoned(&calls).as_slice(), &[format!("peer-a:commit:{committed_id}:etag-new")]); + let err = match TierConfigMgr::acquire_operation_lease(&handle, "COLD-B").await { + Ok(_) => panic!("prepared intent must still block operation leases after reload"), + Err(err) => err, + }; + assert!(err.message.contains("being replaced"), "{err}"); + let reloaded = TierConfigMgr::load_tier_mutation_intents(store) + .await + .expect("intent scan should succeed after mixed reload"); + assert!(reloaded.iter().all(|intent| intent.mutation_id != committed_id)); + assert!(reloaded.iter().any(|intent| intent.mutation_id == prepared_id)); + } + + #[tokio::test] + async fn reload_handle_keeps_committed_block_until_publish_swaps_generation() { + let store = Arc::new(CasConfigStore::default()); + empty_mgr() + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("empty tier config fixture should persist"); + + let mutation_id = uuid::Uuid::from_u128(21); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &intent) + .await + .expect("committed intent fixture should persist"); + let calls = Arc::new(Mutex::new(Vec::new())); + let handle = TierConfigMgr::new(); + { + let mut guard = handle.write().await; + install_lease_backend(&mut guard, "COLD-A", LeaseTestBackend::ready("old")); + } + let old = TierConfigMgr::acquire_operation_lease(&handle, "COLD-A") + .await + .expect("old generation lease should be available before recovery"); + + TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Ok(PeerTierMutationState::Committed), + )], + async { + let reload_handle = handle.clone(); + let reload_store = store.clone(); + let reload = async move { TierConfigMgr::reload_handle_with(&reload_handle, reload_store).await }; + let probe_handle = handle.clone(); + let probe = async move { + tokio::time::timeout(Duration::from_secs(1), async { + while old.inner.accepting.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + }) + .await + .expect("reload publish should revoke the old generation before waiting"); + { + let guard = probe_handle.read().await; + let runtime = registered_tier_driver_runtime(&guard).expect("runtime should remain registered"); + assert_eq!( + lock_unpoisoned(&runtime).prepared_mutation_blocks.get("COLD-A"), + Some(&mutation_id), + "committed replay block must remain until the local publish completes" + ); + } + let blocked = match TierConfigMgr::acquire_operation_lease(&probe_handle, "COLD-A").await { + Ok(_) => panic!("committed replay must block new old-generation leases while publish waits"), + Err(err) => err, + }; + assert!(blocked.message.contains("being replaced"), "{blocked}"); + drop(old); + }; + let (reload_result, ()) = tokio::join!(reload, probe); + reload_result.expect("reload should complete after the old lease drains"); + }, + ) + .await; + + assert_eq!(lock_unpoisoned(&calls).as_slice(), &[format!("peer-a:commit:{mutation_id}:etag-new")]); + assert!( + !handle.read().await.tiers.contains_key("COLD-A"), + "reloaded manager should publish the committed removal after the old lease drains" + ); + let reloaded = TierConfigMgr::load_tier_mutation_intents(store) + .await + .expect("intent scan should succeed after committed removal reload"); + assert!(reloaded.iter().all(|intent| intent.mutation_id != mutation_id)); + } + + #[tokio::test(start_paused = true)] + async fn reload_handle_keeps_committed_intent_when_local_publish_fails() { + let store = Arc::new(CasConfigStore::default()); + empty_mgr() + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("empty tier config fixture should persist"); + + let mutation_id = uuid::Uuid::from_u128(22); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &intent) + .await + .expect("committed intent fixture should persist"); + let calls = Arc::new(Mutex::new(Vec::new())); + let handle = TierConfigMgr::new(); + { + let mut guard = handle.write().await; + install_lease_backend(&mut guard, "COLD-A", LeaseTestBackend::ready("old")); + } + let old = TierConfigMgr::acquire_operation_lease(&handle, "COLD-A") + .await + .expect("old generation lease should be available before recovery"); + + TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Ok(PeerTierMutationState::Committed), + )], + async { + let reload_handle = handle.clone(); + let reload_store = store.clone(); + let reload = async move { TierConfigMgr::reload_handle_with(&reload_handle, reload_store).await }; + let probe = async { + while old.inner.accepting.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + tokio::time::advance(TIER_OPERATION_DRAIN_TIMEOUT).await; + }; + let (reload_result, ()) = tokio::join!(reload, probe); + let err = reload_result.expect_err("reload must fail while the old lease refuses to drain"); + assert!( + err.to_string() + .contains("Timed out waiting for active remote tier operations"), + "{err}" + ); + }, + ) + .await; + + assert_eq!(lock_unpoisoned(&calls).as_slice(), &[format!("peer-a:commit:{mutation_id}:etag-new")]); + let reloaded = TierConfigMgr::load_tier_mutation_intents(store) + .await + .expect("intent scan should succeed after failed local publish"); + assert!( + reloaded.iter().any(|intent| intent.mutation_id == mutation_id), + "failed local publish must keep the committed intent for retry" + ); + { + let guard = handle.read().await; + let runtime = registered_tier_driver_runtime(&guard).expect("runtime should remain registered"); + assert_eq!(lock_unpoisoned(&runtime).prepared_mutation_blocks.get("COLD-A"), Some(&mutation_id)); + } + let blocked = match TierConfigMgr::acquire_operation_lease(&handle, "COLD-A").await { + Ok(_) => panic!("failed local publish must keep blocking old-generation leases"), + Err(err) => err, + }; + assert!(blocked.message.contains("being replaced"), "{blocked}"); + drop(old); + } + + #[tokio::test] + async fn committed_replay_claim_prevents_duplicate_reload_fanout() { + let store = Arc::new(CasConfigStore::default()); + let mut persisted = empty_mgr(); + persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + persisted + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("tier config fixture should persist"); + + let mutation_id = uuid::Uuid::from_u128(20); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &intent) + .await + .expect("committed intent fixture should persist"); + let calls = Arc::new(Mutex::new(Vec::new())); + let left = TierConfigMgr::new(); + let right = TierConfigMgr::new(); + + TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Ok(PeerTierMutationState::Committed), + )], + async { + let (left_result, right_result) = tokio::join!( + TierConfigMgr::reload_handle_with(&left, store.clone()), + TierConfigMgr::reload_handle_with(&right, store.clone()) + ); + left_result.expect("first reload should complete"); + right_result.expect("second reload should observe the cleaned committed intent"); + }, + ) + .await; + + assert_eq!( + lock_unpoisoned(&calls).as_slice(), + &[format!("peer-a:commit:{mutation_id}:etag-new")], + "only the claimed reload should fan out committed replay" + ); + } + + #[tokio::test] + async fn reload_handle_fails_closed_when_committed_replay_fails() { + let store = Arc::new(CasConfigStore::default()); + let mut persisted = empty_mgr(); + persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + persisted + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("tier config fixture should persist"); + + let mutation_id = uuid::Uuid::from_u128(16); + let intent = committed_remove_intent("COLD-A", mutation_id, "etag-new"); + crate::services::tier::tier_mutation_intent::save_tier_mutation_intent_record(store.clone(), &intent) + .await + .expect("committed intent fixture should persist"); + let calls = Arc::new(Mutex::new(Vec::new())); + let handle = TierConfigMgr::new(); + + let err = TIER_MUTATION_TEST_COMMIT_PEERS + .scope( + vec![FakeTierMutationCommitPeer::boxed( + "peer-a", + calls.clone(), + Err("network down"), + )], + async { TierConfigMgr::reload_handle_with(&handle, store.clone()).await }, + ) + .await + .expect_err("reload must fail closed when committed replay fails"); + + assert!(err.to_string().contains("network down"), "{err}"); + assert_eq!(lock_unpoisoned(&calls).as_slice(), &[format!("peer-a:commit:{mutation_id}:etag-new")]); + assert!( + !handle.read().await.tiers.contains_key("COLD-A"), + "local manager must not publish persisted config before peer replay succeeds" + ); + let blocked = match TierConfigMgr::acquire_operation_lease(&handle, "COLD-A").await { + Ok(_) => panic!("failed committed replay must block old tier operations"), + Err(err) => err, + }; + assert!(blocked.message.contains("being replaced"), "{blocked}"); + let reloaded = TierConfigMgr::load_tier_mutation_intents(store) + .await + .expect("intent scan should succeed after failed replay"); + assert!( + reloaded.iter().any(|intent| intent.mutation_id == mutation_id), + "failed committed replay must keep the durable intent for retry" + ); + } + + #[tokio::test] + async fn committed_missing_intent_requires_matching_tier_config_etag() { + let store = Arc::new(CasConfigStore::default()); + let mut persisted = empty_mgr(); + persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + persisted + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("tier config fixture should persist"); + let (_, committed_etag) = load_tier_config_for_update(store.clone()) + .await + .expect("persisted tier config should reload"); + let committed_etag = committed_etag.expect("persisted tier config should carry an ETag"); + + assert!( + tier_config_etag_matches(store.clone(), &committed_etag) + .await + .expect("matching ETag proof should load") + ); + assert!( + !tier_config_etag_matches(store, "other-etag") + .await + .expect("mismatched ETag proof should load"), + "missing peer intent must not be terminal unless the committed config ETag matches" + ); + } + #[tokio::test] async fn prepared_mutation_recovery_blocks_new_operation_lease() { let manager = TierConfigMgr::new(); @@ -6065,7 +6790,7 @@ mod tests { #[derive(Debug)] struct CasConfigStore { - state: tokio::sync::Mutex, String)>>, + objects: tokio::sync::Mutex, String)>>, legacy_state: tokio::sync::Mutex>>, next_etag: AtomicUsize, fail_put: AtomicBool, @@ -6078,7 +6803,7 @@ mod tests { impl Default for CasConfigStore { fn default() -> Self { Self { - state: tokio::sync::Mutex::new(None), + objects: tokio::sync::Mutex::new(HashMap::new()), legacy_state: tokio::sync::Mutex::new(None), next_etag: AtomicUsize::new(0), fail_put: AtomicBool::new(false), @@ -6121,18 +6846,15 @@ mod tests { _headers: Self::HeaderMap, _opts: &Self::ObjectOptions, ) -> Result { - let current_config_path = tier_config_path(TIER_CONFIG_FILE); let legacy_config_path = tier_config_path(TIER_CONFIG_LEGACY_FILE); - let (data, etag) = if object == current_config_path { - let state = self.state.lock().await; - let (data, etag) = state.as_ref().ok_or(Error::ConfigNotFound)?; - (data.clone(), etag.clone()) - } else if object == legacy_config_path { + let (data, etag) = if object == legacy_config_path { let state = self.legacy_state.lock().await; let data = state.as_ref().ok_or(Error::ConfigNotFound)?; (data.clone(), "legacy-config-etag".to_string()) } else { - return Err(Error::ConfigNotFound); + let objects = self.objects.lock().await; + let (data, etag) = objects.get(object).ok_or(Error::ConfigNotFound)?; + (data.clone(), etag.clone()) }; Ok(GetObjectReader { stream: Box::new(Cursor::new(data.clone())), @@ -6161,8 +6883,8 @@ mod tests { } let mut payload = Vec::new(); tokio::io::AsyncReadExt::read_to_end(&mut data.stream, &mut payload).await?; - let mut state = self.state.lock().await; - match state.as_ref() { + let mut objects = self.objects.lock().await; + match objects.get(object) { Some((_, etag)) => opts.precondition_check(&ObjectInfo { etag: Some(etag.clone()), ..Default::default() @@ -6179,7 +6901,7 @@ mod tests { } } let etag = format!("etag-{}", self.next_etag.fetch_add(1, Ordering::SeqCst) + 1); - *state = Some((payload.clone(), etag.clone())); + objects.insert(object.to_string(), (payload.clone(), etag.clone())); Ok(ObjectInfo { bucket: bucket.to_string(), name: object.to_string(), @@ -6191,6 +6913,131 @@ mod tests { } } + #[async_trait::async_trait] + impl ObjectOperations for CasConfigStore { + type Error = Error; + type ObjectInfo = ObjectInfo; + type ObjectOptions = ObjectOptions; + type FileInfo = FileInfo; + type ObjectToDelete = ObjectToDelete; + type DeletedObject = DeletedObject; + + async fn get_object_info(&self, bucket: &str, object: &str, _opts: &Self::ObjectOptions) -> Result { + let objects = self.objects.lock().await; + let (data, etag) = objects + .get(object) + .ok_or(Error::ObjectNotFound(bucket.to_string(), object.to_string()))?; + Ok(ObjectInfo { + bucket: bucket.to_string(), + name: object.to_string(), + size: data.len() as i64, + actual_size: data.len() as i64, + etag: Some(etag.clone()), + ..Default::default() + }) + } + + async fn verify_object_integrity(&self, _bucket: &str, _object: &str, _opts: &Self::ObjectOptions) -> Result<()> { + Err(Error::NotImplemented) + } + + async fn copy_object( + &self, + _src_bucket: &str, + _src_object: &str, + _dst_bucket: &str, + _dst_object: &str, + _src_info: &mut Self::ObjectInfo, + _src_opts: &Self::ObjectOptions, + _dst_opts: &Self::ObjectOptions, + ) -> Result { + Err(Error::NotImplemented) + } + + async fn delete_object_version( + &self, + _bucket: &str, + _object: &str, + _fi: &Self::FileInfo, + _force_del_marker: bool, + ) -> Result<()> { + Err(Error::NotImplemented) + } + + async fn delete_object(&self, bucket: &str, object: &str, _opts: Self::ObjectOptions) -> Result { + let mut objects = self.objects.lock().await; + let (data, etag) = objects + .remove(object) + .ok_or_else(|| Error::ObjectNotFound(bucket.to_string(), object.to_string()))?; + Ok(ObjectInfo { + bucket: bucket.to_string(), + name: object.to_string(), + size: data.len() as i64, + actual_size: data.len() as i64, + etag: Some(etag), + ..Default::default() + }) + } + + async fn delete_objects( + &self, + _bucket: &str, + _objects: Vec, + _opts: Self::ObjectOptions, + ) -> (Vec, Vec>) { + (Vec::new(), vec![Some(Error::NotImplemented)]) + } + + async fn put_object_metadata( + &self, + _bucket: &str, + _object: &str, + _opts: &Self::ObjectOptions, + ) -> Result { + Err(Error::NotImplemented) + } + + async fn get_object_tags(&self, _bucket: &str, _object: &str, _opts: &Self::ObjectOptions) -> Result { + Err(Error::NotImplemented) + } + + async fn put_object_tags( + &self, + _bucket: &str, + _object: &str, + _tags: &str, + _opts: &Self::ObjectOptions, + ) -> Result { + Err(Error::NotImplemented) + } + + async fn delete_object_tags( + &self, + _bucket: &str, + _object: &str, + _opts: &Self::ObjectOptions, + ) -> Result { + Err(Error::NotImplemented) + } + + async fn add_partial(&self, _bucket: &str, _object: &str, _version_id: &str) -> Result<()> { + Err(Error::NotImplemented) + } + + async fn transition_object(&self, _bucket: &str, _object: &str, _opts: &Self::ObjectOptions) -> Result<()> { + Err(Error::NotImplemented) + } + + async fn restore_transitioned_object( + self: Arc, + _bucket: &str, + _object: &str, + _opts: &Self::ObjectOptions, + ) -> Result<()> { + Err(Error::NotImplemented) + } + } + #[async_trait::async_trait] impl BucketOperations for CasConfigStore { type Error = Error; @@ -6254,16 +7101,48 @@ mod tests { async fn list_objects_v2( self: Arc, - _bucket: &str, - _prefix: &str, - _continuation_token: Option, + bucket: &str, + prefix: &str, + continuation_token: Option, _delimiter: Option, - _max_keys: i32, + max_keys: i32, _fetch_owner: bool, _start_after: Option, _incl_deleted: bool, ) -> Result { - Ok(StorageListObjectsV2Info::default()) + let mut objects = if bucket == RUSTFS_META_BUCKET { + self.objects + .lock() + .await + .keys() + .filter(|object| object.starts_with(prefix)) + .map(|object| ObjectInfo { + bucket: bucket.to_string(), + name: object.clone(), + ..Default::default() + }) + .collect::>() + } else { + Vec::new() + }; + objects.sort_by(|left, right| left.name.cmp(&right.name)); + if let Some(marker) = continuation_token { + objects.retain(|object| object.name > marker); + } + let limit = usize::try_from(max_keys).unwrap_or(0); + let is_truncated = objects.len() > limit; + if is_truncated { + objects.truncate(limit); + } + let next_continuation_token = is_truncated + .then(|| objects.last().map(|object| object.name.clone())) + .flatten(); + Ok(StorageListObjectsV2Info { + objects, + is_truncated, + next_continuation_token, + ..Default::default() + }) } async fn list_object_versions( @@ -6699,7 +7578,7 @@ mod tests { other => panic!("{case} should fail before mutation/save, got {other:?}"), } assert!( - store.state.lock().await.is_none(), + !store.objects.lock().await.contains_key(&tier_config_path(TIER_CONFIG_FILE)), "{case} must not create the durable binary tier config after a failed legacy mutation" ); assert!( @@ -6879,7 +7758,7 @@ mod tests { let guard = manager.read().await; let runtime = registered_tier_driver_runtime(&guard).expect("runtime should remain registered"); assert!(!lock_unpoisoned(&runtime).draining.contains_key("COLD-A")); - assert!(store.state.lock().await.is_some(), "detached update must persist its result"); + assert!(!store.objects.lock().await.is_empty(), "detached update must persist its result"); } #[tokio::test] diff --git a/crates/ecstore/src/services/tier/tier_mutation_intent.rs b/crates/ecstore/src/services/tier/tier_mutation_intent.rs index e842c4650..c4c4cae73 100644 --- a/crates/ecstore/src/services/tier/tier_mutation_intent.rs +++ b/crates/ecstore/src/services/tier/tier_mutation_intent.rs @@ -21,10 +21,12 @@ use uuid::Uuid; use crate::config::com; use crate::disk::RUSTFS_META_BUCKET; use crate::error::{Error, Result as EcstoreResult}; -use crate::object_api::ObjectOptions; +use crate::object_api::{ObjectInfo, ObjectOptions}; use crate::services::tier::tier::TierDestinationId; -use crate::storage_api_contracts::{list::ListOperations as _, object::HTTPPreconditions}; -use crate::store::ECStore; +use crate::storage_api_contracts::{ + list::{ListOperations, StorageListObjectsV2Info}, + object::{EcstoreObjectIO, EcstoreObjectOperations, HTTPPreconditions}, +}; pub(crate) const TIER_MUTATION_INTENT_SCHEMA: &str = "rustfs-tier-mutation-intent-v1"; pub(crate) const MAX_TIER_MUTATION_INTENT_SIZE: usize = rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE; @@ -340,16 +342,19 @@ pub(crate) fn tier_mutation_intent_id_from_record_object_name(object: &str) -> R Uuid::parse_str(mutation_key).map_err(|_| TierMutationIntentError::Corrupt("intent record path has invalid uuid")) } -pub(crate) async fn save_tier_mutation_intent_record(api: Arc, intent: &TierMutationIntent) -> EcstoreResult<()> { +pub(crate) async fn save_tier_mutation_intent_record(api: Arc, intent: &TierMutationIntent) -> EcstoreResult<()> +where + S: EcstoreObjectIO, +{ let object = tier_mutation_intent_record_object_name(intent.mutation_id).map_err(tier_mutation_intent_store_error)?; let data = intent.encode().map_err(tier_mutation_intent_store_error)?; com::save_config(api, &object, data).await } -pub(crate) async fn save_tier_mutation_intent_record_if_absent( - api: Arc, - intent: &TierMutationIntent, -) -> EcstoreResult<()> { +pub(crate) async fn save_tier_mutation_intent_record_if_absent(api: Arc, intent: &TierMutationIntent) -> EcstoreResult<()> +where + S: EcstoreObjectIO, +{ let object = tier_mutation_intent_record_object_name(intent.mutation_id).map_err(tier_mutation_intent_store_error)?; let data = intent.encode().map_err(tier_mutation_intent_store_error)?; com::save_config_with_opts( @@ -368,15 +373,21 @@ pub(crate) async fn save_tier_mutation_intent_record_if_absent( .await } -pub(crate) async fn load_tier_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult { +pub(crate) async fn load_tier_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult +where + S: EcstoreObjectIO, +{ let (intent, _) = load_tier_mutation_intent_record_with_etag(api, mutation_id).await?; Ok(intent) } -pub(crate) async fn load_tier_mutation_intent_record_with_etag( - api: Arc, +pub(crate) async fn load_tier_mutation_intent_record_with_etag( + api: Arc, mutation_id: Uuid, -) -> EcstoreResult<(TierMutationIntent, String)> { +) -> EcstoreResult<(TierMutationIntent, String)> +where + S: EcstoreObjectIO, +{ let object = tier_mutation_intent_record_object_name(mutation_id).map_err(tier_mutation_intent_store_error)?; let (data, object_info) = com::read_config_with_metadata(api, &object, &ObjectOptions::default()).await?; let etag = object_info @@ -387,11 +398,14 @@ pub(crate) async fn load_tier_mutation_intent_record_with_etag( Ok((intent, etag)) } -pub(crate) async fn save_tier_mutation_intent_record_if_current( - api: Arc, +pub(crate) async fn save_tier_mutation_intent_record_if_current( + api: Arc, intent: &TierMutationIntent, current_etag: &str, -) -> EcstoreResult<()> { +) -> EcstoreResult<()> +where + S: EcstoreObjectIO, +{ if current_etag.trim().is_empty() { return Err(Error::other("tier mutation intent current ETag is empty")); } @@ -413,7 +427,10 @@ pub(crate) async fn save_tier_mutation_intent_record_if_current( .await } -pub(crate) async fn delete_tier_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult<()> { +pub(crate) async fn delete_tier_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult<()> +where + S: EcstoreObjectOperations, +{ let object = tier_mutation_intent_record_object_name(mutation_id).map_err(tier_mutation_intent_store_error)?; match com::delete_config(api, &object).await { Ok(()) | Err(Error::ConfigNotFound) => Ok(()), @@ -421,12 +438,15 @@ pub(crate) async fn delete_tier_mutation_intent_record(api: Arc, mutati } } -pub(crate) async fn advance_tier_mutation_intent_record_idempotent( - api: Arc, +pub(crate) async fn advance_tier_mutation_intent_record_idempotent( + api: Arc, mutation_id: Uuid, next: TierMutationIntentState, committed_config_etag: Option, -) -> EcstoreResult<(TierMutationIntent, bool)> { +) -> EcstoreResult<(TierMutationIntent, bool)> +where + S: EcstoreObjectIO, +{ let (mut intent, current_etag) = load_tier_mutation_intent_record_with_etag(api.clone(), mutation_id).await?; let advanced = intent .advance_idempotent(next, committed_config_etag) @@ -437,11 +457,14 @@ pub(crate) async fn advance_tier_mutation_intent_record_idempotent( Ok((intent, advanced)) } -pub(crate) async fn list_tier_mutation_intent_records( - api: Arc, +pub(crate) async fn list_tier_mutation_intent_records( + api: Arc, limit: usize, marker: Option, -) -> EcstoreResult { +) -> EcstoreResult +where + S: EcstoreObjectIO + ListOperations>, +{ if limit == 0 { return Err(Error::other("tier mutation intent scan limit must be greater than zero")); } diff --git a/crates/ecstore/src/services/tier/tier_mutation_peer.rs b/crates/ecstore/src/services/tier/tier_mutation_peer.rs index 5ccfd043b..fd408d304 100644 --- a/crates/ecstore/src/services/tier/tier_mutation_peer.rs +++ b/crates/ecstore/src/services/tier/tier_mutation_peer.rs @@ -17,7 +17,7 @@ use std::sync::Arc; use rustfs_protos::{TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase}; use uuid::Uuid; -use super::tier::TierConfigMgr; +use super::tier::{TierConfigMgr, tier_config_etag_matches}; use super::tier_mutation_intent::{ MAX_TIER_MUTATION_INTENT_SIZE, TierMutationIntent, TierMutationIntentState, advance_tier_mutation_intent_record_idempotent, load_tier_mutation_intent_record, save_tier_mutation_intent_record_if_absent, @@ -140,13 +140,27 @@ async fn handle_commit( ) -> TierMutationPeerResult { let committed_config_etag = parse_commit_etag(canonical_payload)?; let tier_config_mgr = api.tier_config_mgr(); - let (intent, applied) = advance_tier_mutation_intent_record_idempotent( - api, + let (intent, applied) = match advance_tier_mutation_intent_record_idempotent( + api.clone(), mutation_id, TierMutationIntentState::Committed, - Some(committed_config_etag), + Some(committed_config_etag.clone()), ) - .await?; + .await + { + Ok(result) => result, + Err(Error::ConfigNotFound) + if tier_config_etag_matches(api, &committed_config_etag) + .await + .map_err(Error::other)? => + { + return Ok(TierMutationPeerOutcome { + state: TierMutationPeerState::Committed, + applied: false, + }); + } + Err(err) => return Err(err.into()), + }; if intent.state == TierMutationIntentState::Committed { TierConfigMgr::clear_prepared_mutation_intent_block(&tier_config_mgr, mutation_id).await; } diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index fecd97a0d..2d25b4e5c 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -546,7 +546,7 @@ mod tests { runtime::{global::set_object_store_resolver, sources as runtime_sources}, services::tier::{ test_util::{MockWarmBackend, TransitionCleanupStoreBarrier, register_mock_tier}, - tier::TierConfigMgr, + tier::{TIER_CONFIG_FILE, TierConfigMgr}, tier_mutation_intent::{ TIER_MUTATION_INTENT_RECORD_PREFIX, TierMutationIntent, TierMutationIntentKind, TierMutationIntentState, TierMutationIntentTarget, advance_tier_mutation_intent_record_idempotent, delete_tier_mutation_intent_record, @@ -1589,6 +1589,47 @@ mod tests { assert_eq!(loaded.state, TierMutationIntentState::Committed); assert_eq!(loaded.committed_config_etag.as_deref(), Some("new-etag")); + store + .tier_config_mgr() + .read() + .await + .save_tiering_config(store.clone()) + .await + .expect("tier config should persist for cleaned intent commit proof"); + let tier_config_info = store + .get_object_info( + RUSTFS_META_BUCKET, + &format!("{}/{}", com::CONFIG_PREFIX, TIER_CONFIG_FILE), + &ObjectOptions::default(), + ) + .await + .expect("tier config object info should load"); + let tier_config_etag = tier_config_info.etag.expect("tier config should carry an ETag"); + delete_tier_mutation_intent_record(store.clone(), mutation_id) + .await + .expect("committed peer intent cleanup should persist"); + let cleaned_commit_retry = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Commit, + mutation_id, + tier_config_etag.as_bytes(), + ) + .await + .expect("commit retry after durable cleanup should be idempotently terminal"); + assert!(!cleaned_commit_retry.applied); + assert_eq!(cleaned_commit_retry.state, TierMutationPeerState::Committed); + let mismatched_cleaned_commit = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Commit, + mutation_id, + b"not-the-current-etag", + ) + .await + .expect_err("missing intent without a matching committed config ETag must fail closed"); + assert!(matches!(mismatched_cleaned_commit, TierMutationPeerError::Store(Error::ConfigNotFound))); + let abort_id = uuid::Uuid::new_v4(); let abort_intent = tier_mutation_peer_test_intent(abort_id, "COLD-B", [4; 32]); let abort_prepare_payload = abort_intent.encode().expect("abort prepare intent should encode");