From 018dae47cfe8823f19b0eb9fba990cae3b605a95 Mon Sep 17 00:00:00 2001 From: ralfkurka Date: Sat, 14 Mar 2026 11:37:35 +0100 Subject: [PATCH] Fix connection returned to pool without rollback after failed commit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PooledConnection.commitTransaction() used doOnSubscribe to set inTransaction=false, which fires immediately when the commit Mono is subscribed — before the database has actually committed. If the commit then fails (e.g. serialization failure, deadlock, transient network error), inTransaction is already false. When close() runs afterwards, it skips the safety-net rollback (line 74: if this.inTransaction) and returns the connection to the pool with an open transaction. The next consumer that acquires this connection inherits the uncommitted transaction state, which can lead to data corruption: their writes become part of the previous transaction, and a subsequent commit or rollback affects both sessions' data. Change doOnSubscribe to doOnSuccess so that inTransaction is only cleared after the commit actually succeeds. On failure, the flag remains true, and close() will issue a rollback before releasing the connection back to the pool. Add a unit test that reproduces the scenario: begin → failed commit → close, and verifies that rollbackTransaction() is called. Co-Authored-By: Claude Opus 4.6 --- .../java/io/r2dbc/pool/PooledConnection.java | 2 +- .../r2dbc/pool/PooledConnectionUnitTests.java | 17 +++++++++++++++++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/src/main/java/io/r2dbc/pool/PooledConnection.java b/src/main/java/io/r2dbc/pool/PooledConnection.java index aff2c09..f691efa 100644 --- a/src/main/java/io/r2dbc/pool/PooledConnection.java +++ b/src/main/java/io/r2dbc/pool/PooledConnection.java @@ -114,7 +114,7 @@ public Mono close() { @Override public Mono commitTransaction() { assertNotClosed(); - return Mono.from(this.connection.commitTransaction()).doOnSubscribe(ignore -> this.inTransaction = false); + return Mono.from(this.connection.commitTransaction()).doOnSuccess(ignore -> this.inTransaction = false); } @Override diff --git a/src/test/java/io/r2dbc/pool/PooledConnectionUnitTests.java b/src/test/java/io/r2dbc/pool/PooledConnectionUnitTests.java index 8e99ed6..da1124b 100644 --- a/src/test/java/io/r2dbc/pool/PooledConnectionUnitTests.java +++ b/src/test/java/io/r2dbc/pool/PooledConnectionUnitTests.java @@ -123,6 +123,23 @@ void committedTransactionLeavesTransactionalStateAsIs() { verify(connectionMock, never()).rollbackTransaction(); } + @Test + void failedCommitShouldRollbackOnClose() { + + AtomicBoolean rollbackCalled = new AtomicBoolean(); + when(connectionMock.commitTransaction()).thenReturn(Mono.error(new RuntimeException("commit failed"))); + when(connectionMock.rollbackTransaction()).thenReturn(Mono.empty().doOnSuccess(o -> rollbackCalled.set(true))); + + PooledConnection connection = new PooledConnection(pooledRefMock); + connection.beginTransaction().as(StepVerifier::create).verifyComplete(); + connection.commitTransaction().as(StepVerifier::create).verifyError(); + + connection.close().as(StepVerifier::create).verifyComplete(); + + verify(connectionMock).rollbackTransaction(); + assertThat(rollbackCalled).isTrue(); + } + @Test void rolledBackTransactionLeavesTransactionalStateAsIs() {