Skip to content
Merged
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
3 changes: 3 additions & 0 deletions pgdog-stats/src/replication.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -124,6 +126,7 @@ impl Default for LsnStats {
timestamp: TimestampTz::default(),
fetched: SystemTime::now(),
aurora: false,
timeline: 0,
}
}
}
Expand Down
8 changes: 4 additions & 4 deletions pgdog/src/backend/pool/lb/mod.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
//! Load balanced connection pool.

use std::{
cmp::Reverse,
sync::{
Arc,
atomic::{AtomicBool, AtomicI64, AtomicUsize, Ordering},
Expand Down Expand Up @@ -185,15 +186,14 @@ impl LoadBalancer {
.map(|target| (target.pool.lsn_stats(), target))
.collect::<Vec<_>>();

// 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()
Expand Down
28 changes: 28 additions & 0 deletions pgdog/src/backend/pool/lb/test/role_detection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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}"
);
}
}
52 changes: 30 additions & 22 deletions pgdog/src/backend/pool/lsn_monitor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,40 +22,44 @@ 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 = "
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.
Expand Down Expand Up @@ -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()
}
Expand Down Expand Up @@ -367,6 +372,7 @@ mod test {
timestamp: TimestampTz::default(),
fetched: SystemTime::now(),
aurora: false,
timeline: 0,
}
.into();

Expand Down Expand Up @@ -473,6 +479,7 @@ mod test {
timestamp: TimestampTz::default(),
fetched: SystemTime::now(),
aurora: true,
timeline: 0,
}
.into();

Expand All @@ -491,6 +498,7 @@ mod test {
timestamp: TimestampTz::default(),
fetched: SystemTime::now(),
aurora: false,
timeline: 0,
}
.into();

Expand Down
1 change: 1 addition & 0 deletions pgdog/src/backend/pool/shard/monitor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -209,6 +209,7 @@ mod test {
timestamp: TimestampTz::decode(timestamp.as_bytes(), Format::Text).unwrap(),
fetched: SystemTime::now(),
aurora: false,
timeline: 0,
}
.into()
}
Expand Down
Loading