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);