diff --git a/libsql-server/src/connection/mod.rs b/libsql-server/src/connection/mod.rs index 167a1a595c..1e5f7dd369 100644 --- a/libsql-server/src/connection/mod.rs +++ b/libsql-server/src/connection/mod.rs @@ -12,8 +12,8 @@ use crate::auth::Authenticated; use crate::error::Error; use crate::http::user::timing::sample_time; use crate::metrics::{ - CONCURRENT_CONNECTIONS_COUNT, CONNECTION_ALIVE_DURATION, CONNECTION_CREATE_TIME, - TOTAL_RESPONSE_SIZE_HIST, + record_connection_alive_duration, record_connection_create_time, + record_total_response_size_before_lock, CONCURRENT_CONNECTIONS_COUNT, }; use crate::namespace::meta_store::MetaStore; use crate::namespace::NamespaceName; @@ -360,8 +360,9 @@ impl MakeConnection for MakeThrottledConnection { "Available semaphore units: {}", self.semaphore.available_permits() ); - TOTAL_RESPONSE_SIZE_HIST - .record(TOTAL_RESPONSE_SIZE.load(std::sync::atomic::Ordering::Relaxed) as f64); + record_total_response_size_before_lock( + TOTAL_RESPONSE_SIZE.load(std::sync::atomic::Ordering::Relaxed) as f64, + ); let units = self.units_to_take(); let waiters_guard = WaitersGuard::new(&self.waiters); if (waiters_guard.waiters.load(Ordering::Relaxed) as u64) >= self.max_concurrent_requests { @@ -389,7 +390,7 @@ impl MakeConnection for MakeThrottledConnection { let inner = self.connection_maker.create().await?; CONCURRENT_CONNECTIONS_COUNT.increment(1.0); - CONNECTION_CREATE_TIME.record(before_create.elapsed()); + record_connection_create_time(before_create.elapsed()); Ok(TrackedConnection { permit, @@ -411,7 +412,7 @@ impl Drop for TrackedConnection { fn drop(&mut self) { sample_time("connection-duration", self.created_at.elapsed()); CONCURRENT_CONNECTIONS_COUNT.decrement(1.0); - CONNECTION_ALIVE_DURATION.record(self.created_at.elapsed()); + record_connection_alive_duration(self.created_at.elapsed()); } } diff --git a/libsql-server/src/http/admin/mod.rs b/libsql-server/src/http/admin/mod.rs index 2d8de1cdd1..ca5ce6a45a 100644 --- a/libsql-server/src/http/admin/mod.rs +++ b/libsql-server/src/http/admin/mod.rs @@ -59,6 +59,15 @@ impl FromRef>> for Metrics { static PROM_HANDLE: Mutex> = Mutex::new(OnceCell::new()); +/// Prometheus recorder configuration shared by the admin server and the metrics tests. +/// +/// Idle metrics (any kind) are evicted `idle_timeout` after their last update. Histograms are +/// therefore recorded through the `metrics` macros (see `crate::metrics`), never through cached +/// handles, so that an evicted histogram is re-registered on its next sample instead of leaking. +pub(crate) fn prometheus_builder(idle_timeout: Duration) -> PrometheusBuilder { + PrometheusBuilder::new().idle_timeout(metrics_util::MetricKindMask::ALL, Some(idle_timeout)) +} + pub async fn run( acceptor: A, user_http_server: Arc, @@ -80,19 +89,18 @@ where let lock = PROM_HANDLE.lock(); let prom_handle = lock.get_or_init(|| { tracing::info!("initializing prometheus metrics"); - let b = PrometheusBuilder::new().idle_timeout( - metrics_util::MetricKindMask::ALL, - Some(Duration::from_secs(120)), - ); + let b = prometheus_builder(Duration::from_secs(120)); - if let Some(app_label) = app_label { + let handle = if let Some(app_label) = app_label { b.add_global_label("app", app_label) .add_global_label("version", ver) .install_recorder() .unwrap() } else { b.install_recorder().unwrap() - } + }; + crate::metrics::describe_histograms(); + handle }); tokio::task::spawn(async move { diff --git a/libsql-server/src/metrics.rs b/libsql-server/src/metrics.rs index e5de216b79..5ad7c45618 100644 --- a/libsql-server/src/metrics.rs +++ b/libsql-server/src/metrics.rs @@ -1,7 +1,16 @@ #![allow(dead_code)] +//! Process-wide metrics. +//! +//! Counters and gauges are cached `Lazy` handles. Histograms MUST NOT be: the Prometheus exporter is +//! configured with an idle timeout (`http::admin`), and when it evicts a histogram key a cached +//! `metrics::Histogram` handle keeps pushing 16-byte samples into an `AtomicBucket` the exporter will +//! never drain again (unbounded memory growth). Record histograms through the `record_*` functions +//! below, which use the `histogram!` macro and therefore re-register the key on every call. +use std::time::Duration; + use metrics::{ - describe_counter, describe_gauge, describe_histogram, register_counter, register_gauge, - register_histogram, Counter, Gauge, Histogram, + describe_counter, describe_gauge, describe_histogram, histogram, register_counter, + register_gauge, Counter, Gauge, }; use once_cell::sync::Lazy; @@ -27,67 +36,38 @@ pub static CONCURRENT_CONNECTIONS_COUNT: Lazy = Lazy::new(|| { describe_gauge!(NAME, "number of concurrent connections"); register_gauge!(NAME) }); -pub static TOTAL_RESPONSE_SIZE_HIST: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_total_response_size_before_lock"; - describe_histogram!(NAME, "total response size value before connection lock"); - register_histogram!(NAME) -}); +/// Total in-flight response size observed before a connection lock is taken. +#[inline] +pub fn record_total_response_size_before_lock(bytes: f64) { + histogram!("libsql_server_total_response_size_before_lock", bytes); +} pub static STREAM_HANDLES_COUNT: Lazy = Lazy::new(|| { const NAME: &str = "libsql_server_stream_handles"; describe_gauge!(NAME, "amount of in-memory stream handles"); register_gauge!(NAME) }); -pub static NAMESPACE_LOAD_LATENCY: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_namespace_load_latency"; - describe_histogram!(NAME, "latency is us when loading a namespace"); - register_histogram!(NAME) -}); -pub static CONNECTION_CREATE_TIME: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_connection_create_time"; - describe_histogram!(NAME, "time to create a connection"); - register_histogram!(NAME) -}); -pub static CONNECTION_ALIVE_DURATION: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_connection_alive_duration"; - describe_histogram!(NAME, "duration for which a connection was kept alive"); - register_histogram!(NAME) -}); -pub static WRITE_TXN_DURATION: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_write_txn_duration"; - describe_histogram!(NAME, "duration for which a write transaction was kept open"); - register_histogram!(NAME) -}); - -pub static STATEMENT_EXECUTION_TIME: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_statement_execution_time"; - describe_histogram!(NAME, "time to execute a statement"); - register_histogram!(NAME) -}); +#[inline] +pub fn record_namespace_load_latency(elapsed: Duration) { + histogram!("libsql_server_namespace_load_latency", elapsed); +} +#[inline] +pub fn record_connection_create_time(elapsed: Duration) { + histogram!("libsql_server_connection_create_time", elapsed); +} +#[inline] +pub fn record_connection_alive_duration(elapsed: Duration) { + histogram!("libsql_server_connection_alive_duration", elapsed); +} pub static VACUUM_COUNT: Lazy = Lazy::new(|| { const NAME: &str = "libsql_server_vacuum_count"; describe_counter!(NAME, "number of vacuum operations"); register_counter!(NAME) }); -pub static WAL_CHECKPOINT_TIME: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_wal_checkpoint_time"; - describe_histogram!(NAME, "time to checkpoint the WAL"); - register_histogram!(NAME) -}); pub static WAL_CHECKPOINT_COUNT: Lazy = Lazy::new(|| { const NAME: &str = "libsql_server_wal_checkpoint_count"; describe_counter!(NAME, "number of WAL checkpoints"); register_counter!(NAME) }); -pub static STATEMENT_MEM_USED_BYTES: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_statement_mem_used_bytes"; - describe_histogram!(NAME, "memory used by a prepared statement"); - register_histogram!(NAME) -}); -pub static RETURNED_BYTES: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_returned_bytes"; - describe_histogram!(NAME, "number of bytes of values returned to the client"); - register_histogram!(NAME) -}); pub static PROGRAM_EXEC_COUNT: Lazy = Lazy::new(|| { const NAME: &str = "libsql_server_libsql_execute_program"; describe_counter!(NAME, "number of hrana program executions"); @@ -127,14 +107,10 @@ pub static DIRTY_STARTUP: Lazy = Lazy::new(|| { ); register_counter!(NAME) }); -pub static REPLICATION_LATENCY: Lazy = Lazy::new(|| { - const NAME: &str = "libsql_server_replication_latency"; - describe_counter!( - NAME, - "Latency between the time a transaction was commited on the primary and the commit frame was received by the replica" - ); - register_histogram!(NAME) -}); +#[inline] +pub fn record_replication_latency(latency: Duration) { + histogram!("libsql_server_replication_latency", latency); +} pub static REPLICATION_LATENCY_OUT_OF_SYNC: Lazy = Lazy::new(|| { const NAME: &str = "libsql_server_replication_latency_out_of_sync"; describe_counter!( @@ -222,3 +198,251 @@ pub static TOKIO_RUNTIME_REMOTE_SCHEDULE_COUNT: Lazy = Lazy::new(|| { describe_gauge!(NAME, "tokio runtime remote_schedule_count"); register_counter!(NAME) }); + +/// Registers HELP text for every histogram. Must run after the global recorder is installed +/// (`describe_histogram!` is a no-op when no recorder is set). +pub(crate) fn describe_histograms() { + describe_histogram!( + "libsql_server_total_response_size_before_lock", + "total response size value before connection lock" + ); + describe_histogram!( + "libsql_server_namespace_load_latency", + "latency is us when loading a namespace" + ); + describe_histogram!( + "libsql_server_connection_create_time", + "time to create a connection" + ); + describe_histogram!( + "libsql_server_connection_alive_duration", + "duration for which a connection was kept alive" + ); + describe_histogram!( + "libsql_server_replication_latency", + "Latency between the time a transaction was commited on the primary and the commit frame was received by the replica" + ); + // Names recorded via `histogram!` elsewhere whose old `Lazy` statics were never forced + // (so they never had HELP text); harmless to describe them here. + describe_histogram!( + "libsql_server_statement_execution_time", + "time to execute a statement" + ); + describe_histogram!( + "libsql_server_statement_mem_used_bytes", + "memory used by a prepared statement" + ); + describe_histogram!( + "libsql_server_wal_checkpoint_time", + "time to checkpoint the WAL" + ); + describe_histogram!( + "libsql_server_returned_bytes", + "number of bytes of values returned to the client" + ); +} + +#[cfg(test)] +mod tests { + use std::sync::{Mutex, MutexGuard, PoisonError}; + use std::time::Duration; + + use metrics_exporter_prometheus::PrometheusHandle; + use once_cell::sync::Lazy; + + const IDLE: Duration = Duration::from_millis(50); + const PAST_IDLE: Duration = Duration::from_millis(150); + /// Samples recorded per production histogram after its eviction. + const BURST: u64 = 10; + + /// `metrics` 0.21 only has a process-global recorder, installable once; every test in this + /// module shares it. Mirrors the production configuration (same builder, same mask, global + /// labels) with a short idle timeout. + static HANDLE: Lazy = Lazy::new(|| { + let handle = crate::http::admin::prometheus_builder(IDLE) + .add_global_label("app", "test") + .add_global_label("version", "0.0.0-test") + .install_recorder() + .expect("another global metrics recorder is installed in this test binary"); + super::describe_histograms(); + handle + }); + + /// The tests in this module run serially. Two concurrent `render()` calls race inside the + /// exporter (one render can re-insert the recency entry of a key the other has just evicted, + /// recreating a distribution nobody visits again), which would make eviction assertions flaky. + /// Other unit tests only record, never render; production render concurrency is a separate + /// exporter limitation. + static SERIAL: Mutex<()> = Mutex::new(()); + + fn setup() -> (MutexGuard<'static, ()>, &'static PrometheusHandle) { + let guard = SERIAL.lock().unwrap_or_else(PoisonError::into_inner); + (guard, &*HANDLE) + } + + /// `_count{...} N` → Some(N), or None when the series is absent from the render. + fn count_of(rendered: &str, name: &str) -> Option { + let prefix = format!("{name}_count"); + rendered.lines().find_map(|line| { + // `{app="test",version="0.0.0-test"} 3` or ` 3` + let rest = line.strip_prefix(prefix.as_str())?; + if !(rest.starts_with('{') || rest.starts_with(' ')) { + return None; + } + rest.rsplit(' ').next()?.parse().ok() + }) + } + + fn has_help(rendered: &str, name: &str) -> bool { + let prefix = format!("# HELP {name} "); + rendered.lines().any(|l| l.starts_with(&prefix)) + } + + /// Renders repeatedly until `name` has been evicted (absent), or panics after `attempts`. + /// Production names may be touched by other unit tests running in the same process, which + /// would keep the key alive; the retry makes that a non-issue in practice. + fn evict(handle: &PrometheusHandle, name: &str, attempts: usize) { + for _ in 0..attempts { + std::thread::sleep(PAST_IDLE); + if count_of(&handle.render(), name).is_none() { + return; + } + } + panic!("{name} was never idle-evicted; is MetricKindMask::HISTOGRAM still in the mask?"); + } + + #[test] + fn cached_histogram_handle_is_orphaned_after_idle_eviction_but_macro_path_recovers() { + let (_serial, handle) = setup(); + const NAME: &str = "libsql_test_detached_histogram"; + + // The pre-fix pattern: a cached handle (what `Lazy` held). + let cached = metrics::register_histogram!(NAME); + cached.record(1.0); + assert_eq!(count_of(&handle.render(), NAME), Some(1)); + + // Two renders > idle_timeout apart with no update → key deleted from the registry. + evict(handle, NAME, 1); + + // Samples recorded through the orphaned handle are never rendered again: this is the leak + // (they accumulate in an AtomicBucket nobody drains). + cached.record(1.0); + cached.record(1.0); + assert_eq!( + count_of(&handle.render(), NAME), + None, + "detached handle must not reappear" + ); + + // The macro re-registers the key; only the new sample is visible, and it is drained. + metrics::histogram!(NAME, 1.0); + assert_eq!(count_of(&handle.render(), NAME), Some(1)); + metrics::histogram!(NAME, 1.0); + metrics::histogram!(NAME, 1.0); + assert_eq!(count_of(&handle.render(), NAME), Some(3)); + // A render with no new samples neither grows nor loses the distribution. + assert_eq!(count_of(&handle.render(), NAME), Some(3)); + } + + #[test] + fn production_histograms_survive_idle_eviction() { + let (_serial, handle) = setup(); + + // (name, record-fn) for every histogram that used to be recorded through a cached handle. + let cases: [(&str, fn()); 5] = [ + ("libsql_server_connection_create_time", || { + super::record_connection_create_time(Duration::from_millis(1)) + }), + ("libsql_server_connection_alive_duration", || { + super::record_connection_alive_duration(Duration::from_millis(1)) + }), + ("libsql_server_total_response_size_before_lock", || { + super::record_total_response_size_before_lock(42.0) + }), + ("libsql_server_namespace_load_latency", || { + super::record_namespace_load_latency(Duration::from_millis(1)) + }), + ("libsql_server_replication_latency", || { + super::record_replication_latency(Duration::from_millis(1)) + }), + ]; + + for (name, record) in cases { + record(); + let rendered = handle.render(); + assert!( + count_of(&rendered, name).unwrap_or(0) >= 1, + "{name} missing after record" + ); + assert!( + has_help(&rendered, name), + "{name} has no HELP line: describe_histograms() out of sync" + ); + + evict(handle, name, 20); + + // Re-registration: samples recorded after the eviction are rendered again. + for _ in 0..BURST { + record(); + } + let after = count_of(&handle.render(), name) + .unwrap_or_else(|| panic!("{name} did not reappear after idle eviction")); + assert!( + after >= BURST, + "{name}: expected >= {BURST} samples after re-registration, got {after}" + ); + + // Drain: an immediate render with no new samples from this test must not count the + // burst again. If the bucket were not cleared, re-reading it would add >= BURST. + // (Exact equality is not asserted because other unit tests in this process may + // record the same production metric concurrently; they cannot plausibly add BURST + // samples between two back-to-back renders. The exact-count drain check for the + // macro path is `cached_histogram_handle_is_orphaned_..._macro_path_recovers`.) + let idle = count_of(&handle.render(), name).unwrap(); + assert!( + idle - after < BURST, + "{name}: samples were re-counted on an idle render ({after} -> {idle}): not drained" + ); + + // Continued recording: one more sample still lands in the registered bucket. + record(); + let next = count_of(&handle.render(), name).unwrap(); + assert!( + next > idle, + "{name}: count did not advance ({idle} -> {next})" + ); + } + } + + #[test] + fn counters_and_gauges_are_still_pruned_when_idle() { + let (_serial, handle) = setup(); + metrics::counter!("libsql_test_pruned_counter", 7); + metrics::gauge!("libsql_test_pruned_gauge", 3.0); + let rendered = handle.render(); + assert!(rendered + .lines() + .any(|l| l.starts_with("libsql_test_pruned_counter{") && l.ends_with(" 7"))); + assert!(rendered + .lines() + .any(|l| l.starts_with("libsql_test_pruned_gauge{"))); + + std::thread::sleep(PAST_IDLE); + let rendered = handle.render(); + assert!( + !rendered.contains("libsql_test_pruned_counter"), + "idle counter must be evicted" + ); + assert!( + !rendered.contains("libsql_test_pruned_gauge"), + "idle gauge must be evicted" + ); + + // Macro-based counters/gauges come back on the next update (value restarts). + metrics::counter!("libsql_test_pruned_counter", 1); + assert!(handle + .render() + .lines() + .any(|l| l.starts_with("libsql_test_pruned_counter{") && l.ends_with(" 1"))); + } +} diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index 86e9438ccd..4e8cf8e920 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -15,7 +15,7 @@ use crate::broadcaster::BroadcastMsg; use crate::connection::config::DatabaseConfig; use crate::database::DatabaseKind; use crate::error::Error; -use crate::metrics::NAMESPACE_LOAD_LATENCY; +use crate::metrics::record_namespace_load_latency; use crate::namespace::{NamespaceBottomlessDbId, NamespaceBottomlessDbIdInit, NamespaceName}; use crate::stats::Stats; @@ -406,7 +406,7 @@ impl NamespaceStore { init.map_ok(|ns| Arc::new(RwLock::new(ns))), ) .await?; - NAMESPACE_LOAD_LATENCY.record(before_load.elapsed()); + record_namespace_load_latency(before_load.elapsed()); Ok(ns) } diff --git a/libsql-server/src/replication/replicator_client.rs b/libsql-server/src/replication/replicator_client.rs index fb8154824d..b98931b5f4 100644 --- a/libsql-server/src/replication/replicator_client.rs +++ b/libsql-server/src/replication/replicator_client.rs @@ -21,7 +21,7 @@ use tonic::{Code, Request, Status}; use crate::connection::config::DatabaseConfig; use crate::metrics::{ - REPLICATION_LATENCY, REPLICATION_LATENCY_CACHE_MISS, REPLICATION_LATENCY_OUT_OF_SYNC, + record_replication_latency, REPLICATION_LATENCY_CACHE_MISS, REPLICATION_LATENCY_OUT_OF_SYNC, }; use crate::namespace::meta_store::MetaStoreHandle; use crate::namespace::{NamespaceName, NamespaceStore}; @@ -221,7 +221,7 @@ impl ReplicatorClient for Client { Ok(lat) => { // we can record negative values if the clocks are out-of-sync. There is not // point in recording those values. - REPLICATION_LATENCY.record(lat); + record_replication_latency(lat); } Err(_) => { REPLICATION_LATENCY_OUT_OF_SYNC.increment(1);