From 28b5b37e40a103afe88b06b28b08a3ccc32b842a Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 6 Oct 2026 11:42:53 -0700 Subject: [PATCH] fix: primary detection needs to account for timeline --- pgdog-stats/src/replication.rs | 3 ++ pgdog/src/backend/pool/lb/mod.rs | 8 +-- .../backend/pool/lb/test/role_detection.rs | 28 ++++++++++ pgdog/src/backend/pool/lsn_monitor.rs | 52 +++++++++++-------- pgdog/src/backend/pool/shard/monitor.rs | 1 + 5 files changed, 66 insertions(+), 26 deletions(-) diff --git a/pgdog-stats/src/replication.rs b/pgdog-stats/src/replication.rs index bc11d58f2..385c762d8 100644 --- a/pgdog-stats/src/replication.rs +++ b/pgdog-stats/src/replication.rs @@ -99,6 +99,8 @@ pub struct LsnStats { pub fetched: SystemTime, /// Running on Aurora. pub aurora: bool, + /// Timeline + pub timeline: i64, } /// Schema-only mirror of `std::time::SystemTime`'s default serde representation. @@ -124,6 +126,7 @@ impl Default for LsnStats { timestamp: TimestampTz::default(), fetched: SystemTime::now(), aurora: false, + timeline: 0, } } } diff --git a/pgdog/src/backend/pool/lb/mod.rs b/pgdog/src/backend/pool/lb/mod.rs index ac6467de9..819a068a8 100644 --- a/pgdog/src/backend/pool/lb/mod.rs +++ b/pgdog/src/backend/pool/lb/mod.rs @@ -1,6 +1,7 @@ //! Load balanced connection pool. use std::{ + cmp::Reverse, sync::{ Arc, atomic::{AtomicBool, AtomicI64, AtomicUsize, Ordering}, @@ -185,15 +186,14 @@ impl LoadBalancer { .map(|target| (target.pool.lsn_stats(), target)) .collect::>(); - // Pick primary by latest data. The one with the most - // up-to-date lsn number and pg_is_in_recovery() = false - // is the new primary. + // Pick the primary with the greatest timeline, then the freshest + // LSN stats, among targets with pg_is_in_recovery() = false. // // The old primary is still part of the config and will be demoted // to replica. If it's down, it will be banned from serving traffic. // let now = SystemTime::now(); - targets.sort_by_cached_key(|target| target.0.lsn_age(now)); + targets.sort_by_cached_key(|target| (Reverse(target.0.timeline), target.0.lsn_age(now))); let primary = 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 c6b9ddce9..f1346d839 100644 --- a/pgdog/src/backend/pool/lb/test/role_detection.rs +++ b/pgdog/src/backend/pool/lb/test/role_detection.rs @@ -95,3 +95,31 @@ async fn test_roles_detected_survives_reload_and_new_targets_remain_unknown() { reloaded.redetect_roles(); assert!(reloaded.roles_detected()); } + +#[test] +fn test_redetect_roles_prefers_highest_timeline_then_freshest_stats() { + for (first_timeline, second_timeline, expected_primary) in [ + (1, 2, "127.0.0.1"), + (2, 2, "localhost"), + (0, 0, "localhost"), + ] { + let lb = auto_pool(&["localhost", "127.0.0.1"]); + let now = SystemTime::now(); + for (target, timeline, age) in [ + (&lb.targets[0], first_timeline, Duration::ZERO), + (&lb.targets[1], second_timeline, Duration::from_secs(10)), + ] { + set_lsn_stats(target, false, 100); + let mut stats = target.pool.inner().lsn_stats.write(); + stats.timeline = timeline; + stats.fetched = now - age; + } + + assert!(lb.redetect_roles()); + assert_eq!( + lb.primary().expect("elected primary").addr().host, + expected_primary, + "timelines: {first_timeline}, {second_timeline}" + ); + } +} diff --git a/pgdog/src/backend/pool/lsn_monitor.rs b/pgdog/src/backend/pool/lsn_monitor.rs index f849a2067..6b50f1794 100644 --- a/pgdog/src/backend/pool/lsn_monitor.rs +++ b/pgdog/src/backend/pool/lsn_monitor.rs @@ -22,32 +22,35 @@ pub(crate) use pgdog_stats::replication::ReplicaLag; static AURORA_DETECTION_QUERY: &str = "SELECT aurora_version()"; static LSN_QUERY: &str = " +WITH recovery AS MATERIALIZED ( + SELECT pg_is_in_recovery() AS replica +), +wal AS MATERIALIZED ( + SELECT + replica, + CASE + WHEN replica THEN + COALESCE( + pg_last_wal_replay_lsn(), + pg_last_wal_receive_lsn() + ) + ELSE + pg_current_wal_lsn() + END AS lsn + FROM recovery +) SELECT - pg_is_in_recovery() AS replica, - CASE - WHEN pg_is_in_recovery() THEN - COALESCE( - pg_last_wal_replay_lsn(), - pg_last_wal_receive_lsn() - ) - ELSE - pg_current_wal_lsn() - END AS lsn, - CASE - WHEN pg_is_in_recovery() THEN - COALESCE( - pg_last_wal_replay_lsn(), - pg_last_wal_receive_lsn() - ) - '0/0'::pg_lsn - ELSE - pg_current_wal_lsn() - '0/0'::pg_lsn - END AS offset_bytes, + replica, + lsn, + lsn - '0/0'::pg_lsn AS offset_bytes, CASE - WHEN pg_is_in_recovery() THEN + WHEN replica THEN COALESCE(pg_last_xact_replay_timestamp(), now()) ELSE now() - END AS timestamp + END AS timestamp, + (pg_control_checkpoint()).timeline_id AS timeline +FROM wal "; static AURORA_LSN_QUERY: &str = " @@ -55,7 +58,8 @@ SELECT pg_is_in_recovery() AS replica, '0/0'::pg_lsn AS lsn, 0::bigint AS offset_bytes, - now() AS timestamp + now() AS timestamp, + (pg_control_checkpoint()).timeline_id AS timeline "; /// LSN information. @@ -105,6 +109,7 @@ impl LsnStats { timestamp: value.get(3, Format::Text).unwrap_or_default(), fetched: SystemTime::now(), aurora, + timeline: value.get(4, Format::Text).unwrap_or_default(), } .into() } @@ -367,6 +372,7 @@ mod test { timestamp: TimestampTz::default(), fetched: SystemTime::now(), aurora: false, + timeline: 0, } .into(); @@ -473,6 +479,7 @@ mod test { timestamp: TimestampTz::default(), fetched: SystemTime::now(), aurora: true, + timeline: 0, } .into(); @@ -491,6 +498,7 @@ mod test { timestamp: TimestampTz::default(), fetched: SystemTime::now(), aurora: false, + timeline: 0, } .into(); diff --git a/pgdog/src/backend/pool/shard/monitor.rs b/pgdog/src/backend/pool/shard/monitor.rs index 3ca7b8bd3..404f50988 100644 --- a/pgdog/src/backend/pool/shard/monitor.rs +++ b/pgdog/src/backend/pool/shard/monitor.rs @@ -209,6 +209,7 @@ mod test { timestamp: TimestampTz::decode(timestamp.as_bytes(), Format::Text).unwrap(), fetched: SystemTime::now(), aurora: false, + timeline: 0, } .into() }