Fix call to block() in CouchbaseCallbackTransactionManager. (#1528)
Closes #1527.
This commit is contained in:
@@ -15,6 +15,7 @@
|
||||
*/
|
||||
package org.springframework.data.couchbase.core;
|
||||
|
||||
import com.couchbase.client.core.transaction.threadlocal.TransactionMarker;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.util.Optional;
|
||||
@@ -26,6 +27,7 @@ import com.couchbase.client.core.error.CasMismatchException;
|
||||
import com.couchbase.client.core.error.transaction.TransactionOperationFailedException;
|
||||
import com.couchbase.client.core.transaction.CoreTransactionAttemptContext;
|
||||
import com.couchbase.client.core.transaction.threadlocal.TransactionMarkerOwner;
|
||||
import reactor.util.context.ContextView;
|
||||
|
||||
/**
|
||||
* Utility methods to support transactions.
|
||||
@@ -50,6 +52,14 @@ public class TransactionalSupport {
|
||||
});
|
||||
}
|
||||
|
||||
public static Optional<CouchbaseResourceHolder> checkForTransactionInThreadLocalStorage(ContextView ctx) {
|
||||
return Optional.ofNullable(ctx.hasKey(TransactionMarker.class) ? new CouchbaseResourceHolder(ctx.get(TransactionMarker.class).context()) : null);
|
||||
}
|
||||
|
||||
//public static Optional<CouchbaseResourceHolder> blockingCheckForTransactionInThreadLocalStorage() {
|
||||
// return TransactionMarkerOwner.marker;
|
||||
// }
|
||||
|
||||
public static Mono<Void> verifyNotInTransaction(String methodName) {
|
||||
return checkForTransactionInThreadLocalStorage().flatMap(s -> {
|
||||
if (s.isPresent()) {
|
||||
|
||||
@@ -44,6 +44,7 @@ import com.couchbase.client.java.transactions.TransactionResult;
|
||||
import com.couchbase.client.java.transactions.config.TransactionOptions;
|
||||
import com.couchbase.client.java.transactions.error.TransactionCommitAmbiguousException;
|
||||
import com.couchbase.client.java.transactions.error.TransactionFailedException;
|
||||
import reactor.util.context.ContextView;
|
||||
|
||||
/**
|
||||
* The Couchbase transaction manager, providing support for @Transactional methods.
|
||||
@@ -73,7 +74,7 @@ public class CouchbaseCallbackTransactionManager implements CallbackPreferringPl
|
||||
|
||||
@Override
|
||||
public <T> T execute(TransactionDefinition definition, TransactionCallback<T> callback) throws TransactionException {
|
||||
boolean createNewTransaction = handlePropagation(definition);
|
||||
boolean createNewTransaction = handlePropagation(definition, null);
|
||||
|
||||
setOptionsFromDefinition(definition);
|
||||
|
||||
@@ -87,8 +88,8 @@ public class CouchbaseCallbackTransactionManager implements CallbackPreferringPl
|
||||
@Stability.Internal
|
||||
<T> Flux<T> executeReactive(TransactionDefinition definition,
|
||||
org.springframework.transaction.reactive.TransactionCallback<T> callback) {
|
||||
return Flux.defer(() -> {
|
||||
boolean createNewTransaction = handlePropagation(definition);
|
||||
return Flux.deferContextual((ctx) -> {
|
||||
boolean createNewTransaction = handlePropagation(definition, ctx);
|
||||
|
||||
setOptionsFromDefinition(definition);
|
||||
|
||||
@@ -187,8 +188,9 @@ public class CouchbaseCallbackTransactionManager implements CallbackPreferringPl
|
||||
}
|
||||
|
||||
// Propagation defines what happens when a @Transactional method is called from another @Transactional method.
|
||||
private boolean handlePropagation(TransactionDefinition definition) {
|
||||
boolean isExistingTransaction = TransactionalSupport.checkForTransactionInThreadLocalStorage().block().isPresent();
|
||||
private boolean handlePropagation(TransactionDefinition definition, ContextView ctx) {
|
||||
boolean isExistingTransaction = ctx != null ? TransactionalSupport.checkForTransactionInThreadLocalStorage(ctx).isPresent() :
|
||||
TransactionalSupport.checkForTransactionInThreadLocalStorage().block().isPresent();
|
||||
|
||||
LOGGER.trace("Deciding propagation behaviour from {} and {}", definition.getPropagationBehavior(),
|
||||
isExistingTransaction);
|
||||
|
||||
Reference in New Issue
Block a user