From 3d63a755a967caff8e5d5ccbb1ea465b564d515d 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 19:00:30 +0800 Subject: [PATCH] fix(tiering): replay exact cleanup journals --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 55 +++++++++++++++++++ .../bucket/lifecycle/tier_delete_journal.rs | 34 ++++++++---- .../src/bucket/lifecycle/tier_sweeper.rs | 12 ++++ crates/ecstore/src/set_disk/mod.rs | 2 + crates/ecstore/src/set_disk/ops/object.rs | 2 +- 5 files changed, 94 insertions(+), 11 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 4369e307e..2d9598339 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -10194,6 +10194,61 @@ mod tests { 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, diff --git a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs index b7cb3cc32..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}; @@ -270,15 +273,26 @@ pub async fn process_tier_delete_journal_entry(api: Arc, je: &Jentry) - 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 } diff --git a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs index f2cb9ed92..3a8d909b0 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs @@ -793,6 +793,18 @@ mod test { 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/object", + "", + "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); } #[test] diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index bdd66c945..8d6c73948 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -690,6 +690,8 @@ mod ops; #[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 2803fc208..d8a2c32dd 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -1807,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,