diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index f6fb5e709..c4065f9bf 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=dbebfbab9b9efd4eff31211e69dd32235dc00e207f2ab0dd919a1b2ac9e724c2 +sha256=db9bd8cdcb0abe43461aa6b36499b17cabd4098e5b34e300b1a0f0d0f34d9884 diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 06cbda8ae..4d0bfc1bd 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -440,6 +440,11 @@ pub mod object { PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError, SnapshotConsistencyError, }; + + #[cfg(feature = "test-util")] + pub mod test_util { + pub use crate::store::DeleteAfterObjectLockSnapshotBarrier; + } } pub mod rebalance { diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 9b175986e..6b637cfac 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -921,7 +921,7 @@ impl ExpiryState { Ok(()) } - pub fn enqueue_free_version(&mut self, oi: ObjectInfo) -> bool { + pub fn enqueue_free_version(&self, oi: ObjectInfo) -> bool { let task = FreeVersionTask(oi); let wrkr = self.get_worker_ch(task.op_hash()); if wrkr.is_none() { @@ -1215,6 +1215,22 @@ impl ExpiryState { } } +pub(crate) async fn enqueue_committed_free_versions(api: &ECStore, free_versions: Vec) -> usize { + if free_versions.is_empty() { + return 0; + } + + let expiry_state = api.ctx.expiry_state(); + let state = expiry_state.read().await; + let mut queued = 0; + for free_version in free_versions { + if state.enqueue_free_version(free_version) { + queued += 1; + } + } + queued +} + async fn enqueue_recovered_free_version_with_state(state: &Arc>, oi: ObjectInfo) -> bool { let task = FreeVersionTask(oi); let hash = task.op_hash(); @@ -6627,7 +6643,7 @@ mod tests { async fn enqueue_free_version_reports_false_without_worker_channel() { let state = ExpiryState::new(); let recovery_notify = Arc::clone(&state.read().await.recovery_notify); - let mut state = state.write().await; + let state = state.write().await; let oi = ObjectInfo { bucket: "bucket".to_string(), name: "object".to_string(), @@ -6819,7 +6835,7 @@ mod tests { }, ..Default::default() }; - let mut state = state.write().await; + let state = state.write().await; assert!(state.enqueue_free_version(oi.clone())); assert!(recovery_notify.notified().now_or_never().is_none()); diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 83fc9a9dc..9127e2c9e 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -6444,6 +6444,9 @@ impl LocalDisk { .abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err) .await); } + if should_fail_after_delete_commit(self.root.as_path(), path) { + return Err(DiskError::Unexpected); + } return Ok(()); } @@ -6514,6 +6517,10 @@ impl LocalDisk { .await); } + if should_fail_after_delete_commit(self.root.as_path(), path) { + return Err(DiskError::Unexpected); + } + Ok(()) } diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index d2a3200fb..f668575b9 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -19,6 +19,7 @@ use crate::storage_api_contracts::{ HTTPPreconditions, ObjectLockRetentionOptions, ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, }, }; +use std::io; use std::sync::atomic::{AtomicBool, AtomicU8, Ordering}; use tokio::sync::{Mutex, Notify, OwnedRwLockReadGuard}; use tokio_util::sync::CancellationToken; @@ -678,6 +679,196 @@ pub struct DecommissionCapacityOptions { pub(crate) mutation_id: Option, } +/// Opaque storage-owned collection point for post-commit tier free-version +/// cleanup receipts. This type is public only because workspace crates build +/// [`ObjectOptions`] with struct literals; callers outside `ecstore` must leave +/// the corresponding option unset. +#[doc(hidden)] +#[derive(Clone)] +pub struct TierFreeVersionReceiptSink { + inner: Arc>, +} + +struct TierFreeVersionReceiptSinkState { + receipts: Option>, +} + +#[derive(PartialEq, Eq, Hash)] +struct TierFreeVersionReceiptIdentity { + bucket: String, + logical_name: String, + tier: String, + remote_name: String, + remote_version_state: TierFreeVersionReceiptVersionState, + remote_version: String, + backend_identity: crate::services::tier::tier::TierDestinationId, +} + +struct TierFreeVersionReceiptPayload { + local_free_version_id: Uuid, + mod_time: Option, +} + +#[derive(Clone, Copy, PartialEq, Eq, Hash)] +enum TierFreeVersionReceiptVersionState { + KnownDisabled, + SuspendedNull, + Exact, +} + +impl TierFreeVersionReceiptSink { + /// Only the delete wrapper may originate a sink. The public type exists so + /// workspace struct literals can carry it, but external crates cannot + /// create an undrainable collector accidentally. + pub(crate) fn new() -> Self { + Self { + inner: Arc::new(parking_lot::Mutex::new(TierFreeVersionReceiptSinkState { + receipts: Some(HashMap::new()), + })), + } + } +} + +impl std::fmt::Debug for TierFreeVersionReceiptSink { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let state = self.inner.lock(); + f.debug_struct("TierFreeVersionReceiptSink") + .field("drained", &state.receipts.is_none()) + .field("receipt_count", &state.receipts.as_ref().map(HashMap::len).unwrap_or_default()) + .finish() + } +} + +impl TierFreeVersionReceiptVersionState { + fn persisted(self) -> rustfs_filemeta::TransitionVersionState { + match self { + Self::KnownDisabled => rustfs_filemeta::TransitionVersionState::KnownDisabled, + Self::SuspendedNull => rustfs_filemeta::TransitionVersionState::SuspendedNull, + Self::Exact => rustfs_filemeta::TransitionVersionState::Exact, + } + } +} + +impl TierFreeVersionReceiptIdentity { + fn into_object_info(self, payload: TierFreeVersionReceiptPayload) -> ObjectInfo { + let mut metadata = HashMap::with_capacity(2); + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::crypto::hex(self.backend_identity), + ); + ObjectInfo { + bucket: self.bucket, + name: self.logical_name, + mod_time: payload.mod_time, + user_defined: Arc::new(metadata), + version_id: Some(payload.local_free_version_id), + delete_marker: true, + transitioned_object: TransitionedObject { + name: self.remote_name, + version_id: self.remote_version, + tier: self.tier, + free_version: true, + status: String::new(), + }, + transition_version_state: self.remote_version_state.persisted(), + ..Default::default() + } + } +} + +fn tier_free_version_scheduling_receipt_from_source( + source: &ObjectInfo, + local_free_version_id: Uuid, +) -> io::Result> { + if source.transitioned_object.status != rustfs_filemeta::TRANSITION_COMPLETE + || source.transitioned_object.free_version + || source.delete_marker + || source.bucket.is_empty() + || source.name.is_empty() + || source.transitioned_object.tier.is_empty() + || source.transitioned_object.name.is_empty() + { + return Ok(None); + } + if local_free_version_id.is_nil() { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "tier free-version receipt has a nil local version identity", + )); + } + + let remote_version = source.transitioned_object.version_id.as_str(); + let remote_version_state = match source.transition_version_state { + rustfs_filemeta::TransitionVersionState::Unknown => return Ok(None), + rustfs_filemeta::TransitionVersionState::KnownDisabled if remote_version.is_empty() => { + TierFreeVersionReceiptVersionState::KnownDisabled + } + rustfs_filemeta::TransitionVersionState::SuspendedNull if remote_version == "null" => { + TierFreeVersionReceiptVersionState::SuspendedNull + } + rustfs_filemeta::TransitionVersionState::Exact if !remote_version.is_empty() && remote_version != "null" => { + TierFreeVersionReceiptVersionState::Exact + } + _ => return Ok(None), + }; + + let Some(backend_identity) = crate::services::tier::tier::tier_destination_id_from_metadata(&source.user_defined) + .map_err(|err| io::Error::new(io::ErrorKind::InvalidData, err))? + else { + return Ok(None); + }; + + Ok(Some(( + TierFreeVersionReceiptIdentity { + bucket: source.bucket.clone(), + logical_name: decode_dir_object(&source.name), + tier: source.transitioned_object.tier.clone(), + remote_name: source.transitioned_object.name.clone(), + remote_version_state, + remote_version: source.transitioned_object.version_id.clone(), + backend_identity, + }, + TierFreeVersionReceiptPayload { + local_free_version_id, + mod_time: source.mod_time, + }, + ))) +} + +impl TierFreeVersionReceiptSink { + /// Record one committed free-version cleanup target. Cloned options share + /// this sink; tuple-equivalent physical copies collapse to one worker task. + /// `false` means the source cannot safely identify a destructive cleanup. + pub(crate) fn record(&self, source: &ObjectInfo, local_free_version_id: Uuid) -> io::Result { + let Some((identity, payload)) = tier_free_version_scheduling_receipt_from_source(source, local_free_version_id)? else { + return Ok(false); + }; + let mut state = self.inner.lock(); + let receipts = state + .receipts + .as_mut() + .ok_or_else(|| io::Error::new(io::ErrorKind::BrokenPipe, "tier free-version receipt sink was already drained"))?; + receipts.entry(identity).or_insert(payload); + Ok(true) + } + + /// Consume every receipt exactly once. A second drain is a caller bug: it + /// could otherwise make two outer wrappers believe they own the same tasks. + pub(crate) fn drain(&self) -> io::Result> { + let mut state = self.inner.lock(); + let receipts = state + .receipts + .take() + .ok_or_else(|| io::Error::new(io::ErrorKind::BrokenPipe, "tier free-version receipt sink was already drained"))?; + drop(state); + Ok(receipts + .into_iter() + .map(|(identity, payload)| identity.into_object_info(payload)) + .collect()) + } +} + #[derive(Default, Clone)] pub struct ObjectOptions { // Use the maximum parity (N/2), used when saving server configuration files @@ -716,6 +907,12 @@ pub struct ObjectOptions { pub skip_rebalancing: bool, pub skip_free_version: bool, + /// Storage-owned, per-request hand-off for committed tier free-version + /// cleanup work. The outer delete wrapper installs and drains it; clones + /// below that boundary share the same opaque sink. + #[doc(hidden)] + pub tier_free_version_receipt_sink: Option, + /// Cooperative cancellation for an owned PutObject before authoritative /// rename begins. Storage ignores it after entering the durable commit. #[doc(hidden)] @@ -851,6 +1048,7 @@ impl std::fmt::Debug for ObjectOptions { .field("skip_decommissioned", &self.skip_decommissioned) .field("skip_rebalancing", &self.skip_rebalancing) .field("skip_free_version", &self.skip_free_version) + .field("tier_free_version_receipt_sink", &self.tier_free_version_receipt_sink) .field("put_object_cancellation", &self.put_object_cancellation.is_some()) .field("scanner_publication_commit_scope", &self.scanner_publication_commit_scope) .field("data_movement", &self.data_movement) @@ -2622,11 +2820,381 @@ mod tests { assert!(default_cloned.parts.is_empty()); } + fn transitioned_receipt_source( + bucket: &str, + object: &str, + remote_version: &str, + version_state: rustfs_filemeta::TransitionVersionState, + identity_hex: Option<&str>, + ) -> ObjectInfo { + let mut metadata = HashMap::new(); + if let Some(identity_hex) = identity_hex { + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + identity_hex.to_string(), + ); + } + ObjectInfo { + bucket: bucket.to_string(), + name: object.to_string(), + mod_time: Some(OffsetDateTime::UNIX_EPOCH), + user_defined: Arc::new(metadata), + transitioned_object: TransitionedObject { + name: format!("remote/{object}"), + version_id: remote_version.to_string(), + tier: "WARM".to_string(), + status: TRANSITION_COMPLETE.to_string(), + ..Default::default() + }, + transition_version_state: version_state, + ..Default::default() + } + } + + #[test] + fn tier_free_version_receipt_matches_persisted_free_version_worker_fields() { + let bucket = "receipt-bucket"; + let object = "archive/object.bin"; + let source_version_id = Uuid::from_u128(1); + let local_free_version_id = Uuid::from_u128(2); + let remote_version_id = Uuid::from_u128(3); + let source_mod_time = OffsetDateTime::UNIX_EPOCH + time::Duration::hours(4); + let identity_hex = "ab".repeat(32); + let mut source_metadata = HashMap::from([ + ("etag".to_string(), "source-etag".to_string()), + ("x-amz-meta-private".to_string(), "must-not-enter-receipt".to_string()), + ]); + rustfs_utils::http::metadata_compat::insert_str( + &mut source_metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + identity_hex.clone(), + ); + let source_file_info = FileInfo { + volume: bucket.to_string(), + name: object.to_string(), + version_id: Some(source_version_id), + transition_status: TRANSITION_COMPLETE.to_string(), + transitioned_objname: "remote/receipt-object".to_string(), + transition_tier: "WARM".to_string(), + transition_version_id: Some(remote_version_id), + transition_version: Some(remote_version_id.to_string()), + transition_version_state: rustfs_filemeta::TransitionVersionState::Exact, + mod_time: Some(source_mod_time), + size: 8192, + data_dir: Some(Uuid::from_u128(4)), + metadata: source_metadata, + ..Default::default() + }; + let source = ObjectInfo::from_file_info(&source_file_info, bucket, object, true); + let mut persisted = FileMeta::new(); + persisted + .add_version(source_file_info) + .expect("transitioned receipt source should be persisted"); + let mut delete_file_info = FileInfo { + volume: bucket.to_string(), + name: object.to_string(), + version_id: Some(source_version_id), + mod_time: Some(source_mod_time + time::Duration::minutes(1)), + ..Default::default() + }; + delete_file_info.set_tier_free_version_id(&local_free_version_id.to_string()); + persisted + .delete_version(&delete_file_info) + .expect("transitioned source delete should create a free-version"); + + let encoded = persisted.marshal_msg().expect("free-version metadata should encode"); + let decoded = FileMeta::load(&encoded).expect("free-version metadata should decode"); + let persisted_free_version = decoded + .get_all_file_info_versions(bucket, object, true) + .expect("decoded free-version should produce FileInfo") + .versions + .into_iter() + .find(|version| version.tier_free_version()) + .expect("decoded metadata should contain the persisted free-version"); + let persisted_object_info = ObjectInfo::from_file_info(&persisted_free_version, bucket, object, true); + + let sink = TierFreeVersionReceiptSink::new(); + assert!( + sink.record(&source, local_free_version_id) + .expect("valid transitioned source should produce a receipt") + ); + let mut receipts = sink.drain().expect("receipt owner should drain exactly once"); + assert_eq!(receipts.len(), 1); + let receipt = receipts.pop().expect("one receipt should be present"); + + assert_eq!(receipt.bucket, persisted_object_info.bucket); + assert_eq!(receipt.name, persisted_object_info.name); + assert_eq!(receipt.version_id, persisted_object_info.version_id); + assert_eq!(receipt.mod_time, persisted_object_info.mod_time); + assert_eq!(receipt.delete_marker, persisted_object_info.delete_marker); + assert_eq!(receipt.transitioned_object.name, persisted_object_info.transitioned_object.name); + assert_eq!( + receipt.transitioned_object.version_id, + persisted_object_info.transitioned_object.version_id + ); + assert_eq!(receipt.transitioned_object.tier, persisted_object_info.transitioned_object.tier); + assert_eq!( + receipt.transitioned_object.free_version, + persisted_object_info.transitioned_object.free_version + ); + assert_eq!(receipt.transitioned_object.status, persisted_object_info.transitioned_object.status); + assert_eq!(receipt.transition_version_state, persisted_object_info.transition_version_state); + assert_eq!( + crate::services::tier::tier::tier_destination_id_from_metadata(&receipt.user_defined) + .expect("receipt identity should decode"), + crate::services::tier::tier::tier_destination_id_from_metadata(&persisted_object_info.user_defined) + .expect("persisted identity should decode") + ); + assert_eq!( + receipt.user_defined.len(), + 2, + "receipt should carry only the two compatibility identity keys" + ); + assert_eq!( + receipt.user_defined.get("x-rustfs-internal-transition-tier-destination-id"), + Some(&identity_hex) + ); + assert_eq!( + receipt.user_defined.get("x-minio-internal-transition-tier-destination-id"), + Some(&identity_hex) + ); + assert!(!receipt.user_defined.contains_key("x-amz-meta-private")); + assert_eq!(receipt.size, 0); + assert_eq!(receipt.actual_size, 0); + assert!(receipt.parts.is_empty()); + assert!(receipt.etag.is_none()); + assert!(receipt.checksum.is_none()); + assert!(receipt.data_dir.is_none()); + } + + #[test] + fn tier_free_version_receipt_sink_deduplicates_remote_target_and_drains_once() { + let identity_hex = "11".repeat(32); + let source = transitioned_receipt_source( + "bucket", + "object", + "remote-version", + rustfs_filemeta::TransitionVersionState::Exact, + Some(&identity_hex), + ); + let other_object = transitioned_receipt_source( + "bucket", + "other-object", + "remote-version", + rustfs_filemeta::TransitionVersionState::Exact, + Some(&identity_hex), + ); + let sink = TierFreeVersionReceiptSink::new(); + let clone = sink.clone(); + + assert!( + sink.record(&source, Uuid::from_u128(10)) + .expect("first physical receipt should record") + ); + assert!( + clone + .record(&source, Uuid::from_u128(11)) + .expect("tuple-equivalent physical receipt should be represented") + ); + assert!( + clone + .record(&other_object, Uuid::from_u128(12)) + .expect("a different logical key should retain its own task") + ); + + let mut receipts = sink.drain().expect("owner should drain shared receipts"); + receipts.sort_by(|left, right| left.name.cmp(&right.name)); + assert_eq!(receipts.len(), 2); + assert_eq!(receipts[0].name, "object"); + assert_eq!(receipts[0].version_id, Some(Uuid::from_u128(10))); + assert_eq!(receipts[1].name, "other-object"); + assert_eq!(receipts[1].version_id, Some(Uuid::from_u128(12))); + assert_eq!( + clone.drain().expect_err("a shared sink must drain only once").kind(), + io::ErrorKind::BrokenPipe + ); + assert_eq!( + clone + .record(&source, Uuid::from_u128(13)) + .expect_err("recording after drain must fail") + .kind(), + io::ErrorKind::BrokenPipe + ); + } + + #[test] + fn tier_free_version_receipt_identity_covers_every_destructive_dimension() { + let identity_hex = "44".repeat(32); + let baseline = transitioned_receipt_source( + "bucket", + "directory/", + "remote-version", + rustfs_filemeta::TransitionVersionState::Exact, + Some(&identity_hex), + ); + let mut encoded_duplicate = baseline.clone(); + encoded_duplicate.name = "directory__XLDIR__".to_string(); + + let mut variants = Vec::new(); + let mut changed = baseline.clone(); + changed.bucket = "other-bucket".to_string(); + variants.push(changed); + let mut changed = baseline.clone(); + changed.name = "other-directory/".to_string(); + variants.push(changed); + let mut changed = baseline.clone(); + changed.transitioned_object.tier = "COLD".to_string(); + variants.push(changed); + let mut changed = baseline.clone(); + changed.transitioned_object.name = "remote/other-directory/".to_string(); + variants.push(changed); + let mut changed = baseline.clone(); + changed.transitioned_object.version_id = "other-remote-version".to_string(); + variants.push(changed); + let mut changed = baseline.clone(); + changed.transition_version_state = rustfs_filemeta::TransitionVersionState::SuspendedNull; + changed.transitioned_object.version_id = "null".to_string(); + variants.push(changed); + let mut changed = baseline.clone(); + rustfs_utils::http::metadata_compat::insert_str( + Arc::make_mut(&mut changed.user_defined), + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + "55".repeat(32), + ); + variants.push(changed); + + let sink = TierFreeVersionReceiptSink::new(); + assert!( + sink.record(&baseline, Uuid::from_u128(20)) + .expect("baseline receipt should record") + ); + assert!( + sink.record(&encoded_duplicate, Uuid::from_u128(21)) + .expect("the encoded spelling of one logical key should deduplicate") + ); + for (offset, variant) in variants.iter().enumerate() { + assert!( + sink.record(variant, Uuid::from_u128(30 + offset as u128)) + .expect("each distinct cleanup identity should record") + ); + } + + let receipts = sink.drain().expect("identity matrix should drain once"); + assert_eq!(receipts.len(), 8, "every destructive identity dimension must prevent deduplication"); + let baseline_receipt = receipts + .iter() + .find(|receipt| { + receipt.bucket == "bucket" + && receipt.name == "directory/" + && receipt.transitioned_object.tier == "WARM" + && receipt.transitioned_object.name == "remote/directory/" + && receipt.transitioned_object.version_id == "remote-version" + && receipt.transition_version_state == rustfs_filemeta::TransitionVersionState::Exact + && crate::services::tier::tier::tier_destination_id_from_metadata(&receipt.user_defined) + .is_ok_and(|identity| identity == Some([0x44; 32])) + }) + .expect("baseline cleanup identity should remain present"); + assert_eq!( + baseline_receipt.version_id, + Some(Uuid::from_u128(20)), + "deduplication must retain the first UUID" + ); + } + + #[test] + fn tier_free_version_receipt_source_validation_fails_closed() { + let identity_hex = "22".repeat(32); + for (state, remote_version) in [ + (rustfs_filemeta::TransitionVersionState::KnownDisabled, ""), + (rustfs_filemeta::TransitionVersionState::SuspendedNull, "null"), + (rustfs_filemeta::TransitionVersionState::Exact, "opaque-version"), + ] { + let source = transitioned_receipt_source("bucket", "object", remote_version, state, Some(&identity_hex)); + assert!( + TierFreeVersionReceiptSink::new() + .record(&source, Uuid::new_v4()) + .expect("canonical remote-version state should be eligible"), + "state={state:?} remote_version={remote_version:?}" + ); + } + + let unknown = transitioned_receipt_source( + "bucket", + "object", + "opaque-version", + rustfs_filemeta::TransitionVersionState::Unknown, + Some(&identity_hex), + ); + assert!( + !TierFreeVersionReceiptSink::new() + .record(&unknown, Uuid::new_v4()) + .expect("unknown remote version state should defer to recovery") + ); + let missing_identity = transitioned_receipt_source( + "bucket", + "object", + "opaque-version", + rustfs_filemeta::TransitionVersionState::Exact, + None, + ); + assert!( + !TierFreeVersionReceiptSink::new() + .record(&missing_identity, Uuid::new_v4()) + .expect("missing durable identity should defer to recovery") + ); + let invalid_exact = transitioned_receipt_source( + "bucket", + "object", + "", + rustfs_filemeta::TransitionVersionState::Exact, + Some(&identity_hex), + ); + assert!( + !TierFreeVersionReceiptSink::new() + .record(&invalid_exact, Uuid::new_v4()) + .expect("conflicting remote state should defer to recovery") + ); + + let mut conflicting = transitioned_receipt_source( + "bucket", + "object", + "opaque-version", + rustfs_filemeta::TransitionVersionState::Exact, + Some(&identity_hex), + ); + Arc::make_mut(&mut conflicting.user_defined) + .insert("x-minio-internal-transition-tier-destination-id".to_string(), "33".repeat(32)); + assert_eq!( + TierFreeVersionReceiptSink::new() + .record(&conflicting, Uuid::new_v4()) + .expect_err("conflicting identity aliases must fail closed") + .kind(), + io::ErrorKind::InvalidData + ); + + let valid = transitioned_receipt_source( + "bucket", + "object", + "opaque-version", + rustfs_filemeta::TransitionVersionState::Exact, + Some(&identity_hex), + ); + assert_eq!( + TierFreeVersionReceiptSink::new() + .record(&valid, Uuid::nil()) + .expect_err("nil local free-version identity must be rejected") + .kind(), + io::ErrorKind::InvalidInput + ); + } + #[test] fn object_options_default_does_not_allocate_lifecycle_delete_all_journal() { let mut opts = ObjectOptions::default(); assert!(opts.lifecycle_delete_all_journal().is_none()); + assert!(opts.tier_free_version_receipt_sink.is_none()); opts.ensure_lifecycle_delete_all_journal(); assert!(opts.lifecycle_delete_all_journal().is_some()); } diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 6392e54a5..988090a72 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -217,7 +217,8 @@ use crate::bucket::lifecycle::{ use crate::bucket::quota::reservation; use crate::bucket::replication::{ DeleteReplicationConfigSnapshot, ReplicationLifecycleBridge, ReplicationStatusType, VersionPurgeStatusType, - replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_to_filemeta, + replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_from_filemeta, + version_purge_status_to_filemeta, }; use crate::data_usage::quota_object_size; use crate::diagnostics::get::GetObjectFailureReason; @@ -283,12 +284,95 @@ fn record_transitioned_delete_cleanup_owner(bucket: &str, object: &str, batch: b ); } -async fn acquire_single_tier_delete_lease( +/// A causal free-version receipt is only valid when the delete request really +/// removes the locked transitioned source. In particular, replication may turn +/// an otherwise successful delete into a metadata-only purge-state update, and +/// a versioned delete without a version ID writes a new delete marker instead +/// of removing the source selected by `goi`. +fn transitioned_delete_publishes_free_version(source: &ObjectInfo, delete_request: &FileInfo, skip_free_version: bool) -> bool { + if source.delete_marker + || source.transitioned_object.status != TRANSITION_COMPLETE + || skip_free_version + || delete_request.skip_tier_free_version() + || delete_request.expire_restored + || delete_request.transition_status == TRANSITION_COMPLETE + || delete_file_info_version_id(source.version_id) != delete_request.version_id + { + return false; + } + + // Keep this predicate aligned with FileMeta::delete_version's Object + // branch: a non-delete-marker request with a nonterminal purge status (or + // mark_deleted with no purge status) updates replication metadata in place + // and never calls MetaObject::init_free_version. + let purge_status = version_purge_status_from_filemeta(delete_request.version_purge_status()); + let metadata_only = !delete_request.deleted + && ((purge_status.is_empty() && delete_request.mark_deleted) + || (!purge_status.is_empty() && purge_status != VersionPurgeStatusType::Complete)); + + !metadata_only +} + +fn record_committed_tier_free_version_receipt( + opts: &ObjectOptions, bucket: &str, object: &str, - opts: &ObjectOptions, source: &ObjectInfo, -) -> Result> { + free_version_id: Uuid, + batch: bool, +) { + if let Some(sink) = opts.tier_free_version_receipt_sink.as_ref() + && let Err(err) = sink.record(source, free_version_id) + { + warn!( + event = EVENT_LIFECYCLE_TRANSITIONED_DELETE_CLEANUP_OWNER, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + bucket, + object, + batch, + error = ?err, + "Failed to retain the in-memory tier free-version scheduling receipt" + ); + } + record_transitioned_delete_cleanup_owner(bucket, object, batch); +} + +struct TierFreeVersionReceiptCandidate { + source: ObjectInfo, + free_version_id: Uuid, +} + +fn committed_tier_free_version_receipt_indices( + versions: &[FileInfoVersions], + delete_errors: &[Option], + candidates: &HashMap, +) -> Vec { + if candidates.is_empty() { + return Vec::new(); + } + + let mut committed = Vec::with_capacity(candidates.len()); + for group in versions { + let should_rollback = group + .versions + .iter() + .any(|version| delete_errors.get(version.idx).is_none_or(|error| error.is_some())); + if should_rollback { + continue; + } + committed.extend( + group + .versions + .iter() + .map(|version| version.idx) + .filter(|idx| candidates.contains_key(idx)), + ); + } + committed +} + +async fn acquire_single_tier_delete_lease(opts: &ObjectOptions, source: &ObjectInfo) -> Result> { let Some(api) = opts.tier_delete_journal_api.as_ref() else { return Ok(None); }; @@ -308,7 +392,6 @@ async fn acquire_single_tier_delete_lease( None => TierConfigMgr::acquire_operation_lease(&api.tier_config_mgr(), &source.transitioned_object.tier).await, } .map_err(Error::other)?; - record_transitioned_delete_cleanup_owner(bucket, object, false); Ok(Some(lease)) } @@ -384,6 +467,168 @@ mod scanner_publication_lease_fence_tests { } } +#[cfg(test)] +mod tier_free_version_receipt_eligibility_tests { + use super::*; + use crate::bucket::replication::ReplicationState; + + fn transitioned_source(version_id: Option) -> ObjectInfo { + let mut source = ObjectInfo { + version_id, + ..Default::default() + }; + source.transitioned_object.status = TRANSITION_COMPLETE.to_string(); + source.transitioned_object.tier = "WARM".to_string(); + source.transitioned_object.name = "remote/object".to_string(); + source + } + + fn delete_request(version_id: Option) -> FileInfo { + let mut request = FileInfo { + version_id, + ..Default::default() + }; + request.set_tier_free_version_id(&Uuid::new_v4().to_string()); + request + } + + #[test] + fn accepts_exact_transitioned_source_removal_and_suspended_null_replacement() { + let version_id = Uuid::new_v4(); + assert!(transitioned_delete_publishes_free_version( + &transitioned_source(Some(version_id)), + &delete_request(Some(version_id)), + false, + )); + + let mut suspended_null_delete = delete_request(None); + suspended_null_delete.deleted = true; + suspended_null_delete.mark_deleted = true; + assert!(transitioned_delete_publishes_free_version( + &transitioned_source(Some(Uuid::nil())), + &suspended_null_delete, + false, + )); + } + + #[test] + fn rejects_new_marker_version_and_non_transitioned_or_delete_marker_sources() { + let source_id = Uuid::new_v4(); + assert!(!transitioned_delete_publishes_free_version( + &transitioned_source(Some(source_id)), + &delete_request(Some(Uuid::new_v4())), + false, + )); + + let mut ordinary = transitioned_source(Some(source_id)); + ordinary.transitioned_object.status.clear(); + assert!(!transitioned_delete_publishes_free_version( + &ordinary, + &delete_request(Some(source_id)), + false, + )); + + let mut delete_marker = transitioned_source(Some(source_id)); + delete_marker.delete_marker = true; + assert!(!transitioned_delete_publishes_free_version( + &delete_marker, + &delete_request(Some(source_id)), + false, + )); + } + + #[test] + fn rejects_skip_restore_and_transition_metadata_updates() { + let version_id = Uuid::new_v4(); + let source = transitioned_source(Some(version_id)); + + assert!(!transitioned_delete_publishes_free_version( + &source, + &delete_request(Some(version_id)), + true, + )); + + let mut skip_request = delete_request(Some(version_id)); + skip_request.set_skip_tier_free_version(); + assert!(!transitioned_delete_publishes_free_version(&source, &skip_request, false)); + + let mut restore_request = delete_request(Some(version_id)); + restore_request.expire_restored = true; + assert!(!transitioned_delete_publishes_free_version(&source, &restore_request, false)); + + let mut transition_update = delete_request(Some(version_id)); + transition_update.transition_status = TRANSITION_COMPLETE.to_string(); + assert!(!transitioned_delete_publishes_free_version(&source, &transition_update, false)); + } + + #[test] + fn rejects_nonterminal_replication_metadata_only_update_but_accepts_complete_purge() { + let version_id = Uuid::new_v4(); + let source = transitioned_source(Some(version_id)); + + let mut pending = delete_request(Some(version_id)); + pending.replication_state_internal = Some(replication_state_to_filemeta(&ReplicationState { + version_purge_status_internal: Some("PENDING".to_string()), + ..Default::default() + })); + assert!(!transitioned_delete_publishes_free_version(&source, &pending, false)); + + let mut mark_deleted = delete_request(Some(version_id)); + mark_deleted.mark_deleted = true; + assert!(!transitioned_delete_publishes_free_version(&source, &mark_deleted, false)); + + let mut complete = delete_request(Some(version_id)); + complete.replication_state_internal = Some(replication_state_to_filemeta(&ReplicationState { + version_purge_status_internal: Some("COMPLETE".to_string()), + ..Default::default() + })); + assert!(transitioned_delete_publishes_free_version(&source, &complete, false)); + } + + #[test] + fn whole_physical_object_group_must_commit_before_any_receipt_is_retained() { + let candidate = || TierFreeVersionReceiptCandidate { + source: transitioned_source(Some(Uuid::new_v4())), + free_version_id: Uuid::new_v4(), + }; + let candidates = HashMap::from([(0, candidate()), (2, candidate())]); + let versions = vec![ + FileInfoVersions { + versions: vec![ + FileInfo { + idx: 0, + ..Default::default() + }, + FileInfo { + idx: 1, + ..Default::default() + }, + ], + ..Default::default() + }, + FileInfoVersions { + versions: vec![FileInfo { + idx: 2, + ..Default::default() + }], + ..Default::default() + }, + ]; + let one_sibling_failed = vec![None, Some(Error::other("injected quorum failure")), None]; + + assert_eq!( + committed_tier_free_version_receipt_indices(&versions, &one_sibling_failed, &candidates), + vec![2], + "a sibling failure must suppress every receipt from the rolled-back xl.meta group" + ); + assert_eq!( + committed_tier_free_version_receipt_indices(&versions, &[None, None, None], &candidates), + vec![0, 2], + "independent fully committed groups should retain their sparse receipts" + ); + } +} + struct PutObjectCommitCancellation { token: CancellationToken, armed: bool, @@ -6977,7 +7222,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }; let mut vers_map: HashMap<&String, FileInfoVersions> = HashMap::new(); let mut tier_reference_leases: Vec<(usize, String, Option)> = Vec::new(); - let mut transitioned_cleanup_items = vec![false; objects.len()]; + let mut tier_free_version_receipt_candidates: HashMap = HashMap::new(); for (i, dobj) in objects.iter().enumerate() { if del_errs[i].is_some() { @@ -7069,7 +7314,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { { match tier_destination_id_from_metadata(&goi.user_defined) { Ok(identity) => { - transitioned_cleanup_items[i] = true; tier_reference_leases.push((i, goi.transitioned_object.tier.clone(), identity)); } Err(err) => { @@ -7110,7 +7354,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { ..Default::default() }; - vr.set_tier_free_version_id(&Uuid::new_v4().to_string()); + let tier_free_version_id = Uuid::new_v4(); + vr.set_tier_free_version_id(&tier_free_version_id.to_string()); // Delete // del_objects[i].object_name.clone_from(&vr.name); @@ -7195,6 +7440,19 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }); } + if opts.tier_free_version_receipt_sink.is_some() + && !dobj.synthetic_version_id + && transitioned_delete_publishes_free_version(&goi, &vr, opts.skip_free_version) + { + tier_free_version_receipt_candidates.insert( + i, + TierFreeVersionReceiptCandidate { + source: goi, + free_version_id: tier_free_version_id, + }, + ); + } + // Only add to vers_map if we hold the lock if locked_objects.contains(&dobj.object_name) { vers_map.insert(&dobj.object_name, v); @@ -7271,12 +7529,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { } } - for (idx, transitioned) in transitioned_cleanup_items.into_iter().enumerate() { - if transitioned && del_errs[idx].is_none() { - record_transitioned_delete_cleanup_owner(bucket, &decode_dir_object(&objects[idx].object_name), true); - } - } - // Keep backend generations pinned through the source mutation, its // free-version write quorum, and any local rollback. Ordinary // single/batch deletes never transfer cleanup ownership to a journal. @@ -7388,6 +7640,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { self.record_capacity_scope_if_needed(opts.capacity_scope_token, &disks); let mut rollback_futures = Vec::new(); + let committed_receipt_indices = + committed_tier_free_version_receipt_indices(&vers, &del_errs, &tier_free_version_receipt_candidates); for fi_vers in &vers { // delete_versions commits one xl.meta per object group, so rollback must use the same boundary. let should_rollback = fi_vers.versions.iter().any(|fi| del_errs[fi.idx].is_some()); @@ -7462,6 +7716,20 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { join_all(rollback_futures).await; + for idx in committed_receipt_indices { + let Some(candidate) = tier_free_version_receipt_candidates.remove(&idx) else { + continue; + }; + record_committed_tier_free_version_receipt( + &opts, + bucket, + &decode_dir_object(&objects[idx].object_name), + &candidate.source, + candidate.free_version_id, + true, + ); + } + // TODO(backlog): support partial object deletion for multi-part objects if dist_erasure { @@ -7813,7 +8081,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?; begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; - let _tier_delete_lease = acquire_single_tier_delete_lease(bucket, object, &opts, &goi).await?; + let _tier_delete_lease = acquire_single_tier_delete_lease(&opts, &goi).await?; if opts.skip_free_version { fi.set_skip_tier_free_version(); } @@ -7823,6 +8091,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { #[cfg(test)] pause_delete_object_commit_after_publish(bucket, object).await; + if opts.tier_free_version_receipt_sink.is_some() + && transitioned_delete_publishes_free_version(&goi, &fi, opts.skip_free_version) + { + record_committed_tier_free_version_receipt(&opts, bucket, object, &goi, find_vid, false); + } + let disks = self.disk_inventory().await; self.record_capacity_scope_if_needed(opts.capacity_scope_token, &disks); @@ -7855,7 +8129,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?; begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; - let _tier_delete_lease = acquire_single_tier_delete_lease(bucket, object, &opts, &goi).await?; + let _tier_delete_lease = acquire_single_tier_delete_lease(&opts, &goi).await?; if opts.skip_free_version { dfi.set_skip_tier_free_version(); } @@ -7865,6 +8139,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { #[cfg(test)] pause_delete_object_commit_after_publish(bucket, object).await; + if opts.tier_free_version_receipt_sink.is_some() + && transitioned_delete_publishes_free_version(&goi, &dfi, opts.skip_free_version) + { + record_committed_tier_free_version_receipt(&opts, bucket, object, &goi, find_vid, false); + } + let disks = self.disk_inventory().await; self.record_capacity_scope_if_needed(opts.capacity_scope_token, &disks); diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index f830ba79e..1343bc475 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -10496,6 +10496,8 @@ mod tests { bucket: &str, object: &str, restore_before_delete: bool, + causal_enqueue: bool, + delete_with_journal: bool, ) { let temp_dir = tempfile::tempdir().expect("create transitioned delete store dir"); let (ctx, store, _shutdown) = @@ -10559,11 +10561,51 @@ mod tests { ); } - backend.set_remove_failure(true); - store - .delete_object_with_tier_delete_journal(bucket, object, ObjectOptions::default()) + if causal_enqueue { + ExpiryState::resize_workers(1, store.clone()).await; + } + backend.set_remove_failure(!causal_enqueue); + if delete_with_journal { + store + .delete_object_with_tier_delete_journal(bucket, object, ObjectOptions::default()) + .await + .expect("transitioned source journal-wrapper delete should commit"); + } else { + store + .delete_object(bucket, object, ObjectOptions::default()) + .await + .expect("transitioned source plain object-layer delete should commit"); + } + + if causal_enqueue { + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let metadata_absent = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("causal free-version cleanup metadata should remain readable") + .is_none(); + if metadata_absent && backend.remove_versions().await.len() == 1 { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) .await - .expect("transitioned source delete should commit"); + .expect("committed free-version should be cleaned without a recovery scan"); + assert_eq!(backend.object_count().await, 0, "causal cleanup should remove the remote object"); + assert_eq!( + tier_delete_journal_count(store.clone()).await, + 0, + "ordinary causal cleanup must not create a journal" + ); + store + .delete_bucket(bucket, &DeleteBucketOptions::default()) + .await + .expect("bucket delete should succeed after causal free-version cleanup"); + return; + } let local_versions = store.pools[0] .get_disks_by_key(object) @@ -10642,6 +10684,8 @@ mod tests { "transitioned-delete-journal-owner-bucket", "transition/archive.bin", false, + false, + true, ) .await; } @@ -10656,6 +10700,40 @@ mod tests { "restored-transitioned-delete-journal-owner-bucket", "transition/archive.bin", true, + false, + true, + ) + .await; + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn transitioned_delete_causally_enqueues_free_version() { + run_transitioned_delete_free_version_owner_case( + "transitioned-delete-causal-enqueue", + "DELETE-CAUSAL", + "transitioned-delete-causal-enqueue-bucket", + "transition/archive.bin", + false, + true, + false, + ) + .await; + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn restored_transitioned_delete_causally_enqueues_free_version() { + run_transitioned_delete_free_version_owner_case( + "restored-transitioned-delete-causal-enqueue", + "RESTORE-DELETE-CAUSAL", + "restored-transitioned-delete-causal-enqueue-bucket", + "transition/archive.bin", + true, + true, + true, ) .await; } @@ -12335,12 +12413,262 @@ mod tests { .is_none(), "free-version recovery must remove the exact cleanup owner" ); + + let causal = "causal.bin"; + let mut causal_reader = PutObjReader::from_vec(vec![b'c'; 1024 * 1024]); + let causal_source = store + .put_object(bucket, causal, &mut causal_reader, &ObjectOptions::default()) + .await + .expect("causal batch source should be written"); + store + .transition_object( + bucket, + causal, + &ObjectOptions { + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: causal_source.etag.clone().expect("causal batch source should have an etag"), + ..Default::default() + }, + mod_time: causal_source.mod_time, + ..Default::default() + }, + ) + .await + .expect("causal batch source transition should commit"); + let (_deleted, errors) = store + .delete_objects( + bucket, + vec![ + ObjectToDelete { + object_name: causal.to_string(), + ..Default::default() + }, + ObjectToDelete { + object_name: causal.to_string(), + ..Default::default() + }, + ], + ObjectOptions::default(), + ) + .await; + assert!( + errors.iter().all(Option::is_none), + "duplicate causal batch deletes should remain idempotent: {errors:?}" + ); + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let metadata_absent = store.pools[0] + .get_disks_by_key(causal) + .load_file_info_versions_exact(bucket, causal) + .await + .expect("causal batch cleanup metadata should remain readable") + .is_none(); + if metadata_absent && backend.remove_versions().await.len() >= 2 { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("batch free-version receipt should converge without another recovery scan"); + assert_eq!( + backend.remove_versions().await.len(), + 2, + "duplicate batch requests must cause only one remote delete for the causal object" + ); + + crate::bucket::metadata_sys::update_in( + &ctx, + bucket, + BUCKET_VERSIONING_CONFIG, + b"Enabled".to_vec(), + ) + .await + .expect("causal batch bucket versioning should be enabled"); + let versioned_causal = "versioned-causal.bin"; + let mut versioned_reader = PutObjReader::from_vec(vec![b'v'; 1024 * 1024]); + let versioned_source = store + .put_object( + bucket, + versioned_causal, + &mut versioned_reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("versioned causal batch source should be written"); + let versioned_source_id = versioned_source + .version_id + .expect("versioned causal batch source should have an identity"); + store + .transition_object( + bucket, + versioned_causal, + &ObjectOptions { + version_id: Some(versioned_source_id.to_string()), + versioned: true, + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: versioned_source + .etag + .clone() + .expect("versioned causal batch source should have an etag"), + ..Default::default() + }, + mod_time: versioned_source.mod_time, + ..Default::default() + }, + ) + .await + .expect("versioned causal batch source transition should commit"); + let (_deleted, errors) = store + .delete_objects( + bucket, + vec![ObjectToDelete { + object_name: versioned_causal.to_string(), + version_id: Some(versioned_source_id), + ..Default::default() + }], + ObjectOptions::default(), + ) + .await; + assert!( + errors.iter().all(Option::is_none), + "explicit-version causal batch delete should commit: {errors:?}" + ); + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let metadata_absent = store.pools[0] + .get_disks_by_key(versioned_causal) + .load_file_info_versions_exact(bucket, versioned_causal) + .await + .expect("versioned causal batch cleanup metadata should remain readable") + .is_none(); + if metadata_absent && backend.remove_versions().await.len() == 3 { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("explicit-version batch receipt should converge without a recovery scan"); store .delete_bucket(bucket, &DeleteBucketOptions::default()) .await .expect("batch source bucket should be physically empty"); } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn batch_transitioned_delete_aggregate_error_still_enqueues_committed_free_version() { + let temp_dir = tempfile::tempdir().expect("create aggregate-error batch delete store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "batch-transitioned-aggregate-error", &[4, 4])) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let tier_name = "BATCH-AGGREGATE-ERROR"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await; + let bucket = "batch-transitioned-aggregate-error-bucket"; + let object = "archive.bin"; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("aggregate-error source bucket should be created"); + let mut reader = PutObjReader::from_vec(vec![b'a'; 1024 * 1024]); + let source = store.pools[0] + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("aggregate-error source should be written"); + store.pools[0] + .transition_object( + bucket, + object, + &ObjectOptions { + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: source.etag.clone().expect("aggregate-error source should have an etag"), + ..Default::default() + }, + mod_time: source.mod_time, + ..Default::default() + }, + ) + .await + .expect("aggregate-error source should transition"); + + // Model a data-movement copy: both pools own the same logical source + // and exact remote tuple, but each batch delete creates its own local + // free-version UUID in the shared request sink. + for disk_index in 0..4 { + let source_meta = temp_dir + .path() + .join(format!("pool0/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}")); + let target_meta = temp_dir + .path() + .join(format!("pool1/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}")); + tokio::fs::create_dir_all(target_meta.parent().expect("target xl.meta should have a parent")) + .await + .expect("second-pool object directory should be created"); + tokio::fs::copy(&source_meta, &target_meta) + .await + .expect("transitioned xl.meta should copy exactly to the second pool"); + } + + ExpiryState::resize_workers(1, store.clone()).await; + let injection = crate::store::object::BatchDeletePoolErrorInjection::install( + bucket, + 1, + vec![(object.to_string(), StorageError::ErasureWriteQuorum)], + ); + let (deleted, errors) = store + .delete_objects( + bucket, + vec![ObjectToDelete { + object_name: object.to_string(), + ..Default::default() + }], + ObjectOptions::default(), + ) + .await; + assert_eq!(injection.observed(), 1, "the second pool should inject one post-commit aggregate error"); + assert_eq!(errors, vec![Some(StorageError::ErasureWriteQuorum)]); + assert!(deleted[0].found, "the aggregate error must retain the committed pool result"); + drop(injection); + + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let mut metadata_absent = true; + for pool in &store.pools { + metadata_absent &= pool + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("aggregate-error cleanup metadata should remain readable") + .is_none(); + } + if metadata_absent && backend.remove_count().await == 1 { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("aggregate failure must not suppress committed receipt dispatch"); + assert_eq!(backend.object_count().await, 0, "the shared remote object should be removed exactly once"); + store + .delete_bucket(bucket, &DeleteBucketOptions::default()) + .await + .expect("aggregate-error bucket should be physically empty"); + shutdown.cancel(); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)] @@ -12572,6 +12900,217 @@ mod tests { .expect("retry should leave the source bucket empty"); } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn batch_transitioned_delete_post_commit_failures_roll_back_without_free_version_receipt() { + let temp_dir = tempfile::tempdir().expect("create failed batch delete store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "batch-delete-local-failure", &[4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let tier_name = "BATCH-DELETE-LOCAL-FAIL"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await; + let bucket = "batch-delete-local-failure-bucket"; + let object = "archive.bin"; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("failed batch source bucket should be created"); + crate::bucket::metadata_sys::update_in( + &ctx, + bucket, + BUCKET_VERSIONING_CONFIG, + b"Enabled".to_vec(), + ) + .await + .expect("failed batch bucket versioning should be enabled"); + let mut transitioned_reader = PutObjReader::from_vec(vec![b't'; 1024 * 1024]); + let transitioned_source = store + .put_object( + bucket, + object, + &mut transitioned_reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed batch transitioned version should be written"); + let transitioned_version_id = transitioned_source + .version_id + .expect("failed batch transitioned source should have a version identity"); + store + .transition_object( + bucket, + object, + &ObjectOptions { + version_id: Some(transitioned_version_id.to_string()), + versioned: true, + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: transitioned_source + .etag + .clone() + .expect("failed batch transitioned source should have an etag"), + ..Default::default() + }, + mod_time: transitioned_source.mod_time, + ..Default::default() + }, + ) + .await + .expect("failed batch source version should transition"); + let mut ordinary_reader = PutObjReader::from_vec(vec![b'o'; 1024 * 1024]); + let ordinary_source = store + .put_object( + bucket, + object, + &mut ordinary_reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed batch ordinary sibling should be written"); + let ordinary_version_id = ordinary_source + .version_id + .expect("failed batch ordinary sibling should have a version identity"); + let delete_requests = || { + vec![ + ObjectToDelete { + object_name: object.to_string(), + version_id: Some(transitioned_version_id), + ..Default::default() + }, + ObjectToDelete { + object_name: object.to_string(), + version_id: Some(ordinary_version_id), + ..Default::default() + }, + ] + }; + + let set = store.pools[0].get_disks_by_key(object); + let disks = set.disks.read().await; + assert_eq!(disks.len(), 4, "the rollback fixture must use four disks"); + // Keep discovery fully online, then make a quorum of disks report an + // error only after their batch metadata commit has completed. + for disk in disks.iter().take(3) { + let disk = disk.as_ref().expect("injected rollback disks should be online"); + crate::disk::local::set_delete_version_fail_after_commit(disk.path().as_path(), object); + } + drop(disks); + let receipt_sink = crate::object_api::TierFreeVersionReceiptSink::new(); + let (_deleted, errors) = store + .delete_objects( + bucket, + delete_requests(), + ObjectOptions { + tier_free_version_receipt_sink: Some(receipt_sink.clone()), + ..Default::default() + }, + ) + .await; + assert_eq!( + errors, + vec![Some(StorageError::Unexpected), Some(StorageError::Unexpected),], + "three post-commit disk errors must fail batch delete before receipts publish" + ); + + assert!( + receipt_sink + .drain() + .expect("the test-owned failed-batch sink should drain exactly once") + .is_empty(), + "a rolled-back physical group must publish no cleanup receipt" + ); + let retained_transitioned = store + .get_object_info( + bucket, + object, + &ObjectOptions { + version_id: Some(transitioned_version_id.to_string()), + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed batch delete must restore the transitioned sibling"); + assert_eq!(retained_transitioned.transitioned_object.status, rustfs_filemeta::TRANSITION_COMPLETE); + let retained_ordinary = store + .get_object_info( + bucket, + object, + &ObjectOptions { + version_id: Some(ordinary_version_id.to_string()), + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed batch delete must restore the ordinary sibling"); + assert_ne!(retained_ordinary.transitioned_object.status, rustfs_filemeta::TRANSITION_COMPLETE); + let retained_versions = set + .load_file_info_versions_exact(bucket, object) + .await + .expect("rolled-back batch metadata should decode") + .expect("rolled-back batch source should remain on disk"); + assert_eq!( + retained_versions + .versions + .iter() + .chain(retained_versions.free_versions.iter()) + .filter(|version| version.tier_free_version()) + .count(), + 0, + "failed batch quorum must not retain a free-version owner" + ); + let retained_version_ids = retained_versions + .versions + .iter() + .filter_map(|version| version.version_id) + .collect::>(); + assert_eq!( + retained_version_ids, + std::collections::HashSet::from([transitioned_version_id, ordinary_version_id]), + "the physical-group rollback must restore both explicit siblings" + ); + assert_eq!(backend.object_count().await, 1, "failed batch commit must retain the remote object"); + assert_eq!(backend.remove_count().await, 0, "failed batch commit must not dispatch remote cleanup"); + + ExpiryState::resize_workers(1, store.clone()).await; + let (_deleted, retry_errors) = store + .delete_objects(bucket, delete_requests(), ObjectOptions::default()) + .await; + assert!( + retry_errors.iter().all(Option::is_none), + "retry after disk recovery should commit: {retry_errors:?}" + ); + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let metadata_absent = set + .load_file_info_versions_exact(bucket, object) + .await + .expect("retry cleanup metadata should remain readable") + .is_none(); + if metadata_absent && backend.remove_count().await == 1 { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("successful batch retry should converge without a recovery scan"); + store + .delete_bucket(bucket, &DeleteBucketOptions::default()) + .await + .expect("successful batch retry should leave the bucket empty"); + shutdown.cancel(); + } + #[cfg(feature = "test-util")] async fn run_multi_pool_same_remote_tuple_delete_case(batch: bool) { let temp_dir = tempfile::tempdir().expect("create shared-tuple multi-pool store dir"); @@ -12645,8 +13184,10 @@ mod tests { ); assert_eq!(backend.object_count().await, 1); + let receipt_sink = crate::object_api::TierFreeVersionReceiptSink::new(); let mut delete_opts = ObjectOptions { tier_delete_journal_api: Some(store.clone()), + tier_free_version_receipt_sink: Some(receipt_sink.clone()), ..Default::default() }; if batch { @@ -12774,7 +13315,20 @@ mod tests { ); backend.set_remove_failure(false); - wait_for_tier_free_version_recovery(store.clone(), &backend, 1).await; + let receipts = receipt_sink + .drain() + .expect("the simulated outer multi-pool wrapper should drain exactly once"); + assert_eq!( + receipts.len(), + 1, + "the same physical key and remote tuple must collapse to one causal task" + ); + assert_eq!( + crate::bucket::lifecycle::bucket_lifecycle_ops::enqueue_committed_free_versions(&store, receipts).await, + 1, + "the committed shared-tuple task should enter the running worker" + ); + wait_for_expiry_workers_idle(&store).await; assert_eq!(backend.remove_count().await, 1, "shared remote tuple should be deleted exactly once"); for pool_idx in 0..2 { assert!( diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index c31e17aea..71a636ae7 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -425,6 +425,8 @@ pub(crate) mod init_format; pub(crate) mod list_objects; mod multipart; mod object; +#[cfg(any(test, feature = "test-util"))] +pub use object::DeleteAfterObjectLockSnapshotBarrier; pub(crate) use object::{ DecommissionFixedReadAnchor, ObjectLockDiagGuard, RemoteTuplePublicationCommitGuard, RemoteTuplePublicationFence, SourceCleanupMutationFence, tiered_data_movement_source_matches, diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 49e211258..413032852 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -14,7 +14,7 @@ use super::*; use crate::bucket::lifecycle::{ - bucket_lifecycle_ops::eval_action_from_lifecycle, + bucket_lifecycle_ops::{enqueue_committed_free_versions, eval_action_from_lifecycle}, get_expiry_configs, tier_delete_journal::{ ActiveTierDeleteDispatch, EVENT_LIFECYCLE_TIER_DELETE_JOURNAL, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_LIFECYCLE, @@ -39,6 +39,7 @@ use crate::core::pools::{DecommissionCapacityOwner, ensure_decommission_capacity use crate::disk::OldCurrentSize; use crate::object_api::{ NamespaceLockFence, ObjectLockConfigSnapshot, ScannerPublicationCommitScopeGuard, ScannerPublicationCommitState, + TierFreeVersionReceiptSink, }; use crate::services::notification_sys::acquire_tier_delete_journal_fleet_proof; use crate::services::tier::tier::{TierConfigMgr, TierDestinationId, TierOperationLease, tier_destination_id_from_metadata}; @@ -71,6 +72,26 @@ const RECURSIVE_DELETE_VERSION_SCAN_PAGE_SIZE: i32 = 1000; #[cfg(test)] const RECURSIVE_DELETE_VERSION_SCAN_PAGE_SIZE: i32 = 2; +fn install_tier_free_version_receipt_sink(opts: &mut ObjectOptions) -> Option { + if opts.tier_free_version_receipt_sink.is_some() || opts.skip_free_version || opts.delete_prefix { + return None; + } + + let sink = TierFreeVersionReceiptSink::new(); + opts.tier_free_version_receipt_sink = Some(sink.clone()); + Some(sink) +} + +async fn enqueue_recorded_tier_free_versions(store: &ECStore, sink: Option) -> usize { + let Some(sink) = sink else { + return 0; + }; + let Ok(receipts) = sink.drain() else { + return 0; + }; + enqueue_committed_free_versions(store, receipts).await +} + fn build_tier_delete_journal_entry( bucket: &str, object: &str, @@ -1293,7 +1314,7 @@ fn should_create_delete_marker_for_missing_object(opts: &ObjectOptions) -> bool (opts.versioned || opts.version_suspended) && opts.version_id.is_none() && !opts.delete_marker && !opts.data_movement } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] struct DeleteAfterObjectLockSnapshotBarrierState { bucket: String, arrived: tokio::sync::Notify, @@ -1302,19 +1323,19 @@ struct DeleteAfterObjectLockSnapshotBarrierState { namespace_acquired: AtomicBool, } -#[cfg(test)] -pub(crate) struct DeleteAfterObjectLockSnapshotBarrier { +#[cfg(any(test, feature = "test-util"))] +pub struct DeleteAfterObjectLockSnapshotBarrier { state: Arc, } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] static DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER: std::sync::OnceLock< std::sync::Mutex>>, > = std::sync::OnceLock::new(); -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] impl DeleteAfterObjectLockSnapshotBarrier { - pub(crate) fn install(bucket: &str) -> Self { + pub fn install(bucket: &str) -> Self { let state = Arc::new(DeleteAfterObjectLockSnapshotBarrierState { bucket: bucket.to_string(), arrived: tokio::sync::Notify::new(), @@ -1331,15 +1352,15 @@ impl DeleteAfterObjectLockSnapshotBarrier { Self { state } } - pub(crate) async fn wait_until_paused(&self) { + pub async fn wait_until_paused(&self) { self.state.arrived.notified().await; } - pub(crate) fn release(&self) { + pub fn release(&self) { self.state.release.notify_one(); } - pub(crate) async fn release_and_wait_until_namespace_pending(&self) { + pub async fn release_and_wait_until_namespace_pending(&self) { let namespace_pending = self.state.namespace_pending.notified(); self.release(); tokio::time::timeout(Duration::from_secs(5), namespace_pending) @@ -1347,12 +1368,12 @@ impl DeleteAfterObjectLockSnapshotBarrier { .expect("delete should proceed to its namespace lock after leaving the snapshot barrier"); } - pub(crate) fn namespace_acquired(&self) -> bool { + pub fn namespace_acquired(&self) -> bool { self.state.namespace_acquired.load(Ordering::Acquire) } } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] impl Drop for DeleteAfterObjectLockSnapshotBarrier { fn drop(&mut self) { self.state.release.notify_one(); @@ -1365,7 +1386,7 @@ impl Drop for DeleteAfterObjectLockSnapshotBarrier { } } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] async fn pause_delete_after_object_lock_snapshot(bucket: &str) { let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) @@ -1377,11 +1398,24 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) { if let Some(state) = state { state.arrived.notify_one(); state.release.notified().await; + } +} + +#[cfg(any(test, feature = "test-util"))] +fn notify_delete_namespace_pending(bucket: &str) { + let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("delete snapshot barrier mutex should not poison") + .as_ref() + .filter(|state| state.bucket == bucket) + .cloned(); + if let Some(state) = state { state.namespace_pending.notify_one(); } } -#[cfg(test)] +#[cfg(any(test, feature = "test-util"))] fn notify_delete_namespace_acquired(bucket: &str) { let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) @@ -2564,6 +2598,10 @@ impl ECStore { let diag_enabled = is_object_lock_diag_enabled(); let ns_lock = self.handle_new_ns_lock(bucket, object).await?; let acquire_start = Instant::now(); + #[cfg(any(test, feature = "test-util"))] + if matches!(op, "delete_object" | "delete_objects") { + notify_delete_namespace_pending(bucket); + } let guard = ns_lock .get_write_lock(get_lock_acquire_timeout()) .await @@ -4155,7 +4193,16 @@ impl ECStore { opts: ObjectOptions, tier_journal_api: Option>, ) -> Result { - Box::pin(self.handle_delete_object_with_journal_inner(bucket, object, opts, tier_journal_api)).await + Box::pin(async move { + let mut opts = opts; + let receipt_sink = install_tier_free_version_receipt_sink(&mut opts); + let result = self + .handle_delete_object_with_journal_inner(bucket, object, opts, tier_journal_api) + .await; + enqueue_recorded_tier_free_versions(self, receipt_sink).await; + result + }) + .await } async fn handle_delete_object_with_journal_inner( @@ -4229,7 +4276,7 @@ impl ECStore { if opts.delete_prefix && opts.expected_bucket_incarnation_id.is_none() { opts.expected_bucket_incarnation_id = current_bucket_incarnation_id; } - #[cfg(test)] + #[cfg(any(test, feature = "test-util"))] pause_delete_after_object_lock_snapshot(bucket).await; if opts.delete_prefix && !opts.delete_prefix_object { @@ -4243,7 +4290,7 @@ impl ECStore { } else { None }; - #[cfg(test)] + #[cfg(any(test, feature = "test-util"))] if _object_lock_guard.is_some() { notify_delete_namespace_acquired(bucket); } @@ -4533,6 +4580,25 @@ impl ECStore { objects: Vec, opts: ObjectOptions, tier_journal_api: Option>, + ) -> (Vec, Vec>, Vec>) { + Box::pin(async move { + let mut opts = opts; + let receipt_sink = install_tier_free_version_receipt_sink(&mut opts); + let result = self + .handle_delete_objects_with_journal_and_accounting_inner(bucket, objects, opts, tier_journal_api) + .await; + enqueue_recorded_tier_free_versions(self, receipt_sink).await; + result + }) + .await + } + + async fn handle_delete_objects_with_journal_and_accounting_inner( + &self, + bucket: &str, + objects: Vec, + opts: ObjectOptions, + tier_journal_api: Option>, ) -> (Vec, Vec>, Vec>) { // encode object name let objects: Vec = objects @@ -4617,7 +4683,7 @@ impl ECStore { StorageError::BucketNotFound(bucket.to_string()), ); } - #[cfg(test)] + #[cfg(any(test, feature = "test-util"))] if current_bucket_incarnation_id.is_some() { pause_delete_after_object_lock_snapshot(bucket).await; } @@ -4625,7 +4691,7 @@ impl ECStore { Ok(guards) => guards, Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err), }; - #[cfg(test)] + #[cfg(any(test, feature = "test-util"))] if !_object_lock_guards.is_empty() { notify_delete_namespace_acquired(bucket); } @@ -7528,6 +7594,15 @@ mod tests { ); drop(unified_future); + let batch_future = + store.handle_delete_objects_with_journal_and_accounting("bucket", Vec::new(), ObjectOptions::default(), None); + let batch_future_size = std::mem::size_of_val(&batch_future); + assert!( + batch_future_size <= 4 * 1024, + "batch delete handler future must remain stack-bounded; measured {batch_future_size} bytes" + ); + drop(batch_future); + let outer_future = store.handle_delete_object("bucket", "object", ObjectOptions::default()); let outer_future_size = std::mem::size_of_val(&outer_future); assert!( diff --git a/crates/scanner/tests/lifecycle_integration_test.rs b/crates/scanner/tests/lifecycle_integration_test.rs index 71e2a9a62..fd81ab3e7 100644 --- a/crates/scanner/tests/lifecycle_integration_test.rs +++ b/crates/scanner/tests/lifecycle_integration_test.rs @@ -40,14 +40,15 @@ use uuid::Uuid; mod storage_api; use storage_api::lifecycle::{ - BUCKET_LIFECYCLE_CONFIG, BucketOperations, BucketOptions, BucketVersioningSys, CompletePart, DiskOption, ECStore, - EcstoreError, Endpoint, EndpointServerPools, Endpoints, IlmAction, LcEvent, LcEventSrc, ListOperations as _, - MakeBucketOptions, MockWarmBackend, MultipartOperations as _, ObjectIO as _, ObjectOperations as _, PoolEndpoints, - STORAGE_FORMAT_FILE, TRANSITION_PENDING, TransitionCleanupStoreBarrier, TransitionOptions, assert_transition_meta_consistent, - enqueue_transition_for_existing_objects, expire_transitioned_object, free_version_count, get_bucket_metadata, - get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys, init_local_disks, is_err_object_not_found, - is_err_version_not_found, new_disk, path2_bucket_object_with_base_path, recover_transition_transaction_records, - register_mock_tier_util, update_bucket_metadata, wait_for_free_version_absence, + BUCKET_LIFECYCLE_CONFIG, BucketOperations, BucketOptions, BucketVersioningSys, CompletePart, + DeleteAfterObjectLockSnapshotBarrier, DiskOption, ECStore, EcstoreError, Endpoint, EndpointServerPools, Endpoints, + ExpiryState, IlmAction, LcEvent, LcEventSrc, ListOperations as _, MakeBucketOptions, MockWarmBackend, + MultipartOperations as _, ObjectIO as _, ObjectOperations as _, PoolEndpoints, STORAGE_FORMAT_FILE, TRANSITION_PENDING, + TransitionCleanupStoreBarrier, TransitionOptions, assert_transition_meta_consistent, enqueue_transition_for_existing_objects, + expire_transitioned_object, free_version_count, get_bucket_metadata, get_global_tier_config_mgr, init_background_expiry, + init_bucket_metadata_sys, init_local_disks, is_err_object_not_found, is_err_version_not_found, new_disk, + path2_bucket_object_with_base_path, recover_transition_transaction_records, register_mock_tier_util, update_bucket_metadata, + wait_for_free_version_absence, }; static GLOBAL_ENV: OnceLock<(Vec, Arc)> = OnceLock::new(); @@ -601,24 +602,21 @@ mod serial_tests { /// persisted free-version recovery -- so no live local metadata ever points /// at an already-removed remote version. /// - /// This test pins the FIXED contract two complementary ways, both - /// revert-proof (reverting to remote-first ordering turns them red): + /// This test pins the fixed contract with deterministic GET and DELETE + /// barriers (reverting to remote-first ordering turns it red): /// - /// 1. Ordering (deterministic): immediately after - /// `expire_transitioned_object` returns, the remote tier object is still - /// present and the mock recorded **zero** remote `remove` calls -- - /// proving the local delete happened with no synchronous remote removal - /// (local-first). Remote-first ordering loses the object and records a - /// `remove`. - /// 2. Concurrent GET (user-visible): a tight GET loop runs concurrently - /// with the expiry; every observation must be either a full, correct - /// body (GET won) or a clean object/version-not-found (expiry won). A - /// tier-fetch failure -- the #3491 symptom -- is never tolerated. + /// 1. A GET that already resolved the transitioned metadata keeps its read + /// lock and returns the complete remote body while expiry waits. + /// 2. Expiry returns after committing the local free-version without + /// waiting for the post-commit worker's remote DELETE. While that DELETE + /// is paused, the durable marker and remote body must both still exist. + /// 3. A later GET observes a clean object/version-not-found, never a tier + /// fetch or read-quorum failure. #[tokio::test(flavor = "multi_thread", worker_threads = 1)] #[serial] #[ignore = "global-state ILM integration test: runs serialized in the CI ILM Integration (serial) lane, see ci.yml test-ilm-integration-serial and rustfs/backlog#1148 (ilm-2)"] async fn test_expire_transitioned_object_never_races_concurrent_get() { - let (_disk_paths, ecstore) = setup_isolated_test_env(false).await; + let (disk_paths, ecstore) = setup_isolated_test_env(false).await; let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); let backend = register_mock_tier(&tier_name).await; @@ -664,59 +662,7 @@ mod serial_tests { "the regression must exercise an unversioned remote tier" ); - // Concurrent GET loop: hammer GET while the expiry runs. Every outcome - // must be a full correct body or a clean not-found -- never a tier-fetch - // failure. - let get_store = ecstore.clone(); - let get_bucket = bucket_name.clone(); - let get_object = object_name.to_string(); - let expected = payload.clone(); - let get_loop = tokio::spawn(async move { - let mut saw_full_body = 0usize; - let mut saw_not_found = 0usize; - for _ in 0..400 { - match get_store - .get_object_reader( - get_bucket.as_str(), - get_object.as_str(), - None, - http::HeaderMap::new(), - &ObjectOptions::default(), - ) - .await - { - Ok(mut reader) => { - let mut data = Vec::new(); - match reader.stream.read_to_end(&mut data).await { - Ok(_) => { - assert_eq!( - data, expected, - "a successful GET during expiry must return the complete, correct body" - ); - saw_full_body += 1; - } - Err(err) => { - panic!("GET during expiry streamed a truncated/failed body (expire/GET race regression): {err:?}") - } - } - } - Err(err) => { - let ec: &EcstoreError = &err; - assert!( - is_err_object_not_found(ec) || is_err_version_not_found(ec), - "GET during expiry may only fail with a clean object/version-not-found (expiry won \ - the race); a tier-fetch failure is the #3491 regression: {err:?}" - ); - saw_not_found += 1; - } - } - tokio::task::yield_now().await; - } - (saw_full_body, saw_not_found) - }); - - // Run the exact expiry action the scanner drives for a transitioned - // current version. + ExpiryState::resize_workers(1, ecstore.clone()).await; let lc_event = LcEvent { action: IlmAction::DeleteAction, ..Default::default() @@ -725,39 +671,127 @@ mod serial_tests { .bucket_incarnation_id(bucket_name.as_str()) .await .expect("read bucket incarnation"); - expire_transitioned_object(ecstore.clone(), &oi, &lc_event, &LcEventSrc::Scanner, bucket_incarnation_id) + + // Pause one real tier GET after it has resolved local transition + // metadata. The reader still owns the object read lock, so local expiry + // cannot commit until this GET finishes. + let get_barrier = backend.arm_get_barrier().await; + let get_store = ecstore.clone(); + let get_bucket = bucket_name.clone(); + let get_object = object_name.to_string(); + let in_flight_get = tokio::spawn(async move { + let mut reader = get_store + .get_object_reader( + get_bucket.as_str(), + get_object.as_str(), + None, + http::HeaderMap::new(), + &ObjectOptions::default(), + ) + .await + .map_err(|err| format!("in-flight GET failed before streaming: {err:?}"))?; + let mut data = Vec::new(); + reader + .stream + .read_to_end(&mut data) + .await + .map_err(|err| format!("in-flight GET returned a failed or truncated stream: {err:?}"))?; + Ok::<_, String>(data) + }); + tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, get_barrier.wait_until_paused()) .await + .expect("the in-flight GET should reach the remote read barrier"); + + // The next remote DELETE pauses and then fails. A correct local-first + // expiry returns while this barrier is still held; synchronous cleanup + // (remote-first or local-first) instead times out here. + let delete_start_barrier = DeleteAfterObjectLockSnapshotBarrier::install(bucket_name.as_str()); + let remove_barrier = backend.arm_failing_remove_barrier().await; + let expiry_store = ecstore.clone(); + let expiry_oi = oi.clone(); + let expiry_event = lc_event.clone(); + let mut expiry = tokio::spawn(async move { + expire_transitioned_object(expiry_store, &expiry_oi, &expiry_event, &LcEventSrc::Scanner, bucket_incarnation_id).await + }); + tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, delete_start_barrier.wait_until_paused()) + .await + .expect("expiry should reach the store delete path while the GET remains paused"); + delete_start_barrier.release_and_wait_until_namespace_pending().await; + assert!( + !delete_start_barrier.namespace_acquired() && !expiry.is_finished(), + "expiry must wait for the in-flight GET's object read lock before committing the local delete" + ); + + get_barrier.release(); + let expiry_outcome = tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, &mut expiry).await; + let delete_lock_acquired_after_get = delete_start_barrier.namespace_acquired(); + drop(delete_start_barrier); + let remove_arrival = tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, remove_barrier.wait_until_paused()).await; + + // Snapshot only lock-free observables while the cleanup worker holds + // the object write lock. Store API reads wait until the barrier is + // released below. + let free_version_persisted = free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0; + let remote_present_at_cleanup = backend.contains(&remote_object).await; + + remove_barrier.release(); + let remove_operation_dropped = if remove_arrival.is_ok() { + Some(tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, remove_barrier.wait_until_operation_dropped()).await) + } else { + None + }; + if expiry_outcome.is_err() && tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, &mut expiry).await.is_err() { + expiry.abort(); + let _ = expiry.await; + } + let in_flight_get_outcome = tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, in_flight_get).await; + let post_expiry_get = tokio::time::timeout( + TRANSITION_WAIT_TIMEOUT, + ecstore.get_object_reader(bucket_name.as_str(), object_name, None, http::HeaderMap::new(), &ObjectOptions::default()), + ) + .await; + + expiry_outcome + .expect("expire_transitioned_object must not wait for asynchronous remote-tier cleanup") + .expect("the expiry task should not panic") .expect("expire_transitioned_object should succeed"); + assert!( + delete_lock_acquired_after_get, + "expiry must acquire the object write lock only after the in-flight GET releases its read lock" + ); + remove_arrival.expect("the post-commit free-version worker should reach the remote DELETE barrier"); + remove_operation_dropped + .expect("the remote DELETE should have reached the barrier") + .expect("the injected remote DELETE should finish after release"); + assert!( + free_version_persisted, + "the durable free-version marker must exist before asynchronous remote cleanup" + ); + assert!( + remote_present_at_cleanup, + "the remote object must remain readable until the paused cleanup DELETE is released" + ); - // --- Ordering contract (deterministic revert-proof) ---------------- - // #3491 defers remote cleanup to free-version recovery, so immediately - // after expiry the remote object is still present and NO synchronous - // remote `remove` was issued. Reverting to remote-first ordering makes - // both assertions fail. + let in_flight_body = in_flight_get_outcome + .expect("the in-flight GET should finish within the test deadline") + .expect("the in-flight GET task should not panic") + .expect("a GET that wins the expiry race must return a complete body"); assert_eq!( - backend.remove_count().await, - 0, - "expire_transitioned_object must NOT issue a synchronous remote-tier removal (local-first \ - ordering, #3491); remote cleanup is deferred to free-version recovery" - ); - assert!( - backend.contains(&remote_object).await, - "remote tier object must still exist immediately after expiry (deferred cleanup, #3491)" + in_flight_body, payload, + "a GET that resolved transitioned metadata before expiry must return the complete, correct body" ); - // Local metadata is gone: the object is atomically unreachable. - assert!( - wait_for_object_absence(&ecstore, bucket_name.as_str(), object_name, Duration::from_secs(5)).await, - "local metadata for the expired transitioned object should be gone" - ); - - // Drain the concurrent GET loop; its internal asserts already guarantee - // no #3491-style tier-fetch failure was ever observed. - let (saw_full_body, saw_not_found) = get_loop.await.expect("concurrent GET loop task panicked"); - assert!( - saw_full_body + saw_not_found > 0, - "the concurrent GET loop should have observed at least one GET outcome" - ); + match post_expiry_get.expect("the post-expiry GET should finish within the test deadline") { + Ok(_) => panic!("the locally expired transitioned object must no longer be readable"), + Err(err) => { + let ec: &EcstoreError = &err; + assert!( + is_err_object_not_found(ec) || is_err_version_not_found(ec), + "a GET after expiry may only fail with a clean object/version-not-found; \ + a tier-fetch or read-quorum failure is the #3491 regression: {err:?}" + ); + } + } } #[test] @@ -1469,10 +1503,18 @@ mod serial_tests { let stale_remote_object = transitioned.transitioned_object.name.clone(); assert!(backend.contains(&stale_remote_object).await); - ecstore - .delete_object(bucket_name.as_str(), object_name, ObjectOptions::default()) + ExpiryState::resize_workers(1, ecstore.clone()).await; + let remove_barrier = backend.arm_failing_remove_barrier().await; + tokio::time::timeout( + Duration::from_secs(5), + ecstore.delete_object(bucket_name.as_str(), object_name, ObjectOptions::default()), + ) + .await + .expect("DeleteObject must not wait for asynchronous remote-tier cleanup") + .expect("Failed to delete transitioned object before scanner fallback"); + tokio::time::timeout(Duration::from_secs(5), remove_barrier.wait_until_paused()) .await - .expect("Failed to delete transitioned object without expiry workers"); + .expect("the immediate free-version worker should reach the injected remote DELETE barrier"); assert!( free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0, @@ -1483,8 +1525,12 @@ mod serial_tests { "stale transitioned remote object should still exist before scanner fallback runs" ); - init_background_expiry(ecstore.clone()).await; + // Queue the scanner fallback while the causal task is still blocked. + // Releasing the barrier fails only that first task, so the queued + // scanner task can prove durable-marker recovery on a healthy backend. scan_object_metadata(&disk_paths[0], bucket_name.as_str(), object_name).await; + remove_barrier.release(); + remove_barrier.wait_until_operation_dropped().await; assert!( backend @@ -1531,10 +1577,18 @@ mod serial_tests { let stale_remote_object = transitioned.transitioned_object.name.clone(); assert!(backend.contains(&stale_remote_object).await); - ecstore - .delete_object(bucket_name.as_str(), object_name, ObjectOptions::default()) + ExpiryState::resize_workers(1, ecstore.clone()).await; + let remove_barrier = backend.arm_failing_remove_barrier().await; + tokio::time::timeout( + Duration::from_secs(5), + ecstore.delete_object(bucket_name.as_str(), object_name, ObjectOptions::default()), + ) + .await + .expect("DeleteObject must not wait for asynchronous remote-tier cleanup") + .expect("Failed to delete transitioned object after compensation-driven transition"); + tokio::time::timeout(Duration::from_secs(5), remove_barrier.wait_until_paused()) .await - .expect("Failed to delete transitioned object after compensation-driven transition"); + .expect("the immediate free-version worker should reach the injected remote DELETE barrier"); assert!( free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0, @@ -1545,8 +1599,12 @@ mod serial_tests { "stale transitioned remote object should still exist before scanner cleanup runs" ); - init_background_expiry(ecstore.clone()).await; + // Enqueue the scanner fallback before the first, causal cleanup task is + // released into its injected failure. This keeps attribution + // deterministic and proves the durable marker drives convergence. scan_object_metadata(&disk_paths[0], bucket_name.as_str(), object_name).await; + remove_barrier.release(); + remove_barrier.wait_until_operation_dropped().await; assert!( backend diff --git a/crates/scanner/tests/storage_api/mod.rs b/crates/scanner/tests/storage_api/mod.rs index 61ace3902..186983c29 100644 --- a/crates/scanner/tests/storage_api/mod.rs +++ b/crates/scanner/tests/storage_api/mod.rs @@ -15,7 +15,9 @@ pub(crate) use rustfs_ecstore::api::bucket::lifecycle::transition_transaction::recover_transition_transaction_records; pub(crate) use rustfs_ecstore::api::bucket::lifecycle::{ bucket_lifecycle_audit::LcEventSrc, - bucket_lifecycle_ops::{enqueue_transition_for_existing_objects, expire_transitioned_object, init_background_expiry}, + bucket_lifecycle_ops::{ + ExpiryState, enqueue_transition_for_existing_objects, expire_transitioned_object, init_background_expiry, + }, lifecycle::{Event as LcEvent, IlmAction, TRANSITION_PENDING, TransitionOptions}, }; pub(crate) use rustfs_ecstore::api::bucket::metadata::BUCKET_LIFECYCLE_CONFIG; @@ -27,6 +29,7 @@ pub(crate) use rustfs_ecstore::api::capacity::path2_bucket_object_with_base_path pub(crate) use rustfs_ecstore::api::disk::{DiskOption, STORAGE_FORMAT_FILE, endpoint::Endpoint, new_disk}; pub(crate) use rustfs_ecstore::api::error::{Error as EcstoreError, is_err_object_not_found, is_err_version_not_found}; pub(crate) use rustfs_ecstore::api::layout::{EndpointServerPools, Endpoints, PoolEndpoints}; +pub(crate) use rustfs_ecstore::api::object::test_util::DeleteAfterObjectLockSnapshotBarrier; pub(crate) use rustfs_ecstore::api::runtime::global_tier_config_mgr as get_global_tier_config_mgr; pub(crate) use rustfs_ecstore::api::storage::{ECStore, init_local_disks}; // Shared lifecycle/tier test utilities (rustfs/backlog#1148 ilm-6). The mock @@ -45,12 +48,12 @@ pub(crate) mod lifecycle { }; pub(crate) use super::{ - BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DiskOption, ECStore, EcstoreError, Endpoint, EndpointServerPools, - Endpoints, IlmAction, LcEvent, LcEventSrc, MockWarmBackend, PoolEndpoints, STORAGE_FORMAT_FILE, TRANSITION_PENDING, - TransitionCleanupStoreBarrier, TransitionOptions, assert_transition_meta_consistent, - enqueue_transition_for_existing_objects, expire_transitioned_object, free_version_count, get_bucket_metadata, - get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys, init_local_disks, is_err_object_not_found, - is_err_version_not_found, new_disk, path2_bucket_object_with_base_path, recover_transition_transaction_records, - register_mock_tier_util, update_bucket_metadata, wait_for_free_version_absence, + BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DeleteAfterObjectLockSnapshotBarrier, DiskOption, ECStore, EcstoreError, + Endpoint, EndpointServerPools, Endpoints, ExpiryState, IlmAction, LcEvent, LcEventSrc, MockWarmBackend, PoolEndpoints, + STORAGE_FORMAT_FILE, TRANSITION_PENDING, TransitionCleanupStoreBarrier, TransitionOptions, + assert_transition_meta_consistent, enqueue_transition_for_existing_objects, expire_transitioned_object, + free_version_count, get_bucket_metadata, get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys, + init_local_disks, is_err_object_not_found, is_err_version_not_found, new_disk, path2_bucket_object_with_base_path, + recover_transition_transaction_records, register_mock_tier_util, update_bucket_metadata, wait_for_free_version_absence, }; }