From 3798b7a622909efafa47451af2afbaa3ba5fb6ce Mon Sep 17 00:00:00 2001 From: Rider Wang Date: Tue, 1 Sep 2026 12:50:19 -0400 Subject: [PATCH 1/3] Add login path observability metrics When the login path saturates, TCP accepts succeed but clients time out mid-startup and PgDog reports nothing: no error, no log line, no metric. Make that failure mode observable on the OpenMetrics endpoint: - logins_in_flight: connections accepted but not yet ReadyForQuery, the leading indicator of login-path saturation - logins: completed client logins - logins_abandoned: clients that started the Postgres handshake but disconnected (or timed out) before login completed; connections that close without sending a single startup byte (load balancer TCP health checks) are excluded - login_duration: accept-to-ReadyForQuery histogram (ms) for completed logins A LoginTimer guard is created at accept, engaged on the first startup message, marked successful when login completes, and disarmed when PgDog ends the login with an explicit response (auth failure, shutdown, pool down) or the connection is a cancel request. Dropping it engaged but unfinished counts an abandoned login, so every exit path is covered without touching each error branch. The scalar metrics also flow through the OTLP exporter; the histogram is rendered only on the OpenMetrics text endpoint via a new overridable OpenMetric::render_measurements (the OTLP renderer currently only speaks gauge/sum). Closes #1470 Co-Authored-By: Claude Fable 5 --- pgdog/src/frontend/client/mod.rs | 11 +- pgdog/src/frontend/client/test/mod.rs | 1 + pgdog/src/frontend/client/test/test_client.rs | 13 +- pgdog/src/frontend/listener.rs | 16 +- pgdog/src/stats/http_server.rs | 11 +- pgdog/src/stats/logins.rs | 362 ++++++++++++++++++ pgdog/src/stats/mod.rs | 2 + pgdog/src/stats/open_metric.rs | 20 +- pgdog/src/stats/otel_exporter.rs | 8 +- 9 files changed, 433 insertions(+), 11 deletions(-) create mode 100644 pgdog/src/stats/logins.rs 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, + "{}{}_bucket{{le=\"{}\"}} {}", + prefix, name, bound, cumulative + )?; + } + let total = cumulative + self.overflow; + writeln!(f, "{}{}_bucket{{le=\"+Inf\"}} {}", prefix, name, total)?; + writeln!( + f, + "{}{}_sum {:.3}", + prefix, + name, + self.sum_micros as f64 / 1_000.0 + )?; + writeln!(f, "{}{}_count {}", prefix, name, 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 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(); From 287d9f6e31077c464b5abf9e37c9660c6438280d Mon Sep 17 00:00:00 2001 From: Rider Wang Date: Tue, 1 Sep 2026 13:29:03 -0400 Subject: [PATCH 2/3] Cover global login stats accessors and success-after-disarm guard Co-Authored-By: Claude Fable 5 --- pgdog/src/stats/logins.rs | 36 ++++++++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/pgdog/src/stats/logins.rs b/pgdog/src/stats/logins.rs index 792a08aa1..98b810b4a 100644 --- a/pgdog/src/stats/logins.rs +++ b/pgdog/src/stats/logins.rs @@ -349,6 +349,42 @@ mod test { 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(); From 42da1fd37aceeb7e26b2703aa3fcd171f784d76f Mon Sep 17 00:00:00 2001 From: Rider Wang Date: Tue, 1 Sep 2026 14:36:27 -0400 Subject: [PATCH 3/3] Collapse histogram writeln macros to single lines The multiline writeln invocations left llvm-cov region artifacts on their continuation lines; captured-identifier format strings keep each write on one line and read cleaner. Co-Authored-By: Claude Fable 5 --- pgdog/src/stats/logins.rs | 19 +++++-------------- 1 file changed, 5 insertions(+), 14 deletions(-) diff --git a/pgdog/src/stats/logins.rs b/pgdog/src/stats/logins.rs index 98b810b4a..d77363927 100644 --- a/pgdog/src/stats/logins.rs +++ b/pgdog/src/stats/logins.rs @@ -248,22 +248,13 @@ impl OpenMetric for LoginDuration { let mut cumulative = 0u64; for (bound, count) in BUCKETS_MS.iter().zip(&self.buckets) { cumulative += count; - writeln!( - f, - "{}{}_bucket{{le=\"{}\"}} {}", - prefix, name, bound, cumulative - )?; + writeln!(f, "{prefix}{name}_bucket{{le=\"{bound}\"}} {cumulative}")?; } let total = cumulative + self.overflow; - writeln!(f, "{}{}_bucket{{le=\"+Inf\"}} {}", prefix, name, total)?; - writeln!( - f, - "{}{}_sum {:.3}", - prefix, - name, - self.sum_micros as f64 / 1_000.0 - )?; - writeln!(f, "{}{}_count {}", prefix, name, total)?; + 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(()) } }