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
80 changes: 79 additions & 1 deletion integration/rust/tests/integration/notify.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};

use crate::setup::admin_sqlx;
use crate::setup::{admin_sqlx, connection_sqlx_direct};
use parking_lot::Mutex;
use rust_decimal::prelude::ToPrimitive;
use sqlx::{Connection, Executor, PgConnection, Pool, Postgres, Row, postgres::PgListener};
Expand All @@ -10,6 +10,84 @@ use tokio::{
sync::Barrier,
time::{sleep, timeout},
};
use uuid::Uuid;

#[tokio::test]
async fn test_listener_restarts_and_resubscribes_after_backend_termination() {
async fn notify_and_receive(
notifier: &mut PgConnection,
listener: &mut PgListener,
channel: &str,
) -> u32 {
let payload = Uuid::new_v4().to_string();
timeout(Duration::from_secs(10), async {
loop {
notifier
.execute(format!("NOTIFY \"{channel}\", '{payload}'").as_str())
.await
.expect("send notification through PgDog");

// Notifications can be lost while the backend reconnects.
// try_recv exposes frontend disconnects instead of hiding them
// behind SQLx's automatic reconnect and resubscribe behavior.
if let Ok(notification) =
timeout(Duration::from_millis(100), listener.try_recv()).await
{
let notification = notification
.expect("receive notification")
.expect("PgDog should keep the frontend connection open");
if notification.payload() == payload {
assert_eq!(notification.channel(), channel);
return notification.process_id();
}
}
}
})
.await
.expect("listener should deliver notifications within ten seconds")
}

let test_id = Uuid::new_v4().simple().to_string();
let channels = [
format!("restart_events_{test_id}"),
format!("restart_jobs_{test_id}"),
];
let mut listener = PgListener::connect("postgres://pgdog:pgdog@127.0.0.1:6432/pgdog")
.await
.expect("connect listener");
listener.eager_reconnect(false);
for channel in &channels {
listener.listen(channel).await.expect("listen on channel");
}
let mut notifier = PgConnection::connect("postgres://pgdog:pgdog@127.0.0.1:6432/pgdog")
.await
.expect("connect notifier");

// NOTIFY through PgDog uses the pub/sub backend, so the notification PID
// identifies the exact connection to terminate.
let old_pid = notify_and_receive(&mut notifier, &mut listener, &channels[0]).await;
assert_eq!(
notify_and_receive(&mut notifier, &mut listener, &channels[1]).await,
old_pid
);

let direct = connection_sqlx_direct().await;
let terminated: bool = sqlx::query_scalar("SELECT pg_terminate_backend($1, 5000)")
.bind(old_pid as i32)
.fetch_one(&direct)
.await
.expect("terminate pub/sub backend and wait for it to exit");
assert!(terminated, "pub/sub backend should be terminated");

// Keep the original frontend and subscriptions: PgDog must restart its
// listener and restore both channels without another LISTEN from the client.
let new_pid = notify_and_receive(&mut notifier, &mut listener, &channels[0]).await;
assert_ne!(new_pid, old_pid, "listener should use a new backend");
assert_eq!(
notify_and_receive(&mut notifier, &mut listener, &channels[1]).await,
new_pid
);
}

#[tokio::test]
async fn test_notify() {
Expand Down
2 changes: 2 additions & 0 deletions pgdog/src/backend/pub_sub/listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,8 @@ impl PubSubListener {
result = Self::run(id, &pool, &pool_key, &mut rx, channels.clone()) => {
if let Err(err) = result {
error!("pub/sub error: {} [{}]", err, pool.addr());
// Request restart of listener.
comms.start.notify_one();
// Don't reconnect for another connect attempt delay
// to avoid connection storms during incidents.
select! {
Expand Down
Loading