Merge branch '6.0.x'
# Conflicts: # spring-r2dbc/src/main/java/org/springframework/r2dbc/connection/R2dbcTransactionManager.java # spring-r2dbc/src/test/java/org/springframework/r2dbc/connection/R2dbcTransactionManagerUnitTests.java
This commit is contained in:
@@ -338,8 +338,8 @@ public class R2dbcTransactionManager extends AbstractReactiveTransactionManager
|
||||
ConnectionFactoryTransactionObject txObject = (ConnectionFactoryTransactionObject) transaction;
|
||||
|
||||
if (txObject.hasSavepoint()) {
|
||||
// Just release the savepoint
|
||||
return Mono.defer(txObject::releaseSavepoint);
|
||||
// Just release the savepoint, keeping the transactional connection.
|
||||
return txObject.releaseSavepoint();
|
||||
}
|
||||
|
||||
// Remove the connection holder from the context, if exposed.
|
||||
@@ -348,30 +348,25 @@ public class R2dbcTransactionManager extends AbstractReactiveTransactionManager
|
||||
}
|
||||
|
||||
// Reset connection.
|
||||
Connection con = txObject.getConnectionHolder().getConnection();
|
||||
|
||||
Mono<Void> afterCleanup = Mono.empty();
|
||||
|
||||
Mono<Void> releaseConnectionStep = Mono.defer(() -> {
|
||||
try {
|
||||
if (txObject.isNewConnectionHolder()) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Releasing R2DBC Connection [" + con + "] after transaction");
|
||||
}
|
||||
Mono<Void> releaseMono = ConnectionFactoryUtils.releaseConnection(con, obtainConnectionFactory());
|
||||
if (logger.isDebugEnabled()) {
|
||||
releaseMono = releaseMono.doOnError(
|
||||
ex -> logger.debug(String.format("Error ignored during cleanup: %s", ex)));
|
||||
}
|
||||
return releaseMono.onErrorComplete();
|
||||
try {
|
||||
if (txObject.isNewConnectionHolder()) {
|
||||
Connection con = txObject.getConnectionHolder().getConnection();
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Releasing R2DBC Connection [" + con + "] after transaction");
|
||||
}
|
||||
Mono<Void> releaseMono = ConnectionFactoryUtils.releaseConnection(con, obtainConnectionFactory());
|
||||
if (logger.isDebugEnabled()) {
|
||||
releaseMono = releaseMono.doOnError(
|
||||
ex -> logger.debug(String.format("Error ignored during cleanup: %s", ex)));
|
||||
}
|
||||
return releaseMono.onErrorComplete();
|
||||
}
|
||||
finally {
|
||||
txObject.getConnectionHolder().clear();
|
||||
}
|
||||
return Mono.empty();
|
||||
});
|
||||
return afterCleanup.then(releaseConnectionStep);
|
||||
}
|
||||
finally {
|
||||
txObject.getConnectionHolder().clear();
|
||||
}
|
||||
|
||||
return Mono.empty();
|
||||
});
|
||||
}
|
||||
|
||||
@@ -517,30 +512,35 @@ public class R2dbcTransactionManager extends AbstractReactiveTransactionManager
|
||||
}
|
||||
|
||||
public boolean hasSavepoint() {
|
||||
return this.savepointName != null;
|
||||
return (this.savepointName != null);
|
||||
}
|
||||
|
||||
public Mono<Void> createSavepoint() {
|
||||
ConnectionHolder holder = getConnectionHolder();
|
||||
this.savepointName = holder.nextSavepoint();
|
||||
return Mono.from(holder.getConnection().createSavepoint(this.savepointName));
|
||||
String currentSavepoint = holder.nextSavepoint();
|
||||
this.savepointName = currentSavepoint;
|
||||
return Mono.from(holder.getConnection().createSavepoint(currentSavepoint));
|
||||
}
|
||||
|
||||
public Mono<Void> releaseSavepoint() {
|
||||
String currentSavepointName = this.savepointName;
|
||||
String currentSavepoint = this.savepointName;
|
||||
if (currentSavepoint == null) {
|
||||
return Mono.empty();
|
||||
}
|
||||
this.savepointName = null;
|
||||
return Mono.from(getConnectionHolder().getConnection().releaseSavepoint(currentSavepointName));
|
||||
return Mono.from(getConnectionHolder().getConnection().releaseSavepoint(currentSavepoint));
|
||||
}
|
||||
|
||||
public Mono<Void> commit() {
|
||||
Connection connection = getConnectionHolder().getConnection();
|
||||
return (this.savepointName != null ? Mono.empty() : Mono.from(connection.commitTransaction()));
|
||||
return (hasSavepoint() ? Mono.empty() :
|
||||
Mono.from(getConnectionHolder().getConnection().commitTransaction()));
|
||||
}
|
||||
|
||||
public Mono<Void> rollback() {
|
||||
Connection connection = getConnectionHolder().getConnection();
|
||||
return (this.savepointName != null ?
|
||||
Mono.from(connection.rollbackTransactionToSavepoint(this.savepointName)) :
|
||||
String currentSavepoint = this.savepointName;
|
||||
return (currentSavepoint != null ?
|
||||
Mono.from(connection.rollbackTransactionToSavepoint(currentSavepoint)) :
|
||||
Mono.from(connection.rollbackTransaction()));
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user