mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 05:06:28 +00:00
fix(ecstore): release GET locks at content length
This commit is contained in:
@@ -55,6 +55,12 @@ test-group = 'ecstore-serial-flaky'
|
|||||||
filter = 'package(rustfs-ecstore) & test(/^set_disk::ops::multipart::tests::crash_consistency::/)'
|
filter = 'package(rustfs-ecstore) & test(/^set_disk::ops::multipart::tests::crash_consistency::/)'
|
||||||
test-group = 'ecstore-serial-flaky'
|
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
|
# Serialize the durable manual-transition checkpoint test across nextest's
|
||||||
# process boundary; it mutates bucket lifecycle metadata and is not quarantined.
|
# process boundary; it mutates bucket lifecycle metadata and is not quarantined.
|
||||||
[[profile.default.overrides]]
|
[[profile.default.overrides]]
|
||||||
@@ -136,6 +142,10 @@ test-group = 'e2e-reliability'
|
|||||||
filter = 'package(rustfs-ecstore) & test(/^set_disk::ops::multipart::tests::crash_consistency::/)'
|
filter = 'package(rustfs-ecstore) & test(/^set_disk::ops::multipart::tests::crash_consistency::/)'
|
||||||
test-group = 'ecstore-serial-flaky'
|
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
|
# Serialize the durable manual-transition checkpoint test under the ci profile
|
||||||
# too. No retries: failures stay visible.
|
# too. No retries: failures stay visible.
|
||||||
[[profile.ci.overrides]]
|
[[profile.ci.overrides]]
|
||||||
|
|||||||
@@ -357,6 +357,7 @@ impl Drop for ObjectLockDiagGuard {
|
|||||||
struct SetDiskLockGuardedReader {
|
struct SetDiskLockGuardedReader {
|
||||||
inner: Box<dyn AsyncRead + Unpin + Send + Sync>,
|
inner: Box<dyn AsyncRead + Unpin + Send + Sync>,
|
||||||
guard: Option<ObjectLockDiagGuard>,
|
guard: Option<ObjectLockDiagGuard>,
|
||||||
|
remaining: Option<usize>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl AsyncRead for SetDiskLockGuardedReader {
|
impl AsyncRead for SetDiskLockGuardedReader {
|
||||||
@@ -364,8 +365,23 @@ impl AsyncRead for SetDiskLockGuardedReader {
|
|||||||
let had_capacity = buf.remaining() > 0;
|
let had_capacity = buf.remaining() > 0;
|
||||||
let filled_before = buf.filled().len();
|
let filled_before = buf.filled().len();
|
||||||
let poll = Pin::new(&mut self.inner).poll_read(cx, buf);
|
let poll = Pin::new(&mut self.inner).poll_read(cx, buf);
|
||||||
if had_capacity && matches!(poll, Poll::Ready(Ok(()))) && buf.filled().len() == filled_before {
|
match &poll {
|
||||||
self.guard.take();
|
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
|
poll
|
||||||
}
|
}
|
||||||
@@ -377,15 +393,17 @@ fn finish_set_disk_read_lock(
|
|||||||
bucket: &str,
|
bucket: &str,
|
||||||
object: &str,
|
object: &str,
|
||||||
) -> GetObjectReader {
|
) -> 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);
|
release_materialized_read_lock(bucket, object, read_lock_guard);
|
||||||
return reader;
|
return reader;
|
||||||
}
|
}
|
||||||
|
|
||||||
if let Some(guard) = read_lock_guard {
|
if let Some(guard) = read_lock_guard {
|
||||||
|
let remaining = usize::try_from(reader.object_info.size).ok();
|
||||||
reader.stream = Box::new(SetDiskLockGuardedReader {
|
reader.stream = Box::new(SetDiskLockGuardedReader {
|
||||||
inner: reader.stream,
|
inner: reader.stream,
|
||||||
guard: Some(guard),
|
guard: Some(guard),
|
||||||
|
remaining,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
reader
|
reader
|
||||||
|
|||||||
@@ -2178,11 +2178,70 @@ mod tests {
|
|||||||
.expect("multipart upload should complete");
|
.expect("multipart upload should complete");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn assert_get_blocks_delete_at_part_boundary(
|
||||||
|
set_disks: &Arc<SetDisks>,
|
||||||
|
bucket: &str,
|
||||||
|
object: &'static str,
|
||||||
|
range: Option<HTTPRangeSpec>,
|
||||||
|
expected: Vec<u8>,
|
||||||
|
) {
|
||||||
|
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]
|
#[tokio::test]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn multipart_get_keeps_delete_blocked_at_part_boundary_in_both_lock_modes() {
|
async fn multipart_get_keeps_delete_blocked_at_part_boundary_in_both_lock_modes() {
|
||||||
use crate::set_disk::read::get_part_boundary_barrier;
|
use crate::set_disk::read::get_part_boundary_barrier;
|
||||||
use http::HeaderMap;
|
|
||||||
use tokio::io::AsyncReadExt as _;
|
use tokio::io::AsyncReadExt as _;
|
||||||
|
|
||||||
for enabled in ["false", "true"] {
|
for enabled in ["false", "true"] {
|
||||||
@@ -2201,44 +2260,20 @@ mod tests {
|
|||||||
let second = vec![0x42; 64 * 1024];
|
let second = vec![0x42; 64 * 1024];
|
||||||
let expected: Vec<u8> = first.iter().chain(&second).copied().collect();
|
let expected: Vec<u8> = first.iter().chain(&second).copied().collect();
|
||||||
put_two_part_object(&set_disks, &bucket, object, &first, &second).await;
|
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 range_object = "range-object";
|
||||||
let mut reader = set_disks
|
put_two_part_object(&set_disks, &bucket, range_object, &first, &second).await;
|
||||||
.get_object_reader(&bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
|
let range_start = i64::try_from(first.len() - 4096).expect("range start should fit i64");
|
||||||
.await
|
let range_end = i64::try_from(first.len() + 4095).expect("range end should fit i64");
|
||||||
.expect("GET should return a body reader");
|
let range = HTTPRangeSpec {
|
||||||
let content_length = usize::try_from(reader.object_info.size).expect("object size should fit usize");
|
is_suffix_length: false,
|
||||||
let read = tokio::spawn(async move {
|
start: range_start,
|
||||||
let mut body = Vec::new();
|
end: range_end,
|
||||||
reader.stream.read_to_end(&mut body).await.map(|_| body)
|
};
|
||||||
});
|
let range_expected = first[first.len() - 4096..].iter().chain(&second[..4096]).copied().collect();
|
||||||
barrier.wait_until_paused().await;
|
assert_get_blocks_delete_at_part_boundary(&set_disks, &bucket, range_object, Some(range), range_expected)
|
||||||
|
.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";
|
let cancelled_object = "cancelled-object";
|
||||||
put_two_part_object(&set_disks, &bucket, cancelled_object, &first, &second).await;
|
put_two_part_object(&set_disks, &bucket, cancelled_object, &first, &second).await;
|
||||||
|
|||||||
@@ -89,6 +89,7 @@ impl PreparedGetObjectReader {
|
|||||||
struct LockGuardedReader {
|
struct LockGuardedReader {
|
||||||
inner: Box<dyn AsyncRead + Unpin + Send + Sync>,
|
inner: Box<dyn AsyncRead + Unpin + Send + Sync>,
|
||||||
guard: Option<ObjectLockDiagGuard>,
|
guard: Option<ObjectLockDiagGuard>,
|
||||||
|
remaining: Option<usize>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl AsyncRead for LockGuardedReader {
|
impl AsyncRead for LockGuardedReader {
|
||||||
@@ -96,8 +97,23 @@ impl AsyncRead for LockGuardedReader {
|
|||||||
let had_capacity = buf.remaining() > 0;
|
let had_capacity = buf.remaining() > 0;
|
||||||
let filled_before = buf.filled().len();
|
let filled_before = buf.filled().len();
|
||||||
let poll = Pin::new(&mut self.inner).poll_read(cx, buf);
|
let poll = Pin::new(&mut self.inner).poll_read(cx, buf);
|
||||||
if had_capacity && matches!(poll, Poll::Ready(Ok(()))) && buf.filled().len() == filled_before {
|
match &poll {
|
||||||
self.guard.take();
|
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
|
poll
|
||||||
}
|
}
|
||||||
@@ -641,14 +657,16 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option<ObjectLockDiagGuard>) -> GetObjectReader {
|
fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option<ObjectLockDiagGuard>) -> GetObjectReader {
|
||||||
if reader.buffered_body.is_some() {
|
if reader.buffered_body.is_some() || reader.object_info.size == 0 {
|
||||||
return reader;
|
return reader;
|
||||||
}
|
}
|
||||||
|
|
||||||
if let Some(guard) = guard {
|
if let Some(guard) = guard {
|
||||||
|
let remaining = usize::try_from(reader.object_info.size).ok();
|
||||||
reader.stream = Box::new(LockGuardedReader {
|
reader.stream = Box::new(LockGuardedReader {
|
||||||
inner: reader.stream,
|
inner: reader.stream,
|
||||||
guard: Some(guard),
|
guard: Some(guard),
|
||||||
|
remaining,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2649,8 +2667,11 @@ mod tests {
|
|||||||
ObjectLockDiagMode::Read,
|
ObjectLockDiagMode::Read,
|
||||||
);
|
);
|
||||||
let reader = GetObjectReader {
|
let reader = GetObjectReader {
|
||||||
stream: Box::new(Cursor::new(Vec::<u8>::new())),
|
stream: Box::new(Cursor::new(vec![1])),
|
||||||
object_info: ObjectInfo::default(),
|
object_info: ObjectInfo {
|
||||||
|
size: 1,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
buffered_body: None,
|
buffered_body: None,
|
||||||
body_source: Default::default(),
|
body_source: Default::default(),
|
||||||
};
|
};
|
||||||
@@ -2690,7 +2711,10 @@ mod tests {
|
|||||||
);
|
);
|
||||||
let reader = GetObjectReader {
|
let reader = GetObjectReader {
|
||||||
stream: Box::new(Cursor::new(vec![1, 2, 3])),
|
stream: Box::new(Cursor::new(vec![1, 2, 3])),
|
||||||
object_info: ObjectInfo::default(),
|
object_info: ObjectInfo {
|
||||||
|
size: 3,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
buffered_body: None,
|
buffered_body: None,
|
||||||
body_source: Default::default(),
|
body_source: Default::default(),
|
||||||
};
|
};
|
||||||
@@ -2747,7 +2771,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[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"] {
|
for enabled in ["false", "true"] {
|
||||||
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some(enabled))], async {
|
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some(enabled))], async {
|
||||||
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
|
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
|
||||||
@@ -2768,19 +2792,26 @@ mod tests {
|
|||||||
);
|
);
|
||||||
let reader = GetObjectReader {
|
let reader = GetObjectReader {
|
||||||
stream: Box::new(Cursor::new(vec![1, 2, 3])),
|
stream: Box::new(Cursor::new(vec![1, 2, 3])),
|
||||||
object_info: ObjectInfo::default(),
|
object_info: ObjectInfo {
|
||||||
|
size: 3,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
buffered_body: None,
|
buffered_body: None,
|
||||||
body_source: Default::default(),
|
body_source: Default::default(),
|
||||||
};
|
};
|
||||||
|
|
||||||
let mut reader = ECStore::attach_read_lock_guard(reader, Some(read_guard));
|
let mut reader = ECStore::attach_read_lock_guard(reader, Some(read_guard));
|
||||||
let mut output = Vec::new();
|
let mut output = [0_u8; 3];
|
||||||
reader.stream.read_to_end(&mut output).await.expect("reader should reach EOF");
|
reader
|
||||||
assert_eq!(output, vec![1, 2, 3]);
|
.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))
|
lock.get_write_lock(key, "writer", Duration::from_secs(1))
|
||||||
.await
|
.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);
|
drop(reader);
|
||||||
})
|
})
|
||||||
.await;
|
.await;
|
||||||
|
|||||||
Reference in New Issue
Block a user