fix(ecstore): retain GET locks through streaming

This commit is contained in:
马登山
2026-07-28 12:22:06 +08:00
parent 5cedab09ab
commit cf1e1a67e0
5 changed files with 248 additions and 62 deletions
+1 -2
View File
@@ -374,11 +374,10 @@ impl AsyncRead for SetDiskLockGuardedReader {
fn finish_set_disk_read_lock(
mut reader: GetObjectReader,
read_lock_guard: Option<ObjectLockDiagGuard>,
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;
}
@@ -2147,6 +2147,128 @@ mod tests {
}
}
async fn put_two_part_object(set_disks: &Arc<SetDisks>, 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<u8> = 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<SetDisks>,
bucket: &str,
+4 -21
View File
@@ -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);
+77
View File
@@ -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<Notify>,
release: Arc<Notify>,
}
fn registry() -> &'static Mutex<HashMap<(String, String), Armed>> {
static REGISTRY: OnceLock<Mutex<HashMap<(String, String), Armed>>> = OnceLock::new();
REGISTRY.get_or_init(|| Mutex::new(HashMap::new()))
}
pub struct BarrierHandle {
key: (String, String),
arrived: Arc<Notify>,
release: Arc<Notify>,
}
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::*;
+44 -39
View File
@@ -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<ObjectLockDiagGuard>) -> 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;
}
}
}