diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index bdd66c945..9731c3bff 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -374,11 +374,10 @@ impl AsyncRead for SetDiskLockGuardedReader { fn finish_set_disk_read_lock( mut reader: GetObjectReader, read_lock_guard: Option, - lock_optimization_enabled: bool, bucket: &str, object: &str, ) -> GetObjectReader { - if lock_optimization_enabled || reader.buffered_body.is_some() { + if reader.buffered_body.is_some() { release_materialized_read_lock(bucket, object, read_lock_guard); return reader; } diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index a3228e079..45064ddcd 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -2147,6 +2147,128 @@ mod tests { } } + async fn put_two_part_object(set_disks: &Arc, bucket: &str, object: &str, first: &[u8], second: &[u8]) { + use crate::storage_api_contracts::multipart::MultipartOperations as _; + + let upload = set_disks + .new_multipart_upload(bucket, object, &ObjectOptions::default()) + .await + .expect("multipart upload should be created"); + let mut completed = Vec::new(); + for (part_number, payload) in [(1, first), (2, second)] { + let payload_len = i64::try_from(payload.len()).expect("test payload length should fit i64"); + let mut reader = PutObjReader::new( + HashReader::from_stream(Cursor::new(payload.to_vec()), payload_len, payload_len, None, None, false) + .expect("part hash reader should be created"), + ); + let part = set_disks + .put_object_part(bucket, object, &upload.upload_id, part_number, &mut reader, &ObjectOptions::default()) + .await + .expect("multipart part should be written"); + completed.push(CompletePart { + part_num: part.part_num, + etag: part.etag, + ..Default::default() + }); + } + set_disks + .clone() + .complete_multipart_upload(bucket, object, &upload.upload_id, completed, &ObjectOptions::default()) + .await + .expect("multipart upload should complete"); + } + + #[tokio::test] + #[serial] + async fn multipart_get_keeps_delete_blocked_at_part_boundary_in_both_lock_modes() { + use crate::set_disk::read::get_part_boundary_barrier; + use http::HeaderMap; + use tokio::io::AsyncReadExt as _; + + for enabled in ["false", "true"] { + temp_env::async_with_vars( + [ + (rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some(enabled)), + (ENV_RUSTFS_GET_CODEC_STREAMING_MULTIPART_ENABLE, Some("false")), + ], + async { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = format!("get-delete-part-boundary-{enabled}"); + let object = "object"; + make_bucket_on_all(&disk_stores, &bucket).await; + + let first = vec![0x41; GLOBAL_MIN_PART_SIZE.as_u64() as usize]; + let second = vec![0x42; 64 * 1024]; + let expected: Vec = first.iter().chain(&second).copied().collect(); + put_two_part_object(&set_disks, &bucket, object, &first, &second).await; + + let barrier = get_part_boundary_barrier::arm(&bucket, object, 0); + let mut reader = set_disks + .get_object_reader(&bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("GET should return a body reader"); + let content_length = usize::try_from(reader.object_info.size).expect("object size should fit usize"); + let read = tokio::spawn(async move { + let mut body = Vec::new(); + reader.stream.read_to_end(&mut body).await.map(|_| body) + }); + barrier.wait_until_paused().await; + + let delete_set = Arc::clone(&set_disks); + let delete_bucket = bucket.clone(); + let delete_started = Arc::new(Notify::new()); + let delete_started_task = Arc::clone(&delete_started); + let delete = tokio::spawn(async move { + delete_started_task.notify_one(); + delete_set + .delete_object(&delete_bucket, object, ObjectOptions::default()) + .await + }); + delete_started.notified().await; + tokio::task::yield_now().await; + assert!(!delete.is_finished(), "DELETE must wait behind the streaming GET read lock"); + + barrier.release(); + let body = read + .await + .expect("GET task should not panic") + .expect("GET stream should reach EOF"); + assert_eq!(body.len(), content_length, "successful GET body must match Content-Length"); + assert_eq!(body, expected, "DELETE must not truncate or mix the multipart snapshot"); + delete + .await + .expect("DELETE task should not panic") + .expect("DELETE should proceed after GET reaches EOF"); + + let cancelled_object = "cancelled-object"; + put_two_part_object(&set_disks, &bucket, cancelled_object, &first, &second).await; + let cancel_barrier = get_part_boundary_barrier::arm(&bucket, cancelled_object, 0); + let mut cancelled_reader = set_disks + .get_object_reader(&bucket, cancelled_object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("cancelled GET should return a body reader"); + let cancelled_read = tokio::spawn(async move { + let mut body = Vec::new(); + cancelled_reader.stream.read_to_end(&mut body).await + }); + cancel_barrier.wait_until_paused().await; + cancelled_read.abort(); + cancelled_read.await.expect_err("aborted GET task should report cancellation"); + cancel_barrier.release(); + + tokio::time::timeout( + Duration::from_secs(10), + set_disks.delete_object(&bucket, cancelled_object, ObjectOptions::default()), + ) + .await + .expect("DELETE should promptly acquire the lock after GET cancellation") + .expect("DELETE after cancellation should succeed"); + }, + ) + .await; + } + } + async fn stage_upload_with_create_opts( set_disks: &Arc, bucket: &str, diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 757af70b4..02169cb35 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -426,13 +426,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { &self.ctx.tier_config_mgr(), ) .await?; - return Ok(finish_set_disk_read_lock( - gr, - read_lock_guard.take(), - lock_optimization_enabled, - bucket, - object, - )); + return Ok(finish_set_disk_read_lock(gr, read_lock_guard.take(), bucket, object)); } // App-layer object data cache probe: metadata (etag/size) is resolved @@ -587,13 +581,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { // Carry the hook probe result so the app layer skips its // now-redundant lookup on the streaming miss path (ODC-16). reader.body_source = body_source; - return Ok(finish_set_disk_read_lock( - reader, - read_lock_guard.take(), - lock_optimization_enabled, - bucket, - object, - )); + return Ok(finish_set_disk_read_lock(reader, read_lock_guard.take(), bucket, object)); } core::io_primitives::GetCodecStreamingReaderBuildOutcome::Fallback(reason) => { record_get_codec_streaming_gate_decision( @@ -632,13 +620,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { let set_index = self.set_index; let pool_index = self.pool_index; let skip_verify = opts.skip_verify_bitrot; - if lock_optimization_enabled { - release_materialized_read_lock(&bucket, &object, read_lock_guard.take()); - debug!(bucket, object, "Lock optimization: released read lock before streaming read"); - } - - // When lock optimization is disabled, keep the read-lock guard in the - // task so it lives for the duration of the streaming read. + // Keep the read lock until the producer reaches EOF or observes that + // the downstream reader was cancelled. tokio::spawn(async move { let _guard = read_lock_guard; let mut writer = GetObjectDownstreamWriter::new(wd); diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 1a33716f4..a8fe25984 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -1117,6 +1117,10 @@ impl SetDisks { total_read += part_length; part_offset = 0; + #[cfg(test)] + if current_part < last_part_index { + get_part_boundary_barrier::checkpoint(bucket, object, current_part).await; + } } // debug!("read end"); @@ -1820,6 +1824,79 @@ fn is_get_object_metadata_cache_request_eligible(bucket: &str, opts: &ObjectOpti get_object_metadata_cache_request_bypass_reason(bucket, opts, read_data).is_none() } +#[cfg(test)] +pub(in crate::set_disk) mod get_part_boundary_barrier { + use std::collections::HashMap; + use std::sync::{Arc, Mutex, OnceLock}; + use tokio::sync::Notify; + + struct Armed { + part_index: usize, + arrived: Arc, + release: Arc, + } + + fn registry() -> &'static Mutex> { + static REGISTRY: OnceLock>> = OnceLock::new(); + REGISTRY.get_or_init(|| Mutex::new(HashMap::new())) + } + + pub struct BarrierHandle { + key: (String, String), + arrived: Arc, + release: Arc, + } + + pub fn arm(bucket: &str, object: &str, part_index: usize) -> BarrierHandle { + let key = (bucket.to_string(), object.to_string()); + let arrived = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let previous = registry().lock().expect("GET part barrier registry poisoned").insert( + key.clone(), + Armed { + part_index, + arrived: Arc::clone(&arrived), + release: Arc::clone(&release), + }, + ); + assert!(previous.is_none(), "GET part barrier already armed for object"); + BarrierHandle { key, arrived, release } + } + + impl BarrierHandle { + pub async fn wait_until_paused(&self) { + self.arrived.notified().await; + } + + pub fn release(&self) { + self.release.notify_one(); + } + } + + impl Drop for BarrierHandle { + fn drop(&mut self) { + self.release.notify_one(); + registry() + .lock() + .expect("GET part barrier registry poisoned") + .remove(&self.key); + } + } + + pub async fn checkpoint(bucket: &str, object: &str, part_index: usize) { + let state = registry() + .lock() + .expect("GET part barrier registry poisoned") + .get(&(bucket.to_string(), object.to_string())) + .filter(|armed| armed.part_index == part_index) + .map(|armed| (Arc::clone(&armed.arrived), Arc::clone(&armed.release))); + if let Some((arrived, release)) = state { + arrived.notify_one(); + release.notified().await; + } + } +} + #[cfg(test)] mod metadata_cache_tests { use super::*; diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 2a8d30ee4..0c875c7f6 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -16,7 +16,7 @@ use super::*; use crate::disk::OldCurrentSize; use crate::set_disk::{ get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold, - is_lock_optimization_enabled, is_object_lock_diag_enabled, + is_object_lock_diag_enabled, }; use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; use rustfs_io_metrics::{ @@ -641,7 +641,7 @@ impl ECStore { } fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option) -> GetObjectReader { - if is_lock_optimization_enabled() || reader.buffered_body.is_some() { + if reader.buffered_body.is_some() { return reader; } @@ -2670,7 +2670,7 @@ mod tests { #[tokio::test] #[serial_test::serial] - async fn reader_lock_is_not_held_for_stream_when_optimization_is_enabled() { + async fn reader_lock_is_held_for_stream_when_optimization_is_enabled() { temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async { let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); let lock = rustfs_lock::NamespaceLock::with_local_manager("test".to_string(), manager); @@ -2697,10 +2697,13 @@ mod tests { let reader = ECStore::attach_read_lock_guard(reader, Some(read_guard)); + lock.get_write_lock(key.clone(), "writer", Duration::from_millis(20)) + .await + .expect_err("streaming readers must retain the read lock when lock optimization is enabled"); + drop(reader); lock.get_write_lock(key, "writer", Duration::from_secs(1)) .await - .expect("lock optimization should release the read lock before returning the stream"); - drop(reader); + .expect("cancelling the reader should release the read lock"); }) .await; } @@ -2744,41 +2747,43 @@ mod tests { #[tokio::test] #[serial_test::serial] - async fn reader_lock_is_released_after_stream_eof() { - temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("false"))], async { - let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); - let lock = rustfs_lock::NamespaceLock::with_local_manager("test".to_string(), manager); - let key = rustfs_lock::ObjectKey::new("bucket", "object"); - let read_guard = lock - .get_read_lock(key.clone(), "reader", Duration::from_secs(1)) - .await - .expect("read lock should be acquired"); - let read_guard = ObjectLockDiagGuard::new( - read_guard, - true, - "test_get_object", - Some("bucket".to_string()), - Some("object".to_string()), - Some("reader".to_string()), - ObjectLockDiagMode::Read, - ); - let reader = GetObjectReader { - stream: Box::new(Cursor::new(vec![1, 2, 3])), - object_info: ObjectInfo::default(), - buffered_body: None, - body_source: Default::default(), - }; + async fn reader_lock_is_released_after_stream_eof_in_both_optimization_states() { + for enabled in ["false", "true"] { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some(enabled))], async { + let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); + let lock = rustfs_lock::NamespaceLock::with_local_manager("test".to_string(), manager); + let key = rustfs_lock::ObjectKey::new("bucket", "object"); + let read_guard = lock + .get_read_lock(key.clone(), "reader", Duration::from_secs(1)) + .await + .expect("read lock should be acquired"); + let read_guard = ObjectLockDiagGuard::new( + read_guard, + true, + "test_get_object", + Some("bucket".to_string()), + Some("object".to_string()), + Some("reader".to_string()), + ObjectLockDiagMode::Read, + ); + let reader = GetObjectReader { + stream: Box::new(Cursor::new(vec![1, 2, 3])), + object_info: ObjectInfo::default(), + buffered_body: None, + body_source: Default::default(), + }; - let mut reader = ECStore::attach_read_lock_guard(reader, Some(read_guard)); - let mut output = Vec::new(); - reader.stream.read_to_end(&mut output).await.expect("reader should reach EOF"); - assert_eq!(output, vec![1, 2, 3]); + let mut reader = ECStore::attach_read_lock_guard(reader, Some(read_guard)); + let mut output = Vec::new(); + reader.stream.read_to_end(&mut output).await.expect("reader should reach EOF"); + assert_eq!(output, vec![1, 2, 3]); - lock.get_write_lock(key, "writer", Duration::from_secs(1)) - .await - .expect("EOF should release the read lock before the reader is dropped"); - drop(reader); - }) - .await; + lock.get_write_lock(key, "writer", Duration::from_secs(1)) + .await + .expect("EOF should release the read lock before the reader is dropped"); + drop(reader); + }) + .await; + } } }