diff --git a/.config/nextest.toml b/.config/nextest.toml index 16418d2e6..5af623d81 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -55,6 +55,12 @@ test-group = 'ecstore-serial-flaky' filter = 'package(rustfs-ecstore) & test(/^set_disk::ops::multipart::tests::crash_consistency::/)' test-group = 'ecstore-serial-flaky' +# This regression test builds three multipart objects on a four-disk set and +# coordinates live GET/DELETE races under both lock modes. +[[profile.default.overrides]] +filter = 'package(rustfs-ecstore) & test(multipart_get_keeps_delete_blocked_at_part_boundary_in_both_lock_modes)' +test-group = 'ecstore-serial-flaky' + # Serialize the durable manual-transition checkpoint test across nextest's # process boundary; it mutates bucket lifecycle metadata and is not quarantined. [[profile.default.overrides]] @@ -136,6 +142,10 @@ test-group = 'e2e-reliability' filter = 'package(rustfs-ecstore) & test(/^set_disk::ops::multipart::tests::crash_consistency::/)' test-group = 'ecstore-serial-flaky' +[[profile.ci.overrides]] +filter = 'package(rustfs-ecstore) & test(multipart_get_keeps_delete_blocked_at_part_boundary_in_both_lock_modes)' +test-group = 'ecstore-serial-flaky' + # Serialize the durable manual-transition checkpoint test under the ci profile # too. No retries: failures stay visible. [[profile.ci.overrides]] diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 9731c3bff..d3e3505ba 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -357,6 +357,7 @@ impl Drop for ObjectLockDiagGuard { struct SetDiskLockGuardedReader { inner: Box, guard: Option, + remaining: Option, } impl AsyncRead for SetDiskLockGuardedReader { @@ -364,8 +365,23 @@ impl AsyncRead for SetDiskLockGuardedReader { let had_capacity = buf.remaining() > 0; let filled_before = buf.filled().len(); let poll = Pin::new(&mut self.inner).poll_read(cx, buf); - if had_capacity && matches!(poll, Poll::Ready(Ok(()))) && buf.filled().len() == filled_before { - self.guard.take(); + match &poll { + Poll::Ready(Ok(())) if buf.filled().len() > filled_before => { + let produced = buf.filled().len() - filled_before; + if let Some(remaining) = self.remaining.as_mut() { + *remaining = remaining.saturating_sub(produced); + if *remaining == 0 { + self.guard.take(); + } + } + } + Poll::Ready(Ok(())) if had_capacity => { + self.guard.take(); + } + Poll::Ready(Err(_)) => { + self.guard.take(); + } + _ => {} } poll } @@ -377,15 +393,17 @@ fn finish_set_disk_read_lock( bucket: &str, object: &str, ) -> GetObjectReader { - if reader.buffered_body.is_some() { + if reader.buffered_body.is_some() || reader.object_info.size == 0 { release_materialized_read_lock(bucket, object, read_lock_guard); return reader; } if let Some(guard) = read_lock_guard { + let remaining = usize::try_from(reader.object_info.size).ok(); reader.stream = Box::new(SetDiskLockGuardedReader { inner: reader.stream, guard: Some(guard), + remaining, }); } reader diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 45064ddcd..db1ed52c7 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -2178,11 +2178,70 @@ mod tests { .expect("multipart upload should complete"); } + async fn assert_get_blocks_delete_at_part_boundary( + set_disks: &Arc, + bucket: &str, + object: &'static str, + range: Option, + expected: Vec, + ) { + use crate::set_disk::read::get_part_boundary_barrier; + use http::HeaderMap; + use tokio::io::AsyncReadExt as _; + + let barrier = get_part_boundary_barrier::arm(bucket, object, 0); + let mut reader = set_disks + .get_object_reader(bucket, object, range, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("GET should return a body reader"); + let body_complete = Arc::new(Notify::new()); + let body_complete_task = Arc::clone(&body_complete); + let allow_reader_drop = Arc::new(Notify::new()); + let allow_reader_drop_task = Arc::clone(&allow_reader_drop); + let expected_len = expected.len(); + let read = tokio::spawn(async move { + let mut body = vec![0_u8; expected_len]; + reader.stream.read_exact(&mut body).await?; + body_complete_task.notify_one(); + allow_reader_drop_task.notified().await; + Ok::<_, std::io::Error>(body) + }); + barrier.wait_until_paused().await; + + let delete_set = Arc::clone(set_disks); + let delete_bucket = bucket.to_string(); + let delete_started = Arc::new(Notify::new()); + let delete_started_task = Arc::clone(&delete_started); + let mut 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(); + body_complete.notified().await; + let delete_result = tokio::time::timeout(Duration::from_secs(2), &mut delete).await; + allow_reader_drop.notify_one(); + delete_result + .expect("DELETE must acquire the lock after the exact Content-Length without waiting for reader drop") + .expect("DELETE task should not panic") + .expect("DELETE should proceed after GET delivers its exact Content-Length"); + let body = read + .await + .expect("GET task should not panic") + .expect("GET stream should deliver its exact Content-Length"); + assert_eq!(body.len(), expected.len(), "successful GET body must match the requested length"); + assert_eq!(body, expected, "DELETE must not truncate or mix the multipart snapshot"); + } + #[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"] { @@ -2201,44 +2260,20 @@ mod tests { 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; + assert_get_blocks_delete_at_part_boundary(&set_disks, &bucket, object, None, expected).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 range_object = "range-object"; + put_two_part_object(&set_disks, &bucket, range_object, &first, &second).await; + let range_start = i64::try_from(first.len() - 4096).expect("range start should fit i64"); + let range_end = i64::try_from(first.len() + 4095).expect("range end should fit i64"); + let range = HTTPRangeSpec { + is_suffix_length: false, + start: range_start, + end: range_end, + }; + let range_expected = first[first.len() - 4096..].iter().chain(&second[..4096]).copied().collect(); + assert_get_blocks_delete_at_part_boundary(&set_disks, &bucket, range_object, Some(range), range_expected) + .await; let cancelled_object = "cancelled-object"; put_two_part_object(&set_disks, &bucket, cancelled_object, &first, &second).await; diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 0c875c7f6..1acc5808e 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -89,6 +89,7 @@ impl PreparedGetObjectReader { struct LockGuardedReader { inner: Box, guard: Option, + remaining: Option, } impl AsyncRead for LockGuardedReader { @@ -96,8 +97,23 @@ impl AsyncRead for LockGuardedReader { let had_capacity = buf.remaining() > 0; let filled_before = buf.filled().len(); let poll = Pin::new(&mut self.inner).poll_read(cx, buf); - if had_capacity && matches!(poll, Poll::Ready(Ok(()))) && buf.filled().len() == filled_before { - self.guard.take(); + match &poll { + Poll::Ready(Ok(())) if buf.filled().len() > filled_before => { + let produced = buf.filled().len() - filled_before; + if let Some(remaining) = self.remaining.as_mut() { + *remaining = remaining.saturating_sub(produced); + if *remaining == 0 { + self.guard.take(); + } + } + } + Poll::Ready(Ok(())) if had_capacity => { + self.guard.take(); + } + Poll::Ready(Err(_)) => { + self.guard.take(); + } + _ => {} } poll } @@ -641,14 +657,16 @@ impl ECStore { } fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option) -> GetObjectReader { - if reader.buffered_body.is_some() { + if reader.buffered_body.is_some() || reader.object_info.size == 0 { return reader; } if let Some(guard) = guard { + let remaining = usize::try_from(reader.object_info.size).ok(); reader.stream = Box::new(LockGuardedReader { inner: reader.stream, guard: Some(guard), + remaining, }); } @@ -2649,8 +2667,11 @@ mod tests { ObjectLockDiagMode::Read, ); let reader = GetObjectReader { - stream: Box::new(Cursor::new(Vec::::new())), - object_info: ObjectInfo::default(), + stream: Box::new(Cursor::new(vec![1])), + object_info: ObjectInfo { + size: 1, + ..Default::default() + }, buffered_body: None, body_source: Default::default(), }; @@ -2690,7 +2711,10 @@ mod tests { ); let reader = GetObjectReader { stream: Box::new(Cursor::new(vec![1, 2, 3])), - object_info: ObjectInfo::default(), + object_info: ObjectInfo { + size: 3, + ..Default::default() + }, buffered_body: None, body_source: Default::default(), }; @@ -2747,7 +2771,7 @@ mod tests { #[tokio::test] #[serial_test::serial] - async fn reader_lock_is_released_after_stream_eof_in_both_optimization_states() { + async fn reader_lock_is_released_after_exact_content_length_without_extra_eof_poll() { 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()); @@ -2768,19 +2792,26 @@ mod tests { ); let reader = GetObjectReader { stream: Box::new(Cursor::new(vec![1, 2, 3])), - object_info: ObjectInfo::default(), + object_info: ObjectInfo { + size: 3, + ..Default::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 output = [0_u8; 3]; + reader + .stream + .read_exact(&mut output) + .await + .expect("reader should deliver the exact content length"); + assert_eq!(output, [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"); + .expect("the final response byte should release the read lock without an extra EOF poll"); drop(reader); }) .await;