Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions pgdog/src/backend/pool/connection/binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -365,7 +365,7 @@ impl Binding {

pub(crate) async fn two_pc_on_guards(
servers: &mut [Guard],
transaction: TwoPcTransaction,
transaction: &TwoPcTransaction,
phase: TwoPcPhase,
ignore_missing: bool,
) -> Result<(), Error> {
Expand Down Expand Up @@ -399,7 +399,7 @@ impl Binding {
/// Execute two-phase commit transaction control statements.
pub(crate) async fn two_pc(
&mut self,
transaction: TwoPcTransaction,
transaction: &TwoPcTransaction,
phase: TwoPcPhase,
ignore_missing: bool,
) -> Result<(), Error> {
Expand Down
50 changes: 29 additions & 21 deletions pgdog/src/backend/pool/connection/binding_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ mod tests {
let mut binding = Binding::Direct(guard, 0);

let result = binding
.two_pc(TwoPcTransaction::new(), TwoPcPhase::Phase1, false)
.two_pc(&TwoPcTransaction::new(), TwoPcPhase::Phase1, false)
.await;

// Should fail with TwoPcMultiShardOnly error
Expand All @@ -97,7 +97,7 @@ mod tests {
let mut binding = Binding::Admin(admin_server);

let result = binding
.two_pc(TwoPcTransaction::new(), TwoPcPhase::Phase1, false)
.two_pc(&TwoPcTransaction::new(), TwoPcPhase::Phase1, false)
.await;

// Should fail with TwoPcMultiShardOnly error
Expand All @@ -113,7 +113,9 @@ mod tests {
let transaction = TwoPcTransaction::new();

// Test Phase1 - PREPARE TRANSACTION
let result = binding.two_pc(transaction, TwoPcPhase::Phase1, false).await;
let result = binding
.two_pc(&transaction, TwoPcPhase::Phase1, false)
.await;

// Should succeed
if let Err(ref error) = result {
Expand All @@ -123,7 +125,7 @@ mod tests {

// Cleanup: Rollback the prepared transaction to avoid leaving dangling transactions
let _cleanup = binding
.two_pc(transaction, TwoPcPhase::Rollback, false)
.two_pc(&transaction, TwoPcPhase::Rollback, false)
.await;
}

Expand All @@ -134,12 +136,14 @@ mod tests {

// First prepare the transaction
binding
.two_pc(transaction, TwoPcPhase::Phase1, false)
.two_pc(&transaction, TwoPcPhase::Phase1, false)
.await
.expect("Phase1 should succeed");

// Then commit it
let result = binding.two_pc(transaction, TwoPcPhase::Phase2, false).await;
let result = binding
.two_pc(&transaction, TwoPcPhase::Phase2, false)
.await;
assert!(result.is_ok());
}

Expand All @@ -150,13 +154,13 @@ mod tests {

// First prepare the transaction
binding
.two_pc(transaction, TwoPcPhase::Phase1, false)
.two_pc(&transaction, TwoPcPhase::Phase1, false)
.await
.expect("Phase1 should succeed");

// Then rollback
let result = binding
.two_pc(transaction, TwoPcPhase::Rollback, false)
.two_pc(&transaction, TwoPcPhase::Rollback, false)
.await;
assert!(result.is_ok());
}
Expand All @@ -168,17 +172,17 @@ mod tests {

// First prepare the transaction
binding
.two_pc(transaction, TwoPcPhase::Phase1, false)
.two_pc(&transaction, TwoPcPhase::Phase1, false)
.await
.expect("Phase1 should succeed");

// Then commit it
binding
.two_pc(transaction, TwoPcPhase::Phase2, true)
.two_pc(&transaction, TwoPcPhase::Phase2, true)
.await
.expect("Phase2 should succeed");

let result = binding.two_pc(transaction, TwoPcPhase::Phase2, true).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Phase2, true).await;
assert!(
result.is_ok(),
"Committing non-existent prepared transaction should be skipped"
Expand All @@ -192,19 +196,19 @@ mod tests {

// First prepare the transaction
binding
.two_pc(transaction, TwoPcPhase::Phase1, true)
.two_pc(&transaction, TwoPcPhase::Phase1, true)
.await
.expect("Phase1 should succeed");

// Then rollback it
binding
.two_pc(transaction, TwoPcPhase::Rollback, true)
.two_pc(&transaction, TwoPcPhase::Rollback, true)
.await
.expect("Rollback should succeed");

// Try to rollback again - should succeed because skip_missing is true for Rollback
let result = binding
.two_pc(transaction, TwoPcPhase::Rollback, true)
.two_pc(&transaction, TwoPcPhase::Rollback, true)
.await;
assert!(
result.is_ok(),
Expand All @@ -221,17 +225,19 @@ mod tests {
let transaction = TwoPcTransaction::new();

// 1. Prepare transaction
let result = binding.two_pc(transaction, TwoPcPhase::Phase1, false).await;
let result = binding
.two_pc(&transaction, TwoPcPhase::Phase1, false)
.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, TwoPcPhase::Phase1, true).await;
let _result = binding.two_pc(&transaction, TwoPcPhase::Phase1, true).await;

// 3. Commit the prepared transaction
let result = binding.two_pc(transaction, TwoPcPhase::Phase2, true).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Phase2, true).await;
assert!(result.is_ok(), "Phase2 commit should succeed");

let result = binding.two_pc(transaction, TwoPcPhase::Phase2, true).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Phase2, true).await;
assert!(
result.is_ok(),
"Committing non-existent transaction should be skipped"
Expand All @@ -244,17 +250,19 @@ mod tests {
let transaction = TwoPcTransaction::new();

// 1. Prepare transaction
let result = binding.two_pc(transaction, TwoPcPhase::Phase1, false).await;
let result = binding
.two_pc(&transaction, TwoPcPhase::Phase1, false)
.await;
assert!(result.is_ok(), "Phase1 preparation should succeed");

// 2. Rollback the prepared transaction
let result = binding
.two_pc(transaction, TwoPcPhase::Rollback, true)
.two_pc(&transaction, TwoPcPhase::Rollback, true)
.await;
assert!(result.is_ok(), "Rollback should succeed");

// 3. Try to commit after rollback
let result = binding.two_pc(transaction, TwoPcPhase::Phase2, true).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Phase2, true).await;
assert!(
result.is_ok(),
"Committing rolled back transaction should be skipped"
Expand Down
2 changes: 1 addition & 1 deletion pgdog/src/backend/replication/logical/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -255,7 +255,7 @@ impl Error {
/// Two-phase commit transaction that still needs manager cleanup, if any.
pub fn two_pc_cleanup_transaction(&self) -> Option<TwoPcTransaction> {
match self {
Self::TwoPcCleanupPending { transaction, .. } => Some(*transaction),
Self::TwoPcCleanupPending { transaction, .. } => Some(transaction.clone()),
_ => None,
}
}
Expand Down
12 changes: 6 additions & 6 deletions pgdog/src/backend/replication/logical/subscriber/copy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -296,16 +296,16 @@ impl CopySubscriber {

async {
let _guard_phase_1 = manager
.transaction_state(txn, &identifier, TwoPcPhase::Phase1)
.transaction_state(txn.clone(), &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.clone(), &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?;
manager.done(txn.clone()).await?;
Ok(())
}
.await
Expand All @@ -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();
Expand Down
4 changes: 2 additions & 2 deletions pgdog/src/frontend/client/query_engine/end_transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -111,15 +111,15 @@ impl QueryEngine {
// If interrupted here, the transaction must be rolled back.
let _guard_phase_1 = self.two_pc.phase_one(&identifier).await?;
self.backend
.two_pc(transaction, TwoPcPhase::Phase1, false)
.two_pc(&transaction, TwoPcPhase::Phase1, false)
.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(transaction, TwoPcPhase::Phase2, false)
.two_pc(&transaction, TwoPcPhase::Phase2, false)
.await?;

debug!("[2pc] phase 2 complete");
Expand Down
16 changes: 8 additions & 8 deletions pgdog/src/frontend/client/query_engine/two_pc/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,7 @@ impl Manager {

/// Two-pc transaction finished.
pub(crate) async fn done(&self, transaction: TwoPcTransaction) -> Result<(), Error> {
if self.remove(transaction).is_some()
if self.remove(transaction.clone()).is_some()
&& let Some(wal) = self.wal.load_full()
{
wal.add(TwoPcRecordRemove { transaction }).await?;
Expand Down Expand Up @@ -174,17 +174,17 @@ impl Manager {
TwoPcPhase::Rollback,
"rollback is derived during recovery and is never written to the WAL"
);
self.set_transaction_state(transaction, identifier, phase);
self.set_transaction_state(transaction.clone(), identifier, phase);

if let Some(wal) = self.wal.load_full() {
if phase == TwoPcPhase::Phase1 {
wal.add(TwoPcRecordIdentity {
transaction,
transaction: transaction.clone(),
identifier: identifier.clone(),
})
.await?;
} else {
wal.add(TwoPcRecordPhase::new(transaction)).await?;
wal.add(TwoPcRecordPhase::new(transaction.clone())).await?;
}
}

Expand Down Expand Up @@ -267,7 +267,7 @@ impl Manager {
.contains_key(&guard.transaction);

if exists {
self.inner.lock().queue.push_back(guard.transaction);
self.inner.lock().queue.push_back(guard.transaction.clone());
self.notify.notify.notify_one();
}
}
Expand Down Expand Up @@ -302,7 +302,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: {}"#,
Expand Down Expand Up @@ -338,10 +338,10 @@ impl Manager {
}

/// Reconnect to cluster if available and close the two-phase transaction.
async fn cleanup_phase(&self, transaction: TwoPcTransaction) -> Result<(), Error> {
async fn cleanup_phase(&self, transaction: &TwoPcTransaction) -> Result<(), Error> {
let (state, in_recovery) = {
let guard = self.inner.lock();
let state = guard.transactions.get(&transaction).cloned();
let state = guard.transactions.get(transaction).cloned();

if let Some(state) = state {
(state, guard.in_recovery)
Expand Down
8 changes: 3 additions & 5 deletions pgdog/src/frontend/client/query_engine/two_pc/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,11 +46,9 @@ impl Default for TwoPc {
impl TwoPc {
/// Get a unique name for the two-pc transaction.
pub(super) fn transaction(&mut self) -> TwoPcTransaction {
if self.transaction.is_none() {
self.transaction = Some(TwoPcTransaction::new());
}

self.transaction.unwrap()
self.transaction
.get_or_insert_with(TwoPcTransaction::new)
.clone()
}

/// Start phase one of two-phase commit.
Expand Down
31 changes: 23 additions & 8 deletions pgdog/src/frontend/client/query_engine/two_pc/statement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ impl TwoPcTransactionOnShard {

/// Get the coordinator transaction.
pub(crate) fn transaction(&self) -> TwoPcTransaction {
self.transaction
self.transaction.clone()
}
}

Expand All @@ -48,11 +48,11 @@ impl FromStr for TwoPcTransactionOnShard {
/// Build `PREPARE TRANSACTION`, `COMMIT PREPARED`, or `ROLLBACK PREPARED`
/// for a shard participant.
pub(crate) fn phase_control(
transaction: TwoPcTransaction,
transaction: &TwoPcTransaction,
shard: usize,
phase: TwoPcPhase,
) -> String {
let txn = TwoPcTransactionOnShard::new(transaction, shard);
let txn = TwoPcTransactionOnShard::new(transaction.clone(), shard);

match phase {
TwoPcPhase::Phase1 => format!("PREPARE TRANSACTION '{txn}'"),
Expand All @@ -70,11 +70,11 @@ mod test {
let transaction = TwoPcTransaction::new();

assert_eq!(
TwoPcTransactionOnShard::new(transaction, 0).to_string(),
TwoPcTransactionOnShard::new(transaction.clone(), 0).to_string(),
format!("{transaction}_0")
);
assert_eq!(
TwoPcTransactionOnShard::new(transaction, 3).to_string(),
TwoPcTransactionOnShard::new(transaction.clone(), 3).to_string(),
format!("{transaction}_3")
);
}
Expand Down Expand Up @@ -106,16 +106,31 @@ mod test {
let transaction = TwoPcTransaction::new();

assert_eq!(
phase_control(transaction, 1, TwoPcPhase::Phase1),
phase_control(&transaction, 1, TwoPcPhase::Phase1),
format!("PREPARE TRANSACTION '{transaction}_1'")
);
assert_eq!(
phase_control(transaction, 1, TwoPcPhase::Phase2),
phase_control(&transaction, 1, TwoPcPhase::Phase2),
format!("COMMIT PREPARED '{transaction}_1'")
);
assert_eq!(
phase_control(transaction, 1, TwoPcPhase::Rollback),
phase_control(&transaction, 1, TwoPcPhase::Rollback),
format!("ROLLBACK PREPARED '{transaction}_1'")
);
}

#[test]
fn phase_control_uses_recovered_gid_verbatim() {
// A transaction recovered from the WAL carries a gid that no longer
// matches this process's rendering (different instance_id). The
// control statement must use it exactly as stored, with the per-shard
// suffix appended.
let recovered = "__pgdog_2pc_oldnode_42"
.parse::<TwoPcTransaction>()
.expect("valid recovered gid");
assert_eq!(
phase_control(&recovered, 3, TwoPcPhase::Rollback),
"ROLLBACK PREPARED '__pgdog_2pc_oldnode_42_3'"
);
}
}
Loading
Loading