diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 21483e742..0bed2372a 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -45,7 +45,7 @@ use lazy_static::lazy_static; use rustfs_common::data_usage::TierStats; use rustfs_common::heal_channel::rep_has_active_rules; use rustfs_common::metrics::{IlmAction, Metrics}; -use rustfs_filemeta::{NULL_VERSION_ID, RestoreStatusOps, is_restored_object_on_disk}; +use rustfs_filemeta::{FileInfo, NULL_VERSION_ID, RestoreStatusOps, is_restored_object_on_disk}; use rustfs_s3_common::EventName; use rustfs_utils::{get_env_i64, get_env_usize, path::encode_dir_object, string::strings_has_prefix_fold}; use s3s::Body; @@ -393,12 +393,81 @@ impl ExpiryState { //delete_object_versions(api, &v.bucket, &v.versions, v.event).await; } else if v.as_any().is::() { - //transitionLogIf(es.ctx, deleteObjectFromRemoteTier(es.ctx, v.ObjName, v.VersionID, v.TierName)) + let v = v.as_any().downcast_ref::().expect("err!"); + if let Err(err) = delete_object_from_remote_tier(&v.obj_name, &v.version_id, &v.tier_name).await { + warn!( + object = %v.obj_name, + version_id = %v.version_id, + tier = %v.tier_name, + error = ?err, + "failed to delete transitioned object from remote tier" + ); + } } else if v.as_any().is::() { let v = v.as_any().downcast_ref::().expect("err!"); - let _oi = v.0.clone(); + let oi = v.0.clone(); + if let Err(err) = delete_object_from_remote_tier( + &oi.transitioned_object.name, + &oi.transitioned_object.version_id, + &oi.transitioned_object.tier, + ) + .await + { + warn!( + bucket = %oi.bucket, + object = %oi.name, + remote_object = %oi.transitioned_object.name, + remote_version_id = %oi.transitioned_object.version_id, + tier = %oi.transitioned_object.tier, + error = ?err, + "failed to sweep transitioned free version from remote tier" + ); + continue; + } + let mut fi = FileInfo { + name: oi.name.clone(), + version_id: oi.version_id, + deleted: true, + ..Default::default() + }; + fi.set_tier_free_version(); + + let mut deleted_locally = false; + for pool in api.pools.iter() { + let set = pool.get_disks_by_key(&oi.name); + match set.delete_object_version(&oi.bucket, &oi.name, &fi, false).await { + Ok(()) => { + deleted_locally = true; + break; + } + Err(err) if is_err_version_not_found(&err) || is_err_object_not_found(&err) => continue, + Err(err) => { + warn!( + bucket = %oi.bucket, + object = %oi.name, + remote_object = %oi.transitioned_object.name, + remote_version_id = %oi.transitioned_object.version_id, + tier = %oi.transitioned_object.tier, + error = ?err, + "failed to delete transitioned free version after remote tier sweep" + ); + break; + } + } + } + + if !deleted_locally { + warn!( + bucket = %oi.bucket, + object = %oi.name, + remote_object = %oi.transitioned_object.name, + remote_version_id = %oi.transitioned_object.version_id, + tier = %oi.transitioned_object.tier, + "transitioned free version was not found during local cleanup" + ); + } } else { //info!("Invalid work type - {:?}", v); diff --git a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs index 8c905776e..896ec8d1f 100644 --- a/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs +++ b/crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs @@ -120,9 +120,9 @@ impl ObjSweeper { #[derive(Debug, Clone)] #[allow(unused_assignments)] pub struct Jentry { - obj_name: String, - version_id: String, - tier_name: String, + pub(crate) obj_name: String, + pub(crate) version_id: String, + pub(crate) tier_name: String, } impl ExpiryOp for Jentry { @@ -147,5 +147,37 @@ pub async fn delete_object_from_remote_tier(obj_name: &str, rv_id: &str, tier_na w.remove(obj_name, rv_id).await } +pub fn transitioned_delete_journal_entry( + version_id: Option, + versioned: bool, + suspended: bool, + transitioned: &TransitionedObject, +) -> Option { + let sweeper = ObjSweeper { + version_id, + versioned, + suspended, + transition_status: transitioned.status.clone(), + transition_tier: transitioned.tier.clone(), + transition_version_id: transitioned.version_id.clone(), + remote_object: transitioned.name.clone(), + ..Default::default() + }; + + sweeper.should_remove_remote_object() +} + +pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject) -> Option { + if transitioned.status != lifecycle::TRANSITION_COMPLETE { + return None; + } + + Some(Jentry { + obj_name: transitioned.name.clone(), + version_id: transitioned.version_id.clone(), + tier_name: transitioned.tier.clone(), + }) +} + #[cfg(test)] mod test {} diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 1237d8cc6..4695f2233 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -1077,7 +1077,7 @@ impl ObjectOperations for SetDisks { let (mut metas, errs) = { if let Some(vid) = &src_opts.version_id { - Self::read_all_fileinfo(&disks, "", src_bucket, src_object, vid, true, false).await? + Self::read_all_fileinfo(&disks, "", src_bucket, src_object, vid, true, false, false).await? } else { Self::read_all_xl(&disks, src_bucket, src_object, true, false).await } @@ -1565,8 +1565,6 @@ impl ObjectOperations for SetDisks { return Ok(oi); } - let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok()); - // Create a single object deletion request let mut dfi = FileInfo { name: object.to_string(), @@ -1673,7 +1671,7 @@ impl ObjectOperations for SetDisks { let (metas, errs) = { if let Some(version_id) = &opts.version_id { - Self::read_all_fileinfo(&disks, "", bucket, object, version_id.to_string().as_str(), false, false).await? + Self::read_all_fileinfo(&disks, "", bucket, object, version_id.to_string().as_str(), false, false, false).await? } else { Self::read_all_xl(&disks, bucket, object, false, false).await } @@ -3279,7 +3277,7 @@ impl HealOperations for SetDisks { let disks = self.disks.read().await; let disks = disks.clone(); - let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false).await?; + let (_, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, false, false, false).await?; if DiskError::is_all_not_found(&errs) { warn!( "heal_object failed, all obj part not found, bucket: {}, obj: {}, version_id: {}", diff --git a/crates/ecstore/src/set_disk/heal.rs b/crates/ecstore/src/set_disk/heal.rs index 3418f6d29..ae672bac7 100644 --- a/crates/ecstore/src/set_disk/heal.rs +++ b/crates/ecstore/src/set_disk/heal.rs @@ -56,7 +56,8 @@ impl SetDisks { } }; - let (mut parts_metadata, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, version_id, true, true).await?; + let (mut parts_metadata, errs) = + Self::read_all_fileinfo(&disks, "", bucket, object, version_id, true, true, false).await?; info!( parts_count = parts_metadata.len(), diff --git a/crates/ecstore/src/set_disk/multipart.rs b/crates/ecstore/src/set_disk/multipart.rs index 0efb96f26..491180827 100644 --- a/crates/ecstore/src/set_disk/multipart.rs +++ b/crates/ecstore/src/set_disk/multipart.rs @@ -105,7 +105,8 @@ impl SetDisks { let disks = disks.clone(); let (parts_metadata, errs) = - Self::read_all_fileinfo(&disks, bucket, RUSTFS_META_MULTIPART_BUCKET, &upload_id_path, "", false, false).await?; + Self::read_all_fileinfo(&disks, bucket, RUSTFS_META_MULTIPART_BUCKET, &upload_id_path, "", false, false, false) + .await?; let map_err_notfound = |err: DiskError| { if err == DiskError::FileNotFound { diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index e8306f91c..f377a7593 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -131,6 +131,7 @@ impl SetDisks { Ok(ret) } + #[allow(clippy::too_many_arguments)] #[tracing::instrument(level = "debug", skip(disks))] pub(super) async fn read_all_fileinfo( disks: &[Option], @@ -140,13 +141,14 @@ impl SetDisks { version_id: &str, read_data: bool, healing: bool, + incl_free_versions: bool, ) -> disk::error::Result<(Vec, Vec>)> { let mut ress = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len()); let opts = Arc::new(ReadOptions { + incl_free_versions, read_data, healing, - ..Default::default() }); let org_bucket = Arc::new(org_bucket.to_string()); let bucket = Arc::new(bucket.to_string()); @@ -474,7 +476,8 @@ impl SetDisks { let vid = opts.version_id.clone().unwrap_or_default(); // TODO: optimize concurrency and break once enough slots are available - let (parts_metadata, errs) = Self::read_all_fileinfo(&disks, "", bucket, object, vid.as_str(), read_data, false).await?; + let (parts_metadata, errs) = + Self::read_all_fileinfo(&disks, "", bucket, object, vid.as_str(), read_data, false, opts.incl_free_versions).await?; // warn!("get_object_fileinfo parts_metadata {:?}", &parts_metadata); // warn!("get_object_fileinfo {}/{} errs {:?}", bucket, object, &errs); @@ -541,6 +544,9 @@ impl SetDisks { } if fi.deleted { + if opts.incl_free_versions && fi.tier_free_version() && opts.version_id.is_some() { + return (oi, write_quorum, None); + } return if opts.version_id.is_none() || opts.delete_marker { (oi, write_quorum, Some(to_object_err(StorageError::FileNotFound, vec![bucket, object]))) } else { diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 63dda34e1..a964556ea 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -286,7 +286,9 @@ impl ECStore { } if !errs.is_empty() && !opts.versioned && !opts.version_suspended { - return self.delete_object_from_all_pools(bucket, object, &opts, errs).await; + let mut obj = self.delete_object_from_all_pools(bucket, object, &opts, errs).await?; + obj.name = decode_dir_object(object); + return Ok(obj); } for pool in self.pools.iter() { diff --git a/crates/ecstore/src/store_api/types.rs b/crates/ecstore/src/store_api/types.rs index 9acc239df..829f81ce5 100644 --- a/crates/ecstore/src/store_api/types.rs +++ b/crates/ecstore/src/store_api/types.rs @@ -47,6 +47,7 @@ pub struct ObjectOptions { pub versioned: bool, pub version_suspended: bool, + pub incl_free_versions: bool, pub skip_decommissioned: bool, pub skip_rebalancing: bool, diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 3b24bfc1c..35f1352f6 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -23,6 +23,7 @@ use rand::seq::SliceRandom as _; use rustfs_common::heal_channel::HealScanMode; use rustfs_common::metrics::{Metric, Metrics, emit_scan_bucket_drive_complete}; use rustfs_ecstore::bucket::bucket_target_sys::BucketTargetSys; +use rustfs_ecstore::bucket::lifecycle::bucket_lifecycle_ops::GLOBAL_ExpiryState; use rustfs_ecstore::bucket::lifecycle::lifecycle::Lifecycle; use rustfs_ecstore::bucket::metadata_sys::{get_lifecycle_config, get_object_lock_config, get_replication_config}; use rustfs_ecstore::bucket::replication::{ReplicationConfig, ReplicationConfigurationExt}; @@ -530,6 +531,11 @@ impl ScannerIODisk for Disk { .iter() .map(|v| ObjectInfo::from_file_info(v, item.bucket.as_str(), item.object_path().as_str(), versioned)) .collect::>(); + let free_version_infos = fivs + .free_versions + .iter() + .map(|v| ObjectInfo::from_file_info(v, item.bucket.as_str(), item.object_path().as_str(), versioned)) + .collect::>(); let mut size_summary = SizeSummary::default(); @@ -563,9 +569,14 @@ impl ScannerIODisk for Disk { item.apply_actions(ecstore, object_infos, lock_config, &mut size_summary) .await; - done_object(); + if !free_version_infos.is_empty() { + let mut expiry_state = GLOBAL_ExpiryState.write().await; + for oi in free_version_infos { + expiry_state.enqueue_free_version(oi).await; + } + } - // TODO: enqueueFreeVersion + done_object(); Ok(size_summary) } diff --git a/crates/scanner/tests/lifecycle_integration_test.rs b/crates/scanner/tests/lifecycle_integration_test.rs index a3264edd0..b8e8304f7 100644 --- a/crates/scanner/tests/lifecycle_integration_test.rs +++ b/crates/scanner/tests/lifecycle_integration_test.rs @@ -18,8 +18,10 @@ use rustfs_ecstore::{ bucket::{lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects, metadata_sys}, client::transition_api::{ReadCloser, ReaderImpl}, disk::endpoint::Endpoint, + disk::{DiskAPI, DiskOption, STORAGE_FORMAT_FILE, new_disk}, endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}, global::GLOBAL_TierConfigMgr, + pools::path2_bucket_object_with_base_path, store::ECStore, store_api::{ BucketOperations, MakeBucketOptions, MultipartOperations, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, @@ -29,13 +31,17 @@ use rustfs_ecstore::{ warm_backend::{WarmBackend, WarmBackendGetOpts, build_transition_put_options}, }, }; +use rustfs_filemeta::FileMeta; use rustfs_scanner::scanner::init_data_scanner; +use rustfs_scanner::scanner_folder::ScannerItem; +use rustfs_scanner::scanner_io::ScannerIODisk; +use rustfs_utils::path::path_join_buf; use s3s::dto::RestoreRequest; use serial_test::serial; use std::{ collections::HashMap, io::Cursor, - path::PathBuf, + path::{Path, PathBuf}, sync::{Arc, Once, OnceLock}, time::Duration, }; @@ -137,6 +143,70 @@ async fn setup_test_env() -> (Vec, Arc) { (disk_paths, ecstore) } +async fn setup_isolated_test_env(init_expiry: bool) -> (Vec, Arc) { + init_tracing(); + + let test_base_dir = format!("/tmp/rustfs_scanner_lifecycle_test_{}", uuid::Uuid::new_v4()); + let temp_dir = std::path::PathBuf::from(&test_base_dir); + if temp_dir.exists() { + fs::remove_dir_all(&temp_dir).await.ok(); + } + fs::create_dir_all(&temp_dir).await.unwrap(); + + let disk_paths = vec![ + temp_dir.join("disk1"), + temp_dir.join("disk2"), + temp_dir.join("disk3"), + temp_dir.join("disk4"), + ]; + + for disk_path in &disk_paths { + fs::create_dir_all(disk_path).await.unwrap(); + } + + let mut endpoints = Vec::new(); + for (i, disk_path) in disk_paths.iter().enumerate() { + let mut endpoint = Endpoint::try_from(disk_path.to_str().unwrap()).unwrap(); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(i); + endpoints.push(endpoint); + } + + let pool_endpoints = PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 4, + endpoints: Endpoints::from(endpoints), + cmd_line: "test".to_string(), + platform: format!("OS: {} | Arch: {}", std::env::consts::OS, std::env::consts::ARCH), + }; + + let endpoint_pools = EndpointServerPools(vec![pool_endpoints]); + rustfs_ecstore::store::init_local_disks(endpoint_pools.clone()).await.unwrap(); + + let server_addr: std::net::SocketAddr = "127.0.0.1:0".parse().unwrap(); + let ecstore = ECStore::new(server_addr, endpoint_pools, CancellationToken::new()) + .await + .unwrap(); + + let buckets_list = ecstore + .list_bucket(&rustfs_ecstore::store_api::BucketOptions { + no_metadata: true, + ..Default::default() + }) + .await + .unwrap(); + let buckets = buckets_list.into_iter().map(|v| v.name).collect(); + rustfs_ecstore::bucket::metadata_sys::init_bucket_metadata_sys(ecstore.clone(), buckets).await; + + if init_expiry { + rustfs_ecstore::bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await; + } + + (disk_paths, ecstore) +} + /// Test helper: Create a test bucket #[allow(dead_code)] async fn create_test_bucket(ecstore: &Arc, bucket_name: &str) { @@ -374,6 +444,85 @@ async fn wait_for_object_absence(ecstore: &Arc, bucket: &str, object: & } } +async fn wait_for_remote_absence(backend: &MockWarmBackend, object: &str, timeout: Duration) -> bool { + let deadline = tokio::time::Instant::now() + timeout; + + loop { + if !backend.objects.lock().await.contains_key(object) { + return true; + } + + if tokio::time::Instant::now() >= deadline { + return false; + } + + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn free_version_count(disk_path: &Path, bucket: &str, object: &str) -> usize { + let mut endpoint = Endpoint::try_from(disk_path.to_str().unwrap()).unwrap(); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(0); + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("failed to open local disk"); + let data = disk + .read_metadata(bucket, &path_join_buf(&[object, STORAGE_FORMAT_FILE])) + .await; + let Ok(data) = data else { + return 0; + }; + let meta = FileMeta::load(&data).expect("failed to load file metadata"); + meta.get_file_info_versions(bucket, object, false) + .expect("failed to decode file info versions") + .free_versions + .len() +} + +async fn scan_object_metadata(disk_path: &Path, bucket: &str, object: &str) { + let mut endpoint = Endpoint::try_from(disk_path.to_str().unwrap()).unwrap(); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(0); + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("failed to open local disk"); + let metadata_path = disk_path.join(bucket).join(object).join(STORAGE_FORMAT_FILE); + let relative_path = metadata_path.to_string_lossy().to_string(); + let (_, scanner_path) = path2_bucket_object_with_base_path(disk_path.to_string_lossy().as_ref(), relative_path.as_str()); + let file_type = fs::metadata(&metadata_path) + .await + .expect("failed to stat object metadata") + .file_type(); + let item = ScannerItem { + path: scanner_path.clone(), + bucket: bucket.to_string(), + prefix: object.to_string(), + object_name: STORAGE_FORMAT_FILE.to_string(), + file_type, + lifecycle: None, + replication: None, + heal_enabled: false, + heal_bitrot: false, + debug: false, + }; + disk.get_size(item).await.expect("scanner get_size should succeed"); +} + #[derive(Clone, Default)] struct MockStoredObject { bytes: Vec, @@ -480,6 +629,15 @@ async fn register_mock_tier(tier_name: &str) -> MockWarmBackend { version: "v1".to_string(), tier_type: TierType::MinIO, name: tier_name.to_string(), + minio: Some(TierMinIO { + access_key: "minioadmin".to_string(), + secret_key: "minioadmin".to_string(), + bucket: "mock-tier".to_string(), + endpoint: "http://127.0.0.1:0".to_string(), + prefix: format!("mock/{}/", Uuid::new_v4()), + region: String::new(), + ..Default::default() + }), ..Default::default() }, ); @@ -913,4 +1071,64 @@ mod serial_tests { .expect("Failed to consume restored object stream"); assert_eq!(data, expected); } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + #[serial] + #[ignore = "requires isolated global object layer state"] + async fn test_scanner_enqueues_free_version_cleanup_for_stale_transitioned_object() { + 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; + + let bucket_name = format!("test-scanner-free-version-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_name = "test/object.txt"; + let initial_payload = b"scanner should clean stale transitioned null version"; + create_test_bucket(&ecstore, bucket_name.as_str()).await; + set_bucket_lifecycle_transition_with_tier(bucket_name.as_str(), &tier_name) + .await + .expect("Failed to set lifecycle configuration"); + + upload_test_object(&ecstore, bucket_name.as_str(), object_name, initial_payload).await; + enqueue_transition_for_existing_objects(ecstore.clone(), bucket_name.as_str()) + .await + .expect("Failed to enqueue transitioned object"); + + let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should transition before overwrite"); + let stale_remote_object = transitioned.transitioned_object.name.clone(); + assert!(backend.objects.lock().await.contains_key(&stale_remote_object)); + + ecstore + .delete_object(bucket_name.as_str(), object_name, ObjectOptions::default()) + .await + .expect("Failed to delete transitioned object without expiry workers"); + + assert!( + free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0, + "deleting a transitioned null version should leave a free version for async cleanup" + ); + assert!( + backend.objects.lock().await.contains_key(&stale_remote_object), + "stale transitioned remote object should still exist before scanner fallback runs" + ); + + rustfs_ecstore::bucket::lifecycle::bucket_lifecycle_ops::init_background_expiry(ecstore.clone()).await; + scan_object_metadata(&disk_paths[0], bucket_name.as_str(), object_name).await; + + assert!( + wait_for_remote_absence(&backend, &stale_remote_object, TRANSITION_WAIT_TIMEOUT).await, + "scanner should enqueue stale free-version cleanup for the transitioned remote object" + ); + assert_eq!( + free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await, + 0, + "free-version metadata should be removed after scanner-triggered cleanup" + ); + assert!( + wait_for_object_absence(&ecstore, bucket_name.as_str(), object_name, Duration::from_secs(1)).await, + "deleted object should remain absent after scanner cleanup" + ); + } } diff --git a/rustfs/src/app/lifecycle_transition_api_test.rs b/rustfs/src/app/lifecycle_transition_api_test.rs index 74d13a45e..ae144b880 100644 --- a/rustfs/src/app/lifecycle_transition_api_test.rs +++ b/rustfs/src/app/lifecycle_transition_api_test.rs @@ -35,6 +35,7 @@ use rustfs_ecstore::{ warm_backend::{WarmBackend, WarmBackendGetOpts}, }, }; +use rustfs_utils::http::{SUFFIX_FORCE_DELETE, insert_header}; use s3s::{S3Request, dto::*}; use serial_test::serial; use std::{ @@ -141,12 +142,17 @@ async fn create_test_bucket(ecstore: &Arc, bucket_name: &str) { .expect("Failed to create test bucket"); } -async fn upload_test_object(ecstore: &Arc, bucket: &str, object: &str, data: &[u8]) { +async fn upload_test_object( + ecstore: &Arc, + bucket: &str, + object: &str, + data: &[u8], +) -> rustfs_ecstore::store_api::ObjectInfo { let mut reader = PutObjReader::from_vec(data.to_vec()); (**ecstore) .put_object(bucket, object, &mut reader, &ObjectOptions::default()) .await - .expect("Failed to upload test object"); + .expect("Failed to upload test object") } async fn set_bucket_lifecycle_transition_with_tier( @@ -282,6 +288,42 @@ async fn wait_for_transition( } } +async fn wait_for_remote_absence(backend: &MockWarmBackend, object: &str, timeout: Duration) -> bool { + let deadline = tokio::time::Instant::now() + timeout; + + loop { + if !backend.objects.lock().await.contains_key(object) { + return true; + } + + if tokio::time::Instant::now() >= deadline { + return false; + } + + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn wait_for_object_absence(ecstore: &Arc, bucket: &str, object: &str, timeout: Duration) -> bool { + let deadline = tokio::time::Instant::now() + timeout; + + loop { + if ecstore + .get_object_info(bucket, object, &ObjectOptions::default()) + .await + .is_err() + { + return true; + } + + if tokio::time::Instant::now() >= deadline { + return false; + } + + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + fn build_request(input: T, method: Method) -> S3Request { S3Request { input, @@ -353,7 +395,7 @@ async fn put_and_copy_object_transition_immediately_via_usecases() { set_bucket_lifecycle_transition_with_tier(dst_bucket.as_str(), &tier_name) .await .expect("Failed to set destination lifecycle configuration"); - upload_test_object(&ecstore, src_bucket.as_str(), src_object, copy_payload).await; + let _ = upload_test_object(&ecstore, src_bucket.as_str(), src_object, copy_payload).await; let copy_input = CopyObjectInput::builder() .copy_source(CopySource::Bucket { @@ -437,3 +479,63 @@ async fn complete_multipart_upload_transitions_immediately_via_usecase() { assert_eq!(info.transitioned_object.tier, tier_name); assert!(backend.objects.lock().await.contains_key(&info.transitioned_object.name)); } + +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +#[serial] +#[ignore = "requires isolated global object layer state"] +async fn delete_transitioned_object_removes_remote_tier_copy_via_usecase() { + let (_disk_paths, ecstore) = setup_test_env().await; + let usecase = DefaultObjectUsecase::without_context(); + + let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&tier_name).await; + + let bucket = format!("test-api-delete-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/object.txt"; + let payload = b"delete transitioned object through delete API"; + + create_test_bucket(&ecstore, bucket.as_str()).await; + set_bucket_lifecycle_transition_with_tier(bucket.as_str(), &tier_name) + .await + .expect("Failed to set lifecycle configuration"); + let _ = upload_test_object(&ecstore, bucket.as_str(), object, payload).await; + + rustfs_ecstore::bucket::lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects( + ecstore.clone(), + bucket.as_str(), + ) + .await + .expect("Failed to enqueue transitioned object"); + + let transitioned = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should transition before delete usecase runs"); + let remote_object = transitioned.transitioned_object.name.clone(); + + assert!(backend.objects.lock().await.contains_key(&remote_object)); + + let mut req = build_request( + DeleteObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .build() + .unwrap(), + Method::DELETE, + ); + insert_header(&mut req.headers, SUFFIX_FORCE_DELETE, "true"); + + usecase + .execute_delete_object(req) + .await + .expect("Failed to delete object through usecase"); + + assert!( + wait_for_object_absence(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT).await, + "object should be removed from hot tier after delete usecase" + ); + + assert!( + wait_for_remote_absence(&backend, &remote_object, TRANSITION_WAIT_TIMEOUT).await, + "transitioned object should be removed from remote tier after delete usecase" + ); +} diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index f576aba1b..daaf01b2f 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -142,6 +142,42 @@ impl Drop for DeadlockRequestGuard { } } +async fn enqueue_transitioned_delete_cleanup(bucket: &str, object: &str, opts: &ObjectOptions, existing: Option<&ObjectInfo>) { + let Some(existing) = existing else { + return; + }; + + let je = if opts.delete_prefix { + rustfs_ecstore::bucket::lifecycle::tier_sweeper::transitioned_force_delete_journal_entry(&existing.transitioned_object) + } else { + let version_id = opts.version_id.as_ref().and_then(|v| Uuid::parse_str(v).ok()); + rustfs_ecstore::bucket::lifecycle::tier_sweeper::transitioned_delete_journal_entry( + version_id, + opts.versioned, + opts.version_suspended, + &existing.transitioned_object, + ) + }; + let Some(je) = je else { + return; + }; + + let mut expiry_state = rustfs_ecstore::bucket::lifecycle::bucket_lifecycle_ops::GLOBAL_ExpiryState + .write() + .await; + if let Err(err) = expiry_state.enqueue_tier_journal_entry(&je).await { + warn!( + bucket, + object, + remote_object = %existing.transitioned_object.name, + remote_version_id = %existing.transitioned_object.version_id, + tier = %existing.transitioned_object.tier, + error = ?err, + "failed to enqueue transitioned object cleanup" + ); + } +} + pin_project! { struct ExtractArchiveEtagReader { #[pin] @@ -2794,6 +2830,7 @@ impl DefaultObjectUsecase { let mut object_to_delete = Vec::new(); let mut object_to_delete_idx = Vec::new(); let mut object_sizes = Vec::new(); + let mut existing_object_infos = Vec::new(); for (idx, obj_id) in delete.objects.iter().enumerate() { let raw_version_id = obj_id.version_id.clone(); let (version_id, version_uuid) = match normalize_delete_objects_version_id(raw_version_id.clone()) { @@ -2893,6 +2930,7 @@ impl DefaultObjectUsecase { object_to_delete_idx.push(idx); object_to_delete.push(object); + existing_object_infos.push(gerr.is_none().then_some(goi)); } let (mut dobjs, errs) = store @@ -2944,6 +2982,18 @@ impl DefaultObjectUsecase { dobjs[i].replication_state = Some(object_to_delete[i].replication_state()); } delete_results[didx].delete_object = Some(dobjs[i].clone()); + enqueue_transitioned_delete_cleanup( + &bucket, + &object_to_delete[i].object_name, + &ObjectOptions { + version_id: object_to_delete[i].version_id.map(|v| v.to_string()), + versioned: version_cfg.prefix_enabled(object_to_delete[i].object_name.as_str()), + version_suspended: version_cfg.suspended(), + ..Default::default() + }, + existing_object_infos[i].as_ref(), + ) + .await; let size = object_sizes[i]; if size > 0 { rustfs_ecstore::data_usage::decrement_bucket_usage_memory(&bucket, size as u64).await; @@ -3121,7 +3171,7 @@ impl DefaultObjectUsecase { .await .map_err(ApiError::from)?; - match store.get_object_info(&bucket, &key, &get_opts).await { + let existing_object_info = match store.get_object_info(&bucket, &key, &get_opts).await { Ok(obj_info) => { // Check for bypass governance retention header (permission already verified in access.rs) let bypass_governance = has_bypass_governance_header(&req.headers); @@ -3129,17 +3179,19 @@ impl DefaultObjectUsecase { if let Some(block_reason) = check_object_lock_for_deletion(&bucket, &obj_info, bypass_governance).await { return Err(S3Error::with_message(S3ErrorCode::AccessDenied, block_reason.error_message())); } + Some(obj_info) } Err(err) => { // If object not found, allow deletion to proceed (will return 204 No Content) if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) { return Err(ApiError::from(err).into()); } + None } - } + }; let obj_info = { - match store.delete_object(&bucket, &key, opts).await { + match store.delete_object(&bucket, &key, opts.clone()).await { Ok(obj) => obj, Err(err) => { if is_err_bucket_not_found(&err) { @@ -3157,6 +3209,8 @@ impl DefaultObjectUsecase { } }; + enqueue_transitioned_delete_cleanup(&bucket, &key, &opts, existing_object_info.as_ref()).await; + // Fast in-memory update for immediate quota consistency rustfs_ecstore::data_usage::decrement_bucket_usage_memory(&bucket, obj_info.size as u64).await;