use std::collections::HashMap;
use std::time::Instant;
use parking_lot::Mutex;
#[derive(Clone)]
struct HistogramData {
sum: f64,
count: u64,
buckets: Vec<(f64, u64)>,
}
impl HistogramData {
fn new() -> Self {
let boundaries = vec![
0.0001, 0.0005, 0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0,
30.0, 60.0,
];
let buckets = boundaries.into_iter().map(|b| (b, 0u64)).collect();
Self {
sum: 0.0,
count: 0,
buckets,
}
}
fn observe(&mut self, value: f64) {
self.sum += value;
self.count += 1;
for (boundary, count) in &mut self.buckets {
if value <= *boundary {
*count += 1;
}
}
}
fn render(&self, name: &str, labels: &str, help: &str) -> String {
let mut out = String::new();
out.push_str(&format!("# HELP {} {}\n", name, help));
out.push_str(&format!("# TYPE {} histogram\n", name));
for (boundary, count) in &self.buckets {
out.push_str(&format!(
"{}_bucket{{{},le=\"{}\"}} {}\n",
name, labels, boundary, count
));
}
out.push_str(&format!(
"{}_bucket{{{},le=\"+Inf\"}} {}\n",
name, labels, self.count
));
out.push_str(&format!("{}_sum{{{}}} {}\n", name, labels, self.sum));
out.push_str(&format!("{}_count{{{}}} {}\n", name, labels, self.count));
out
}
}
pub struct MetricsStore {
handler_durations: Mutex<HashMap<String, HistogramData>>,
lock_waits: Mutex<HashMap<String, HistogramData>>,
request_counts: Mutex<HashMap<String, u64>>,
recall_rejected_counts: Mutex<HashMap<&'static str, u64>>,
recall_in_flight_gauge: std::sync::atomic::AtomicI64,
expansion_concurrent_gauge: std::sync::atomic::AtomicI64,
raft_term_changes: std::sync::atomic::AtomicU64,
raft_elections: Mutex<HashMap<&'static str, u64>>,
raft_heartbeat_lag: Mutex<HistogramData>,
raft_task_poll_latency: Mutex<HistogramData>,
recall_request_counts: Mutex<HashMap<(&'static str, bool), u64>>,
recall_request_top_k: Mutex<HashMap<&'static str, HistogramData>>,
embedder_failures: Mutex<HashMap<&'static str, u64>>,
null_embedding_counts: Mutex<HashMap<i64, i64>>,
enrichment_paused: Mutex<HashMap<String, u64>>,
enrichment_resumed: Mutex<HashMap<String, u64>>,
enrichment_pending_at_pause: Mutex<HistogramData>,
maintenance_aggs: Mutex<HashMap<String, MaintenanceAgg>>,
maintenance_skipped: Mutex<HashMap<(String, &'static str), u64>>,
maintenance_duration_ms: Mutex<HistogramData>,
}
#[derive(Default, Clone)]
pub struct MaintenanceAgg {
pub runs: u64,
pub failures: u64,
pub consolidations: u64,
pub conflicts_resolved: u64,
pub triggers_pruned: u64,
pub entities_linked: u64,
pub relations_upserted: u64,
pub pass_errors: u64,
pub last_duration_ms: f64,
}
impl MetricsStore {
pub fn new() -> Self {
Self {
handler_durations: Mutex::new(HashMap::new()),
lock_waits: Mutex::new(HashMap::new()),
request_counts: Mutex::new(HashMap::new()),
recall_rejected_counts: Mutex::new(HashMap::new()),
recall_in_flight_gauge: std::sync::atomic::AtomicI64::new(0),
expansion_concurrent_gauge: std::sync::atomic::AtomicI64::new(0),
raft_term_changes: std::sync::atomic::AtomicU64::new(0),
raft_elections: Mutex::new(HashMap::new()),
raft_heartbeat_lag: Mutex::new(HistogramData::new()),
raft_task_poll_latency: Mutex::new(HistogramData::new()),
recall_request_counts: Mutex::new(HashMap::new()),
recall_request_top_k: Mutex::new(HashMap::new()),
embedder_failures: Mutex::new(HashMap::new()),
null_embedding_counts: Mutex::new(HashMap::new()),
enrichment_paused: Mutex::new(HashMap::new()),
enrichment_resumed: Mutex::new(HashMap::new()),
enrichment_pending_at_pause: Mutex::new(HistogramData::new()),
maintenance_aggs: Mutex::new(HashMap::new()),
maintenance_skipped: Mutex::new(HashMap::new()),
maintenance_duration_ms: Mutex::new(HistogramData::new()),
}
}
pub fn record_handler_duration(&self, handler: &str, duration_secs: f64) {
let mut map = self.handler_durations.lock();
map.entry(handler.to_string())
.or_insert_with(HistogramData::new)
.observe(duration_secs);
}
pub fn record_lock_wait(&self, lock_name: &str, duration_secs: f64) {
let mut map = self.lock_waits.lock();
map.entry(lock_name.to_string())
.or_insert_with(HistogramData::new)
.observe(duration_secs);
}
pub fn increment_request(&self, handler: &str) {
let mut map = self.request_counts.lock();
*map.entry(handler.to_string()).or_insert(0) += 1;
}
pub fn render_prometheus(&self) -> String {
let mut out = String::with_capacity(4096);
{
let map = self.handler_durations.lock();
for (handler, hist) in map.iter() {
out.push_str(&hist.render(
"yantrikdb_handler_duration_seconds",
&format!("handler=\"{}\"", handler),
"Duration of HTTP handler execution in seconds",
));
}
}
{
let map = self.lock_waits.lock();
for (lock_name, hist) in map.iter() {
out.push_str(&hist.render(
"yantrikdb_lock_wait_seconds",
&format!("lock=\"{}\"", lock_name),
"Time spent waiting to acquire a lock in seconds",
));
}
}
{
let map = self.request_counts.lock();
if !map.is_empty() {
out.push_str("# HELP yantrikdb_requests_total Total HTTP requests per handler\n");
out.push_str("# TYPE yantrikdb_requests_total counter\n");
for (handler, count) in map.iter() {
out.push_str(&format!(
"yantrikdb_requests_total{{handler=\"{}\"}} {}\n",
handler, count,
));
}
}
}
{
let map = self.recall_rejected_counts.lock();
if !map.is_empty() {
out.push_str(
"# HELP yantrikdb_recall_rejected_total Recall requests rejected by admission control, by reason\n",
);
out.push_str("# TYPE yantrikdb_recall_rejected_total counter\n");
for (reason, count) in map.iter() {
out.push_str(&format!(
"yantrikdb_recall_rejected_total{{reason=\"{}\"}} {}\n",
reason, count
));
}
}
}
out.push_str("# HELP yantrikdb_recall_in_flight Current in-flight recalls (any kind)\n");
out.push_str("# TYPE yantrikdb_recall_in_flight gauge\n");
out.push_str(&format!(
"yantrikdb_recall_in_flight {}\n",
self.recall_in_flight_gauge
.load(std::sync::atomic::Ordering::Relaxed)
));
out.push_str("# HELP yantrikdb_expansion_concurrent Current concurrent expanded recalls\n");
out.push_str("# TYPE yantrikdb_expansion_concurrent gauge\n");
out.push_str(&format!(
"yantrikdb_expansion_concurrent {}\n",
self.expansion_concurrent_gauge
.load(std::sync::atomic::Ordering::Relaxed)
));
out.push_str(
"# HELP yantrikdb_raft_term_changes_total Raft term increments (new election or stepdown)\n",
);
out.push_str("# TYPE yantrikdb_raft_term_changes_total counter\n");
out.push_str(&format!(
"yantrikdb_raft_term_changes_total {}\n",
self.raft_term_changes
.load(std::sync::atomic::Ordering::Relaxed)
));
{
let map = self.raft_elections.lock();
if !map.is_empty() {
out.push_str("# HELP yantrikdb_raft_elections_total Raft elections by outcome\n");
out.push_str("# TYPE yantrikdb_raft_elections_total counter\n");
for (result, count) in map.iter() {
out.push_str(&format!(
"yantrikdb_raft_elections_total{{result=\"{}\"}} {}\n",
result, count
));
}
}
}
{
let hist = self.raft_heartbeat_lag.lock();
if hist.count > 0 {
out.push_str(&hist.render(
"yantrikdb_raft_heartbeat_lag_seconds",
"",
"Heartbeat round-trip lag in seconds",
));
}
}
{
let hist = self.raft_task_poll_latency.lock();
if hist.count > 0 {
out.push_str(&hist.render(
"yantrikdb_raft_task_poll_latency_seconds",
"",
"Control-runtime task scheduling latency in seconds (acceptance gate signal)",
));
}
}
{
let map = self.recall_request_counts.lock();
if !map.is_empty() {
out.push_str("# HELP yantrikdb_recall_requests_total Recall requests received\n");
out.push_str("# TYPE yantrikdb_recall_requests_total counter\n");
for ((version, expand), count) in map.iter() {
out.push_str(&format!(
"yantrikdb_recall_requests_total{{api_version=\"{}\",expand=\"{}\"}} {}\n",
version, expand, count
));
}
}
}
{
let map = self.recall_request_top_k.lock();
for (version, hist) in map.iter() {
if hist.count > 0 {
out.push_str(&hist.render(
"yantrikdb_recall_request_top_k",
&format!("api_version=\"{}\"", version),
"Distribution of requested top_k values",
));
}
}
}
render_version_gauges_if_set(&mut out);
out
}
}
pub fn increment_recall_rejected(reason: &'static str) {
let mut map = global().recall_rejected_counts.lock();
*map.entry(reason).or_insert(0) += 1;
}
pub fn increment_embedder_failure(handler: &'static str) {
let mut map = global().embedder_failures.lock();
*map.entry(handler).or_insert(0) += 1;
}
pub fn embedder_failure_counts() -> Vec<(&'static str, u64)> {
let map = global().embedder_failures.lock();
map.iter().map(|(k, v)| (*k, *v)).collect()
}
pub fn set_null_embedding_count(tenant_id: i64, count: i64) {
let mut map = global().null_embedding_counts.lock();
map.insert(tenant_id, count);
}
pub fn null_embedding_counts_snapshot() -> Vec<(i64, i64)> {
let map = global().null_embedding_counts.lock();
map.iter().map(|(k, v)| (*k, *v)).collect()
}
pub fn record_enrichment_paused(db_name: &str, pending: u64) {
{
let mut map = global().enrichment_paused.lock();
*map.entry(db_name.to_string()).or_insert(0) += 1;
}
global()
.enrichment_pending_at_pause
.lock()
.observe(pending as f64);
}
pub fn record_enrichment_resumed(db_name: &str) {
let mut map = global().enrichment_resumed.lock();
*map.entry(db_name.to_string()).or_insert(0) += 1;
}
#[allow(clippy::too_many_arguments)]
pub fn record_maintenance_cycle(
db_name: &str,
duration_ms: f64,
consolidations: u64,
conflicts_resolved: u64,
triggers_pruned: u64,
entities_linked: u64,
relations_upserted: u64,
pass_errors: u64,
) {
{
let mut map = global().maintenance_aggs.lock();
let agg = map.entry(db_name.to_string()).or_default();
agg.runs += 1;
agg.consolidations += consolidations;
agg.conflicts_resolved += conflicts_resolved;
agg.triggers_pruned += triggers_pruned;
agg.entities_linked += entities_linked;
agg.relations_upserted += relations_upserted;
agg.pass_errors += pass_errors;
agg.last_duration_ms = duration_ms;
}
global().maintenance_duration_ms.lock().observe(duration_ms);
}
pub fn record_maintenance_skipped(db_name: &str, reason: &'static str) {
let mut map = global().maintenance_skipped.lock();
*map.entry((db_name.to_string(), reason)).or_insert(0) += 1;
}
pub fn record_maintenance_failed(db_name: &str) {
let mut map = global().maintenance_aggs.lock();
map.entry(db_name.to_string()).or_default().failures += 1;
}
pub fn maintenance_aggs_snapshot() -> Vec<(String, MaintenanceAgg)> {
let map = global().maintenance_aggs.lock();
map.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
}
pub fn maintenance_skipped_snapshot() -> Vec<(String, &'static str, u64)> {
let map = global().maintenance_skipped.lock();
map.iter()
.map(|((db, reason), v)| (db.clone(), *reason, *v))
.collect()
}
pub fn maintenance_duration_totals() -> (u64, f64) {
let h = global().maintenance_duration_ms.lock();
(h.count, h.sum)
}
pub fn enrichment_paused_snapshot() -> Vec<(String, u64)> {
let map = global().enrichment_paused.lock();
map.iter().map(|(k, v)| (k.clone(), *v)).collect()
}
pub fn enrichment_resumed_snapshot() -> Vec<(String, u64)> {
let map = global().enrichment_resumed.lock();
map.iter().map(|(k, v)| (k.clone(), *v)).collect()
}
pub fn enrichment_pending_at_pause_totals() -> (u64, f64) {
let h = global().enrichment_pending_at_pause.lock();
(h.count, h.sum)
}
pub fn set_recall_in_flight_gauge(value: i64) {
global()
.recall_in_flight_gauge
.store(value, std::sync::atomic::Ordering::Relaxed);
}
pub fn set_expansion_concurrent_gauge(value: i64) {
global()
.expansion_concurrent_gauge
.store(value, std::sync::atomic::Ordering::Relaxed);
}
#[allow(dead_code)]
pub fn increment_raft_term_changes() {
global()
.raft_term_changes
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
#[allow(dead_code)]
pub fn record_raft_election(result: &'static str) {
debug_assert!(
matches!(result, "won" | "lost" | "stepped_down"),
"raft election result must be one of won|lost|stepped_down"
);
let mut map = global().raft_elections.lock();
*map.entry(result).or_insert(0) += 1;
}
#[allow(dead_code)]
pub fn record_raft_heartbeat_lag(duration: std::time::Duration) {
global()
.raft_heartbeat_lag
.lock()
.observe(duration.as_secs_f64());
}
pub fn record_raft_task_poll_latency(duration: std::time::Duration) {
global()
.raft_task_poll_latency
.lock()
.observe(duration.as_secs_f64());
}
pub fn record_recall_request(api_version: &'static str, expand: bool) {
let mut map = global().recall_request_counts.lock();
*map.entry((api_version, expand)).or_insert(0) += 1;
}
pub fn record_recall_top_k(api_version: &'static str, top_k: usize) {
let mut map = global().recall_request_top_k.lock();
map.entry(api_version)
.or_insert_with(HistogramData::new)
.observe(top_k as f64);
}
pub fn increment_version_rejection(reason: &'static str) {
let mut map = global().recall_rejected_counts.lock();
*map.entry(reason).or_insert(0) += 1;
}
pub type VersionGaugeRenderer = fn(&mut String);
static VERSION_GAUGE_RENDERER: std::sync::OnceLock<VersionGaugeRenderer> =
std::sync::OnceLock::new();
pub fn set_version_gauge_renderer(f: VersionGaugeRenderer) {
let _ = VERSION_GAUGE_RENDERER.set(f);
}
fn render_version_gauges_if_set(out: &mut String) {
if let Some(renderer) = VERSION_GAUGE_RENDERER.get() {
renderer(out);
}
}
#[cfg(test)]
pub fn raft_task_poll_latency_p99() -> f64 {
let hist = global().raft_task_poll_latency.lock();
if hist.count == 0 {
return 0.0;
}
let target = (hist.count as f64 * 0.99) as u64;
let mut acc = 0u64;
let mut last_boundary = 0.0;
for (boundary, bucket_count) in &hist.buckets {
acc = *bucket_count;
if acc >= target {
return *boundary;
}
last_boundary = *boundary;
}
last_boundary.max(hist.sum / hist.count.max(1) as f64)
}
static METRICS: std::sync::OnceLock<MetricsStore> = std::sync::OnceLock::new();
pub fn global() -> &'static MetricsStore {
METRICS.get_or_init(MetricsStore::new)
}
pub struct HandlerTimer {
handler: &'static str,
start: Instant,
}
impl HandlerTimer {
pub fn new(handler: &'static str) -> Self {
global().increment_request(handler);
Self {
handler,
start: Instant::now(),
}
}
}
impl Drop for HandlerTimer {
fn drop(&mut self) {
global().record_handler_duration(self.handler, self.start.elapsed().as_secs_f64());
}
}
pub fn record_engine_lock_wait(duration: std::time::Duration) {
global().record_lock_wait("engine", duration.as_secs_f64());
}
pub struct LockHoldTimer {
op: &'static str,
start: std::time::Instant,
}
impl LockHoldTimer {
pub fn start(op: &'static str) -> Self {
Self {
op,
start: std::time::Instant::now(),
}
}
}
impl Drop for LockHoldTimer {
fn drop(&mut self) {
record_engine_lock_hold(self.op, self.start.elapsed());
}
}
pub fn record_engine_lock_hold(op: &str, duration: std::time::Duration) {
let secs = duration.as_secs_f64();
global().record_lock_wait(&format!("engine_hold_{op}"), secs);
let slow_ms: u128 = std::env::var("YANTRIKDB_SLOW_LOCK_MS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(50);
let elapsed_ms = duration.as_millis();
if elapsed_ms > slow_ms {
tracing::warn!(
op = %op,
hold_ms = %elapsed_ms,
threshold_ms = %slow_ms,
"engine lock held longer than threshold (slow-holder)"
);
}
}
#[allow(dead_code)]
pub fn record_control_lock_wait(duration: std::time::Duration) {
global().record_lock_wait("control", duration.as_secs_f64());
}
#[allow(dead_code)]
#[cfg(debug_assertions)]
pub mod lock_rank {
pub const CONTROL: u8 = 0;
pub const TENANT_POOL: u8 = 1;
pub const ENGINE: u8 = 2;
pub const CONN: u8 = 3;
pub const VEC_INDEX: u8 = 4;
pub const GRAPH_INDEX: u8 = 5;
pub const SCORING_CACHE: u8 = 6;
pub const ACTIVE_SESSIONS: u8 = 7;
pub const HLC: u8 = 8;
}
#[allow(dead_code)]
#[cfg(debug_assertions)]
pub fn check_lock_order(rank: u8, lock_name: &str) {
thread_local! {
static HELD_RANKS: std::cell::RefCell<Vec<(u8, &'static str)>> = const { std::cell::RefCell::new(Vec::new()) };
}
HELD_RANKS.with(|held| {
let held = held.borrow();
for &(held_rank, held_name) in held.iter() {
if held_rank > rank {
panic!(
"LOCK ORDER VIOLATION: trying to acquire '{}' (rank {}) \
while holding '{}' (rank {}). See CONCURRENCY.md Rule 3.",
lock_name, rank, held_name, held_rank,
);
}
}
});
}
#[allow(dead_code)]
#[cfg(debug_assertions)]
pub fn push_lock(rank: u8, lock_name: &'static str) {
thread_local! {
static HELD_RANKS: std::cell::RefCell<Vec<(u8, &'static str)>> = const { std::cell::RefCell::new(Vec::new()) };
}
HELD_RANKS.with(|held| {
held.borrow_mut().push((rank, lock_name));
});
}
#[allow(dead_code)]
#[cfg(debug_assertions)]
pub fn pop_lock(rank: u8) {
thread_local! {
static HELD_RANKS: std::cell::RefCell<Vec<(u8, &'static str)>> = const { std::cell::RefCell::new(Vec::new()) };
}
HELD_RANKS.with(|held| {
let mut held = held.borrow_mut();
if let Some(pos) = held.iter().rposition(|(r, _)| *r == rank) {
held.remove(pos);
}
});
}
#[cfg(not(debug_assertions))]
pub fn check_lock_order(_rank: u8, _lock_name: &str) {}
#[cfg(not(debug_assertions))]
pub fn push_lock(_rank: u8, _lock_name: &'static str) {}
#[cfg(not(debug_assertions))]
pub fn pop_lock(_rank: u8) {}