diff --git a/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs b/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs index e3f3036e8..cedd04659 100644 --- a/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs +++ b/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs @@ -700,7 +700,7 @@ async fn recover_unknown_upload_outcome( .map_err(Error::other)?; match lease - .probe_transition_candidate(&transaction.remote_object) + .probe_transition_candidate_for(&transaction.remote_object, transaction.transaction_id) .await .map_err(Error::other)? { diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 5ead6a9d9..ce0bc97e0 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -236,6 +236,7 @@ struct TierPublishTransition { struct PreparedTierDriver { tier_name: String, + tier_config: TierConfig, config_fingerprint: TierDriverFingerprint, backend_identity: TierDestinationId, exact_get_delete: bool, @@ -1342,6 +1343,7 @@ fn tier_exact_get_delete(config: &TierConfig) -> bool { struct TierDriverGeneration { tier_name: Arc, + tier_config: TierConfig, generation: DriverRevision, // Process-local only: this may reflect credential changes and must never be persisted or logged. config_fingerprint: TierDriverFingerprint, @@ -1349,6 +1351,9 @@ struct TierDriverGeneration { backend_identity: TierDestinationId, exact_get_delete: bool, driver: SharedWarmBackend, + reconciler: tokio::sync::OnceCell< + Option>, + >, accepting: AtomicBool, active_leases: AtomicUsize, drained: tokio::sync::Notify, @@ -1473,6 +1478,35 @@ impl TierOperationLease { self.inner.backend_identity } + pub(crate) async fn probe_transition_candidate_for( + &self, + object: &str, + transaction_id: uuid::Uuid, + ) -> io::Result { + let Some(reconciler) = self + .inner + .reconciler + .get_or_try_init(|| async { + crate::services::tier::warm_backend::new_transition_candidate_reconciler(&self.inner.tier_config) + .await + .map(|reconciler| reconciler.map(Arc::from)) + }) + .await + .map_err(|err| io::Error::other(err.message))? + else { + return self.inner.driver.probe_transition_candidate(object).await; + }; + reconciler + .probe_transition_candidate_for( + object, + crate::services::tier::warm_backend::TransitionCandidateIdentity { + transaction_id, + destination_id: self.backend_identity(), + }, + ) + .await + } + pub(crate) fn validate_remote_version_id(&self, remote_version_id: &str) -> io::Result<()> { self.inner.driver.validate_remote_version_id(remote_version_id)?; if !remote_version_id.is_empty() && !self.inner.exact_get_delete { @@ -2833,6 +2867,7 @@ impl TierConfigMgr { let exact_get_delete = tier_exact_get_delete(config); Some(PreparedTierDriver { tier_name: tier_name.to_string(), + tier_config: config.clone(), config_fingerprint, backend_identity, exact_get_delete, @@ -2861,11 +2896,13 @@ impl TierConfigMgr { })?; let entry = Arc::new(TierDriverGeneration { tier_name: Arc::from(prepared.tier_name.as_str()), + tier_config: prepared.tier_config.clone(), generation, config_fingerprint: prepared.config_fingerprint, backend_identity: prepared.backend_identity, exact_get_delete: prepared.exact_get_delete, driver: prepared.driver.clone(), + reconciler: tokio::sync::OnceCell::new(), accepting: AtomicBool::new(true), active_leases: AtomicUsize::new(0), drained: tokio::sync::Notify::new(), @@ -3609,11 +3646,13 @@ impl TierConfigMgr { let driver: SharedWarmBackend = Arc::from(driver); let entry = Arc::new(TierDriverGeneration { tier_name: Arc::from(tier_name), + tier_config: config.clone(), generation, config_fingerprint, backend_identity, exact_get_delete, driver: driver.clone(), + reconciler: tokio::sync::OnceCell::new(), accepting: AtomicBool::new(true), active_leases: AtomicUsize::new(0), drained: tokio::sync::Notify::new(), diff --git a/crates/ecstore/src/services/tier/warm_backend.rs b/crates/ecstore/src/services/tier/warm_backend.rs index 1c667c843..464dd41a9 100644 --- a/crates/ecstore/src/services/tier/warm_backend.rs +++ b/crates/ecstore/src/services/tier/warm_backend.rs @@ -73,6 +73,21 @@ pub enum TransitionCandidateProbe { Unsupported, } +#[derive(Clone, Copy)] +pub(crate) struct TransitionCandidateIdentity { + pub transaction_id: uuid::Uuid, + pub destination_id: [u8; 32], +} + +#[async_trait::async_trait] +pub(crate) trait TransitionCandidateReconciler { + async fn probe_transition_candidate_for( + &self, + object: &str, + identity: TransitionCandidateIdentity, + ) -> Result; +} + #[async_trait::async_trait] pub trait WarmBackend { async fn validate(&self) -> Result<(), std::io::Error> { @@ -189,6 +204,20 @@ pub fn build_transition_put_options(storage_class: String, mut metadata: HashMap metadata.remove(key); } + for suffix in [ + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TRANSACTION_ID, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + ] { + for key in [ + rustfs_utils::http::metadata_compat::internal_key_rustfs(suffix), + format!("{}{}", rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX, suffix), + ] { + if let Some(value) = metadata.remove(&key) { + metadata.insert(format!("x-amz-meta-{key}"), value); + } + } + } + opts.user_metadata = metadata; opts } @@ -446,6 +475,51 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result Result>, AdminError> { + let reconciler: Box = match tier.tier_type { + TierType::S3 => Box::new( + WarmBackendS3::new(tier.s3.as_ref().ok_or_else(|| ERR_TIER_INVALID_CONFIG.clone())?, &tier.name) + .await + .map_err(|err| { + let mut admin_err = ERR_TIER_INVALID_CONFIG.clone(); + admin_err.message = err.to_string(); + admin_err + })?, + ), + TierType::MinIO => Box::new( + WarmBackendMinIO::new(tier.minio.as_ref().ok_or_else(|| ERR_TIER_INVALID_CONFIG.clone())?, &tier.name) + .await + .map_err(|err| { + let mut admin_err = ERR_TIER_INVALID_CONFIG.clone(); + admin_err.message = err.to_string(); + admin_err + })?, + ), + TierType::RustFS => Box::new( + WarmBackendRustFS::new(tier.rustfs.as_ref().ok_or_else(|| ERR_TIER_INVALID_CONFIG.clone())?, &tier.name) + .await + .map_err(|err| { + let mut admin_err = ERR_TIER_INVALID_CONFIG.clone(); + admin_err.message = err.to_string(); + admin_err + })?, + ), + TierType::R2 => Box::new( + WarmBackendR2::new(tier.r2.as_ref().ok_or_else(|| ERR_TIER_INVALID_CONFIG.clone())?, &tier.name) + .await + .map_err(|err| { + let mut admin_err = ERR_TIER_INVALID_CONFIG.clone(); + admin_err.message = err.to_string(); + admin_err + })?, + ), + _ => return Ok(None), + }; + Ok(Some(reconciler)) +} + #[cfg(test)] mod tests { use super::*; @@ -775,6 +849,37 @@ mod tests { assert!(!opts.user_metadata.contains_key(X_AMZ_REPLICATION_STATUS.as_str())); } + #[test] + fn build_transition_put_options_persists_both_candidate_identity_keys_as_s3_metadata() { + let mut metadata = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TRANSACTION_ID, + "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa".to_string(), + ); + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + "5a".repeat(32), + ); + + let opts = build_transition_put_options("COLD".to_string(), metadata); + + for suffix in [ + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TRANSACTION_ID, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + ] { + assert!(opts.user_metadata.contains_key(&format!( + "x-amz-meta-{}", + rustfs_utils::http::metadata_compat::internal_key_rustfs(suffix) + ))); + assert!(opts.user_metadata.contains_key(&format!( + "x-amz-meta-{}{suffix}", + rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX + ))); + } + } + #[test] fn build_transition_put_options_requests_no_checksum_and_content_md5() { // Regression for rustfs/rustfs#4811: transition uploads must leave the diff --git a/crates/ecstore/src/services/tier/warm_backend_minio.rs b/crates/ecstore/src/services/tier/warm_backend_minio.rs index fa62df7cb..a97a01baf 100644 --- a/crates/ecstore/src/services/tier/warm_backend_minio.rs +++ b/crates/ecstore/src/services/tier/warm_backend_minio.rs @@ -135,6 +135,20 @@ impl WarmBackend for WarmBackendMinIO { } } +#[async_trait::async_trait] +impl crate::services::tier::warm_backend::TransitionCandidateReconciler for WarmBackendMinIO { + async fn probe_transition_candidate_for( + &self, + object: &str, + identity: crate::services::tier::warm_backend::TransitionCandidateIdentity, + ) -> Result { + crate::services::tier::warm_backend::TransitionCandidateReconciler::probe_transition_candidate_for( + &self.0, object, identity, + ) + .await + } +} + fn optimal_part_size(object_size: i64) -> Result { let mut object_size = object_size; if object_size == -1 { diff --git a/crates/ecstore/src/services/tier/warm_backend_r2.rs b/crates/ecstore/src/services/tier/warm_backend_r2.rs index e93c1ad9c..cee1e1108 100644 --- a/crates/ecstore/src/services/tier/warm_backend_r2.rs +++ b/crates/ecstore/src/services/tier/warm_backend_r2.rs @@ -135,6 +135,20 @@ impl WarmBackend for WarmBackendR2 { } } +#[async_trait::async_trait] +impl crate::services::tier::warm_backend::TransitionCandidateReconciler for WarmBackendR2 { + async fn probe_transition_candidate_for( + &self, + object: &str, + identity: crate::services::tier::warm_backend::TransitionCandidateIdentity, + ) -> Result { + crate::services::tier::warm_backend::TransitionCandidateReconciler::probe_transition_candidate_for( + &self.0, object, identity, + ) + .await + } +} + fn optimal_part_size(object_size: i64) -> Result { let mut object_size = object_size; if object_size == -1 { diff --git a/crates/ecstore/src/services/tier/warm_backend_rustfs.rs b/crates/ecstore/src/services/tier/warm_backend_rustfs.rs index db03c9041..503f1238e 100644 --- a/crates/ecstore/src/services/tier/warm_backend_rustfs.rs +++ b/crates/ecstore/src/services/tier/warm_backend_rustfs.rs @@ -132,6 +132,20 @@ impl WarmBackend for WarmBackendRustFS { } } +#[async_trait::async_trait] +impl crate::services::tier::warm_backend::TransitionCandidateReconciler for WarmBackendRustFS { + async fn probe_transition_candidate_for( + &self, + object: &str, + identity: crate::services::tier::warm_backend::TransitionCandidateIdentity, + ) -> Result { + crate::services::tier::warm_backend::TransitionCandidateReconciler::probe_transition_candidate_for( + &self.0, object, identity, + ) + .await + } +} + fn optimal_part_size(object_size: i64) -> Result { let mut object_size = object_size; if object_size == -1 { diff --git a/crates/ecstore/src/services/tier/warm_backend_s3.rs b/crates/ecstore/src/services/tier/warm_backend_s3.rs index cae31d853..cda27acf0 100644 --- a/crates/ecstore/src/services/tier/warm_backend_s3.rs +++ b/crates/ecstore/src/services/tier/warm_backend_s3.rs @@ -37,7 +37,10 @@ use crate::error::ErrorResponse; use crate::error::error_resp_to_object_err; use crate::services::tier::{ tier_config::TierS3, - warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options}, + warm_backend::{ + TransitionCandidateIdentity, TransitionCandidateProbe, TransitionCandidateReconciler, WarmBackend, WarmBackendGetOpts, + build_transition_put_options, + }, }; use http::HeaderMap; use rustfs_utils::egress::validate_outbound_url; @@ -219,6 +222,92 @@ impl WarmBackendS3 { advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?; } } + + async fn probe_transition_candidate_identity( + &self, + object: &str, + identity: TransitionCandidateIdentity, + bucket_versioning: RemoteBucketVersioning, + ) -> Result { + let remote_object = self.get_dest(object); + let mut opts = ListObjectsOptions::default(); + opts.set("prefix", &remote_object); + opts.set("max-keys", "1000"); + let mut key_marker = String::new(); + let mut version_id_marker = String::new(); + let mut matched_version = None; + let mut saw_unproven_candidate = false; + + loop { + let versions = self + .client + .list_object_versions_query(&self.bucket, &opts, &key_marker, &version_id_marker, "") + .await?; + for version in versions.versions.iter().filter(|version| version.key == remote_object) { + let mut stat_opts = GetObjectOptions::default(); + stat_opts.version_id.clone_from(&version.version_id); + let info = self.client.stat_object(&self.bucket, &remote_object, &stat_opts).await?; + let mut metadata = info.user_metadata; + for (name, value) in &info.metadata { + if (name + .as_str() + .starts_with(rustfs_utils::http::metadata_compat::RUSTFS_INTERNAL_PREFIX) + || name + .as_str() + .starts_with(rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX)) + && let Ok(value) = value.to_str() + { + metadata.insert(name.as_str().to_string(), value.to_string()); + } + } + if transition_candidate_metadata_matches(&metadata, identity)? { + if matched_version.is_some() { + return Ok(TransitionCandidateProbe::Ambiguous); + } + matched_version = Some(version.version_id.clone()); + } else { + saw_unproven_candidate = true; + } + } + if !versions.is_truncated { + if matched_version.is_none() && saw_unproven_candidate { + return Ok(TransitionCandidateProbe::Unsupported); + } + let candidates = TransitionCandidateVersions { + version_id: matched_version, + ambiguous: false, + }; + return classify_transition_candidates(candidates, bucket_versioning); + } + advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?; + } + } +} + +fn transition_candidate_metadata_matches( + metadata: &HashMap, + identity: TransitionCandidateIdentity, +) -> Result { + use rustfs_utils::http::metadata_compat::{ + SUFFIX_TRANSITION_TIER_DESTINATION_ID, SUFFIX_TRANSITION_TRANSACTION_ID, contains_key_str, get_consistent_str, + }; + + let transaction_id = get_consistent_str(metadata, SUFFIX_TRANSITION_TRANSACTION_ID); + let destination_id = get_consistent_str(metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID); + if transaction_id.is_none() || destination_id.is_none() { + if contains_key_str(metadata, SUFFIX_TRANSITION_TRANSACTION_ID) + || contains_key_str(metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) + { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "transition candidate identity metadata is empty or conflicting", + )); + } + return Ok(false); + } + let expected_transaction_id = identity.transaction_id.to_string(); + let expected_destination_id = rustfs_utils::crypto::hex(identity.destination_id); + Ok(transaction_id == Some(expected_transaction_id.as_str()) && destination_id == Some(expected_destination_id.as_str())) } fn classify_transition_candidates( @@ -342,6 +431,60 @@ mod tests { candidates.classify(bucket_versioning) } + fn candidate_identity() -> TransitionCandidateIdentity { + TransitionCandidateIdentity { + transaction_id: uuid::Uuid::parse_str("aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa").unwrap(), + destination_id: [0x5a; 32], + } + } + + fn candidate_metadata(identity: TransitionCandidateIdentity) -> HashMap { + let mut metadata = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TRANSACTION_ID, + identity.transaction_id.to_string(), + ); + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::crypto::hex(identity.destination_id), + ); + metadata + } + + #[test] + fn transition_candidate_identity_requires_exact_compatible_metadata() { + let identity = candidate_identity(); + let metadata = candidate_metadata(identity); + assert!(transition_candidate_metadata_matches(&metadata, identity).unwrap()); + + let mut adjacent = metadata.clone(); + rustfs_utils::http::metadata_compat::insert_str( + &mut adjacent, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TRANSACTION_ID, + uuid::Uuid::parse_str("bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb") + .unwrap() + .to_string(), + ); + assert!(!transition_candidate_metadata_matches(&adjacent, identity).unwrap()); + + assert!(!transition_candidate_metadata_matches(&HashMap::new(), identity).unwrap()); + } + + #[test] + fn transition_candidate_identity_rejects_conflicting_compatibility_keys() { + let identity = candidate_identity(); + let mut metadata = candidate_metadata(identity); + metadata.insert( + rustfs_utils::http::metadata_compat::internal_key_rustfs( + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TRANSACTION_ID, + ), + uuid::Uuid::new_v4().to_string(), + ); + assert!(transition_candidate_metadata_matches(&metadata, identity).is_err()); + } + #[test] fn transition_candidate_probe_classifier_is_fail_closed() { assert_eq!( @@ -503,3 +646,16 @@ impl WarmBackend for WarmBackendS3 { Ok(result.common_prefixes.len() > 0 || result.contents.len() > 0) } } + +#[async_trait::async_trait] +impl TransitionCandidateReconciler for WarmBackendS3 { + async fn probe_transition_candidate_for( + &self, + object: &str, + identity: TransitionCandidateIdentity, + ) -> Result { + let bucket_versioning = self.remote_bucket_versioning().await?; + self.probe_transition_candidate_identity(object, identity, bucket_versioning) + .await + } +} diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 122cc5d8c..03997eef2 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -3841,6 +3841,16 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { let dest_obj = transaction.remote_object.clone(); let mut transition_meta = (*oi.user_defined).clone(); transition_meta.insert("name".to_string(), object.to_string()); + rustfs_utils::http::metadata_compat::insert_str( + &mut transition_meta, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TRANSACTION_ID, + transaction.transaction_id.to_string(), + ); + rustfs_utils::http::metadata_compat::insert_str( + &mut transition_meta, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::crypto::hex(transaction.backend_fingerprint), + ); if let Some(content_type) = oi.content_type.as_ref().filter(|value| !value.is_empty()) { transition_meta.insert(CONTENT_TYPE.to_ascii_lowercase(), content_type.clone()); diff --git a/crates/utils/src/http/metadata_compat.rs b/crates/utils/src/http/metadata_compat.rs index c95f5b3d7..ab680c953 100644 --- a/crates/utils/src/http/metadata_compat.rs +++ b/crates/utils/src/http/metadata_compat.rs @@ -40,6 +40,7 @@ pub const SUFFIX_TRANSITIONED_VERSION_ID: &str = "transitioned-versionID"; pub const SUFFIX_TRANSITIONED_VERSION_STATE: &str = "transitioned-version-state"; pub const SUFFIX_TRANSITION_TIER: &str = "transition-tier"; pub const SUFFIX_TRANSITION_TIER_DESTINATION_ID: &str = "transition-tier-destination-id"; +pub const SUFFIX_TRANSITION_TRANSACTION_ID: &str = "transition-transaction-id"; pub const SUFFIX_RESTORE_OPERATION_ID: &str = "restore-operation-id"; pub const SUFFIX_FREE_VERSION: &str = "free-version"; pub const SUFFIX_PURGESTATUS: &str = "purgestatus";