From 4839096440bc0bd324e8b243a4ef057ad4cc27e9 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 18:18:44 +0800 Subject: [PATCH] fix(tiering): preserve remote version state on delete --- .../src/bucket/lifecycle/tier_sweeper.rs | 66 +++++++++++++++++-- rustfs/src/app/object_usecase.rs | 38 ++++++----- rustfs/src/app/storage_api.rs | 4 ++ 3 files changed, 85 insertions(+), 23 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs index c5fbc29a9..113263eb9 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs @@ -252,7 +252,10 @@ 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, }); } @@ -465,6 +468,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, @@ -473,6 +477,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() }; @@ -480,8 +485,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; } @@ -490,8 +500,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_state: rustfs_filemeta::TransitionVersionState::Unknown, + version_id_exact: matches!( + transition_version_state, + rustfs_filemeta::TransitionVersionState::SuspendedNull | rustfs_filemeta::TransitionVersionState::Exact + ), + version_state: transition_version_state, }) } @@ -502,9 +515,11 @@ mod test { use super::{ ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED, RemoteDeleteBreaker, RemoteTierDeleteOutcome, 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}; @@ -548,6 +563,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() { diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 6e04bef5c..a9828919e 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,21 @@ 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)); } 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, ) } }