use std::sync::atomic::{AtomicI64, Ordering};
use std::sync::OnceLock;
use dashmap::DashMap;
#[derive(Default)]
pub struct Counter {
v: AtomicI64,
}
impl Counter {
#[inline]
pub fn inc(&self, n: i64) {
self.v.fetch_add(n, Ordering::Relaxed);
}
#[inline]
pub fn get(&self) -> i64 {
self.v.load(Ordering::Relaxed)
}
}
#[derive(Default)]
pub struct Gauge {
v: AtomicI64,
}
impl Gauge {
#[inline]
pub fn set(&self, val: i64) {
self.v.store(val, Ordering::Relaxed);
}
#[inline]
pub fn get(&self) -> i64 {
self.v.load(Ordering::Relaxed)
}
}
pub(crate) struct Registry {
pub(crate) counters: DashMap<String, std::sync::Arc<Counter>>,
pub(crate) gauges: DashMap<String, std::sync::Arc<Gauge>>,
}
impl Default for Registry {
fn default() -> Self {
Self {
counters: DashMap::new(),
gauges: DashMap::new(),
}
}
}
pub(crate) static REGISTRY: OnceLock<Registry> = OnceLock::new();
pub(crate) fn get_registry() -> &'static Registry {
REGISTRY.get_or_init(Registry::default)
}
pub fn counter(name: &str) -> std::sync::Arc<Counter> {
let registry = get_registry();
if let Some(c) = registry.counters.get(name) {
return c.value().clone();
}
let c = std::sync::Arc::new(Counter::default());
registry
.counters
.entry(name.to_string())
.or_insert(c.clone());
registry
.counters
.get(name)
.map(|entry| entry.value().clone())
.unwrap_or(c)
}
pub fn gauge(name: &str) -> std::sync::Arc<Gauge> {
let registry = get_registry();
if let Some(g) = registry.gauges.get(name) {
return g.value().clone();
}
let g = std::sync::Arc::new(Gauge::default());
registry.gauges.entry(name.to_string()).or_insert(g.clone());
registry
.gauges
.get(name)
.map(|entry| entry.value().clone())
.unwrap_or(g)
}
pub mod name {
pub const CLIENT_BYTES_READ_LOCAL: &str = "Client.BytesReadLocal";
pub const CLIENT_BYTES_WRITTEN_LOCAL: &str = "Client.BytesWrittenLocal";
pub const CLIENT_BYTES_WRITTEN_UFS: &str = "Client.BytesWrittenUfs";
pub const CLIENT_READ_OPS_TOTAL: &str = "Client.ReadOpsTotal";
pub const CLIENT_WRITE_OPS_TOTAL: &str = "Client.WriteOpsTotal";
pub const CLIENT_GET_STATUS_OPS: &str = "Client.GetStatusOps";
pub const CLIENT_LIST_STATUS_OPS: &str = "Client.ListStatusOps";
pub const CLIENT_CREATE_FILE_OPS: &str = "Client.CreateFileOps";
pub const CLIENT_CREATE_DIR_OPS: &str = "Client.CreateDirOps";
pub const CLIENT_DELETE_OPS: &str = "Client.DeleteOps";
pub const CLIENT_RENAME_OPS: &str = "Client.RenameOps";
pub const CLIENT_RPC_ERRORS_TOTAL: &str = "Client.RpcErrorsTotal";
pub const CLIENT_RPC_AUTH_ERRORS: &str = "Client.RpcAuthErrors";
pub const CLIENT_RPC_UNAVAILABLE_ERRORS: &str = "Client.RpcUnavailableErrors";
pub const CLIENT_READ_FAILURES: &str = "Client.ReadFailures";
pub const CLIENT_WRITE_FAILURES: &str = "Client.WriteFailures";
pub const CLIENT_READ_LATENCY_US: &str = "Client.ReadLatencyUs";
pub const CLIENT_WRITE_LATENCY_US: &str = "Client.WriteLatencyUs";
pub const CLIENT_GET_STATUS_LATENCY_US: &str = "Client.GetStatusLatencyUs";
pub const CLIENT_LIST_STATUS_LATENCY_US: &str = "Client.ListStatusLatencyUs";
pub const CLIENT_WORKER_CONNECTIONS_ACTIVE: &str = "Client.WorkerConnectionsActive";
pub const CLIENT_WORKER_RECONNECTS_TOTAL: &str = "Client.WorkerReconnectsTotal";
pub const CLIENT_WORKER_RECONNECTS_COALESCED: &str = "Client.WorkerReconnectsCoalesced";
pub const CLIENT_BLOCKS_READ_IN_PROGRESS: &str = "Client.BlocksReadInProgress";
pub const CLIENT_BLOCKS_WRITTEN_IN_PROGRESS: &str = "Client.BlocksWrittenInProgress";
pub const CLIENT_BLOCKS_READ_TOTAL: &str = "Client.BlocksReadTotal";
pub const CLIENT_BLOCKS_WRITTEN_TOTAL: &str = "Client.BlocksWrittenTotal";
pub const CLIENT_SC_OPEN_SUCCESS: &str = "Client.ShortCircuitOpenSuccess";
pub const CLIENT_SC_OPENLOCAL_FAIL: &str = "Client.ShortCircuitOpenLocalFail";
pub const CLIENT_SC_FILE_OPEN_FAIL: &str = "Client.ShortCircuitFileOpenFail";
pub const CLIENT_SC_MMAP_FAIL: &str = "Client.ShortCircuitMmapFail";
pub const CLIENT_SC_READ_BYTES: &str = "Client.ShortCircuitReadBytes";
pub const CLIENT_SC_READ_CALLS: &str = "Client.ShortCircuitReadCalls";
pub const CLIENT_SC_CACHE_HITS: &str = "Client.ShortCircuitCacheHits";
pub const CLIENT_SC_CACHE_EVICTIONS: &str = "Client.ShortCircuitCacheEvictions";
pub const CLIENT_SC_NEG_CACHE_HITS: &str = "Client.ShortCircuitNegCacheHits";
pub const CLIENT_SC_ACTIVE_READERS: &str = "Client.ShortCircuitActiveReaders";
pub const CLIENT_SC_PREFETCH_CALLS: &str = "Client.ShortCircuitPrefetchCalls";
pub const CLIENT_SC_PREFETCH_BYTES: &str = "Client.ShortCircuitPrefetchBytes";
pub const CLIENT_SC_PREFETCH_MADVISE: &str = "Client.ShortCircuitPrefetchMadvise";
pub const CLIENT_SC_DECISION_HIT: &str = "Client.ShortCircuitDecisionHit";
pub const CLIENT_SC_DECISION_SKIPPED: &str = "Client.ShortCircuitDecisionSkipped";
pub const CLIENT_SC_DECISION_FALLBACK_OPEN: &str = "Client.ShortCircuitDecisionFallbackOpen";
pub const CLIENT_SC_DECISION_FALLBACK_READ: &str = "Client.ShortCircuitDecisionFallbackRead";
pub const CLIENT_SC_DECISION_SEMANTIC_ERROR: &str = "Client.ShortCircuitDecisionSemanticError";
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn counter_inc_and_get() {
let c = Counter::default();
assert_eq!(c.get(), 0);
c.inc(42);
assert_eq!(c.get(), 42);
c.inc(8);
assert_eq!(c.get(), 50);
c.inc(-10);
assert_eq!(c.get(), 40);
}
#[test]
fn counter_negative_increment() {
let c = Counter::default();
c.inc(-5);
assert_eq!(c.get(), -5);
}
#[test]
fn gauge_set_and_get() {
let g = Gauge::default();
assert_eq!(g.get(), 0);
g.set(99);
assert_eq!(g.get(), 99);
g.set(-10);
assert_eq!(g.get(), -10);
}
#[test]
fn registry_counter_factory() {
let c1 = counter("my_counter");
assert_eq!(c1.get(), 0);
c1.inc(10);
let c2 = counter("my_counter");
assert_eq!(c2.get(), 10);
let c3 = counter("other_counter");
assert_eq!(c3.get(), 0);
}
#[test]
fn registry_gauge_factory() {
let g1 = gauge("my_gauge");
assert_eq!(g1.get(), 0);
g1.set(55);
let g2 = gauge("my_gauge");
assert_eq!(g2.get(), 55);
let g3 = gauge("other_gauge");
assert_eq!(g3.get(), 0);
}
#[test]
fn registry_counter_concurrent() {
use std::thread;
let c = counter("concurrent_counter");
let mut handles = vec![];
for _ in 0..10 {
let c = c.clone();
let handle = thread::spawn(move || {
for _ in 0..1000 {
c.inc(1);
}
});
handles.push(handle);
}
for h in handles {
h.join().unwrap();
}
assert_eq!(c.get(), 10_000);
}
#[test]
fn name_constants() {
assert_eq!(name::CLIENT_BYTES_READ_LOCAL, "Client.BytesReadLocal");
assert_eq!(name::CLIENT_BYTES_WRITTEN_LOCAL, "Client.BytesWrittenLocal");
assert_eq!(name::CLIENT_BYTES_WRITTEN_UFS, "Client.BytesWrittenUfs");
}
#[test]
fn name_constants_rpc_ops() {
assert_eq!(name::CLIENT_READ_OPS_TOTAL, "Client.ReadOpsTotal");
assert_eq!(name::CLIENT_WRITE_OPS_TOTAL, "Client.WriteOpsTotal");
assert_eq!(name::CLIENT_GET_STATUS_OPS, "Client.GetStatusOps");
assert_eq!(name::CLIENT_LIST_STATUS_OPS, "Client.ListStatusOps");
assert_eq!(name::CLIENT_CREATE_FILE_OPS, "Client.CreateFileOps");
assert_eq!(name::CLIENT_CREATE_DIR_OPS, "Client.CreateDirOps");
assert_eq!(name::CLIENT_DELETE_OPS, "Client.DeleteOps");
assert_eq!(name::CLIENT_RENAME_OPS, "Client.RenameOps");
}
#[test]
fn name_constants_errors() {
assert_eq!(name::CLIENT_RPC_ERRORS_TOTAL, "Client.RpcErrorsTotal");
assert_eq!(name::CLIENT_RPC_AUTH_ERRORS, "Client.RpcAuthErrors");
assert_eq!(
name::CLIENT_RPC_UNAVAILABLE_ERRORS,
"Client.RpcUnavailableErrors"
);
assert_eq!(name::CLIENT_READ_FAILURES, "Client.ReadFailures");
assert_eq!(name::CLIENT_WRITE_FAILURES, "Client.WriteFailures");
}
#[test]
fn name_constants_latency() {
assert_eq!(name::CLIENT_READ_LATENCY_US, "Client.ReadLatencyUs");
assert_eq!(name::CLIENT_WRITE_LATENCY_US, "Client.WriteLatencyUs");
assert_eq!(
name::CLIENT_GET_STATUS_LATENCY_US,
"Client.GetStatusLatencyUs"
);
assert_eq!(
name::CLIENT_LIST_STATUS_LATENCY_US,
"Client.ListStatusLatencyUs"
);
}
#[test]
fn name_constants_pool_and_blocks() {
assert_eq!(
name::CLIENT_WORKER_CONNECTIONS_ACTIVE,
"Client.WorkerConnectionsActive"
);
assert_eq!(
name::CLIENT_WORKER_RECONNECTS_TOTAL,
"Client.WorkerReconnectsTotal"
);
assert_eq!(
name::CLIENT_WORKER_RECONNECTS_COALESCED,
"Client.WorkerReconnectsCoalesced"
);
assert_eq!(
name::CLIENT_BLOCKS_READ_IN_PROGRESS,
"Client.BlocksReadInProgress"
);
assert_eq!(
name::CLIENT_BLOCKS_WRITTEN_IN_PROGRESS,
"Client.BlocksWrittenInProgress"
);
assert_eq!(name::CLIENT_BLOCKS_READ_TOTAL, "Client.BlocksReadTotal");
assert_eq!(
name::CLIENT_BLOCKS_WRITTEN_TOTAL,
"Client.BlocksWrittenTotal"
);
}
#[test]
fn name_constants_sc_decision_histogram() {
assert_eq!(
name::CLIENT_SC_DECISION_HIT,
"Client.ShortCircuitDecisionHit"
);
assert_eq!(
name::CLIENT_SC_DECISION_SKIPPED,
"Client.ShortCircuitDecisionSkipped"
);
assert_eq!(
name::CLIENT_SC_DECISION_FALLBACK_OPEN,
"Client.ShortCircuitDecisionFallbackOpen"
);
assert_eq!(
name::CLIENT_SC_DECISION_FALLBACK_READ,
"Client.ShortCircuitDecisionFallbackRead"
);
assert_eq!(
name::CLIENT_SC_DECISION_SEMANTIC_ERROR,
"Client.ShortCircuitDecisionSemanticError"
);
}
#[test]
fn sc_decision_counters_are_registered_and_shared() {
for cname in [
name::CLIENT_SC_DECISION_HIT,
name::CLIENT_SC_DECISION_SKIPPED,
name::CLIENT_SC_DECISION_FALLBACK_OPEN,
name::CLIENT_SC_DECISION_FALLBACK_READ,
name::CLIENT_SC_DECISION_SEMANTIC_ERROR,
] {
let c1 = counter(cname);
let c2 = counter(cname);
let base = c1.get();
c2.inc(1);
assert_eq!(
c1.get(),
base + 1,
"counter '{}' must be process-wide shared",
cname
);
}
}
#[test]
fn sc_decision_counter_names_are_pairwise_distinct() {
let names = [
name::CLIENT_SC_DECISION_HIT,
name::CLIENT_SC_DECISION_SKIPPED,
name::CLIENT_SC_DECISION_FALLBACK_OPEN,
name::CLIENT_SC_DECISION_FALLBACK_READ,
name::CLIENT_SC_DECISION_SEMANTIC_ERROR,
];
for i in 0..names.len() {
for j in (i + 1)..names.len() {
assert_ne!(names[i], names[j], "duplicate SC decision name");
}
}
}
}