From 63663b8d9cc522b55eb1e5c62b77afa7bf5bbc5d Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 6 Oct 2026 15:10:54 -0700 Subject: [PATCH 1/2] fix: demote broken primary to replica without waiting for all replicas to report --- pgdog/src/backend/pool/lb/mod.rs | 11 +++++++---- .../src/backend/pool/lb/test/role_detection.rs | 18 ++++++++++++++++-- 2 files changed, 23 insertions(+), 6 deletions(-) diff --git a/pgdog/src/backend/pool/lb/mod.rs b/pgdog/src/backend/pool/lb/mod.rs index 819a068a8..70a6a5a50 100644 --- a/pgdog/src/backend/pool/lb/mod.rs +++ b/pgdog/src/backend/pool/lb/mod.rs @@ -216,11 +216,14 @@ impl LoadBalancer { .for_each(|(_, target)| { target.1.set_role(Role::Replica); }); - } else if targets.iter().all(|target| target.0.valid()) { + } else { // All targets are replicas until we get a primary. - targets.iter().for_each(|target| { - target.1.set_role(Role::Replica); - }); + targets + .iter() + .filter(|target| target.0.valid()) + .for_each(|target| { + target.1.set_role(Role::Replica); + }); } self.elected_primary.send_replace(self.primary().cloned()); diff --git a/pgdog/src/backend/pool/lb/test/role_detection.rs b/pgdog/src/backend/pool/lb/test/role_detection.rs index f1346d839..78743d27c 100644 --- a/pgdog/src/backend/pool/lb/test/role_detection.rs +++ b/pgdog/src/backend/pool/lb/test/role_detection.rs @@ -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()); @@ -74,7 +85,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 From 8efc9b614fbe677bb8a30f57afe329b2545a546d Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 6 Oct 2026 15:26:09 -0700 Subject: [PATCH 2/2] fix read_only --- pgdog/src/backend/pool/lb/mod.rs | 25 +++++++----- .../backend/pool/lb/test/role_detection.rs | 39 +++++++++++++++++++ .../parser/query/test/test_replica_only.rs | 16 ++++++++ 3 files changed, 71 insertions(+), 9 deletions(-) diff --git a/pgdog/src/backend/pool/lb/mod.rs b/pgdog/src/backend/pool/lb/mod.rs index 70a6a5a50..b40245578 100644 --- a/pgdog/src/backend/pool/lb/mod.rs +++ b/pgdog/src/backend/pool/lb/mod.rs @@ -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); @@ -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 } @@ -217,13 +224,13 @@ impl LoadBalancer { target.1.set_role(Role::Replica); }); } else { - // All targets are replicas until we get a primary. - targets - .iter() - .filter(|target| target.0.valid()) - .for_each(|target| { - target.1.set_role(Role::Replica); - }); + 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()); @@ -308,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() diff --git a/pgdog/src/backend/pool/lb/test/role_detection.rs b/pgdog/src/backend/pool/lb/test/role_detection.rs index 78743d27c..e385e0ae9 100644 --- a/pgdog/src/backend/pool/lb/test/role_detection.rs +++ b/pgdog/src/backend/pool/lb/test/role_detection.rs @@ -66,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"]); diff --git a/pgdog/src/frontend/router/parser/query/test/test_replica_only.rs b/pgdog/src/frontend/router/parser/query/test/test_replica_only.rs index 4c4b6b237..f0ef4c9ca 100644 --- a/pgdog/src/frontend/router/parser/query/test/test_replica_only.rs +++ b/pgdog/src/frontend/router/parser/query/test/test_replica_only.rs @@ -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");