From 5f9e451db29c5efbfb694ca5bba8028401a5a246 Mon Sep 17 00:00:00 2001 From: Mohammad Roshangara Date: Tue, 29 Sep 2026 01:39:18 +0200 Subject: [PATCH] fix(lb): a write waiting for the automatic primary follows the election With `role = "auto"`, get_primary() checked out from the pool elected when the write arrived and waited there for up to checkout_timeout. If that primary went down, the write stayed in its queue after the shard monitor had elected another server, and failed with a checkout timeout, or got a new connection to the old primary if it came back as a replica ("cannot execute ... in a read-only transaction"). A write now follows the election for up to checkout_timeout: while no server is the primary it waits, and when the election changes it drops its pending checkout and checks out from the new primary. The election channel is written only when the elected pool changes, so a monitor tick that confirms the same primary doesn't wake waiting writes, and a configuration reload publishes the primary carried over to the new pools. Co-Authored-By: Claude Opus 5.5 --- pgdog/src/backend/pool/lb/mod.rs | 78 ++++++++++++----- pgdog/src/backend/pool/lb/test/mod.rs | 119 +++++++++++++++++++++++++- 2 files changed, 170 insertions(+), 27 deletions(-) diff --git a/pgdog/src/backend/pool/lb/mod.rs b/pgdog/src/backend/pool/lb/mod.rs index 819a068a8..8f0a8fed8 100644 --- a/pgdog/src/backend/pool/lb/mod.rs +++ b/pgdog/src/backend/pool/lb/mod.rs @@ -2,6 +2,7 @@ use std::{ cmp::Reverse, + future::pending, sync::{ Arc, atomic::{AtomicBool, AtomicI64, AtomicUsize, Ordering}, @@ -10,6 +11,7 @@ use std::{ }; use rand::seq::SliceRandom; +use tokio::select; use tokio::sync::watch; use tokio_util::sync::CancellationToken; use tracing::warn; @@ -199,8 +201,6 @@ impl LoadBalancer { .iter() .position(|target| !target.0.replica && target.0.valid()); - self.elected_primary.send_replace(None); - if let Some(primary) = primary { promoted = targets[primary].1.set_role(Role::Primary); @@ -223,11 +223,25 @@ impl LoadBalancer { }); } - self.elected_primary.send_replace(self.primary().cloned()); + self.publish_primary(); promoted } + /// Tell writes waiting in [`Self::get_primary`] which pool is the primary. + /// Only a change wakes them up. + fn publish_primary(&self) { + let elected = self.primary().cloned(); + + self.elected_primary.send_if_modified(|current| { + let changed = current.as_ref().map(Pool::id) != elected.as_ref().map(Pool::id); + if changed { + *current = elected; + } + changed + }); + } + /// Launch replica pools and start the monitor. pub(crate) fn launch(&self) { self.targets.iter().for_each(|target| target.pool.launch()); @@ -285,6 +299,7 @@ impl LoadBalancer { } } destination.require_healthcheck_for_new_targets(&self.targets); + destination.publish_primary(); Ok(moved) } @@ -332,34 +347,51 @@ impl LoadBalancer { result } - /// Wait until automatic role detection elects a primary. + /// Check out a connection from the primary. /// - /// Fails with [`Error::CheckoutTimeout`] if no election happens in time. - async fn wait_primary(&self) -> Result { - let mut receiver = self.elected_primary.subscribe(); + /// Static replica-only configurations fail immediately with + /// [`Error::NoPrimary`]. With automatic roles, the caller follows the + /// election for up to `checkout_timeout`: it waits while there is no + /// primary, and a checkout waiting on a primary that loses the election + /// moves to the new one. + pub(super) async fn get_primary(&self, request: &Request) -> Result { + if !self.role_detection_enabled() { + return match self.primary() { + Some(pool) => pool.get(request).await, + None => Err(Error::NoPrimary), + }; + } - safe_timeout(self.checkout_timeout, receiver.wait_for(|p| p.is_some())) + safe_timeout(self.checkout_timeout, self.follow_election(request)) .await .map_err(|_| Error::CheckoutTimeout)? - .ok() - .and_then(|elected| elected.as_ref().cloned()) - .ok_or(Error::NoPrimary) } - /// Check out a connection from the primary. - /// - /// In automatic mode, the caller waits for an election. Static - /// replica-only configurations fail immediately with [`Error::NoPrimary`]. - pub(super) async fn get_primary(&self, request: &Request) -> Result { - if let Some(pool) = self.primary() { - return pool.get(request).await; - } + /// Check out a connection from the elected primary, starting over on + /// the new primary whenever the election changes. + async fn follow_election(&self, request: &Request) -> Result { + let mut elections = self.elected_primary.subscribe(); - if !self.role_detection_enabled() { - return Err(Error::NoPrimary); - } + loop { + let elected = elections.borrow_and_update().clone(); + let elected_id = elected.as_ref().map(Pool::id); - self.wait_primary().await?.get(request).await + let checkout = async { + match &elected { + Some(pool) => pool.get(request).await, + None => pending().await, + } + }; + let new_election = + elections.wait_for(|primary| primary.as_ref().map(Pool::id) != elected_id); + + select! { + result = checkout => return result, + changed = new_election => { + changed.map_err(|_| Error::NoPrimary)?; + } + } + } } async fn get_internal(&self, request: &Request) -> Result { diff --git a/pgdog/src/backend/pool/lb/test/mod.rs b/pgdog/src/backend/pool/lb/test/mod.rs index 473857e5a..2bfa8fe13 100644 --- a/pgdog/src/backend/pool/lb/test/mod.rs +++ b/pgdog/src/backend/pool/lb/test/mod.rs @@ -1592,19 +1592,130 @@ async fn test_election_channel_no_race_condition() { // make sure there is no primary assert!(lb.primary().is_none()); - // initialize the primary before the wait_primary is called + // initialize the primary before get_primary is called set_lsn_stats(&lb.targets[0], false, 100); assert!(lb.redetect_roles()); - let elected = timeout(Duration::ZERO, lb.wait_primary()) + let primary = timeout(Duration::from_secs(1), lb.get_primary(&Request::default())) .await - .expect("wait_primary should resolve") + .expect("an election that already happened is not waited for") .expect("primary should be elected"); - assert_eq!(elected.addr().host, "127.0.0.1"); + assert_eq!(primary.pool.addr().host, "127.0.0.1"); + drop(primary); lb.shutdown(); } +/// Auto targets that wait up to 2s for a connection. Port 1 refuses +/// connections: a primary that went down. +fn auto_lb_with_checkout_timeout(ports: &[u16]) -> LoadBalancer { + let configs: Vec<_> = ports + .iter() + .map(|port| { + let mut config = create_auto_test_pool_config("127.0.0.1", *port); + config.config.checkout_timeout = Duration::from_secs(2); + config.config.connect_timeout = Duration::from_millis(100); + config + }) + .collect(); + + LoadBalancer::new( + &None, + &configs, + LoadBalancingStrategy::Random, + ReadWriteSplit::IncludePrimary, + Default::default(), + ) +} + +#[tokio::test] +async fn test_waiting_write_moves_to_new_primary() { + let lb = auto_lb_with_checkout_timeout(&[1, 5432]); + lb.launch(); + + // The elected primary is down: a write waits on its pool. + set_lsn_stats(&lb.targets[0], false, 500); + set_lsn_stats(&lb.targets[1], true, 500); + assert!(lb.redetect_roles()); + + let started = Instant::now(); + let write = { + let lb = lb.clone(); + tokio::spawn(async move { lb.get_primary(&Request::default()).await }) + }; + + sleep(Duration::from_millis(300)).await; + assert!(!write.is_finished(), "the write waits for the primary"); + + // Failover: the replica is promoted. + set_lsn_stats(&lb.targets[1], false, 600); + set_lsn_stats(&lb.targets[0], true, 500); + assert!(lb.redetect_roles()); + + let primary = timeout(Duration::from_secs(1), write) + .await + .expect("the write moves to the new primary at once") + .unwrap() + .expect("a connection to the new primary"); + assert_eq!(primary.pool.addr().port, 5432); + assert!(started.elapsed() < Duration::from_secs(2)); + drop(primary); + + lb.shutdown(); +} + +#[tokio::test] +async fn test_unchanged_election_does_not_wake_waiting_writes() { + let lb = auto_lb_with_checkout_timeout(&[5432, 1]); + set_lsn_stats(&lb.targets[0], false, 500); + set_lsn_stats(&lb.targets[1], true, 500); + + let mut elections = lb.elected_primary.subscribe(); + assert!(lb.redetect_roles()); + assert!( + elections.has_changed().unwrap(), + "a new primary is published" + ); + elections.borrow_and_update(); + + // The LSN check runs again and nothing changed. Waiting writes keep + // their place in the primary's queue. + set_lsn_stats(&lb.targets[0], false, 700); + assert!(!lb.redetect_roles()); + assert!(!elections.has_changed().unwrap()); + + // The primary is gone: that is a change. + set_lsn_stats(&lb.targets[0], true, 700); + lb.redetect_roles(); + assert!(elections.has_changed().unwrap()); + assert!(elections.borrow_and_update().is_none()); +} + +#[tokio::test] +async fn test_write_after_reload_goes_to_elected_primary() { + let old = auto_lb_with_checkout_timeout(&[5432, 1]); + set_lsn_stats(&old.targets[0], false, 500); + set_lsn_stats(&old.targets[1], true, 500); + assert!(old.redetect_roles()); + + let new = auto_lb_with_checkout_timeout(&[5432, 1]); + old.move_conns_to(&new).unwrap(); + new.launch(); + + // The new load balancer knows the primary before its first election. + let primary = timeout( + Duration::from_millis(500), + new.get_primary(&Request::default()), + ) + .await + .expect("no wait for an election after a reload") + .expect("a connection to the primary"); + assert_eq!(primary.pool.addr().port, 5432); + drop(primary); + + new.shutdown(); +} + #[tokio::test] async fn test_static_replica_only_does_not_wait_for_primary() { let config = create_test_pool_config("127.0.0.1", 5432);