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
7 changes: 4 additions & 3 deletions integration/two_pc_crash_safety/wal_helper/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,8 @@ async fn main() {
std::process::exit(2);
}
let dir = PathBuf::from(&args[1]);
let txn = TwoPcTransaction::from_str(&args[2])
.unwrap_or_else(|_| panic!("invalid gid {:?}", args[2]));
let gid = args[2].clone();
let txn = TwoPcTransaction::from_str(&gid).unwrap_or_else(|_| panic!("invalid gid {:?}", gid));
let user = args[3].clone();
let database = args[4].clone();

Expand All @@ -34,9 +34,10 @@ async fn main() {

let mut buf = BytesMut::new();
Record::Begin(BeginPayload {
txn,
txn: txn.clone(),
user,
database,
gid,
})
.encode(&mut buf)
.expect("encode begin");
Expand Down
9 changes: 7 additions & 2 deletions pgdog/src/backend/pool/connection/binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -364,7 +364,7 @@ impl Binding {

pub(crate) async fn two_pc_on_guards(
servers: &mut [Guard],
transaction: TwoPcTransaction,
transaction: &TwoPcTransaction,
phase: TwoPcPhase,
) -> Result<(), Error> {
let skip_missing = matches!(phase, TwoPcPhase::Phase2 | TwoPcPhase::Rollback);
Expand Down Expand Up @@ -397,9 +397,14 @@ impl Binding {
}

/// Execute two-phase commit transaction control statements.
///
/// The per-shard participant names are derived from `transaction`. For a
/// transaction rebuilt by WAL recovery this reproduces the exact gid it
/// was prepared with, even after a restart with a different
/// `instance_id` (see [`TwoPcTransaction`]).
pub(crate) async fn two_pc(
&mut self,
transaction: TwoPcTransaction,
transaction: &TwoPcTransaction,
phase: TwoPcPhase,
) -> Result<(), Error> {
match self {
Expand Down
42 changes: 21 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)
.two_pc(&TwoPcTransaction::new(), TwoPcPhase::Phase1)
.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)
.two_pc(&TwoPcTransaction::new(), TwoPcPhase::Phase1)
.await;

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

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

// Should succeed
if let Err(ref error) = result {
Expand All @@ -122,7 +122,7 @@ mod tests {
assert!(result.is_ok());

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

#[tokio::test]
Expand All @@ -132,12 +132,12 @@ mod tests {

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

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

Expand All @@ -148,12 +148,12 @@ mod tests {

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

// Then rollback
let result = binding.two_pc(transaction, TwoPcPhase::Rollback).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Rollback).await;
assert!(result.is_ok());
}

Expand All @@ -164,18 +164,18 @@ mod tests {

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

// Then commit it
binding
.two_pc(transaction, 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, TwoPcPhase::Phase2).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Phase2).await;
assert!(
result.is_ok(),
"Committing non-existent prepared transaction should be skipped"
Expand All @@ -189,18 +189,18 @@ mod tests {

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

// Then rollback it
binding
.two_pc(transaction, 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, 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"
Expand All @@ -216,19 +216,19 @@ mod tests {
let transaction = TwoPcTransaction::new();

// 1. Prepare transaction
let result = binding.two_pc(transaction, 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, 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, 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, TwoPcPhase::Phase2).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Phase2).await;
assert!(
result.is_ok(),
"Committing non-existent transaction should be skipped"
Expand All @@ -241,15 +241,15 @@ mod tests {
let transaction = TwoPcTransaction::new();

// 1. Prepare transaction
let result = binding.two_pc(transaction, 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, 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, TwoPcPhase::Phase2).await;
let result = binding.two_pc(&transaction, TwoPcPhase::Phase2).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
10 changes: 5 additions & 5 deletions pgdog/src/backend/replication/logical/subscriber/copy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -296,14 +296,14 @@ 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?;
Ok(())
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
8 changes: 6 additions & 2 deletions pgdog/src/frontend/client/query_engine/end_transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,13 +110,17 @@ 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).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(transaction, TwoPcPhase::Phase2).await?;
self.backend
.two_pc(&transaction, TwoPcPhase::Phase2)
.await?;

debug!("[2pc] phase 2 complete");

Expand Down
24 changes: 14 additions & 10 deletions pgdog/src/frontend/client/query_engine/two_pc/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -208,7 +208,7 @@ impl Manager {
let prior = {
let mut guard = self.inner.lock();
let prior = guard.transactions.get(&transaction).cloned();
let entry = guard.transactions.entry(transaction).or_default();
let entry = guard.transactions.entry(transaction.clone()).or_default();
entry.identifier = identifier.clone();
entry.phase = phase;
prior
Expand All @@ -218,13 +218,14 @@ impl Manager {
let result = match phase {
TwoPcPhase::Phase1 => {
wal.append_begin(
transaction,
transaction.clone(),
identifier.user.clone(),
identifier.database.clone(),
transaction.to_string(),
)
.await
}
TwoPcPhase::Phase2 => wal.append_committing(transaction).await,
TwoPcPhase::Phase2 => wal.append_committing(transaction.clone()).await,
TwoPcPhase::Rollback => {
unreachable!("rollback is not a state transition; it's the cleanup direction")
}
Expand All @@ -233,7 +234,7 @@ impl Manager {
let mut guard = self.inner.lock();
match prior {
Some(prior) => {
guard.transactions.insert(transaction, prior);
guard.transactions.insert(transaction.clone(), prior);
}
None => {
guard.transactions.remove(&transaction);
Expand Down Expand Up @@ -269,7 +270,7 @@ impl Manager {
let mut guard = self.inner.lock();
guard
.transactions
.insert(transaction, TransactionInfo { phase, identifier });
.insert(transaction.clone(), TransactionInfo { phase, identifier });
guard.queue.push_back(transaction);
}
self.stats.incr_recovered();
Expand All @@ -284,7 +285,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 All @@ -309,7 +310,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,15 +339,15 @@ impl Manager {
async fn remove(&self, transaction: &TwoPcTransaction) {
self.inner.lock().transactions.remove(transaction);
if let Some(wal) = self.wal.load_full()
&& let Err(err) = wal.append_end(*transaction).await
&& let Err(err) = wal.append_end(transaction.clone()).await
{
warn!("[2pc] wal end record failed for {}: {}", transaction, err);
}
}

/// 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(());
Expand Down Expand Up @@ -379,6 +380,9 @@ impl Manager {
&Route::write(ShardWithPriority::new_override_transaction(Shard::All)),
)
.await?;
// `transaction` carries its own gid: for a recovered transaction
// that's the exact value it was prepared with, so this resolves the
// orphan on Postgres even after a restart with a new instance_id.
connection.two_pc(transaction, phase).await?;
connection.disconnect();

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
Loading
Loading