diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 9033fe275..5f270668f 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(()) @@ -4201,6 +4202,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, @@ -4210,6 +4228,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) => { @@ -5506,6 +5525,7 @@ mod tests { tier: tier.clone(), ..Default::default() }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Exact, ..Default::default() }; @@ -5569,6 +5589,7 @@ mod tests { tier, ..Default::default() }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Exact, ..Default::default() }; @@ -5591,6 +5612,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() { @@ -5640,6 +5725,7 @@ mod tests { oi.transitioned_object.tier = "WARM".to_string(); oi.transitioned_object.name = "remote/object".to_string(); oi.transitioned_object.version_id = "remote-version".to_string(); + oi.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact; let local_delete_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0)); let legacy_err = delete_free_version_remote_object_then(&oi, &manager, { @@ -5650,7 +5736,8 @@ mod tests { }) .await .expect_err("legacy free-version without identity must be retained"); - assert!(legacy_err.to_string().contains("no durable backend identity")); + assert_eq!(legacy_err.kind(), std::io::ErrorKind::Other); + assert_eq!(old_backend.remove_count().await, 0); assert_eq!(local_delete_calls.load(Ordering::Relaxed), 0); let mut invalid_metadata = HashMap::new(); @@ -5781,6 +5868,7 @@ mod tests { oi.transitioned_object.tier = "WARM".to_string(); oi.transitioned_object.name = "remote/object".to_string(); oi.transitioned_object.version_id = "remote-version".to_string(); + oi.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact; let err = match get_transitioned_object_reader_with_tier_manager( "bucket", @@ -5796,7 +5884,12 @@ mod tests { Ok(_) => panic!("identity-bound GET must reject a same-name tier rebind"), Err(err) => err, }; - assert!(err.to_string().contains("identity no longer matches")); + assert_eq!(err.kind(), std::io::ErrorKind::Other); + let admin_err = err + .get_ref() + .and_then(|source| source.downcast_ref::()) + .expect("identity mismatch should retain the typed tier error"); + assert_eq!(admin_err.code, crate::services::tier::tier::ERR_TIER_INVALID_CONFIG.code); assert_eq!(new_backend.get_count().await, 0); oi.user_defined = Arc::new(HashMap::new()); @@ -5902,7 +5995,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 @@ -6013,7 +6107,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 @@ -10081,6 +10176,87 @@ 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); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn journal_replay_deletes_confirmed_exact_provider_token() { + let (_disk_paths, ecstore) = setup_test_env().await; + let (backend, _) = register_recovery_mock_tier(&ecstore).await; + let lease = TierConfigMgr::acquire_operation_lease(&ecstore.tier_config_mgr(), "WARM") + .await + .expect("mock tier lease should be available"); + let identity = lease.backend_identity(); + backend + .set_put_remote_version(Some("provider-version-token".to_string())) + .await; + lease + .put( + "remote/object", + crate::client::transition_api::ReaderImpl::Body(bytes::Bytes::from_static(b"candidate")), + 9, + ) + .await + .expect("confirmed remote candidate should be seeded"); + backend.set_remove_failure(true); + backend.set_reject_non_empty_remote_versions(true); + let je = Jentry { + obj_name: "remote/object".to_string(), + version_id: "provider-version-token".to_string(), + tier_name: "WARM".to_string(), + backend_identity: Some(identity), + version_id_exact: true, + version_state: rustfs_filemeta::TransitionVersionState::Exact, + }; + + crate::set_disk::cleanup_rejected_transition_upload_durably( + &lease, + &je.obj_name, + &je.version_id, + true, + Some(ecstore.clone()), + ) + .await + .expect("failed immediate cleanup should remain durable in the journal"); + assert!(backend.contains(&je.obj_name).await); + + backend.set_remove_failure(false); + crate::bucket::lifecycle::tier_delete_journal::process_tier_delete_journal_entry(ecstore, &je) + .await + .expect("identity-bound exact journal must retry confirmed candidate cleanup"); + + assert!(!backend.contains(&je.obj_name).await); + assert_eq!(backend.exact_remove_count(), 2); + assert_eq!( + backend.remove_versions().await, + vec![("remote/object".to_string(), "provider-version-token".to_string())] + ); + } + async fn seed_recoverable_free_version( disk_paths: &[PathBuf], bucket: &str, @@ -10100,6 +10276,7 @@ mod tests { identity, ); } + let transition_version_id = Uuid::new_v4(); let mut metadata = FileMeta::new(); metadata .add_version(FileInfo { @@ -10108,7 +10285,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..a217d4978 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs @@ -20,7 +20,10 @@ use tokio_util::sync::CancellationToken; use tracing::{debug, warn}; use crate::bucket::lifecycle::config_boundary; -use crate::bucket::lifecycle::tier_sweeper::{Jentry, delete_object_from_remote_tier_idempotent_with_manager_and_identity}; +use crate::bucket::lifecycle::tier_sweeper::{ + Jentry, delete_confirmed_transition_candidate_exact_with_manager_and_identity, + delete_object_from_remote_tier_idempotent_with_manager_and_identity, +}; use crate::disk::RUSTFS_META_BUCKET; use crate::error::{Error, Result}; use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader}; @@ -42,6 +45,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 +59,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 +99,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 +127,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 +153,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 +214,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,18 +264,35 @@ 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"))?; - delete_object_from_remote_tier_idempotent_with_manager_and_identity( - &je.obj_name, - &je.version_id, - &je.tier_name, - backend_identity, - &api.tier_config_mgr(), - je.version_id_exact, - ) - .await?; + if je.version_id_exact { + delete_confirmed_transition_candidate_exact_with_manager_and_identity( + &je.obj_name, + &je.version_id, + &je.tier_name, + backend_identity, + &api.tier_config_mgr(), + ) + .await?; + } else { + delete_object_from_remote_tier_idempotent_with_manager_and_identity( + &je.obj_name, + &je.version_id, + &je.tier_name, + backend_identity, + &api.tier_config_mgr(), + false, + ) + .await?; + } remove_tier_delete_journal_entry(api, je).await } @@ -406,8 +478,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 +493,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 +510,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 +525,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 +588,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 +645,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 +692,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..c4ca3b809 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; } @@ -249,7 +252,11 @@ impl ObjSweeper { version_id: self.transition_version_id.clone(), tier_name: self.transition_tier.clone(), backend_identity: None, - version_id_exact: false, + version_id_exact: matches!( + self.transition_version_state, + rustfs_filemeta::TransitionVersionState::SuspendedNull | rustfs_filemeta::TransitionVersionState::Exact + ), + version_state: self.transition_version_state, }); } None @@ -286,6 +293,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 { @@ -330,7 +338,7 @@ async fn delete_object_from_remote_tier_raw_with_manager( let lease = TierConfigMgr::acquire_operation_lease(&tier_config_mgr, tier_name) .await .map_err(std::io::Error::other)?; - delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease, false).await + delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease, false, true).await } async fn delete_object_from_remote_tier_raw_with_lease( @@ -338,8 +346,11 @@ async fn delete_object_from_remote_tier_raw_with_lease( rv_id: &str, lease: &TierOperationLease, version_id_exact: bool, + validate_remote_version_id: bool, ) -> Result<(), std::io::Error> { - lease.validate_remote_version_id(rv_id)?; + if validate_remote_version_id { + lease.validate_remote_version_id(rv_id)?; + } if remote_delete_breaker_is_open(Instant::now()).await { metrics::counter!(METRIC_DELETE_REMOTE_BREAKER_TOTAL).increment(1); @@ -435,7 +446,53 @@ pub(crate) async fn delete_object_from_remote_tier_with_lease_idempotent( lease: &TierOperationLease, version_id_exact: bool, ) -> Result { - match delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, lease, version_id_exact).await { + delete_object_from_remote_tier_with_lease_idempotent_inner(obj_name, rv_id, lease, version_id_exact, true).await +} + +pub(crate) async fn delete_confirmed_transition_candidate_exact_with_lease_idempotent( + obj_name: &str, + rv_id: &str, + lease: &TierOperationLease, +) -> Result { + if rv_id.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "confirmed versioned transition candidate requires a non-empty remote version", + )); + } + #[cfg(test)] + if obj_name == "remote/empty-guard-probe" { + CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + } + delete_object_from_remote_tier_with_lease_idempotent_inner(obj_name, rv_id, lease, true, false).await +} + +#[cfg(test)] +static CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0); + +pub(crate) async fn delete_confirmed_transition_candidate_exact_with_manager_and_identity( + obj_name: &str, + rv_id: &str, + tier_name: &str, + backend_identity: TierDestinationId, + tier_config_mgr: &Arc>, +) -> Result { + let lease = TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, tier_name, backend_identity) + .await + .map_err(std::io::Error::other)?; + delete_confirmed_transition_candidate_exact_with_lease_idempotent(obj_name, rv_id, &lease).await +} + +async fn delete_object_from_remote_tier_with_lease_idempotent_inner( + obj_name: &str, + rv_id: &str, + lease: &TierOperationLease, + version_id_exact: bool, + validate_remote_version_id: bool, +) -> Result { + match delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, lease, version_id_exact, validate_remote_version_id) + .await + { Ok(()) => Ok(RemoteTierDeleteOutcome::Deleted), Err(err) if is_remote_tier_not_found_error(&err) => Ok(RemoteTierDeleteOutcome::AlreadyRemoved), Err(err) => { @@ -460,6 +517,7 @@ pub fn transitioned_delete_journal_entry( versioned: bool, suspended: bool, transitioned: &TransitionedObject, + transition_version_state: rustfs_filemeta::TransitionVersionState, ) -> Option { let sweeper = ObjSweeper { version_id, @@ -468,6 +526,7 @@ pub fn transitioned_delete_journal_entry( transition_status: transitioned.status.clone(), transition_tier: transitioned.tier.clone(), transition_version_id: transitioned.version_id.clone(), + transition_version_state, remote_object: transitioned.name.clone(), ..Default::default() }; @@ -475,8 +534,13 @@ pub fn transitioned_delete_journal_entry( sweeper.should_remove_remote_object() } -pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject) -> Option { - if transitioned.status != lifecycle::TRANSITION_COMPLETE { +pub fn transitioned_force_delete_journal_entry( + transitioned: &TransitionedObject, + transition_version_state: rustfs_filemeta::TransitionVersionState, +) -> Option { + if transitioned.status != lifecycle::TRANSITION_COMPLETE + || transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown + { return None; } @@ -485,7 +549,11 @@ pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject version_id: transitioned.version_id.clone(), tier_name: transitioned.tier.clone(), backend_identity: None, - version_id_exact: false, + version_id_exact: matches!( + transition_version_state, + rustfs_filemeta::TransitionVersionState::SuspendedNull | rustfs_filemeta::TransitionVersionState::Exact + ), + version_state: transition_version_state, }) } @@ -494,11 +562,14 @@ mod test { use crate::client::signer_error::invalid_utf8_header_error; use super::{ - ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED, RemoteDeleteBreaker, RemoteTierDeleteOutcome, + CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES, ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED, + RemoteDeleteBreaker, RemoteTierDeleteOutcome, delete_confirmed_transition_candidate_exact_with_manager_and_identity, delete_object_from_remote_tier_idempotent, delete_object_from_remote_tier_idempotent_with_manager_and_identity, - is_remote_tier_not_found_error, is_signer_header_error, set_remote_tier_delete_test_hook, - should_record_remote_delete_failure, + is_remote_tier_not_found_error, is_signer_header_error, lifecycle, set_remote_tier_delete_test_hook, + should_record_remote_delete_failure, transitioned_delete_journal_entry, transitioned_force_delete_journal_entry, }; + use crate::storage_api_contracts::lifecycle::TransitionedObject; + use rustfs_filemeta::TransitionVersionState; use std::io::{Error, ErrorKind}; use std::time::{Duration, Instant}; @@ -542,6 +613,43 @@ mod test { assert!(should_record_remote_delete_failure(&Error::other("NoSuchVersion"))); } + #[test] + fn transitioned_delete_journal_preserves_remote_version_state() { + let cases = [ + (TransitionVersionState::Unknown, "legacy-version", None), + (TransitionVersionState::KnownDisabled, "", Some(false)), + (TransitionVersionState::SuspendedNull, "null", Some(true)), + (TransitionVersionState::Exact, "opaque-version", Some(true)), + ]; + + for (state, version_id, expected_exact) in cases { + let transitioned = TransitionedObject { + name: "remote/object".to_string(), + version_id: version_id.to_string(), + tier: "WARM".to_string(), + status: lifecycle::TRANSITION_COMPLETE.to_string(), + ..Default::default() + }; + let regular = transitioned_delete_journal_entry(None, false, false, &transitioned, state); + let forced = transitioned_force_delete_journal_entry(&transitioned, state); + + match expected_exact { + Some(expected_exact) => { + let regular = regular.expect("known version state should produce a regular delete journal entry"); + assert_eq!(regular.version_state, state); + assert_eq!(regular.version_id_exact, expected_exact); + let forced = forced.expect("known version state should produce a forced delete journal entry"); + assert_eq!(forced.version_state, state); + assert_eq!(forced.version_id_exact, expected_exact); + } + None => { + assert!(regular.is_none()); + assert!(forced.is_none()); + } + } + } + } + #[tokio::test] #[serial_test::serial] async fn idempotent_remote_delete_treats_hooked_nosuchversion_as_already_removed() { @@ -664,6 +772,55 @@ mod test { assert_eq!(backend.remove_versions().await, vec![("remote/object".to_string(), String::new())]); } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial] + async fn confirmed_transition_cleanup_deletes_exact_provider_token() { + CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES.store(0, std::sync::atomic::Ordering::Relaxed); + let manager = crate::services::tier::tier::TierConfigMgr::new(); + let backend = crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await; + let lease = crate::services::tier::tier::TierConfigMgr::acquire_operation_lease(&manager, "WARM") + .await + .expect("test tier lease should be available"); + let identity = lease.backend_identity(); + drop(lease); + backend.set_reject_non_empty_remote_versions(true); + + let outcome = delete_confirmed_transition_candidate_exact_with_manager_and_identity( + "remote/object", + "provider-version-token", + "WARM", + identity, + &manager, + ) + .await + .expect("confirmed upload compensation should delete the exact provider token"); + + assert_eq!(outcome, RemoteTierDeleteOutcome::Deleted); + assert_eq!(backend.exact_remove_count(), 1); + assert_eq!( + backend.remove_versions().await, + vec![("remote/object".to_string(), "provider-version-token".to_string())] + ); + + let err = delete_confirmed_transition_candidate_exact_with_manager_and_identity( + "remote/empty-guard-probe", + "", + "WARM", + identity, + &manager, + ) + .await + .expect_err("confirmed versioned cleanup must reject an empty token"); + assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput); + assert_eq!(backend.remove_count().await, 1); + assert_eq!( + CONFIRMED_TRANSITION_EMPTY_GUARD_DISPATCHES.load(std::sync::atomic::Ordering::Relaxed), + 0, + "empty remote versions must be rejected before exact cleanup dispatch" + ); + } + #[test] fn breaker_opens_at_threshold_and_recovers_after_window() { let mut breaker = RemoteDeleteBreaker::new(3, Duration::from_secs(30)); diff --git a/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs b/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs index 823ae1d2a..e3f3036e8 100644 --- a/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs +++ b/crates/ecstore/src/bucket/lifecycle/transition_transaction.rs @@ -22,7 +22,10 @@ use uuid::Uuid; use crate::bucket::lifecycle::config_boundary; use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE; -use crate::bucket::lifecycle::tier_sweeper::delete_object_from_remote_tier_idempotent_with_manager_and_identity; +use crate::bucket::lifecycle::tier_sweeper::{ + delete_confirmed_transition_candidate_exact_with_manager_and_identity, + delete_object_from_remote_tier_idempotent_with_manager_and_identity, +}; use crate::disk::RUSTFS_META_BUCKET; use crate::error::{Error, Result as EcstoreResult}; use crate::object_api::ObjectOptions; @@ -708,6 +711,21 @@ async fn recover_unknown_upload_outcome( TransitionCandidateProbe::UnversionedPresent => { cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::unversioned()).await } + TransitionCandidateProbe::VersionedPresent(version_id) + if Uuid::parse_str(&version_id).is_ok_and(|version_id| version_id.is_nil()) => + { + delete_confirmed_transition_candidate_exact_with_manager_and_identity( + &transaction.remote_object, + &version_id, + &transaction.tier_name, + transaction.backend_fingerprint, + &api.tier_config_mgr(), + ) + .await + .map_err(Error::other)?; + delete_transition_transaction_record(api, transaction.transaction_id).await?; + Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted) + } TransitionCandidateProbe::VersionedPresent(version_id) => { cleanup_recovered_unknown_upload_candidate(api, transaction, TransitionRemoteVersion::versioned(version_id)).await } 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 c9a0e560b..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(), @@ -464,11 +466,11 @@ impl ObjectInfo { let transitioned_object = TransitionedObject { name: fi.transitioned_objname.clone(), - version_id: if let Some(transition_version_id) = fi.transition_version_id { - transition_version_id.to_string() - } else { - "".to_string() - }, + version_id: fi + .transition_version + .clone() + .or_else(|| fi.transition_version_id.map(|version_id| version_id.to_string())) + .unwrap_or_default(), status: fi.transition_status.clone(), free_version: fi.tier_free_version(), tier: fi.transition_tier.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/test_util.rs b/crates/ecstore/src/services/tier/test_util.rs index 1a2b37d26..751f2e92b 100644 --- a/crates/ecstore/src/services/tier/test_util.rs +++ b/crates/ecstore/src/services/tier/test_util.rs @@ -931,7 +931,10 @@ pub async fn read_transition_meta(disk_path: &Path, bucket: &str, object: &str) status: fi.transition_status.clone(), tier: fi.transition_tier.clone(), remote_object: fi.transitioned_objname.clone(), - remote_version_id: fi.transition_version_id.map(|id| id.to_string()), + remote_version_id: fi + .transition_version + .clone() + .or_else(|| fi.transition_version_id.map(|id| id.to_string())), free_version_count, }) } 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 4972274f1..b8fd83ced 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -547,6 +547,8 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo && left.transitioned_objname == right.transitioned_objname && 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 136a6cf7e..7eb4c57ad 100644 --- a/crates/ecstore/src/set_disk/metadata.rs +++ b/crates/ecstore/src/set_disk/metadata.rs @@ -578,6 +578,13 @@ impl SetDisks { Self::update_hash_str(hasher, &meta.transition_tier); 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/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 692571243..16accb2dd 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -692,6 +692,8 @@ pub(crate) use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause}; #[cfg(feature = "test-util")] pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier; pub(crate) use ops::object::body_cache_plaintext_len; +#[cfg(test)] +pub(crate) use ops::object::cleanup_rejected_transition_upload_durably; mod read; mod replication; pub(crate) mod shard_source; diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 7725c2292..41775da3e 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -25,7 +25,10 @@ use crate::set_disk::read::GetObjectDownstreamWriter; use crate::bucket::lifecycle::{ tier_delete_journal::{persist_tier_delete_journal_entry, remove_tier_delete_journal_entry}, - tier_sweeper::{Jentry, RemoteTierDeleteOutcome, delete_object_from_remote_tier_with_lease_idempotent}, + tier_sweeper::{ + Jentry, RemoteTierDeleteOutcome, delete_confirmed_transition_candidate_exact_with_lease_idempotent, + delete_object_from_remote_tier_with_lease_idempotent, + }, transition_transaction::{ TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction, TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record, @@ -1663,7 +1666,11 @@ pub(crate) async fn cleanup_uncommitted_transition_upload( cleanup_version: &str, version_id_exact: bool, ) -> std::io::Result { - delete_object_from_remote_tier_with_lease_idempotent(object, cleanup_version, lease, version_id_exact).await + if version_id_exact { + delete_confirmed_transition_candidate_exact_with_lease_idempotent(object, cleanup_version, lease).await + } else { + delete_object_from_remote_tier_with_lease_idempotent(object, cleanup_version, lease, false).await + } } fn log_transition_upload_cleanup_failure(lease: &TierOperationLease, object: &str, cleanup_version: &str, err: &std::io::Error) { @@ -1800,7 +1807,7 @@ impl Drop for TransitionUploadCleanup { } } -async fn cleanup_rejected_transition_upload_durably( +pub(crate) async fn cleanup_rejected_transition_upload_durably( lease: &TierOperationLease, object: &str, cleanup_version: &str, @@ -1813,6 +1820,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() { @@ -1972,12 +1986,83 @@ async fn advance_and_save_transition_transaction( next: TransitionTransactionState, remote_version: Option, ) -> Result<()> { + #[cfg(test)] + record_transition_uploaded_save_attempt(transaction, next); transaction .advance(transaction.fence(), next, remote_version) .map_err(Error::other)?; save_transition_transaction_if_available(api, transaction).await } +#[cfg(test)] +struct TransitionUploadedSaveProbeState { + bucket: String, + object: String, + attempts: std::sync::atomic::AtomicUsize, +} + +#[cfg(test)] +struct TransitionUploadedSaveProbe { + state: Arc, +} + +#[cfg(test)] +static TRANSITION_UPLOADED_SAVE_PROBE: std::sync::OnceLock>>> = + std::sync::OnceLock::new(); + +#[cfg(test)] +impl TransitionUploadedSaveProbe { + fn install(bucket: &str, object: &str) -> Self { + let state = Arc::new(TransitionUploadedSaveProbeState { + bucket: bucket.to_string(), + object: object.to_string(), + attempts: std::sync::atomic::AtomicUsize::new(0), + }); + let mut slot = TRANSITION_UPLOADED_SAVE_PROBE + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("transition uploaded-save probe mutex should not poison"); + assert!(slot.is_none(), "transition uploaded-save probe must be installed by one test at a time"); + *slot = Some(Arc::clone(&state)); + drop(slot); + Self { state } + } + + fn attempts(&self) -> usize { + self.state.attempts.load(std::sync::atomic::Ordering::Acquire) + } +} + +#[cfg(test)] +impl Drop for TransitionUploadedSaveProbe { + fn drop(&mut self) { + let mut slot = TRANSITION_UPLOADED_SAVE_PROBE + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("transition uploaded-save probe mutex should not poison"); + if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *slot = None; + } + } +} + +#[cfg(test)] +fn record_transition_uploaded_save_attempt(transaction: &TransitionTransaction, next: TransitionTransactionState) { + if next != TransitionTransactionState::Uploaded { + return; + } + let state = TRANSITION_UPLOADED_SAVE_PROBE + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("transition uploaded-save probe mutex should not poison") + .as_ref() + .filter(|state| state.bucket == transaction.source.bucket && state.object == transaction.source.object) + .cloned(); + if let Some(state) = state { + state.attempts.fetch_add(1, std::sync::atomic::Ordering::AcqRel); + } +} + async fn delete_transition_transaction_if_available(api: Option<&Arc>, transaction_id: Uuid) -> Result<()> { if let Some(api) = api { return delete_transition_transaction_record(api.clone(), transaction_id).await; @@ -2228,11 +2313,25 @@ async fn pause_transition_commit(bucket: &str, object: &str, pause: TransitionCo } } -fn parse_transition_version_id(remote_version: &str) -> std::result::Result, uuid::Error> { +fn persisted_transition_version( + remote_version: &str, +) -> std::io::Result<(Option, rustfs_filemeta::TransitionVersionState)> { if remote_version.is_empty() { - return Ok(None); + return Ok((None, rustfs_filemeta::TransitionVersionState::KnownDisabled)); } - Uuid::parse_str(remote_version).map(|version_id| (!version_id.is_nil()).then_some(version_id)) + 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)] @@ -2470,16 +2569,17 @@ 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("").expect("empty remote version should be valid"), None); assert_eq!( - parse_transition_version_id(&Uuid::nil().to_string()).expect("nil remote version should be valid"), - None + persisted_transition_version("").expect("empty remote version identifies an unversioned tier"), + (None, TransitionVersionState::KnownDisabled) ); + 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); @@ -2491,12 +2591,14 @@ mod transition_version_id_tests { } #[test] - fn preserves_valid_remote_id_and_rejects_invalid_text() { + fn preserves_uuid_and_gates_opaque_remote_ids() { let version_id = Uuid::new_v4(); assert_eq!( - parse_transition_version_id(&version_id.to_string()).expect("UUID remote version should be valid"), - Some(version_id) + 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() @@ -2505,7 +2607,6 @@ mod transition_version_id_tests { TransitionUploadCandidate::from_put_response("opaque-version-token".to_string()).cleanup_version(), "opaque-version-token" ); - assert!(parse_transition_version_id("not-a-uuid").is_err()); } } @@ -3775,6 +3876,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.into()); } + 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()); + } + }; if let Err(err) = advance_and_save_transition_transaction( transaction_api.as_ref(), &mut transaction, @@ -3792,16 +3907,6 @@ 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 = match parse_transition_version_id(candidate.remote_version()) { - Ok(version_id) => version_id, - Err(err) => { - if upload_cleanup.cleanup().await.is_ok() { - 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; @@ -3859,7 +3964,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { current_fi.transition_status = TRANSITION_COMPLETE.to_string(); current_fi.transitioned_objname = dest_obj; current_fi.transition_tier = opts.transition.tier.clone(); - current_fi.transition_version_id = transition_version_id; + current_fi.transition_version_id = transition_version_id + .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, @@ -4668,6 +4777,33 @@ mod transition_commit_failure_tests { } #[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())] + ); + } + } + #[serial_test::serial(restore_multipart_failure_point)] async fn multipart_restore_aborts_every_post_create_failure() { let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; @@ -6615,6 +6751,45 @@ mod transition_upload_integrity_tests { ); } + #[tokio::test] + #[serial_test::serial] + async fn unversioned_remote_version_is_persisted_without_version_id() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "transition-unversioned-tier-bucket"; + let object = "object.bin"; + let payload = b"unversioned remote tier must commit without a version id".repeat(1024); + let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await; + let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await; + backend.set_put_remote_version(Some(String::new())).await; + let save_probe = TransitionUploadedSaveProbe::install(bucket, object); + + set_disks + .transition_object(bucket, object, &transition_options(&original, tier_name)) + .await + .expect("an unversioned remote version must commit"); + let (fi, _, _) = set_disks + .get_object_fileinfo( + bucket, + object, + &ObjectOptions { + no_lock: true, + metadata_cache_safe: false, + ..Default::default() + }, + true, + false, + ) + .await + .expect("committed unversioned transition metadata should be readable"); + assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version, None); + assert_eq!(fi.transition_version_state, rustfs_filemeta::TransitionVersionState::KnownDisabled); + assert_eq!(save_probe.attempts(), 1); + assert_eq!(backend.remove_count().await, 0); + assert_eq!(backend.object_count().await, 1); + } + #[tokio::test] #[serial_test::serial] async fn opaque_remote_version_is_cleaned_before_parse_failure() { @@ -6630,7 +6805,7 @@ mod transition_upload_integrity_tests { set_disks .transition_object(bucket, object, &transition_options(&original, tier_name)) .await - .expect_err("an unparseable remote version must fail closed"); + .expect_err("an opaque remote version must fail closed until the capability gate is active"); let removed_versions = backend.remove_versions().await; assert_eq!(removed_versions.len(), 1); assert_eq!(removed_versions[0].1, "opaque-version-token"); @@ -6638,6 +6813,38 @@ mod transition_upload_integrity_tests { assert_local_source_intact(&set_disks, bucket, object, &payload).await; } + #[tokio::test] + #[serial_test::serial] + async fn nil_remote_version_is_cleaned_exactly_before_transaction_persistence() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "transition-nil-version-bucket"; + let object = "object.bin"; + let payload = b"nil remote version must retain local data".repeat(1024); + let original = write_source(&set_disks, &disk_stores, bucket, object, &payload).await; + let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let remote_version = Uuid::nil().to_string(); + let backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await; + backend.set_put_remote_version(Some(remote_version.clone())).await; + let save_probe = TransitionUploadedSaveProbe::install(bucket, object); + + set_disks + .transition_object(bucket, object, &transition_options(&original, tier_name)) + .await + .expect_err("a nil remote version must fail closed before transaction persistence"); + let put_versions = backend.put_versions().await; + let removed_versions = backend.remove_versions().await; + assert_eq!(removed_versions, put_versions); + assert_eq!(removed_versions.len(), 1); + assert_eq!( + removed_versions.first().map(|(_, version)| version.as_str()), + Some(remote_version.as_str()) + ); + assert_eq!(save_probe.attempts(), 0, "nil remote version must be rejected before saving Uploaded"); + assert_eq!(backend.exact_remove_count(), 1); + assert_eq!(backend.object_count().await, 0); + assert_local_source_intact(&set_disks, bucket, object, &payload).await; + } + #[tokio::test] #[serial_test::serial] async fn authoritative_read_failure_after_upload_cleans_exact_candidate_and_preserves_source() { @@ -6985,23 +7192,36 @@ mod transition_source_identity_matrix_tests { let object = format!("identity-{index}.bin"); let payload = vec![u8::try_from(index + 1).expect("matrix index should fit u8"); 1024 * 1024]; let mut reader = PutObjReader::from_vec(payload); + let source_version_id = Uuid::new_v4(); + let source_opts = ObjectOptions { + version_id: Some(source_version_id.to_string()), + versioned: true, + ..Default::default() + }; let original = set_disks - .put_object(bucket, &object, &mut reader, &ObjectOptions::default()) + .put_object(bucket, &object, &mut reader, &source_opts) .await .expect("source object should be written"); let (source, _, _) = set_disks - .get_object_fileinfo(bucket, &object, &ObjectOptions::default(), true, false) + .get_object_fileinfo(bucket, &object, &source_opts, true, false) .await .expect("source metadata should resolve"); + assert_eq!(source.version_id, Some(source_version_id)); + assert_eq!( + transition_source_identity(bucket, &object, &source, &source_opts, &get_raw_etag(&source.metadata)) + .expect("persisted versioned source identity should build") + .version_mode, + TransitionSourceVersionMode::Versioned + ); let opts = ObjectOptions { no_lock: true, + versioned: true, transition: TransitionOptions { status: TRANSITION_PENDING.to_string(), tier: tier_name.clone(), etag: original.etag.clone().unwrap_or_default(), ..Default::default() }, - version_id: original.version_id.map(|version| version.to_string()), mod_time: original.mod_time, ..Default::default() }; @@ -7014,7 +7234,10 @@ mod transition_source_identity_matrix_tests { let mut changed = source.clone(); match field { - IdentityField::VersionId => changed.version_id = Some(Uuid::new_v4()), + IdentityField::VersionId => { + changed.version_id = Some(Uuid::new_v4()); + changed.fresh = true; + } IdentityField::DataDir => changed.data_dir = Some(Uuid::new_v4()), IdentityField::ModTime => { changed.mod_time = changed.mod_time.map(|value| value + time::Duration::nanoseconds(1)); @@ -7032,12 +7255,19 @@ mod transition_source_identity_matrix_tests { .await .expect("single-field metadata drift should be written"); } + let persisted_opts = ObjectOptions { + version_id: changed.version_id.map(|version_id| version_id.to_string()), + versioned: true, + ..Default::default() + }; + let (persisted, _, _) = set_disks + .get_object_fileinfo(bucket, &object, &persisted_opts, true, false) + .await + .expect("drifted source metadata should resolve"); put_barrier.release(); - transition - .await - .expect("transition task should not panic") - .expect_err("transition must reject a source whose identity changed after upload"); + let result = transition.await.expect("transition task should not panic"); + assert!(result.is_err(), "transition must reject {field:?} drift"); let expected_attempts = index + 1; assert_eq!(backend.put_count().await, expected_attempts); assert_eq!(backend.remove_count().await, expected_attempts); @@ -7048,26 +7278,26 @@ mod transition_source_identity_matrix_tests { ); match field { - IdentityField::VersionId => assert_ne!(source.version_id, changed.version_id), - IdentityField::DataDir => assert_ne!(source.data_dir, changed.data_dir), - IdentityField::ModTime => assert_ne!(source.mod_time, changed.mod_time), - IdentityField::Size => assert_ne!(source.size, changed.size), - IdentityField::Etag => assert_ne!(get_raw_etag(&source.metadata), get_raw_etag(&changed.metadata)), + IdentityField::VersionId => assert_ne!(source.version_id, persisted.version_id), + IdentityField::DataDir => assert_ne!(source.data_dir, persisted.data_dir), + IdentityField::ModTime => assert_ne!(source.mod_time, persisted.mod_time), + IdentityField::Size => assert_ne!(source.size, persisted.size), + IdentityField::Etag => assert_ne!(get_raw_etag(&source.metadata), get_raw_etag(&persisted.metadata)), } if !matches!(field, IdentityField::VersionId) { - assert_eq!(source.version_id, changed.version_id); + assert_eq!(source.version_id, persisted.version_id); } if !matches!(field, IdentityField::DataDir) { - assert_eq!(source.data_dir, changed.data_dir); + assert_eq!(source.data_dir, persisted.data_dir); } if !matches!(field, IdentityField::ModTime) { - assert_eq!(source.mod_time, changed.mod_time); + assert_eq!(source.mod_time, persisted.mod_time); } if !matches!(field, IdentityField::Size) { - assert_eq!(source.size, changed.size); + assert_eq!(source.size, persisted.size); } if !matches!(field, IdentityField::Etag) { - assert_eq!(get_raw_etag(&source.metadata), get_raw_etag(&changed.metadata)); + assert_eq!(get_raw_etag(&source.metadata), get_raw_etag(&persisted.metadata)); } } } diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 5c410839e..1ce3bbd7f 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) @@ -2720,10 +2722,12 @@ mod tests { #[serial_test::serial(storage_class_env)] async fn transition_transaction_recovery_deletes_provider_recovered_unknown_upload() { let versioned_remote = uuid::Uuid::new_v4().to_string(); + let nil_remote = uuid::Uuid::nil().to_string(); for (case, tier_name, remote_version) in [ ("missing", "TXPROBEMISSING", None), ("unversioned", "TXPROBEUNVERSIONED", Some(String::new())), ("versioned", "TXPROBEVERSIONED", Some(versioned_remote)), + ("nil-version", "TXPROBENILVERSION", Some(nil_remote)), ] { let temp_dir = tempfile::tempdir().expect("create temp store dir"); let (ctx, store, _shutdown) = without_storage_class_env(build_isolated_test_store( diff --git a/crates/filemeta/src/fileinfo.rs b/crates/filemeta/src/fileinfo.rs index 2f1714e01..e1f1f2e1b 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, @@ -230,6 +240,10 @@ pub struct FileInfo { pub transitioned_objname: String, pub transition_tier: String, 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, @@ -459,6 +473,10 @@ impl FileInfo { if self.mod_time.is_none_or(|mod_time| mod_time <= OffsetDateTime::UNIX_EPOCH) || (!allow_nil_version_id && self.version_id.is_some_and(|version_id| version_id.is_nil())) || self.transition_version_id.is_some_and(|version_id| version_id.is_nil()) + || self + .transition_version + .as_ref() + .is_some_and(|version_id| version_id.is_empty()) || self.size != 0 || self.data_dir.is_some() || self.mode.is_some() @@ -492,6 +510,7 @@ impl FileInfo { || !self.transitioned_objname.is_empty() || !self.transition_tier.is_empty() || self.transition_version_id.is_some() + || self.transition_version.is_some() || self.expire_restored || self.size != 0 || self.data_dir.is_some() @@ -536,6 +555,25 @@ impl FileInfo { /// return `None`. pub fn validate(&self, mode: ValidationMode) -> Result> { self.validate_collection_bounds()?; + if let (Some(version), Some(version_id)) = (&self.transition_version, self.transition_version_id) + && Uuid::parse_str(version).ok() != Some(version_id) + { + 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()?), @@ -832,6 +870,8 @@ impl FileInfo { && self.transition_tier == other.transition_tier && 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 @@ -1351,6 +1391,15 @@ mod tests { assert_file_corrupt(&fi, ValidationMode::DeleteOnly); } + #[test] + fn metadata_read_validation_rejects_conflicting_transition_versions() { + let mut fi = one_shard_validation_fileinfo(1); + fi.transition_version_id = Some(Uuid::new_v4()); + fi.transition_version = Some(Uuid::new_v4().to_string()); + + assert_file_corrupt(&fi, ValidationMode::RequireErasure); + } + #[test] fn metadata_read_validation_requires_canonical_delete_marker_shape() { let marker = FileInfo { @@ -1722,6 +1771,12 @@ mod tests { transitioned_objname, 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 f491e75b9..091bfb527 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; @@ -43,6 +44,7 @@ const MSGPACK_FIXEXT8: u8 = 0xd7; const MSGPACK_TIME_EXT_LEGACY: i8 = 5; const MSGPACK_TIME_EXT_OFFICIAL: i8 = -1; const MSGPACK_TIME_LEN: u8 = 12; +const MAX_TRANSITION_VERSION_LEN: usize = 1024; /// Sentinel signature returned when a version has no computable body (invalid / /// missing inner object). Mirrors MinIO's `signatureErr` so such versions never @@ -251,23 +253,93 @@ fn parse_legacy_uuid_bytes(bytes: &[u8], field: &str) -> Result> { /// Decode a stored transitioned-version-id from a version's `meta_sys`. /// -/// RustFS writes it as 16 raw UUID bytes; MinIO-migrated tiered objects store -/// the remote tier's version id as a UUID *string*. Accept both, and treat any -/// absent / nil / otherwise-unparseable value as "no tier version" (matching the -/// tolerant pre-hardening behavior) rather than failing the whole object read — -/// a malformed tier id must not make an otherwise-readable object unreadable. -fn transitioned_version_id_from_meta_sys(meta_sys: &HashMap>) -> Option { - let value = get_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID)?; +/// 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>) -> 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_some(id); + return Ok((!id.is_nil()).then(|| id.to_string())); } - std::str::from_utf8(&value) + 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()) + { + Ok(None) + } else { + 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) .ok() - .and_then(|s| Uuid::parse_str(s.trim()).ok()) - .filter(|id| !id.is_nil()) + .flatten() + .and_then(|value| Uuid::parse_str(&value).ok()) +} + +fn transitioned_version_bytes(fi: &FileInfo) -> Option> { + fi.transition_version + .as_ref() + .map(|version| version.as_bytes().to_vec()) + .or_else(|| fi.transition_version_id.map(|version_id| version_id.as_bytes().to_vec())) } fn parse_legacy_erasure_algo(value: &str) -> ErasureAlgo { @@ -2398,7 +2470,9 @@ 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_id = transitioned_version_id_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()) .unwrap_or_default(); @@ -2419,6 +2493,8 @@ impl MetaObject { transition_status, transitioned_objname, transition_version_id, + transition_version, + transition_version_state, transition_tier, ..Default::default() }) @@ -2431,13 +2507,12 @@ impl MetaObject { SUFFIX_TRANSITIONED_OBJECTNAME, fi.transitioned_objname.as_bytes().to_vec(), ); - if let Some(transition_version_id) = fi.transition_version_id.as_ref() { - insert_bytes( - &mut self.meta_sys, - SUFFIX_TRANSITIONED_VERSION_ID, - transition_version_id.as_bytes().to_vec(), - ); + if let Some(transition_version) = transitioned_version_bytes(fi) { + insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); + } else { + remove_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID); } + 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()); @@ -2501,6 +2576,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); @@ -2562,8 +2638,11 @@ impl From for MetaObject { ); } - if let Some(vid) = &value.transition_version_id { - insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, vid.as_bytes().to_vec()); + 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() { @@ -2706,7 +2785,11 @@ impl MetaDeleteMarker { .map(|v| String::from_utf8_lossy(&v).to_string()) .unwrap_or_default(); - fi.transition_version_id = transitioned_version_id_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 @@ -2859,8 +2942,11 @@ impl From for MetaDeleteMarker { value.transitioned_objname.as_bytes().to_vec(), ); } - if let Some(version_id) = value.transition_version_id { - insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, version_id.as_bytes().to_vec()); + 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()); @@ -3412,7 +3498,7 @@ mod tests { .insert("x-rustfs-internal-healing".to_string(), "true".to_string()); marker.metadata.insert("content-type".to_string(), "text/plain".to_string()); let remote_version_id = Uuid::new_v4(); - marker.transition_version_id = Some(remote_version_id); + marker.transition_version = Some(remote_version_id.to_string()); let converted = MetaDeleteMarker::from(marker); @@ -3420,7 +3506,19 @@ mod tests { assert_eq!(converted.meta_sys.get("x-minio-internal-purgestatus"), Some(&b"pending".to_vec())); assert_eq!( get_bytes(&converted.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID), - Some(remote_version_id.as_bytes().to_vec()) + Some(remote_version_id.to_string().into_bytes()) + ); + assert_eq!( + converted + .meta_sys + .get(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_ID}")), + Some(&remote_version_id.to_string().into_bytes()) + ); + assert_eq!( + converted + .meta_sys + .get(&format!("{}{SUFFIX_TRANSITIONED_VERSION_ID}", rustfs_utils::http::MINIO_INTERNAL_PREFIX)), + Some(&remote_version_id.to_string().into_bytes()) ); assert!(!converted.meta_sys.contains_key("x-rustfs-internal-healing")); assert!(!converted.meta_sys.contains_key("content-type")); @@ -4097,19 +4195,129 @@ mod tests { .into_fileinfo("b", "k", false) .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] - fn meta_object_transition_version_id_unparseable_stays_readable_as_none() { - // A non-UUID / non-16-byte tier version id must NOT make the object - // unreadable; it is tolerated as "no tier version" (compat with - // pre-hardening behavior and foreign/edge metadata). + fn meta_object_transition_version_id_opaque_text_is_preserved() { let mut sys = HashMap::new(); - insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"not-a-uuid".to_vec()); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"opaque-generation-42".to_vec()); let fi = make_meta_object_with_sys(sys) .into_fileinfo("b", "k", false) - .expect("unparseable transition version id must not fail the object read"); + .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 set_transition_known_disabled_removes_stale_version_dual_keys() { + let mut meta_sys = HashMap::new(); + insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, b"stale-legacy-version".to_vec()); + let mut object = make_meta_object_with_sys(meta_sys); + object.set_transition(&FileInfo { + transition_status: TRANSITION_COMPLETE.to_string(), + transitioned_objname: "remote/object".to_string(), + transition_version_state: TransitionVersionState::KnownDisabled, + transition_tier: "WARM".to_string(), + ..Default::default() + }); + + assert_eq!(get_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID), None); + assert!( + !object + .meta_sys + .contains_key(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_ID}")) + ); + assert!( + !object + .meta_sys + .contains_key(&format!("{}{SUFFIX_TRANSITIONED_VERSION_ID}", rustfs_utils::http::MINIO_INTERNAL_PREFIX)) + ); + let decoded = object + .into_fileinfo("b", "k", false) + .expect("known-disabled transition must remain readable after replacing stale metadata"); + assert_eq!(decoded.transition_version, None); + assert_eq!(decoded.transition_version_state, TransitionVersionState::KnownDisabled); + } + + #[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] + fn meta_object_transition_version_id_invalid_utf8_yields_none() { + let mut sys = HashMap::new(); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, vec![0xff]); + let fi = make_meta_object_with_sys(sys) + .into_fileinfo("b", "k", false) + .expect("invalid transition version bytes must not fail the object read"); + assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version, None); + } + + #[test] + fn meta_object_transition_version_id_unsafe_text_yields_none() { + for value in [b"opaque\0version".to_vec(), vec![b'x'; MAX_TRANSITION_VERSION_LEN + 1]] { + let mut sys = HashMap::new(); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, value); + let fi = make_meta_object_with_sys(sys) + .into_fileinfo("b", "k", false) + .expect("unsafe transition version text must not fail the object read"); + assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version, None); + } } #[test] @@ -4123,6 +4331,7 @@ mod tests { .into_fileinfo("b", "k", false) .expect("string-form transition version id must decode"); assert_eq!(fi.transition_version_id, Some(id)); + assert_eq!(fi.transition_version, Some(id.to_string())); } #[test] @@ -4152,16 +4361,14 @@ mod tests { } .into_fileinfo("b", "k", false); assert_eq!(fi.transition_version_id, Some(id)); + assert_eq!(fi.transition_version, Some(id.to_string())); } #[test] - fn delete_marker_free_version_transition_version_id_unparseable_stays_readable() { - // A malformed tier version id must not make a free-version record corrupt: - // it decodes to None and stays readable. Otherwise free-version expiry - // fails and the remote-tier object leaks. + fn delete_marker_free_version_transition_version_id_opaque_text_is_preserved() { let mut sys = HashMap::new(); insert_bytes(&mut sys, SUFFIX_FREE_VERSION, vec![]); - insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"not-a-uuid".to_vec()); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"opaque-generation-42".to_vec()); insert_bytes(&mut sys, SUFFIX_TRANSITION_TIER, b"WARM".to_vec()); insert_bytes(&mut sys, SUFFIX_TRANSITIONED_OBJECTNAME, b"remote-object".to_vec()); let fi = MetaDeleteMarker { @@ -4172,8 +4379,9 @@ mod tests { .into_fileinfo("b", "k", false); assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version.as_deref(), Some("opaque-generation-42")); fi.validate_for_metadata_read() - .expect("free-version record with an unparseable tier id must remain readable"); + .expect("free-version record with an opaque tier id must remain readable"); } #[test] @@ -4193,6 +4401,7 @@ mod tests { .into_fileinfo("b", "k", false); assert_eq!(fi.transition_version_id, Some(id)); + assert_eq!(fi.transition_version, Some(id.to_string())); } #[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"; diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 6e04bef5c..e77f7c0e6 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -763,7 +763,7 @@ async fn enqueue_transitioned_delete_cleanup( let _activity_guard = DeleteTailActivityGuard::new(DeleteTailStage::Cleanup); let je = if opts.delete_prefix { - tier_sweeper::transitioned_force_delete_journal_entry(&existing.transitioned_object) + tier_sweeper::transitioned_force_delete_journal_entry(&existing.transitioned_object, existing.transition_version_state) } else { let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok()); tier_sweeper::transitioned_delete_journal_entry( @@ -771,6 +771,7 @@ async fn enqueue_transitioned_delete_cleanup( opts.versioned, opts.version_suspended, &existing.transitioned_object, + existing.transition_version_state, ) }; let Some(mut je) = je else { @@ -9683,7 +9684,7 @@ mod tests { #[tokio::test] #[serial_test::serial] - async fn transitioned_delete_cleanup_persists_identity_bound_and_legacy_journals() { + async fn transitioned_delete_cleanup_persists_known_state_and_rejects_unknown_state() { let store = crate::app::gating_test_env::shared_gating_ecstore().await; if current_app_context().is_none() { crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await; @@ -9703,8 +9704,9 @@ mod tests { current.transitioned_object.tier = "WARM".to_string(); current.transitioned_object.name = "remote/identity-bound".to_string(); current.transitioned_object.version_id = "remote-version".to_string(); + current.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact; - let journal_name = |remote_object: &str, backend_identity: Option<[u8; 32]>| { + let journal_name = |remote_object: &str, backend_identity: Option<[u8; 32]>, version_id_exact: bool| { use sha2::{Digest, Sha256}; let mut hasher = Sha256::new(); @@ -9717,6 +9719,10 @@ mod tests { hasher.update([0]); hasher.update(backend_identity); } + if version_id_exact { + hasher.update([0]); + hasher.update(b"exact-version-id"); + } format!("ilm/tier-delete-journal/{}.json", rustfs_utils::crypto::hex(hasher.finalize().as_slice())) }; @@ -9726,7 +9732,7 @@ mod tests { let mut identity_bound = store .get_object_reader( ".rustfs.sys", - &journal_name("remote/identity-bound", Some(identity)), + &journal_name("remote/identity-bound", Some(identity), true), None, http::HeaderMap::new(), &ObjectOptions::default(), @@ -9739,11 +9745,14 @@ mod tests { .expect("identity-bound journal body should be readable"); let identity_bound: serde_json::Value = serde_json::from_slice(&identity_bound_data).expect("identity-bound journal should decode as JSON"); - assert_eq!(identity_bound["version"], serde_json::json!(2)); + assert_eq!(identity_bound["version"], serde_json::json!(4)); assert_eq!(identity_bound["backend_identity"], serde_json::json!(identity)); + assert_eq!(identity_bound["version_id_exact"], serde_json::json!(true)); + assert_eq!(identity_bound["version_state"], serde_json::json!("exact")); current.user_defined = Arc::new(HashMap::new()); current.transitioned_object.name = "remote/legacy".to_string(); + current.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown; enqueue_transitioned_delete_cleanup( store.clone(), "bucket", @@ -9755,24 +9764,31 @@ mod tests { Some(¤t), ) .await - .expect("legacy force-delete cleanup should persist a fail-closed v1 journal"); - let mut legacy = store + .expect("unknown force-delete cleanup should fail closed without a journal"); + let legacy_err = match store .get_object_reader( ".rustfs.sys", - &journal_name("remote/legacy", None), + &journal_name("remote/legacy", None, false), None, http::HeaderMap::new(), &ObjectOptions::default(), ) .await - .expect("legacy journal should be readable"); - let mut legacy_data = Vec::new(); - tokio::io::AsyncReadExt::read_to_end(&mut legacy.stream, &mut legacy_data) - .await - .expect("legacy journal body should be readable"); - let legacy: serde_json::Value = serde_json::from_slice(&legacy_data).expect("legacy journal should decode as JSON"); - assert_eq!(legacy["version"], serde_json::json!(1)); - assert_eq!(legacy["backend_identity"], serde_json::Value::Null); + { + Ok(_) => panic!("unknown remote version state must not persist a delete journal"), + Err(err) => err, + }; + assert!( + matches!( + &legacy_err, + StorageError::FileNotFound + | StorageError::ObjectNotFound(_, _) + | StorageError::FileVersionNotFound + | StorageError::VersionNotFound(_, _, _) + | StorageError::VolumeNotFound + ), + "unknown remote version state must leave no journal, got {legacy_err:?}" + ); } async fn put_real_cold_fill_object(store: &Arc, bucket: &str, object: &str, body: &[u8]) -> ObjectInfo { diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index 06765ac96..b97477b9a 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -409,20 +409,24 @@ pub(crate) mod bucket { versioned: bool, suspended: bool, transitioned: &super::super::super::storage_contracts::TransitionedObject, + transition_version_state: rustfs_filemeta::TransitionVersionState, ) -> Option { crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_delete_journal_entry( version_id, versioned, suspended, transitioned, + transition_version_state, ) } pub(crate) fn transitioned_force_delete_journal_entry( transitioned: &super::super::super::storage_contracts::TransitionedObject, + transition_version_state: rustfs_filemeta::TransitionVersionState, ) -> Option { crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_force_delete_journal_entry( transitioned, + transition_version_state, ) } }