diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs
index 326b798ee..abdca65c1 100644
--- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs
+++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs
@@ -4994,6 +4994,9 @@ pub async fn apply_expiry_on_transitioned_object(
src: &LcEventSrc,
bucket_incarnation_id: Uuid,
) -> bool {
+ if lc_event.action.delete_all() {
+ return apply_expiry_on_non_transitioned_objects(api, oi, lc_event, src, bucket_incarnation_id).await;
+ }
let time_ilm = Metrics::time_ilm(lc_event.action);
if let Err(_err) = expire_transitioned_object(api, oi, lc_event, src, bucket_incarnation_id).await {
return false;
@@ -5047,13 +5050,24 @@ pub async fn apply_expiry_on_non_transitioned_objects(
if lc_event.action.delete_all() {
opts.delete_prefix = true;
opts.delete_prefix_object = true;
+ opts.lifecycle_delete_all = Some(crate::object_api::LifecycleDeleteAllRequest {
+ version_id: oi.version_id.filter(|version_id| !version_id.is_nil()),
+ delete_marker: oi.delete_marker,
+ action: lc_event.action,
+ rule_id: lc_event.rule_id.clone(),
+ phase: crate::object_api::LifecycleDeleteAllPhase::Preflight,
+ });
+ opts.ensure_lifecycle_delete_all_journal();
}
let time_ilm = Metrics::time_ilm(lc_event.action);
//debug!("lc_event.action: {:?}", lc_event.action);
debug!("expiry_on_non_transitioned_objects opts: {:?}", opts);
- let mut dobj = match api.delete_object(&oi.bucket, &encode_dir_object(&oi.name), opts).await {
+ let mut dobj = match api
+ .delete_object_with_tier_delete_journal(&oi.bucket, &encode_dir_object(&oi.name), opts)
+ .await
+ {
Ok(dobj) => dobj,
Err(e) => {
error!(
@@ -5283,7 +5297,7 @@ mod tests {
};
use crate::bucket::lifecycle::tier_last_day_stats::LastDayTierStats;
use crate::bucket::lifecycle::tier_sweeper::Jentry;
- use crate::bucket::metadata::BUCKET_LIFECYCLE_CONFIG;
+ use crate::bucket::metadata::{BUCKET_LIFECYCLE_CONFIG, BUCKET_VERSIONING_CONFIG};
use crate::bucket::metadata_sys;
#[cfg(feature = "test-util")]
use crate::client::transition_api::ReaderImpl;
@@ -5304,6 +5318,7 @@ mod tests {
use crate::storage_api_contracts::{
bucket::{BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions},
lifecycle::ExpirationOptions,
+ list::ListOperations as _,
multipart::MultipartOperations as _,
object::{ObjectIO as _, ObjectOperations as _},
};
@@ -10917,6 +10932,199 @@ mod tests {
);
}
+ #[tokio::test]
+ #[serial]
+ async fn queued_delete_all_rechecks_a_same_id_rule_moved_into_the_future() {
+ let (_disk_paths, ecstore) = setup_test_env().await;
+ let bucket = format!("stale-delete-all-rule-{}", Uuid::new_v4().simple());
+ let object = "object";
+ create_test_bucket(&ecstore, &bucket).await;
+ metadata_sys::update(
+ &bucket,
+ BUCKET_VERSIONING_CONFIG,
+ b"Enabled".to_vec(),
+ )
+ .await
+ .expect("bucket versioning should be enabled");
+ let lifecycle_xml = |days| {
+ format!(
+ r#"
+
+ delete-marker-history
+ Enabled
+
+ {days}
+
+"#
+ )
+ };
+ metadata_sys::update(&bucket, BUCKET_LIFECYCLE_CONFIG, lifecycle_xml(1).into_bytes())
+ .await
+ .expect("initial lifecycle rule should be stored");
+
+ let old_time = OffsetDateTime::now_utc() - time::Duration::days(3);
+ let mut reader = PutObjReader::from_vec(b"old version".to_vec());
+ ecstore
+ .put_object(
+ &bucket,
+ object,
+ &mut reader,
+ &ObjectOptions {
+ versioned: true,
+ mod_time: Some(old_time - time::Duration::hours(1)),
+ ..Default::default()
+ },
+ )
+ .await
+ .expect("old version should be stored");
+ let marker = ecstore
+ .delete_object(
+ &bucket,
+ object,
+ ObjectOptions {
+ versioned: true,
+ mod_time: Some(old_time),
+ ..Default::default()
+ },
+ )
+ .await
+ .expect("delete marker should be created");
+
+ let queued_event = crate::bucket::lifecycle::lifecycle::Event {
+ action: IlmAction::DelMarkerDeleteAllVersionsAction,
+ rule_id: "delete-marker-history".to_string(),
+ ..Default::default()
+ };
+ metadata_sys::update(&bucket, BUCKET_LIFECYCLE_CONFIG, lifecycle_xml(30).into_bytes())
+ .await
+ .expect("updated lifecycle rule should be stored");
+ let incarnation = ecstore
+ .bucket_incarnation_id_from_disk(&bucket)
+ .await
+ .expect("bucket incarnation should be available");
+
+ let deleted = super::apply_expiry_on_non_transitioned_objects(
+ ecstore.clone(),
+ &marker,
+ &queued_event,
+ &LcEventSrc::Scanner,
+ incarnation,
+ )
+ .await;
+ assert!(!deleted, "the stale queued rule must be rejected");
+ let versions = ecstore
+ .clone()
+ .list_object_versions(&bucket, object, None, None, None, 10)
+ .await
+ .expect("remaining versions should be listable");
+ assert_eq!(versions.objects.iter().filter(|version| version.name == object).count(), 2);
+
+ metadata_sys::update(&bucket, BUCKET_LIFECYCLE_CONFIG, lifecycle_xml(1).into_bytes())
+ .await
+ .expect("due lifecycle rule should be restored");
+ let deleted = super::apply_expiry_on_non_transitioned_objects(
+ ecstore.clone(),
+ &marker,
+ &queued_event,
+ &LcEventSrc::Scanner,
+ incarnation,
+ )
+ .await;
+ assert!(deleted, "the current due rule should purge marker and history");
+ let versions = ecstore
+ .clone()
+ .list_object_versions(&bucket, object, None, None, None, 10)
+ .await
+ .expect("purged versions should be listable");
+ assert_eq!(versions.objects.iter().filter(|version| version.name == object).count(), 0);
+ }
+
+ #[tokio::test]
+ #[serial]
+ async fn queued_expired_object_all_versions_purges_history_through_transitioned_dispatch() {
+ let (_disk_paths, ecstore) = setup_test_env().await;
+ let bucket = format!("expired-all-versions-{}", Uuid::new_v4().simple());
+ let object = "object";
+ create_test_bucket(&ecstore, &bucket).await;
+ metadata_sys::update(
+ &bucket,
+ BUCKET_VERSIONING_CONFIG,
+ b"Enabled".to_vec(),
+ )
+ .await
+ .expect("bucket versioning should be enabled");
+ metadata_sys::update(
+ &bucket,
+ BUCKET_LIFECYCLE_CONFIG,
+ br#"
+
+ delete-all-versions
+ Enabled
+
+ 1true
+
+"#
+ .to_vec(),
+ )
+ .await
+ .expect("delete-all lifecycle rule should be stored");
+
+ let old_time = OffsetDateTime::now_utc() - time::Duration::days(3);
+ let mut old_reader = PutObjReader::from_vec(b"old version".to_vec());
+ ecstore
+ .put_object(
+ &bucket,
+ object,
+ &mut old_reader,
+ &ObjectOptions {
+ versioned: true,
+ mod_time: Some(old_time - time::Duration::hours(1)),
+ ..Default::default()
+ },
+ )
+ .await
+ .expect("old version should be stored");
+ let mut current_reader = PutObjReader::from_vec(b"current version".to_vec());
+ let mut current = ecstore
+ .put_object(
+ &bucket,
+ object,
+ &mut current_reader,
+ &ObjectOptions {
+ versioned: true,
+ mod_time: Some(old_time),
+ ..Default::default()
+ },
+ )
+ .await
+ .expect("current version should be stored");
+ current.transitioned_object.status = crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string();
+
+ let incarnation = ecstore
+ .bucket_incarnation_id_from_disk(&bucket)
+ .await
+ .expect("bucket incarnation should be available");
+ let deleted = super::apply_expiry_on_transitioned_object(
+ ecstore.clone(),
+ ¤t,
+ &crate::bucket::lifecycle::lifecycle::Event {
+ action: IlmAction::DeleteAllVersionsAction,
+ rule_id: "delete-all-versions".to_string(),
+ ..Default::default()
+ },
+ &LcEventSrc::Scanner,
+ incarnation,
+ )
+ .await;
+
+ assert!(deleted, "delete-all must not degrade to transitioned single-version expiry");
+ let versions = ecstore
+ .list_object_versions(&bucket, object, None, None, None, 10)
+ .await
+ .expect("purged versions should be listable");
+ assert_eq!(versions.objects.iter().filter(|version| version.name == object).count(), 0);
+ }
+
#[tokio::test]
async fn existing_object_lifecycle_skips_current_expiration_for_explicit_legal_hold() {
let lc = latest_expiration_lifecycle();
diff --git a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs
index ee05be7c7..6308b3767 100644
--- a/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs
+++ b/crates/ecstore/src/bucket/lifecycle/tier_delete_journal.rs
@@ -435,12 +435,88 @@ async fn process_committed_tier_delete_journal_entry(api: Arc, je: &Jen
remove_tier_delete_journal_entry(api, je).await
}
-async fn reconcile_prepared_tier_delete_journal_entry(api: Arc, je: &Jentry) -> std::io::Result<()> {
- let (data, metadata) =
- config_boundary::read_config_with_metadata(api.clone(), &tier_delete_journal_object_name(je), &ObjectOptions::default())
+fn object_info_references_tier_delete(info: &ObjectInfo, je: &Jentry) -> std::io::Result {
+ if info.transitioned_object.status != rustfs_filemeta::TRANSITION_COMPLETE
+ || info.transitioned_object.name != je.obj_name
+ || info.transitioned_object.tier != je.tier_name
+ {
+ return Ok(false);
+ }
+ let source_backend_identity = tier_destination_id_from_metadata(&info.user_defined)?;
+ if source_backend_identity.is_some() && source_backend_identity != je.backend_identity {
+ return Ok(false);
+ }
+ if !je.version_id_exact {
+ return Ok(true);
+ }
+ Ok(match info.transition_version_state {
+ rustfs_filemeta::TransitionVersionState::Unknown => true,
+ rustfs_filemeta::TransitionVersionState::KnownDisabled => false,
+ rustfs_filemeta::TransitionVersionState::SuspendedNull | rustfs_filemeta::TransitionVersionState::Exact => {
+ info.transitioned_object.version_id == je.version_id
+ }
+ })
+}
+
+async fn prepared_tier_delete_has_live_source(
+ api: &ECStore,
+ source: &TierDeleteSourceIdentity,
+ je: &Jentry,
+) -> std::io::Result<(bool, Vec)> {
+ let lock_object = rustfs_utils::path::encode_dir_object(&source.object);
+ let mut lock_opts = ObjectOptions::default();
+ let read_guards = api
+ .acquire_all_object_read_locks("tier_delete_journal_recovery", &source.bucket, &lock_object, &mut lock_opts)
+ .await
+ .map_err(std::io::Error::other)?;
+ if api.ctx.lock_manager().is_disabled() {
+ return Err(std::io::Error::new(
+ std::io::ErrorKind::WouldBlock,
+ "tier delete journal recovery requires namespace locking",
+ ));
+ }
+ let mut has_live_source = false;
+ for pool in &api.pools {
+ let set = pool.get_disks_by_key(&lock_object);
+ let Some(versions) = set
+ .load_file_info_versions_exact(&source.bucket, &source.object)
.await
- .map_err(std::io::Error::other)?;
+ .map_err(std::io::Error::other)?
+ else {
+ continue;
+ };
+ for version in versions.versions.iter().filter(|version| !version.tier_free_version()) {
+ let info = ObjectInfo::from_file_info(version, &source.bucket, &source.object, source.versioned);
+ if object_info_references_tier_delete(&info, je)? {
+ has_live_source = true;
+ break;
+ }
+ }
+ if has_live_source {
+ break;
+ }
+ }
+ if read_guards.iter().any(crate::store::ObjectLockDiagGuard::is_lock_lost) {
+ return Err(std::io::Error::new(
+ std::io::ErrorKind::WouldBlock,
+ "tier delete journal recovery object read lock was lost",
+ ));
+ }
+ Ok((has_live_source, read_guards))
+}
+
+async fn reconcile_prepared_tier_delete_journal_entry(api: Arc, je: &Jentry) -> std::io::Result<()> {
+ let journal_name = tier_delete_journal_object_name(je);
+ let (data, metadata) = config_boundary::read_config_with_metadata(api.clone(), &journal_name, &ObjectOptions::default())
+ .await
+ .map_err(std::io::Error::other)?;
let current = decode_tier_delete_journal_entry(&data).map_err(std::io::Error::other)?;
+ if tier_delete_journal_object_name(¤t) != journal_name {
+ return Err(std::io::Error::new(
+ std::io::ErrorKind::InvalidData,
+ "prepared tier delete journal content does not match its object name",
+ ));
+ }
if current.state != TierDeleteJournalState::Prepared {
return Err(std::io::Error::new(
std::io::ErrorKind::WouldBlock,
@@ -453,16 +529,21 @@ async fn reconcile_prepared_tier_delete_journal_entry(api: Arc, je: &Je
"prepared tier delete journal has no entity tag",
));
};
- let source = je
+ let source = current
.source
.as_ref()
.ok_or_else(|| std::io::Error::new(std::io::ErrorKind::InvalidData, "prepared tier delete journal has no source"))?;
- match api
- .get_object_info(&source.bucket, &source.object, &source.lookup_options())
- .await
- {
- Ok(info) if source.matches(&info) => {
- match config_boundary::delete_config_if_match(api, &tier_delete_journal_object_name(¤t), &etag).await {
+ match prepared_tier_delete_has_live_source(&api, source, ¤t).await {
+ Ok((true, read_guards)) => {
+ if read_guards.iter().any(crate::store::ObjectLockDiagGuard::is_lock_lost) {
+ return Err(std::io::Error::new(
+ std::io::ErrorKind::WouldBlock,
+ "tier delete journal recovery object read lock was lost before abort",
+ ));
+ }
+ let result = config_boundary::delete_config_if_match(api, &tier_delete_journal_object_name(¤t), &etag).await;
+ drop(read_guards);
+ match result {
Ok(()) => Ok(()),
Err(Error::PreconditionFailed) => Err(std::io::Error::new(
std::io::ErrorKind::WouldBlock,
@@ -471,17 +552,32 @@ async fn reconcile_prepared_tier_delete_journal_entry(api: Arc, je: &Je
Err(err) => Err(std::io::Error::other(err)),
}
}
- Ok(_info) if source.has_stable_identity() => {
- commit_prepared_tier_delete_journal_entry_if_current(api, current, etag).await
+ Ok((false, read_guards)) if source.has_stable_identity() => {
+ if read_guards.iter().any(crate::store::ObjectLockDiagGuard::is_lock_lost) {
+ return Err(std::io::Error::new(
+ std::io::ErrorKind::WouldBlock,
+ "tier delete journal recovery object read lock was lost before commit",
+ ));
+ }
+ let mut commit_opts = ObjectOptions::default();
+ for signal in read_guards
+ .iter()
+ .filter_map(crate::store::ObjectLockDiagGuard::lock_lost_signal)
+ {
+ commit_opts.add_namespace_lock_lost_signal(signal);
+ }
+ let committed =
+ commit_prepared_tier_delete_journal_entry_if_current(api.clone(), current, etag, &commit_opts).await?;
+ // Keep namespace locks only through the journal CAS. Remote-tier IO
+ // must not block writers for the object during recovery.
+ drop(read_guards);
+ process_committed_tier_delete_journal_entry(api, &committed).await
}
- Ok(_) => Err(std::io::Error::new(
+ Ok((false, _read_guards)) => Err(std::io::Error::new(
std::io::ErrorKind::WouldBlock,
"prepared tier delete journal source identity is not sufficient to confirm deletion",
)),
- Err(Error::ObjectNotFound(_, _)) | Err(Error::FileNotFound) | Err(Error::FileVersionNotFound) => {
- commit_prepared_tier_delete_journal_entry_if_current(api, current, etag).await
- }
- Err(err) => Err(std::io::Error::other(err)),
+ Err(err) => Err(err),
}
}
@@ -489,7 +585,8 @@ async fn commit_prepared_tier_delete_journal_entry_if_current(
api: Arc,
mut committed: Jentry,
etag: String,
-) -> std::io::Result<()> {
+ lock_opts: &ObjectOptions,
+) -> std::io::Result {
committed.state = TierDeleteJournalState::Committed;
let data = encode_tier_delete_journal_entry(&committed).map_err(std::io::Error::other)?;
match config_boundary::save_config_with_opts(
@@ -502,12 +599,13 @@ async fn commit_prepared_tier_delete_journal_entry_if_current(
if_match: Some(etag),
..Default::default()
}),
+ namespace_lock_fence: lock_opts.namespace_lock_fence.clone(),
..Default::default()
},
)
.await
{
- Ok(()) => process_committed_tier_delete_journal_entry(api, &committed).await,
+ Ok(()) => Ok(committed),
Err(Error::PreconditionFailed) => Err(std::io::Error::new(
std::io::ErrorKind::WouldBlock,
"prepared tier delete journal changed before commit",
@@ -582,6 +680,18 @@ pub async fn recover_tier_delete_journal_entries(
}
};
+ if tier_delete_journal_object_name(&je) != object.name {
+ stats.failed += 1;
+ warn!(
+ event = EVENT_LIFECYCLE_TIER_DELETE_JOURNAL,
+ component = LOG_COMPONENT_ECSTORE,
+ subsystem = LOG_SUBSYSTEM_LIFECYCLE,
+ journal_object = %object.name,
+ "Tier delete journal content does not match its object name and will be retained"
+ );
+ continue;
+ }
+
if je.backend_identity.is_none() {
stats.failed += 1;
warn!(
@@ -699,16 +809,14 @@ where
mod tests {
use super::{
TIER_DELETE_JOURNAL_EXACT_VERSION, TIER_DELETE_JOURNAL_STATE_VERSION, await_tier_delete_journal_recovery,
- decode_tier_delete_journal_entry, encode_tier_delete_journal_entry, record_tier_delete_journal_backend_identity,
- tier_delete_journal_object_name,
+ decode_tier_delete_journal_entry, encode_tier_delete_journal_entry, object_info_references_tier_delete,
+ record_tier_delete_journal_backend_identity, tier_delete_journal_object_name,
};
use crate::bucket::lifecycle::tier_sweeper::{Jentry, TierDeleteJournalState, TierDeleteSourceIdentity};
use crate::error::Result;
use crate::object_api::ObjectInfo;
use std::time::Duration;
- use time::OffsetDateTime;
use tokio_util::sync::CancellationToken;
- use uuid::Uuid;
fn journal_entry() -> Jentry {
Jentry {
@@ -738,6 +846,16 @@ mod tests {
assert_eq!(decoded.version_state, je.version_state);
}
+ #[test]
+ fn tier_delete_journal_object_name_binds_persisted_content() {
+ let original = journal_entry();
+ let original_name = tier_delete_journal_object_name(&original);
+ let mut replaced = original;
+ replaced.obj_name = "remote/replaced".to_string();
+
+ assert_ne!(tier_delete_journal_object_name(&replaced), original_name);
+ }
+
#[test]
fn tier_delete_transaction_roundtrips_prepared_source_identity() {
let mut je = journal_entry();
@@ -765,26 +883,35 @@ mod tests {
}
#[test]
- fn tier_delete_source_identity_rejects_recreated_object() {
- let version_id = Uuid::from_u128(1);
- let data_dir = Uuid::from_u128(2);
- let mod_time = OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(1);
- let info = ObjectInfo {
- bucket: "bucket".to_string(),
- name: "object".to_string(),
- version_id: Some(version_id),
- data_dir: Some(data_dir),
- mod_time: Some(mod_time),
+ fn prepared_recovery_blocks_any_live_reference_to_the_remote_version() {
+ let je = journal_entry();
+ let mut metadata = std::collections::HashMap::new();
+ rustfs_utils::http::metadata_compat::insert_str(
+ &mut metadata,
+ rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
+ rustfs_utils::crypto::hex(je.backend_identity.expect("test journal should bind a backend")),
+ );
+ let mut info = ObjectInfo {
+ user_defined: std::sync::Arc::new(metadata),
+ transitioned_object: crate::storage_api_contracts::lifecycle::TransitionedObject {
+ name: je.obj_name.clone(),
+ version_id: je.version_id.clone(),
+ tier: je.tier_name.clone(),
+ status: rustfs_filemeta::TRANSITION_COMPLETE.to_string(),
+ ..Default::default()
+ },
+ transition_version_state: rustfs_filemeta::TransitionVersionState::Exact,
..Default::default()
};
- let source = TierDeleteSourceIdentity::from_object_info("bucket", "object", &info, true, false);
- assert!(source.matches(&info));
- let recreated = ObjectInfo {
- data_dir: Some(Uuid::from_u128(3)),
- ..info
- };
- assert!(!source.matches(&recreated));
+ assert!(object_info_references_tier_delete(&info, &je).expect("matching reference should be valid"));
+ info.transitioned_object.version_id = "other-version".to_string();
+ assert!(!object_info_references_tier_delete(&info, &je).expect("different exact version should be valid"));
+ info.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown;
+ assert!(
+ object_info_references_tier_delete(&info, &je).expect("legacy unknown reference should fail closed"),
+ "an unknown live source may still reference the journaled remote version"
+ );
}
#[test]
diff --git a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs
index 0bde1a23c..b2786fb18 100644
--- a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs
+++ b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs
@@ -334,32 +334,6 @@ impl TierDeleteSourceIdentity {
}
}
- pub(crate) fn lookup_options(&self) -> crate::object_api::ObjectOptions {
- crate::object_api::ObjectOptions {
- version_id: self.version_id.clone(),
- versioned: self.versioned,
- version_suspended: self.version_suspended,
- ..Default::default()
- }
- }
-
- pub(crate) fn matches(&self, info: &ObjectInfo) -> bool {
- if self.bucket != info.bucket {
- return false;
- }
- if let Some(version_id) = &self.version_id {
- return info.version_id.map(|id| id.to_string()).as_deref() == Some(version_id.as_str())
- && self.data_dir == info.data_dir.map(|id| id.to_string());
- }
- if self.data_dir.is_some() {
- return self.data_dir == info.data_dir.map(|id| id.to_string());
- }
- self.etag.is_some()
- && self.etag == info.etag
- && self.mod_time.is_some()
- && self.mod_time == info.mod_time.map(|time| time.to_string())
- }
-
pub(crate) fn has_stable_identity(&self) -> bool {
self.version_id.is_some() || self.data_dir.is_some() || (self.etag.is_some() && self.mod_time.is_some())
}
diff --git a/crates/ecstore/src/bucket/replication/replication_lifecycle_bridge.rs b/crates/ecstore/src/bucket/replication/replication_lifecycle_bridge.rs
index c35bf99c1..a89e89519 100644
--- a/crates/ecstore/src/bucket/replication/replication_lifecycle_bridge.rs
+++ b/crates/ecstore/src/bucket/replication/replication_lifecycle_bridge.rs
@@ -23,6 +23,8 @@ use super::replication_queue_boundary::DeletedObjectReplicationInfo;
use super::replication_storage_boundary::{
DeletedObject, ObjectInfo, ObjectOptions, ObjectToDelete, deleted_object_for_replication,
};
+#[cfg(test)]
+use std::sync::Mutex;
#[allow(
dead_code,
@@ -32,6 +34,9 @@ pub(crate) type ReplicationLifecycleConfig = ReplicationConfig;
pub(crate) struct ReplicationLifecycleBridge;
+#[cfg(test)]
+static SCHEDULED_DELETE_OBJECTS: Mutex> = Mutex::new(Vec::new());
+
impl ReplicationLifecycleBridge {
#[allow(
dead_code,
@@ -85,6 +90,13 @@ impl ReplicationLifecycleBridge {
}
pub(crate) async fn schedule_delete(bucket: String, delete_object: DeletedObject) {
+ #[cfg(test)]
+ {
+ SCHEDULED_DELETE_OBJECTS
+ .lock()
+ .expect("scheduled delete test hook lock should not poison")
+ .push(delete_object.clone());
+ }
super::replication_pool::schedule_replication_delete(DeletedObjectReplicationInfo {
delete_object: deleted_object_for_replication(delete_object),
bucket,
@@ -93,6 +105,15 @@ impl ReplicationLifecycleBridge {
})
.await;
}
+
+ #[cfg(test)]
+ pub(crate) fn take_scheduled_deletes_for_test() -> Vec {
+ std::mem::take(
+ &mut *SCHEDULED_DELETE_OBJECTS
+ .lock()
+ .expect("scheduled delete test hook lock should not poison"),
+ )
+ }
}
#[cfg(test)]
diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs
index 99ec6c83e..1b97a6df1 100644
--- a/crates/ecstore/src/object_api/types.rs
+++ b/crates/ecstore/src/object_api/types.rs
@@ -218,6 +218,63 @@ pub struct QuotaAdmission {
quota_limit: u64,
}
+#[doc(hidden)]
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub struct LifecycleDeleteAllRequest {
+ pub(crate) version_id: Option,
+ pub(crate) delete_marker: bool,
+ pub(crate) action: rustfs_common::metrics::IlmAction,
+ pub(crate) rule_id: String,
+ pub(crate) phase: LifecycleDeleteAllPhase,
+}
+
+#[doc(hidden)]
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum LifecycleDeleteAllPhase {
+ Preflight,
+ History,
+ FinalPreflight,
+ Trigger,
+}
+
+#[doc(hidden)]
+#[derive(Default)]
+pub struct LifecycleDeleteAllJournalState {
+ prepared: HashMap,
+ mutation_started: bool,
+}
+
+impl Debug for LifecycleDeleteAllJournalState {
+ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ f.debug_struct("LifecycleDeleteAllJournalState")
+ .field("prepared_count", &self.prepared.len())
+ .field("mutation_started", &self.mutation_started)
+ .finish()
+ }
+}
+
+impl LifecycleDeleteAllJournalState {
+ pub(crate) fn contains(&self, name: &str) -> bool {
+ self.prepared.contains_key(name)
+ }
+
+ pub(crate) fn insert(&mut self, name: String, entry: crate::bucket::lifecycle::tier_sweeper::Jentry) {
+ self.prepared.insert(name, entry);
+ }
+
+ pub(crate) fn prepared_entries(&self) -> Vec {
+ self.prepared.values().cloned().collect()
+ }
+
+ pub(crate) fn mark_mutation_started(&mut self) {
+ self.mutation_started = true;
+ }
+
+ pub(crate) fn mutation_started(&self) -> bool {
+ self.mutation_started
+ }
+}
+
impl QuotaAdmission {
pub(crate) fn current_usage(self) -> u64 {
self.current_usage
@@ -242,6 +299,11 @@ pub struct ObjectOptions {
pub delete_prefix: bool,
pub delete_prefix_object: bool,
pub version_id: Option,
+ /// Lifecycle-only staged purge request checked under the object write lock.
+ #[doc(hidden)]
+ pub lifecycle_delete_all: Option,
+ #[doc(hidden)]
+ pub lifecycle_delete_all_journal: Option>>,
/// RustFS-only compare-and-set condition checked under the object write lock.
pub expected_current_version_id: Option,
/// Persisted bucket incarnation observed before authorization.
@@ -349,6 +411,15 @@ impl ObjectOptions {
self.namespace_lock_fence.get_or_insert_with(NamespaceLockFence::new);
}
+ pub(crate) fn ensure_lifecycle_delete_all_journal(&mut self) {
+ self.lifecycle_delete_all_journal
+ .get_or_insert_with(|| Arc::new(parking_lot::Mutex::new(LifecycleDeleteAllJournalState::default())));
+ }
+
+ pub(crate) fn lifecycle_delete_all_journal(&self) -> Option<&Arc>> {
+ self.lifecycle_delete_all_journal.as_ref()
+ }
+
pub fn add_namespace_lock_guard(&mut self, guard: &rustfs_lock::NamespaceLockGuard) {
if let Some(signal) = guard.lock_lost_signal() {
self.add_namespace_lock_lost_signal(signal);
@@ -1890,4 +1961,13 @@ mod tests {
assert!(default_cloned.user_tags.is_empty());
assert!(default_cloned.parts.is_empty());
}
+
+ #[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());
+ opts.ensure_lifecycle_delete_all_journal();
+ assert!(opts.lifecycle_delete_all_journal().is_some());
+ }
}
diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs
index 5a56ea36e..d4c535461 100644
--- a/crates/ecstore/src/set_disk/mod.rs
+++ b/crates/ecstore/src/set_disk/mod.rs
@@ -4578,10 +4578,15 @@ impl SetDisks {
)?;
let fi = build_tiered_decommission_file_info(bucket, object, fi, layout);
let write_quorum = layout.write_quorum;
- if opts
- .bucket_lifecycle_lock_fence
- .as_ref()
- .is_some_and(NamespaceLockFence::is_lock_lost)
+ if _lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
+ || opts
+ .namespace_lock_fence
+ .as_ref()
+ .is_some_and(NamespaceLockFence::is_lock_lost)
+ || opts
+ .bucket_lifecycle_lock_fence
+ .as_ref()
+ .is_some_and(NamespaceLockFence::is_lock_lost)
|| bucket_lifecycle_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs
index 929ac4d5c..5430b5d1d 100644
--- a/crates/ecstore/src/set_disk/ops/object.rs
+++ b/crates/ecstore/src/set_disk/ops/object.rs
@@ -27,12 +27,12 @@ use crate::set_disk::read::GetObjectDownstreamWriter;
use crate::bucket::lifecycle::{
tier_delete_journal::{
enqueue_committed_tier_delete_journal_entry, persist_tier_delete_journal_entry,
- record_tier_delete_journal_backend_identity, remove_tier_delete_journal_entry,
+ record_tier_delete_journal_backend_identity, remove_tier_delete_journal_entry, tier_delete_journal_object_name,
},
tier_sweeper::{
- Jentry, RemoteTierDeleteOutcome, TierDeleteJournalState,
+ Jentry, RemoteTierDeleteOutcome, TierDeleteJournalState, attach_tier_delete_source,
delete_confirmed_transition_candidate_exact_with_lease_idempotent, delete_object_from_remote_tier_with_lease_idempotent,
- transitioned_delete_journal_entry_for_source,
+ transitioned_delete_journal_entry_for_source, transitioned_force_delete_journal_entry,
},
transition_transaction::{
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
@@ -42,7 +42,8 @@ use crate::bucket::lifecycle::{
};
use crate::bucket::quota::reservation;
use crate::bucket::replication::{
- DeleteReplicationConfigSnapshot, VersionPurgeStatusType, replication_state_to_filemeta, version_purge_status_to_filemeta,
+ DeleteReplicationConfigSnapshot, ReplicationLifecycleBridge, ReplicationStatusType, VersionPurgeStatusType,
+ replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_to_filemeta,
};
use crate::diagnostics::get::GetObjectFailureReason;
use crate::disk::{DataDirDeleteStatus, OldCurrentSize};
@@ -97,6 +98,636 @@ fn duration_millis_f64(duration: std::time::Duration) -> f64 {
duration.as_secs_f64() * 1000.0
}
+struct LifecycleDeleteAllPlan<'a> {
+ history: Vec<&'a FileInfo>,
+ trigger: Option<&'a FileInfo>,
+}
+
+impl<'a> LifecycleDeleteAllPlan<'a> {
+ fn trigger_only(&self) -> Result