diff --git a/.github/actions/setup/action.yml b/.github/actions/setup/action.yml index a08c61f0b..8a3ae8fc4 100644 --- a/.github/actions/setup/action.yml +++ b/.github/actions/setup/action.yml @@ -58,13 +58,13 @@ runs: - name: Install protoc uses: arduino/setup-protoc@v3 with: - version: "33.1" + version: "34.1" repo-token: ${{ inputs.github-token }} - name: Install flatc - uses: Nugine/setup-flatc@v1.2.4 + uses: Nugine/setup-flatc@v1 with: - version: "25.9.23" + version: "25.12.19" - name: Install Rust toolchain uses: dtolnay/rust-toolchain@stable diff --git a/crates/config/src/constants/runtime.rs b/crates/config/src/constants/runtime.rs index b164a03b3..7bfc96ac4 100644 --- a/crates/config/src/constants/runtime.rs +++ b/crates/config/src/constants/runtime.rs @@ -59,6 +59,25 @@ pub const DEFAULT_RUNTIME_DIAL9_ROTATION_COUNT: usize = 10; pub const DEFAULT_RUNTIME_DIAL9_SAMPLING_RATE: f64 = 1.0; // 100% sampling // Note: S3 bucket/prefix have no default; absence means upload is disabled (modeled as Option) +/// Maximum transition workers used as a local fallback when runtime env is unset. +pub const DEFAULT_TRANSITION_WORKERS_CAP: i64 = 16; +/// Absolute upper bound for transition workers accepted from runtime env. +pub const DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX: i64 = 32; +/// Default capacity for the transition queue. +pub const DEFAULT_TRANSITION_QUEUE_CAPACITY: usize = 1000; +/// Default send timeout for transition queue enqueue attempts, in milliseconds. +pub const DEFAULT_TRANSITION_QUEUE_SEND_TIMEOUT_MS: usize = 100; +/// Test-only fault injection env var that forces the immediate transition enqueue timeout path. +pub const ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT: &str = "RUSTFS_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT"; +/// Runtime env var controlling the transition worker count. +pub const ENV_TRANSITION_WORKERS: &str = "RUSTFS_MAX_TRANSITION_WORKERS"; +/// Runtime env var controlling the absolute maximum transition workers. +pub const ENV_TRANSITION_WORKERS_ABSOLUTE_MAX: &str = "RUSTFS_ABSOLUTE_MAX_WORKERS"; +/// Runtime env var controlling the transition queue capacity. +pub const ENV_TRANSITION_QUEUE_CAPACITY: &str = "RUSTFS_TRANSITION_QUEUE_CAPACITY"; +/// Runtime env var controlling the transition queue send timeout in milliseconds. +pub const ENV_TRANSITION_QUEUE_SEND_TIMEOUT_MS: &str = "RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS"; + // Allocator reclaim configuration pub const ENV_ALLOCATOR_RECLAIM_ENABLED: &str = "RUSTFS_ALLOCATOR_RECLAIM_ENABLED"; pub const ENV_ALLOCATOR_RECLAIM_INTERVAL_SECS: &str = "RUSTFS_ALLOCATOR_RECLAIM_INTERVAL_SECS"; diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 0d5666d8d..dd4e4c53e 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -47,6 +47,11 @@ 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_config::{ + DEFAULT_TRANSITION_QUEUE_CAPACITY, DEFAULT_TRANSITION_QUEUE_SEND_TIMEOUT_MS, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX, + DEFAULT_TRANSITION_WORKERS_CAP, ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT, ENV_TRANSITION_QUEUE_CAPACITY, + ENV_TRANSITION_QUEUE_SEND_TIMEOUT_MS, ENV_TRANSITION_WORKERS, ENV_TRANSITION_WORKERS_ABSOLUTE_MAX, +}; use rustfs_filemeta::{ FileInfo, FileInfoOpts, NULL_VERSION_ID, REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicationState, RestoreStatusOps, VersionPurgeStatusType, get_file_info, is_restored_object_on_disk, @@ -60,7 +65,7 @@ use s3s::dto::{ use s3s::header::{X_AMZ_RESTORE, X_AMZ_SERVER_SIDE_ENCRYPTION}; use sha2::{Digest, Sha256}; use std::any::Any; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::env; use std::pin::Pin; use std::sync::atomic::{AtomicI64, Ordering}; @@ -99,6 +104,48 @@ lazy_static! { pub static ref GLOBAL_TransitionState: Arc = TransitionState::new(); } +fn resolve_transition_worker_count() -> (i64, i64, i64) { + let fallback = std::cmp::min(num_cpus::get() as i64, DEFAULT_TRANSITION_WORKERS_CAP); + let configured = env::var(ENV_TRANSITION_WORKERS) + .ok() + .and_then(|value| value.parse::().ok()) + .filter(|value| *value > 0) + .unwrap_or(fallback); + let mut effective = configured; + let absolute_max = get_env_i64(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX); + effective = std::cmp::min(effective, absolute_max); + (configured, absolute_max, effective) +} + +fn resolve_transition_queue_capacity() -> usize { + get_env_usize(ENV_TRANSITION_QUEUE_CAPACITY, DEFAULT_TRANSITION_QUEUE_CAPACITY).max(1) +} + +fn resolve_transition_queue_send_timeout() -> StdDuration { + StdDuration::from_millis( + get_env_usize(ENV_TRANSITION_QUEUE_SEND_TIMEOUT_MS, DEFAULT_TRANSITION_QUEUE_SEND_TIMEOUT_MS).max(1) as u64, + ) +} + +fn is_immediate_transition_source(src: &LcEventSrc) -> bool { + matches!( + src, + LcEventSrc::S3PutObject | LcEventSrc::S3CopyObject | LcEventSrc::S3CompleteMultipartUpload + ) +} + +#[cfg(any(test, debug_assertions))] +fn should_force_immediate_transition_enqueue_timeout() -> bool { + env::var(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT) + .ok() + .is_some_and(|value| value == "1") +} + +#[cfg(not(any(test, debug_assertions)))] +fn should_force_immediate_transition_enqueue_timeout() -> bool { + false +} + pub struct LifecycleSys; impl LifecycleSys { @@ -537,15 +584,34 @@ pub struct TransitionState { pub num_workers: AtomicI64, kill_tx: A_Sender<()>, kill_rx: A_Receiver<()>, + transition_queue_capacity: usize, + transition_queue_send_timeout: StdDuration, active_tasks: AtomicI64, missed_immediate_tasks: AtomicI64, + queue_full_tasks: AtomicI64, + queue_send_timeout_tasks: AtomicI64, + compensation_scheduled_tasks: AtomicI64, + compensation_running_tasks: AtomicI64, + compensation_buckets: Arc>>, last_day_stats: Arc>>, } +enum ImmediateEnqueueFailure { + ForcedTimeout, + QueueClosed { timeout_ms: Option }, + QueueSendTimedOut { timeout_ms: u64 }, +} + impl TransitionState { #[allow(clippy::new_ret_no_self)] pub fn new() -> Arc { - let (tx1, rx1) = bounded(1000); + Self::new_with_capacity(resolve_transition_queue_capacity()) + } + + fn new_with_capacity(capacity: usize) -> Arc { + let capacity = capacity.max(1); + let queue_send_timeout = resolve_transition_queue_send_timeout(); + let (tx1, rx1) = bounded(capacity); let (tx2, rx2) = bounded(1); Arc::new(Self { transition_tx: tx1, @@ -553,39 +619,191 @@ impl TransitionState { num_workers: AtomicI64::new(0), kill_tx: tx2, kill_rx: rx2, + transition_queue_capacity: capacity, + transition_queue_send_timeout: queue_send_timeout, active_tasks: AtomicI64::new(0), missed_immediate_tasks: AtomicI64::new(0), + queue_full_tasks: AtomicI64::new(0), + queue_send_timeout_tasks: AtomicI64::new(0), + compensation_scheduled_tasks: AtomicI64::new(0), + compensation_running_tasks: AtomicI64::new(0), + compensation_buckets: Arc::new(Mutex::new(HashSet::new())), last_day_stats: Arc::new(Mutex::new(HashMap::new())), }) } - pub async fn queue_transition_task(&self, oi: &ObjectInfo, event: &lifecycle::Event, src: &LcEventSrc) { + fn schedule_bucket_compensation(self: &Arc, bucket: &str) -> bool { + let mut scheduled = self.compensation_buckets.lock().unwrap(); + if !scheduled.insert(bucket.to_string()) { + return false; + } + Self::inc_counter(&self.compensation_scheduled_tasks); + let bucket = bucket.to_string(); + let scheduled = Arc::clone(&self.compensation_buckets); + let state = Arc::clone(self); + tokio::spawn(async move { + Self::inc_counter(&state.compensation_running_tasks); + let Some(api) = crate::new_object_layer_fn() else { + scheduled.lock().unwrap().remove(&bucket); + Self::add_counter(&state.compensation_running_tasks, -1); + warn!(bucket = %bucket, "transition compensation skipped because object layer is unavailable"); + return; + }; + + if let Err(err) = enqueue_transition_for_existing_objects(api, &bucket).await { + warn!(bucket = %bucket, error = ?err, "transition compensation backfill failed"); + } else { + info!(bucket = %bucket, "transition compensation backfill completed"); + } + + scheduled.lock().unwrap().remove(&bucket); + Self::add_counter(&state.compensation_running_tasks, -1); + }); + true + } + + #[inline] + fn inc_counter(counter: &AtomicI64) { + Self::add_counter(counter, 1); + } + + #[inline] + fn add_counter(counter: &AtomicI64, delta: i64) { + counter.fetch_add(delta, Ordering::Relaxed); + } + + #[inline] + fn counter_value(counter: &AtomicI64) -> i64 { + counter.load(Ordering::Relaxed) + } + + fn handle_immediate_enqueue_failure(self: &Arc, oi: &ObjectInfo, src: &LcEventSrc, failure: ImmediateEnqueueFailure) { + Self::inc_counter(&self.missed_immediate_tasks); + let scheduled = self.schedule_bucket_compensation(&oi.bucket); + match failure { + ImmediateEnqueueFailure::ForcedTimeout => { + Self::inc_counter(&self.queue_send_timeout_tasks); + warn!( + bucket = %oi.bucket, + object = %oi.name, + source = ?src, + compensation_scheduled = scheduled, + "transition enqueue forced into timeout path for test fault injection" + ); + } + ImmediateEnqueueFailure::QueueClosed { timeout_ms } => match timeout_ms { + Some(timeout_ms) => { + warn!( + bucket = %oi.bucket, + object = %oi.name, + source = ?src, + timeout_ms, + compensation_scheduled = scheduled, + "transition enqueue failed because the queue is closed" + ); + } + None => { + warn!( + bucket = %oi.bucket, + object = %oi.name, + source = ?src, + compensation_scheduled = scheduled, + "transition enqueue failed because the queue is closed" + ); + } + }, + ImmediateEnqueueFailure::QueueSendTimedOut { timeout_ms } => { + Self::inc_counter(&self.queue_send_timeout_tasks); + warn!( + bucket = %oi.bucket, + object = %oi.name, + source = ?src, + timeout_ms, + compensation_scheduled = scheduled, + "transition enqueue timed out under backpressure" + ); + } + } + } + + pub async fn queue_transition_task(self: &Arc, oi: &ObjectInfo, event: &lifecycle::Event, src: &LcEventSrc) { + if is_immediate_transition_source(src) && should_force_immediate_transition_enqueue_timeout() { + self.handle_immediate_enqueue_failure(oi, src, ImmediateEnqueueFailure::ForcedTimeout); + return; + } + let task = TransitionTask { obj_info: oi.clone(), src: src.clone(), event: event.clone(), }; - select! { - //_ -> t.ctx.Done() => (), - _ = self.transition_tx.send(Some(task)) => (), - else => { - match src { - LcEventSrc::S3PutObject | LcEventSrc::S3CopyObject | LcEventSrc::S3CompleteMultipartUpload => { - self.missed_immediate_tasks.fetch_add(1, Ordering::SeqCst); + if is_immediate_transition_source(src) { + match self.transition_tx.try_send(Some(task)) { + Ok(()) => {} + Err(async_channel::TrySendError::Full(task)) => { + Self::inc_counter(&self.queue_full_tasks); + let send_timeout = self.transition_queue_send_timeout; + match tokio::time::timeout(send_timeout, self.transition_tx.send(task)).await { + Ok(Ok(())) => {} + Ok(Err(_)) => { + self.handle_immediate_enqueue_failure( + oi, + src, + ImmediateEnqueueFailure::QueueClosed { + timeout_ms: Some(send_timeout.as_millis() as u64), + }, + ); + } + Err(_) => { + self.handle_immediate_enqueue_failure( + oi, + src, + ImmediateEnqueueFailure::QueueSendTimedOut { + timeout_ms: send_timeout.as_millis() as u64, + }, + ); + } } - _ => () } - }, + Err(async_channel::TrySendError::Closed(_task)) => { + self.handle_immediate_enqueue_failure(oi, src, ImmediateEnqueueFailure::QueueClosed { timeout_ms: None }); + } + } + return; + } + + if let Err(err) = self.transition_tx.try_send(Some(task)) { + match err { + async_channel::TrySendError::Full(_) => { + debug!( + bucket = %oi.bucket, + object = %oi.name, + source = ?src, + "transition queue is full; deferring to scanner/backfill" + ); + } + async_channel::TrySendError::Closed(_) => { + warn!( + bucket = %oi.bucket, + object = %oi.name, + source = ?src, + "transition enqueue failed because the queue is closed" + ); + } + } } } pub async fn init(api: Arc) { - let max_workers = get_env_i64("RUSTFS_MAX_TRANSITION_WORKERS", std::cmp::min(num_cpus::get() as i64, 16)); - let mut n = max_workers; - let tw = 8; //globalILMConfig.getTransitionWorkers(); - if tw > 0 { - n = tw; - } + let (configured, absolute_max, n) = resolve_transition_worker_count(); + info!( + configured_transition_workers = configured, + absolute_max_workers = absolute_max, + effective_transition_workers = n, + transition_queue_capacity = GLOBAL_TransitionState.transition_queue_capacity, + transition_queue_send_timeout_ms = GLOBAL_TransitionState.transition_queue_send_timeout.as_millis() as u64, + "transition worker count resolved" + ); //let mut transition_state = GLOBAL_TransitionState.write().await; //self.objAPI = objAPI @@ -599,11 +817,27 @@ impl TransitionState { } pub fn active_tasks(&self) -> i64 { - self.active_tasks.load(Ordering::SeqCst) + Self::counter_value(&self.active_tasks) } pub fn missed_immediate_tasks(&self) -> i64 { - self.missed_immediate_tasks.load(Ordering::SeqCst) + Self::counter_value(&self.missed_immediate_tasks) + } + + pub fn queue_full_tasks(&self) -> i64 { + Self::counter_value(&self.queue_full_tasks) + } + + pub fn queue_send_timeout_tasks(&self) -> i64 { + Self::counter_value(&self.queue_send_timeout_tasks) + } + + pub fn compensation_scheduled_tasks(&self) -> i64 { + Self::counter_value(&self.compensation_scheduled_tasks) + } + + pub fn compensation_running_tasks(&self) -> i64 { + Self::counter_value(&self.compensation_running_tasks) } pub async fn worker(api: Arc) { @@ -626,7 +860,7 @@ impl TransitionState { if task.as_any().is::() { let task = task.as_any().downcast_ref::().expect("TransitionTask downcast failed"); - GLOBAL_TransitionState.active_tasks.fetch_add(1, Ordering::SeqCst); + TransitionState::inc_counter(&GLOBAL_TransitionState.active_tasks); let obj_info_for_event = ObjectInfo { bucket: task.obj_info.bucket.clone(), @@ -671,7 +905,7 @@ impl TransitionState { ..Default::default() }); } - GLOBAL_TransitionState.active_tasks.fetch_add(-1, Ordering::SeqCst); + TransitionState::add_counter(&GLOBAL_TransitionState.active_tasks, -1); } } else => () @@ -702,15 +936,17 @@ impl TransitionState { pub async fn update_workers_inner(api: Arc, n: i64) { let mut n = n; + let requested = n; if n == 0 { - let max_workers = get_env_i64("RUSTFS_MAX_TRANSITION_WORKERS", std::cmp::min(num_cpus::get() as i64, 16)); - n = max_workers; + let (_, _, effective) = resolve_transition_worker_count(); + n = effective; } // Allow environment override of maximum workers - let absolute_max = get_env_i64("RUSTFS_ABSOLUTE_MAX_WORKERS", 32); + let absolute_max = get_env_i64(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX); n = std::cmp::min(n, absolute_max); - let mut num_workers = GLOBAL_TransitionState.num_workers.load(Ordering::SeqCst); + let previous_num_workers = GLOBAL_TransitionState.num_workers.load(Ordering::SeqCst); + let mut num_workers = previous_num_workers; while num_workers < n { let clone_api = api.clone(); tokio::spawn(async move { @@ -727,6 +963,15 @@ impl TransitionState { num_workers -= 1; GLOBAL_TransitionState.num_workers.fetch_add(-1, Ordering::SeqCst); } + + info!( + requested_transition_workers = requested, + effective_transition_workers = n, + absolute_max_workers = absolute_max, + previous_transition_workers = previous_num_workers, + current_transition_workers = GLOBAL_TransitionState.num_workers.load(Ordering::SeqCst), + "transition workers updated" + ); } } @@ -2004,11 +2249,13 @@ pub async fn apply_lifecycle_action(event: &lifecycle::Event, src: &LcEventSrc, #[cfg(test)] mod tests { use super::{ - DATE_EXPIRY_EXISTING_OBJECTS_GRACE_SECS, StaleMultipartUploadCandidate, cleanup_empty_multipart_sha_dirs_on_local_disks, - cleanup_stale_multipart_uploads_once_at, lifecycle_deleted_object, lifecycle_rule_has_date_expiration, - lifecycle_version_purge_state_from_completed_targets, mark_delete_opts_skip_decommissioned_on_remote_success, - merge_stale_multipart_candidate, replication_state_for_delete, should_defer_date_expiry_for_recent_config_update, - should_reuse_lifecycle_delete_replication_state, + DATE_EXPIRY_EXISTING_OBJECTS_GRACE_SECS, DEFAULT_TRANSITION_QUEUE_CAPACITY, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX, + DEFAULT_TRANSITION_WORKERS_CAP, GLOBAL_TransitionState, StaleMultipartUploadCandidate, TransitionState, + cleanup_empty_multipart_sha_dirs_on_local_disks, cleanup_stale_multipart_uploads_once_at, lifecycle_deleted_object, + lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets, + mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate, replication_state_for_delete, + resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count, + should_defer_date_expiry_for_recent_config_update, should_reuse_lifecycle_delete_replication_state, }; use crate::bucket::metadata::BUCKET_LIFECYCLE_CONFIG; use crate::bucket::metadata_sys; @@ -2021,12 +2268,16 @@ mod tests { use crate::store_api::{ BucketOperations, BucketOptions, MakeBucketOptions, MultipartOperations, ObjectInfo, ObjectOptions, PutObjReader, }; + use futures::FutureExt; + use rustfs_config::ENV_TRANSITION_WORKERS_ABSOLUTE_MAX; use rustfs_filemeta::{ReplicateDecision, VersionPurgeStatusType}; use s3s::dto::{BucketLifecycleConfiguration, ExpirationStatus, LifecycleExpiration, LifecycleRule, Timestamp}; use serial_test::serial; use sha2::{Digest, Sha256}; use std::collections::HashMap; + use std::env; use std::path::PathBuf; + use std::sync::atomic::Ordering; use std::sync::{Arc, OnceLock}; use std::time::Duration as StdDuration; use time::OffsetDateTime; @@ -2043,6 +2294,204 @@ mod tests { assert!(opts.skip_decommissioned); } + #[allow(unsafe_code)] + fn with_transition_worker_env(transition: Option<&str>, absolute: Option<&str>, test_fn: F) + where + F: FnOnce(), + { + let original_transition = env::var_os("RUSTFS_MAX_TRANSITION_WORKERS"); + let original_absolute = env::var_os(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX); + + match transition { + Some(value) => unsafe { + env::set_var("RUSTFS_MAX_TRANSITION_WORKERS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_MAX_TRANSITION_WORKERS"); + }, + } + match absolute { + Some(value) => unsafe { + env::set_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX, value); + }, + None => unsafe { + env::remove_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX); + }, + } + + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(test_fn)); + + match original_transition { + Some(value) => unsafe { + env::set_var("RUSTFS_MAX_TRANSITION_WORKERS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_MAX_TRANSITION_WORKERS"); + }, + } + match original_absolute { + Some(value) => unsafe { + env::set_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX, value); + }, + None => unsafe { + env::remove_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX); + }, + } + + if let Err(e) = result { + std::panic::resume_unwind(e); + } + } + + #[allow(unsafe_code)] + async fn with_transition_worker_env_async(transition: Option<&str>, absolute: Option<&str>, test_fn: F) + where + F: FnOnce() -> Fut, + Fut: std::future::Future, + { + let original_transition = env::var_os("RUSTFS_MAX_TRANSITION_WORKERS"); + let original_absolute = env::var_os(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX); + + match transition { + Some(value) => unsafe { + env::set_var("RUSTFS_MAX_TRANSITION_WORKERS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_MAX_TRANSITION_WORKERS"); + }, + } + match absolute { + Some(value) => unsafe { + env::set_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX, value); + }, + None => unsafe { + env::remove_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX); + }, + } + + let result = std::panic::AssertUnwindSafe(test_fn()).catch_unwind().await; + + match original_transition { + Some(value) => unsafe { + env::set_var("RUSTFS_MAX_TRANSITION_WORKERS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_MAX_TRANSITION_WORKERS"); + }, + } + match original_absolute { + Some(value) => unsafe { + env::set_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX, value); + }, + None => unsafe { + env::remove_var(ENV_TRANSITION_WORKERS_ABSOLUTE_MAX); + }, + } + + if let Err(e) = result { + std::panic::resume_unwind(e); + } + } + + #[allow(unsafe_code)] + fn with_transition_queue_env(capacity: Option<&str>, timeout_ms: Option<&str>, test_fn: F) + where + F: FnOnce(), + { + let original_capacity = env::var_os("RUSTFS_TRANSITION_QUEUE_CAPACITY"); + let original_timeout = env::var_os("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS"); + + match capacity { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_CAPACITY", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_CAPACITY"); + }, + } + match timeout_ms { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS"); + }, + } + + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(test_fn)); + + match original_capacity { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_CAPACITY", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_CAPACITY"); + }, + } + match original_timeout { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS"); + }, + } + + if let Err(e) = result { + std::panic::resume_unwind(e); + } + } + + #[allow(unsafe_code)] + async fn with_transition_queue_env_async(capacity: Option<&str>, timeout_ms: Option<&str>, test_fn: F) + where + F: FnOnce() -> Fut, + Fut: std::future::Future, + { + let original_capacity = env::var_os("RUSTFS_TRANSITION_QUEUE_CAPACITY"); + let original_timeout = env::var_os("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS"); + + match capacity { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_CAPACITY", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_CAPACITY"); + }, + } + match timeout_ms { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS"); + }, + } + + let result = std::panic::AssertUnwindSafe(test_fn()).catch_unwind().await; + + match original_capacity { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_CAPACITY", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_CAPACITY"); + }, + } + match original_timeout { + Some(value) => unsafe { + env::set_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS", value); + }, + None => unsafe { + env::remove_var("RUSTFS_TRANSITION_QUEUE_SEND_TIMEOUT_MS"); + }, + } + + if let Err(e) = result { + std::panic::resume_unwind(e); + } + } + #[test] fn lifecycle_rule_has_date_expiration_detects_enabled_date_rule() { let lc = BucketLifecycleConfiguration { @@ -2068,6 +2517,114 @@ mod tests { assert!(!lifecycle_rule_has_date_expiration(&lc, "missing-rule")); } + #[test] + #[serial] + fn resolve_transition_worker_count_uses_fallback_when_env_missing() { + with_transition_worker_env(None, None, || { + let (configured, absolute_max, effective) = resolve_transition_worker_count(); + + let fallback = std::cmp::min(num_cpus::get() as i64, DEFAULT_TRANSITION_WORKERS_CAP); + assert_eq!(configured, fallback); + assert_eq!(absolute_max, DEFAULT_TRANSITION_WORKERS_ABSOLUTE_MAX); + assert_eq!(effective, fallback); + }); + } + + #[test] + #[serial] + fn resolve_transition_worker_count_honors_positive_env_value() { + with_transition_worker_env(Some("4"), Some("32"), || { + let (configured, absolute_max, effective) = resolve_transition_worker_count(); + + assert_eq!(configured, 4); + assert_eq!(absolute_max, 32); + assert_eq!(effective, 4); + }); + } + + #[test] + #[serial] + fn resolve_transition_worker_count_clamps_to_absolute_max() { + with_transition_worker_env(Some("64"), Some("16"), || { + let (configured, absolute_max, effective) = resolve_transition_worker_count(); + + assert_eq!(configured, 64); + assert_eq!(absolute_max, 16); + assert_eq!(effective, 16); + }); + } + + #[test] + #[serial] + fn resolve_transition_worker_count_falls_back_for_zero_value() { + with_transition_worker_env(Some("0"), Some("32"), || { + let (configured, absolute_max, effective) = resolve_transition_worker_count(); + + let fallback = std::cmp::min(num_cpus::get() as i64, DEFAULT_TRANSITION_WORKERS_CAP); + assert_eq!(configured, fallback); + assert_eq!(absolute_max, 32); + assert_eq!(effective, fallback); + }); + } + + #[test] + #[serial] + fn resolve_transition_queue_capacity_uses_default_when_env_missing() { + with_transition_queue_env(None, None, || { + assert_eq!(resolve_transition_queue_capacity(), DEFAULT_TRANSITION_QUEUE_CAPACITY); + }); + } + + #[test] + #[serial] + fn resolve_transition_queue_capacity_honors_positive_env_value() { + with_transition_queue_env(Some("128"), None, || { + assert_eq!(resolve_transition_queue_capacity(), 128); + }); + } + + #[test] + #[serial] + fn resolve_transition_queue_send_timeout_honors_positive_env_value() { + with_transition_queue_env(None, Some("250"), || { + assert_eq!(resolve_transition_queue_send_timeout(), StdDuration::from_millis(250)); + }); + } + + #[tokio::test(flavor = "current_thread")] + async fn schedule_bucket_compensation_deduplicates_same_bucket() { + let state = TransitionState::new_with_capacity(1); + + let first = state.schedule_bucket_compensation("bucket-a"); + let second = state.schedule_bucket_compensation("bucket-a"); + + assert!(first); + assert!(!second); + assert_eq!(state.compensation_scheduled_tasks(), 1); + } + + #[tokio::test] + #[serial] + async fn transition_state_init_honors_runtime_configured_worker_count() { + let (_paths, ecstore) = setup_test_env().await; + let original_workers = GLOBAL_TransitionState.num_workers.load(Ordering::SeqCst); + with_transition_worker_env_async(Some("3"), Some("8"), || async { + TransitionState::update_workers(ecstore.clone(), 0).await; + assert_eq!(GLOBAL_TransitionState.num_workers.load(Ordering::SeqCst), 3); + }) + .await; + + let current_workers = GLOBAL_TransitionState.num_workers.load(Ordering::SeqCst); + if original_workers > 0 { + TransitionState::update_workers(ecstore, original_workers).await; + } else { + for _ in 0..current_workers { + let _ = GLOBAL_TransitionState.kill_tx.send(()).await; + GLOBAL_TransitionState.num_workers.fetch_add(-1, Ordering::SeqCst); + } + } + } + #[test] fn should_defer_date_expiry_for_recent_config_update_respects_grace_window() { let now = OffsetDateTime::now_utc(); diff --git a/crates/obs/src/metrics/collectors/ilm.rs b/crates/obs/src/metrics/collectors/ilm.rs index d19b7f404..73d7dd0ba 100644 --- a/crates/obs/src/metrics/collectors/ilm.rs +++ b/crates/obs/src/metrics/collectors/ilm.rs @@ -36,6 +36,14 @@ pub struct IlmStats { pub transition_pending_tasks: u64, /// Number of missed immediate ILM transition tasks pub transition_missed_immediate_tasks: u64, + /// Number of ILM transition tasks that initially hit full queue backpressure + pub transition_queue_full_tasks: u64, + /// Number of ILM transition tasks that timed out waiting for queue capacity + pub transition_queue_send_timeout_tasks: u64, + /// Number of bucket-level compensation tasks scheduled after immediate enqueue failure + pub transition_compensation_scheduled_tasks: u64, + /// Number of bucket-level compensation tasks currently running + pub transition_compensation_running_tasks: u64, /// Total number of object versions scanned for ILM pub versions_scanned: u64, } @@ -53,6 +61,19 @@ pub fn collect_ilm_metrics(stats: &IlmStats) -> Vec { &ILM_TRANSITION_MISSED_IMMEDIATE_TASKS_MD, stats.transition_missed_immediate_tasks as f64, ), + PrometheusMetric::from_descriptor(&ILM_TRANSITION_QUEUE_FULL_TASKS_MD, stats.transition_queue_full_tasks as f64), + PrometheusMetric::from_descriptor( + &ILM_TRANSITION_QUEUE_SEND_TIMEOUT_TASKS_MD, + stats.transition_queue_send_timeout_tasks as f64, + ), + PrometheusMetric::from_descriptor( + &ILM_TRANSITION_COMPENSATION_SCHEDULED_TASKS_MD, + stats.transition_compensation_scheduled_tasks as f64, + ), + PrometheusMetric::from_descriptor( + &ILM_TRANSITION_COMPENSATION_RUNNING_TASKS_MD, + stats.transition_compensation_running_tasks as f64, + ), PrometheusMetric::from_descriptor(&ILM_VERSIONS_SCANNED_MD, stats.versions_scanned as f64), ] } @@ -68,12 +89,16 @@ mod tests { transition_active_tasks: 5, transition_pending_tasks: 50, transition_missed_immediate_tasks: 10, + transition_queue_full_tasks: 2, + transition_queue_send_timeout_tasks: 3, + transition_compensation_scheduled_tasks: 4, + transition_compensation_running_tasks: 1, versions_scanned: 1000000, }; let metrics = collect_ilm_metrics(&stats); - assert_eq!(metrics.len(), 5); + assert_eq!(metrics.len(), 9); let pending = metrics.iter().find(|m| m.value == 100.0); assert!(pending.is_some()); @@ -87,7 +112,7 @@ mod tests { let stats = IlmStats::default(); let metrics = collect_ilm_metrics(&stats); - assert_eq!(metrics.len(), 5); + assert_eq!(metrics.len(), 9); for metric in &metrics { assert_eq!(metric.value, 0.0); assert!(metric.labels.is_empty()); diff --git a/crates/obs/src/metrics/schema/entry/metric_name.rs b/crates/obs/src/metrics/schema/entry/metric_name.rs index a7cc2a812..6ab12af1f 100644 --- a/crates/obs/src/metrics/schema/entry/metric_name.rs +++ b/crates/obs/src/metrics/schema/entry/metric_name.rs @@ -248,6 +248,10 @@ pub enum MetricName { IlmTransitionActiveTasks, IlmTransitionPendingTasks, IlmTransitionMissedImmediateTasks, + IlmTransitionQueueFullTasks, + IlmTransitionQueueSendTimeoutTasks, + IlmTransitionCompensationScheduledTasks, + IlmTransitionCompensationRunningTasks, IlmVersionsScanned, // Copy the relevant metrics @@ -584,6 +588,10 @@ impl MetricName { Self::IlmTransitionActiveTasks => "transition_active_tasks".to_string(), Self::IlmTransitionPendingTasks => "transition_pending_tasks".to_string(), Self::IlmTransitionMissedImmediateTasks => "transition_missed_immediate_tasks".to_string(), + Self::IlmTransitionQueueFullTasks => "transition_queue_full_tasks".to_string(), + Self::IlmTransitionQueueSendTimeoutTasks => "transition_queue_send_timeout_tasks".to_string(), + Self::IlmTransitionCompensationScheduledTasks => "transition_compensation_scheduled_tasks".to_string(), + Self::IlmTransitionCompensationRunningTasks => "transition_compensation_running_tasks".to_string(), Self::IlmVersionsScanned => "versions_scanned".to_string(), // Copy the relevant metrics diff --git a/crates/obs/src/metrics/schema/ilm.rs b/crates/obs/src/metrics/schema/ilm.rs index bfa4914c1..253464552 100644 --- a/crates/obs/src/metrics/schema/ilm.rs +++ b/crates/obs/src/metrics/schema/ilm.rs @@ -53,6 +53,42 @@ pub static ILM_TRANSITION_MISSED_IMMEDIATE_TASKS_MD: LazyLock ) }); +pub static ILM_TRANSITION_QUEUE_FULL_TASKS_MD: LazyLock = LazyLock::new(|| { + new_counter_md( + MetricName::IlmTransitionQueueFullTasks, + "Number of ILM transition tasks that initially hit full transition queue backpressure", + &[], + subsystems::ILM, + ) +}); + +pub static ILM_TRANSITION_QUEUE_SEND_TIMEOUT_TASKS_MD: LazyLock = LazyLock::new(|| { + new_counter_md( + MetricName::IlmTransitionQueueSendTimeoutTasks, + "Number of ILM transition tasks that timed out waiting for queue capacity", + &[], + subsystems::ILM, + ) +}); + +pub static ILM_TRANSITION_COMPENSATION_SCHEDULED_TASKS_MD: LazyLock = LazyLock::new(|| { + new_counter_md( + MetricName::IlmTransitionCompensationScheduledTasks, + "Number of bucket-level ILM transition compensation tasks scheduled after enqueue failure", + &[], + subsystems::ILM, + ) +}); + +pub static ILM_TRANSITION_COMPENSATION_RUNNING_TASKS_MD: LazyLock = LazyLock::new(|| { + new_gauge_md( + MetricName::IlmTransitionCompensationRunningTasks, + "Number of bucket-level ILM transition compensation tasks currently running", + &[], + subsystems::ILM, + ) +}); + pub static ILM_VERSIONS_SCANNED_MD: LazyLock = LazyLock::new(|| { new_counter_md( MetricName::IlmVersionsScanned, diff --git a/crates/obs/src/metrics/stats_collector.rs b/crates/obs/src/metrics/stats_collector.rs index 7cfa4991f..d75333b8a 100644 --- a/crates/obs/src/metrics/stats_collector.rs +++ b/crates/obs/src/metrics/stats_collector.rs @@ -840,6 +840,10 @@ pub async fn collect_ilm_metric_stats() -> Option { let transition_active_tasks = GLOBAL_TransitionState.active_tasks().max(0) as u64; let transition_pending_tasks = GLOBAL_TransitionState.pending_tasks() as u64; let transition_missed_immediate_tasks = GLOBAL_TransitionState.missed_immediate_tasks().max(0) as u64; + let transition_queue_full_tasks = GLOBAL_TransitionState.queue_full_tasks().max(0) as u64; + let transition_queue_send_timeout_tasks = GLOBAL_TransitionState.queue_send_timeout_tasks().max(0) as u64; + let transition_compensation_scheduled_tasks = GLOBAL_TransitionState.compensation_scheduled_tasks().max(0) as u64; + let transition_compensation_running_tasks = GLOBAL_TransitionState.compensation_running_tasks().max(0) as u64; let metrics = global_metrics().report().await; let versions_scanned = metrics.life_time_ilm.values().copied().sum(); @@ -848,6 +852,10 @@ pub async fn collect_ilm_metric_stats() -> Option { transition_active_tasks, transition_pending_tasks, transition_missed_immediate_tasks, + transition_queue_full_tasks, + transition_queue_send_timeout_tasks, + transition_compensation_scheduled_tasks, + transition_compensation_running_tasks, versions_scanned, }) } diff --git a/crates/scanner/tests/lifecycle_integration_test.rs b/crates/scanner/tests/lifecycle_integration_test.rs index fc6c5e9e1..40b0d6fd4 100644 --- a/crates/scanner/tests/lifecycle_integration_test.rs +++ b/crates/scanner/tests/lifecycle_integration_test.rs @@ -12,10 +12,15 @@ // See the License for the specific language governing permissions and // limitations under the License. +use futures::FutureExt; +use rustfs_config::ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT; use rustfs_ecstore::{ bucket::lifecycle::lifecycle::TransitionOptions, bucket::metadata::BUCKET_LIFECYCLE_CONFIG, - bucket::{lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects, metadata_sys}, + bucket::{ + lifecycle::bucket_lifecycle_ops::enqueue_transition_for_existing_objects, metadata_sys, + versioning_sys::BucketVersioningSys, + }, client::transition_api::{ReadCloser, ReaderImpl}, disk::endpoint::Endpoint, disk::{DiskAPI, DiskOption, STORAGE_FORMAT_FILE, new_disk}, @@ -41,6 +46,7 @@ use s3s::dto::RestoreRequest; use serial_test::serial; use std::{ collections::HashMap, + env, io::Cursor, path::{Path, PathBuf}, sync::{Arc, Once, OnceLock}, @@ -245,6 +251,14 @@ async fn upload_test_object(ecstore: &Arc, bucket: &str, object: &str, info!("Uploaded test object: {}/{} ({} bytes)", bucket, object, object_info.size); } +async fn modeled_versioned_delete_opts(bucket: &str, object: &str) -> ObjectOptions { + ObjectOptions { + versioned: BucketVersioningSys::prefix_enabled(bucket, object).await, + version_suspended: BucketVersioningSys::prefix_suspended(bucket, object).await, + ..Default::default() + } +} + /// Test helper: Set bucket lifecycle configuration #[allow(dead_code)] async fn set_bucket_lifecycle(bucket_name: &str) -> Result<(), Box> { @@ -271,7 +285,8 @@ async fn set_bucket_lifecycle(bucket_name: &str) -> Result<(), Box Result<(), Box> { - // Create a simple lifecycle configuration XML with 0 days expiry for immediate testing + // Create lifecycle rule that targets delete-marker cleanup only. + // Keep Expiration.Days unset to avoid expiring live transitioned object versions. let lifecycle_xml = r#" @@ -281,7 +296,6 @@ async fn set_bucket_lifecycle_deletemarker(bucket_name: &str) -> Result<(), Box< test/ - 0 true @@ -292,6 +306,29 @@ async fn set_bucket_lifecycle_deletemarker(bucket_name: &str) -> Result<(), Box< Ok(()) } +#[allow(dead_code)] +async fn set_bucket_lifecycle_delmarker_expiration(bucket_name: &str, days: i64) -> Result<(), Box> { + let lifecycle_xml = format!( + r#" + + + test-rule + Enabled + + test/ + + + {days} + + +"# + ); + + metadata_sys::update(bucket_name, BUCKET_LIFECYCLE_CONFIG, lifecycle_xml.into_bytes()).await?; + + Ok(()) +} + #[allow(dead_code)] async fn set_bucket_lifecycle_transition(bucket_name: &str) -> Result<(), Box> { set_bucket_lifecycle_transition_with_tier(bucket_name, "COLDTIER44").await @@ -529,6 +566,22 @@ async fn wait_for_version_count(ecstore: &Arc, bucket: &str, object: &s } } +async fn wait_for_remote_object_count(backend: &MockWarmBackend, expected: usize, timeout: Duration) -> bool { + let deadline = tokio::time::Instant::now() + timeout; + + loop { + if backend.objects.lock().await.len() == expected { + return true; + } + + if tokio::time::Instant::now() >= deadline { + return false; + } + + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + async fn scan_object_with_lifecycle(disk_path: &Path, bucket: &str, object: &str) { let mut endpoint = Endpoint::try_from(disk_path.to_str().unwrap()).unwrap(); endpoint.set_pool_index(0); @@ -754,6 +807,30 @@ async fn wait_for_transition( } } +#[allow(unsafe_code)] +async fn with_forced_immediate_enqueue_timeout(test_fn: F) +where + F: FnOnce() -> Fut, + Fut: std::future::Future, +{ + let original = env::var_os(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT); + unsafe { + env::set_var(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT, "1"); + } + let result = std::panic::AssertUnwindSafe(test_fn()).catch_unwind().await; + match original { + Some(value) => unsafe { + env::set_var(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT, value); + }, + None => unsafe { + env::remove_var(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT); + }, + } + if let Err(err) = result { + std::panic::resume_unwind(err); + } +} + mod serial_tests { use super::*; @@ -1217,6 +1294,398 @@ mod serial_tests { ); } + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + #[serial] + #[ignore = "requires isolated global object layer state"] + async fn test_scanner_cleanup_still_works_after_immediate_compensation_transition() { + 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-after-compensation-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_name = "test/object.txt"; + let payload = b"scanner cleanup should still work after immediate compensation"; + + 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"); + + with_forced_immediate_enqueue_timeout(|| async { + upload_test_object(&ecstore, bucket_name.as_str(), object_name, payload).await; + }) + .await; + + let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should transition after compensation backfill"); + 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 after compensation-driven transition"); + + assert!( + free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0, + "deleting a compensation-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 cleanup 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 clean stale remote object even after immediate compensation transitioned it" + ); + assert_eq!( + free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await, + 0, + "free-version metadata should be removed after scanner cleanup" + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + #[serial] + #[ignore = "requires isolated global object layer state"] + async fn test_existing_object_backfill_is_idempotent_after_immediate_compensation_transition() { + 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-backfill-after-compensation-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_name = "test/object.txt"; + let payload = b"existing-object backfill should be idempotent after compensation transition"; + + 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"); + + with_forced_immediate_enqueue_timeout(|| async { + upload_test_object(&ecstore, bucket_name.as_str(), object_name, payload).await; + }) + .await; + + let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should transition after immediate compensation backfill"); + let remote_object = transitioned.transitioned_object.name.clone(); + assert!(backend.objects.lock().await.contains_key(&remote_object)); + + enqueue_transition_for_existing_objects(ecstore.clone(), bucket_name.as_str()) + .await + .expect("existing-object backfill should succeed after compensation transition"); + + let info = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should remain transitioned after existing-object backfill rerun"); + + assert_eq!(info.transitioned_object.status, "complete"); + assert_eq!(info.transitioned_object.tier, tier_name); + assert_eq!(info.transitioned_object.name, remote_object); + assert!(backend.objects.lock().await.contains_key(&remote_object)); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + #[serial] + #[ignore = "requires isolated global object layer state"] + async fn test_noncurrent_expiry_still_works_after_immediate_compensation_transition() { + let (disk_paths, ecstore) = setup_isolated_test_env(true).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-versioned-compensation-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_name = "test/object.txt"; + + create_test_lock_bucket(&ecstore, bucket_name.as_str()).await; + + let lifecycle_xml = format!( + r#" + + + test-rule + Enabled + + test/ + + + 0 + {tier_name} + + + 0 + + +"# + ); + metadata_sys::update(bucket_name.as_str(), BUCKET_LIFECYCLE_CONFIG, lifecycle_xml.into_bytes()) + .await + .expect("Failed to set lifecycle configuration"); + + let mut reader = PutObjReader::from_vec(b"v1".to_vec()); + ecstore + .put_object( + bucket_name.as_str(), + object_name, + &mut reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed to upload v1"); + + with_forced_immediate_enqueue_timeout(|| async { + let mut reader = PutObjReader::from_vec(b"v2".to_vec()); + ecstore + .put_object( + bucket_name.as_str(), + object_name, + &mut reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed to upload v2"); + }) + .await; + + let info = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("current version should transition after compensation backfill"); + + assert_eq!(info.transitioned_object.status, "complete"); + assert_eq!(info.transitioned_object.tier, tier_name); + assert!(backend.objects.lock().await.contains_key(&info.transitioned_object.name)); + + scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await; + + assert!( + wait_for_version_count(&ecstore, bucket_name.as_str(), object_name, 1, Duration::from_secs(3)).await, + "noncurrent expiry should still remove the previous version after compensation transition" + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + #[serial] + #[ignore = "requires isolated global object layer state"] + async fn test_noncurrent_transition_still_works_after_immediate_compensation_transition() { + let (disk_paths, ecstore) = setup_isolated_test_env(true).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-noncurrent-transition-comp-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_name = "test/object.txt"; + + create_test_lock_bucket(&ecstore, bucket_name.as_str()).await; + + let lifecycle_xml = format!( + r#" + + + test-rule + Enabled + + test/ + + + 0 + {tier_name} + + + 0 + {tier_name} + + +"# + ); + metadata_sys::update(bucket_name.as_str(), BUCKET_LIFECYCLE_CONFIG, lifecycle_xml.into_bytes()) + .await + .expect("Failed to set lifecycle configuration"); + + let mut reader = PutObjReader::from_vec(b"v1".to_vec()); + ecstore + .put_object( + bucket_name.as_str(), + object_name, + &mut reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed to upload v1"); + + with_forced_immediate_enqueue_timeout(|| async { + let mut reader = PutObjReader::from_vec(b"v2".to_vec()); + ecstore + .put_object( + bucket_name.as_str(), + object_name, + &mut reader, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect("failed to upload v2"); + }) + .await; + + let info = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("current version should transition after compensation backfill"); + assert_eq!(info.transitioned_object.status, "complete"); + assert_eq!(info.transitioned_object.tier, tier_name); + + scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await; + + assert!( + wait_for_remote_object_count(&backend, 2, TRANSITION_WAIT_TIMEOUT).await, + "noncurrent transition should still move the previous version into the remote tier" + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + #[serial] + #[ignore = "requires isolated global object layer state"] + async fn test_modeled_versioned_delete_creates_delete_marker_after_immediate_compensation_transition() { + let (_disk_paths, ecstore) = setup_isolated_test_env(true).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-modeled-versioned-delete-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_name = "test/object.txt"; + let payload = b"modeled versioned delete should create delete marker after compensation"; + + create_test_lock_bucket(&ecstore, bucket_name.as_str()).await; + set_bucket_lifecycle_transition_with_tier(bucket_name.as_str(), &tier_name) + .await + .expect("Failed to set transition lifecycle configuration"); + + with_forced_immediate_enqueue_timeout(|| async { + upload_test_object(&ecstore, bucket_name.as_str(), object_name, payload).await; + }) + .await; + + let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("current version should transition after compensation backfill"); + let remote_object = transitioned.transitioned_object.name.clone(); + assert!(backend.objects.lock().await.contains_key(&remote_object)); + + ecstore + .delete_object( + bucket_name.as_str(), + object_name, + modeled_versioned_delete_opts(bucket_name.as_str(), object_name).await, + ) + .await + .expect("modeled versioned delete should succeed"); + + assert!( + object_is_delete_marker(&ecstore, bucket_name.as_str(), object_name).await, + "versioned delete modeled with versioned flags should create a delete marker" + ); + assert!( + backend.objects.lock().await.contains_key(&remote_object), + "creating a delete marker should not remove the transitioned remote object version" + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] + #[serial] + #[ignore = "requires isolated global object layer state"] + async fn test_modeled_delete_marker_cleanup_after_immediate_compensation_transition() { + let (disk_paths, ecstore) = setup_isolated_test_env(true).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-modeled-del-marker-cleanup-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_name = "test/object.txt"; + let payload = b"modeled delete-marker cleanup should converge after compensation transition"; + + create_test_lock_bucket(&ecstore, bucket_name.as_str()).await; + set_bucket_lifecycle_transition_with_tier(bucket_name.as_str(), &tier_name) + .await + .expect("Failed to set transition lifecycle configuration"); + + with_forced_immediate_enqueue_timeout(|| async { + upload_test_object(&ecstore, bucket_name.as_str(), object_name, payload).await; + }) + .await; + + let transitioned = wait_for_transition(&ecstore, bucket_name.as_str(), object_name, TRANSITION_WAIT_TIMEOUT) + .await + .expect("current version should transition after compensation backfill"); + let remote_object = transitioned.transitioned_object.name.clone(); + assert!(backend.objects.lock().await.contains_key(&remote_object)); + + ecstore + .delete_object( + bucket_name.as_str(), + object_name, + modeled_versioned_delete_opts(bucket_name.as_str(), object_name).await, + ) + .await + .expect("modeled versioned delete should succeed"); + + assert!( + object_is_delete_marker(&ecstore, bucket_name.as_str(), object_name).await, + "modeled versioned delete should create delete marker before cleanup" + ); + assert!( + backend.objects.lock().await.contains_key(&remote_object), + "delete marker creation should not remove transitioned remote object" + ); + + set_bucket_lifecycle_delmarker_expiration(bucket_name.as_str(), 1) + .await + .expect("Failed to set delete marker expiration lifecycle configuration"); + + scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await; + + assert!( + object_is_delete_marker(&ecstore, bucket_name.as_str(), object_name).await, + "delete marker should remain before DelMarkerExpiration due time" + ); + assert!( + backend.objects.lock().await.contains_key(&remote_object), + "pre-due delete marker lifecycle scan should not remove transitioned remote object" + ); + + set_bucket_lifecycle_deletemarker(bucket_name.as_str()) + .await + .expect("Failed to set expired object delete marker lifecycle configuration"); + scan_object_with_lifecycle(&disk_paths[0], bucket_name.as_str(), object_name).await; + + assert!( + wait_for_object_absence(&ecstore, bucket_name.as_str(), object_name, Duration::from_secs(5)).await, + "expired object delete marker lifecycle should eventually clean up the delete marker" + ); + assert!( + backend.objects.lock().await.contains_key(&remote_object), + "delete marker lifecycle cleanup should not remove transitioned remote object" + ); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] #[serial] #[ignore = "requires isolated global object layer state"] diff --git a/rustfs/src/app/lifecycle_transition_api_test.rs b/rustfs/src/app/lifecycle_transition_api_test.rs index 3e759121c..02ee25543 100644 --- a/rustfs/src/app/lifecycle_transition_api_test.rs +++ b/rustfs/src/app/lifecycle_transition_api_test.rs @@ -16,8 +16,10 @@ use super::{multipart_usecase::DefaultMultipartUsecase, object_usecase::DefaultO use crate::app::bucket_usecase::DefaultBucketUsecase; use crate::storage::ecfs::FS; use bytes::Bytes; +use futures::FutureExt; use futures::stream; use http::{Extensions, HeaderMap, Method, Uri}; +use rustfs_config::ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT; use rustfs_ecstore::{ bucket::metadata::BUCKET_LIFECYCLE_CONFIG, bucket::metadata_sys, @@ -43,7 +45,7 @@ use serial_test::serial; use std::{ collections::HashMap, convert::Infallible, - fs as stdfs, + env, fs as stdfs, io::Cursor, path::PathBuf, sync::{Arc, Once, OnceLock}, @@ -325,6 +327,30 @@ async fn wait_for_transition( } } +#[allow(unsafe_code)] +async fn with_forced_immediate_enqueue_timeout(test_fn: F) +where + F: FnOnce() -> Fut, + Fut: std::future::Future, +{ + let original = env::var_os(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT); + unsafe { + env::set_var(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT, "1"); + } + let result = std::panic::AssertUnwindSafe(test_fn()).catch_unwind().await; + match original { + Some(value) => unsafe { + env::set_var(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT, value); + }, + None => unsafe { + env::remove_var(ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT); + }, + } + if let Err(err) = result { + std::panic::resume_unwind(err); + } +} + async fn wait_for_remote_absence(backend: &MockWarmBackend, object: &str, timeout: Duration) -> bool { let deadline = tokio::time::Instant::now() + timeout; @@ -637,6 +663,334 @@ async fn lifecycle_transition_marks_dirty_disks_for_capacity_manager() { assert_eq!(actual_paths, expected_paths); } +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +#[serial] +#[ignore = "requires isolated global object layer state"] +async fn immediate_transition_timeout_eventually_completes_via_compensation() { + let (_disk_paths, ecstore) = setup_test_env().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 = format!("test-compensation-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/object.txt"; + let payload = b"transition compensation should eventually complete"; + + 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"); + + with_forced_immediate_enqueue_timeout(|| async { + let _ = upload_test_object(&ecstore, bucket.as_str(), object, payload).await; + }) + .await; + + let info = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should eventually transition after compensation backfill"); + + assert_eq!(info.transitioned_object.status, "complete"); + 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 compensation_driven_copy_still_completes_transition() { + 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 src_bucket = format!("test-comp-copy-src-{}", &Uuid::new_v4().simple().to_string()[..8]); + let dst_bucket = format!("test-comp-copy-dst-{}", &Uuid::new_v4().simple().to_string()[..8]); + let src_object = "test/source.txt"; + let dst_object = "test/copied.txt"; + let payload = b"copy object should still transition after compensation"; + + create_test_bucket(&ecstore, src_bucket.as_str()).await; + create_test_bucket(&ecstore, dst_bucket.as_str()).await; + set_bucket_lifecycle_transition_with_tier(dst_bucket.as_str(), &tier_name) + .await + .expect("Failed to set destination lifecycle configuration"); + let _ = upload_test_object(&ecstore, src_bucket.as_str(), src_object, payload).await; + + let copy_input = CopyObjectInput::builder() + .copy_source(CopySource::Bucket { + bucket: src_bucket.clone().into(), + key: src_object.to_string().into(), + version_id: None, + }) + .bucket(dst_bucket.clone()) + .key(dst_object.to_string()) + .build() + .unwrap(); + + with_forced_immediate_enqueue_timeout(|| async { + Box::pin(usecase.execute_copy_object(build_request(copy_input, Method::PUT))) + .await + .expect("Failed to copy object through usecase"); + }) + .await; + + let info = wait_for_transition(&ecstore, dst_bucket.as_str(), dst_object, TRANSITION_WAIT_TIMEOUT) + .await + .expect("copied object should eventually transition after compensation backfill"); + + assert_eq!(info.transitioned_object.status, "complete"); + 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 compensation_driven_complete_multipart_upload_still_transitions() { + let (_disk_paths, ecstore) = setup_test_env().await; + let usecase = DefaultMultipartUsecase::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-comp-mpu-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/multipart.txt"; + let payload = b"multipart should still transition after compensation"; + + 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 = ecstore + .new_multipart_upload(bucket.as_str(), object, &ObjectOptions::default()) + .await + .expect("Failed to create multipart upload"); + + let mut reader = PutObjReader::from_vec(payload.to_vec()); + let uploaded_part = ecstore + .put_object_part(bucket.as_str(), object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default()) + .await + .expect("Failed to upload multipart part"); + + let complete_input = CompleteMultipartUploadInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .upload_id(upload.upload_id.clone()) + .multipart_upload(Some(CompletedMultipartUpload { + parts: Some(vec![CompletedPart { + part_number: Some(1), + e_tag: uploaded_part.etag.clone().map(|etag| to_s3s_etag(&etag)), + ..Default::default() + }]), + })) + .build() + .unwrap(); + + with_forced_immediate_enqueue_timeout(|| async { + Box::pin(usecase.execute_complete_multipart_upload(build_request(complete_input, Method::POST))) + .await + .expect("Failed to complete multipart upload through usecase"); + }) + .await; + + let info = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT) + .await + .expect("multipart object should eventually transition after compensation backfill"); + + assert_eq!(info.transitioned_object.status, "complete"); + 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 compensation_driven_transition_still_cleans_remote_tier_on_delete() { + 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-compensation-delete-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/object.txt"; + let payload = b"compensation should still preserve delete cleanup"; + + 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"); + + with_forced_immediate_enqueue_timeout(|| async { + let _ = upload_test_object(&ecstore, bucket.as_str(), object, payload).await; + }) + .await; + + let transitioned = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should eventually transition after compensation backfill"); + 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"); + + Box::pin(usecase.execute_delete_object(req)) + .await + .expect("Failed to delete object through usecase after compensation-driven transition"); + + 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" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +#[serial] +#[ignore = "requires isolated global object layer state"] +async fn compensation_driven_versioned_delete_still_creates_delete_marker() { + 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-comp-versioned-delete-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/object.txt"; + let payload = b"versioned delete should preserve transitioned remote version behind delete marker"; + + 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"); + + with_forced_immediate_enqueue_timeout(|| async { + let _ = upload_test_object(&ecstore, bucket.as_str(), object, payload).await; + }) + .await; + + let transitioned = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should eventually transition after compensation backfill"); + let remote_object = transitioned.transitioned_object.name.clone(); + + assert!(backend.objects.lock().await.contains_key(&remote_object)); + + let req = build_request( + DeleteObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .build() + .unwrap(), + Method::DELETE, + ); + + Box::pin(usecase.execute_delete_object(req)) + .await + .expect("Failed to issue versioned delete after compensation-driven transition"); + + assert!( + wait_for_delete_marker(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT).await, + "versioned delete should create a delete marker after compensation-driven transition" + ); + assert!( + backend.objects.lock().await.contains_key(&remote_object), + "creating a delete marker should not remove the transitioned remote object version" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +#[serial] +#[ignore = "requires isolated global object layer state"] +async fn compensation_driven_delete_marker_still_honors_lifecycle_cleanup() { + let (_disk_paths, ecstore) = setup_test_env().await; + let usecase = DefaultObjectUsecase::without_context(); + let bucket_usecase = DefaultBucketUsecase::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-comp-del-marker-cleanup-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/object.txt"; + let payload = b"delete marker lifecycle should still clean up after compensation-driven transition"; + + create_test_bucket(&ecstore, bucket.as_str()).await; + set_bucket_lifecycle_transition_with_tier(bucket.as_str(), &tier_name) + .await + .expect("Failed to set transition lifecycle configuration"); + + with_forced_immediate_enqueue_timeout(|| async { + let _ = upload_test_object(&ecstore, bucket.as_str(), object, payload).await; + }) + .await; + + let transitioned = wait_for_transition(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT) + .await + .expect("object should eventually transition after compensation backfill"); + let remote_object = transitioned.transitioned_object.name.clone(); + + assert!(backend.objects.lock().await.contains_key(&remote_object)); + + let req = build_request( + DeleteObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .build() + .unwrap(), + Method::DELETE, + ); + + Box::pin(usecase.execute_delete_object(req)) + .await + .expect("Failed to issue versioned delete after compensation-driven transition"); + + assert!( + wait_for_delete_marker(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT).await, + "versioned delete should create a delete marker before lifecycle cleanup" + ); + assert!( + backend.objects.lock().await.contains_key(&remote_object), + "delete marker creation should keep the transitioned remote object version" + ); + + let req = build_request( + PutBucketLifecycleConfigurationInput::builder() + .bucket(bucket.clone()) + .lifecycle_configuration(Some(expiration_lifecycle_configuration("test/"))) + .build() + .unwrap(), + Method::PUT, + ); + bucket_usecase + .execute_put_bucket_lifecycle_configuration(req) + .await + .expect("Failed to update lifecycle configuration for delete marker cleanup"); + + assert!( + wait_for_delete_marker(&ecstore, bucket.as_str(), object, TRANSITION_WAIT_TIMEOUT).await, + "delete marker should remain visible after lifecycle update until cleanup completes" + ); + assert!( + backend.objects.lock().await.contains_key(&remote_object), + "delete marker lifecycle cleanup should not remove the transitioned remote object version" + ); +} + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] #[serial] #[ignore = "requires isolated global object layer state"]