fix(tiering): preserve remote version state on delete

This commit is contained in:
马登山
2026-07-28 18:18:44 +08:00
parent 4d8088ddbd
commit 4839096440
3 changed files with 85 additions and 23 deletions
@@ -252,7 +252,10 @@ impl ObjSweeper {
version_id: self.transition_version_id.clone(), version_id: self.transition_version_id.clone(),
tier_name: self.transition_tier.clone(), tier_name: self.transition_tier.clone(),
backend_identity: None, 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, version_state: self.transition_version_state,
}); });
} }
@@ -465,6 +468,7 @@ pub fn transitioned_delete_journal_entry(
versioned: bool, versioned: bool,
suspended: bool, suspended: bool,
transitioned: &TransitionedObject, transitioned: &TransitionedObject,
transition_version_state: rustfs_filemeta::TransitionVersionState,
) -> Option<Jentry> { ) -> Option<Jentry> {
let sweeper = ObjSweeper { let sweeper = ObjSweeper {
version_id, version_id,
@@ -473,6 +477,7 @@ pub fn transitioned_delete_journal_entry(
transition_status: transitioned.status.clone(), transition_status: transitioned.status.clone(),
transition_tier: transitioned.tier.clone(), transition_tier: transitioned.tier.clone(),
transition_version_id: transitioned.version_id.clone(), transition_version_id: transitioned.version_id.clone(),
transition_version_state,
remote_object: transitioned.name.clone(), remote_object: transitioned.name.clone(),
..Default::default() ..Default::default()
}; };
@@ -480,8 +485,13 @@ pub fn transitioned_delete_journal_entry(
sweeper.should_remove_remote_object() sweeper.should_remove_remote_object()
} }
pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject) -> Option<Jentry> { pub fn transitioned_force_delete_journal_entry(
if transitioned.status != lifecycle::TRANSITION_COMPLETE { transitioned: &TransitionedObject,
transition_version_state: rustfs_filemeta::TransitionVersionState,
) -> Option<Jentry> {
if transitioned.status != lifecycle::TRANSITION_COMPLETE
|| transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown
{
return None; return None;
} }
@@ -490,8 +500,11 @@ pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject
version_id: transitioned.version_id.clone(), version_id: transitioned.version_id.clone(),
tier_name: transitioned.tier.clone(), tier_name: transitioned.tier.clone(),
backend_identity: None, backend_identity: None,
version_id_exact: false, version_id_exact: matches!(
version_state: rustfs_filemeta::TransitionVersionState::Unknown, transition_version_state,
rustfs_filemeta::TransitionVersionState::SuspendedNull | rustfs_filemeta::TransitionVersionState::Exact
),
version_state: transition_version_state,
}) })
} }
@@ -502,9 +515,11 @@ mod test {
use super::{ use super::{
ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED, RemoteDeleteBreaker, RemoteTierDeleteOutcome, 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, 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, is_remote_tier_not_found_error, is_signer_header_error, lifecycle, set_remote_tier_delete_test_hook,
should_record_remote_delete_failure, 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::io::{Error, ErrorKind};
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
@@ -548,6 +563,43 @@ mod test {
assert!(should_record_remote_delete_failure(&Error::other("NoSuchVersion"))); 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] #[tokio::test]
#[serial_test::serial] #[serial_test::serial]
async fn idempotent_remote_delete_treats_hooked_nosuchversion_as_already_removed() { async fn idempotent_remote_delete_treats_hooked_nosuchversion_as_already_removed() {
+22 -16
View File
@@ -763,7 +763,7 @@ async fn enqueue_transitioned_delete_cleanup(
let _activity_guard = DeleteTailActivityGuard::new(DeleteTailStage::Cleanup); let _activity_guard = DeleteTailActivityGuard::new(DeleteTailStage::Cleanup);
let je = if opts.delete_prefix { 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 { } else {
let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok()); let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok());
tier_sweeper::transitioned_delete_journal_entry( tier_sweeper::transitioned_delete_journal_entry(
@@ -771,6 +771,7 @@ async fn enqueue_transitioned_delete_cleanup(
opts.versioned, opts.versioned,
opts.version_suspended, opts.version_suspended,
&existing.transitioned_object, &existing.transitioned_object,
existing.transition_version_state,
) )
}; };
let Some(mut je) = je else { let Some(mut je) = je else {
@@ -9683,7 +9684,7 @@ mod tests {
#[tokio::test] #[tokio::test]
#[serial_test::serial] #[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; let store = crate::app::gating_test_env::shared_gating_ecstore().await;
if current_app_context().is_none() { if current_app_context().is_none() {
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await; 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.tier = "WARM".to_string();
current.transitioned_object.name = "remote/identity-bound".to_string(); current.transitioned_object.name = "remote/identity-bound".to_string();
current.transitioned_object.version_id = "remote-version".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}; use sha2::{Digest, Sha256};
let mut hasher = Sha256::new(); let mut hasher = Sha256::new();
@@ -9717,6 +9719,10 @@ mod tests {
hasher.update([0]); hasher.update([0]);
hasher.update(backend_identity); 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())) 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 let mut identity_bound = store
.get_object_reader( .get_object_reader(
".rustfs.sys", ".rustfs.sys",
&journal_name("remote/identity-bound", Some(identity)), &journal_name("remote/identity-bound", Some(identity), true),
None, None,
http::HeaderMap::new(), http::HeaderMap::new(),
&ObjectOptions::default(), &ObjectOptions::default(),
@@ -9739,11 +9745,14 @@ mod tests {
.expect("identity-bound journal body should be readable"); .expect("identity-bound journal body should be readable");
let identity_bound: serde_json::Value = let identity_bound: serde_json::Value =
serde_json::from_slice(&identity_bound_data).expect("identity-bound journal should decode as JSON"); 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["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.user_defined = Arc::new(HashMap::new());
current.transitioned_object.name = "remote/legacy".to_string(); current.transitioned_object.name = "remote/legacy".to_string();
current.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown;
enqueue_transitioned_delete_cleanup( enqueue_transitioned_delete_cleanup(
store.clone(), store.clone(),
"bucket", "bucket",
@@ -9755,24 +9764,21 @@ mod tests {
Some(&current), Some(&current),
) )
.await .await
.expect("legacy force-delete cleanup should persist a fail-closed v1 journal"); .expect("unknown force-delete cleanup should fail closed without a journal");
let mut legacy = store let legacy_err = match store
.get_object_reader( .get_object_reader(
".rustfs.sys", ".rustfs.sys",
&journal_name("remote/legacy", None), &journal_name("remote/legacy", None, false),
None, None,
http::HeaderMap::new(), http::HeaderMap::new(),
&ObjectOptions::default(), &ObjectOptions::default(),
) )
.await .await
.expect("legacy journal should be readable"); {
let mut legacy_data = Vec::new(); Ok(_) => panic!("unknown remote version state must not persist a delete journal"),
tokio::io::AsyncReadExt::read_to_end(&mut legacy.stream, &mut legacy_data) Err(err) => err,
.await };
.expect("legacy journal body should be readable"); assert!(matches!(legacy_err, StorageError::FileNotFound));
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);
} }
async fn put_real_cold_fill_object(store: &Arc<ECStore>, bucket: &str, object: &str, body: &[u8]) -> ObjectInfo { async fn put_real_cold_fill_object(store: &Arc<ECStore>, bucket: &str, object: &str, body: &[u8]) -> ObjectInfo {
+4
View File
@@ -409,20 +409,24 @@ pub(crate) mod bucket {
versioned: bool, versioned: bool,
suspended: bool, suspended: bool,
transitioned: &super::super::super::storage_contracts::TransitionedObject, transitioned: &super::super::super::storage_contracts::TransitionedObject,
transition_version_state: rustfs_filemeta::TransitionVersionState,
) -> Option<Jentry> { ) -> Option<Jentry> {
crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_delete_journal_entry( crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_delete_journal_entry(
version_id, version_id,
versioned, versioned,
suspended, suspended,
transitioned, transitioned,
transition_version_state,
) )
} }
pub(crate) fn transitioned_force_delete_journal_entry( pub(crate) fn transitioned_force_delete_journal_entry(
transitioned: &super::super::super::storage_contracts::TransitionedObject, transitioned: &super::super::super::storage_contracts::TransitionedObject,
transition_version_state: rustfs_filemeta::TransitionVersionState,
) -> Option<Jentry> { ) -> Option<Jentry> {
crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_force_delete_journal_entry( crate::storage::storage_api::ecstore_bucket::lifecycle::tier_sweeper::transitioned_force_delete_journal_entry(
transitioned, transitioned,
transition_version_state,
) )
} }
} }