fix(ecstore): hold read locks for streaming readers (#4195)

This commit is contained in:
cxymds
2026-07-02 22:36:12 +08:00
committed by GitHub
parent 91a02914cc
commit cf33bcc742
3 changed files with 137 additions and 36 deletions
+2 -3
View File
@@ -191,9 +191,8 @@ pub const DEFAULT_OBJECT_IO_BUFFER_SIZE: usize = 128 * 1024;
/// Environment variable to enable/disable lock optimization.
///
/// When enabled, read locks are released immediately after metadata
/// is read, rather than being held for the entire data transfer.
/// This significantly reduces lock contention under high concurrency.
/// When enabled, fully materialized reads may release read locks before the
/// reader is returned. Streaming reads keep the lock until EOF or drop.
///
/// Default: true (enabled, can be overridden by `RUSTFS_OBJECT_LOCK_OPTIMIZATION_ENABLE`).
pub const ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE: &str = "RUSTFS_OBJECT_LOCK_OPTIMIZATION_ENABLE";
+58 -31
View File
@@ -132,7 +132,9 @@ use s3s::header::{X_AMZ_OBJECT_LOCK_LEGAL_HOLD, X_AMZ_OBJECT_LOCK_MODE, X_AMZ_OB
use sha2::{Digest, Sha256};
use std::hash::Hash;
use std::mem::{self};
use std::pin::Pin;
use std::sync::OnceLock;
use std::task::{Context, Poll};
use std::time::{Instant, SystemTime, UNIX_EPOCH};
use std::{
collections::{HashMap, HashSet},
@@ -143,7 +145,7 @@ use std::{
};
use time::OffsetDateTime;
use tokio::{
io::{AsyncReadExt, AsyncWrite, AsyncWriteExt, BufReader},
io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, BufReader, ReadBuf},
sync::{RwLock, broadcast},
};
use tokio::{
@@ -251,6 +253,44 @@ impl Drop for ObjectLockDiagGuard {
}
}
struct SetDiskLockGuardedReader {
inner: Box<dyn AsyncRead + Unpin + Send + Sync>,
guard: Option<ObjectLockDiagGuard>,
}
impl AsyncRead for SetDiskLockGuardedReader {
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
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();
}
poll
}
}
fn attach_set_disk_read_lock_guard(mut reader: GetObjectReader, read_lock_guard: Option<ObjectLockDiagGuard>) -> GetObjectReader {
if let Some(guard) = read_lock_guard
&& reader.buffered_body.is_none()
{
reader.stream = Box::new(SetDiskLockGuardedReader {
inner: reader.stream,
guard: Some(guard),
});
}
reader
}
fn release_materialized_read_lock(bucket: &str, object: &str, read_lock_guard: Option<ObjectLockDiagGuard>) {
if read_lock_guard.is_some() {
let lock_id = format!("{}:{}", bucket, object);
record_lock_release(bucket, object, &lock_id, "read");
metrics::counter!("rustfs.lock.release.early.total", "type" => "read").increment(1);
}
drop(read_lock_guard);
}
pub(crate) fn strip_internal_multipart_metadata(metadata: &mut HashMap<String, String>) {
metadata.remove(RUSTFS_MULTIPART_BUCKET_KEY);
metadata.remove(RUSTFS_MULTIPART_OBJECT_KEY);
@@ -437,8 +477,8 @@ pub fn get_object_lock_diag_slow_hold_threshold() -> Duration {
}
/// Check if lock optimization is enabled.
/// When enabled, read locks are released after metadata read instead of
/// being held for the entire data transfer duration.
/// When enabled, fully materialized reads may release the read lock before
/// returning to the caller. Streaming reads keep the lock until EOF or drop.
///
/// **Note**: Cached via `OnceLock` in production — env var changes require
/// process restart. In test builds the env var is read directly so that
@@ -1999,12 +2039,11 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
opts: &ObjectOptions,
) -> Result<GetObjectReader> {
let stage_metrics_enabled = rustfs_io_metrics::get_stage_metrics_enabled();
// Check if lock optimization is enabled
// When enabled, read locks are released after metadata read
// Check if lock optimization is enabled for reads that are fully materialized in memory.
let lock_optimization_enabled = is_lock_optimization_enabled();
// Acquire a shared read-lock early to protect read consistency
let read_lock_guard = if !opts.no_lock {
let mut read_lock_guard = if !opts.no_lock {
let acquire_start = Instant::now();
let lock_stage_start = get_stage_timer_if_enabled(stage_metrics_enabled);
@@ -2275,28 +2314,9 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
opts.part_number = Some(1);
}
let gr = get_transitioned_object_reader(bucket, object, &range, &h, &object_info, &opts).await?;
return Ok(gr);
return Ok(attach_set_disk_read_lock_guard(gr, read_lock_guard.take()));
}
// Lock optimization: release read lock after metadata read if enabled
// This reduces lock contention by not holding the lock during data transfer
let read_lock_guard = if lock_optimization_enabled {
// Record lock release for deadlock detection
if read_lock_guard.is_some() {
let lock_id = format!("{}:{}", bucket, object);
record_lock_release(bucket, object, &lock_id, "read");
// Record early lock release statistics
metrics::counter!("rustfs.lock.release.early.total", "type" => "read").increment(1);
}
// Explicitly drop the lock guard to release the lock early
drop(read_lock_guard);
debug!(bucket, object, "Lock optimization: released read lock after metadata read");
None
} else {
read_lock_guard
};
if is_get_small_object_direct_memory_eligible(&range, &object_info, &fi, opts) {
let object_size = usize::try_from(object_info.size)
.map_err(|_| to_object_err(Error::other("direct-memory GET object size is invalid"), vec![bucket, object]))?;
@@ -2325,6 +2345,10 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
object_info,
buffered_body: Some(body),
};
if lock_optimization_enabled {
release_materialized_read_lock(bucket, object, read_lock_guard.take());
debug!(bucket, object, "Lock optimization: released read lock after direct-memory read");
}
return Ok(reader);
}
@@ -2362,6 +2386,10 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
object_info,
buffered_body: Some(body),
};
if lock_optimization_enabled {
release_materialized_read_lock(bucket, object, read_lock_guard.take());
debug!(bucket, object, "Lock optimization: released read lock after direct-memory read");
}
return Ok(reader);
}
@@ -2390,7 +2418,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
);
record_get_object_reader_path_observation(GET_OBJECT_PATH_CODEC_STREAMING, object_class, size_bucket);
let (reader, _offset, _length) = GetObjectReader::new(stream, range, &object_info, opts, &h).await?;
return Ok(reader);
return Ok(attach_set_disk_read_lock_guard(reader, read_lock_guard.take()));
}
read::GetCodecStreamingReaderBuildOutcome::Fallback(reason) => {
record_get_codec_streaming_gate_decision(
@@ -2426,11 +2454,10 @@ 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;
// Move the read-lock guard into the task so it lives for the duration of the read
// Note: when lock optimization is enabled, read_lock_guard is None
// let _guard_to_hold = _read_lock_guard; // moved into closure below
// Move the read-lock guard into the task so it lives for the duration of the read.
// Fully materialized paths release it before returning; streaming paths keep it.
tokio::spawn(async move {
let _guard = read_lock_guard; // keep guard alive until task ends (None if optimization enabled)
let _guard = read_lock_guard;
let mut writer = wd;
// Do not wrap the entire read+write pipeline in `disk_read_timeout`.
// `get_object_with_fileinfo` also waits on `writer`, so an outer timeout
+77 -2
View File
@@ -15,7 +15,7 @@
use super::*;
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::{
@@ -442,7 +442,7 @@ impl ECStore {
}
fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option<ObjectLockDiagGuard>) -> GetObjectReader {
if is_lock_optimization_enabled() {
if reader.buffered_body.is_some() {
return reader;
}
@@ -1932,6 +1932,81 @@ mod tests {
.await;
}
#[tokio::test]
#[serial_test::serial]
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);
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,
};
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 reader should hold the read lock");
drop(reader);
lock.get_write_lock(key, "writer", Duration::from_secs(1))
.await
.expect("dropping the reader should release the read lock");
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn reader_lock_is_not_held_for_buffered_body_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);
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: Some(Bytes::from_static(b"123")),
};
let reader = ECStore::attach_read_lock_guard(reader, Some(read_guard));
lock.get_write_lock(key, "writer", Duration::from_secs(1))
.await
.expect("buffered reader should release the read lock immediately");
drop(reader);
})
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn reader_lock_is_released_after_stream_eof() {