use crossbeam_channel::{bounded, Receiver as CbReceiver, RecvTimeoutError, Sender as CbSender};
use hdrhistogram::Histogram;
use std::collections::HashMap;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use crate::lib_on::{meta_rw_lock, MetaRwLock};
use crate::batch::{EventProducer, EventQueueRegistry};
use crate::instant::Instant;
use crate::lib_on::hotpath_guard::DRAIN_INTERVAL_MS;
use crate::lib_on::START_TIME;
use crate::metrics_server::METRICS_SERVER_PORT;
pub(crate) mod wrapper;
pub use wrapper::std::{RwLock, RwLockReadGuard, RwLockWriteGuard};
static RW_LOCK_ID_COUNTER: AtomicU32 = AtomicU32::new(1);
fn next_rw_lock_id() -> u32 {
RW_LOCK_ID_COUNTER.fetch_add(1, Ordering::Relaxed)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RwLockKind {
Read,
Write,
}
#[derive(Debug)]
pub(crate) enum RwLockEvent {
Created {
id: u32,
source: &'static str,
label: Option<String>,
type_name: &'static str,
},
Released {
id: u32,
kind: RwLockKind,
wait_nanos: Option<u64>,
acquire_nanos: Option<u64>,
},
}
#[derive(Debug, Clone)]
pub(crate) struct RwLockEntry {
pub(crate) id: u32,
pub(crate) source: &'static str,
pub(crate) label: Option<String>,
pub(crate) type_name: &'static str,
pub(crate) read_count: u64,
pub(crate) write_count: u64,
pub(crate) read_sampled_count: u64,
pub(crate) write_sampled_count: u64,
pub(crate) read_wait_total_nanos: u64,
pub(crate) write_wait_total_nanos: u64,
pub(crate) read_acquire_total_nanos: u64,
pub(crate) write_acquire_total_nanos: u64,
read_wait_hist: Option<Histogram<u64>>,
write_wait_hist: Option<Histogram<u64>>,
read_acquire_hist: Option<Histogram<u64>>,
write_acquire_hist: Option<Histogram<u64>>,
pub(crate) iter: u32,
}
impl RwLockEntry {
const LOW_NS: u64 = 1;
const HIGH_NS: u64 = 1_000_000_000_000; const SIGFIGS: u8 = 3;
fn new_histogram() -> Histogram<u64> {
Histogram::<u64>::new_with_bounds(Self::LOW_NS, Self::HIGH_NS, Self::SIGFIGS)
.expect("hdrhistogram init")
}
#[inline]
fn record(hist: &mut Option<Histogram<u64>>, nanos: u64) {
if let Some(ref mut hist) = hist {
hist.record(nanos.clamp(Self::LOW_NS, Self::HIGH_NS))
.unwrap();
}
}
pub(crate) fn count(&self, kind: RwLockKind) -> u64 {
match kind {
RwLockKind::Read => self.read_count,
RwLockKind::Write => self.write_count,
}
}
pub(crate) fn sampled_count(&self, kind: RwLockKind) -> u64 {
match kind {
RwLockKind::Read => self.read_sampled_count,
RwLockKind::Write => self.write_sampled_count,
}
}
pub(crate) fn wait_avg_nanos(&self, kind: RwLockKind) -> u64 {
let total = match kind {
RwLockKind::Read => self.read_wait_total_nanos,
RwLockKind::Write => self.write_wait_total_nanos,
};
total.checked_div(self.sampled_count(kind)).unwrap_or(0)
}
pub(crate) fn acquire_avg_nanos(&self, kind: RwLockKind) -> u64 {
let total = match kind {
RwLockKind::Read => self.read_acquire_total_nanos,
RwLockKind::Write => self.write_acquire_total_nanos,
};
total.checked_div(self.sampled_count(kind)).unwrap_or(0)
}
fn percentile(hist: &Option<Histogram<u64>>, count: u64, p: f64) -> u64 {
match hist {
Some(hist) if count > 0 => hist.value_at_percentile(p.clamp(0.0, 100.0)),
_ => 0,
}
}
pub(crate) fn wait_percentile_nanos(&self, kind: RwLockKind, p: f64) -> u64 {
let hist = match kind {
RwLockKind::Read => &self.read_wait_hist,
RwLockKind::Write => &self.write_wait_hist,
};
Self::percentile(hist, self.sampled_count(kind), p)
}
pub(crate) fn acquire_percentile_nanos(&self, kind: RwLockKind, p: f64) -> u64 {
let hist = match kind {
RwLockKind::Read => &self.read_acquire_hist,
RwLockKind::Write => &self.write_acquire_hist,
};
Self::percentile(hist, self.sampled_count(kind), p)
}
}
pub(crate) struct RwLocksInternalState {
pub(crate) stats: HashMap<u32, RwLockEntry>,
}
pub(crate) struct RwLocksState {
pub(crate) inner: Arc<MetaRwLock<RwLocksInternalState>>,
pub(crate) shutdown_tx: Mutex<Option<CbSender<()>>>,
pub(crate) completion_rx: Mutex<Option<CbReceiver<()>>>,
}
pub(crate) static RW_LOCKS_STATE: OnceLock<RwLocksState> = OnceLock::new();
pub(crate) fn get_sorted_rw_lock_entries() -> Vec<RwLockEntry> {
let Some(state) = RW_LOCKS_STATE.get() else {
return Vec::new();
};
let guard = state.inner.read().unwrap();
let mut stats: Vec<RwLockEntry> = guard.stats.values().cloned().collect();
stats.sort_by(compare_rw_lock_entries);
stats
}
pub(crate) fn get_rw_locks_json() -> crate::json::JsonRwLocksList {
let entries = get_sorted_rw_lock_entries();
let elapsed = std::time::Duration::from_nanos(crate::lib_on::current_elapsed_ns());
crate::lib_on::report::collect_rw_locks_json(
&entries,
elapsed,
&crate::lib_on::hotpath_guard::configured_percentiles(),
)
}
#[inline]
pub(crate) fn elapsed_nanos(start: Instant) -> u64 {
start.elapsed().as_nanos() as u64
}
#[inline]
pub(crate) fn wait_stamp() -> Option<Instant> {
crate::lib_on::sampling::rw_locks_should_time().then(Instant::now)
}
#[inline]
pub(crate) fn cancel_wait_stamp() {
crate::lib_on::sampling::rw_locks_untime();
}
static EVENT_QUEUES: EventQueueRegistry<RwLockEvent> = EventQueueRegistry::new();
thread_local! {
static EVENT_PRODUCER: EventProducer<RwLockEvent> = EVENT_QUEUES.register();
}
#[inline]
pub(crate) fn send_rw_lock_event(event: RwLockEvent) {
if !EVENT_QUEUES.is_active() {
return;
}
let _suspend = crate::lib_on::SuspendAllocTracking::new();
let _ = EVENT_PRODUCER.try_with(|producer| producer.push(event));
}
pub(crate) fn stop_rw_lock_events() {
EVENT_QUEUES.set_active(false);
}
fn placeholder_rw_lock_entry(id: u32) -> RwLockEntry {
RwLockEntry {
id,
source: "",
label: None,
type_name: "",
read_count: 0,
write_count: 0,
read_sampled_count: 0,
write_sampled_count: 0,
read_wait_total_nanos: 0,
write_wait_total_nanos: 0,
read_acquire_total_nanos: 0,
write_acquire_total_nanos: 0,
read_wait_hist: Some(RwLockEntry::new_histogram()),
write_wait_hist: Some(RwLockEntry::new_histogram()),
read_acquire_hist: Some(RwLockEntry::new_histogram()),
write_acquire_hist: Some(RwLockEntry::new_histogram()),
iter: 0,
}
}
fn process_rw_lock_event(state: &mut RwLocksInternalState, event: RwLockEvent) {
match event {
RwLockEvent::Created {
id,
source,
label,
type_name,
} => {
let iter = state.stats.values().filter(|s| s.source == source).count() as u32;
let entry = state
.stats
.entry(id)
.or_insert_with(|| placeholder_rw_lock_entry(id));
entry.source = source;
entry.label = label;
entry.type_name = type_name;
entry.iter = iter;
}
RwLockEvent::Released {
id,
kind,
wait_nanos,
acquire_nanos,
} => {
let entry = state
.stats
.entry(id)
.or_insert_with(|| placeholder_rw_lock_entry(id));
let sampled = match (wait_nanos, acquire_nanos) {
(Some(wait), Some(acquire)) => Some((wait, acquire)),
_ => None,
};
match kind {
RwLockKind::Read => {
entry.read_count += 1;
if let Some((wait_nanos, acquire_nanos)) = sampled {
entry.read_sampled_count += 1;
entry.read_wait_total_nanos += wait_nanos;
entry.read_acquire_total_nanos += acquire_nanos;
RwLockEntry::record(&mut entry.read_wait_hist, wait_nanos);
RwLockEntry::record(&mut entry.read_acquire_hist, acquire_nanos);
}
}
RwLockKind::Write => {
entry.write_count += 1;
if let Some((wait_nanos, acquire_nanos)) = sampled {
entry.write_sampled_count += 1;
entry.write_wait_total_nanos += wait_nanos;
entry.write_acquire_total_nanos += acquire_nanos;
RwLockEntry::record(&mut entry.write_wait_hist, wait_nanos);
RwLockEntry::record(&mut entry.write_acquire_hist, acquire_nanos);
}
}
}
}
}
}
pub(crate) fn register_rw_lock<T>(source: &'static str, label: Option<String>) -> u32 {
let type_name = std::any::type_name::<T>();
init_rw_locks_state();
let id = next_rw_lock_id();
send_rw_lock_event(RwLockEvent::Created {
id,
source,
label,
type_name,
});
id
}
fn flush_rw_lock_buffer(
buffer: &mut Vec<RwLockEvent>,
inner: &Arc<MetaRwLock<RwLocksInternalState>>,
) {
if buffer.is_empty() {
return;
}
if let Ok(mut shared) = inner.write() {
for e in buffer.drain(..) {
process_rw_lock_event(&mut shared, e);
}
}
}
pub(crate) fn init_rw_locks_state() -> &'static RwLocksState {
RW_LOCKS_STATE.get_or_init(|| {
START_TIME.get_or_init(Instant::now);
let (shutdown_tx, shutdown_rx) = bounded::<()>(1);
let (completion_tx, completion_rx) = bounded::<()>(1);
let inner = Arc::new(meta_rw_lock!(
"rw_locks_state",
RwLocksInternalState {
stats: HashMap::new(),
},
));
let inner_clone = Arc::clone(&inner);
EVENT_QUEUES.set_active(true);
std::thread::Builder::new()
.name("hp-rw-locks".into())
.spawn(move || {
let flush_interval = std::time::Duration::from_millis(*DRAIN_INTERVAL_MS);
let mut swept: Vec<RwLockEvent> = Vec::new();
loop {
let shutdown = !matches!(
shutdown_rx.recv_timeout(flush_interval),
Err(RecvTimeoutError::Timeout)
);
if shutdown {
EVENT_QUEUES.drain_all(&mut swept);
flush_rw_lock_buffer(&mut swept, &inner_clone);
break;
}
EVENT_QUEUES.sweep(&mut swept);
flush_rw_lock_buffer(&mut swept, &inner_clone);
}
let _ = completion_tx.send(());
})
.expect("Failed to spawn rw_lock-stats-collector thread");
crate::metrics_server::start_metrics_server_once(*METRICS_SERVER_PORT);
RwLocksState {
inner,
shutdown_tx: Mutex::new(Some(shutdown_tx)),
completion_rx: Mutex::new(Some(completion_rx)),
}
})
}
pub(crate) fn compare_rw_lock_entries(a: &RwLockEntry, b: &RwLockEntry) -> std::cmp::Ordering {
match (a.label.is_some(), b.label.is_some()) {
(true, false) => std::cmp::Ordering::Less,
(false, true) => std::cmp::Ordering::Greater,
(true, true) => a
.label
.as_ref()
.unwrap()
.cmp(b.label.as_ref().unwrap())
.then_with(|| a.iter.cmp(&b.iter)),
(false, false) => a.source.cmp(b.source).then_with(|| a.iter.cmp(&b.iter)),
}
}
#[doc(hidden)]
pub trait InstrumentRwLock {
type Output;
fn instrument(self, source: &'static str, label: Option<String>) -> Self::Output;
}
#[macro_export]
macro_rules! rw_lock {
($expr:expr) => {{
const RW_LOCK_ID: &'static str = concat!(file!(), ":", line!());
$crate::InstrumentRwLock::instrument($expr, RW_LOCK_ID, None)
}};
($expr:expr, label = $label:expr) => {{
const RW_LOCK_ID: &'static str = concat!(file!(), ":", line!());
$crate::InstrumentRwLock::instrument($expr, RW_LOCK_ID, Some($label.to_string()))
}};
}