#[derive(Debug, Clone, Default)]
pub struct LockMetricsSnapshot {
pub name: &'static str,
pub acquisitions: u64,
pub contentions: u64,
pub wait_ns: u64,
pub hold_ns: u64,
pub max_wait_ns: u64,
pub max_hold_ns: u64,
pub wait_percentile_sample_count: u64,
pub p95_wait_ns: u64,
pub p999_wait_ns: u64,
pub hold_percentile_sample_count: u64,
pub p95_hold_ns: u64,
pub p999_hold_ns: u64,
pub instrumentation_mode: &'static str,
}
#[cfg(feature = "lock-metrics")]
mod inner {
use super::LockMetricsSnapshot;
use crate::sync::lock_ordering::{self, LockModule, LockRank};
use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{LockResult, Mutex, MutexGuard, PoisonError};
use std::time::Instant;
#[derive(Debug, Default)]
#[repr(C, align(64))]
struct Metrics {
acquisitions: AtomicU64,
contentions: AtomicU64,
wait_ns: AtomicU64,
max_wait_ns: AtomicU64,
_pad: [u8; 32],
hold_ns: AtomicU64,
max_hold_ns: AtomicU64,
wait_samples: Mutex<VecDeque<u64>>,
hold_samples: Mutex<VecDeque<u64>>,
}
const MAX_SAMPLES: usize = 10_000;
impl Metrics {
fn update_max(current: &AtomicU64, value: u64) {
current.fetch_max(value, Ordering::Relaxed);
}
fn record_sample(samples: &mut VecDeque<u64>, sample: u64) {
if samples.len() == MAX_SAMPLES {
let evicted = samples.pop_front();
debug_assert!(evicted.is_some());
}
samples.push_back(sample);
}
fn record_acquire(&self, wait_ns: u64, contended: bool) {
let mut samples = self
.wait_samples
.lock()
.unwrap_or_else(PoisonError::into_inner);
self.acquisitions.fetch_add(1, Ordering::Relaxed);
self.wait_ns.fetch_add(wait_ns, Ordering::Relaxed);
Self::update_max(&self.max_wait_ns, wait_ns);
if contended {
self.contentions.fetch_add(1, Ordering::Relaxed);
}
Self::record_sample(&mut samples, wait_ns);
}
fn record_hold(&self, hold_ns: u64) {
let mut samples = self
.hold_samples
.lock()
.unwrap_or_else(PoisonError::into_inner);
self.hold_ns.fetch_add(hold_ns, Ordering::Relaxed);
Self::update_max(&self.max_hold_ns, hold_ns);
Self::record_sample(&mut samples, hold_ns);
}
fn percentile_from_sorted(sorted: &[u64], numerator: usize, denominator: usize) -> u64 {
if sorted.is_empty() {
return 0;
}
let last_index = sorted.len() - 1;
let rank = last_index
.saturating_mul(numerator)
.saturating_add(denominator / 2)
/ denominator;
sorted[rank.min(last_index)]
}
fn snapshot(&self, name: &'static str) -> LockMetricsSnapshot {
let acquisitions;
let contentions;
let wait_ns;
let max_wait_ns;
let mut wait_frozen: Vec<u64>;
{
let samples = self
.wait_samples
.lock()
.unwrap_or_else(PoisonError::into_inner);
acquisitions = self.acquisitions.load(Ordering::Relaxed);
contentions = self.contentions.load(Ordering::Relaxed);
wait_ns = self.wait_ns.load(Ordering::Relaxed);
max_wait_ns = self.max_wait_ns.load(Ordering::Relaxed);
wait_frozen = samples.iter().copied().collect();
}
let wait_percentile_sample_count = u64::try_from(wait_frozen.len()).unwrap_or(u64::MAX);
wait_frozen.sort_unstable();
let p95_wait_ns = Self::percentile_from_sorted(&wait_frozen, 95, 100);
let p999_wait_ns = Self::percentile_from_sorted(&wait_frozen, 999, 1000);
let hold_ns;
let max_hold_ns;
let mut hold_frozen: Vec<u64>;
{
let samples = self
.hold_samples
.lock()
.unwrap_or_else(PoisonError::into_inner);
hold_ns = self.hold_ns.load(Ordering::Relaxed);
max_hold_ns = self.max_hold_ns.load(Ordering::Relaxed);
hold_frozen = samples.iter().copied().collect();
}
let hold_percentile_sample_count = u64::try_from(hold_frozen.len()).unwrap_or(u64::MAX);
hold_frozen.sort_unstable();
let p95_hold_ns = Self::percentile_from_sorted(&hold_frozen, 95, 100);
let p999_hold_ns = Self::percentile_from_sorted(&hold_frozen, 999, 1000);
LockMetricsSnapshot {
name,
acquisitions,
contentions,
wait_ns,
hold_ns,
max_wait_ns,
max_hold_ns,
wait_percentile_sample_count,
p95_wait_ns,
p999_wait_ns,
hold_percentile_sample_count,
p95_hold_ns,
p999_hold_ns,
instrumentation_mode: "opt_in_lock_metrics",
}
}
fn reset(&self) {
{
let mut samples = self
.wait_samples
.lock()
.unwrap_or_else(PoisonError::into_inner);
self.acquisitions.store(0, Ordering::Relaxed);
self.contentions.store(0, Ordering::Relaxed);
self.wait_ns.store(0, Ordering::Relaxed);
self.max_wait_ns.store(0, Ordering::Relaxed);
samples.clear();
}
{
let mut samples = self
.hold_samples
.lock()
.unwrap_or_else(PoisonError::into_inner);
self.hold_ns.store(0, Ordering::Relaxed);
self.max_hold_ns.store(0, Ordering::Relaxed);
samples.clear();
}
}
}
#[cfg(test)]
mod tests {
use super::{MAX_SAMPLES, Metrics};
#[test]
fn percentile_horizon_reports_retained_suffix_and_all_history_counters() {
let metrics = Metrics::default();
for _ in 0..2_500 {
metrics.record_acquire(1_000, true);
metrics.record_hold(1_000);
}
for _ in 0..7_500 {
metrics.record_acquire(1, false);
metrics.record_hold(1);
}
let full = metrics.snapshot("percentile_horizon");
assert_eq!(full.wait_percentile_sample_count, 10_000);
assert_eq!(full.hold_percentile_sample_count, 10_000);
assert_eq!(full.p95_wait_ns, 1_000);
assert_eq!(full.p999_wait_ns, 1_000);
assert_eq!(full.p95_hold_ns, 1_000);
assert_eq!(full.p999_hold_ns, 1_000);
for _ in 0..2_500 {
metrics.record_acquire(1, false);
metrics.record_hold(1);
}
let evicted = metrics.snapshot("percentile_horizon");
assert_eq!(evicted.wait_percentile_sample_count, 10_000);
assert_eq!(evicted.hold_percentile_sample_count, 10_000);
assert_eq!(evicted.p95_wait_ns, 1);
assert_eq!(evicted.p999_wait_ns, 1);
assert_eq!(evicted.p95_hold_ns, 1);
assert_eq!(evicted.p999_hold_ns, 1);
assert_eq!(evicted.acquisitions, 12_500);
assert_eq!(evicted.contentions, 2_500);
assert_eq!(evicted.wait_ns, 2_510_000);
assert_eq!(evicted.max_wait_ns, 1_000);
assert_eq!(evicted.hold_ns, 2_510_000);
assert_eq!(evicted.max_hold_ns, 1_000);
}
#[test]
fn sample_rings_replace_fifo_one_at_a_time_across_wraps() {
let metrics = Metrics::default();
let max_samples = u64::try_from(MAX_SAMPLES).expect("sample bound fits u64");
for sample in 0..max_samples {
metrics.record_acquire(sample, false);
metrics.record_hold(sample);
}
let wait_capacity = metrics
.wait_samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.capacity();
let hold_capacity = metrics
.hold_samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.capacity();
metrics.record_acquire(max_samples, false);
metrics.record_hold(max_samples);
{
let samples = metrics
.wait_samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
assert_eq!(samples.len(), MAX_SAMPLES);
assert_eq!(samples.capacity(), wait_capacity);
assert_eq!(samples.front().copied(), Some(1));
assert_eq!(samples.back().copied(), Some(max_samples));
}
{
let samples = metrics
.hold_samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
assert_eq!(samples.len(), MAX_SAMPLES);
assert_eq!(samples.capacity(), hold_capacity);
assert_eq!(samples.front().copied(), Some(1));
assert_eq!(samples.back().copied(), Some(max_samples));
}
let total_samples = max_samples * 3 + 17;
for sample in (max_samples + 1)..total_samples {
metrics.record_acquire(sample, false);
metrics.record_hold(sample);
}
let expected_front = total_samples - max_samples;
let expected_back = total_samples - 1;
{
let samples = metrics
.wait_samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
assert_eq!(samples.len(), MAX_SAMPLES);
assert_eq!(samples.capacity(), wait_capacity);
assert_eq!(samples.front().copied(), Some(expected_front));
assert_eq!(samples.back().copied(), Some(expected_back));
}
{
let samples = metrics
.hold_samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
assert_eq!(samples.len(), MAX_SAMPLES);
assert_eq!(samples.capacity(), hold_capacity);
assert_eq!(samples.front().copied(), Some(expected_front));
assert_eq!(samples.back().copied(), Some(expected_back));
}
}
}
#[derive(Debug)]
pub struct ContendedMutex<T> {
inner: Mutex<T>,
metrics: Metrics,
name: &'static str,
rank: Option<LockRank>,
module: LockModule,
}
impl<T> ContendedMutex<T> {
pub fn new(name: &'static str, value: T) -> Self {
let policy = lock_ordering::enforce_lock_name_policy(name);
Self {
inner: Mutex::new(value),
metrics: Metrics::default(),
name,
rank: policy.rank(),
module: policy.module(),
}
}
pub fn lock(&self) -> LockResult<ContendedMutexGuard<'_, T>> {
if let Some(rank) = self.rank {
lock_ordering::check_acquire_with_module(self.name, rank, self.module);
}
let start = Instant::now();
let (result, contended) = match self.inner.try_lock() {
Ok(guard) => (Ok(guard), false),
Err(std::sync::TryLockError::Poisoned(poison)) => (Err(poison), false),
Err(std::sync::TryLockError::WouldBlock) => (self.inner.lock(), true),
};
let acquired_at = Instant::now();
let wait_ns =
u64::try_from(acquired_at.duration_since(start).as_nanos()).unwrap_or(u64::MAX);
self.metrics.record_acquire(wait_ns, contended);
if let Some(rank) = self.rank {
lock_ordering::record_acquire_with_module(self.name, rank, self.module);
}
match result {
Ok(guard) => Ok(ContendedMutexGuard {
guard: Some(guard),
acquired_at,
metrics: &self.metrics,
name: self.name,
rank: self.rank,
module: self.module,
}),
Err(poison) => Err(PoisonError::new(ContendedMutexGuard {
guard: Some(poison.into_inner()),
acquired_at,
metrics: &self.metrics,
name: self.name,
rank: self.rank,
module: self.module,
})),
}
}
pub fn try_lock(
&self,
) -> Result<ContendedMutexGuard<'_, T>, std::sync::TryLockError<ContendedMutexGuard<'_, T>>>
{
match self.inner.try_lock() {
Ok(guard) => {
if let Some(rank) = self.rank {
let (name, module) = (self.name, self.module);
if let Err(payload) =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
lock_ordering::check_acquire_with_module(name, rank, module);
}))
{
drop(guard);
std::panic::resume_unwind(payload);
}
lock_ordering::record_acquire_with_module(name, rank, module);
}
let acquired_at = Instant::now();
self.metrics.record_acquire(0, false);
Ok(ContendedMutexGuard {
guard: Some(guard),
acquired_at,
metrics: &self.metrics,
name: self.name,
rank: self.rank,
module: self.module,
})
}
Err(std::sync::TryLockError::WouldBlock) => {
Err(std::sync::TryLockError::WouldBlock)
}
Err(std::sync::TryLockError::Poisoned(poison)) => {
if let Some(rank) = self.rank {
lock_ordering::check_acquire_with_module(self.name, rank, self.module);
lock_ordering::record_acquire_with_module(self.name, rank, self.module);
}
let acquired_at = Instant::now();
self.metrics.record_acquire(0, false);
Err(std::sync::TryLockError::Poisoned(PoisonError::new(
ContendedMutexGuard {
guard: Some(poison.into_inner()),
acquired_at,
metrics: &self.metrics,
name: self.name,
rank: self.rank,
module: self.module,
},
)))
}
}
}
pub fn snapshot(&self) -> LockMetricsSnapshot {
self.metrics.snapshot(self.name)
}
pub fn reset_metrics(&self) {
self.metrics.reset();
}
pub fn name(&self) -> &'static str {
self.name
}
}
pub struct ContendedMutexGuard<'a, T> {
guard: Option<MutexGuard<'a, T>>,
acquired_at: Instant,
metrics: &'a Metrics,
name: &'static str,
rank: Option<LockRank>,
module: LockModule,
}
impl<T> std::ops::Deref for ContendedMutexGuard<'_, T> {
type Target = T;
fn deref(&self) -> &T {
self.guard.as_ref().expect("guard used after drop")
}
}
impl<T> std::ops::DerefMut for ContendedMutexGuard<'_, T> {
fn deref_mut(&mut self) -> &mut T {
self.guard.as_mut().expect("guard used after drop")
}
}
impl<T> Drop for ContendedMutexGuard<'_, T> {
fn drop(&mut self) {
let hold_ns = u64::try_from(self.acquired_at.elapsed().as_nanos()).unwrap_or(u64::MAX);
drop(self.guard.take());
if let Some(rank) = self.rank {
lock_ordering::record_release_with_module(self.name, rank, self.module);
}
self.metrics.record_hold(hold_ns);
}
}
impl<T: std::fmt::Debug> std::fmt::Debug for ContendedMutexGuard<'_, T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ContendedMutexGuard")
.field("data", &self.guard)
.finish()
}
}
}
#[cfg(not(feature = "lock-metrics"))]
mod inner {
use super::LockMetricsSnapshot;
use crate::sync::lock_ordering::{self, LockRank};
use std::sync::{LockResult, Mutex, MutexGuard, PoisonError};
#[derive(Debug)]
pub struct ContendedMutex<T> {
inner: Mutex<T>,
name: &'static str,
rank: Option<LockRank>,
}
impl<T> ContendedMutex<T> {
#[inline]
pub fn new(name: &'static str, value: T) -> Self {
let rank = lock_ordering::rank_for_lock_name(name);
Self {
inner: Mutex::new(value),
name,
rank,
}
}
#[inline]
pub fn lock(&self) -> LockResult<ContendedMutexGuard<'_, T>> {
if let Some(rank) = self.rank {
lock_ordering::check_acquire(self.name, rank);
}
match self.inner.lock() {
Ok(guard) => {
if let Some(rank) = self.rank {
lock_ordering::record_acquire(self.name, rank);
}
Ok(ContendedMutexGuard {
guard,
name: self.name,
rank: self.rank,
})
}
Err(poison) => {
if let Some(rank) = self.rank {
lock_ordering::record_acquire(self.name, rank);
}
Err(PoisonError::new(ContendedMutexGuard {
guard: poison.into_inner(),
name: self.name,
rank: self.rank,
}))
}
}
}
pub fn try_lock(
&self,
) -> Result<ContendedMutexGuard<'_, T>, std::sync::TryLockError<ContendedMutexGuard<'_, T>>>
{
match self.inner.try_lock() {
Ok(guard) => {
if let Some(rank) = self.rank {
let name = self.name;
if let Err(payload) =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
lock_ordering::check_acquire(name, rank);
}))
{
drop(guard);
std::panic::resume_unwind(payload);
}
lock_ordering::record_acquire(name, rank);
}
Ok(ContendedMutexGuard {
guard,
name: self.name,
rank: self.rank,
})
}
Err(std::sync::TryLockError::WouldBlock) => {
Err(std::sync::TryLockError::WouldBlock)
}
Err(std::sync::TryLockError::Poisoned(poison)) => {
if let Some(rank) = self.rank {
lock_ordering::check_acquire(self.name, rank);
lock_ordering::record_acquire(self.name, rank);
}
Err(std::sync::TryLockError::Poisoned(PoisonError::new(
ContendedMutexGuard {
guard: poison.into_inner(),
name: self.name,
rank: self.rank,
},
)))
}
}
}
pub fn snapshot(&self) -> LockMetricsSnapshot {
LockMetricsSnapshot {
name: self.name,
instrumentation_mode: "disabled",
..Default::default()
}
}
pub fn reset_metrics(&self) {}
pub fn name(&self) -> &'static str {
self.name
}
}
pub struct ContendedMutexGuard<'a, T> {
guard: MutexGuard<'a, T>,
name: &'static str,
rank: Option<LockRank>,
}
impl<T> std::ops::Deref for ContendedMutexGuard<'_, T> {
type Target = T;
#[inline]
fn deref(&self) -> &T {
&self.guard
}
}
impl<T> std::ops::DerefMut for ContendedMutexGuard<'_, T> {
#[inline]
fn deref_mut(&mut self) -> &mut T {
&mut self.guard
}
}
impl<T> Drop for ContendedMutexGuard<'_, T> {
fn drop(&mut self) {
if let Some(rank) = self.rank {
lock_ordering::record_release(self.name, rank);
}
}
}
impl<T: std::fmt::Debug> std::fmt::Debug for ContendedMutexGuard<'_, T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ContendedMutexGuard")
.field("data", &*self.guard)
.finish()
}
}
}
pub use inner::{ContendedMutex, ContendedMutexGuard};
#[cfg(test)]
#[allow(clippy::significant_drop_tightening)]
mod tests {
use super::*;
#[cfg(feature = "lock-metrics")]
use crate::sync::lock_ordering;
use std::sync::Arc;
#[cfg(feature = "lock-metrics")]
use std::thread;
fn init_test(name: &str) {
crate::test_utils::init_test_logging();
crate::test_phase!(name);
}
#[test]
fn basic_lock_unlock() {
init_test("basic_lock_unlock");
let m = ContendedMutex::new("unknown", 42);
{
let guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
crate::assert_with_log!(*guard == 42, "value", 42, *guard);
drop(guard);
}
crate::test_complete!("basic_lock_unlock");
}
#[test]
fn mutate_through_guard() {
init_test("mutate_through_guard");
let m = ContendedMutex::new("unknown", 0);
{
let mut guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
*guard = 99;
}
let guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
crate::assert_with_log!(*guard == 99, "mutated value", 99, *guard);
drop(guard);
crate::test_complete!("mutate_through_guard");
}
#[test]
fn try_lock_succeeds_when_free() {
init_test("try_lock_succeeds_when_free");
let m = ContendedMutex::new("unknown", 42);
let guard = m.try_lock().expect("should succeed");
crate::assert_with_log!(*guard == 42, "try_lock value", 42, *guard);
drop(guard);
crate::test_complete!("try_lock_succeeds_when_free");
}
#[test]
fn try_lock_fails_when_held() {
init_test("try_lock_fails_when_held");
let m = ContendedMutex::new("unknown", 42);
let _guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let is_err = m.try_lock().is_err();
crate::assert_with_log!(is_err, "try_lock fails", true, is_err);
crate::test_complete!("try_lock_fails_when_held");
}
#[test]
fn snapshot_returns_name() {
init_test("snapshot_returns_name");
let m = ContendedMutex::new("unknown", 0);
let snap = m.snapshot();
crate::assert_with_log!(snap.name == "unknown", "name", "unknown", snap.name);
crate::test_complete!("snapshot_returns_name");
}
#[test]
fn name_accessor() {
init_test("name_accessor");
let m = ContendedMutex::new("tasks", 0);
crate::assert_with_log!(m.name() == "tasks", "name", "tasks", m.name());
crate::test_complete!("name_accessor");
}
#[test]
fn reset_metrics_no_panic() {
init_test("reset_metrics_no_panic");
let m = ContendedMutex::new("unknown", 0);
{
let _g = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
}
m.reset_metrics();
let snap = m.snapshot();
crate::assert_with_log!(
snap.acquisitions == 0,
"acquisitions after reset",
0u64,
snap.acquisitions
);
crate::test_complete!("reset_metrics_no_panic");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn metrics_track_acquisitions() {
init_test("metrics_track_acquisitions");
let m = ContendedMutex::new("unknown", 0);
for _ in 0..10 {
let _g = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
}
let snap = m.snapshot();
crate::assert_with_log!(
snap.acquisitions == 10,
"acquisitions",
10u64,
snap.acquisitions
);
crate::test_complete!("metrics_track_acquisitions");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn metrics_track_hold_time() {
init_test("metrics_track_hold_time");
let m = ContendedMutex::new("unknown", 0);
{
let _g = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
std::thread::sleep(std::time::Duration::from_millis(5));
}
let snap = m.snapshot();
crate::assert_with_log!(
snap.hold_ns >= 4_000_000,
"hold_ns >= 4ms",
true,
snap.hold_ns >= 4_000_000
);
crate::assert_with_log!(
snap.max_hold_ns >= 4_000_000,
"max_hold_ns >= 4ms",
true,
snap.max_hold_ns >= 4_000_000
);
crate::test_complete!("metrics_track_hold_time");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn metrics_track_contention() {
init_test("metrics_track_contention");
let m = Arc::new(ContendedMutex::new("unknown", 0));
let guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let m2 = Arc::clone(&m);
let handle = thread::spawn(move || {
let _g = m2.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
});
thread::sleep(std::time::Duration::from_millis(10));
drop(guard);
handle.join().expect("thread panicked");
let snap = m.snapshot();
crate::assert_with_log!(
snap.contentions >= 1,
"contentions >= 1",
true,
snap.contentions >= 1
);
crate::assert_with_log!(snap.wait_ns > 0, "wait_ns > 0", true, snap.wait_ns > 0);
crate::test_complete!("metrics_track_contention");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn reset_clears_all_metrics() {
init_test("reset_clears_all_metrics");
let m = ContendedMutex::new("unknown", 0);
{
let _g = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
}
let before = m.snapshot();
crate::assert_with_log!(
before.acquisitions == 1,
"before reset",
1u64,
before.acquisitions
);
m.reset_metrics();
let after = m.snapshot();
crate::assert_with_log!(
after.acquisitions == 0,
"after reset acquisitions",
0u64,
after.acquisitions
);
crate::assert_with_log!(
after.hold_ns == 0,
"after reset hold_ns",
0u64,
after.hold_ns
);
crate::test_complete!("reset_clears_all_metrics");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn poisoned_lock_does_not_count_as_contention() {
init_test("poisoned_lock_does_not_count_as_contention");
let m = Arc::new(ContendedMutex::new("unknown", 0u8));
let m2 = Arc::clone(&m);
let poisoner = thread::spawn(move || {
let _guard = m2.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
panic!("intentional poison");
});
let _ = poisoner.join();
let poison_err = m.lock().expect_err("lock should be poisoned");
drop(poison_err.into_inner());
let snap = m.snapshot();
crate::assert_with_log!(
snap.contentions == 0,
"poison is not contention",
0u64,
snap.contentions
);
crate::test_complete!("poisoned_lock_does_not_count_as_contention");
}
#[cfg(any(debug_assertions, feature = "lock-metrics"))]
#[test]
fn try_lock_order_violation_does_not_poison_lower_mutex() {
use crate::sync::lock_ordering;
init_test("try_lock_order_violation_does_not_poison_lower_mutex");
lock_ordering::clear_held_locks();
let high = ContendedMutex::new("tasks", 100u32); let low = ContendedMutex::new("regions_table", 7u32); let high_guard = high.lock().expect("acquire higher rank");
let inverted = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let _ = low.try_lock();
}));
assert!(
inverted.is_err(),
"inverted try_lock must raise the lock-order diagnostic"
);
drop(high_guard);
lock_ordering::clear_held_locks();
match low.try_lock() {
Ok(guard) => assert_eq!(*guard, 7u32, "lower mutex data unchanged"),
Err(std::sync::TryLockError::Poisoned(_)) => {
panic!("lower mutex was poisoned by the lock-order panic")
}
Err(std::sync::TryLockError::WouldBlock) => panic!("unexpected WouldBlock"),
}
lock_ordering::clear_held_locks();
crate::test_complete!("try_lock_order_violation_does_not_poison_lower_mutex");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn poisoned_ranked_lock_release_clears_lock_order_state() {
init_test("poisoned_ranked_lock_release_clears_lock_order_state");
lock_ordering::clear_held_locks();
let m = Arc::new(ContendedMutex::new("tasks", 0u8));
let m2 = Arc::clone(&m);
let poisoner = thread::spawn(move || {
let _guard = m2.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
panic!("intentional poison");
});
let _ = poisoner.join();
let poison_err = m.lock().expect_err("lock should be poisoned");
let guard = poison_err.into_inner();
let held_locks = lock_ordering::current_held_locks();
crate::assert_with_log!(
held_locks
.get(&lock_ordering::LockRank::Tasks)
.map_or(0, Vec::len)
== 1,
"poisoned guard acquire is tracked",
1usize,
held_locks
.get(&lock_ordering::LockRank::Tasks)
.map_or(0, Vec::len)
);
drop(guard);
crate::assert_with_log!(
lock_ordering::current_held_locks().is_empty(),
"poisoned guard drop clears held locks",
true,
lock_ordering::current_held_locks().is_empty()
);
crate::assert_with_log!(
lock_ordering::current_held_ranks().is_empty(),
"poisoned guard drop clears held ranks",
true,
lock_ordering::current_held_ranks().is_empty()
);
crate::test_complete!("poisoned_ranked_lock_release_clears_lock_order_state");
}
#[test]
fn lock_metrics_snapshot_debug_clone_default() {
let snap = LockMetricsSnapshot::default();
let dbg = format!("{snap:?}");
assert!(dbg.contains("LockMetricsSnapshot"));
assert_eq!(snap.acquisitions, 0);
assert_eq!(snap.contentions, 0);
assert_eq!(snap.wait_ns, 0);
assert_eq!(snap.hold_ns, 0);
assert_eq!(snap.max_wait_ns, 0);
assert_eq!(snap.max_hold_ns, 0);
assert_eq!(snap.wait_percentile_sample_count, 0);
assert_eq!(snap.p95_wait_ns, 0);
assert_eq!(snap.p999_wait_ns, 0);
assert_eq!(snap.hold_percentile_sample_count, 0);
assert_eq!(snap.p95_hold_ns, 0);
assert_eq!(snap.p999_hold_ns, 0);
let cloned = snap.clone();
assert_eq!(cloned.name, snap.name);
}
#[cfg(feature = "lock-metrics")]
#[test]
fn metrics_snapshot_reports_tail_latencies() {
init_test("metrics_snapshot_reports_tail_latencies");
let m = ContendedMutex::new("tasks", 0u32);
for _ in 0..4 {
let mut guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
*guard += 1;
}
let snap = m.snapshot();
crate::assert_with_log!(
snap.instrumentation_mode == "opt_in_lock_metrics",
"instrumentation mode",
"opt_in_lock_metrics",
snap.instrumentation_mode
);
crate::assert_with_log!(
snap.acquisitions == 4,
"acquisitions",
4u64,
snap.acquisitions
);
crate::assert_with_log!(
snap.p95_wait_ns <= snap.max_wait_ns,
"p95 wait <= max wait",
true,
snap.p95_wait_ns <= snap.max_wait_ns
);
crate::assert_with_log!(
snap.p999_wait_ns <= snap.max_wait_ns,
"p999 wait <= max wait",
true,
snap.p999_wait_ns <= snap.max_wait_ns
);
crate::assert_with_log!(
snap.p95_hold_ns <= snap.max_hold_ns,
"p95 hold <= max hold",
true,
snap.p95_hold_ns <= snap.max_hold_ns
);
crate::assert_with_log!(
snap.p999_hold_ns <= snap.max_hold_ns,
"p999 hold <= max hold",
true,
snap.p999_hold_ns <= snap.max_hold_ns
);
crate::test_complete!("metrics_snapshot_reports_tail_latencies");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn metrics_snapshot_and_reset_coherent_under_concurrency() {
use std::sync::atomic::{AtomicBool, Ordering as AtOrd};
init_test("metrics_snapshot_and_reset_coherent_under_concurrency");
let m = Arc::new(ContendedMutex::new("tasks", 0u64));
let stop = Arc::new(AtomicBool::new(false));
let mut recorders = Vec::new();
for _ in 0..4 {
let m = Arc::clone(&m);
let stop = Arc::clone(&stop);
recorders.push(thread::spawn(move || {
while !stop.load(AtOrd::Relaxed) {
let mut guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
*guard = guard.wrapping_add(1);
drop(guard);
}
}));
}
let resetter = {
let m = Arc::clone(&m);
let stop = Arc::clone(&stop);
thread::spawn(move || {
let mut i = 0u32;
while !stop.load(AtOrd::Relaxed) {
i = i.wrapping_add(1);
if i.is_multiple_of(64) {
m.reset_metrics();
}
std::hint::spin_loop();
}
})
};
for _ in 0..2000 {
let snap = m.snapshot();
crate::assert_with_log!(
snap.p95_wait_ns <= snap.p999_wait_ns,
"p95_wait <= p999_wait under concurrency",
true,
snap.p95_wait_ns <= snap.p999_wait_ns
);
crate::assert_with_log!(
snap.p999_wait_ns <= snap.max_wait_ns,
"p999_wait <= max_wait under concurrency",
true,
snap.p999_wait_ns <= snap.max_wait_ns
);
crate::assert_with_log!(
snap.p95_hold_ns <= snap.p999_hold_ns,
"p95_hold <= p999_hold under concurrency",
true,
snap.p95_hold_ns <= snap.p999_hold_ns
);
crate::assert_with_log!(
snap.p999_hold_ns <= snap.max_hold_ns,
"p999_hold <= max_hold under concurrency",
true,
snap.p999_hold_ns <= snap.max_hold_ns
);
}
stop.store(true, AtOrd::Relaxed);
for h in recorders {
let _ = h.join();
}
let _ = resetter.join();
m.reset_metrics();
let snap = m.snapshot();
crate::assert_with_log!(
snap.acquisitions == 0,
"reset zeroes acquisitions",
0u64,
snap.acquisitions
);
crate::assert_with_log!(
snap.max_wait_ns == 0,
"reset zeroes max_wait_ns",
0u64,
snap.max_wait_ns
);
crate::assert_with_log!(
snap.max_hold_ns == 0,
"reset zeroes max_hold_ns",
0u64,
snap.max_hold_ns
);
crate::assert_with_log!(
snap.p999_wait_ns == 0,
"reset clears wait samples",
0u64,
snap.p999_wait_ns
);
crate::assert_with_log!(
snap.p999_hold_ns == 0,
"reset clears hold samples",
0u64,
snap.p999_hold_ns
);
crate::test_complete!("metrics_snapshot_and_reset_coherent_under_concurrency");
}
#[test]
fn contended_mutex_debug() {
let m = ContendedMutex::new("unknown", 42_i32);
let dbg = format!("{m:?}");
assert!(dbg.contains("ContendedMutex"));
}
#[test]
fn contended_mutex_guard_debug() {
let m = ContendedMutex::new("unknown", 42_i32);
let guard = m.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let dbg = format!("{guard:?}");
assert!(dbg.contains("ContendedMutexGuard"));
drop(guard);
}
#[test]
fn try_lock_returns_poisoned_after_panic() {
init_test("try_lock_returns_poisoned_after_panic");
let m = Arc::new(ContendedMutex::new("unknown", 7u32));
let m2 = Arc::clone(&m);
let poisoner = std::thread::spawn(move || {
let _guard = m2.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
panic!("deliberate poison");
});
let _ = poisoner.join();
let result = m.try_lock();
let is_poisoned = matches!(result, Err(std::sync::TryLockError::Poisoned(_)));
crate::assert_with_log!(is_poisoned, "try_lock returns Poisoned", true, is_poisoned);
if let Err(std::sync::TryLockError::Poisoned(pe)) = m.try_lock() {
let guard = pe.into_inner();
crate::assert_with_log!(*guard == 7, "data preserved", 7u32, *guard);
}
crate::test_complete!("try_lock_returns_poisoned_after_panic");
}
#[cfg(feature = "lock-metrics")]
#[test]
fn hold_time_recorded_on_panic_in_critical_section() {
init_test("hold_time_recorded_on_panic_in_critical_section");
let m = Arc::new(ContendedMutex::new("unknown", 0u32));
let m2 = Arc::clone(&m);
let handle = std::thread::spawn(move || {
let _guard = m2.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
std::thread::sleep(std::time::Duration::from_millis(5));
panic!("panic while holding guard");
});
let _ = handle.join();
let snap = m.snapshot();
crate::assert_with_log!(
snap.hold_ns >= 4_000_000,
"hold_ns recorded despite panic",
true,
snap.hold_ns >= 4_000_000
);
crate::test_complete!("hold_time_recorded_on_panic_in_critical_section");
}
}