From 58c4ff402a2ea3228c14bc18b45686352127b6d2 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 21 Jul 2026 09:58:09 -0700 Subject: [PATCH 01/22] feat: add instance id to txn name --- pgdog/src/backend/pool/connection/binding.rs | 19 +++++++++++---- .../replication/logical/subscriber/copy.rs | 12 +++++----- .../client/query_engine/end_transaction.rs | 6 ++--- .../client/query_engine/two_pc/manager.rs | 24 +++++++++---------- .../client/query_engine/two_pc/mod.rs | 4 ++-- .../client/query_engine/two_pc/statement.rs | 6 +++-- .../client/query_engine/two_pc/transaction.rs | 4 +++- 7 files changed, 44 insertions(+), 31 deletions(-) diff --git a/pgdog/src/backend/pool/connection/binding.rs b/pgdog/src/backend/pool/connection/binding.rs index cf3fb0d77..09503afc6 100644 --- a/pgdog/src/backend/pool/connection/binding.rs +++ b/pgdog/src/backend/pool/connection/binding.rs @@ -3,7 +3,10 @@ use crate::{ frontend::{ ClientRequest, - client::query_engine::{TwoPcPhase, two_pc::statement::phase_control}, + client::query_engine::{ + TwoPcPhase, + two_pc::{TwoPcTransaction, statement::phase_control}, + }, }, net::{FrontendPid, ProtocolMessage, Query, parameter::Parameters}, state::State, @@ -361,14 +364,14 @@ impl Binding { pub(crate) async fn two_pc_on_guards( servers: &mut [Guard], - name: &str, + transaction: TwoPcTransaction, phase: TwoPcPhase, ) -> Result<(), Error> { let skip_missing = matches!(phase, TwoPcPhase::Phase2 | TwoPcPhase::Rollback); let mut futures = Vec::new(); for (shard, server) in servers.iter_mut().enumerate() { - let query = phase_control(name, shard, phase); + let query = phase_control(transaction, shard, phase); futures.push(server.execute(query)); } @@ -394,9 +397,15 @@ impl Binding { } /// Execute two-phase commit transaction control statements. - pub async fn two_pc(&mut self, name: &str, phase: TwoPcPhase) -> Result<(), Error> { + pub async fn two_pc( + &mut self, + transaction: TwoPcTransaction, + phase: TwoPcPhase, + ) -> Result<(), Error> { match self { - Binding::MultiShard(servers, _) => Self::two_pc_on_guards(servers, name, phase).await, + Binding::MultiShard(servers, _) => { + Self::two_pc_on_guards(servers, transaction, phase).await + } _ => Err(Error::TwoPcMultiShardOnly), } diff --git a/pgdog/src/backend/replication/logical/subscriber/copy.rs b/pgdog/src/backend/replication/logical/subscriber/copy.rs index eeeccdb86..6b7c5eff0 100644 --- a/pgdog/src/backend/replication/logical/subscriber/copy.rs +++ b/pgdog/src/backend/replication/logical/subscriber/copy.rs @@ -296,14 +296,14 @@ impl CopySubscriber { async { let _guard_phase_1 = manager - .transaction_state(&txn, &identifier, TwoPcPhase::Phase1) + .transaction_state(txn, &identifier, TwoPcPhase::Phase1) .await?; - self.two_pc_on_shards(&txn, TwoPcPhase::Phase1).await?; + self.two_pc_on_shards(txn, TwoPcPhase::Phase1).await?; let _guard_phase_2 = manager - .transaction_state(&txn, &identifier, TwoPcPhase::Phase2) + .transaction_state(txn, &identifier, TwoPcPhase::Phase2) .await?; - self.two_pc_on_shards(&txn, TwoPcPhase::Phase2).await?; + self.two_pc_on_shards(txn, TwoPcPhase::Phase2).await?; manager.done(&txn).await?; Ok(()) @@ -317,7 +317,7 @@ impl CopySubscriber { async fn two_pc_on_shards( &mut self, - txn: &TwoPcTransaction, + txn: TwoPcTransaction, phase: TwoPcPhase, ) -> Result<(), Error> { let mut futures = Vec::new(); @@ -329,7 +329,7 @@ impl CopySubscriber { // via binding.rs using the same phase_control() helper. let query = match phase { TwoPcPhase::Rollback => unreachable!(), - phase => phase_control(&txn.to_string(), shard, phase), + phase => phase_control(txn, shard, phase), }; futures.push(Self::send_and_confirm(server, Query::new(query).into())); } diff --git a/pgdog/src/frontend/client/query_engine/end_transaction.rs b/pgdog/src/frontend/client/query_engine/end_transaction.rs index fc74b149f..affcc9053 100644 --- a/pgdog/src/frontend/client/query_engine/end_transaction.rs +++ b/pgdog/src/frontend/client/query_engine/end_transaction.rs @@ -106,17 +106,17 @@ impl QueryEngine { } let identifier = cluster.identifier(); - let name = self.two_pc.transaction().to_string(); + let transaction = self.two_pc.transaction(); // If interrupted here, the transaction must be rolled back. let _guard_phase_1 = self.two_pc.phase_one(&identifier).await?; - self.backend.two_pc(&name, TwoPcPhase::Phase1).await?; + self.backend.two_pc(transaction, TwoPcPhase::Phase1).await?; debug!("[2pc] phase 1 complete"); // If interrupted here, the transaction must be committed. let _guard_phase_2 = self.two_pc.phase_two(&identifier).await?; - self.backend.two_pc(&name, TwoPcPhase::Phase2).await?; + self.backend.two_pc(transaction, TwoPcPhase::Phase2).await?; debug!("[2pc] phase 2 complete"); diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index 94dd9fa36..ff22ffb23 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -201,14 +201,14 @@ impl Manager { /// the WAL exists to prevent. pub(crate) async fn transaction_state( &self, - transaction: &TwoPcTransaction, + transaction: TwoPcTransaction, identifier: &Arc, phase: TwoPcPhase, ) -> Result { let prior = { let mut guard = self.inner.lock(); - let prior = guard.transactions.get(transaction).cloned(); - let entry = guard.transactions.entry(*transaction).or_default(); + let prior = guard.transactions.get(&transaction).cloned(); + let entry = guard.transactions.entry(transaction).or_default(); entry.identifier = identifier.clone(); entry.phase = phase; prior @@ -218,13 +218,13 @@ impl Manager { let result = match phase { TwoPcPhase::Phase1 => { wal.append_begin( - *transaction, + transaction, identifier.user.clone(), identifier.database.clone(), ) .await } - TwoPcPhase::Phase2 => wal.append_committing(*transaction).await, + TwoPcPhase::Phase2 => wal.append_committing(transaction).await, TwoPcPhase::Rollback => { unreachable!("rollback is not a state transition; it's the cleanup direction") } @@ -233,10 +233,10 @@ impl Manager { let mut guard = self.inner.lock(); match prior { Some(prior) => { - guard.transactions.insert(*transaction, prior); + guard.transactions.insert(transaction, prior); } None => { - guard.transactions.remove(transaction); + guard.transactions.remove(&transaction); } } warn!( @@ -248,7 +248,7 @@ impl Manager { } Ok(TwoPcGuard { - transaction: *transaction, + transaction, manager: Self::get(), }) } @@ -309,7 +309,7 @@ impl Manager { r#"[2pc] cleaning up transaction "{}""#, transaction.to_string() ); - match manager.cleanup_phase(&transaction).await { + match manager.cleanup_phase(transaction).await { Err(err) => { error!( r#"[2pc] error cleaning up "{}" transaction: {}"#, @@ -345,8 +345,8 @@ impl Manager { } /// Reconnect to cluster if available and rollback the two-phase transaction. - async fn cleanup_phase(&self, transaction: &TwoPcTransaction) -> Result<(), Error> { - let state = match self.inner.lock().transactions.get(transaction).cloned() { + async fn cleanup_phase(&self, transaction: TwoPcTransaction) -> Result<(), Error> { + let state = match self.inner.lock().transactions.get(&transaction).cloned() { Some(state) => state, _ => { return Ok(()); @@ -379,7 +379,7 @@ impl Manager { &Route::write(ShardWithPriority::new_override_transaction(Shard::All)), ) .await?; - connection.two_pc(&transaction.to_string(), phase).await?; + connection.two_pc(transaction, phase).await?; connection.disconnect(); Ok(()) diff --git a/pgdog/src/frontend/client/query_engine/two_pc/mod.rs b/pgdog/src/frontend/client/query_engine/two_pc/mod.rs index 1024d12ab..c6d109284 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/mod.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/mod.rs @@ -58,7 +58,7 @@ impl TwoPc { pub(super) async fn phase_one(&mut self, cluster: &Arc) -> Result { let transaction = self.transaction(); self.manager - .transaction_state(&transaction, cluster, TwoPcPhase::Phase1) + .transaction_state(transaction, cluster, TwoPcPhase::Phase1) .await } @@ -68,7 +68,7 @@ impl TwoPc { pub(super) async fn phase_two(&mut self, cluster: &Arc) -> Result { let transaction = self.transaction(); self.manager - .transaction_state(&transaction, cluster, TwoPcPhase::Phase2) + .transaction_state(transaction, cluster, TwoPcPhase::Phase2) .await } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs index 4456db029..fc06c2511 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs @@ -1,15 +1,17 @@ //! Per-shard two-phase commit transaction names and control statements. +use crate::frontend::client::query_engine::two_pc::TwoPcTransaction; + use super::TwoPcPhase; /// Prepared transaction name for a coordinator transaction on one shard. -pub fn shard_name(transaction: &str, shard: usize) -> String { +pub fn shard_name(transaction: TwoPcTransaction, shard: usize) -> String { format!("{transaction}_{shard}") } /// Build `PREPARE TRANSACTION`, `COMMIT PREPARED`, or `ROLLBACK PREPARED` /// for a shard participant. -pub fn phase_control(transaction: &str, shard: usize, phase: TwoPcPhase) -> String { +pub fn phase_control(transaction: TwoPcTransaction, shard: usize, phase: TwoPcPhase) -> String { let name = shard_name(transaction, shard); match phase { diff --git a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs index be09e9d00..343f0e9c8 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs @@ -2,6 +2,8 @@ use rand::{Rng, rng}; use serde::{Deserialize, Serialize}; use std::{fmt::Display, str::FromStr}; +use crate::util::instance_id; + #[derive(Debug, Clone, Copy, PartialEq, Hash, Eq, Serialize, Deserialize)] pub struct TwoPcTransaction(usize); @@ -17,7 +19,7 @@ impl TwoPcTransaction { impl Display for TwoPcTransaction { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{PREFIX}{}", self.0) + write!(f, "{PREFIX}_{}_{}", instance_id(), self.0) } } From 4b79ec0e5ff4f2e5b89c2d05a1ef7a58f0ab1ab3 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 21 Jul 2026 10:33:50 -0700 Subject: [PATCH 02/22] tests --- .../backend/pool/connection/binding_test.rs | 62 ++++++++++--------- .../client/query_engine/two_pc/statement.rs | 20 +++--- .../client/query_engine/two_pc/test.rs | 8 +-- .../client/query_engine/two_pc/transaction.rs | 16 ++++- 4 files changed, 61 insertions(+), 45 deletions(-) diff --git a/pgdog/src/backend/pool/connection/binding_test.rs b/pgdog/src/backend/pool/connection/binding_test.rs index 8fb3f300e..0a6258b6e 100644 --- a/pgdog/src/backend/pool/connection/binding_test.rs +++ b/pgdog/src/backend/pool/connection/binding_test.rs @@ -8,7 +8,7 @@ mod tests { server::test::test_server, }, frontend::{ - client::query_engine::TwoPcPhase, + client::query_engine::{TwoPcPhase, two_pc::TwoPcTransaction}, router::{ Route, parser::{Shard, ShardWithPriority}, @@ -76,7 +76,9 @@ mod tests { let guard = crate::backend::pool::Guard::new(pool, server, Instant::now()); let mut binding = Binding::Direct(guard, 0); - let result = binding.two_pc("test", TwoPcPhase::Phase1).await; + let result = binding + .two_pc(TwoPcTransaction::new(), TwoPcPhase::Phase1) + .await; // Should fail with TwoPcMultiShardOnly error assert!(result.is_err()); @@ -94,7 +96,9 @@ mod tests { let admin_server = AdminServer::default(); let mut binding = Binding::Admin(admin_server); - let result = binding.two_pc("test", TwoPcPhase::Phase1).await; + let result = binding + .two_pc(TwoPcTransaction::new(), TwoPcPhase::Phase1) + .await; // Should fail with TwoPcMultiShardOnly error assert!(result.is_err()); @@ -106,10 +110,10 @@ mod tests { #[tokio::test] async fn test_two_pc_phase1_prepare() { let mut binding = create_multishard_binding().await; - let transaction_name = "test_txn"; + let transaction = TwoPcTransaction::new(); // Test Phase1 - PREPARE TRANSACTION - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; // Should succeed if let Err(ref error) = result { @@ -118,60 +122,60 @@ mod tests { assert!(result.is_ok()); // Cleanup: Rollback the prepared transaction to avoid leaving dangling transactions - let _cleanup = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let _cleanup = binding.two_pc(transaction, TwoPcPhase::Rollback).await; } #[tokio::test] async fn test_two_pc_phase2_commit() { let mut binding = create_multishard_binding().await; - let transaction_name = "test_commit_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then commit it - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!(result.is_ok()); } #[tokio::test] async fn test_two_pc_rollback() { let mut binding = create_multishard_binding().await; - let transaction_name = "test_rollback_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then rollback - let result = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let result = binding.two_pc(transaction, TwoPcPhase::Rollback).await; assert!(result.is_ok()); } #[tokio::test] async fn test_two_pc_commit_after_prepare_and_commit() { let mut binding = create_multishard_binding().await; - let transaction_name = "committed_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then commit it binding - .two_pc(transaction_name, TwoPcPhase::Phase2) + .two_pc(transaction, TwoPcPhase::Phase2) .await .expect("Phase2 should succeed"); // Try to commit again - should succeed because skip_missing is true for Phase2 - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!( result.is_ok(), "Committing non-existent prepared transaction should be skipped" @@ -181,22 +185,22 @@ mod tests { #[tokio::test] async fn test_two_pc_rollback_after_prepare_and_rollback() { let mut binding = create_multishard_binding().await; - let transaction_name = "rolled_back_txn"; + let transaction = TwoPcTransaction::new(); // First prepare the transaction binding - .two_pc(transaction_name, TwoPcPhase::Phase1) + .two_pc(transaction, TwoPcPhase::Phase1) .await .expect("Phase1 should succeed"); // Then rollback it binding - .two_pc(transaction_name, TwoPcPhase::Rollback) + .two_pc(transaction, TwoPcPhase::Rollback) .await .expect("Rollback should succeed"); // Try to rollback again - should succeed because skip_missing is true for Rollback - let result = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let result = binding.two_pc(transaction, TwoPcPhase::Rollback).await; assert!( result.is_ok(), "Rolling back non-existent prepared transaction should be skipped" @@ -209,22 +213,22 @@ mod tests { #[tokio::test] async fn test_two_pc_transaction_lifecycle() { let mut binding = create_multishard_binding().await; - let transaction_name = "lifecycle_test"; + let transaction = TwoPcTransaction::new(); // 1. Prepare transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; assert!(result.is_ok(), "Phase1 preparation should succeed"); // 2. Try to prepare the same transaction again - PostgreSQL behavior may vary - let _result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let _result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; // Note: PostgreSQL behavior for duplicate PREPARE TRANSACTION can vary depending on context // 3. Commit the prepared transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!(result.is_ok(), "Phase2 commit should succeed"); // 4. Try to commit again - should succeed (skip_missing = true) - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!( result.is_ok(), "Committing non-existent transaction should be skipped" @@ -234,18 +238,18 @@ mod tests { #[tokio::test] async fn test_two_pc_prepare_then_rollback() { let mut binding = create_multishard_binding().await; - let transaction_name = "prepare_rollback_test"; + let transaction = TwoPcTransaction::new(); // 1. Prepare transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase1).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase1).await; assert!(result.is_ok(), "Phase1 preparation should succeed"); // 2. Rollback the prepared transaction - let result = binding.two_pc(transaction_name, TwoPcPhase::Rollback).await; + let result = binding.two_pc(transaction, TwoPcPhase::Rollback).await; assert!(result.is_ok(), "Rollback should succeed"); // 3. Try to commit after rollback - should succeed (skip_missing = true) - let result = binding.two_pc(transaction_name, TwoPcPhase::Phase2).await; + let result = binding.two_pc(transaction, TwoPcPhase::Phase2).await; assert!( result.is_ok(), "Committing rolled back transaction should be skipped" diff --git a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs index fc06c2511..9cc689107 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs @@ -27,23 +27,27 @@ mod test { #[test] fn shard_name_appends_index() { - assert_eq!(shard_name("__pgdog_2pc_42", 0), "__pgdog_2pc_42_0"); - assert_eq!(shard_name("test_txn", 3), "test_txn_3"); + let transaction = TwoPcTransaction::new(); + + assert_eq!(shard_name(transaction, 0), format!("{transaction}_0")); + assert_eq!(shard_name(transaction, 3), format!("{transaction}_3")); } #[test] fn phase_control_statements() { + let transaction = TwoPcTransaction::new(); + assert_eq!( - phase_control("test", 1, TwoPcPhase::Phase1), - "PREPARE TRANSACTION 'test_1'" + phase_control(transaction, 1, TwoPcPhase::Phase1), + format!("PREPARE TRANSACTION '{transaction}_1'") ); assert_eq!( - phase_control("test", 1, TwoPcPhase::Phase2), - "COMMIT PREPARED 'test_1'" + phase_control(transaction, 1, TwoPcPhase::Phase2), + format!("COMMIT PREPARED '{transaction}_1'") ); assert_eq!( - phase_control("test", 1, TwoPcPhase::Rollback), - "ROLLBACK PREPARED 'test_1'" + phase_control(transaction, 1, TwoPcPhase::Rollback), + format!("ROLLBACK PREPARED '{transaction}_1'") ); } } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/test.rs b/pgdog/src/frontend/client/query_engine/two_pc/test.rs index 8d6e8fb32..0ad10f27d 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/test.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/test.rs @@ -38,9 +38,7 @@ async fn test_cleanup_transaction_phase_one() { let info = Manager::get().transaction(&transaction).unwrap(); assert_eq!(info.phase, TwoPcPhase::Phase1); - conn.two_pc(&transaction.to_string(), TwoPcPhase::Phase1) - .await - .unwrap(); + conn.two_pc(transaction, TwoPcPhase::Phase1).await.unwrap(); let two_pc = conn .execute("SELECT * FROM pg_prepared_xacts") @@ -110,9 +108,7 @@ async fn test_cleanup_transaction_phase_two() { let info = Manager::get().transaction(&transaction).unwrap(); assert_eq!(info.phase, TwoPcPhase::Phase1); - conn.two_pc(&transaction.to_string(), TwoPcPhase::Phase1) - .await - .unwrap(); + conn.two_pc(transaction, TwoPcPhase::Phase1).await.unwrap(); let txns = conn .execute("SELECT * FROM pg_prepared_xacts") diff --git a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs index 343f0e9c8..2902fab75 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs @@ -19,7 +19,7 @@ impl TwoPcTransaction { impl Display for TwoPcTransaction { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{PREFIX}_{}_{}", instance_id(), self.0) + write!(f, "{PREFIX}{}_{}", instance_id(), self.0) } } @@ -27,7 +27,7 @@ impl FromStr for TwoPcTransaction { type Err = (); fn from_str(s: &str) -> Result { - let id = s.split(PREFIX).last().map(|id| id.parse()); + let id = s.rsplit("_").next().map(|id| id.parse()); if let Some(Ok(id)) = id { Ok(Self(id)) @@ -48,4 +48,16 @@ mod test { let reverse = TwoPcTransaction::from_str(transaction.to_string().as_str()).unwrap(); assert_eq!(reverse.0, transaction.0); } + + #[test] + fn test_instance_id() { + for id in [1024, 11111111, usize::MAX, usize::MIN] { + let transaction = TwoPcTransaction(id); + let instance_id = instance_id(); // Generate it, it's a singleton. + assert_eq!( + format!("__pgdog_2pc_{instance_id}_{id}"), + transaction.to_string() + ); + } + } } From f7f5e5787bde1d59bb2ff924685d69210d9e37dc Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 21 Jul 2026 17:08:12 -0700 Subject: [PATCH 03/22] feat: cleanup abandoned 2pc transactions --- pgdog-config/src/general.rs | 12 ++++ .../client/query_engine/two_pc/manager.rs | 70 ++++++++++++++++++- .../two_pc/server_transactions.rs | 42 ++++++++--- .../client/query_engine/two_pc/transaction.rs | 33 ++++++++- pgdog/src/main.rs | 10 ++- pgdog/src/util.rs | 16 +++++ 6 files changed, 166 insertions(+), 17 deletions(-) diff --git a/pgdog-config/src/general.rs b/pgdog-config/src/general.rs index 5ad4cc018..ce0384deb 100644 --- a/pgdog-config/src/general.rs +++ b/pgdog-config/src/general.rs @@ -622,6 +622,13 @@ pub struct General { #[serde(default = "General::two_phase_commit_wal_checkpoint_interval")] pub two_phase_commit_wal_checkpoint_interval: u64, + /// Rollback abandoned transactions. + /// + /// WARNING: Data loss will occur, enable this only if you don't care about consistency + /// and are not using the 2pc WAL. + #[serde(default = "General::two_phase_commit_rollback_abandoned")] + pub two_phase_commit_rollback_abandoned: bool, + /// Enable expanded (`\x`) output for `EXPLAIN` results returned by PgDog's built-in query plan aggregation. #[serde(default = "General::expanded_explain")] pub expanded_explain: bool, @@ -884,6 +891,7 @@ impl Default for General { two_phase_commit_wal_fsync_interval: Self::two_phase_commit_wal_fsync_interval(), two_phase_commit_wal_checkpoint_interval: Self::two_phase_commit_wal_checkpoint_interval(), + two_phase_commit_rollback_abandoned: Self::two_phase_commit_rollback_abandoned(), expanded_explain: Self::expanded_explain(), server_lifetime: Self::server_lifetime(), server_lifetime_jitter: Self::server_lifetime_jitter(), @@ -1062,6 +1070,10 @@ impl General { Self::env_or_default("PGDOG_TWO_PHASE_COMMIT_WAL_CHECKPOINT_INTERVAL", 60) } + fn two_phase_commit_rollback_abandoned() -> bool { + Self::env_bool_or_default("PGDOG_TWO_PHASE_COMMIT_ROLLBACK_ABANDONED", false) + } + fn idle_timeout() -> u64 { Self::env_or_default( "PGDOG_IDLE_TIMEOUT", diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index ff22ffb23..f089a9610 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -21,14 +21,15 @@ use tracing::{debug, error, info, warn}; use crate::{ backend::{ - databases::User, + Cluster, + databases::{User, databases}, pool::{Connection, Request}, }, config::config, frontend::{ client::query_engine::{ TwoPcPhase, - two_pc::{TwoPcGuard, TwoPcStats, TwoPcTransaction, wal::Wal}, + two_pc::{TwoPcGuard, TwoPcStats, TwoPcTransaction, TwoPcTransactions, wal::Wal}, }, router::{ Route, @@ -38,7 +39,7 @@ use crate::{ tasks, }; -use super::Error; +use super::{Error, server_transactions::TwoPcServerTransaction, statement::phase_control}; static MANAGER: Lazy = Lazy::new(Manager::init); static MAINTENANCE: Duration = Duration::from_millis(333); @@ -295,6 +296,8 @@ impl Manager { debug!("[2pc] monitor started"); + // Cleanup orphaned transactions. + loop { // Wake up either because it's time to check // or manager told us to. @@ -385,6 +388,67 @@ impl Manager { Ok(()) } + /// Drop abandoned two-phase commit transactions. + /// + /// WARNING: This only happens if durability for 2pc is off. Running this + /// will cause data loss. + /// + pub async fn cleanup_abandoned(&self) -> Result<(), Error> { + let mut cleaned_up = 0; + for cluster in databases().all().values() { + cleaned_up += self.cleanup_abandoned_for_cluster(cluster).await?; + } + + if cleaned_up > 0 { + warn!("[2pc] rolled back up {} abandoned transactions", cleaned_up); + } else { + info!("[2pc] no abandoned transactions found"); + } + + Ok(()) + } + + async fn cleanup_abandoned_for_cluster(&self, cluster: &Cluster) -> Result { + use crate::backend::Error as BackendError; + use crate::backend::pool::Error as PoolError; + let mut cleaned_up = 0; + + for (number, shard) in cluster.shards().iter().enumerate() { + let mut conn = match shard.primary(&Request::default()).await { + Ok(conn) => conn, + Err(PoolError::NoPrimary) => continue, + Err(err) => return Err(BackendError::Pool(err).into()), + }; + + let txns = TwoPcTransactions::load(&mut conn).await?; + + for txn in txns.iter() { + if let TwoPcServerTransaction::Ours { txn, user, .. } = txn { + // Postgres transactions can be only be rolled back by their owners. + if user == cluster.user() { + if txn.is_mine() { + match conn + .execute(phase_control(*txn, number, TwoPcPhase::Rollback)) + .await + { + Ok(_) => cleaned_up += 1, + Err(BackendError::ExecutionError(err)) => { + warn!( + "[2pc] error cleaning abandoned transaction \"{}\": {}", + txn, err + ); + } + Err(err) => return Err(err.into()), + } + } + } + } + } + } + + Ok(cleaned_up) + } + /// Shutdown manager and wait for all transactions to be cleaned up. /// Once the monitor has drained the cleanup queue, the WAL is shut /// down too so any final End records make it to disk before exit. diff --git a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs index 8e95be6f0..39aafc332 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs @@ -9,8 +9,16 @@ use super::TwoPcTransaction; #[derive(Debug, Clone)] pub enum TwoPcServerTransaction { - Ours(TwoPcTransaction), - Other { name: String }, + Ours { + txn: TwoPcTransaction, + user: String, + database: String, + }, + Other { + name: String, + user: String, + database: String, + }, } pub(crate) struct TwoPcTransactions { @@ -28,20 +36,32 @@ impl Deref for TwoPcTransactions { impl TwoPcTransactions { /// Load two phase transactions from the server. pub(crate) async fn load(server: &mut Server) -> Result { - let records: Vec = server.fetch_all("SELECT * FROM pg_prepared_xacts").await?; + let records: Vec = server + .fetch_all("SELECT gid, owner, database FROM pg_prepared_xacts") + .await?; let mut transactions = vec![]; for record in records { - let transaction = record.get_text(1).map(|name| { - if let Ok(ours) = TwoPcTransaction::from_str(&name) { - TwoPcServerTransaction::Ours(ours) + let gid = record.get_text(0); + let user = record.get_text(1).unwrap_or_default(); + let database = record.get_text(2).unwrap_or_default(); + + if let Some(gid) = gid { + let txn = if let Ok(txn) = TwoPcTransaction::from_str(&gid) { + TwoPcServerTransaction::Ours { + txn, + user, + database, + } } else { - TwoPcServerTransaction::Other { name } - } - }); + TwoPcServerTransaction::Other { + name: gid, + user, + database, + } + }; - if let Some(transaction) = transaction { - transactions.push(transaction); + transactions.push(txn); } } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs index 2902fab75..9b5daed31 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs @@ -2,7 +2,7 @@ use rand::{Rng, rng}; use serde::{Deserialize, Serialize}; use std::{fmt::Display, str::FromStr}; -use crate::util::instance_id; +use crate::util::{deployment_id, instance_id}; #[derive(Debug, Clone, Copy, PartialEq, Hash, Eq, Serialize, Deserialize)] pub struct TwoPcTransaction(usize); @@ -15,11 +15,30 @@ impl TwoPcTransaction { // so multiple instances of PgDog don't create an identical transaction. Self(rng().random_range(0..usize::MAX)) } + + /// This transaction was created by this process. + pub(crate) fn is_mine(&self) -> bool { + self.to_string().starts_with(&Self::global_prefix()) + } + + /// A prefix to identify two-phase commit transactions generated + /// by this PgDog process. + pub(crate) fn global_prefix() -> String { + format!( + "{PREFIX}{}{}_", + if let Some(cluster_id) = deployment_id() { + format!("{}_", cluster_id) + } else { + "".into() + }, + instance_id(), + ) + } } impl Display for TwoPcTransaction { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{PREFIX}{}_{}", instance_id(), self.0) + write!(f, "{}{}", Self::global_prefix(), self.0) } } @@ -39,6 +58,8 @@ impl FromStr for TwoPcTransaction { #[cfg(test)] mod test { + use crate::test_utils::set_env_var; + use super::*; #[test] @@ -60,4 +81,12 @@ mod test { ); } } + + #[test] + fn test_deployment_id() { + let _guard = set_env_var("DEPLOYMENT_ID", "1"); + let txn = TwoPcTransaction(1678); + let instance_id = instance_id(); // Generate it, it's a singleton. + assert_eq!(format!("__pgdog_2pc_1_{instance_id}_1678"), txn.to_string()); + } } diff --git a/pgdog/src/main.rs b/pgdog/src/main.rs index f7a995023..5864692eb 100644 --- a/pgdog/src/main.rs +++ b/pgdog/src/main.rs @@ -11,10 +11,10 @@ use pgdog::config::{self, config}; use pgdog::frontend::client::query_engine::two_pc::Manager; use pgdog::frontend::listener::Listener; use pgdog::frontend::prepared_statements; -use pgdog::plugin; use pgdog::stats; use pgdog::util::pgdog_version; use pgdog::{healthcheck, net}; +use pgdog::{plugin, tasks}; use tokio::runtime::Builder; use tracing::{error, info, warn}; @@ -157,6 +157,14 @@ async fn pgdog(command: Option) -> Result<(), Box &'static str { /// Get an externally assigned, unique, node identifier /// for this instance of PgDog. +/// +/// This assumes the NODE ID follows the following format: +/// +/// - +/// pub fn node_id() -> Result { // split always returns at least one element. instance_id().split("-").last().unwrap().parse() } +static DEPLOYMENT_ID: Lazy> = Lazy::new(|| env::var("DEPLOYMENT_ID").ok()); + +/// Get the ID of this PgDog deployment. +/// +/// This should be _globally_ unique +/// and is used to differentiate 2pc transactions. +/// +pub fn deployment_id() -> Option<&'static str> { + DEPLOYMENT_ID.as_deref() +} + static HOSTNAME: Lazy = Lazy::new(|| { let hostname = env::var("HOSTNAME").unwrap_or_default(); let host = env::var("HOST").unwrap_or_default(); From 585000740aba4fa4abae98b6e00059bfa77e48a8 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 21 Jul 2026 17:10:28 -0700 Subject: [PATCH 04/22] schema --- .schema/pgdog.schema.json | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/.schema/pgdog.schema.json b/.schema/pgdog.schema.json index 3900e97cd..8c894fb0b 100644 --- a/.schema/pgdog.schema.json +++ b/.schema/pgdog.schema.json @@ -117,6 +117,7 @@ "tls_verify": "prefer", "two_phase_commit": false, "two_phase_commit_auto": null, + "two_phase_commit_rollback_abandoned": false, "two_phase_commit_wal_checkpoint_interval": 60, "two_phase_commit_wal_dir": null, "two_phase_commit_wal_fsync_interval": 2, @@ -1191,6 +1192,11 @@ ], "default": null }, + "two_phase_commit_rollback_abandoned": { + "description": "Rollback abandoned transactions.\n\nWARNING: Data loss will occur, enable this only if you don't care about consistency\nand are not using the 2pc WAL.", + "type": "boolean", + "default": false + }, "two_phase_commit_wal_checkpoint_interval": { "description": "How often, in seconds, to write a checkpoint record to the two-phase commit WAL and garbage-collect old segments.\n\n_Default:_ `60`\n\n", "type": "integer", From 8c7b1f07d60ed430f315adb2ba75cbf0153fce9c Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Tue, 21 Jul 2026 17:13:30 -0700 Subject: [PATCH 05/22] cippeh --- .../client/query_engine/two_pc/manager.rs | 26 +++++++++---------- 1 file changed, 12 insertions(+), 14 deletions(-) diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index f089a9610..7ccc25421 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -425,21 +425,19 @@ impl Manager { for txn in txns.iter() { if let TwoPcServerTransaction::Ours { txn, user, .. } = txn { // Postgres transactions can be only be rolled back by their owners. - if user == cluster.user() { - if txn.is_mine() { - match conn - .execute(phase_control(*txn, number, TwoPcPhase::Rollback)) - .await - { - Ok(_) => cleaned_up += 1, - Err(BackendError::ExecutionError(err)) => { - warn!( - "[2pc] error cleaning abandoned transaction \"{}\": {}", - txn, err - ); - } - Err(err) => return Err(err.into()), + if user == cluster.user() && txn.is_mine() { + match conn + .execute(phase_control(*txn, number, TwoPcPhase::Rollback)) + .await + { + Ok(_) => cleaned_up += 1, + Err(BackendError::ExecutionError(err)) => { + warn!( + "[2pc] error cleaning abandoned transaction \"{}\": {}", + txn, err + ); } + Err(err) => return Err(err.into()), } } } From 6f3473c970e1a728106a4d8192f514a97986097f Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 06:49:47 -0700 Subject: [PATCH 06/22] Update pgdog/src/frontend/client/query_engine/two_pc/statement.rs Co-authored-by: Sage Griffin --- pgdog/src/frontend/client/query_engine/two_pc/statement.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs index 9cc689107..88a3da198 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs @@ -5,7 +5,7 @@ use crate::frontend::client::query_engine::two_pc::TwoPcTransaction; use super::TwoPcPhase; /// Prepared transaction name for a coordinator transaction on one shard. -pub fn shard_name(transaction: TwoPcTransaction, shard: usize) -> String { +pub(crate) fn shard_name(transaction: TwoPcTransaction, shard: usize) -> String { format!("{transaction}_{shard}") } From 72b846e073aeef1b468ef22f4290f62acc46fded Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 06:49:56 -0700 Subject: [PATCH 07/22] Update pgdog/src/util.rs Co-authored-by: Sage Griffin --- pgdog/src/util.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pgdog/src/util.rs b/pgdog/src/util.rs index 7f7eea202..8ea220448 100644 --- a/pgdog/src/util.rs +++ b/pgdog/src/util.rs @@ -155,7 +155,7 @@ static DEPLOYMENT_ID: Lazy> = Lazy::new(|| env::var("DEPLOYMENT_I /// This should be _globally_ unique /// and is used to differentiate 2pc transactions. /// -pub fn deployment_id() -> Option<&'static str> { +pub(crate) fn deployment_id() -> Option<&'static str> { DEPLOYMENT_ID.as_deref() } From 3f08c65debf6e1b4429c87b424b4640289b79e3e Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 06:50:05 -0700 Subject: [PATCH 08/22] Update pgdog/src/frontend/client/query_engine/two_pc/transaction.rs Co-authored-by: Sage Griffin --- pgdog/src/frontend/client/query_engine/two_pc/transaction.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs index 9b5daed31..f91e03082 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs @@ -23,7 +23,7 @@ impl TwoPcTransaction { /// A prefix to identify two-phase commit transactions generated /// by this PgDog process. - pub(crate) fn global_prefix() -> String { + fn global_prefix() -> String { format!( "{PREFIX}{}{}_", if let Some(cluster_id) = deployment_id() { From 3b3f9afe6b611c9b8310bb7127d615b4f8291a89 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 08:28:28 -0700 Subject: [PATCH 09/22] tests and parallelization and remove the user check which was incorrect --- pgdog-config/src/general.rs | 11 +++ pgdog/src/config/mod.rs | 13 +++- .../client/query_engine/two_pc/manager.rs | 65 ++++++++++------- .../two_pc/server_transactions.rs | 10 ++- .../client/query_engine/two_pc/test.rs | 71 +++++++++++++++++++ .../client/query_engine/two_pc/transaction.rs | 6 +- pgdog/src/main.rs | 23 +++++- 7 files changed, 166 insertions(+), 33 deletions(-) diff --git a/pgdog-config/src/general.rs b/pgdog-config/src/general.rs index ce0384deb..af4669099 100644 --- a/pgdog-config/src/general.rs +++ b/pgdog-config/src/general.rs @@ -629,6 +629,11 @@ pub struct General { #[serde(default = "General::two_phase_commit_rollback_abandoned")] pub two_phase_commit_rollback_abandoned: bool, + /// Maximum amount of time to block startup in order to rollback abandoned + /// transactions. + #[serde(default = "General::two_phase_commit_rollback_abandoned_timeout")] + pub two_phase_commit_rollback_abandoned_timeout: u64, + /// Enable expanded (`\x`) output for `EXPLAIN` results returned by PgDog's built-in query plan aggregation. #[serde(default = "General::expanded_explain")] pub expanded_explain: bool, @@ -892,6 +897,8 @@ impl Default for General { two_phase_commit_wal_checkpoint_interval: Self::two_phase_commit_wal_checkpoint_interval(), two_phase_commit_rollback_abandoned: Self::two_phase_commit_rollback_abandoned(), + two_phase_commit_rollback_abandoned_timeout: + Self::two_phase_commit_rollback_abandoned_timeout(), expanded_explain: Self::expanded_explain(), server_lifetime: Self::server_lifetime(), server_lifetime_jitter: Self::server_lifetime_jitter(), @@ -1074,6 +1081,10 @@ impl General { Self::env_bool_or_default("PGDOG_TWO_PHASE_COMMIT_ROLLBACK_ABANDONED", false) } + fn two_phase_commit_rollback_abandoned_timeout() -> u64 { + Self::env_or_default("PGDOG_TWO_PHASE_COMMIT_ROLLBACK_ABANDONED_TIMEOUT", 15_000) + } + fn idle_timeout() -> u64 { Self::env_or_default( "PGDOG_IDLE_TIMEOUT", diff --git a/pgdog/src/config/mod.rs b/pgdog/src/config/mod.rs index 22cc79d10..4aa140ae9 100644 --- a/pgdog/src/config/mod.rs +++ b/pgdog/src/config/mod.rs @@ -154,6 +154,16 @@ pub fn load_test() { #[cfg(test)] pub fn load_test_with_pooler_mode(pooler_mode: PoolerMode) { + load_test_with_user_and_pooler_mode("pgdog", pooler_mode, Role::default()) +} + +#[cfg(test)] +pub fn load_test_with_user(user: &str) { + load_test_with_user_and_pooler_mode(user, PoolerMode::Transaction, Role::Primary) +} + +#[cfg(test)] +fn load_test_with_user_and_pooler_mode(user: &str, pooler_mode: PoolerMode, role: Role) { use crate::backend::databases::init; let mut config = ConfigAndUsers::default(); @@ -162,10 +172,11 @@ pub fn load_test_with_pooler_mode(pooler_mode: PoolerMode) { host: "127.0.0.1".into(), port: 5432, pooler_mode: Some(pooler_mode), + role, ..Default::default() }]; config.users.users = vec![User { - name: "pgdog".into(), + name: user.into(), database: "pgdog".into(), password: Some("pgdog".into()), pooler_mode: Some(pooler_mode), diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index 7ccc25421..1f086adf9 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -1,6 +1,7 @@ //! Global two-phase commit transaction manager. use arc_swap::ArcSwapOption; use fnv::FnvHashMap as HashMap; +use futures::future::try_join_all; use once_cell::sync::Lazy; use parking_lot::Mutex; use std::{ @@ -394,6 +395,8 @@ impl Manager { /// will cause data loss. /// pub async fn cleanup_abandoned(&self) -> Result<(), Error> { + info!("[2pc] rolling back abandoned transactions"); + let mut cleaned_up = 0; for cluster in databases().all().values() { cleaned_up += self.cleanup_abandoned_for_cluster(cluster).await?; @@ -409,36 +412,50 @@ impl Manager { } async fn cleanup_abandoned_for_cluster(&self, cluster: &Cluster) -> Result { + let cleaned_up = try_join_all( + cluster + .shards() + .iter() + .enumerate() + .map(|(number, shard)| Self::cleanup_abandoned_for_shard(number, shard)), + ) + .await?; + + Ok(cleaned_up.into_iter().sum()) + } + + async fn cleanup_abandoned_for_shard( + number: usize, + shard: &crate::backend::pool::Shard, + ) -> Result { use crate::backend::Error as BackendError; use crate::backend::pool::Error as PoolError; let mut cleaned_up = 0; - for (number, shard) in cluster.shards().iter().enumerate() { - let mut conn = match shard.primary(&Request::default()).await { - Ok(conn) => conn, - Err(PoolError::NoPrimary) => continue, - Err(err) => return Err(BackendError::Pool(err).into()), - }; + let mut conn = match shard.primary(&Request::default()).await { + Ok(conn) => conn, + Err(PoolError::NoPrimary) => return Ok(cleaned_up), + Err(err) => return Err(BackendError::Pool(err).into()), + }; - let txns = TwoPcTransactions::load(&mut conn).await?; - - for txn in txns.iter() { - if let TwoPcServerTransaction::Ours { txn, user, .. } = txn { - // Postgres transactions can be only be rolled back by their owners. - if user == cluster.user() && txn.is_mine() { - match conn - .execute(phase_control(*txn, number, TwoPcPhase::Rollback)) - .await - { - Ok(_) => cleaned_up += 1, - Err(BackendError::ExecutionError(err)) => { - warn!( - "[2pc] error cleaning abandoned transaction \"{}\": {}", - txn, err - ); - } - Err(err) => return Err(err.into()), + let txns = TwoPcTransactions::load(&mut conn).await?; + + for txn in txns.iter() { + if let TwoPcServerTransaction::Ours { txn, .. } = txn { + // Postgres transactions can be only be rolled back by their owners. + if txn.is_created_by_this_process() { + match conn + .execute(phase_control(*txn, number, TwoPcPhase::Rollback)) + .await + { + Ok(_) => cleaned_up += 1, + Err(BackendError::ExecutionError(err)) => { + warn!( + "[2pc] error cleaning abandoned transaction \"{}\": {}", + txn, err + ); } + Err(err) => return Err(err.into()), } } } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs index 39aafc332..46af54af9 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs @@ -37,7 +37,9 @@ impl TwoPcTransactions { /// Load two phase transactions from the server. pub(crate) async fn load(server: &mut Server) -> Result { let records: Vec = server - .fetch_all("SELECT gid, owner, database FROM pg_prepared_xacts") + .fetch_all( + "SELECT gid, owner, database FROM pg_prepared_xacts WHERE owner = current_user", + ) .await?; let mut transactions = vec![]; @@ -47,7 +49,11 @@ impl TwoPcTransactions { let database = record.get_text(2).unwrap_or_default(); if let Some(gid) = gid { - let txn = if let Ok(txn) = TwoPcTransaction::from_str(&gid) { + let transaction = gid + .rsplit_once('_') + .map(|(transaction, _)| transaction) + .unwrap_or(&gid); + let txn = if let Ok(txn) = TwoPcTransaction::from_str(transaction) { TwoPcServerTransaction::Ours { txn, user, diff --git a/pgdog/src/frontend/client/query_engine/two_pc/test.rs b/pgdog/src/frontend/client/query_engine/two_pc/test.rs index 0ad10f27d..6db21efe7 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/test.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/test.rs @@ -1,5 +1,6 @@ use crate::{ backend::{ + Server, databases::databases, pool::{Connection, Request}, }, @@ -13,6 +14,76 @@ use crate::{ }; use super::*; +use super::{server_transactions::TwoPcServerTransaction, statement::phase_control}; + +async fn server_has_transaction(server: &mut Server, transaction: TwoPcTransaction) -> bool { + TwoPcTransactions::load(server) + .await + .unwrap() + .iter() + .any(|server_transaction| { + matches!( + server_transaction, + TwoPcServerTransaction::Ours { txn, .. } if *txn == transaction + ) + }) +} + +#[tokio::test] +async fn test_cleanup_abandoned() { + config::load_test_with_user("pgdog"); + let cluster = databases().all().iter().next().unwrap().1.clone(); + let transaction = TwoPcTransaction::new(); + let mut conn = cluster.shards()[0] + .primary(&Request::default()) + .await + .unwrap(); + + conn.execute("BEGIN").await.unwrap(); + conn.execute(phase_control(transaction, 0, TwoPcPhase::Phase1)) + .await + .unwrap(); + + assert!(server_has_transaction(&mut conn, transaction).await); + + Manager::get().cleanup_abandoned().await.unwrap(); + + assert!(!server_has_transaction(&mut conn, transaction).await); +} + +#[tokio::test] +async fn test_cleanup_abandoned_different_user() { + config::load_test_with_user("pgdog1"); + let cluster = databases().all().iter().next().unwrap().1.clone(); + let transaction = TwoPcTransaction::new(); + let mut conn = cluster.shards()[0] + .primary(&Request::default()) + .await + .unwrap(); + + conn.execute("BEGIN").await.unwrap(); + conn.execute(phase_control(transaction, 0, TwoPcPhase::Phase1)) + .await + .unwrap(); + assert!(server_has_transaction(&mut conn, transaction).await); + drop(conn); + + config::load_test_with_user("pgdog2"); + Manager::get().cleanup_abandoned().await.unwrap(); + + config::load_test_with_user("pgdog1"); + let cluster = databases().all().iter().next().unwrap().1.clone(); + let mut conn = cluster.shards()[0] + .primary(&Request::default()) + .await + .unwrap(); + let was_preserved = server_has_transaction(&mut conn, transaction).await; + conn.execute(phase_control(transaction, 0, TwoPcPhase::Rollback)) + .await + .unwrap(); + + assert!(was_preserved); +} #[tokio::test] async fn test_cleanup_transaction_phase_one() { diff --git a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs index f91e03082..0e21ac396 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs @@ -17,7 +17,7 @@ impl TwoPcTransaction { } /// This transaction was created by this process. - pub(crate) fn is_mine(&self) -> bool { + pub(crate) fn is_created_by_this_process(&self) -> bool { self.to_string().starts_with(&Self::global_prefix()) } @@ -74,7 +74,7 @@ mod test { fn test_instance_id() { for id in [1024, 11111111, usize::MAX, usize::MIN] { let transaction = TwoPcTransaction(id); - let instance_id = instance_id(); // Generate it, it's a singleton. + let instance_id = instance_id(); // It's a singleton. assert_eq!( format!("__pgdog_2pc_{instance_id}_{id}"), transaction.to_string() @@ -86,7 +86,7 @@ mod test { fn test_deployment_id() { let _guard = set_env_var("DEPLOYMENT_ID", "1"); let txn = TwoPcTransaction(1678); - let instance_id = instance_id(); // Generate it, it's a singleton. + let instance_id = instance_id(); // It's a singleton. assert_eq!(format!("__pgdog_2pc_1_{instance_id}_1678"), txn.to_string()); } } diff --git a/pgdog/src/main.rs b/pgdog/src/main.rs index 5864692eb..54a6beb23 100644 --- a/pgdog/src/main.rs +++ b/pgdog/src/main.rs @@ -3,6 +3,7 @@ use std::fs::read_to_string; use std::path::Path; use std::process::exit; +use std::time::Duration; use clap::Parser; use pgdog::backend::databases; @@ -16,6 +17,7 @@ use pgdog::util::pgdog_version; use pgdog::{healthcheck, net}; use pgdog::{plugin, tasks}; use tokio::runtime::Builder; +use tokio::time::timeout; use tracing::{error, info, warn}; fn main() -> Result<(), Box> { @@ -158,11 +160,26 @@ async fn pgdog(command: Option) -> Result<(), Box { + error!( + "[2pc] abandoned transactions cleanup timed out after {}ms", + abandoned_timeout.as_millis() + ); + } + + Ok(Err(err)) => { error!("[2pc] abandoned transactions cleanup error: {}", err); } - }); + + Ok(_) => (), + }; } let mut listener = Listener::new(format!("{}:{}", general.host, general.port)); From c6f5feb0518d1e6b027eed6a408fd41762bc1987 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 14:49:52 -0700 Subject: [PATCH 10/22] refactor to better rust --- .../client/query_engine/two_pc/manager.rs | 14 ++- .../client/query_engine/two_pc/mod.rs | 1 + .../two_pc/server_transactions.rs | 12 +-- .../client/query_engine/two_pc/statement.rs | 90 ++++++++++++++++--- 4 files changed, 94 insertions(+), 23 deletions(-) diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index 1f086adf9..ac649ffd4 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -397,10 +397,16 @@ impl Manager { pub async fn cleanup_abandoned(&self) -> Result<(), Error> { info!("[2pc] rolling back abandoned transactions"); - let mut cleaned_up = 0; - for cluster in databases().all().values() { - cleaned_up += self.cleanup_abandoned_for_cluster(cluster).await?; - } + let databases = databases(); + let cleaned_up = try_join_all( + databases + .all() + .values() + .map(|cluster| self.cleanup_abandoned_for_cluster(cluster)), + ) + .await? + .into_iter() + .sum::(); if cleaned_up > 0 { warn!("[2pc] rolled back up {} abandoned transactions", cleaned_up); diff --git a/pgdog/src/frontend/client/query_engine/two_pc/mod.rs b/pgdog/src/frontend/client/query_engine/two_pc/mod.rs index c6d109284..b0cba45aa 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/mod.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/mod.rs @@ -18,6 +18,7 @@ pub use guard::TwoPcGuard; pub use manager::Manager; pub use phase::TwoPcPhase; pub(crate) use server_transactions::TwoPcTransactions; +pub(crate) use statement::TwoPcTransactionOnShard; pub use stats::TwoPcStats; pub use transaction::TwoPcTransaction; diff --git a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs index 46af54af9..986e35890 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs @@ -1,11 +1,11 @@ -use std::{ops::Deref, str::FromStr}; +use std::ops::Deref; use crate::{ backend::{Error, Server}, net::DataRow, }; -use super::TwoPcTransaction; +use super::{TwoPcTransaction, TwoPcTransactionOnShard}; #[derive(Debug, Clone)] pub enum TwoPcServerTransaction { @@ -49,13 +49,9 @@ impl TwoPcTransactions { let database = record.get_text(2).unwrap_or_default(); if let Some(gid) = gid { - let transaction = gid - .rsplit_once('_') - .map(|(transaction, _)| transaction) - .unwrap_or(&gid); - let txn = if let Ok(txn) = TwoPcTransaction::from_str(transaction) { + let txn = if let Ok(txn) = gid.parse::() { TwoPcServerTransaction::Ours { - txn, + txn: txn.transaction(), user, database, } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs index 88a3da198..9c59824e4 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/statement.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/statement.rs @@ -1,23 +1,63 @@ //! Per-shard two-phase commit transaction names and control statements. +use std::{fmt::Display, str::FromStr}; + use crate::frontend::client::query_engine::two_pc::TwoPcTransaction; use super::TwoPcPhase; -/// Prepared transaction name for a coordinator transaction on one shard. -pub(crate) fn shard_name(transaction: TwoPcTransaction, shard: usize) -> String { - format!("{transaction}_{shard}") +/// 2pc transaction executed on a shard. We +/// make them unique per shard in case two or more +/// shards are located on the same postgres server. +pub(crate) struct TwoPcTransactionOnShard { + transaction: TwoPcTransaction, + shard: usize, +} + +impl TwoPcTransactionOnShard { + /// Create new 2pc transaction on shard x. + pub(crate) fn new(transaction: TwoPcTransaction, shard: usize) -> Self { + Self { transaction, shard } + } + + /// Get the coordinator transaction. + pub(crate) fn transaction(&self) -> TwoPcTransaction { + self.transaction + } +} + +impl Display for TwoPcTransactionOnShard { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{}_{}", self.transaction, self.shard) + } +} + +impl FromStr for TwoPcTransactionOnShard { + type Err = (); + + fn from_str(s: &str) -> Result { + let (transaction, shard) = s.rsplit_once('_').ok_or(())?; + + Ok(Self { + transaction: transaction.parse()?, + shard: shard.parse().map_err(|_| ())?, + }) + } } /// Build `PREPARE TRANSACTION`, `COMMIT PREPARED`, or `ROLLBACK PREPARED` /// for a shard participant. -pub fn phase_control(transaction: TwoPcTransaction, shard: usize, phase: TwoPcPhase) -> String { - let name = shard_name(transaction, shard); +pub(crate) fn phase_control( + transaction: TwoPcTransaction, + shard: usize, + phase: TwoPcPhase, +) -> String { + let txn = TwoPcTransactionOnShard::new(transaction, shard); match phase { - TwoPcPhase::Phase1 => format!("PREPARE TRANSACTION '{name}'"), - TwoPcPhase::Phase2 => format!("COMMIT PREPARED '{name}'"), - TwoPcPhase::Rollback => format!("ROLLBACK PREPARED '{name}'"), + TwoPcPhase::Phase1 => format!("PREPARE TRANSACTION '{txn}'"), + TwoPcPhase::Phase2 => format!("COMMIT PREPARED '{txn}'"), + TwoPcPhase::Rollback => format!("ROLLBACK PREPARED '{txn}'"), } } @@ -26,11 +66,39 @@ mod test { use super::*; #[test] - fn shard_name_appends_index() { + fn transaction_on_shard_appends_index() { let transaction = TwoPcTransaction::new(); - assert_eq!(shard_name(transaction, 0), format!("{transaction}_0")); - assert_eq!(shard_name(transaction, 3), format!("{transaction}_3")); + assert_eq!( + TwoPcTransactionOnShard::new(transaction, 0).to_string(), + format!("{transaction}_0") + ); + assert_eq!( + TwoPcTransactionOnShard::new(transaction, 3).to_string(), + format!("{transaction}_3") + ); + } + + #[test] + fn parse_transaction_on_shard() { + let transaction = TwoPcTransaction::new(); + let parsed: TwoPcTransactionOnShard = format!("{transaction}_3") + .parse() + .expect("valid transaction on shard"); + + assert_eq!(parsed.transaction, transaction); + assert_eq!(parsed.shard, 3); + } + + #[test] + fn reject_invalid_transaction_on_shard() { + assert!("invalid".parse::().is_err()); + assert!("invalid_0".parse::().is_err()); + assert!( + "__pgdog_2pc_1_invalid" + .parse::() + .is_err() + ); } #[test] From 92e18884dcebc61691cc744d452b2a8fe01bbdfc Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 14:53:20 -0700 Subject: [PATCH 11/22] warn --- .schema/pgdog.schema.json | 8 ++++++++ pgdog/src/main.rs | 2 +- 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/.schema/pgdog.schema.json b/.schema/pgdog.schema.json index 8c894fb0b..a28b11526 100644 --- a/.schema/pgdog.schema.json +++ b/.schema/pgdog.schema.json @@ -118,6 +118,7 @@ "two_phase_commit": false, "two_phase_commit_auto": null, "two_phase_commit_rollback_abandoned": false, + "two_phase_commit_rollback_abandoned_timeout": 15000, "two_phase_commit_wal_checkpoint_interval": 60, "two_phase_commit_wal_dir": null, "two_phase_commit_wal_fsync_interval": 2, @@ -1197,6 +1198,13 @@ "type": "boolean", "default": false }, + "two_phase_commit_rollback_abandoned_timeout": { + "description": "Maximum amount of time to block startup in order to rollback abandoned\ntransactions.", + "type": "integer", + "format": "uint64", + "default": 15000, + "minimum": 0 + }, "two_phase_commit_wal_checkpoint_interval": { "description": "How often, in seconds, to write a checkpoint record to the two-phase commit WAL and garbage-collect old segments.\n\n_Default:_ `60`\n\n", "type": "integer", diff --git a/pgdog/src/main.rs b/pgdog/src/main.rs index 54a6beb23..6d3c66579 100644 --- a/pgdog/src/main.rs +++ b/pgdog/src/main.rs @@ -12,10 +12,10 @@ use pgdog::config::{self, config}; use pgdog::frontend::client::query_engine::two_pc::Manager; use pgdog::frontend::listener::Listener; use pgdog::frontend::prepared_statements; +use pgdog::plugin; use pgdog::stats; use pgdog::util::pgdog_version; use pgdog::{healthcheck, net}; -use pgdog::{plugin, tasks}; use tokio::runtime::Builder; use tokio::time::timeout; use tracing::{error, info, warn}; From c140bb6f68b29ef7586000556a397b9f468593b6 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 15:02:04 -0700 Subject: [PATCH 12/22] pub(crate) --- pgdog/src/backend/pool/connection/binding.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pgdog/src/backend/pool/connection/binding.rs b/pgdog/src/backend/pool/connection/binding.rs index 09503afc6..a762b248d 100644 --- a/pgdog/src/backend/pool/connection/binding.rs +++ b/pgdog/src/backend/pool/connection/binding.rs @@ -397,7 +397,7 @@ impl Binding { } /// Execute two-phase commit transaction control statements. - pub async fn two_pc( + pub(crate) async fn two_pc( &mut self, transaction: TwoPcTransaction, phase: TwoPcPhase, From 6655ac99d6b803430c1be8fe3373f23eb6903b53 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 18:44:57 -0700 Subject: [PATCH 13/22] refactor --- pgdog-config/src/general.rs | 4 +- pgdog/src/backend/databases.rs | 47 ++++ pgdog/src/backend/pool/cluster.rs | 228 +++++++---------- pgdog/src/backend/pool/cluster_launch.rs | 233 ++++++++++++++++++ pgdog/src/backend/pool/mod.rs | 1 + .../replication/logical/orchestrator.rs | 2 +- .../client/query_engine/route_query.rs | 6 +- .../client/query_engine/two_pc/manager.rs | 5 +- pgdog/src/main.rs | 27 +- 9 files changed, 387 insertions(+), 166 deletions(-) create mode 100644 pgdog/src/backend/pool/cluster_launch.rs diff --git a/pgdog-config/src/general.rs b/pgdog-config/src/general.rs index af4669099..09deb1ba6 100644 --- a/pgdog-config/src/general.rs +++ b/pgdog-config/src/general.rs @@ -629,8 +629,8 @@ pub struct General { #[serde(default = "General::two_phase_commit_rollback_abandoned")] pub two_phase_commit_rollback_abandoned: bool, - /// Maximum amount of time to block startup in order to rollback abandoned - /// transactions. + /// Maximum amount of time, in milliseconds, each cluster waits for abandoned + /// transactions rollback on boot before allowing traffic. #[serde(default = "General::two_phase_commit_rollback_abandoned_timeout")] pub two_phase_commit_rollback_abandoned_timeout: u64, diff --git a/pgdog/src/backend/databases.rs b/pgdog/src/backend/databases.rs index f3d6b762b..d71b6e04d 100644 --- a/pgdog/src/backend/databases.rs +++ b/pgdog/src/backend/databases.rs @@ -467,6 +467,9 @@ impl Databases { cluster.user(), cluster.name() ); + // No boot-time maintenance will run, don't block + // readiness waiters. Checkouts will fail instead. + cluster.mark_ready(); } else { cluster.launch(); } @@ -1486,6 +1489,50 @@ mod tests { ); } + #[test] + fn test_two_pc_rollback_abandoned_config_passed_to_cluster() { + let config = Config { + databases: vec![Database { + name: "db1".to_string(), + host: "localhost".to_string(), + port: 5432, + role: Role::Primary, + ..Default::default() + }], + general: General { + two_phase_commit_rollback_abandoned: true, + two_phase_commit_rollback_abandoned_timeout: 1234, + ..Default::default() + }, + ..Default::default() + }; + + let users = crate::config::Users { + users: vec![crate::config::User { + name: "user".to_string(), + database: "db1".to_string(), + password: Some("pass".to_string()), + ..Default::default() + }], + ..Default::default() + }; + + let databases = from_config(&ConfigAndUsers { + config, + users, + config_path: std::path::PathBuf::new(), + users_path: std::path::PathBuf::new(), + ..Default::default() + }); + + let cluster = databases.cluster(("user", "db1")).unwrap(); + assert!(cluster.two_pc_rollback_abandoned()); + assert_eq!( + cluster.two_pc_rollback_abandoned_timeout(), + std::time::Duration::from_millis(1234) + ); + } + #[test] fn test_user_all_databases_creates_pools_for_all_dbs() { let config = Config { diff --git a/pgdog/src/backend/pool/cluster.rs b/pgdog/src/backend/pool/cluster.rs index 78f39af2e..8bfc33a1e 100644 --- a/pgdog/src/backend/pool/cluster.rs +++ b/pgdog/src/backend/pool/cluster.rs @@ -6,22 +6,13 @@ use pgdog_config::{ LoadSchema, PreparedStatements, QueryParser, QueryParserEngine, QueryParserLevel, Rewrite, RewriteMode, users::PasswordKind, }; -use std::{ - sync::{ - Arc, - atomic::{AtomicBool, Ordering}, - }, - time::Duration, -}; -use tracing::error; +use std::{sync::Arc, time::Duration}; use crate::frontend::router::sharding::ShardedTable; -use crate::tasks; use crate::{ backend::{ Schema, ShardedTables, databases::{User as DatabaseUser, databases}, - pool::ee::schema_changed_hook, replication::{ReplicationConfig, ShardedSchemas}, }, config::{ @@ -31,7 +22,10 @@ use crate::{ net::{Query, messages::FrontendPid}, }; -use super::{Address, Config, Error, Guard, MirrorStats, Request, Shard, ShardConfig}; +use super::{ + Address, Config, Error, Guard, MirrorStats, Request, Shard, ShardConfig, + cluster_launch::Readiness, +}; use crate::config::LoadBalancingStrategy; #[derive(Clone, Debug, Default)] @@ -43,11 +37,6 @@ pub struct PoolConfig { pub(crate) config: Config, } -#[derive(Default, Debug)] -struct Readiness { - online: AtomicBool, -} - /// A collection of sharded replicas and primaries /// belonging to the same database cluster. #[derive(Clone, Default, Debug)] @@ -67,7 +56,9 @@ pub struct Cluster { cross_shard_disabled: bool, two_phase_commit: bool, two_phase_commit_auto: bool, - readiness: Arc, + two_pc_rollback_abandoned: bool, + two_pc_rollback_abandoned_timeout: Duration, + pub(super) readiness: Arc, rewrite: Rewrite, prepared_statements: PreparedStatements, dry_run: bool, @@ -152,6 +143,8 @@ pub struct ClusterConfig<'a> { pub cross_shard_disabled: bool, pub two_pc: bool, pub two_pc_auto: bool, + pub two_pc_rollback_abandoned: bool, + pub two_pc_rollback_abandoned_timeout: u64, pub sharded_schemas: ShardedSchemas, pub rewrite: &'a Rewrite, pub prepared_statements: &'a PreparedStatements, @@ -215,6 +208,8 @@ impl<'a> ClusterConfig<'a> { two_pc_auto: user .two_phase_commit_auto .unwrap_or(general.two_phase_commit_auto.unwrap_or(false)), // Disable by default. + two_pc_rollback_abandoned: general.two_phase_commit_rollback_abandoned, + two_pc_rollback_abandoned_timeout: general.two_phase_commit_rollback_abandoned_timeout, sharded_schemas, rewrite, prepared_statements: &general.prepared_statements, @@ -262,6 +257,8 @@ impl Cluster { cross_shard_disabled, two_pc, two_pc_auto, + two_pc_rollback_abandoned, + two_pc_rollback_abandoned_timeout, sharded_schemas, rewrite, prepared_statements, @@ -323,6 +320,10 @@ impl Cluster { cross_shard_disabled, two_phase_commit: two_pc && shards.len() > 1, two_phase_commit_auto: two_pc_auto && shards.len() > 1, + two_pc_rollback_abandoned, + two_pc_rollback_abandoned_timeout: Duration::from_millis( + two_pc_rollback_abandoned_timeout, + ), readiness: Arc::new(Readiness::default()), rewrite: rewrite.clone(), prepared_statements: *prepared_statements, @@ -561,7 +562,7 @@ impl Cluster { self.reload_schema_on_ddl && self.load_schema() } - fn load_schema(&self) -> bool { + pub(super) fn load_schema(&self) -> bool { match self.load_schema { LoadSchema::On => true, LoadSchema::Off => false, @@ -608,6 +609,17 @@ impl Cluster { self.two_phase_commit_auto && self.two_pc_enabled() } + /// Rollback abandoned two-phase commit transactions on launch. + pub(crate) fn two_pc_rollback_abandoned(&self) -> bool { + self.two_pc_rollback_abandoned + } + + /// Maximum time abandoned transactions cleanup can take on launch + /// before traffic is allowed through. + pub(crate) fn two_pc_rollback_abandoned_timeout(&self) -> Duration { + self.two_pc_rollback_abandoned_timeout + } + /// How many parallel COPY commands can we /// run to re-shard this cluster. pub fn resharding_parallel_copies(&self) -> usize { @@ -635,65 +647,6 @@ impl Cluster { self.resharding_replication_retry_min_delay } - /// Launch the connection pools. - pub(crate) fn launch(&self) { - for shard in self.shards() { - shard.launch(); - } - - self.readiness.online.store(true, Ordering::Relaxed); - - if !self.load_schema() { - for shard in &self.shards { - shard.schema_not_needed(); - } - return; - } - - for shard in self.shards() { - let identifier = self.identifier(); - let shard = shard.clone(); - - tasks::spawn("cluster schema sync", async move { - use tokio::time::sleep; - - loop { - match shard.load_schema().await { - Ok(true) => { - schema_changed_hook(&shard.schema(), &identifier, &shard); - return; - } - Ok(false) => return, - Err(err) => { - if shard.online() { - error!( - "error loading schema for shard {}: {}", - shard.number(), - err - ); - sleep(Duration::from_millis(100)).await; - } else { - // Cluster is shutting down: unblock any - // wait_schema_loaded callers. - shard.schema_not_needed(); - return; - } - } - } - } - }); - } - } - - /// Shutdown the connection pools. - pub(crate) fn shutdown(&self) { - for shard in self.shards() { - shard.shutdown(); - } - - self.readiness.online.store(false, Ordering::Relaxed); - } - /// Send a cancellation request for all running queries. pub(crate) async fn cancel_all(&self) -> Result<(), Error> { let pools: Vec<_> = self @@ -709,22 +662,6 @@ impl Cluster { Ok(()) } - /// Is the cluster online? - pub(crate) fn online(&self) -> bool { - self.readiness.online.load(Ordering::Relaxed) - } - - /// Schema loaded for all shards? - pub(crate) async fn wait_schema_loaded(&self) { - if !self.load_schema() { - return; - } - - for shard in &self.shards { - shard.wait_schema_loaded().await; - } - } - /// Execute a query on every primary in the cluster. pub async fn execute( &self, @@ -1075,83 +1012,106 @@ mod test { } #[tokio::test] - async fn test_launch_schema_loading_idempotent() { - use tokio::time::{Duration, sleep}; + async fn test_launch_marks_ready() { + let config = ConfigAndUsers::default(); + let cluster = Cluster::new_test(&config); + + // Pre-populate per-shard schemas so launch doesn't touch a real database. + for shard in &cluster.shards { + shard.schema_not_needed(); + } + + assert!(!cluster.ready()); + cluster.launch(); + cluster.wait_ready().await; + assert!(cluster.ready()); + } + + #[tokio::test] + async fn test_wait_ready_waits_for_two_pc_cleanup() { + use tokio::time::{Duration, timeout}; let config = ConfigAndUsers::default(); let mut cluster = Cluster::new_test(&config); - cluster.sharded_schemas = ShardedSchemas::default(); - - assert!(cluster.load_schema()); + cluster.two_pc_rollback_abandoned = true; + // Expire immediately, don't touch a real database. + cluster.two_pc_rollback_abandoned_timeout = Duration::ZERO; - // Pre-populate per-shard schemas so launch's spawned loaders become - // no-ops (Shard::load_schema returns Ok(false) when already initialized). - // This avoids touching a real database in the test. + // Pre-populate per-shard schemas, same reason. for shard in &cluster.shards { shard.schema_not_needed(); } cluster.launch(); - cluster.wait_schema_loaded().await; - // Second launch must be safe: per-shard OnceCell prevents any reload. - cluster.launch(); - sleep(Duration::from_millis(50)).await; - cluster.wait_schema_loaded().await; + // Cleanup runs in the background, readiness can't be set yet. + assert!(!cluster.ready()); + + let result = timeout(Duration::from_millis(500), cluster.wait_ready()).await; + assert!(result.is_ok()); + assert!(cluster.ready()); } #[tokio::test] - async fn test_wait_schema_loaded_returns_immediately_when_not_needed() { + async fn test_shutdown_releases_readiness_waiters() { + use tokio::time::{Duration, timeout}; + let config = ConfigAndUsers::default(); - let cluster = Cluster::new_test_single_shard(&config); + let cluster = Cluster::new_test(&config); - // load_schema() returns false for single shard without multi_tenant - assert!(!cluster.load_schema()); + assert!(!cluster.ready()); - // Should return immediately without waiting - cluster.wait_schema_loaded().await; + let waiter = cluster.clone(); + let handle = tokio::spawn(async move { + waiter.wait_ready().await; + }); + + cluster.shutdown(); + + let result = timeout(Duration::from_millis(200), handle).await; + assert!(result.is_ok()); } #[tokio::test] - async fn test_wait_schema_loaded_fast_path_when_already_loaded() { + async fn test_launch_schema_loading_idempotent() { + use tokio::time::{Duration, sleep}; + let config = ConfigAndUsers::default(); let mut cluster = Cluster::new_test(&config); cluster.sharded_schemas = ShardedSchemas::default(); assert!(cluster.load_schema()); - // Mark every shard's schema as already loaded. + // Pre-populate per-shard schemas so launch's spawned loaders become + // no-ops (Shard::load_schema returns Ok(false) when already initialized). + // This avoids touching a real database in the test. for shard in &cluster.shards { shard.schema_not_needed(); } - // Each shard's wait_schema_loaded() takes the fast path and returns - // immediately because schema.initialized() is true. - cluster.wait_schema_loaded().await; + cluster.launch(); + cluster.wait_ready().await; + + // Second launch must be safe: per-shard OnceCell prevents any reload. + cluster.launch(); + sleep(Duration::from_millis(50)).await; + cluster.wait_ready().await; } #[tokio::test] - async fn test_wait_schema_loaded_waits_for_notification() { - use tokio::time::{Duration, timeout}; - + async fn test_wait_ready_returns_immediately_when_schema_not_needed() { let config = ConfigAndUsers::default(); - let mut cluster = Cluster::new_test(&config); - cluster.sharded_schemas = ShardedSchemas::default(); + let cluster = Cluster::new_test_single_shard(&config); - assert!(cluster.load_schema()); + // load_schema() returns false for single shard without multi_tenant + assert!(!cluster.load_schema()); - // Trigger schema_not_needed on each shard after a short delay so the - // waiter wakes up via the per-shard schema_waiter notification. - let shards: Vec<_> = cluster.shards.to_vec(); - tokio::spawn(async move { - tokio::time::sleep(Duration::from_millis(10)).await; - for shard in &shards { - shard.schema_not_needed(); - } - }); + cluster.launch(); - let result = timeout(Duration::from_millis(200), cluster.wait_schema_loaded()).await; - assert!(result.is_ok()); + // Should return without waiting: no schema to load, + // two-pc cleanup disabled. + cluster.wait_ready().await; + assert!(cluster.ready()); } #[test] diff --git a/pgdog/src/backend/pool/cluster_launch.rs b/pgdog/src/backend/pool/cluster_launch.rs new file mode 100644 index 000000000..82ef2f8e3 --- /dev/null +++ b/pgdog/src/backend/pool/cluster_launch.rs @@ -0,0 +1,233 @@ +//! Cluster startup and shutdown primitives. +//! +//! Launching and shutting down the connection pools, and gating +//! traffic until boot-time maintenance (two-phase commit cleanup, +//! schema sync) has completed. + +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; + +use tokio::time::{sleep, timeout}; +use tokio_util::sync::CancellationToken; +use tracing::{error, warn}; + +use crate::backend::pool::ee::schema_changed_hook; +use crate::frontend::client::query_engine::two_pc::Manager; +use crate::tasks; + +use super::Cluster; + +/// Cluster readiness state. +/// +/// `online` is set once the pools are launched and cleared on shutdown. +/// `launch_waiter` is cancelled once boot-time maintenance — two-phase +/// commit cleanup and schema loading — is done and the cluster can +/// serve traffic. +#[derive(Default, Debug)] +pub(super) struct Readiness { + online: AtomicBool, + launch_waiter: CancellationToken, + two_pc_waiter: CancellationToken, +} + +impl Readiness { + fn set_online(&self, online: bool) { + self.online.store(online, Ordering::Relaxed); + } + + fn online(&self) -> bool { + self.online.load(Ordering::Relaxed) + } + + fn mark_ready(&self) { + self.launch_waiter.cancel(); + } + + #[cfg(test)] + fn ready(&self) -> bool { + self.launch_waiter.is_cancelled() + } + + async fn wait_ready(&self) { + self.launch_waiter.cancelled().await; + } + + fn mark_two_pc_cleaned_up(&self) { + self.two_pc_waiter.cancel(); + } + + async fn wait_two_pc_cleaned_up(&self) { + self.two_pc_waiter.cancelled().await; + } +} + +impl Cluster { + /// Launch the connection pools. + pub(crate) fn launch(&self) { + for shard in self.shards() { + shard.launch(); + } + + self.readiness.set_online(true); + + self.launch_two_pc_cleanup(); + self.launch_schema_sync(); + self.launch_readiness_monitor(); + } + + /// Shutdown the connection pools. + pub(crate) fn shutdown(&self) { + for shard in self.shards() { + shard.shutdown(); + } + + self.readiness.set_online(false); + + // Release readiness waiters, the cluster is going away. + self.mark_ready(); + } + + /// Is the cluster online? + pub(crate) fn online(&self) -> bool { + self.readiness.online() + } + + /// Wait until boot-time maintenance — two-phase commit cleanup + /// and schema loading — is done and the cluster can serve traffic. + /// + /// Two-phase commit cleanup is bounded by + /// `two_phase_commit_rollback_abandoned_timeout`; schema loading + /// retries until success, so callers should apply their own timeout. + pub(crate) async fn wait_ready(&self) { + self.readiness.wait_ready().await; + } + + /// Boot-time maintenance is done and the cluster can serve traffic. + #[cfg(test)] + pub(crate) fn ready(&self) -> bool { + self.readiness.ready() + } + + /// Mark the cluster as ready to serve traffic. + pub(crate) fn mark_ready(&self) { + self.readiness.mark_ready(); + } + + /// Rollback abandoned two-phase commit transactions before + /// serving traffic. Queries wait on [`Cluster::wait_ready`] + /// until this is done. + fn launch_two_pc_cleanup(&self) { + if !self.two_pc_rollback_abandoned() { + self.readiness.mark_two_pc_cleaned_up(); + return; + } + + let cluster = self.clone(); + tasks::spawn("two-pc abandoned transactions cleanup", async move { + cluster.rollback_abandoned_two_pc().await; + cluster.readiness.mark_two_pc_cleaned_up(); + }); + } + + /// Mark the cluster ready once two-phase commit cleanup + /// and schema loading are done. + fn launch_readiness_monitor(&self) { + let cluster = self.clone(); + tasks::spawn("cluster readiness monitor", async move { + cluster.readiness.wait_two_pc_cleaned_up().await; + cluster.wait_schema_loaded().await; + cluster.mark_ready(); + }); + } + + /// Wait for the schema to load on all shards. + async fn wait_schema_loaded(&self) { + if !self.load_schema() { + return; + } + + for shard in self.shards() { + shard.wait_schema_loaded().await; + } + } + + /// Load database schema in the background, retrying + /// until success or shutdown. + fn launch_schema_sync(&self) { + if !self.load_schema() { + for shard in self.shards() { + shard.schema_not_needed(); + } + return; + } + + for shard in self.shards() { + let identifier = self.identifier(); + let shard = shard.clone(); + + tasks::spawn("cluster schema sync", async move { + loop { + match shard.load_schema().await { + Ok(true) => { + schema_changed_hook(&shard.schema(), &identifier, &shard); + return; + } + Ok(false) => return, + Err(err) => { + if shard.online() { + error!( + "error loading schema for shard {}: {}", + shard.number(), + err + ); + sleep(Duration::from_millis(100)).await; + } else { + // Cluster is shutting down: unblock any + // wait_schema_loaded callers. + shard.schema_not_needed(); + return; + } + } + } + } + }); + } + } + + /// Rollback abandoned two-phase commit transactions on all shards. + /// + /// Bounded by `two_phase_commit_rollback_abandoned_timeout`. + async fn rollback_abandoned_two_pc(&self) { + let result = timeout( + self.two_pc_rollback_abandoned_timeout(), + Manager::get().cleanup_abandoned_for_cluster(self), + ) + .await; + + match result { + Ok(Ok(cleaned_up)) => { + if cleaned_up > 0 { + warn!( + r#"[2pc] rolled back {} abandoned transactions on database "{}""#, + cleaned_up, + self.name(), + ); + } + } + Ok(Err(err)) => { + error!( + r#"[2pc] abandoned transactions cleanup error on database "{}": {}"#, + self.name(), + err + ); + } + Err(_) => { + error!( + r#"[2pc] abandoned transactions cleanup on database "{}" timed out after {}ms"#, + self.name(), + self.two_pc_rollback_abandoned_timeout().as_millis() + ); + } + } + } +} diff --git a/pgdog/src/backend/pool/mod.rs b/pgdog/src/backend/pool/mod.rs index ce5635800..ef7ecce01 100644 --- a/pgdog/src/backend/pool/mod.rs +++ b/pgdog/src/backend/pool/mod.rs @@ -3,6 +3,7 @@ pub mod address; pub mod cleanup; pub mod cluster; +pub mod cluster_launch; pub mod comms; pub mod config; pub mod connection; diff --git a/pgdog/src/backend/replication/logical/orchestrator.rs b/pgdog/src/backend/replication/logical/orchestrator.rs index 598c1fb82..e09a20497 100644 --- a/pgdog/src/backend/replication/logical/orchestrator.rs +++ b/pgdog/src/backend/replication/logical/orchestrator.rs @@ -121,7 +121,7 @@ impl Orchestrator { self.destination = databases().schema_owner(&self.destination.identifier().database)?; self.source = databases().schema_owner(&self.source.identifier().database)?; - self.destination.wait_schema_loaded().await; + self.destination.wait_ready().await; self.refresh_publisher(); diff --git a/pgdog/src/frontend/client/query_engine/route_query.rs b/pgdog/src/frontend/client/query_engine/route_query.rs index 292dfde7f..090036ff8 100644 --- a/pgdog/src/frontend/client/query_engine/route_query.rs +++ b/pgdog/src/frontend/client/query_engine/route_query.rs @@ -50,12 +50,12 @@ impl QueryEngine { }; if let Ok(ClusterCheck::Ok) = res { - // Make sure schema is loaded before we throw traffic - // at it. This matters for sharded deployments only. + // Wait for boot-time maintenance (two-phase commit cleanup, + // schema load) before we throw traffic at the cluster. if let Ok(cluster) = self.backend.cluster() { safe_timeout( context.timeouts.query_timeout(&State::Active), - cluster.wait_schema_loaded(), + cluster.wait_ready(), ) .await .map_err(|_| Error::SchemaLoad)?; diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index ac649ffd4..625ce79cd 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -417,7 +417,10 @@ impl Manager { Ok(()) } - async fn cleanup_abandoned_for_cluster(&self, cluster: &Cluster) -> Result { + pub(crate) async fn cleanup_abandoned_for_cluster( + &self, + cluster: &Cluster, + ) -> Result { let cleaned_up = try_join_all( cluster .shards() diff --git a/pgdog/src/main.rs b/pgdog/src/main.rs index 6d3c66579..e4d8030e1 100644 --- a/pgdog/src/main.rs +++ b/pgdog/src/main.rs @@ -3,7 +3,6 @@ use std::fs::read_to_string; use std::path::Path; use std::process::exit; -use std::time::Duration; use clap::Parser; use pgdog::backend::databases; @@ -17,7 +16,6 @@ use pgdog::stats; use pgdog::util::pgdog_version; use pgdog::{healthcheck, net}; use tokio::runtime::Builder; -use tokio::time::timeout; use tracing::{error, info, warn}; fn main() -> Result<(), Box> { @@ -159,29 +157,8 @@ async fn pgdog(command: Option) -> Result<(), Box { - error!( - "[2pc] abandoned transactions cleanup timed out after {}ms", - abandoned_timeout.as_millis() - ); - } - - Ok(Err(err)) => { - error!("[2pc] abandoned transactions cleanup error: {}", err); - } - - Ok(_) => (), - }; - } - + // Abandoned two-phase commit transactions are cleaned up + // by each cluster on launch, before it serves traffic. let mut listener = Listener::new(format!("{}:{}", general.host, general.port)); listener.listen().await?; } From d176b639a7a40fc3e67b07c3a93a8d59d54d09ab Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Wed, 22 Jul 2026 19:03:22 -0700 Subject: [PATCH 14/22] fmt --- .schema/pgdog.schema.json | 2 +- pgdog/src/backend/pool/cluster_launch.rs | 22 ++++++++++++---------- 2 files changed, 13 insertions(+), 11 deletions(-) diff --git a/.schema/pgdog.schema.json b/.schema/pgdog.schema.json index a28b11526..e28b86550 100644 --- a/.schema/pgdog.schema.json +++ b/.schema/pgdog.schema.json @@ -1199,7 +1199,7 @@ "default": false }, "two_phase_commit_rollback_abandoned_timeout": { - "description": "Maximum amount of time to block startup in order to rollback abandoned\ntransactions.", + "description": "Maximum amount of time, in milliseconds, each cluster waits for abandoned\ntransactions rollback on boot before allowing traffic.", "type": "integer", "format": "uint64", "default": 15000, diff --git a/pgdog/src/backend/pool/cluster_launch.rs b/pgdog/src/backend/pool/cluster_launch.rs index 82ef2f8e3..dcc12b3d3 100644 --- a/pgdog/src/backend/pool/cluster_launch.rs +++ b/pgdog/src/backend/pool/cluster_launch.rs @@ -9,7 +9,7 @@ use std::time::Duration; use tokio::time::{sleep, timeout}; use tokio_util::sync::CancellationToken; -use tracing::{error, warn}; +use tracing::{error, info, warn}; use crate::backend::pool::ee::schema_changed_hook; use crate::frontend::client::query_engine::two_pc::Manager; @@ -204,28 +204,30 @@ impl Cluster { ) .await; + let identifier = self.identifier(); + match result { Ok(Ok(cleaned_up)) => { if cleaned_up > 0 { warn!( - r#"[2pc] rolled back {} abandoned transactions on database "{}""#, - cleaned_up, - self.name(), + "[2pc] rolled back {} abandoned transactions [{}]", + cleaned_up, identifier, ); + } else { + info!("[2pc] no abandoned transactions found [{}]", identifier,); } } Ok(Err(err)) => { error!( - r#"[2pc] abandoned transactions cleanup error on database "{}": {}"#, - self.name(), - err + "[2pc] abandoned transactions cleanup error: {} [{}]", + err, identifier ); } Err(_) => { error!( - r#"[2pc] abandoned transactions cleanup on database "{}" timed out after {}ms"#, - self.name(), - self.two_pc_rollback_abandoned_timeout().as_millis() + "[2pc] abandoned transactions cleanup timed out after {}ms [{}]", + self.two_pc_rollback_abandoned_timeout().as_millis(), + identifier, ); } } From f6a46edbfeffee69d19be39129f17e37412d4248 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Thu, 23 Jul 2026 05:20:27 -0700 Subject: [PATCH 15/22] remove two pc cleanup --- integration/pgdog.toml | 2 +- pgdog-config/src/general.rs | 23 ----- pgdog/src/backend/databases.rs | 44 --------- pgdog/src/backend/pool/cluster.rs | 73 ++++++--------- pgdog/src/backend/pool/cluster_launch.rs | 73 +-------------- .../client/query_engine/two_pc/manager.rs | 92 +------------------ .../two_pc/server_transactions.rs | 4 +- .../client/query_engine/two_pc/test.rs | 71 -------------- .../client/query_engine/two_pc/transaction.rs | 5 - 9 files changed, 35 insertions(+), 352 deletions(-) diff --git a/integration/pgdog.toml b/integration/pgdog.toml index d5fd64d75..61b43acfd 100644 --- a/integration/pgdog.toml +++ b/integration/pgdog.toml @@ -16,7 +16,7 @@ expanded_explain = true dns_ttl = 1_000 query_cache_limit = 500 pub_sub_channel_size = 4098 -two_phase_commit = false +two_phase_commit = true healthcheck_port = 8080 tls_certificate = "integration/tls/cert.pem" tls_private_key = "integration/tls/key.pem" diff --git a/pgdog-config/src/general.rs b/pgdog-config/src/general.rs index 09deb1ba6..5ad4cc018 100644 --- a/pgdog-config/src/general.rs +++ b/pgdog-config/src/general.rs @@ -622,18 +622,6 @@ pub struct General { #[serde(default = "General::two_phase_commit_wal_checkpoint_interval")] pub two_phase_commit_wal_checkpoint_interval: u64, - /// Rollback abandoned transactions. - /// - /// WARNING: Data loss will occur, enable this only if you don't care about consistency - /// and are not using the 2pc WAL. - #[serde(default = "General::two_phase_commit_rollback_abandoned")] - pub two_phase_commit_rollback_abandoned: bool, - - /// Maximum amount of time, in milliseconds, each cluster waits for abandoned - /// transactions rollback on boot before allowing traffic. - #[serde(default = "General::two_phase_commit_rollback_abandoned_timeout")] - pub two_phase_commit_rollback_abandoned_timeout: u64, - /// Enable expanded (`\x`) output for `EXPLAIN` results returned by PgDog's built-in query plan aggregation. #[serde(default = "General::expanded_explain")] pub expanded_explain: bool, @@ -896,9 +884,6 @@ impl Default for General { two_phase_commit_wal_fsync_interval: Self::two_phase_commit_wal_fsync_interval(), two_phase_commit_wal_checkpoint_interval: Self::two_phase_commit_wal_checkpoint_interval(), - two_phase_commit_rollback_abandoned: Self::two_phase_commit_rollback_abandoned(), - two_phase_commit_rollback_abandoned_timeout: - Self::two_phase_commit_rollback_abandoned_timeout(), expanded_explain: Self::expanded_explain(), server_lifetime: Self::server_lifetime(), server_lifetime_jitter: Self::server_lifetime_jitter(), @@ -1077,14 +1062,6 @@ impl General { Self::env_or_default("PGDOG_TWO_PHASE_COMMIT_WAL_CHECKPOINT_INTERVAL", 60) } - fn two_phase_commit_rollback_abandoned() -> bool { - Self::env_bool_or_default("PGDOG_TWO_PHASE_COMMIT_ROLLBACK_ABANDONED", false) - } - - fn two_phase_commit_rollback_abandoned_timeout() -> u64 { - Self::env_or_default("PGDOG_TWO_PHASE_COMMIT_ROLLBACK_ABANDONED_TIMEOUT", 15_000) - } - fn idle_timeout() -> u64 { Self::env_or_default( "PGDOG_IDLE_TIMEOUT", diff --git a/pgdog/src/backend/databases.rs b/pgdog/src/backend/databases.rs index d71b6e04d..11d553d94 100644 --- a/pgdog/src/backend/databases.rs +++ b/pgdog/src/backend/databases.rs @@ -1489,50 +1489,6 @@ mod tests { ); } - #[test] - fn test_two_pc_rollback_abandoned_config_passed_to_cluster() { - let config = Config { - databases: vec![Database { - name: "db1".to_string(), - host: "localhost".to_string(), - port: 5432, - role: Role::Primary, - ..Default::default() - }], - general: General { - two_phase_commit_rollback_abandoned: true, - two_phase_commit_rollback_abandoned_timeout: 1234, - ..Default::default() - }, - ..Default::default() - }; - - let users = crate::config::Users { - users: vec![crate::config::User { - name: "user".to_string(), - database: "db1".to_string(), - password: Some("pass".to_string()), - ..Default::default() - }], - ..Default::default() - }; - - let databases = from_config(&ConfigAndUsers { - config, - users, - config_path: std::path::PathBuf::new(), - users_path: std::path::PathBuf::new(), - ..Default::default() - }); - - let cluster = databases.cluster(("user", "db1")).unwrap(); - assert!(cluster.two_pc_rollback_abandoned()); - assert_eq!( - cluster.two_pc_rollback_abandoned_timeout(), - std::time::Duration::from_millis(1234) - ); - } - #[test] fn test_user_all_databases_creates_pools_for_all_dbs() { let config = Config { diff --git a/pgdog/src/backend/pool/cluster.rs b/pgdog/src/backend/pool/cluster.rs index 8bfc33a1e..204f5dd1a 100644 --- a/pgdog/src/backend/pool/cluster.rs +++ b/pgdog/src/backend/pool/cluster.rs @@ -56,8 +56,6 @@ pub struct Cluster { cross_shard_disabled: bool, two_phase_commit: bool, two_phase_commit_auto: bool, - two_pc_rollback_abandoned: bool, - two_pc_rollback_abandoned_timeout: Duration, pub(super) readiness: Arc, rewrite: Rewrite, prepared_statements: PreparedStatements, @@ -143,8 +141,6 @@ pub struct ClusterConfig<'a> { pub cross_shard_disabled: bool, pub two_pc: bool, pub two_pc_auto: bool, - pub two_pc_rollback_abandoned: bool, - pub two_pc_rollback_abandoned_timeout: u64, pub sharded_schemas: ShardedSchemas, pub rewrite: &'a Rewrite, pub prepared_statements: &'a PreparedStatements, @@ -208,8 +204,6 @@ impl<'a> ClusterConfig<'a> { two_pc_auto: user .two_phase_commit_auto .unwrap_or(general.two_phase_commit_auto.unwrap_or(false)), // Disable by default. - two_pc_rollback_abandoned: general.two_phase_commit_rollback_abandoned, - two_pc_rollback_abandoned_timeout: general.two_phase_commit_rollback_abandoned_timeout, sharded_schemas, rewrite, prepared_statements: &general.prepared_statements, @@ -257,8 +251,6 @@ impl Cluster { cross_shard_disabled, two_pc, two_pc_auto, - two_pc_rollback_abandoned, - two_pc_rollback_abandoned_timeout, sharded_schemas, rewrite, prepared_statements, @@ -320,10 +312,6 @@ impl Cluster { cross_shard_disabled, two_phase_commit: two_pc && shards.len() > 1, two_phase_commit_auto: two_pc_auto && shards.len() > 1, - two_pc_rollback_abandoned, - two_pc_rollback_abandoned_timeout: Duration::from_millis( - two_pc_rollback_abandoned_timeout, - ), readiness: Arc::new(Readiness::default()), rewrite: rewrite.clone(), prepared_statements: *prepared_statements, @@ -609,16 +597,7 @@ impl Cluster { self.two_phase_commit_auto && self.two_pc_enabled() } - /// Rollback abandoned two-phase commit transactions on launch. - pub(crate) fn two_pc_rollback_abandoned(&self) -> bool { - self.two_pc_rollback_abandoned - } - /// Maximum time abandoned transactions cleanup can take on launch - /// before traffic is allowed through. - pub(crate) fn two_pc_rollback_abandoned_timeout(&self) -> Duration { - self.two_pc_rollback_abandoned_timeout - } /// How many parallel COPY commands can we /// run to re-shard this cluster. @@ -1027,31 +1006,6 @@ mod test { assert!(cluster.ready()); } - #[tokio::test] - async fn test_wait_ready_waits_for_two_pc_cleanup() { - use tokio::time::{Duration, timeout}; - - let config = ConfigAndUsers::default(); - let mut cluster = Cluster::new_test(&config); - cluster.two_pc_rollback_abandoned = true; - // Expire immediately, don't touch a real database. - cluster.two_pc_rollback_abandoned_timeout = Duration::ZERO; - - // Pre-populate per-shard schemas, same reason. - for shard in &cluster.shards { - shard.schema_not_needed(); - } - - cluster.launch(); - - // Cleanup runs in the background, readiness can't be set yet. - assert!(!cluster.ready()); - - let result = timeout(Duration::from_millis(500), cluster.wait_ready()).await; - assert!(result.is_ok()); - assert!(cluster.ready()); - } - #[tokio::test] async fn test_shutdown_releases_readiness_waiters() { use tokio::time::{Duration, timeout}; @@ -1098,6 +1052,33 @@ mod test { cluster.wait_ready().await; } + #[tokio::test] + async fn test_wait_ready_waits_for_schema_notification() { + use tokio::time::{Duration, sleep, timeout}; + + let config = ConfigAndUsers::default(); + let cluster = Cluster::new_test(&config); + + cluster.launch(); + + // Schemas not loaded yet, readiness is pending. + assert!(!cluster.ready()); + + // Simulate schema load finishing on each shard: the readiness + // monitor wakes up via the per-shard schema_waiter notification. + let shards: Vec<_> = cluster.shards.to_vec(); + tokio::spawn(async move { + sleep(Duration::from_millis(10)).await; + for shard in &shards { + shard.schema_not_needed(); + } + }); + + let result = timeout(Duration::from_millis(500), cluster.wait_ready()).await; + assert!(result.is_ok()); + assert!(cluster.ready()); + } + #[tokio::test] async fn test_wait_ready_returns_immediately_when_schema_not_needed() { let config = ConfigAndUsers::default(); diff --git a/pgdog/src/backend/pool/cluster_launch.rs b/pgdog/src/backend/pool/cluster_launch.rs index dcc12b3d3..4876ba2a5 100644 --- a/pgdog/src/backend/pool/cluster_launch.rs +++ b/pgdog/src/backend/pool/cluster_launch.rs @@ -7,12 +7,11 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; -use tokio::time::{sleep, timeout}; +use tokio::time::sleep; use tokio_util::sync::CancellationToken; -use tracing::{error, info, warn}; +use tracing::error; use crate::backend::pool::ee::schema_changed_hook; -use crate::frontend::client::query_engine::two_pc::Manager; use crate::tasks; use super::Cluster; @@ -27,7 +26,6 @@ use super::Cluster; pub(super) struct Readiness { online: AtomicBool, launch_waiter: CancellationToken, - two_pc_waiter: CancellationToken, } impl Readiness { @@ -51,14 +49,6 @@ impl Readiness { async fn wait_ready(&self) { self.launch_waiter.cancelled().await; } - - fn mark_two_pc_cleaned_up(&self) { - self.two_pc_waiter.cancel(); - } - - async fn wait_two_pc_cleaned_up(&self) { - self.two_pc_waiter.cancelled().await; - } } impl Cluster { @@ -70,7 +60,6 @@ impl Cluster { self.readiness.set_online(true); - self.launch_two_pc_cleanup(); self.launch_schema_sync(); self.launch_readiness_monitor(); } @@ -113,28 +102,11 @@ impl Cluster { self.readiness.mark_ready(); } - /// Rollback abandoned two-phase commit transactions before - /// serving traffic. Queries wait on [`Cluster::wait_ready`] - /// until this is done. - fn launch_two_pc_cleanup(&self) { - if !self.two_pc_rollback_abandoned() { - self.readiness.mark_two_pc_cleaned_up(); - return; - } - - let cluster = self.clone(); - tasks::spawn("two-pc abandoned transactions cleanup", async move { - cluster.rollback_abandoned_two_pc().await; - cluster.readiness.mark_two_pc_cleaned_up(); - }); - } - /// Mark the cluster ready once two-phase commit cleanup /// and schema loading are done. fn launch_readiness_monitor(&self) { let cluster = self.clone(); tasks::spawn("cluster readiness monitor", async move { - cluster.readiness.wait_two_pc_cleaned_up().await; cluster.wait_schema_loaded().await; cluster.mark_ready(); }); @@ -165,7 +137,7 @@ impl Cluster { let identifier = self.identifier(); let shard = shard.clone(); - tasks::spawn("cluster schema sync", async move { + tasks::spawn("shard schema sync", async move { loop { match shard.load_schema().await { Ok(true) => { @@ -193,43 +165,4 @@ impl Cluster { }); } } - - /// Rollback abandoned two-phase commit transactions on all shards. - /// - /// Bounded by `two_phase_commit_rollback_abandoned_timeout`. - async fn rollback_abandoned_two_pc(&self) { - let result = timeout( - self.two_pc_rollback_abandoned_timeout(), - Manager::get().cleanup_abandoned_for_cluster(self), - ) - .await; - - let identifier = self.identifier(); - - match result { - Ok(Ok(cleaned_up)) => { - if cleaned_up > 0 { - warn!( - "[2pc] rolled back {} abandoned transactions [{}]", - cleaned_up, identifier, - ); - } else { - info!("[2pc] no abandoned transactions found [{}]", identifier,); - } - } - Ok(Err(err)) => { - error!( - "[2pc] abandoned transactions cleanup error: {} [{}]", - err, identifier - ); - } - Err(_) => { - error!( - "[2pc] abandoned transactions cleanup timed out after {}ms [{}]", - self.two_pc_rollback_abandoned_timeout().as_millis(), - identifier, - ); - } - } - } } diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index 625ce79cd..a7da1f3fa 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -1,7 +1,6 @@ //! Global two-phase commit transaction manager. use arc_swap::ArcSwapOption; use fnv::FnvHashMap as HashMap; -use futures::future::try_join_all; use once_cell::sync::Lazy; use parking_lot::Mutex; use std::{ @@ -22,15 +21,14 @@ use tracing::{debug, error, info, warn}; use crate::{ backend::{ - Cluster, - databases::{User, databases}, + databases::User, pool::{Connection, Request}, }, config::config, frontend::{ client::query_engine::{ TwoPcPhase, - two_pc::{TwoPcGuard, TwoPcStats, TwoPcTransaction, TwoPcTransactions, wal::Wal}, + two_pc::{TwoPcGuard, TwoPcStats, TwoPcTransaction, wal::Wal}, }, router::{ Route, @@ -40,7 +38,7 @@ use crate::{ tasks, }; -use super::{Error, server_transactions::TwoPcServerTransaction, statement::phase_control}; +use super::Error; static MANAGER: Lazy = Lazy::new(Manager::init); static MAINTENANCE: Duration = Duration::from_millis(333); @@ -389,90 +387,6 @@ impl Manager { Ok(()) } - /// Drop abandoned two-phase commit transactions. - /// - /// WARNING: This only happens if durability for 2pc is off. Running this - /// will cause data loss. - /// - pub async fn cleanup_abandoned(&self) -> Result<(), Error> { - info!("[2pc] rolling back abandoned transactions"); - - let databases = databases(); - let cleaned_up = try_join_all( - databases - .all() - .values() - .map(|cluster| self.cleanup_abandoned_for_cluster(cluster)), - ) - .await? - .into_iter() - .sum::(); - - if cleaned_up > 0 { - warn!("[2pc] rolled back up {} abandoned transactions", cleaned_up); - } else { - info!("[2pc] no abandoned transactions found"); - } - - Ok(()) - } - - pub(crate) async fn cleanup_abandoned_for_cluster( - &self, - cluster: &Cluster, - ) -> Result { - let cleaned_up = try_join_all( - cluster - .shards() - .iter() - .enumerate() - .map(|(number, shard)| Self::cleanup_abandoned_for_shard(number, shard)), - ) - .await?; - - Ok(cleaned_up.into_iter().sum()) - } - - async fn cleanup_abandoned_for_shard( - number: usize, - shard: &crate::backend::pool::Shard, - ) -> Result { - use crate::backend::Error as BackendError; - use crate::backend::pool::Error as PoolError; - let mut cleaned_up = 0; - - let mut conn = match shard.primary(&Request::default()).await { - Ok(conn) => conn, - Err(PoolError::NoPrimary) => return Ok(cleaned_up), - Err(err) => return Err(BackendError::Pool(err).into()), - }; - - let txns = TwoPcTransactions::load(&mut conn).await?; - - for txn in txns.iter() { - if let TwoPcServerTransaction::Ours { txn, .. } = txn { - // Postgres transactions can be only be rolled back by their owners. - if txn.is_created_by_this_process() { - match conn - .execute(phase_control(*txn, number, TwoPcPhase::Rollback)) - .await - { - Ok(_) => cleaned_up += 1, - Err(BackendError::ExecutionError(err)) => { - warn!( - "[2pc] error cleaning abandoned transaction \"{}\": {}", - txn, err - ); - } - Err(err) => return Err(err.into()), - } - } - } - } - - Ok(cleaned_up) - } - /// Shutdown manager and wait for all transactions to be cleaned up. /// Once the monitor has drained the cleanup queue, the WAL is shut /// down too so any final End records make it to disk before exit. diff --git a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs index 986e35890..b2a3445ef 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/server_transactions.rs @@ -37,9 +37,7 @@ impl TwoPcTransactions { /// Load two phase transactions from the server. pub(crate) async fn load(server: &mut Server) -> Result { let records: Vec = server - .fetch_all( - "SELECT gid, owner, database FROM pg_prepared_xacts WHERE owner = current_user", - ) + .fetch_all("SELECT gid, owner, database FROM pg_prepared_xacts") .await?; let mut transactions = vec![]; diff --git a/pgdog/src/frontend/client/query_engine/two_pc/test.rs b/pgdog/src/frontend/client/query_engine/two_pc/test.rs index 6db21efe7..0ad10f27d 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/test.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/test.rs @@ -1,6 +1,5 @@ use crate::{ backend::{ - Server, databases::databases, pool::{Connection, Request}, }, @@ -14,76 +13,6 @@ use crate::{ }; use super::*; -use super::{server_transactions::TwoPcServerTransaction, statement::phase_control}; - -async fn server_has_transaction(server: &mut Server, transaction: TwoPcTransaction) -> bool { - TwoPcTransactions::load(server) - .await - .unwrap() - .iter() - .any(|server_transaction| { - matches!( - server_transaction, - TwoPcServerTransaction::Ours { txn, .. } if *txn == transaction - ) - }) -} - -#[tokio::test] -async fn test_cleanup_abandoned() { - config::load_test_with_user("pgdog"); - let cluster = databases().all().iter().next().unwrap().1.clone(); - let transaction = TwoPcTransaction::new(); - let mut conn = cluster.shards()[0] - .primary(&Request::default()) - .await - .unwrap(); - - conn.execute("BEGIN").await.unwrap(); - conn.execute(phase_control(transaction, 0, TwoPcPhase::Phase1)) - .await - .unwrap(); - - assert!(server_has_transaction(&mut conn, transaction).await); - - Manager::get().cleanup_abandoned().await.unwrap(); - - assert!(!server_has_transaction(&mut conn, transaction).await); -} - -#[tokio::test] -async fn test_cleanup_abandoned_different_user() { - config::load_test_with_user("pgdog1"); - let cluster = databases().all().iter().next().unwrap().1.clone(); - let transaction = TwoPcTransaction::new(); - let mut conn = cluster.shards()[0] - .primary(&Request::default()) - .await - .unwrap(); - - conn.execute("BEGIN").await.unwrap(); - conn.execute(phase_control(transaction, 0, TwoPcPhase::Phase1)) - .await - .unwrap(); - assert!(server_has_transaction(&mut conn, transaction).await); - drop(conn); - - config::load_test_with_user("pgdog2"); - Manager::get().cleanup_abandoned().await.unwrap(); - - config::load_test_with_user("pgdog1"); - let cluster = databases().all().iter().next().unwrap().1.clone(); - let mut conn = cluster.shards()[0] - .primary(&Request::default()) - .await - .unwrap(); - let was_preserved = server_has_transaction(&mut conn, transaction).await; - conn.execute(phase_control(transaction, 0, TwoPcPhase::Rollback)) - .await - .unwrap(); - - assert!(was_preserved); -} #[tokio::test] async fn test_cleanup_transaction_phase_one() { diff --git a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs index 0e21ac396..4466fa650 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/transaction.rs @@ -16,11 +16,6 @@ impl TwoPcTransaction { Self(rng().random_range(0..usize::MAX)) } - /// This transaction was created by this process. - pub(crate) fn is_created_by_this_process(&self) -> bool { - self.to_string().starts_with(&Self::global_prefix()) - } - /// A prefix to identify two-phase commit transactions generated /// by this PgDog process. fn global_prefix() -> String { From 2d2ab52255e31cac16876d7c980104936720510d Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Thu, 23 Jul 2026 05:21:17 -0700 Subject: [PATCH 16/22] schema --- .schema/pgdog.schema.json | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/.schema/pgdog.schema.json b/.schema/pgdog.schema.json index e28b86550..3900e97cd 100644 --- a/.schema/pgdog.schema.json +++ b/.schema/pgdog.schema.json @@ -117,8 +117,6 @@ "tls_verify": "prefer", "two_phase_commit": false, "two_phase_commit_auto": null, - "two_phase_commit_rollback_abandoned": false, - "two_phase_commit_rollback_abandoned_timeout": 15000, "two_phase_commit_wal_checkpoint_interval": 60, "two_phase_commit_wal_dir": null, "two_phase_commit_wal_fsync_interval": 2, @@ -1193,18 +1191,6 @@ ], "default": null }, - "two_phase_commit_rollback_abandoned": { - "description": "Rollback abandoned transactions.\n\nWARNING: Data loss will occur, enable this only if you don't care about consistency\nand are not using the 2pc WAL.", - "type": "boolean", - "default": false - }, - "two_phase_commit_rollback_abandoned_timeout": { - "description": "Maximum amount of time, in milliseconds, each cluster waits for abandoned\ntransactions rollback on boot before allowing traffic.", - "type": "integer", - "format": "uint64", - "default": 15000, - "minimum": 0 - }, "two_phase_commit_wal_checkpoint_interval": { "description": "How often, in seconds, to write a checkpoint record to the two-phase commit WAL and garbage-collect old segments.\n\n_Default:_ `60`\n\n", "type": "integer", From 4de458a0a2e3b4c64a74309f1c7480d6b5033f1b Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Thu, 23 Jul 2026 05:38:23 -0700 Subject: [PATCH 17/22] fix --- integration/pgdog.toml | 2 +- pgdog/src/backend/pool/cluster.rs | 2 -- pgdog/src/backend/pool/cluster_launch.rs | 30 ++++++++++++++---------- 3 files changed, 18 insertions(+), 16 deletions(-) diff --git a/integration/pgdog.toml b/integration/pgdog.toml index 61b43acfd..d5fd64d75 100644 --- a/integration/pgdog.toml +++ b/integration/pgdog.toml @@ -16,7 +16,7 @@ expanded_explain = true dns_ttl = 1_000 query_cache_limit = 500 pub_sub_channel_size = 4098 -two_phase_commit = true +two_phase_commit = false healthcheck_port = 8080 tls_certificate = "integration/tls/cert.pem" tls_private_key = "integration/tls/key.pem" diff --git a/pgdog/src/backend/pool/cluster.rs b/pgdog/src/backend/pool/cluster.rs index 204f5dd1a..22328d881 100644 --- a/pgdog/src/backend/pool/cluster.rs +++ b/pgdog/src/backend/pool/cluster.rs @@ -597,8 +597,6 @@ impl Cluster { self.two_phase_commit_auto && self.two_pc_enabled() } - /// Maximum time abandoned transactions cleanup can take on launch - /// How many parallel COPY commands can we /// run to re-shard this cluster. pub fn resharding_parallel_copies(&self) -> usize { diff --git a/pgdog/src/backend/pool/cluster_launch.rs b/pgdog/src/backend/pool/cluster_launch.rs index 4876ba2a5..16e98d7b3 100644 --- a/pgdog/src/backend/pool/cluster_launch.rs +++ b/pgdog/src/backend/pool/cluster_launch.rs @@ -1,13 +1,12 @@ //! Cluster startup and shutdown primitives. //! //! Launching and shutting down the connection pools, and gating -//! traffic until boot-time maintenance (two-phase commit cleanup, -//! schema sync) has completed. +//! traffic until boot-time maintenance has completed. use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; -use tokio::time::sleep; +use tokio::{select, time::sleep}; use tokio_util::sync::CancellationToken; use tracing::error; @@ -81,12 +80,7 @@ impl Cluster { self.readiness.online() } - /// Wait until boot-time maintenance — two-phase commit cleanup - /// and schema loading — is done and the cluster can serve traffic. - /// - /// Two-phase commit cleanup is bounded by - /// `two_phase_commit_rollback_abandoned_timeout`; schema loading - /// retries until success, so callers should apply their own timeout. + /// Wait until boot-time maintenance is done and the cluster can serve traffic. pub(crate) async fn wait_ready(&self) { self.readiness.wait_ready().await; } @@ -102,12 +96,15 @@ impl Cluster { self.readiness.mark_ready(); } - /// Mark the cluster ready once two-phase commit cleanup - /// and schema loading are done. + /// Mark the cluster ready schema loading is done. fn launch_readiness_monitor(&self) { let cluster = self.clone(); tasks::spawn("cluster readiness monitor", async move { - cluster.wait_schema_loaded().await; + let shutdown = tasks::shutdown_signal(); + select! { + _ = cluster.wait_schema_loaded() => {} + _ = shutdown.cancelled() => {} + } cluster.mark_ready(); }); } @@ -136,10 +133,17 @@ impl Cluster { for shard in self.shards() { let identifier = self.identifier(); let shard = shard.clone(); + let shutdown = tasks::shutdown_signal(); tasks::spawn("shard schema sync", async move { loop { - match shard.load_schema().await { + let loader = shard.load_schema(); + let result = select! { + _ = shutdown.cancelled() => break, + result = loader => { result }, + }; + + match result { Ok(true) => { schema_changed_hook(&shard.schema(), &identifier, &shard); return; From 4e8eb9498bf22b873dab34b5d9799ada87d0fda9 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Thu, 23 Jul 2026 05:50:12 -0700 Subject: [PATCH 18/22] wrong comments --- pgdog/src/backend/pool/cluster.rs | 3 +-- pgdog/src/backend/pool/cluster_launch.rs | 5 ----- pgdog/src/frontend/client/query_engine/two_pc/manager.rs | 2 -- 3 files changed, 1 insertion(+), 9 deletions(-) diff --git a/pgdog/src/backend/pool/cluster.rs b/pgdog/src/backend/pool/cluster.rs index 22328d881..ff843a61c 100644 --- a/pgdog/src/backend/pool/cluster.rs +++ b/pgdog/src/backend/pool/cluster.rs @@ -1087,8 +1087,7 @@ mod test { cluster.launch(); - // Should return without waiting: no schema to load, - // two-pc cleanup disabled. + // Should return without waiting: no schema to load. cluster.wait_ready().await; assert!(cluster.ready()); } diff --git a/pgdog/src/backend/pool/cluster_launch.rs b/pgdog/src/backend/pool/cluster_launch.rs index 16e98d7b3..2f89f6223 100644 --- a/pgdog/src/backend/pool/cluster_launch.rs +++ b/pgdog/src/backend/pool/cluster_launch.rs @@ -16,11 +16,6 @@ use crate::tasks; use super::Cluster; /// Cluster readiness state. -/// -/// `online` is set once the pools are launched and cleared on shutdown. -/// `launch_waiter` is cancelled once boot-time maintenance — two-phase -/// commit cleanup and schema loading — is done and the cluster can -/// serve traffic. #[derive(Default, Debug)] pub(super) struct Readiness { online: AtomicBool, diff --git a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs index a7da1f3fa..ff22ffb23 100644 --- a/pgdog/src/frontend/client/query_engine/two_pc/manager.rs +++ b/pgdog/src/frontend/client/query_engine/two_pc/manager.rs @@ -295,8 +295,6 @@ impl Manager { debug!("[2pc] monitor started"); - // Cleanup orphaned transactions. - loop { // Wake up either because it's time to check // or manager told us to. From 01a03a4744a630b1ab6c23a9c0e94e2265c60577 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Thu, 23 Jul 2026 05:51:48 -0700 Subject: [PATCH 19/22] correct naming --- pgdog/src/frontend/client/query_engine/route_query.rs | 5 ++--- pgdog/src/frontend/error.rs | 4 ++-- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/pgdog/src/frontend/client/query_engine/route_query.rs b/pgdog/src/frontend/client/query_engine/route_query.rs index 090036ff8..74d6be3e1 100644 --- a/pgdog/src/frontend/client/query_engine/route_query.rs +++ b/pgdog/src/frontend/client/query_engine/route_query.rs @@ -50,15 +50,14 @@ impl QueryEngine { }; if let Ok(ClusterCheck::Ok) = res { - // Wait for boot-time maintenance (two-phase commit cleanup, - // schema load) before we throw traffic at the cluster. + // Wait for boot-time maintenance before we throw traffic at the cluster. if let Ok(cluster) = self.backend.cluster() { safe_timeout( context.timeouts.query_timeout(&State::Active), cluster.wait_ready(), ) .await - .map_err(|_| Error::SchemaLoad)?; + .map_err(|_| Error::ClusterStart)?; } res } else { diff --git a/pgdog/src/frontend/error.rs b/pgdog/src/frontend/error.rs index 7469fbf7b..f0325badc 100644 --- a/pgdog/src/frontend/error.rs +++ b/pgdog/src/frontend/error.rs @@ -46,8 +46,8 @@ pub enum Error { #[error("query timeout")] Timeout(#[from] tokio::time::error::Elapsed), - #[error("schema load timeout")] - SchemaLoad, + #[error("cluster start timeout")] + ClusterStart, #[error("join error")] Join(#[from] tokio::task::JoinError), From 369f1078603dd539ae135c2fee8bb140d72cf197 Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Thu, 23 Jul 2026 05:52:20 -0700 Subject: [PATCH 20/22] remove bs comment --- pgdog/src/main.rs | 2 -- 1 file changed, 2 deletions(-) diff --git a/pgdog/src/main.rs b/pgdog/src/main.rs index e4d8030e1..f7a995023 100644 --- a/pgdog/src/main.rs +++ b/pgdog/src/main.rs @@ -157,8 +157,6 @@ async fn pgdog(command: Option) -> Result<(), Box Date: Thu, 23 Jul 2026 05:54:52 -0700 Subject: [PATCH 21/22] mark cluster offline before shutting down shards --- pgdog/src/backend/pool/cluster_launch.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pgdog/src/backend/pool/cluster_launch.rs b/pgdog/src/backend/pool/cluster_launch.rs index 2f89f6223..a8536913f 100644 --- a/pgdog/src/backend/pool/cluster_launch.rs +++ b/pgdog/src/backend/pool/cluster_launch.rs @@ -60,12 +60,12 @@ impl Cluster { /// Shutdown the connection pools. pub(crate) fn shutdown(&self) { + self.readiness.set_online(false); + for shard in self.shards() { shard.shutdown(); } - self.readiness.set_online(false); - // Release readiness waiters, the cluster is going away. self.mark_ready(); } From ae77c993cf057dbe89cc6056bfd0167794c20b0e Mon Sep 17 00:00:00 2001 From: Lev Kokotov Date: Thu, 23 Jul 2026 07:40:32 -0700 Subject: [PATCH 22/22] try to get coverage for integration runs --- .github/workflows/ci-new-parser.yml | 21 ++--- .github/workflows/ci.yml | 83 ++++++++++++++------ integration/ci/install-deps.sh | 7 ++ integration/ci/prepare-instrumented-pgdog.sh | 24 +++++- pgdog/src/main.rs | 22 ++++++ 5 files changed, 115 insertions(+), 42 deletions(-) diff --git a/.github/workflows/ci-new-parser.yml b/.github/workflows/ci-new-parser.yml index 76b8131e9..2e2a11c38 100644 --- a/.github/workflows/ci-new-parser.yml +++ b/.github/workflows/ci-new-parser.yml @@ -25,9 +25,7 @@ jobs: id: pgdog-bin uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} - uses: Swatinem/rust-cache@v2 if: steps.pgdog-bin.outputs.cache-hit != 'true' @@ -36,16 +34,11 @@ jobs: - name: Build (debug) if: steps.pgdog-bin.outputs.cache-hit != 'true' run: cargo build --no-default-features --features new_parser --bin pgdog - - name: Build (release) - if: steps.pgdog-bin.outputs.cache-hit != 'true' - run: cargo build --release --no-default-features --features new_parser --bin pgdog - name: Save pgdog binaries if: steps.pgdog-bin.outputs.cache-hit != 'true' uses: actions/cache/save@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} ci: @@ -79,7 +72,7 @@ jobs: # cached pgdog binary. - { name: plugins, script: integration/plugins/run.sh, needs_rust_cache: true, continue_on_error: true } env: - PGDOG_BIN: ${{ github.workspace }}/target/release/pgdog + PGDOG_BIN: ${{ github.workspace }}/target/debug/pgdog PGDOG_PLUGIN_FEATURES: new_parser steps: - uses: actions/checkout@v6 @@ -89,8 +82,8 @@ jobs: id: cache-key run: echo "key=pgdog-bin-new-parser-${{ runner.os }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" # rust-cache must run before the binary restore: it lays down a - # stale target/ that can otherwise wipe target/release/pgdog when - # cargo reconciles fingerprints during plugin builds. + # stale target/ that can otherwise wipe target binaries when cargo + # reconciles fingerprints during plugin builds. - name: Restore Rust cache for plugin builds if: matrix.needs_rust_cache uses: Swatinem/rust-cache@v2 @@ -99,9 +92,7 @@ jobs: - name: Restore pgdog binaries uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} fail-on-cache-miss: true - name: Setup dependencies diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index c2544e4d9..6f803e378 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -7,8 +7,25 @@ on: types: [opened, synchronize, reopened] jobs: + # Build both pgdog artifacts in parallel: + # - release: plain release binary + # - coverage: llvm-cov instrumented debug binary used by the + # integration suites to produce line coverage for Codecov build: runs-on: blacksmith-4vcpu-ubuntu-2404 + strategy: + fail-fast: false + matrix: + include: + - artifact: release + path: target/release/pgdog + rust_prefix: build-v1 + - artifact: coverage + path: target/llvm-cov-target/debug/pgdog + # Separate prefix: instrumented artifacts live in + # target/llvm-cov-target and must not race with the plain + # build's cache. + rust_prefix: build-cov-debug-v1 steps: - uses: actions/checkout@v6 - name: Install CI deps @@ -17,35 +34,31 @@ jobs: # target/ build outputs that would otherwise change the key. Use the # split cache actions so pgdog binaries are saved immediately after the # build, before rust-cache's post-job cleanup prunes workspace binaries - # from target/{debug,release}. + # from target/. - name: Compute cache key id: cache-key - run: echo "key=pgdog-bin-${{ runner.os }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" - - name: Restore pgdog binaries + run: echo "key=pgdog-bin-${{ runner.os }}-${{ matrix.artifact }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" + - name: Restore pgdog binary id: pgdog-bin uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: ${{ matrix.path }} key: ${{ steps.cache-key.outputs.key }} - uses: Swatinem/rust-cache@v2 if: steps.pgdog-bin.outputs.cache-hit != 'true' with: - prefix-key: build-v1 - - name: Build (debug) - if: steps.pgdog-bin.outputs.cache-hit != 'true' - run: cargo build --bin pgdog + prefix-key: ${{ matrix.rust_prefix }} - name: Build (release) - if: steps.pgdog-bin.outputs.cache-hit != 'true' + if: steps.pgdog-bin.outputs.cache-hit != 'true' && matrix.artifact == 'release' run: cargo build --release --bin pgdog - - name: Save pgdog binaries + - name: Build (instrumented debug, for coverage) + if: steps.pgdog-bin.outputs.cache-hit != 'true' && matrix.artifact == 'coverage' + run: bash integration/ci/prepare-instrumented-pgdog.sh debug + - name: Save pgdog binary if: steps.pgdog-bin.outputs.cache-hit != 'true' uses: actions/cache/save@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: ${{ matrix.path }} key: ${{ steps.cache-key.outputs.key }} ci: @@ -79,30 +92,39 @@ jobs: # cached pgdog binary. - { name: plugins, script: integration/plugins/run.sh, needs_rust_cache: true, continue_on_error: true } env: - PGDOG_BIN: ${{ github.workspace }}/target/release/pgdog + # All suites run against the llvm-cov instrumented debug binary so we + # get line coverage from integration tests, not just unit tests. + PGDOG_BIN: ${{ github.workspace }}/target/llvm-cov-target/debug/pgdog + # NB: cargo llvm-cov report only globs profraw files at the top of + # target/llvm-cov-target, so LLVM_PROFILE_FILE must point there. + LLVM_PROFILE_FILE: ${{ github.workspace }}/target/llvm-cov-target/${{ matrix.name }}-%p-%m.profraw steps: - uses: actions/checkout@v6 - name: Install CI deps run: bash integration/ci/install-deps.sh - name: Compute cache key id: cache-key - run: echo "key=pgdog-bin-${{ runner.os }}-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" + run: echo "key=pgdog-bin-${{ runner.os }}-coverage-$(bash integration/ci/cache-key.sh)" >> "$GITHUB_OUTPUT" # rust-cache must run before the binary restore: it lays down a - # stale target/ that can otherwise wipe target/release/pgdog when - # cargo reconciles fingerprints during plugin builds. + # stale target/ that can otherwise wipe target binaries when cargo + # reconciles fingerprints during plugin builds. - name: Restore Rust cache for plugin builds if: matrix.needs_rust_cache uses: Swatinem/rust-cache@v2 with: prefix-key: build-v1 - - name: Restore pgdog binaries + - name: Restore instrumented pgdog binary uses: actions/cache/restore@v5 with: - path: | - target/debug/pgdog - target/release/pgdog + path: target/llvm-cov-target/debug/pgdog key: ${{ steps.cache-key.outputs.key }} fail-on-cache-miss: true + # The Rust cache (plugins job) can lay down stale .profraw files that + # don't match the restored instrumented binary; start clean. + - name: Prepare coverage profile dir + run: | + mkdir -p target/llvm-cov-target + rm -f target/llvm-cov-target/*.profraw target/llvm-cov-target/*.profdata - name: Setup dependencies run: bash integration/ci/setup.sh --with-toxi - name: Run ${{ matrix.name }} @@ -110,3 +132,18 @@ jobs: - name: Ensure PgDog stopped if: always() run: bash integration/ci/ensure-pgdog-stopped.sh + - name: Generate coverage report + if: always() + run: cargo llvm-cov report --package pgdog --lcov --output-path ${{ matrix.name }}.lcov + # Codecov merges uploads per flag server-side, so each suite uploads + # its own lcov under the shared "integration" flag. + - name: Upload coverage to Codecov + if: always() + uses: codecov/codecov-action@v4 + env: + CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }} + with: + files: ${{ matrix.name }}.lcov + flags: integration + name: integration-${{ matrix.name }} + fail_ci_if_error: false diff --git a/integration/ci/install-deps.sh b/integration/ci/install-deps.sh index 8716c7a25..b30cea9f0 100755 --- a/integration/ci/install-deps.sh +++ b/integration/ci/install-deps.sh @@ -54,3 +54,10 @@ if ! command -v cargo-llvm-cov >/dev/null; then curl -LsSf "https://github.com/taiki-e/cargo-llvm-cov/releases/download/v${LLVM_COV_VERSION}/cargo-llvm-cov-x86_64-unknown-linux-gnu.tar.gz" \ | tar zxf - -C "$CARGO_BIN" fi + +# cargo llvm-cov report needs llvm-profdata/llvm-cov from llvm-tools-preview. +if command -v rustup >/dev/null 2>&1; then + if ! rustup component list --installed 2>/dev/null | grep -q llvm-tools; then + rustup component add llvm-tools-preview || true + fi +fi diff --git a/integration/ci/prepare-instrumented-pgdog.sh b/integration/ci/prepare-instrumented-pgdog.sh index eeea163fe..38760a4d5 100755 --- a/integration/ci/prepare-instrumented-pgdog.sh +++ b/integration/ci/prepare-instrumented-pgdog.sh @@ -1,20 +1,36 @@ #!/usr/bin/env bash # Build an llvm-cov-instrumented pgdog binary for the integration job and # export its path via $GITHUB_ENV so later steps can invoke it. +# +# Usage: prepare-instrumented-pgdog.sh [debug|release] (default: release) set -euo pipefail -export RUSTFLAGS="-C link-dead-code" +PROFILE="${1:-release}" +if [[ "${PROFILE}" != "debug" && "${PROFILE}" != "release" ]]; then + echo "Usage: $0 [debug|release]" >&2 + exit 1 +fi + +PROFILE_FLAG="" +if [[ "${PROFILE}" == "release" ]]; then + PROFILE_FLAG="--release" +fi + +# RUSTFLAGS overrides .cargo/config.toml rustflags entirely, so +# tokio_unstable (set there) must be repeated here. +export RUSTFLAGS="--cfg tokio_unstable -C link-dead-code" cargo llvm-cov clean --workspace mkdir -p target/llvm-cov-target/profiles -cargo llvm-cov run --no-report --release --package pgdog --bin pgdog -- --help +# shellcheck disable=SC2086 +cargo llvm-cov run --no-report ${PROFILE_FLAG} --package pgdog --bin pgdog -- --help rm -f target/llvm-cov-target/profiles/*.profraw rm -f target/llvm-cov-target/profiles/.last_snapshot rm -rf target/llvm-cov-target/reports -BIN_PATH=$(find target/llvm-cov-target -type f -path '*/release/pgdog' | head -n 1) +BIN_PATH=$(find target/llvm-cov-target -type f -path "*/${PROFILE}/pgdog" | head -n 1) if [ -z "$BIN_PATH" ]; then - echo "Instrumented PgDog binary not found" >&2 + echo "Instrumented PgDog binary (${PROFILE}) not found" >&2 exit 1 fi echo "Using instrumented binary at $BIN_PATH" diff --git a/pgdog/src/main.rs b/pgdog/src/main.rs index f7a995023..684c87e55 100644 --- a/pgdog/src/main.rs +++ b/pgdog/src/main.rs @@ -105,6 +105,10 @@ fn main() -> Result<(), Box> { } async fn pgdog(command: Option) -> Result<(), Box> { + // Run atexit handlers on SIGTERM (e.g. llvm-cov profile flushing). + #[cfg(unix)] + install_sigterm_handler(); + // Preload TLS. Resulting primitives // are async, so doing this after Tokio launched seems prudent. net::tls::load()?; @@ -236,6 +240,24 @@ async fn pgdog(command: Option) -> Result<(), Box std::io::Result { match workers { 0 => Builder::new_current_thread()