feat(ecstore): add object lock diagnostics (#3178)

* feat(ecstore): add object lock diagnostics

Add configurable namespace lock diagnostics for object operations so production contention can be traced by operation, owner, and key.

Wrap object read/write lock acquisition in diagnostic guards across get, head, put, delete, copy, and multipart flows, and log slow acquisition and long hold durations behind new RUSTFS_OBJECT_LOCK_DIAG_* settings.

Verification:

- make pre-commit

* feat(obs): expose object lock diagnostics metrics

Add Prometheus metrics and Grafana panels for object namespace lock diagnostics, covering slow acquire counts, slow hold counts, acquire duration, hold duration, and the diagnostics-enabled state.

Adopt PR review feedback by keeping diagnostic guards alive through the guarded operation so long-hold warnings and metrics are emitted, reusing shared env helpers, and reducing default-path overhead when diagnostics are disabled.

Verification:

- cargo check -p rustfs-ecstore -p rustfs-io-metrics

- cargo test -p rustfs-ecstore store::object -- --nocapture

- make pre-commit

* perf(obs): reduce object lock diag overhead

Avoid repeated environment parsing on hot object-lock paths by caching the diagnostics-enabled flag, and stop allocating label strings for object lock metrics by recording static labels directly.

Strengthen io-metrics tests by using a local recorder and asserting that the expected object lock diagnostic counters, gauges, and histograms are emitted.

Verification:

- cargo check -p rustfs-ecstore -p rustfs-io-metrics

- cargo test -p rustfs-io-metrics -- --nocapture

- make pre-commit
This commit is contained in:
houseme
2026-06-03 10:10:27 +08:00
committed by GitHub
parent 0dbf0b13a8
commit 29fbdc2dbf
7 changed files with 1037 additions and 92 deletions
+198 -11
View File
@@ -13,16 +13,25 @@
// limitations under the License.
use super::*;
use crate::set_disk::{get_lock_acquire_timeout, is_lock_optimization_enabled};
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,
};
use rustfs_io_metrics::{
record_object_lock_diag_acquire_duration, record_object_lock_diag_hold_duration, record_object_lock_diag_slow_acquire,
record_object_lock_diag_slow_hold,
};
use std::{
fmt,
pin::Pin,
task::{Context, Poll},
time::{Duration, Instant},
};
use tokio::io::{AsyncRead, ReadBuf};
struct LockGuardedReader {
inner: Box<dyn AsyncRead + Unpin + Send + Sync>,
guard: Option<rustfs_lock::NamespaceLockGuard>,
guard: Option<ObjectLockDiagGuard>,
}
impl AsyncRead for LockGuardedReader {
@@ -37,6 +46,118 @@ impl AsyncRead for LockGuardedReader {
}
}
#[derive(Clone, Copy, Debug)]
enum ObjectLockDiagMode {
Read,
Write,
}
impl ObjectLockDiagMode {
fn as_str(self) -> &'static str {
match self {
Self::Read => "read",
Self::Write => "write",
}
}
}
impl fmt::Display for ObjectLockDiagMode {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
struct ObjectLockDiagGuard {
guard: rustfs_lock::NamespaceLockGuard,
enabled: bool,
op: &'static str,
bucket: Option<String>,
object: Option<String>,
owner: Option<String>,
mode: ObjectLockDiagMode,
acquired_at: Instant,
}
impl ObjectLockDiagGuard {
fn new(
guard: rustfs_lock::NamespaceLockGuard,
enabled: bool,
op: &'static str,
bucket: Option<String>,
object: Option<String>,
owner: Option<String>,
mode: ObjectLockDiagMode,
) -> Self {
Self {
guard,
enabled,
op,
bucket,
object,
owner,
mode,
acquired_at: Instant::now(),
}
}
}
impl Drop for ObjectLockDiagGuard {
fn drop(&mut self) {
if !self.enabled || self.guard.is_released() {
return;
}
let hold = self.acquired_at.elapsed();
record_object_lock_diag_hold_duration(self.op, self.mode.as_str(), hold);
let threshold = get_object_lock_diag_slow_hold_threshold();
if hold >= threshold {
record_object_lock_diag_slow_hold(self.op, self.mode.as_str());
warn!(
target: "rustfs_ecstore::object_lock_diag",
op = self.op,
bucket = %self.bucket.as_deref().unwrap_or_default(),
object = %self.object.as_deref().unwrap_or_default(),
mode = %self.mode,
owner = %self.owner.as_deref().unwrap_or_default(),
hold_ms = hold.as_millis(),
threshold_ms = threshold.as_millis(),
"object namespace lock held longer than threshold"
);
}
}
}
fn log_object_lock_acquire_if_slow(
op: &'static str,
bucket: &str,
object: &str,
owner: Option<&str>,
mode: ObjectLockDiagMode,
elapsed: Duration,
diag_enabled: bool,
) {
if !diag_enabled {
return;
}
let threshold = get_object_lock_diag_slow_acquire_threshold();
record_object_lock_diag_acquire_duration(op, mode.as_str(), elapsed);
if elapsed >= threshold {
record_object_lock_diag_slow_acquire(op, mode.as_str());
warn!(
target: "rustfs_ecstore::object_lock_diag",
op,
bucket,
object,
mode = %mode,
owner = owner.unwrap_or_default(),
acquire_ms = elapsed.as_millis(),
threshold_ms = threshold.as_millis(),
"object namespace lock acquisition exceeded threshold"
);
}
}
fn select_data_movement_target_pool(
existing_pool_idx: Result<usize>,
src_pool_idx: usize,
@@ -128,45 +249,87 @@ impl ECStore {
async fn acquire_object_write_lock_if_needed(
&self,
op: &'static str,
bucket: &str,
object: &str,
opts: &mut ObjectOptions,
) -> Result<Option<rustfs_lock::NamespaceLockGuard>> {
) -> Result<Option<ObjectLockDiagGuard>> {
if opts.no_lock {
return Ok(None);
}
let diag_enabled = is_object_lock_diag_enabled();
let ns_lock = self.handle_new_ns_lock(bucket, object).await?;
let acquire_start = Instant::now();
let guard = ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(|err| Self::map_namespace_lock_error(bucket, object, "write", err))?;
let owner = diag_enabled.then(|| ns_lock.owner().to_string());
log_object_lock_acquire_if_slow(
op,
bucket,
object,
owner.as_deref(),
ObjectLockDiagMode::Write,
acquire_start.elapsed(),
diag_enabled,
);
opts.no_lock = true;
Ok(Some(guard))
Ok(Some(ObjectLockDiagGuard::new(
guard,
diag_enabled,
op,
diag_enabled.then(|| bucket.to_string()),
diag_enabled.then(|| object.to_string()),
owner,
ObjectLockDiagMode::Write,
)))
}
async fn acquire_object_read_lock_if_needed(
&self,
op: &'static str,
bucket: &str,
object: &str,
opts: &mut ObjectOptions,
) -> Result<Option<rustfs_lock::NamespaceLockGuard>> {
) -> Result<Option<ObjectLockDiagGuard>> {
if opts.no_lock {
return Ok(None);
}
let diag_enabled = is_object_lock_diag_enabled();
let ns_lock = self.handle_new_ns_lock(bucket, object).await?;
let acquire_start = Instant::now();
let guard = ns_lock
.get_read_lock(get_lock_acquire_timeout())
.await
.map_err(|err| Self::map_namespace_lock_error(bucket, object, "read", err))?;
let owner = diag_enabled.then(|| ns_lock.owner().to_string());
log_object_lock_acquire_if_slow(
op,
bucket,
object,
owner.as_deref(),
ObjectLockDiagMode::Read,
acquire_start.elapsed(),
diag_enabled,
);
opts.no_lock = true;
Ok(Some(guard))
Ok(Some(ObjectLockDiagGuard::new(
guard,
diag_enabled,
op,
diag_enabled.then(|| bucket.to_string()),
diag_enabled.then(|| object.to_string()),
owner,
ObjectLockDiagMode::Read,
)))
}
fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option<rustfs_lock::NamespaceLockGuard>) -> GetObjectReader {
fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option<ObjectLockDiagGuard>) -> GetObjectReader {
if is_lock_optimization_enabled() {
return reader;
}
@@ -286,7 +449,9 @@ impl ECStore {
let object = encode_dir_object(object);
let mut opts = opts.clone();
let read_lock_guard = self.acquire_object_read_lock_if_needed(bucket, &object, &mut opts).await?;
let read_lock_guard = self
.acquire_object_read_lock_if_needed("get_object", bucket, &object, &mut opts)
.await?;
let reader = if self.single_pool() {
self.pools[0]
@@ -346,7 +511,9 @@ impl ECStore {
let object = encode_dir_object(object);
let mut opts = opts.clone();
let _object_lock_guard = self.acquire_object_read_lock_if_needed(bucket, &object, &mut opts).await?;
let _object_lock_guard = self
.acquire_object_read_lock_if_needed("get_object_info", bucket, &object, &mut opts)
.await?;
let info = if self.single_pool() {
self.pools[0].get_object_info(bucket, object.as_str(), &opts).await?
@@ -381,7 +548,7 @@ impl ECStore {
let mut dst_opts = dst_opts.clone();
let _dst_lock_guard = if cp_src_dst_same {
self.acquire_object_write_lock_if_needed(dst_bucket, &dst_object, &mut dst_opts)
self.acquire_object_write_lock_if_needed("copy_object", dst_bucket, &dst_object, &mut dst_opts)
.await?
} else {
None
@@ -464,7 +631,9 @@ impl ECStore {
return Ok(ObjectInfo::default());
}
let _object_lock_guard = self.acquire_object_write_lock_if_needed(bucket, object, &mut opts).await?;
let _object_lock_guard = self
.acquire_object_write_lock_if_needed("delete_object", bucket, object, &mut opts)
.await?;
if opts.delete_prefix {
self.delete_prefix(bucket, object, &opts).await?;
@@ -1116,6 +1285,15 @@ mod tests {
.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::<u8>::new())),
object_info: ObjectInfo::default(),
@@ -1145,6 +1323,15 @@ mod tests {
.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(),