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
24 changes: 17 additions & 7 deletions pgdog/src/backend/pool/lb/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@ impl Target {
self.role.role()
}

/// Set a known role, including when detection confirms the initial replica role.
/// Set a role and mark detection resolved for this target.
pub(super) fn set_role(&self, role: Role) -> bool {
let lb = self.role.set_role(role);
let pool = self.pool.set_role(role);
Expand All @@ -83,6 +83,13 @@ impl Target {
lb && pool
}

/// The role for this target is no longer known, e.g.,
/// we never received its LSN stats and we haven't detected
/// a primary either.
pub(super) fn unset_role(&self) {
self.role_detected.store(false, Ordering::Release)
}

pub(super) fn health(&self) -> &TargetHealth {
&self.pool.inner().health
}
Expand Down Expand Up @@ -216,11 +223,14 @@ impl LoadBalancer {
.for_each(|(_, target)| {
target.1.set_role(Role::Replica);
});
} else if targets.iter().all(|target| target.0.valid()) {
// All targets are replicas until we get a primary.
targets.iter().for_each(|target| {
target.1.set_role(Role::Replica);
});
} else {
for (stats, target) in &targets {
if stats.valid() {
target.set_role(Role::Replica);
} else if self.role_detection_enabled() {
target.unset_role();
}
}
}

self.elected_primary.send_replace(self.primary().cloned());
Expand Down Expand Up @@ -305,7 +315,7 @@ impl LoadBalancer {
.all(|target| target.pool.config().role_detection)
}

/// True once every target has a configured or detected role.
/// True if every target has a configured or detected role.
pub(crate) fn roles_detected(&self) -> bool {
self.targets
.iter()
Expand Down
57 changes: 55 additions & 2 deletions pgdog/src/backend/pool/lb/test/role_detection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,11 +20,22 @@ fn test_roles_detected_waits_for_all_replicas() {
assert!(lb.has_replicas(), "reads can start before role detection");
assert!(!lb.roles_detected());
assert!(!lb.redetect_roles());
assert!(!lb.roles_detected());
assert!(
!lb.roles_detected(),
"targets without LSN stats remain unknown"
);

set_lsn_stats(&lb.targets[0], true, 100);
assert!(!lb.redetect_roles());
assert!(!lb.roles_detected(), "the other target could be a primary");
assert!(
lb.targets[0].role_detected.load(Ordering::Acquire),
"valid replica stats resolve its role before other targets"
);
assert!(
!lb.targets[1].role_detected.load(Ordering::Acquire),
"target without LSN stats remains unknown"
);

set_lsn_stats(&lb.targets[1], true, 100);
assert!(!lb.redetect_roles());
Expand Down Expand Up @@ -55,6 +66,45 @@ fn test_roles_detected_when_primary_found_before_other_targets() {
assert!(lb.primary().is_some());
}

#[tokio::test(start_paused = true)]
async fn test_primary_demotion_waits_for_unknown_target_election() {
let lb = auto_pool(&["localhost", "127.0.0.1", "unknown-replica"]);
set_lsn_stats(&lb.targets[0], false, 500);
set_lsn_stats(&lb.targets[1], true, 500);
assert!(lb.redetect_roles());

let reloaded = auto_pool(&["localhost", "127.0.0.1", "unknown-replica"]);
lb.move_conns_to(&reloaded).expect("transfer pools");
assert!(
reloaded.primary().is_some(),
"elected primary survives reload"
);

set_lsn_stats(&reloaded.targets[0], true, 510);
assert!(!reloaded.redetect_roles());
assert!(reloaded.primary().is_none());
assert_eq!(reloaded.targets[0].role(), Role::Replica);
assert!(!reloaded.roles_detected());

let start = Instant::now();
assert!(matches!(
reloaded.get_primary(&Request::default()).await,
Err(Error::CheckoutTimeout)
));
assert!(start.elapsed() >= reloaded.checkout_timeout);

let (primary, ()) = tokio::join!(reloaded.wait_primary(), async {
sleep(Duration::from_millis(1)).await;
set_lsn_stats(&reloaded.targets[2], false, 520);
assert!(reloaded.redetect_roles());
});
assert_eq!(
primary.expect("new primary elected").addr().host,
"unknown-replica"
);
assert!(reloaded.roles_detected());
}

#[tokio::test]
async fn test_roles_detected_survives_reload_and_new_targets_remain_unknown() {
let old = auto_pool(&["localhost", "127.0.0.1"]);
Expand All @@ -74,7 +124,10 @@ async fn test_roles_detected_survives_reload_and_new_targets_remain_unknown() {
.expect("transfer existing pools");
assert!(!expanded.roles_detected(), "new target needs detection");
expanded.redetect_roles();
assert!(!expanded.roles_detected());
assert!(
!expanded.roles_detected(),
"new target still needs LSN stats"
);

let reloaded = auto_pool(&["localhost", "127.0.0.1", "new-replica"]);
expanded
Expand Down
16 changes: 16 additions & 0 deletions pgdog/src/frontend/router/parser/query/test/test_replica_only.rs
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,22 @@ fn test_replica_only_transactions_after_role_detection() {
shard.redetect_roles();
assert!(!cluster.read_only(), "one role is still unknown");

set_replica(&pools[0], false);
shard.redetect_roles();
assert!(!cluster.read_only(), "a primary was elected");

set_replica(&pools[0], true);
shard.redetect_roles();
assert!(
shard.has_primary(),
"primary routing remains pending while another target is unknown"
);
assert!(
!cluster.read_only(),
"demotion does not resolve the unknown role"
);
assert_transaction_route(&cluster, false);

set_replica(&pools[1], true);
shard.redetect_roles();
assert!(cluster.read_only(), "all configured servers are replicas");
Expand Down
Loading