From 198a07d3faeeefd05b24722ef7fac178364d8c06 Mon Sep 17 00:00:00 2001 From: houseme Date: Thu, 20 Aug 2026 21:47:30 +0800 Subject: [PATCH] feat(ecstore): observe put commit lock wait (#6310) --- Cargo.lock | 1 + crates/ecstore/Cargo.toml | 1 + crates/ecstore/src/set_disk/mod.rs | 98 ++++++++++++++++++++++++++++++ crates/io-metrics/src/lib.rs | 13 ++-- 4 files changed, 108 insertions(+), 5 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index a7657e2d0..503b2d3d8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9432,6 +9432,7 @@ dependencies = [ "md-5 0.11.0", "memmap2", "metrics", + "metrics-util", "moka", "num_cpus", "opentelemetry", diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index 46fec8cfc..517adb44c 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -267,6 +267,7 @@ tracing-subscriber = { workspace = true, features = ["json", "env-filter", "time # dispatcher to keep tracing's process-global callsite-interest cache honest. tracing-core = { workspace = true } serial_test = { workspace = true } +metrics-util = { workspace = true, features = ["debugging"] } opentelemetry_sdk = { workspace = true, features = ["rt-tokio"] } proptest = "1" rcgen.workspace = true diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index ca05840db..5e46b511b 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -3264,6 +3264,17 @@ impl SetDisks { self.get_object_metadata_cache.invalidate_all(); } + #[inline(always)] + fn record_put_object_commit_namespace_lock_wait(op: &'static str, acquire_start: Instant) { + if op != "put_object_commit" || !rustfs_io_metrics::put_stage_metrics_enabled() { + return; + } + rustfs_io_metrics::record_put_object_stage_duration_from( + rustfs_io_metrics::PUT_STAGE_PUT_OBJECT_COMMIT_NAMESPACE_LOCK_WAIT, + Some(acquire_start), + ); + } + async fn acquire_read_lock_diag(&self, op: &'static str, bucket: &str, object: &str) -> Result { crate::hp_guard!("SetDisks::acquire_read_lock"); let diag_enabled = is_object_lock_diag_enabled(); @@ -3295,6 +3306,7 @@ impl SetDisks { .get_write_lock(get_lock_acquire_timeout()) .await .map_err(|e| self.map_namespace_lock_error(bucket, object, "write", e))?; + Self::record_put_object_commit_namespace_lock_wait(op, acquire_start); let owner = diag_enabled.then(|| ns_lock.owner().to_string()); self.log_object_lock_acquire_if_slow( op, @@ -3342,6 +3354,7 @@ impl SetDisks { }) .await .map_err(|e| self.map_namespace_lock_error(bucket, object, "write", e))?; + Self::record_put_object_commit_namespace_lock_wait(op, acquire_start); let owner = diag_enabled.then(|| ns_lock.owner().to_string()); self.log_object_lock_acquire_if_slow( op, @@ -5460,6 +5473,7 @@ mod tests { }; use crate::store::init_format::save_format_file; use crate::store::list_objects::ListPathOptions; + use metrics_util::debugging::{DebugValue, DebuggingRecorder}; use rustfs_filemeta::ErasureInfo; use rustfs_filemeta::FileMeta; use rustfs_filemeta::MetaCacheEntry; @@ -5695,6 +5709,90 @@ mod tests { assert_eq!(Arc::strong_count(&set.set_lock_namespace), before); } + fn put_object_commit_namespace_lock_wait_sample_count(snapshotter: &metrics_util::debugging::Snapshotter) -> usize { + snapshotter + .snapshot() + .into_vec() + .into_iter() + .filter(|(composite, _, _, _)| { + composite.key().name() == "rustfs_s3_put_object_stage_duration_ms" + && composite.key().labels().any(|label| { + label.key().to_string() == "stage" + && label.value().to_string() == rustfs_io_metrics::PUT_STAGE_PUT_OBJECT_COMMIT_NAMESPACE_LOCK_WAIT + }) + }) + .map(|(_, _, _, value)| match value { + DebugValue::Histogram(samples) => samples.len(), + _ => 0, + }) + .sum() + } + + #[test] + #[serial] + fn put_object_commit_namespace_lock_wait_metric_is_wired_to_both_write_lock_paths() { + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("test runtime should start"); + + metrics::with_local_recorder(&recorder, || { + runtime.block_on(async { + let ctx = Arc::new(InstanceContext::new()); + ctx.update_erasure_type(SetupType::Erasure).await; + let set = make_test_set_disks_with_ctx(Vec::new(), ctx).await; + let bucket = "bucket"; + let object = "object"; + + rustfs_io_metrics::set_put_stage_metrics_enabled(false); + let guard = set + .acquire_write_lock_diag("put_object_commit", bucket, object) + .await + .expect("disabled metrics acquire should succeed"); + drop(guard); + assert_eq!(put_object_commit_namespace_lock_wait_sample_count(&snapshotter), 0); + + rustfs_io_metrics::set_put_stage_metrics_enabled(true); + let guard = set + .acquire_write_lock_diag("put_object_commit", bucket, object) + .await + .expect("normal PUT commit acquire should succeed"); + drop(guard); + assert_eq!(put_object_commit_namespace_lock_wait_sample_count(&snapshotter), 1); + + let guard = set + .acquire_write_lock_diag("complete_multipart_upload_commit", bucket, object) + .await + .expect("non-PUT commit acquire should succeed"); + drop(guard); + assert_eq!(put_object_commit_namespace_lock_wait_sample_count(&snapshotter), 0); + + let held_guard = set + .acquire_write_lock_diag("put_object_commit", bucket, object) + .await + .expect("holder acquire should succeed"); + assert_eq!(put_object_commit_namespace_lock_wait_sample_count(&snapshotter), 1); + + let (pending_tx, pending_rx) = tokio::sync::oneshot::channel(); + let pending_acquire = + set.acquire_write_lock_diag_with_pending_hook("put_object_commit", bucket, object, move || { + let _ = pending_tx.send(()); + }); + let release_holder = async { + pending_rx.await.expect("pending hook should fire"); + drop(held_guard); + }; + let (pending_guard, ()) = tokio::join!(pending_acquire, release_holder); + drop(pending_guard.expect("pending-hook PUT commit acquire should succeed")); + assert_eq!(put_object_commit_namespace_lock_wait_sample_count(&snapshotter), 1); + + rustfs_io_metrics::set_put_stage_metrics_enabled(false); + }); + }); + } + #[tokio::test] async fn new_ns_lock_shares_clients_without_changing_quorum() { let healthy: Arc = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new()))); diff --git a/crates/io-metrics/src/lib.rs b/crates/io-metrics/src/lib.rs index 3aa15c382..a52f5adb6 100644 --- a/crates/io-metrics/src/lib.rs +++ b/crates/io-metrics/src/lib.rs @@ -109,6 +109,7 @@ pub fn put_stage_timer() -> Option { put_stage_metrics_enabled().then(std::time::Instant::now) } +pub const PUT_STAGE_PUT_OBJECT_COMMIT_NAMESPACE_LOCK_WAIT: &str = "put_object_commit_namespace_lock_wait"; pub const PUT_STAGE_SET_DISK_RENAME_QUORUM_WAIT: &str = "set_disk_rename_quorum_wait"; pub const PUT_STAGE_SET_DISK_RENAME_DISK_WAIT: &str = "set_disk_rename_disk_wait"; pub const PUT_STAGE_SET_DISK_RENAME_FILE_SYNC_PERMIT_WAIT: &str = "set_disk_rename_file_sync_permit_wait"; @@ -3147,7 +3148,9 @@ mod tests { #[test] fn put_stage_sync_tail_labels_are_static_and_gated() { let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner()); + assert_eq!(PUT_STAGE_PUT_OBJECT_COMMIT_NAMESPACE_LOCK_WAIT, "put_object_commit_namespace_lock_wait"); let stages = [ + PUT_STAGE_PUT_OBJECT_COMMIT_NAMESPACE_LOCK_WAIT, PUT_STAGE_SET_DISK_RENAME_QUORUM_WAIT, PUT_STAGE_SET_DISK_RENAME_DISK_WAIT, PUT_STAGE_SET_DISK_RENAME_FILE_SYNC_PERMIT_WAIT, @@ -3161,11 +3164,11 @@ mod tests { ]; let unique = stages.iter().copied().collect::>(); assert_eq!(unique.len(), stages.len()); - assert!( - stages - .iter() - .all(|stage| stage.starts_with("set_disk_rename_") && !stage.contains('/') && !stage.contains('{')) - ); + assert!(stages.iter().all(|stage| { + (stage.starts_with("set_disk_rename_") || *stage == PUT_STAGE_PUT_OBJECT_COMMIT_NAMESPACE_LOCK_WAIT) + && !stage.contains('/') + && !stage.contains('{') + })); let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter();