Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
78 changes: 55 additions & 23 deletions pgdog/src/backend/pool/lb/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

use std::{
cmp::Reverse,
future::pending,
sync::{
Arc,
atomic::{AtomicBool, AtomicI64, AtomicUsize, Ordering},
Expand All @@ -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;
Expand Down Expand Up @@ -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);

Expand All @@ -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());
Expand Down Expand Up @@ -285,6 +299,7 @@ impl LoadBalancer {
}
}
destination.require_healthcheck_for_new_targets(&self.targets);
destination.publish_primary();

Ok(moved)
}
Expand Down Expand Up @@ -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<Pool, Error> {
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<Guard, Error> {
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<Guard, Error> {
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<Guard, Error> {
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<Guard, Error> {
Expand Down
119 changes: 115 additions & 4 deletions pgdog/src/backend/pool/lb/test/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down