diff --git a/integration/rust/tests/integration/notify.rs b/integration/rust/tests/integration/notify.rs index 8baf8cad7..36d326d67 100644 --- a/integration/rust/tests/integration/notify.rs +++ b/integration/rust/tests/integration/notify.rs @@ -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}; @@ -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() { diff --git a/pgdog/src/backend/pub_sub/listener.rs b/pgdog/src/backend/pub_sub/listener.rs index 6ed8139cc..09bad8f1e 100644 --- a/pgdog/src/backend/pub_sub/listener.rs +++ b/pgdog/src/backend/pub_sub/listener.rs @@ -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! {