diff --git a/pgdog/src/frontend/client/mod.rs b/pgdog/src/frontend/client/mod.rs index 4f59aba9f..c836f2a63 100644 --- a/pgdog/src/frontend/client/mod.rs +++ b/pgdog/src/frontend/client/mod.rs @@ -30,6 +30,7 @@ use crate::net::messages::{ }; use crate::net::{MessageBuffer, ProtocolMessage, Stream, parameter::Parameters}; use crate::state::State; +use crate::stats::logins::LoginTimer; use crate::stats::memory::MemoryUsage; use crate::util::{safe_timeout, user_database_from_params}; @@ -138,6 +139,7 @@ impl Client { addr: SocketAddr, config: Arc, protocol_version: ProtocolVersion, + mut login_timer: LoginTimer, ) -> Result<(), Error> { let login_timeout = Duration::from_millis(config.config.general.client_login_timeout); @@ -148,6 +150,7 @@ impl Client { .await { Ok(Ok(Some(mut client))) => { + login_timer.success(); if client.admin { // Admin clients are not waited on during shutdown. spawn(async move { @@ -160,10 +163,16 @@ impl Client { Ok(()) } Err(_) => { + // Login timeouts are logged; the drop still counts the + // connection as an abandoned login. error!("client login timeout [{}]", addr); Ok(()) } - Ok(Ok(None)) => Ok(()), + Ok(Ok(None)) => { + // Login was rejected with an explicit response. + login_timer.disarm(); + Ok(()) + } Ok(Err(err)) => Err(err), } } diff --git a/pgdog/src/frontend/client/test/mod.rs b/pgdog/src/frontend/client/test/mod.rs index 5b9906827..b650e3eb6 100644 --- a/pgdog/src/frontend/client/test/mod.rs +++ b/pgdog/src/frontend/client/test/mod.rs @@ -487,6 +487,7 @@ async fn test_client_login_timeout() { addr, crate::config::config(), ProtocolVersion::V3_0, + crate::stats::logins::LoginTimer::new(), ) .await }); diff --git a/pgdog/src/frontend/client/test/test_client.rs b/pgdog/src/frontend/client/test/test_client.rs index 11fb6214e..a71586509 100644 --- a/pgdog/src/frontend/client/test/test_client.rs +++ b/pgdog/src/frontend/client/test/test_client.rs @@ -341,9 +341,16 @@ impl SpawnedClient { let handle = tokio::spawn(async move { let (stream, addr) = listener.accept().await.unwrap(); let stream = Stream::plain(stream, 4096); - Client::spawn(stream, params, addr, config(), ProtocolVersion::V3_0) - .await - .unwrap(); + Client::spawn( + stream, + params, + addr, + config(), + ProtocolVersion::V3_0, + crate::stats::logins::LoginTimer::new(), + ) + .await + .unwrap(); }); let conn = TcpStream::connect(format!("127.0.0.1:{}", port)) diff --git a/pgdog/src/frontend/listener.rs b/pgdog/src/frontend/listener.rs index bf3935fac..1a93d9344 100644 --- a/pgdog/src/frontend/listener.rs +++ b/pgdog/src/frontend/listener.rs @@ -10,6 +10,7 @@ use crate::net::messages::{FrontendPid, NegotiateProtocolVersion, Startup, hello use crate::net::tls::{acceptor, peer_certificate_present, peer_identity}; use crate::net::{self, Stream, tweak}; use crate::sighup::Sighup; +use crate::stats::logins::LoginTimer; use tokio::net::{TcpListener, TcpSocket, TcpStream, lookup_host}; use tokio::signal::ctrl_c; use tokio::{select, spawn}; @@ -204,6 +205,7 @@ impl Listener { let mut stream = Stream::plain(stream, config.config.memory.net_buffer); let tls = acceptor(); + let mut login_timer = LoginTimer::new(); loop { let startup = match Startup::from_stream(&mut stream).await { @@ -222,6 +224,7 @@ impl Listener { match startup { Startup::Ssl => { + login_timer.engage(); if let Some(tls) = tls.as_ref() { stream.send_flush(&SslReply::Yes).await?; let plain = stream.take()?; @@ -248,6 +251,7 @@ impl Listener { } Startup::GssEnc => { + login_timer.engage(); // GSS encryption is not yet supported; reject and wait for a normal startup. stream.send_flush(&SslReply::No).await?; } @@ -257,6 +261,7 @@ impl Listener { params, unrecognized_options, } => { + login_timer.engage(); let negotiated = version .negotiated() .ok_or_else(|| net::Error::UnsupportedStartup(version.as_i32()))?; @@ -273,11 +278,20 @@ impl Listener { .await?; } - Box::pin(Client::spawn(stream, params, addr, config, negotiated)).await?; + Box::pin(Client::spawn( + stream, + params, + addr, + config, + negotiated, + login_timer, + )) + .await?; break; } Startup::Cancel { ref id } => { + login_timer.disarm(); if comms().verify_cancel(id) { let _ = databases().cancel(FrontendPid::from(id)).await; } diff --git a/pgdog/src/stats/http_server.rs b/pgdog/src/stats/http_server.rs index 88cdc4617..fdcbff8b3 100644 --- a/pgdog/src/stats/http_server.rs +++ b/pgdog/src/stats/http_server.rs @@ -11,7 +11,8 @@ use tokio::select; use tracing::{info, warn}; use super::{ - Clients, ClientsLocked, Listeners, LookupMetrics, MirrorStatsMetrics, Pools, QueryCache, TwoPc, + Clients, ClientsLocked, Listeners, Logins, LookupMetrics, MirrorStatsMetrics, Pools, + QueryCache, TwoPc, }; use crate::tasks; @@ -41,6 +42,12 @@ async fn metrics(_: Request) -> Result = Logins::load() + .into_iter() + .chain(std::iter::once(Logins::histogram())) + .map(|m| m.to_string()) + .collect(); + let logins = logins.join("\n"); let metrics_data = clients.to_string() + "\n" + &clients_locked.to_string() @@ -55,6 +62,8 @@ async fn metrics(_: Request) -> Result Self { + Self { + in_flight: AtomicI64::new(0), + completed: AtomicU64::new(0), + abandoned: AtomicU64::new(0), + duration_sum_micros: AtomicU64::new(0), + duration_buckets: [const { AtomicU64::new(0) }; BUCKETS_MS.len()], + duration_overflow: AtomicU64::new(0), + } + } + + fn record_duration(&self, micros: u64) { + self.duration_sum_micros + .fetch_add(micros, Ordering::Relaxed); + let millis = micros / 1_000; + match BUCKETS_MS.iter().position(|bound| millis <= *bound) { + Some(bucket) => &self.duration_buckets[bucket], + None => &self.duration_overflow, + } + .fetch_add(1, Ordering::Relaxed); + } +} + +static LOGINS: LoginStats = LoginStats::new(); + +/// Global login stats. +pub(crate) fn logins() -> &'static LoginStats { + &LOGINS +} + +/// Tracks one connection through the login path. +/// +/// Create it at accept. Call [`LoginTimer::engage`] once the connection sends +/// startup bytes (i.e., it's a Postgres client, not a TCP health check), +/// [`LoginTimer::success`] when the client reaches ReadyForQuery, and +/// [`LoginTimer::disarm`] when PgDog itself ends the login with an explicit +/// response (auth failure, shutdown, pool down) or the connection turns out +/// not to be a login (cancel request). Dropping the timer engaged but +/// unfinished counts the login as abandoned. +pub(crate) struct LoginTimer { + stats: &'static LoginStats, + start: Instant, + engaged: bool, + finished: bool, +} + +impl LoginTimer { + pub(crate) fn new() -> Self { + Self::with_stats(logins()) + } + + fn with_stats(stats: &'static LoginStats) -> Self { + stats.in_flight.fetch_add(1, Ordering::Relaxed); + Self { + stats, + start: Instant::now(), + engaged: false, + finished: false, + } + } + + /// The connection sent startup bytes: it's a real client, not a probe. + pub(crate) fn engage(&mut self) { + self.engaged = true; + } + + /// Login ended with an explicit response (or wasn't a login at all); + /// don't count it as abandoned. + pub(crate) fn disarm(&mut self) { + self.finished = true; + } + + /// The client reached ReadyForQuery. + pub(crate) fn success(&mut self) { + if self.finished { + return; + } + self.finished = true; + self.stats.completed.fetch_add(1, Ordering::Relaxed); + self.stats + .record_duration(self.start.elapsed().as_micros() as u64); + } +} + +impl Drop for LoginTimer { + fn drop(&mut self) { + self.stats.in_flight.fetch_sub(1, Ordering::Relaxed); + if self.engaged && !self.finished { + self.stats.abandoned.fetch_add(1, Ordering::Relaxed); + } + } +} + +/// Login metrics for the OpenMetrics endpoint. +pub(crate) struct Logins; + +impl Logins { + /// Scalar login metrics (gauge + counters). + pub(crate) fn load() -> Vec { + Self::load_from(logins()) + } + + fn load_from(stats: &LoginStats) -> Vec { + vec![ + Metric::new(LoginMetric { + name: "logins_in_flight".into(), + metric_type: "gauge".into(), + help: "Connections accepted but not yet ReadyForQuery.".into(), + value: stats.in_flight.load(Ordering::Relaxed).into(), + }), + Metric::new(LoginMetric { + name: "logins".into(), + metric_type: "counter".into(), + help: "Total number of completed client logins.".into(), + value: stats.completed.load(Ordering::Relaxed).into(), + }), + Metric::new(LoginMetric { + name: "logins_abandoned".into(), + metric_type: "counter".into(), + help: "Clients that started the handshake but disconnected before login completed." + .into(), + value: stats.abandoned.load(Ordering::Relaxed).into(), + }), + ] + } + + /// Login duration histogram (OpenMetrics text format only). + pub(crate) fn histogram() -> Metric { + Self::histogram_from(logins()) + } + + fn histogram_from(stats: &LoginStats) -> Metric { + let buckets = stats + .duration_buckets + .iter() + .map(|bucket| bucket.load(Ordering::Relaxed)) + .collect(); + Metric::new(LoginDuration { + buckets, + overflow: stats.duration_overflow.load(Ordering::Relaxed), + sum_micros: stats.duration_sum_micros.load(Ordering::Relaxed), + }) + } +} + +struct LoginMetric { + name: String, + metric_type: String, + help: String, + value: MeasurementType, +} + +impl OpenMetric for LoginMetric { + fn name(&self) -> String { + self.name.clone() + } + + fn measurements(&self) -> Vec { + vec![Measurement { + labels: vec![], + measurement: self.value.clone(), + }] + } + + fn metric_type(&self) -> String { + self.metric_type.clone() + } + + fn help(&self) -> Option { + Some(self.help.clone()) + } +} + +struct LoginDuration { + buckets: Vec, + overflow: u64, + sum_micros: u64, +} + +impl OpenMetric for LoginDuration { + fn name(&self) -> String { + "login_duration".into() + } + + fn metric_type(&self) -> String { + "histogram".into() + } + + fn unit(&self) -> Option { + Some("milliseconds".into()) + } + + fn help(&self) -> Option { + Some("Accept-to-ReadyForQuery latency of completed client logins.".into()) + } + + fn measurements(&self) -> Vec { + // Unused: histograms render through `render_measurements`. + vec![] + } + + fn render_measurements( + &self, + f: &mut std::fmt::Formatter<'_>, + prefix: &str, + name: &str, + ) -> std::fmt::Result { + let mut cumulative = 0u64; + for (bound, count) in BUCKETS_MS.iter().zip(&self.buckets) { + cumulative += count; + writeln!(f, "{prefix}{name}_bucket{{le=\"{bound}\"}} {cumulative}")?; + } + let total = cumulative + self.overflow; + let sum_ms = self.sum_micros as f64 / 1_000.0; + writeln!(f, "{prefix}{name}_bucket{{le=\"+Inf\"}} {total}")?; + writeln!(f, "{prefix}{name}_sum {sum_ms:.3}")?; + writeln!(f, "{prefix}{name}_count {total}")?; + Ok(()) + } +} + +#[cfg(test)] +mod test { + use super::*; + + fn test_stats() -> &'static LoginStats { + Box::leak(Box::new(LoginStats::new())) + } + + #[test] + fn success_records_completion_and_duration() { + let stats = test_stats(); + let mut timer = LoginTimer::with_stats(stats); + assert_eq!(stats.in_flight.load(Ordering::Relaxed), 1); + timer.engage(); + timer.success(); + drop(timer); + + assert_eq!(stats.in_flight.load(Ordering::Relaxed), 0); + assert_eq!(stats.completed.load(Ordering::Relaxed), 1); + assert_eq!(stats.abandoned.load(Ordering::Relaxed), 0); + let observed: u64 = stats + .duration_buckets + .iter() + .map(|bucket| bucket.load(Ordering::Relaxed)) + .sum::() + + stats.duration_overflow.load(Ordering::Relaxed); + assert_eq!(observed, 1); + } + + #[test] + fn engaged_drop_counts_abandoned() { + let stats = test_stats(); + let mut timer = LoginTimer::with_stats(stats); + timer.engage(); + drop(timer); + + assert_eq!(stats.in_flight.load(Ordering::Relaxed), 0); + assert_eq!(stats.abandoned.load(Ordering::Relaxed), 1); + assert_eq!(stats.completed.load(Ordering::Relaxed), 0); + } + + #[test] + fn health_check_drop_is_not_abandoned() { + let stats = test_stats(); + let timer = LoginTimer::with_stats(stats); + drop(timer); + + assert_eq!(stats.in_flight.load(Ordering::Relaxed), 0); + assert_eq!(stats.abandoned.load(Ordering::Relaxed), 0); + } + + #[test] + fn disarmed_drop_is_not_abandoned() { + let stats = test_stats(); + let mut timer = LoginTimer::with_stats(stats); + timer.engage(); + timer.disarm(); + drop(timer); + + assert_eq!(stats.abandoned.load(Ordering::Relaxed), 0); + assert_eq!(stats.completed.load(Ordering::Relaxed), 0); + } + + #[test] + fn histogram_renders_cumulative_buckets() { + let stats = test_stats(); + stats.record_duration(1_500); // 1.5ms -> le=2 + stats.record_duration(1_500); + stats.record_duration(600_000); // 600ms -> le=1000 + stats.record_duration(60_000_000); // 60s -> overflow + + let rendered = Logins::histogram_from(stats).to_string(); + assert!(rendered.contains("# TYPE login_duration histogram")); + assert!(rendered.contains("login_duration_bucket{le=\"2\"} 2")); + assert!(rendered.contains("login_duration_bucket{le=\"1000\"} 3")); + assert!(rendered.contains("login_duration_bucket{le=\"30000\"} 3")); + assert!(rendered.contains("login_duration_bucket{le=\"+Inf\"} 4")); + assert!(rendered.contains("login_duration_count 4")); + assert!(rendered.contains("login_duration_sum 60603.000")); + } + + #[test] + fn success_after_disarm_records_nothing() { + let stats = test_stats(); + let mut timer = LoginTimer::with_stats(stats); + timer.engage(); + timer.disarm(); + timer.success(); + drop(timer); + + assert_eq!(stats.completed.load(Ordering::Relaxed), 0); + assert_eq!(stats.abandoned.load(Ordering::Relaxed), 0); + } + + #[test] + fn global_stats_flow_end_to_end() { + let completed_before = logins().completed.load(Ordering::Relaxed); + + let mut timer = LoginTimer::new(); + timer.engage(); + timer.success(); + drop(timer); + + // Other tests share the global stats, so assert deltas only. + assert!(logins().completed.load(Ordering::Relaxed) > completed_before); + + let names: Vec<_> = Logins::load().iter().map(|metric| metric.name()).collect(); + assert_eq!(names, ["logins_in_flight", "logins", "logins_abandoned"]); + + let histogram = Logins::histogram(); + assert_eq!(histogram.metric_type(), "histogram"); + // Histograms render through render_measurements; the trait method + // returns nothing. + assert!(histogram.measurements().is_empty()); + assert!(histogram.to_string().contains("login_duration_count")); + } + + #[test] + fn scalar_metrics_have_expected_names_and_types() { + let stats = test_stats(); + let metrics = Logins::load_from(stats); + let names: Vec<_> = metrics.iter().map(|metric| metric.name()).collect(); + assert_eq!(names, ["logins_in_flight", "logins", "logins_abandoned"]); + assert_eq!(metrics[0].metric_type(), "gauge"); + assert_eq!(metrics[1].metric_type(), "counter"); + assert_eq!(metrics[2].metric_type(), "counter"); + } +} diff --git a/pgdog/src/stats/mod.rs b/pgdog/src/stats/mod.rs index 8d1d6f2e5..cfca3b4c9 100644 --- a/pgdog/src/stats/mod.rs +++ b/pgdog/src/stats/mod.rs @@ -11,6 +11,7 @@ pub(crate) mod pools; pub(crate) use open_metric::*; pub(crate) mod listeners; pub(crate) mod logger; +pub(crate) mod logins; pub(crate) mod memory; pub(crate) mod query_cache; pub(crate) mod two_pc; @@ -19,6 +20,7 @@ pub(crate) use clients::Clients; pub(crate) use clients_locked::ClientsLocked; pub(crate) use listeners::Listeners; pub(crate) use logger::Logger as StatsLogger; +pub(crate) use logins::Logins; pub(crate) use lookup::LookupMetrics; pub(crate) use mirror_stats::MirrorStatsMetrics; pub(crate) use pools::Pools; diff --git a/pgdog/src/stats/open_metric.rs b/pgdog/src/stats/open_metric.rs index 161badab2..4afb5b033 100644 --- a/pgdog/src/stats/open_metric.rs +++ b/pgdog/src/stats/open_metric.rs @@ -19,6 +19,21 @@ pub(crate) trait OpenMetric: Send + Sync { fn help(&self) -> Option { None } + + /// Render measurement lines. The default renders one line per + /// measurement; multi-series families (histograms) override this to + /// control the sample names (`_bucket`, `_sum`, `_count`). + fn render_measurements( + &self, + f: &mut std::fmt::Formatter<'_>, + prefix: &str, + name: &str, + ) -> std::fmt::Result { + for measurement in self.measurements() { + writeln!(f, "{}{}", prefix, measurement.render(name))?; + } + Ok(()) + } } #[derive(Debug, Clone)] @@ -127,10 +142,7 @@ impl std::fmt::Display for Metric { writeln!(f, "# HELP {}{} {}", prefix, name, help)?; } - for measurement in self.measurements() { - writeln!(f, "{}{}", prefix, measurement.render(&name))?; - } - Ok(()) + self.render_measurements(f, prefix, &name) } } diff --git a/pgdog/src/stats/otel_exporter.rs b/pgdog/src/stats/otel_exporter.rs index 4e1ea9a88..0bb2188d1 100644 --- a/pgdog/src/stats/otel_exporter.rs +++ b/pgdog/src/stats/otel_exporter.rs @@ -8,7 +8,9 @@ use std::time::Duration; use tracing::{info, warn}; use super::otel; -use super::{Clients, ClientsLocked, Listeners, MirrorStatsMetrics, Pools, QueryCache, TwoPc}; +use super::{ + Clients, ClientsLocked, Listeners, Logins, MirrorStatsMetrics, Pools, QueryCache, TwoPc, +}; use crate::util::safe_sleep; use crate::{config::config, tasks}; @@ -47,12 +49,16 @@ pub(crate) async fn run() { let listeners = Listeners::load(); let query_cache = QueryCache::load().metrics(); let two_pc = TwoPc::load(); + // Scalar login metrics only: the OTLP renderer speaks gauge/sum, so + // the login_duration histogram stays on the OpenMetrics endpoint. + let logins = Logins::load(); let mut all: Vec<&super::Metric> = vec![&clients, &clients_locked, &two_pc]; all.extend(pools.iter()); all.extend(mirror.iter()); all.extend(listeners.iter()); all.extend(query_cache.iter()); + all.extend(logins.iter()); let now = otel::now_nanos();