From 821e056f706641b505d59756d10db56c4f78c87f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Tue, 28 Jul 2026 15:10:53 +0800 Subject: [PATCH] fix(tiering): gate remote version state safely --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 124 +++++++++++++- .../bucket/lifecycle/tier_delete_journal.rs | 135 ++++++++++++++-- .../src/bucket/lifecycle/tier_sweeper.rs | 8 +- crates/ecstore/src/config/com.rs | 1 + crates/ecstore/src/object_api/types.rs | 3 + crates/ecstore/src/services/tier/tier.rs | 3 +- .../src/set_disk/core/io_primitives.rs | 1 + crates/ecstore/src/set_disk/metadata.rs | 6 + crates/ecstore/src/set_disk/ops/object.rs | 93 +++++++++-- crates/ecstore/src/store/init.rs | 6 +- crates/filemeta/src/fileinfo.rs | 32 ++++ crates/filemeta/src/filemeta/version.rs | 151 ++++++++++++++++-- crates/utils/src/http/metadata_compat.rs | 1 + 13 files changed, 515 insertions(+), 49 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index f0d26f72a..0e04bca35 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -554,6 +554,7 @@ async fn delete_free_version_remote_object( oi: &ObjectInfo, tier_config_mgr: &Arc>, ) -> Result<(), std::io::Error> { + let version_id_exact = validate_transition_remote_version(oi)?; let identity = tier_destination_id_from_metadata(&oi.user_defined)? .ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?; delete_object_from_remote_tier_idempotent_with_manager_and_identity( @@ -562,7 +563,7 @@ async fn delete_free_version_remote_object( &oi.transitioned_object.tier, identity, tier_config_mgr, - false, + version_id_exact, ) .await?; Ok(()) @@ -4135,6 +4136,23 @@ pub async fn get_transitioned_object_reader( get_transitioned_object_reader_with_tier_manager(bucket, object, rs, h, oi, opts, &tier_config_mgr).await } +fn validate_transition_remote_version(oi: &ObjectInfo) -> Result { + let version = oi.transitioned_object.version_id.as_str(); + match oi.transition_version_state { + rustfs_filemeta::TransitionVersionState::Unknown => Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object version state is unknown", + )), + rustfs_filemeta::TransitionVersionState::KnownDisabled if version.is_empty() => Ok(false), + rustfs_filemeta::TransitionVersionState::SuspendedNull if version == "null" => Ok(true), + rustfs_filemeta::TransitionVersionState::Exact if !version.is_empty() && version != "null" => Ok(true), + _ => Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object version state conflicts with its version ID", + )), + } +} + pub(crate) async fn get_transitioned_object_reader_with_tier_manager( bucket: &str, object: &str, @@ -4144,6 +4162,7 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager( opts: &ObjectOptions, tier_config_mgr: &Arc>, ) -> Result { + validate_transition_remote_version(oi)?; let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?; let lease = match expected_identity { Some(identity) => { @@ -5438,6 +5457,7 @@ mod tests { tier: tier.clone(), ..Default::default() }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Exact, ..Default::default() }; @@ -5501,6 +5521,7 @@ mod tests { tier, ..Default::default() }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Exact, ..Default::default() }; @@ -5523,6 +5544,70 @@ mod tests { assert_eq!(backend.get_count().await, 0); } + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_rejects_unknown_version_state_before_backend_io() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 1, + transitioned_object: TransitionedObject { + name: "remote/object".to_string(), + version_id: String::new(), + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier, + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + ..Default::default() + }; + + let err = match get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &None, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + ) + .await + { + Ok(_) => panic!("unknown remote version state must fail before backend IO"), + Err(err) => err, + }; + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(backend.get_count().await, 0); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_delete_rejects_unknown_version_state_before_backend_io() { + let manager = TierConfigMgr::new(); + let backend = register_mock_tier(&manager, "WARM").await; + let object_info = ObjectInfo { + transitioned_object: TransitionedObject { + name: "remote/object".to_string(), + version_id: "legacy-version".to_string(), + tier: "WARM".to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + ..Default::default() + }; + + let err = super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect_err("unknown remote version state must fail before backend IO"); + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(backend.remove_count().await, 0); + } + #[cfg(feature = "test-util")] #[tokio::test] async fn free_version_remote_delete_requires_persisted_destination_identity() { @@ -5834,7 +5919,8 @@ mod tests { version_id: "remote-version".to_string(), tier_name: "WARM".to_string(), backend_identity: Some([1; 32]), - version_id_exact: false, + version_id_exact: true, + version_state: rustfs_filemeta::TransitionVersionState::Exact, }; let err = state @@ -5945,7 +6031,8 @@ mod tests { version_id: "remote-version".to_string(), tier_name: "WARM".to_string(), backend_identity: Some([1; 32]), - version_id_exact: false, + version_id_exact: true, + version_state: rustfs_filemeta::TransitionVersionState::Exact, }; state @@ -9866,6 +9953,32 @@ mod tests { (backend, identity_hex) } + #[cfg(feature = "test-util")] + #[tokio::test] + async fn journal_replay_rejects_unknown_version_state_before_backend_io() { + let (_disk_paths, ecstore) = setup_test_env().await; + let (backend, _) = register_recovery_mock_tier(&ecstore).await; + let identity = TierConfigMgr::acquire_operation_lease(&ecstore.tier_config_mgr(), "WARM") + .await + .expect("mock tier lease should be available") + .backend_identity(); + let je = Jentry { + obj_name: "remote/object".to_string(), + version_id: "legacy-version".to_string(), + tier_name: "WARM".to_string(), + backend_identity: Some(identity), + version_id_exact: false, + version_state: rustfs_filemeta::TransitionVersionState::Unknown, + }; + + let err = crate::bucket::lifecycle::tier_delete_journal::process_tier_delete_journal_entry(ecstore, &je) + .await + .expect_err("unknown journal state must fail before backend IO"); + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(backend.remove_count().await, 0); + } + async fn seed_recoverable_free_version( disk_paths: &[PathBuf], bucket: &str, @@ -9885,6 +9998,7 @@ mod tests { identity, ); } + let transition_version_id = Uuid::new_v4(); let mut metadata = FileMeta::new(); metadata .add_version(FileInfo { @@ -9893,7 +10007,9 @@ mod tests { version_id: Some(object_version_id), transition_status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), transitioned_objname: format!("remote/{bucket}/{object}"), - transition_version_id: Some(Uuid::new_v4()), + transition_version_id: Some(transition_version_id), + transition_version: Some(transition_version_id.to_string()), + transition_version_state: rustfs_filemeta::TransitionVersionState::Exact, transition_tier: "WARM".to_string(), mod_time: Some(OffsetDateTime::now_utc()), metadata: transitioned_metadata, diff --git a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs index 39d07ed05..b7cb3cc32 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs @@ -42,6 +42,7 @@ const TIER_DELETE_JOURNAL_RECOVERY_INTERVAL: Duration = Duration::from_secs(60); const TIER_DELETE_JOURNAL_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300); const TIER_DELETE_JOURNAL_VERSION: u8 = 2; const TIER_DELETE_JOURNAL_EXACT_VERSION: u8 = 3; +const TIER_DELETE_JOURNAL_STATE_VERSION: u8 = 4; pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = "ilm/tier-delete-journal/"; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -55,24 +56,35 @@ struct PersistedTierDeleteJournalEntry { backend_identity: Option<[u8; 32]>, #[serde(default, skip_serializing_if = "Option::is_none")] version_id_exact: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + version_state: Option, } impl PersistedTierDeleteJournalEntry { - fn from_jentry(je: &Jentry) -> Self { - Self { - version: if je.version_id_exact { - TIER_DELETE_JOURNAL_EXACT_VERSION - } else if je.backend_identity.is_some() { + fn from_jentry(je: &Jentry) -> Result { + validate_version_state(je.version_state, &je.version_id, je.version_id_exact)?; + let legacy_unknown = je.version_state == rustfs_filemeta::TransitionVersionState::Unknown; + let version = if legacy_unknown { + if je.backend_identity.is_some() { TIER_DELETE_JOURNAL_VERSION } else { 1 - }, + } + } else { + if je.backend_identity.is_none() { + return Err(Error::other("new tier delete journal entry is missing its backend identity")); + } + TIER_DELETE_JOURNAL_STATE_VERSION + }; + Ok(Self { + version, obj_name: je.obj_name.clone(), version_id: je.version_id.clone(), tier_name: je.tier_name.clone(), backend_identity: je.backend_identity, version_id_exact: je.version_id_exact.then_some(true), - } + version_state: (!legacy_unknown).then_some(je.version_state), + }) } fn into_jentry(self) -> Result { @@ -84,19 +96,23 @@ impl PersistedTierDeleteJournalEntry { if self.obj_name.is_empty() || self.tier_name.is_empty() { return Err(Error::other("tier delete journal entry is incomplete")); } - if self.version != TIER_DELETE_JOURNAL_EXACT_VERSION && self.version_id_exact.unwrap_or(false) { + if self.version != TIER_DELETE_JOURNAL_EXACT_VERSION + && self.version != TIER_DELETE_JOURNAL_STATE_VERSION + && self.version_id_exact.unwrap_or(false) + { return Err(Error::other( "legacy tier delete journal entry has an unsupported exact version constraint", )); } - let (backend_identity, version_id_exact) = match self.version { - 1 => (None, false), + let (backend_identity, version_id_exact, version_state) = match self.version { + 1 => (None, false, rustfs_filemeta::TransitionVersionState::Unknown), TIER_DELETE_JOURNAL_VERSION => ( Some( self.backend_identity .ok_or_else(|| Error::other("tier delete journal v2 entry is missing its backend identity"))?, ), false, + rustfs_filemeta::TransitionVersionState::Unknown, ), TIER_DELETE_JOURNAL_EXACT_VERSION => { if self.version_id.is_empty() || self.version_id_exact != Some(true) { @@ -108,6 +124,22 @@ impl PersistedTierDeleteJournalEntry { .ok_or_else(|| Error::other("tier delete journal v3 entry is missing its backend identity"))?, ), true, + rustfs_filemeta::TransitionVersionState::Exact, + ) + } + TIER_DELETE_JOURNAL_STATE_VERSION => { + let state = self + .version_state + .ok_or_else(|| Error::other("tier delete journal v4 entry is missing its version state"))?; + let exact = self.version_id_exact.unwrap_or(false); + validate_version_state(state, &self.version_id, exact)?; + ( + Some( + self.backend_identity + .ok_or_else(|| Error::other("tier delete journal v4 entry is missing its backend identity"))?, + ), + exact, + state, ) } version => return Err(Error::other(format!("unsupported tier delete journal version {version}"))), @@ -118,10 +150,30 @@ impl PersistedTierDeleteJournalEntry { tier_name: self.tier_name, backend_identity, version_id_exact, + version_state, }) } } +fn validate_version_state( + state: rustfs_filemeta::TransitionVersionState, + version_id: &str, + version_id_exact: bool, +) -> Result<()> { + use rustfs_filemeta::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull, Unknown}; + + let valid = match state { + Unknown => !version_id_exact, + KnownDisabled => version_id.is_empty() && !version_id_exact, + SuspendedNull => version_id == "null" && version_id_exact, + Exact => !version_id.is_empty() && version_id != "null" && version_id_exact, + }; + if !valid { + return Err(Error::other("tier delete journal version state conflicts with its version id")); + } + Ok(()) +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct TierDeleteJournalRecoveryStats { pub scanned: usize, @@ -159,7 +211,7 @@ pub(crate) fn decode_tier_delete_journal_entry(data: &[u8]) -> Result { } pub(crate) fn encode_tier_delete_journal_entry(je: &Jentry) -> Result> { - serde_json::to_vec(&PersistedTierDeleteJournalEntry::from_jentry(je)) + serde_json::to_vec(&PersistedTierDeleteJournalEntry::from_jentry(je)?) .map_err(|err| Error::other(format!("encode tier delete journal failed: {err}"))) } @@ -209,6 +261,12 @@ where } pub async fn process_tier_delete_journal_entry(api: Arc, je: &Jentry) -> std::io::Result<()> { + if je.version_state == rustfs_filemeta::TransitionVersionState::Unknown { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "tier delete journal remote version state is unknown", + )); + } let backend_identity = je .backend_identity .ok_or_else(|| std::io::Error::other("legacy tier delete journal has no durable backend identity"))?; @@ -406,8 +464,9 @@ where #[cfg(test)] mod tests { use super::{ - TIER_DELETE_JOURNAL_EXACT_VERSION, await_tier_delete_journal_recovery, decode_tier_delete_journal_entry, - encode_tier_delete_journal_entry, record_tier_delete_journal_backend_identity, tier_delete_journal_object_name, + TIER_DELETE_JOURNAL_EXACT_VERSION, TIER_DELETE_JOURNAL_STATE_VERSION, await_tier_delete_journal_recovery, + decode_tier_delete_journal_entry, encode_tier_delete_journal_entry, record_tier_delete_journal_backend_identity, + tier_delete_journal_object_name, }; use crate::bucket::lifecycle::tier_sweeper::Jentry; use crate::error::Result; @@ -420,7 +479,8 @@ mod tests { version_id: "remote-version".to_string(), tier_name: "WARM".to_string(), backend_identity: Some([7; 32]), - version_id_exact: false, + version_id_exact: true, + version_state: rustfs_filemeta::TransitionVersionState::Exact, } } @@ -436,6 +496,7 @@ mod tests { assert_eq!(decoded.tier_name, je.tier_name); assert_eq!(decoded.backend_identity, je.backend_identity); assert_eq!(decoded.version_id_exact, je.version_id_exact); + assert_eq!(decoded.version_state, je.version_state); } #[test] @@ -450,7 +511,7 @@ mod tests { let persisted: serde_json::Value = serde_json::from_slice(&encoded).expect("exact journal JSON should decode"); let decoded = decode_tier_delete_journal_entry(&encoded).expect("exact journal entry should decode"); - assert_eq!(persisted["version"], TIER_DELETE_JOURNAL_EXACT_VERSION); + assert_eq!(persisted["version"], TIER_DELETE_JOURNAL_STATE_VERSION); assert_eq!(persisted["version_id_exact"], true); assert!(decoded.version_id_exact); assert_ne!(tier_delete_journal_object_name(&exact), tier_delete_journal_object_name(&normalized)); @@ -513,6 +574,46 @@ mod tests { } } + #[test] + fn tier_delete_journal_rejects_conflicting_v4_version_states() { + let identity = vec![7_u8; 32]; + let invalid = [ + ("known-disabled", "unexpected", false), + ("suspended-null", "", true), + ("suspended-null", "null", false), + ("exact", "", true), + ("exact", "null", true), + ("exact", "version", false), + ("unknown", "version", true), + ]; + + for (state, version_id, exact) in invalid { + let persisted = serde_json::json!({ + "version": TIER_DELETE_JOURNAL_STATE_VERSION, + "obj_name": "remote/object", + "version_id": version_id, + "tier_name": "WARM", + "backend_identity": identity, + "version_id_exact": exact.then_some(true), + "version_state": state, + }); + let encoded = serde_json::to_vec(&persisted).expect("invalid journal fixture should encode"); + decode_tier_delete_journal_entry(&encoded).expect_err("conflicting v4 version state must fail closed"); + } + } + + #[test] + fn legacy_journals_decode_with_unknown_version_state() { + let v1 = br#"{"version":1,"obj_name":"remote/object","version_id":"opaque","tier_name":"WARM"}"#; + let v2 = br#"{"version":2,"obj_name":"remote/object","version_id":"opaque","tier_name":"WARM","backend_identity":[7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7,7]}"#; + + for payload in [v1.as_slice(), v2.as_slice()] { + let decoded = decode_tier_delete_journal_entry(payload).expect("legacy journal should decode"); + assert_eq!(decoded.version_state, rustfs_filemeta::TransitionVersionState::Unknown); + assert!(!decoded.version_id_exact); + } + } + #[test] fn tier_delete_journal_path_is_stable_and_sanitized() { let je = journal_entry(); @@ -530,6 +631,8 @@ mod tests { fn tier_delete_journal_paths_separate_legacy_and_backend_identities() { let mut legacy = journal_entry(); legacy.backend_identity = None; + legacy.version_id_exact = false; + legacy.version_state = rustfs_filemeta::TransitionVersionState::Unknown; let mut backend_a = journal_entry(); backend_a.backend_identity = Some([1; 32]); let mut backend_b = journal_entry(); @@ -575,6 +678,8 @@ mod tests { fn tier_delete_journal_without_transition_identity_stays_legacy() { let mut je = journal_entry(); je.backend_identity = None; + je.version_id_exact = false; + je.version_state = rustfs_filemeta::TransitionVersionState::Unknown; let encoded = encode_tier_delete_journal_entry(&je).expect("legacy journal should remain encodable"); let persisted: serde_json::Value = serde_json::from_slice(&encoded).expect("journal JSON should decode"); diff --git a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs index 67ef6d15b..c5fbc29a9 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs @@ -185,6 +185,7 @@ struct ObjSweeper { transition_status: String, transition_tier: String, transition_version_id: String, + transition_version_state: rustfs_filemeta::TransitionVersionState, remote_object: String, } @@ -231,7 +232,9 @@ impl ObjSweeper { } pub fn should_remove_remote_object(&self) -> Option { - if self.transition_status != lifecycle::TRANSITION_COMPLETE { + if self.transition_status != lifecycle::TRANSITION_COMPLETE + || self.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown + { return None; } @@ -250,6 +253,7 @@ impl ObjSweeper { tier_name: self.transition_tier.clone(), backend_identity: None, version_id_exact: false, + version_state: self.transition_version_state, }); } None @@ -286,6 +290,7 @@ pub struct Jentry { pub(crate) tier_name: String, pub(crate) backend_identity: Option, pub(crate) version_id_exact: bool, + pub(crate) version_state: rustfs_filemeta::TransitionVersionState, } impl ExpiryOp for Jentry { @@ -486,6 +491,7 @@ pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject tier_name: transitioned.tier.clone(), backend_identity: None, version_id_exact: false, + version_state: rustfs_filemeta::TransitionVersionState::Unknown, }) } diff --git a/crates/ecstore/src/config/com.rs b/crates/ecstore/src/config/com.rs index 0f3eb9ab5..539464173 100644 --- a/crates/ecstore/src/config/com.rs +++ b/crates/ecstore/src/config/com.rs @@ -2543,6 +2543,7 @@ mod tests { data_dir: None, delete_marker: false, transitioned_object: Default::default(), + transition_version_state: Default::default(), restore_ongoing: false, restore_expires: None, user_tags: Arc::new(String::new()), diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index 33f53347b..dfd403e18 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -180,6 +180,7 @@ pub struct ObjectInfo { pub data_dir: Option, pub delete_marker: bool, pub transitioned_object: TransitionedObject, + pub transition_version_state: rustfs_filemeta::TransitionVersionState, pub restore_ongoing: bool, pub restore_expires: Option, pub user_tags: Arc, @@ -220,6 +221,7 @@ impl Clone for ObjectInfo { data_dir: self.data_dir, delete_marker: self.delete_marker, transitioned_object: self.transitioned_object.clone(), + transition_version_state: self.transition_version_state, restore_ongoing: self.restore_ongoing, restore_expires: self.restore_expires, user_tags: self.user_tags.clone(), @@ -537,6 +539,7 @@ impl ObjectInfo { inlined, user_defined: Arc::new(metadata), transitioned_object, + transition_version_state: fi.transition_version_state, checksum: fi.checksum.clone(), storage_class, restore_ongoing, diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index c501d52ec..5ead6a9d9 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -9274,7 +9274,8 @@ mod tests { version_id: "v1".to_string(), tier_name: "COLD-A".to_string(), backend_identity: Some(current_identity), - version_id_exact: false, + version_id_exact: true, + version_state: rustfs_filemeta::TransitionVersionState::Exact, }; journal_store .insert_config_object( diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index ceaff1996..c10020177 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -545,6 +545,7 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo && left.transition_tier == right.transition_tier && left.transition_version_id == right.transition_version_id && left.transition_version == right.transition_version + && left.transition_version_state == right.transition_version_state && left.expire_restored == right.expire_restored && left.size == right.size && left.mod_time == right.mod_time diff --git a/crates/ecstore/src/set_disk/metadata.rs b/crates/ecstore/src/set_disk/metadata.rs index 419055f7a..7eb4c57ad 100644 --- a/crates/ecstore/src/set_disk/metadata.rs +++ b/crates/ecstore/src/set_disk/metadata.rs @@ -579,6 +579,12 @@ impl SetDisks { Self::update_hash_str(hasher, &meta.transitioned_objname); Self::update_hash_optional_uuid(hasher, meta.transition_version_id); Self::update_hash_optional_str(hasher, meta.transition_version.as_deref()); + hasher.update([match meta.transition_version_state { + rustfs_filemeta::TransitionVersionState::Unknown => 0, + rustfs_filemeta::TransitionVersionState::KnownDisabled => 1, + rustfs_filemeta::TransitionVersionState::SuspendedNull => 2, + rustfs_filemeta::TransitionVersionState::Exact => 3, + }]); Self::update_hash_optional_u32(hasher, meta.mode); Self::update_hash_optional_u64(hasher, meta.written_by_version); diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 3adf2ffa8..4d554298e 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -1692,6 +1692,13 @@ async fn cleanup_rejected_transition_upload_durably( tier_name: lease.tier_name().to_string(), backend_identity: Some(lease.backend_identity()), version_id_exact, + version_state: if !version_id_exact { + rustfs_filemeta::TransitionVersionState::KnownDisabled + } else if cleanup_version == "null" { + rustfs_filemeta::TransitionVersionState::SuspendedNull + } else { + rustfs_filemeta::TransitionVersionState::Exact + }, }; let journal_error = if let Some(api) = api.as_ref() { @@ -2107,12 +2114,28 @@ async fn pause_transition_commit(bucket: &str, object: &str, pause: TransitionCo } } -fn parse_transition_version_id(remote_version: &str) -> Option { - if remote_version.is_empty() || Uuid::parse_str(remote_version).is_ok_and(|version_id| version_id.is_nil()) { - None - } else { - Some(remote_version.to_string()) +fn persisted_transition_version( + remote_version: &str, +) -> std::io::Result<(Option, rustfs_filemeta::TransitionVersionState)> { + if remote_version.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "a missing remote tier object version remains unknown until the cluster capability gate is active", + )); } + let version_id = Uuid::parse_str(remote_version).map_err(|_| { + std::io::Error::new( + std::io::ErrorKind::Unsupported, + "opaque remote tier versions require the cluster capability gate", + ) + })?; + if version_id.is_nil() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier returned a nil object version ID", + )); + } + Ok((Some(remote_version.to_string()), rustfs_filemeta::TransitionVersionState::Exact)) } #[cfg(test)] @@ -2274,13 +2297,14 @@ mod transition_upload_completion_tests { #[cfg(test)] mod transition_version_id_tests { - use super::{TransitionUploadCandidate, parse_transition_version_id}; + use super::{TransitionUploadCandidate, persisted_transition_version}; + use rustfs_filemeta::TransitionVersionState; use uuid::Uuid; #[test] fn normalizes_persisted_unversioned_ids_and_preserves_put_constraints() { - assert_eq!(parse_transition_version_id(""), None); - assert_eq!(parse_transition_version_id(&Uuid::nil().to_string()), None); + assert!(persisted_transition_version("").is_err()); + assert!(persisted_transition_version(&Uuid::nil().to_string()).is_err()); let nil_put_response = Uuid::nil().to_string(); let nil_candidate = TransitionUploadCandidate::from_put_response(nil_put_response.clone()); assert_eq!(nil_candidate.cleanup_version(), nil_put_response); @@ -2292,13 +2316,14 @@ mod transition_version_id_tests { } #[test] - fn preserves_uuid_and_opaque_remote_ids() { + fn preserves_uuid_and_gates_opaque_remote_ids() { let version_id = Uuid::new_v4(); - assert_eq!(parse_transition_version_id(&version_id.to_string()), Some(version_id.to_string())); assert_eq!( - parse_transition_version_id("opaque-version-token"), - Some("opaque-version-token".to_string()) + persisted_transition_version(&version_id.to_string()).expect("UUID remote version"), + (Some(version_id.to_string()), TransitionVersionState::Exact) ); + assert!(persisted_transition_version("null").is_err()); + assert!(persisted_transition_version("opaque-version-token").is_err()); assert_eq!( TransitionUploadCandidate::from_put_response(version_id.to_string()).cleanup_version(), version_id.to_string() @@ -3552,7 +3577,20 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await; return Err(err); } - let transition_version_id = parse_transition_version_id(candidate.remote_version()); + let (transition_version_id, transition_version_state) = match persisted_transition_version(candidate.remote_version()) { + Ok(version) => version, + Err(err) => { + let cleanup_api = transition_cleanup_store(&self.ctx).await; + if let Err(cleanup_err) = upload_cleanup.cleanup_rejected_upload(cleanup_api).await { + return Err(StorageError::Io(std::io::Error::other(format!( + "{err}; rejected remote upload cleanup failed: {cleanup_err}" + )))); + } + delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object) + .await; + return Err(err.into()); + } + }; let mut commit_opts = opts.clone(); commit_opts.no_lock = true; @@ -3614,6 +3652,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { .as_deref() .and_then(|version_id| Uuid::parse_str(version_id).ok()); current_fi.transition_version = transition_version_id; + current_fi.transition_version_state = transition_version_state; rustfs_utils::http::metadata_compat::insert_str( &mut current_fi.metadata, rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, @@ -4429,6 +4468,34 @@ mod transition_commit_failure_tests { metadata } + #[tokio::test] + async fn rejected_unsupported_remote_versions_are_cleaned_up() { + for remote_version in ["", "null", "opaque-version-token"] { + let manager = TierConfigMgr::new(); + let backend = register_mock_tier(&manager, "WARM").await; + let lease = TierConfigMgr::acquire_operation_lease(&manager, "WARM") + .await + .expect("mock tier lease should be available"); + let candidate = TransitionUploadCandidate::from_put_response(remote_version.to_string()); + + persisted_transition_version(candidate.remote_version()).expect_err("unsupported writer version must fail closed"); + cleanup_rejected_transition_upload_durably( + &lease, + "remote/object", + candidate.cleanup_version(), + candidate.cleanup_version_is_exact(), + None, + ) + .await + .expect("rejected remote upload must be cleaned up"); + + assert_eq!( + backend.remove_versions().await, + vec![("remote/object".to_string(), candidate.cleanup_version().to_string())] + ); + } + } + #[tokio::test] #[serial_test::serial] async fn local_commit_failure_returns_error_and_preserves_remote_candidate() { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 5c410839e..801ccd411 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -1311,14 +1311,16 @@ mod tests { version_id: "version-a".to_string(), tier_name: tier_a.to_string(), backend_identity: Some(identity_a), - version_id_exact: false, + version_id_exact: true, + version_state: rustfs_filemeta::TransitionVersionState::Exact, }; let entry_b = Jentry { obj_name: "remote-b".to_string(), version_id: "version-b".to_string(), tier_name: tier_b.to_string(), backend_identity: Some(identity_b), - version_id_exact: false, + version_id_exact: true, + version_state: rustfs_filemeta::TransitionVersionState::Exact, }; let remove_a = backend_a.arm_failing_remove_barrier().await; persist_tier_delete_journal_entry(store_a.clone(), &entry_a) diff --git a/crates/filemeta/src/fileinfo.rs b/crates/filemeta/src/fileinfo.rs index 012adccb2..bf2841e3a 100644 --- a/crates/filemeta/src/fileinfo.rs +++ b/crates/filemeta/src/fileinfo.rs @@ -219,6 +219,16 @@ impl ErasureInfo { } // #[derive(Debug, Clone)] +#[derive(Serialize, Deserialize, Debug, PartialEq, Eq, Clone, Copy, Default)] +#[serde(rename_all = "kebab-case")] +pub enum TransitionVersionState { + #[default] + Unknown, + KnownDisabled, + SuspendedNull, + Exact, +} + #[derive(Serialize, Deserialize, Debug, PartialEq, Clone, Default)] pub struct FileInfo { pub volume: String, @@ -232,6 +242,8 @@ pub struct FileInfo { pub transition_version_id: Option, #[serde(default)] pub transition_version: Option, + #[serde(default)] + pub transition_version_state: TransitionVersionState, pub expire_restored: bool, pub data_dir: Option, pub mod_time: Option, @@ -544,6 +556,20 @@ impl FileInfo { { return Err(Error::FileCorrupt); } + let transition_state_valid = match self.transition_version_state { + TransitionVersionState::Unknown => true, + TransitionVersionState::KnownDisabled => self.transition_version.is_none() && self.transition_version_id.is_none(), + TransitionVersionState::SuspendedNull => { + self.transition_version.as_deref() == Some("null") && self.transition_version_id.is_none() + } + TransitionVersionState::Exact => self + .transition_version + .as_deref() + .is_some_and(|version| version != "null" && !version.is_empty()), + }; + if !transition_state_valid { + return Err(Error::FileCorrupt); + } let erasure_layout = match mode { ValidationMode::RequireErasure => Some(self.validate_erasure_geometry()?), @@ -837,6 +863,7 @@ impl FileInfo { && self.transitioned_objname == other.transitioned_objname && self.transition_version_id == other.transition_version_id && self.transition_version == other.transition_version + && self.transition_version_state == other.transition_version_state } /// Check if metadata maps are equal @@ -1722,6 +1749,11 @@ mod tests { transition_tier, transition_version_id, transition_version: transition_version_id.map(|version_id| version_id.to_string()), + transition_version_state: if transition_version_id.is_some() { + TransitionVersionState::Exact + } else { + TransitionVersionState::Unknown + }, expire_restored, data_dir, mod_time, diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index fa3cf6c0d..3ecc4dfcc 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -26,13 +26,14 @@ use super::msgp_decode::{ PrependByteReader, prealloc_hint, read_exact_vec, read_nil_or_array_len, read_nil_or_map_len, skip_msgp_value, }; use super::*; -use crate::ChecksumInfo; +use crate::{ChecksumInfo, TransitionVersionState}; use rustfs_utils::HashAlgorithm; use rustfs_utils::http::{ RUSTFS_INTERNAL_PREFIX, SUFFIX_CRC, SUFFIX_FREE_VERSION, SUFFIX_INLINE_DATA, SUFFIX_PURGESTATUS, SUFFIX_TIER_FV_ID, SUFFIX_TIER_FV_MARKER, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID, - SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, contains_key_bytes, get_bytes, get_consistent_bytes, get_str, - has_internal_suffix, insert_bytes, is_internal_key, remove_bytes, strip_internal_prefix, + SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, contains_key_bytes, + get_bytes, get_consistent_bytes, get_str, has_internal_suffix, insert_bytes, is_internal_key, remove_bytes, + strip_internal_prefix, }; const MSGPACK_EXT8: u8 = 0xc7; @@ -255,28 +256,83 @@ fn parse_legacy_uuid_bytes(bytes: &[u8], field: &str) -> Result> { /// Legacy RustFS writes used 16 raw UUID bytes. New writes and MinIO-migrated /// records use the provider's exact UTF-8 version text. Empty, nil UUID, and /// malformed bytes are not usable remote versions. -fn transitioned_version_from_meta_sys(meta_sys: &HashMap>) -> Option { - let value = get_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID)?; +fn transitioned_version_from_meta_sys(meta_sys: &HashMap>) -> Result> { + if !contains_key_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID) { + return Ok(None); + } + let Some(value) = get_consistent_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID) else { + return Ok(None); + }; + let value = value.to_vec(); if value.is_empty() { - return None; + return Ok(None); } if let Ok(id) = Uuid::from_slice(&value) { - return (!id.is_nil()).then(|| id.to_string()); + return Ok((!id.is_nil()).then(|| id.to_string())); } - let value = String::from_utf8(value).ok()?; + let Ok(value) = String::from_utf8(value) else { + return Ok(None); + }; if value.is_empty() || value.len() > MAX_TRANSITION_VERSION_LEN || value.chars().any(char::is_control) || Uuid::parse_str(&value).is_ok_and(|id| id.is_nil()) { - None + Ok(None) } else { - Some(value) + Ok(Some(value)) + } +} + +fn transition_version_state_from_meta_sys( + meta_sys: &HashMap>, + version: Option<&str>, +) -> Result { + if !contains_key_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE) { + return Ok(TransitionVersionState::Unknown); + } + let value = get_consistent_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE).ok_or(Error::FileCorrupt)?; + let state = match value { + b"known-disabled" => TransitionVersionState::KnownDisabled, + b"suspended-null" => TransitionVersionState::SuspendedNull, + b"exact" => TransitionVersionState::Exact, + b"unknown" => TransitionVersionState::Unknown, + _ => return Err(Error::FileCorrupt), + }; + let valid = match state { + TransitionVersionState::Unknown | TransitionVersionState::KnownDisabled => version.is_none(), + TransitionVersionState::SuspendedNull => version == Some("null"), + TransitionVersionState::Exact => version.is_some_and(|value| value != "null"), + }; + valid.then_some(state).ok_or(Error::FileCorrupt) +} + +fn transition_version_state_bytes(state: TransitionVersionState) -> &'static [u8] { + match state { + TransitionVersionState::Unknown => b"unknown", + TransitionVersionState::KnownDisabled => b"known-disabled", + TransitionVersionState::SuspendedNull => b"suspended-null", + TransitionVersionState::Exact => b"exact", + } +} + +fn set_transition_version_state(meta_sys: &mut HashMap>, state: TransitionVersionState) { + if state == TransitionVersionState::Unknown { + remove_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE); + } else { + insert_bytes( + meta_sys, + SUFFIX_TRANSITIONED_VERSION_STATE, + transition_version_state_bytes(state).to_vec(), + ); } } fn legacy_transitioned_version_id_from_meta_sys(meta_sys: &HashMap>) -> Option { - transitioned_version_from_meta_sys(meta_sys).and_then(|value| Uuid::parse_str(&value).ok()) + transitioned_version_from_meta_sys(meta_sys) + .ok() + .flatten() + .and_then(|value| Uuid::parse_str(&value).ok()) } fn transitioned_version_bytes(fi: &FileInfo) -> Option> { @@ -2414,7 +2470,8 @@ impl MetaObject { let transitioned_objname = get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_OBJECTNAME) .map(|v| String::from_utf8_lossy(&v).to_string()) .unwrap_or_default(); - let transition_version = transitioned_version_from_meta_sys(&self.meta_sys); + let transition_version = transitioned_version_from_meta_sys(&self.meta_sys)?; + let transition_version_state = transition_version_state_from_meta_sys(&self.meta_sys, transition_version.as_deref())?; let transition_version_id = transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); let transition_tier = get_bytes(&self.meta_sys, SUFFIX_TRANSITION_TIER) .map(|v| String::from_utf8_lossy(&v).to_string()) @@ -2437,6 +2494,7 @@ impl MetaObject { transitioned_objname, transition_version_id, transition_version, + transition_version_state, transition_tier, ..Default::default() }) @@ -2452,6 +2510,7 @@ impl MetaObject { if let Some(transition_version) = transitioned_version_bytes(fi) { insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } + set_transition_version_state(&mut self.meta_sys, fi.transition_version_state); insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec()); if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) { insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes()); @@ -2515,6 +2574,7 @@ impl MetaObject { SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, + SUFFIX_TRANSITIONED_VERSION_STATE, ] { if let Some(v) = get_bytes(&self.meta_sys, suffix) { insert_bytes(&mut delete_marker.meta_sys, suffix, v); @@ -2579,6 +2639,9 @@ impl From for MetaObject { if let Some(transition_version) = transitioned_version_bytes(&value) { insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } + if !value.transition_status.is_empty() { + set_transition_version_state(&mut meta_sys, value.transition_version_state); + } if !value.transition_tier.is_empty() { insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec()); @@ -2720,8 +2783,11 @@ impl MetaDeleteMarker { .map(|v| String::from_utf8_lossy(&v).to_string()) .unwrap_or_default(); - fi.transition_version = transitioned_version_from_meta_sys(&self.meta_sys); + fi.transition_version = transitioned_version_from_meta_sys(&self.meta_sys).ok().flatten(); fi.transition_version_id = legacy_transitioned_version_id_from_meta_sys(&self.meta_sys); + fi.transition_version_state = + transition_version_state_from_meta_sys(&self.meta_sys, fi.transition_version.as_deref()) + .unwrap_or(TransitionVersionState::Unknown); } fi @@ -2877,6 +2943,9 @@ impl From for MetaDeleteMarker { if let Some(transition_version) = transitioned_version_bytes(&value) { insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } + if !value.transition_status.is_empty() || value.tier_free_version() { + set_transition_version_state(&mut meta_sys, value.transition_version_state); + } if !value.transition_tier.is_empty() { insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec()); } @@ -4125,6 +4194,7 @@ mod tests { .expect("into_fileinfo"); assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); + assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown); } #[test] @@ -4136,6 +4206,61 @@ mod tests { .expect("opaque transition version id must decode"); assert_eq!(fi.transition_version_id, None); assert_eq!(fi.transition_version.as_deref(), Some("opaque-generation-42")); + assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown); + } + + #[test] + fn meta_object_transition_version_state_exact_round_trips_dual_keys() { + let id = sample_version_id(); + let expected_version = id.to_string(); + let fi = FileInfo { + transition_status: "complete".to_string(), + transition_version: Some(expected_version.clone()), + transition_version_state: TransitionVersionState::Exact, + ..Default::default() + }; + + let object = MetaObject::from(fi); + assert_eq!( + object + .meta_sys + .get(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_STATE}")) + .map(Vec::as_slice), + Some(b"exact".as_slice()) + ); + assert_eq!( + object + .meta_sys + .get(&format!( + "{}{SUFFIX_TRANSITIONED_VERSION_STATE}", + rustfs_utils::http::MINIO_INTERNAL_PREFIX + )) + .map(Vec::as_slice), + Some(b"exact".as_slice()) + ); + assert_eq!( + legacy_transitioned_version_id_from_meta_sys(&object.meta_sys), + Some(id), + "UUID exact writes must remain readable by the legacy UUID consumer" + ); + let decoded = object.into_fileinfo("b", "k", false).expect("exact state should round trip"); + assert_eq!(decoded.transition_version_state, TransitionVersionState::Exact); + assert_eq!(decoded.transition_version.as_deref(), Some(expected_version.as_str())); + } + + #[test] + fn meta_object_transition_version_state_conflict_fails_closed() { + let mut sys = HashMap::new(); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, sample_version_id().as_bytes().to_vec()); + sys.insert(format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_STATE}"), b"exact".to_vec()); + sys.insert( + format!("{}{SUFFIX_TRANSITIONED_VERSION_STATE}", rustfs_utils::http::MINIO_INTERNAL_PREFIX), + b"known-disabled".to_vec(), + ); + + make_meta_object_with_sys(sys) + .into_fileinfo("b", "k", false) + .expect_err("conflicting state keys must fail closed"); } #[test] diff --git a/crates/utils/src/http/metadata_compat.rs b/crates/utils/src/http/metadata_compat.rs index 9dc6ba175..c95f5b3d7 100644 --- a/crates/utils/src/http/metadata_compat.rs +++ b/crates/utils/src/http/metadata_compat.rs @@ -37,6 +37,7 @@ pub const SUFFIX_CRC: &str = "crc"; pub const SUFFIX_TRANSITION_STATUS: &str = "transition-status"; pub const SUFFIX_TRANSITIONED_OBJECTNAME: &str = "transitioned-object"; 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_RESTORE_OPERATION_ID: &str = "restore-operation-id";