diff --git a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs index 113263eb9..f2cb9ed92 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs @@ -338,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( @@ -346,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); @@ -443,7 +446,46 @@ 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", + )); + } + delete_object_from_remote_tier_with_lease_idempotent_inner(obj_name, rv_id, lease, true, false).await +} + +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) => { @@ -514,9 +556,10 @@ 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, lifecycle, set_remote_tier_delete_test_hook, - should_record_remote_delete_failure, transitioned_delete_journal_entry, transitioned_force_delete_journal_entry, + 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, 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; @@ -722,6 +765,36 @@ mod test { assert_eq!(backend.remove_versions().await, vec![("remote/object".to_string(), String::new())]); } + #[cfg(feature = "test-util")] + #[tokio::test] + async fn confirmed_transition_cleanup_deletes_exact_provider_token() { + 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())] + ); + } + #[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/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 6c1474d83..2803fc208 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) { @@ -1979,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; @@ -2239,10 +2317,7 @@ fn persisted_transition_version( remote_version: &str, ) -> std::io::Result<(Option, rustfs_filemeta::TransitionVersionState)> { if remote_version.is_empty() { - return Err(std::io::Error::new( - std::io::ErrorKind::Unsupported, - "a missing remote tier object version remains unknown until the cluster capability gate is active", - )); + return Ok((None, rustfs_filemeta::TransitionVersionState::KnownDisabled)); } let version_id = Uuid::parse_str(remote_version).map_err(|_| { std::io::Error::new( @@ -2500,7 +2575,10 @@ mod transition_version_id_tests { #[test] fn normalizes_persisted_unversioned_ids_and_preserves_put_constraints() { - assert!(persisted_transition_version("").is_err()); + assert_eq!( + 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()); @@ -3798,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, @@ -3815,20 +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, transition_version_state) = match persisted_transition_version(candidate.remote_version()) { - Ok(version) => version, - Err(err) => { - let cleanup_api = transition_cleanup_store(&self.ctx).await; - if let Err(cleanup_err) = upload_cleanup.cleanup_rejected_upload(cleanup_api).await { - return Err(StorageError::Io(std::io::Error::other(format!( - "{err}; rejected remote upload cleanup failed: {cleanup_err}" - )))); - } - delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object) - .await; - return Err(err.into()); - } - }; let mut commit_opts = opts.clone(); commit_opts.no_lock = true; @@ -4700,7 +4778,7 @@ mod transition_commit_failure_tests { #[tokio::test] async fn rejected_unsupported_remote_versions_are_cleaned_up() { - for remote_version in ["", "null", "opaque-version-token"] { + 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") @@ -6675,20 +6753,21 @@ mod transition_upload_integrity_tests { #[tokio::test] #[serial_test::serial] - async fn opaque_remote_version_is_persisted_exactly() { + 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-unknown-version-bucket"; + let bucket = "transition-unversioned-tier-bucket"; let object = "object.bin"; - let payload = b"unknown remote version must retain local data".repeat(1024); + 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("opaque-version-token".to_string())).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 accepted opaque remote version must commit"); + .expect("an unversioned remote version must commit"); let (fi, _, _) = set_disks .get_object_fileinfo( bucket, @@ -6702,13 +6781,70 @@ mod transition_upload_integrity_tests { false, ) .await - .expect("committed opaque transition metadata should be readable"); + .expect("committed unversioned transition metadata should be readable"); assert_eq!(fi.transition_version_id, None); - assert_eq!(fi.transition_version.as_deref(), Some("opaque-version-token")); + 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() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "transition-unknown-version-bucket"; + let object = "object.bin"; + let payload = b"unknown 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 backend = register_mock_tier(&runtime_sources::global_tier_config_mgr(), &tier_name).await; + backend.set_put_remote_version(Some("opaque-version-token".to_string())).await; + + set_disks + .transition_object(bucket, object, &transition_options(&original, tier_name)) + .await + .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"); + assert_eq!(backend.object_count().await, 0); + 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() { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 801ccd411..1ce3bbd7f 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -2722,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/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index 3ecc4dfcc..091bfb527 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -2509,6 +2509,8 @@ impl MetaObject { ); 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()); @@ -4248,6 +4250,37 @@ mod tests { 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();