diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index c0df19feb..27e7ed643 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -81,8 +81,10 @@ use s3s::S3ErrorCode; 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, - delete_tier_mutation_intent_record, list_tier_mutation_intent_records, + TierMutationDigest, TierMutationIntent, TierMutationIntentKind, TierMutationIntentState, TierMutationIntentTarget, + advance_tier_coordinator_mutation_intent_record_idempotent, delete_tier_coordinator_mutation_intent_record, + delete_tier_mutation_intent_record, list_tier_coordinator_mutation_intent_records, list_tier_mutation_intent_records, + save_tier_coordinator_mutation_intent_record_if_absent, }, warm_backend::WarmBackendImpl, }; @@ -92,6 +94,7 @@ const TIER_OPERATION_DRAIN_TIMEOUT: Duration = Duration::from_secs(30); const TIER_REMOTE_VALIDATION_TIMEOUT: Duration = Duration::from_secs(30); const TIER_REFERENCE_PROOF_LIST_LIMIT: i32 = 1000; const TIER_MUTATION_INTENT_RECOVERY_SCAN_LIMIT: usize = 1000; +const TIER_MUTATION_INTENT_TTL: Duration = Duration::from_secs(15 * 60); const TIER_MUTATION_BLOCK_MESSAGE: &str = "Remote tier configuration is being replaced"; type TierReferenceProofWalkOptions = StorageWalkOptions bool>; @@ -208,6 +211,7 @@ struct TierDriverRuntime { generations: HashMap>, draining: HashMap, prepared_mutation_blocks: HashMap, + prepared_mutation_blocks_revision: u64, next_generation: u64, next_drain_epoch: u64, admin_updates: Arc>, @@ -802,6 +806,91 @@ async fn commit_tier_mutation_peers( Ok(()) } +fn tier_config_candidate_digest(candidate: &TierConfigMgr) -> io::Result { + let bytes = encode_external_tiering_config_blob(candidate)?; + let mut digest = Sha256::new(); + digest.update(&bytes); + Ok(digest.finalize().into()) +} + +fn tier_mutation_intent_expiry_unix_nanos() -> i64 { + let now = OffsetDateTime::now_utc().unix_timestamp_nanos(); + let ttl = i128::try_from(TIER_MUTATION_INTENT_TTL.as_nanos()).unwrap_or(i128::MAX); + i64::try_from(now.saturating_add(ttl)).unwrap_or(i64::MAX) +} + +fn build_coordinator_tier_mutation_intent( + kind: TierMutationIntentKind, + old_config_etag: Option, + candidate: &TierConfigMgr, + affected_targets: Vec, +) -> io::Result> { + if affected_targets.is_empty() || (old_config_etag.is_none() && kind != TierMutationIntentKind::Add) { + return Ok(None); + } + Ok(Some(TierMutationIntent { + mutation_id: uuid::Uuid::new_v4(), + revision: 1, + kind, + state: TierMutationIntentState::Prepared, + old_config_etag, + committed_config_etag: None, + candidate_digest: tier_config_candidate_digest(candidate)?, + affected_targets, + expires_at_unix_nanos: tier_mutation_intent_expiry_unix_nanos(), + })) +} + +async fn save_coordinator_tier_mutation_intent(api: Arc, intent: Option<&TierMutationIntent>) -> io::Result<()> +where + S: TierReferenceProofStore, +{ + if let Some(intent) = intent { + save_tier_coordinator_mutation_intent_record_if_absent(api, intent) + .await + .map_err(io::Error::other)?; + } + Ok(()) +} + +async fn commit_coordinator_tier_mutation_intent( + api: Arc, + intent: Option<&TierMutationIntent>, + committed_config_etag: &str, +) -> io::Result<()> +where + S: TierReferenceProofStore, +{ + if let Some(intent) = intent { + advance_tier_coordinator_mutation_intent_record_idempotent( + api, + intent.mutation_id, + TierMutationIntentState::Committed, + Some(committed_config_etag.to_string()), + ) + .await + .map_err(io::Error::other)?; + } + Ok(()) +} + +async fn abort_coordinator_tier_mutation_intent(api: Arc, intent: Option<&TierMutationIntent>) +where + S: TierReferenceProofStore, +{ + if let Some(intent) = intent + && let Err(err) = advance_tier_coordinator_mutation_intent_record_idempotent( + api, + intent.mutation_id, + TierMutationIntentState::Aborted, + None, + ) + .await + { + warn!(mutation_id = %intent.mutation_id, error = ?err, "failed to abort coordinator tier mutation intent after save failure"); + } +} + async fn apply_tier_candidate_mutation( mutation: TierCandidateMutation, candidate: &mut TierConfigMgr, @@ -2717,8 +2806,28 @@ impl TierConfigMgr { ensure_no_authoritative_tier_object_references(api.clone(), &affected_targets) .await .map_err(TierConfigUpdateError::Publish)?; - candidate - .save_tiering_config_if_current(api, version.as_deref()) + let coordinator_intent = + build_coordinator_tier_mutation_intent(mutation_kind, version.clone(), &candidate, affected_targets) + .map_err(TierConfigUpdateError::Save)?; + save_coordinator_tier_mutation_intent(api.clone(), coordinator_intent.as_ref()) + .await + .map_err(TierConfigUpdateError::Save)?; + let saved = match candidate + .save_tiering_config_if_current_with_info(api.clone(), version.as_deref()) + .await + { + Ok(saved) => saved, + Err(err) => { + abort_coordinator_tier_mutation_intent(api.clone(), coordinator_intent.as_ref()).await; + return Err(TierConfigUpdateError::Save(err)); + } + }; + let committed_config_etag = saved + .etag + .filter(|etag| !etag.trim().is_empty()) + .ok_or_else(|| io::Error::other("tier configuration object is missing an ETag after save")) + .map_err(TierConfigUpdateError::Save)?; + commit_coordinator_tier_mutation_intent(api, coordinator_intent.as_ref(), &committed_config_etag) .await .map_err(TierConfigUpdateError::Save)?; Self::publish_candidate_after_drain(&handle, candidate, driver_tier.as_deref(), transition) @@ -2756,11 +2865,13 @@ impl TierConfigMgr { + NamespaceLocking + ListOperations>, { - let mut mutation_intents = Self::load_tier_mutation_intents(api.clone()).await?; + let mut peer_mutation_intents = Self::load_tier_mutation_intents(api.clone()).await?; + let mut coordinator_mutation_intents = Self::load_coordinator_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 + let config_lock = if peer_mutation_intents .iter() .any(|intent| intent.state == TierMutationIntentState::Committed) + || !coordinator_mutation_intents.is_empty() { let guard = Self::acquire_tier_config_write_lock(api.clone()) .await @@ -2768,14 +2879,18 @@ impl TierConfigMgr { TierConfigUpdateError::Load(err) => err, other => io::Error::other(format!("{other:?}")), })?; - mutation_intents = Self::load_tier_mutation_intents(api.clone()).await?; + peer_mutation_intents = Self::load_tier_mutation_intents(api.clone()).await?; + coordinator_mutation_intents = Self::load_coordinator_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 + Self::recover_committed_coordinator_mutation_intents(api.clone(), &mut coordinator_mutation_intents).await?; + let mut all_mutation_intents = peer_mutation_intents.clone(); + all_mutation_intents.extend(coordinator_mutation_intents.iter().cloned()); + Self::reconcile_committed_mutation_intent_blocks(handle, &all_mutation_intents).await?; + Self::replay_committed_mutation_intents(&peer_mutation_intents).await?; + let allowed_mutation_blocks = all_mutation_intents .iter() .filter(|intent| matches!(intent.state, TierMutationIntentState::Prepared | TierMutationIntentState::Committed)) .map(|intent| intent.mutation_id) @@ -2785,11 +2900,9 @@ impl TierConfigMgr { 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)?; + Self::delete_replayed_committed_mutation_intents(api.clone(), &peer_mutation_intents).await?; + Self::delete_finished_coordinator_mutation_intents(api.clone(), &coordinator_mutation_intents).await?; + Self::reconcile_prepared_mutation_intents_from_store(handle, api).await?; drop(config_lock); Ok(()) } @@ -2821,6 +2934,72 @@ impl TierConfigMgr { } } + async fn load_prepared_coordinator_mutation_intents(api: Arc) -> io::Result> + where + S: EcstoreObjectIO + ListOperations>, + { + let mut intents = Self::load_coordinator_mutation_intents(api).await?; + intents.retain(|intent| intent.state == TierMutationIntentState::Prepared); + Ok(intents) + } + + async fn load_coordinator_mutation_intents(api: Arc) -> io::Result> + where + S: EcstoreObjectIO + ListOperations>, + { + let mut marker = None; + let mut intents = Vec::new(); + loop { + let scan = + list_tier_coordinator_mutation_intent_records(api.clone(), TIER_MUTATION_INTENT_RECOVERY_SCAN_LIMIT, marker) + .await + .map_err(io::Error::other)?; + if scan.failed != 0 { + return Err(io::Error::other(format!( + "failed to scan {} coordinator tier mutation intent record(s) during recovery", + scan.failed + ))); + } + intents.extend(scan.intents); + if !scan.truncated { + return Ok(intents); + } + let next_marker = scan.next_marker.ok_or_else(|| { + io::Error::other("coordinator tier mutation intent recovery scan was truncated without a continuation marker") + })?; + marker = Some(next_marker); + } + } + + async fn recover_committed_coordinator_mutation_intents(api: Arc, intents: &mut [TierMutationIntent]) -> io::Result<()> + where + S: EcstoreObjectIO, + { + if !intents.iter().any(|intent| intent.state == TierMutationIntentState::Prepared) { + return Ok(()); + } + let (candidate, config_etag) = load_tier_config_for_update(api.clone()).await?; + let Some(config_etag) = config_etag else { + return Ok(()); + }; + let current_digest = tier_config_candidate_digest(&candidate)?; + for intent in intents + .iter_mut() + .filter(|intent| intent.state == TierMutationIntentState::Prepared && intent.candidate_digest == current_digest) + { + let (advanced, _) = advance_tier_coordinator_mutation_intent_record_idempotent( + api.clone(), + intent.mutation_id, + TierMutationIntentState::Committed, + Some(config_etag.clone()), + ) + .await + .map_err(io::Error::other)?; + *intent = advanced; + } + Ok(()) + } + async fn replay_committed_mutation_intents(intents: &[TierMutationIntent]) -> io::Result<()> { if !intents .iter() @@ -2857,10 +3036,60 @@ impl TierConfigMgr { Ok(()) } + async fn delete_finished_coordinator_mutation_intents(api: Arc, intents: &[TierMutationIntent]) -> io::Result<()> + where + S: EcstoreObjectOperations, + { + for intent in intents + .iter() + .filter(|intent| matches!(intent.state, TierMutationIntentState::Committed | TierMutationIntentState::Aborted)) + { + delete_tier_coordinator_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> { + Self::reconcile_prepared_mutation_intents_with_revision(handle, intents, None) + .await + .map(|_| ()) + } + + async fn reconcile_prepared_mutation_intents_from_store(handle: &Arc>, api: Arc) -> io::Result<()> + where + S: EcstoreObjectIO + ListOperations>, + { + for _ in 0..3 { + let prepared_blocks_revision = Self::prepared_mutation_blocks_revision(handle).await; + let mut prepared_intents = Self::load_tier_mutation_intents(api.clone()).await?; + prepared_intents.retain(|intent| intent.state == TierMutationIntentState::Prepared); + prepared_intents.extend(Self::load_prepared_coordinator_mutation_intents(api.clone()).await?); + if Self::reconcile_prepared_mutation_intents_with_revision(handle, &prepared_intents, Some(prepared_blocks_revision)) + .await + .map_err(io::Error::other)? + { + return Ok(()); + } + } + Err(io::Error::other("tier mutation prepared block recovery changed repeatedly during reload")) + } + + async fn prepared_mutation_blocks_revision(handle: &Arc>) -> u64 { + let manager = handle.read().await; + let runtime = tier_driver_runtime(handle, &manager); + lock_unpoisoned(&runtime).prepared_mutation_blocks_revision + } + + async fn reconcile_prepared_mutation_intents_with_revision( + handle: &Arc>, + intents: &[TierMutationIntent], + base_revision: Option, + ) -> std::result::Result { let mut prepared_mutation_blocks = HashMap::new(); for intent in intents { Self::collect_mutation_intent_block(&mut prepared_mutation_blocks, intent, TierMutationIntentState::Prepared)?; @@ -2868,8 +3097,15 @@ impl TierConfigMgr { let manager = handle.read().await; let runtime = tier_driver_runtime(handle, &manager); let mut runtime = lock_unpoisoned(&runtime); + if base_revision.is_some_and(|revision| runtime.prepared_mutation_blocks_revision != revision) { + return Ok(false); + } + let changed = runtime.prepared_mutation_blocks != prepared_mutation_blocks; runtime.prepared_mutation_blocks = prepared_mutation_blocks; - Ok(()) + if changed { + runtime.prepared_mutation_blocks_revision = runtime.prepared_mutation_blocks_revision.saturating_add(1); + } + Ok(true) } async fn reconcile_committed_mutation_intent_blocks( @@ -2897,7 +3133,11 @@ impl TierConfigMgr { let mut runtime = lock_unpoisoned(&runtime); let mut prepared_mutation_blocks = runtime.prepared_mutation_blocks.clone(); Self::collect_mutation_intent_block(&mut prepared_mutation_blocks, intent, TierMutationIntentState::Prepared)?; + let changed = prepared_mutation_blocks != runtime.prepared_mutation_blocks; runtime.prepared_mutation_blocks = prepared_mutation_blocks; + if changed { + runtime.prepared_mutation_blocks_revision = runtime.prepared_mutation_blocks_revision.saturating_add(1); + } Ok(()) } @@ -2906,9 +3146,14 @@ impl TierConfigMgr { let Some(runtime) = registered_tier_driver_runtime(&manager) else { return; }; - lock_unpoisoned(&runtime) + let mut runtime = lock_unpoisoned(&runtime); + let before = runtime.prepared_mutation_blocks.len(); + runtime .prepared_mutation_blocks .retain(|_, blocked_mutation_id| *blocked_mutation_id != mutation_id); + if runtime.prepared_mutation_blocks.len() != before { + runtime.prepared_mutation_blocks_revision = runtime.prepared_mutation_blocks_revision.saturating_add(1); + } } fn collect_mutation_intent_block( @@ -3309,6 +3554,26 @@ impl TierConfigMgr { api: Arc, version: Option<&str>, ) -> std::result::Result<(), std::io::Error> + where + S: ObjectIO< + Error = Error, + RangeSpec = HTTPRangeSpec, + HeaderMap = HeaderMap, + ObjectOptions = ObjectOptions, + ObjectInfo = ObjectInfo, + GetObjectReader = GetObjectReader, + PutObjectReader = PutObjReader, + >, + { + self.save_tiering_config_if_current_with_info(api, version).await?; + Ok(()) + } + + async fn save_tiering_config_if_current_with_info( + &self, + api: Arc, + version: Option<&str>, + ) -> std::result::Result where S: ObjectIO< Error = Error, @@ -3333,7 +3598,7 @@ impl TierConfigMgr { }, }; - self.save_config_with_opts( + self.save_config_with_opts_info( api, &config_file, data, @@ -3377,6 +3642,28 @@ impl TierConfigMgr { data: Bytes, opts: &ObjectOptions, ) -> std::result::Result<(), std::io::Error> + where + S: ObjectIO< + Error = Error, + RangeSpec = HTTPRangeSpec, + HeaderMap = HeaderMap, + ObjectOptions = ObjectOptions, + ObjectInfo = ObjectInfo, + GetObjectReader = GetObjectReader, + PutObjectReader = PutObjReader, + >, + { + self.save_config_with_opts_info(api, file, data, opts).await?; + Ok(()) + } + + async fn save_config_with_opts_info( + &self, + api: Arc, + file: &str, + data: Bytes, + opts: &ObjectOptions, + ) -> std::result::Result where S: ObjectIO< Error = Error, @@ -3390,8 +3677,9 @@ impl TierConfigMgr { { debug!("save tier config:{}", file); let mut put_data = PutObjReader::from_vec(data.to_vec()); - let _ = api.put_object(RUSTFS_META_BUCKET, file, &mut put_data, opts).await?; - Ok(()) + api.put_object(RUSTFS_META_BUCKET, file, &mut put_data, opts) + .await + .map_err(io::Error::other) } pub async fn refresh_tier_config(&mut self, api: Arc) { @@ -5646,6 +5934,140 @@ mod tests { ); } + #[tokio::test] + async fn coordinator_mutation_recovery_scan_blocks_new_operation_lease() { + let store = Arc::new(CasConfigStore::default()); + let mutation_id = uuid::Uuid::from_u128(23); + let intent = prepared_remove_intent("COLD-A", mutation_id); + save_tier_coordinator_mutation_intent_record_if_absent(store.clone(), &intent) + .await + .expect("coordinator intent fixture should persist"); + + let loaded_prepared = TierConfigMgr::load_prepared_coordinator_mutation_intents(store.clone()) + .await + .expect("coordinator intent scan should load prepared record"); + assert_eq!(loaded_prepared.len(), 1); + assert_eq!(loaded_prepared[0].mutation_id, mutation_id); + + let manager = TierConfigMgr::new(); + TierConfigMgr::reconcile_prepared_mutation_intents_from_store(&manager, store) + .await + .expect("coordinator prepared recovery should reconcile"); + let err = match TierConfigMgr::acquire_operation_lease(&manager, "COLD-A").await { + Ok(_) => panic!("coordinator prepared intent must block new operation leases"), + Err(err) => err, + }; + assert!(err.message.contains("being replaced"), "{err}"); + } + + #[tokio::test] + async fn coordinator_mutation_recovery_commits_saved_config_digest() { + let store = Arc::new(CasConfigStore::default()); + let mut current = empty_mgr(); + current.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + current + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("current tier config fixture should persist"); + let (_, old_etag) = load_tier_config_for_update(store.clone()) + .await + .expect("current tier config fixture should have metadata"); + + let candidate = empty_mgr(); + let intent = build_coordinator_tier_mutation_intent( + TierMutationIntentKind::Remove, + old_etag.clone(), + &candidate, + build_tier_mutation_affected_targets( + TierMutationIntentKind::Remove, + HashSet::from(["COLD-A".to_string()]), + ¤t, + &candidate, + ) + .expect("remove fixture should produce an affected target"), + ) + .expect("coordinator intent fixture should build") + .expect("remove fixture should require a coordinator intent"); + save_tier_coordinator_mutation_intent_record_if_absent(store.clone(), &intent) + .await + .expect("coordinator intent fixture should persist"); + let mut aborted = prepared_remove_intent("COLD-Z", uuid::Uuid::from_u128(26)); + aborted + .advance(TierMutationIntentState::Aborted, None) + .expect("aborted coordinator fixture should advance"); + save_tier_coordinator_mutation_intent_record_if_absent(store.clone(), &aborted) + .await + .expect("aborted coordinator fixture should persist"); + candidate + .save_tiering_config_if_current(store.clone(), old_etag.as_deref()) + .await + .expect("candidate fixture should persist before coordinator commit"); + + let handle = TierConfigMgr::new(); + TierConfigMgr::reload_handle_with(&handle, store.clone()) + .await + .expect("reload should commit coordinator intent whose candidate digest is already saved"); + + assert!( + !handle.read().await.tiers.contains_key("COLD-A"), + "reload should publish the already saved coordinator candidate" + ); + let reloaded = TierConfigMgr::load_coordinator_mutation_intents(store) + .await + .expect("coordinator intent scan should succeed after recovery cleanup"); + assert!( + reloaded.iter().all(|record| record.mutation_id != intent.mutation_id), + "committed coordinator intent should be removed after local publish" + ); + assert!( + reloaded.iter().all(|record| record.mutation_id != aborted.mutation_id), + "aborted coordinator intent should be removed after local publish" + ); + } + + #[tokio::test] + async fn prepared_mutation_recovery_retries_when_runtime_blocks_change() { + let manager = TierConfigMgr::new(); + let first = prepared_remove_intent("COLD-A", uuid::Uuid::from_u128(24)); + let second = prepared_remove_intent("COLD-B", uuid::Uuid::from_u128(25)); + + TierConfigMgr::reconcile_prepared_mutation_intents(&manager, std::slice::from_ref(&first)) + .await + .expect("initial prepared block should reconcile"); + let base_revision = TierConfigMgr::prepared_mutation_blocks_revision(&manager).await; + TierConfigMgr::apply_prepared_mutation_intent_block(&manager, &second) + .await + .expect("concurrent peer prepare should update runtime blocks"); + + assert!( + !TierConfigMgr::reconcile_prepared_mutation_intents_with_revision( + &manager, + std::slice::from_ref(&first), + Some(base_revision), + ) + .await + .expect("stale revision should be detected without overwriting") + ); + { + let guard = manager.read().await; + let runtime = registered_tier_driver_runtime(&guard).expect("runtime should be registered"); + let blocks = &lock_unpoisoned(&runtime).prepared_mutation_blocks; + assert_eq!(blocks.get("COLD-A"), Some(&first.mutation_id)); + assert_eq!(blocks.get("COLD-B"), Some(&second.mutation_id)); + } + + let retry_revision = TierConfigMgr::prepared_mutation_blocks_revision(&manager).await; + assert!( + TierConfigMgr::reconcile_prepared_mutation_intents_with_revision( + &manager, + &[first.clone(), second.clone()], + Some(retry_revision), + ) + .await + .expect("fresh revision should reconcile") + ); + } + #[tokio::test] async fn reload_handle_fails_closed_when_committed_replay_fails() { let store = Arc::new(CasConfigStore::default()); diff --git a/crates/ecstore/src/services/tier/tier_mutation_intent.rs b/crates/ecstore/src/services/tier/tier_mutation_intent.rs index c4c4cae73..d409069d4 100644 --- a/crates/ecstore/src/services/tier/tier_mutation_intent.rs +++ b/crates/ecstore/src/services/tier/tier_mutation_intent.rs @@ -31,6 +31,7 @@ use crate::storage_api_contracts::{ 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; pub(crate) const TIER_MUTATION_INTENT_RECORD_PREFIX: &str = "tier/mutation-intents/records"; +pub(crate) const TIER_COORDINATOR_MUTATION_INTENT_RECORD_PREFIX: &str = "tier/mutation-intents/coordinators"; pub(crate) type TierMutationDigest = [u8; 32]; pub(crate) type Result = std::result::Result; @@ -295,21 +296,23 @@ impl TierMutationIntent { } pub(crate) fn tier_mutation_intent_record_object_name(mutation_id: Uuid) -> Result { + tier_mutation_intent_record_object_name_with_prefix(TIER_MUTATION_INTENT_RECORD_PREFIX, mutation_id) +} + +fn tier_mutation_intent_record_object_name_with_prefix(prefix: &str, mutation_id: Uuid) -> Result { if mutation_id.is_nil() { return Err(TierMutationIntentError::Corrupt("mutation_id is nil")); } let mutation_key = mutation_id.simple().to_string(); - Ok(format!( - "{}/{}/{}/{}.json", - TIER_MUTATION_INTENT_RECORD_PREFIX, - &mutation_key[..2], - &mutation_key[2..4], - mutation_key - )) + Ok(format!("{}/{}/{}/{}.json", prefix, &mutation_key[..2], &mutation_key[2..4], mutation_key)) } pub(crate) fn tier_mutation_intent_id_from_record_object_name(object: &str) -> Result { - let prefix = format!("{TIER_MUTATION_INTENT_RECORD_PREFIX}/"); + tier_mutation_intent_id_from_record_object_name_with_prefix(TIER_MUTATION_INTENT_RECORD_PREFIX, object) +} + +fn tier_mutation_intent_id_from_record_object_name_with_prefix(prefix: &str, object: &str) -> Result { + let prefix = format!("{prefix}/"); let suffix = object .strip_prefix(&prefix) .ok_or(TierMutationIntentError::Corrupt("intent record path has wrong prefix"))?; @@ -346,16 +349,39 @@ pub(crate) async fn save_tier_mutation_intent_record(api: Arc, intent: &Ti where S: EcstoreObjectIO, { - let object = tier_mutation_intent_record_object_name(intent.mutation_id).map_err(tier_mutation_intent_store_error)?; + let object = tier_mutation_intent_record_object_name_with_prefix(TIER_MUTATION_INTENT_RECORD_PREFIX, 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_coordinator_mutation_intent_record_if_absent( + api: Arc, + intent: &TierMutationIntent, +) -> EcstoreResult<()> +where + S: EcstoreObjectIO, +{ + save_tier_mutation_intent_record_if_absent_with_prefix(api, TIER_COORDINATOR_MUTATION_INTENT_RECORD_PREFIX, intent).await +} + 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)?; + save_tier_mutation_intent_record_if_absent_with_prefix(api, TIER_MUTATION_INTENT_RECORD_PREFIX, intent).await +} + +async fn save_tier_mutation_intent_record_if_absent_with_prefix( + api: Arc, + prefix: &str, + intent: &TierMutationIntent, +) -> EcstoreResult<()> +where + S: EcstoreObjectIO, +{ + let object = tier_mutation_intent_record_object_name_with_prefix(prefix, 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( api, @@ -388,7 +414,19 @@ pub(crate) async fn load_tier_mutation_intent_record_with_etag( where S: EcstoreObjectIO, { - let object = tier_mutation_intent_record_object_name(mutation_id).map_err(tier_mutation_intent_store_error)?; + load_tier_mutation_intent_record_with_etag_at_prefix(api, TIER_MUTATION_INTENT_RECORD_PREFIX, mutation_id).await +} + +async fn load_tier_mutation_intent_record_with_etag_at_prefix( + api: Arc, + prefix: &str, + mutation_id: Uuid, +) -> EcstoreResult<(TierMutationIntent, String)> +where + S: EcstoreObjectIO, +{ + let object = + tier_mutation_intent_record_object_name_with_prefix(prefix, 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 .etag @@ -409,7 +447,23 @@ where if current_etag.trim().is_empty() { return Err(Error::other("tier mutation intent current ETag is empty")); } - let object = tier_mutation_intent_record_object_name(intent.mutation_id).map_err(tier_mutation_intent_store_error)?; + save_tier_mutation_intent_record_if_current_with_prefix(api, TIER_MUTATION_INTENT_RECORD_PREFIX, intent, current_etag).await +} + +async fn save_tier_mutation_intent_record_if_current_with_prefix( + api: Arc, + prefix: &str, + intent: &TierMutationIntent, + current_etag: &str, +) -> EcstoreResult<()> +where + S: EcstoreObjectIO, +{ + if current_etag.trim().is_empty() { + return Err(Error::other("tier mutation intent current ETag is empty")); + } + let object = tier_mutation_intent_record_object_name_with_prefix(prefix, 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( api, @@ -431,7 +485,22 @@ pub(crate) async fn delete_tier_mutation_intent_record(api: Arc, mutation_ where S: EcstoreObjectOperations, { - let object = tier_mutation_intent_record_object_name(mutation_id).map_err(tier_mutation_intent_store_error)?; + delete_tier_mutation_intent_record_with_prefix(api, TIER_MUTATION_INTENT_RECORD_PREFIX, mutation_id).await +} + +pub(crate) async fn delete_tier_coordinator_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult<()> +where + S: EcstoreObjectOperations, +{ + delete_tier_mutation_intent_record_with_prefix(api, TIER_COORDINATOR_MUTATION_INTENT_RECORD_PREFIX, mutation_id).await +} + +async fn delete_tier_mutation_intent_record_with_prefix(api: Arc, prefix: &str, mutation_id: Uuid) -> EcstoreResult<()> +where + S: EcstoreObjectOperations, +{ + let object = + tier_mutation_intent_record_object_name_with_prefix(prefix, mutation_id).map_err(tier_mutation_intent_store_error)?; match com::delete_config(api, &object).await { Ok(()) | Err(Error::ConfigNotFound) => Ok(()), Err(err) => Err(err), @@ -447,12 +516,52 @@ pub(crate) async fn advance_tier_mutation_intent_record_idempotent( where S: EcstoreObjectIO, { - let (mut intent, current_etag) = load_tier_mutation_intent_record_with_etag(api.clone(), mutation_id).await?; + advance_tier_mutation_intent_record_idempotent_at_prefix( + api, + TIER_MUTATION_INTENT_RECORD_PREFIX, + mutation_id, + next, + committed_config_etag, + ) + .await +} + +pub(crate) async fn advance_tier_coordinator_mutation_intent_record_idempotent( + api: Arc, + mutation_id: Uuid, + next: TierMutationIntentState, + committed_config_etag: Option, +) -> EcstoreResult<(TierMutationIntent, bool)> +where + S: EcstoreObjectIO, +{ + advance_tier_mutation_intent_record_idempotent_at_prefix( + api, + TIER_COORDINATOR_MUTATION_INTENT_RECORD_PREFIX, + mutation_id, + next, + committed_config_etag, + ) + .await +} + +async fn advance_tier_mutation_intent_record_idempotent_at_prefix( + api: Arc, + prefix: &str, + mutation_id: Uuid, + next: TierMutationIntentState, + committed_config_etag: Option, +) -> EcstoreResult<(TierMutationIntent, bool)> +where + S: EcstoreObjectIO, +{ + let (mut intent, current_etag) = + load_tier_mutation_intent_record_with_etag_at_prefix(api.clone(), prefix, mutation_id).await?; let advanced = intent .advance_idempotent(next, committed_config_etag) .map_err(tier_mutation_intent_store_error)?; if advanced { - save_tier_mutation_intent_record_if_current(api, &intent, ¤t_etag).await?; + save_tier_mutation_intent_record_if_current_with_prefix(api, prefix, &intent, ¤t_etag).await?; } Ok((intent, advanced)) } @@ -462,6 +571,29 @@ pub(crate) async fn list_tier_mutation_intent_records( limit: usize, marker: Option, ) -> EcstoreResult +where + S: EcstoreObjectIO + ListOperations>, +{ + list_tier_mutation_intent_records_at_prefix(api, TIER_MUTATION_INTENT_RECORD_PREFIX, limit, marker).await +} + +pub(crate) async fn list_tier_coordinator_mutation_intent_records( + api: Arc, + limit: usize, + marker: Option, +) -> EcstoreResult +where + S: EcstoreObjectIO + ListOperations>, +{ + list_tier_mutation_intent_records_at_prefix(api, TIER_COORDINATOR_MUTATION_INTENT_RECORD_PREFIX, limit, marker).await +} + +async fn list_tier_mutation_intent_records_at_prefix( + api: Arc, + prefix: &str, + limit: usize, + marker: Option, +) -> EcstoreResult where S: EcstoreObjectIO + ListOperations>, { @@ -473,7 +605,7 @@ where .clone() .list_objects_v2( RUSTFS_META_BUCKET, - TIER_MUTATION_INTENT_RECORD_PREFIX, + prefix, marker, None, i32::try_from(limit).map_or(i32::MAX, |value| value), @@ -492,15 +624,15 @@ where for object in list.objects { scan.scanned += 1; - let mutation_id = match tier_mutation_intent_id_from_record_object_name(&object.name) { + let mutation_id = match tier_mutation_intent_id_from_record_object_name_with_prefix(prefix, &object.name) { Ok(mutation_id) => mutation_id, Err(_) => { scan.failed += 1; continue; } }; - match load_tier_mutation_intent_record(api.clone(), mutation_id).await { - Ok(intent) => scan.intents.push(intent), + match load_tier_mutation_intent_record_with_etag_at_prefix(api.clone(), prefix, mutation_id).await { + Ok((intent, _)) => scan.intents.push(intent), Err(Error::ConfigNotFound) => {} Err(_) => scan.failed += 1, }