From d721579655a54b795a4f4a57c3c954a49b3e1524 Mon Sep 17 00:00:00 2001 From: Christoph Strobl Date: Thu, 15 Feb 2024 14:09:04 +0100 Subject: [PATCH] Introduce `MongoTransactionResolver`. See #1628 Original pull request: #4552 --- ...efaultMongoTransactionOptionsResolver.java | 62 +++++ .../data/mongodb/MongoTransactionManager.java | 24 +- .../data/mongodb/MongoTransactionOptions.java | 205 ++++++++++++++++ .../MongoTransactionOptionsResolver.java | 114 +++++++++ .../data/mongodb/MongoTransactionUtils.java | 98 -------- .../ReactiveMongoTransactionManager.java | 26 +- .../SimpleMongoTransactionOptions.java | 154 ++++++++++++ .../data/mongodb/TransactionMetadata.java | 34 +++ .../mongodb/TransactionOptionResolver.java | 29 +++ .../data/mongodb/core/WriteConcernAware.java | 41 ++++ .../CapturingTransactionOptionsResolver.java | 64 +++++ .../MongoTransactionOptionsUnitTests.java | 116 +++++++++ .../MongoTransactionUtilsUnitTests.java | 227 ------------------ .../ReactiveTransactionIntegrationTests.java | 59 ++++- ...ReactiveTransactionOptionsTestService.java | 5 + ...goTransactionOptionsResolverUnitTests.java | 131 ++++++++++ .../core/MongoTemplateTransactionTests.java | 55 ++++- .../core/TransactionOptionsTestService.java | 6 + 18 files changed, 1092 insertions(+), 358 deletions(-) create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/DefaultMongoTransactionOptionsResolver.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptions.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptionsResolver.java delete mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionUtils.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/SimpleMongoTransactionOptions.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionMetadata.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionOptionResolver.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/WriteConcernAware.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/CapturingTransactionOptionsResolver.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionOptionsUnitTests.java delete mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionUtilsUnitTests.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/SimpleMongoTransactionOptionsResolverUnitTests.java diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/DefaultMongoTransactionOptionsResolver.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/DefaultMongoTransactionOptionsResolver.java new file mode 100644 index 000000000..02447ff0e --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/DefaultMongoTransactionOptionsResolver.java @@ -0,0 +1,62 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import java.util.Map; +import java.util.Set; + +import org.springframework.data.util.Lazy; +import org.springframework.lang.Nullable; + +/** + * Default implementation of {@link MongoTransactionOptions} using {@literal mongo:} as {@link #getLabelPrefix() label + * prefix} creating {@link SimpleMongoTransactionOptions} out of a given argument {@link Map}. Uses + * {@link SimpleMongoTransactionOptions#KNOWN_KEYS} to validate entries in arguments to resolve and errors on unknown + * entries. + * + * @author Christoph Strobl + * @since 4.3 + */ +class DefaultMongoTransactionOptionsResolver implements MongoTransactionOptionsResolver { + + static final Lazy INSTANCE = Lazy.of(DefaultMongoTransactionOptionsResolver::new); + + private static final String PREFIX = "mongo:"; + + private DefaultMongoTransactionOptionsResolver() {} + + @Override + public MongoTransactionOptions convert(Map options) { + + validateKeys(options.keySet()); + return SimpleMongoTransactionOptions.of(options); + } + + @Nullable + @Override + public String getLabelPrefix() { + return PREFIX; + } + + private static void validateKeys(Set keys) { + + if (!keys.stream().allMatch(SimpleMongoTransactionOptions.KNOWN_KEYS::contains)) { + + throw new IllegalArgumentException("Transaction labels contained invalid values. Has to be one of %s" + .formatted(SimpleMongoTransactionOptions.KNOWN_KEYS)); + } + } +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionManager.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionManager.java index 4b1ad5617..895297b3f 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionManager.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionManager.java @@ -65,7 +65,8 @@ public class MongoTransactionManager extends AbstractPlatformTransactionManager implements ResourceTransactionManager, InitializingBean { private @Nullable MongoDatabaseFactory dbFactory; - private @Nullable TransactionOptions options; + private MongoTransactionOptions options; + private MongoTransactionOptionsResolver transactionOptionsResolver; /** * Create a new {@link MongoTransactionManager} for bean-style usage. @@ -99,11 +100,25 @@ public class MongoTransactionManager extends AbstractPlatformTransactionManager * @param options can be {@literal null}. */ public MongoTransactionManager(MongoDatabaseFactory dbFactory, @Nullable TransactionOptions options) { + this(dbFactory, MongoTransactionOptionsResolver.defaultResolver(), MongoTransactionOptions.of(options)); + } + + /** + * Create a new {@link MongoTransactionManager} obtaining sessions from the given {@link MongoDatabaseFactory} + * applying the given {@link TransactionOptions options}, if present, when starting a new transaction. + * + * @param dbFactory must not be {@literal null}. + * @param transactionOptionsResolver + * @param defaultTransactionOptions can be {@literal null}. + * @since 4.3 + */ + public MongoTransactionManager(MongoDatabaseFactory dbFactory, MongoTransactionOptionsResolver transactionOptionsResolver, MongoTransactionOptions defaultTransactionOptions) { Assert.notNull(dbFactory, "DbFactory must not be null"); this.dbFactory = dbFactory; - this.options = options; + this.transactionOptionsResolver = transactionOptionsResolver; + this.options = defaultTransactionOptions; } @Override @@ -134,7 +149,8 @@ public class MongoTransactionManager extends AbstractPlatformTransactionManager } try { - mongoTransactionObject.startTransaction(MongoTransactionUtils.extractOptions(definition, options)); + MongoTransactionOptions mongoTransactionOptions = transactionOptionsResolver.resolve(definition).mergeWith(options); + mongoTransactionObject.startTransaction(mongoTransactionOptions.toDriverOptions()); } catch (MongoException ex) { throw new TransactionSystemException(String.format("Could not start Mongo transaction for session %s.", debugString(mongoTransactionObject.getSession())), ex); @@ -276,7 +292,7 @@ public class MongoTransactionManager extends AbstractPlatformTransactionManager * @param options can be {@literal null}. */ public void setOptions(@Nullable TransactionOptions options) { - this.options = options; + this.options = MongoTransactionOptions.of(options); } /** diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptions.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptions.java new file mode 100644 index 000000000..4c9957b6e --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptions.java @@ -0,0 +1,205 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import java.time.Duration; +import java.util.concurrent.TimeUnit; + +import org.springframework.data.mongodb.core.ReadConcernAware; +import org.springframework.data.mongodb.core.ReadPreferenceAware; +import org.springframework.data.mongodb.core.WriteConcernAware; +import org.springframework.lang.Nullable; + +import com.mongodb.Function; +import com.mongodb.ReadConcern; +import com.mongodb.ReadPreference; +import com.mongodb.TransactionOptions; +import com.mongodb.WriteConcern; + +/** + * Options to be applied within a specific transaction scope. + * + * @author Christoph Strobl + * @since 4.3 + */ +public interface MongoTransactionOptions + extends TransactionMetadata, ReadConcernAware, ReadPreferenceAware, WriteConcernAware { + + /** + * Value Object representing empty options enforcing client defaults. Returns {@literal null} for all getter methods. + */ + MongoTransactionOptions NONE = new MongoTransactionOptions() { + + @Nullable + @Override + public Duration getMaxCommitTime() { + return null; + } + + @Nullable + @Override + public ReadConcern getReadConcern() { + return null; + } + + @Nullable + @Override + public ReadPreference getReadPreference() { + return null; + } + + @Nullable + @Override + public WriteConcern getWriteConcern() { + return null; + } + }; + + /** + * Merge current options with given ones. Will return first non {@literal null} value from getters whereas the + * {@literal this} has precedence over the given fallbackOptions. + * + * @param fallbackOptions can be {@literal null}. + * @return new instance of {@link MongoTransactionOptions} or this if {@literal fallbackOptions} is {@literal null} or + * {@link #NONE}. + */ + default MongoTransactionOptions mergeWith(@Nullable MongoTransactionOptions fallbackOptions) { + + if (fallbackOptions == null || MongoTransactionOptions.NONE.equals(fallbackOptions)) { + return this; + } + + return new MongoTransactionOptions() { + + @Nullable + @Override + public Duration getMaxCommitTime() { + return MongoTransactionOptions.this.hasMaxCommitTime() ? MongoTransactionOptions.this.getMaxCommitTime() + : fallbackOptions.getMaxCommitTime(); + } + + @Nullable + @Override + public ReadConcern getReadConcern() { + return MongoTransactionOptions.this.hasReadConcern() ? MongoTransactionOptions.this.getReadConcern() + : fallbackOptions.getReadConcern(); + } + + @Nullable + @Override + public ReadPreference getReadPreference() { + return MongoTransactionOptions.this.hasReadPreference() ? MongoTransactionOptions.this.getReadPreference() + : fallbackOptions.getReadPreference(); + } + + @Nullable + @Override + public WriteConcern getWriteConcern() { + return MongoTransactionOptions.this.hasWriteConcern() ? MongoTransactionOptions.this.getWriteConcern() + : fallbackOptions.getWriteConcern(); + } + }; + } + + /** + * Map the current options using the given mapping {@link Function}. + * + * @param mappingFunction + * @return instance of T. + * @param + */ + default T as(Function mappingFunction) { + return mappingFunction.apply(this); + } + + /** + * @return MongoDB driver native {@link TransactionOptions}. + * @see MongoTransactionOptions#as(Function) + */ + @Nullable + default TransactionOptions toDriverOptions() { + + return as(it -> { + + if (MongoTransactionOptions.NONE.equals(it)) { + return null; + } + + TransactionOptions.Builder builder = TransactionOptions.builder(); + if (it.hasMaxCommitTime()) { + builder.maxCommitTime(it.getMaxCommitTime().toMillis(), TimeUnit.MILLISECONDS); + } + if (it.hasReadConcern()) { + builder.readConcern(it.getReadConcern()); + } + if (it.hasReadPreference()) { + builder.readPreference(it.getReadPreference()); + } + if (it.hasWriteConcern()) { + builder.writeConcern(it.getWriteConcern()); + } + return builder.build(); + }); + } + + /** + * Factory method to wrap given MongoDB driver native {@link TransactionOptions} into {@link MongoTransactionOptions}. + * + * @param options + * @return {@link MongoTransactionOptions#NONE} if given object is {@literal null}. + */ + static MongoTransactionOptions of(@Nullable TransactionOptions options) { + + if (options == null) { + return NONE; + } + + return new MongoTransactionOptions() { + + @Nullable + @Override + public Duration getMaxCommitTime() { + + Long millis = options.getMaxCommitTime(TimeUnit.MILLISECONDS); + return millis != null ? Duration.ofMillis(millis) : null; + } + + @Nullable + @Override + public ReadConcern getReadConcern() { + return options.getReadConcern(); + } + + @Nullable + @Override + public ReadPreference getReadPreference() { + return options.getReadPreference(); + } + + @Nullable + @Override + public WriteConcern getWriteConcern() { + return options.getWriteConcern(); + } + + @Nullable + @Override + public TransactionOptions toDriverOptions() { + return options; + } + }; + } +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptionsResolver.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptionsResolver.java new file mode 100644 index 000000000..78ffc9774 --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionOptionsResolver.java @@ -0,0 +1,114 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import java.util.Map; +import java.util.stream.Collectors; + +import org.springframework.lang.Nullable; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.interceptor.TransactionAttribute; +import org.springframework.util.Assert; +import org.springframework.util.StringUtils; + +/** + * A {@link TransactionOptionResolver} reading MongoDB specific {@link MongoTransactionOptions transaction options} from + * a {@link TransactionDefinition}. Implementations of {@link MongoTransactionOptions} may choose a specific + * {@link #getLabelPrefix() prefix} for {@link TransactionAttribute#getLabels() transaction attribute labels} to avoid + * evaluating non store specific ones. + *

+ * {@link TransactionAttribute#getLabels()} evaluated by default should follow the property style using {@code =} to + * separate key and value pairs. + *

+ * By default {@link #resolve(TransactionDefinition)} will filter labels by the {@link #getLabelPrefix() prefix} and + * strip the prefix from the label before handing the pruned {@link Map} to the {@link #convert(Map)} function. + *

+ * A transaction definition with labels targeting MongoDB may look like the following: + *

+ * + * @Transactional(label = { "mongo:readConcern=majority" }) + * + * + * @author Christoph Strobl + * @since 4.3 + */ +public interface MongoTransactionOptionsResolver extends TransactionOptionResolver { + + /** + * Obtain the default {@link MongoTransactionOptionsResolver} implementation using a {@literal mongo:} + * {@link #getLabelPrefix() prefix}. + * + * @return instance of default {@link MongoTransactionOptionsResolver} implementation. + */ + static MongoTransactionOptionsResolver defaultResolver() { + return DefaultMongoTransactionOptionsResolver.INSTANCE.get(); + } + + /** + * Get the prefix used to filter applicable {@link TransactionAttribute#getLabels() labels}. + * + * @return {@literal null} if no label defined. + */ + @Nullable + String getLabelPrefix(); + + /** + * Resolve {@link MongoTransactionOptions} from a given {@link TransactionDefinition} by evaluating + * {@link TransactionAttribute#getLabels()} labels if possible. + *

+ * Splits applicable labels property style using {@literal =} as deliminator and removes a potential + * {@link #getLabelPrefix() prefix} before calling {@link #convert(Map)} with filtered label values. + * + * @param txDefinition + * @return {@link MongoTransactionOptions#NONE} in case the given {@link TransactionDefinition} is not a + * {@link TransactionAttribute} if no matching {@link TransactionAttribute#getLabels() labels} could be found. + * @throws IllegalArgumentException for options that do not map to valid transactions options or malformatted labels. + */ + @Override + default MongoTransactionOptions resolve(TransactionDefinition txDefinition) { + + if (!(txDefinition instanceof TransactionAttribute attribute)) { + return MongoTransactionOptions.NONE; + } + + if (attribute.getLabels().isEmpty()) { + return MongoTransactionOptions.NONE; + } + + Map attributeMap = attribute.getLabels().stream() + .filter(it -> !StringUtils.hasText(getLabelPrefix()) || it.startsWith(getLabelPrefix())) + .map(it -> StringUtils.hasText(getLabelPrefix()) ? it.substring(getLabelPrefix().length()) : it).map(it -> { + + String[] kvPair = StringUtils.split(it, "="); + Assert.isTrue(kvPair != null && kvPair.length == 2, + () -> "No value present for transaction option %s".formatted(kvPair != null ? kvPair[0] : it)); + return kvPair; + }) + + .collect(Collectors.toMap(it -> it[0].trim(), it -> it[1].trim())); + + return attributeMap.isEmpty() ? MongoTransactionOptions.NONE : convert(attributeMap); + } + + /** + * Convert the given {@link Map} into an instance of {@link MongoTransactionOptions}. + * + * @param options never {@literal null}. + * @return never {@literal null}. + * @throws IllegalArgumentException for invalid options. + */ + MongoTransactionOptions convert(Map options); +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionUtils.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionUtils.java deleted file mode 100644 index 13c4c259b..000000000 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/MongoTransactionUtils.java +++ /dev/null @@ -1,98 +0,0 @@ -/* - * Copyright 2023 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.data.mongodb; - -import java.time.Duration; -import java.util.concurrent.TimeUnit; - -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; -import org.springframework.lang.Nullable; -import org.springframework.transaction.TransactionDefinition; -import org.springframework.transaction.interceptor.TransactionAttribute; - -import com.mongodb.ReadConcern; -import com.mongodb.ReadConcernLevel; -import com.mongodb.ReadPreference; -import com.mongodb.TransactionOptions; -import com.mongodb.WriteConcern; - -/** - * Helper class for translating @Transactional labels into Mongo-specific {@link TransactionOptions}. - * - * @author Yan Kardziyaka - */ -public final class MongoTransactionUtils { - private static final Log LOGGER = LogFactory.getLog(MongoTransactionUtils.class); - - private static final String MAX_COMMIT_TIME = "mongo:maxCommitTime"; - - private static final String READ_CONCERN_OPTION = "mongo:readConcern"; - - private static final String READ_PREFERENCE_OPTION = "mongo:readPreference"; - - private static final String WRITE_CONCERN_OPTION = "mongo:writeConcern"; - - private MongoTransactionUtils() {} - - @Nullable - public static TransactionOptions extractOptions(TransactionDefinition transactionDefinition, - @Nullable TransactionOptions fallbackOptions) { - if (transactionDefinition instanceof TransactionAttribute transactionAttribute) { - TransactionOptions.Builder builder = null; - for (String label : transactionAttribute.getLabels()) { - String[] tokens = label.split("=", 2); - builder = tokens.length == 2 ? enhanceWithProperty(builder, tokens[0], tokens[1]) : builder; - } - if (builder == null) { - return fallbackOptions; - } - TransactionOptions options = builder.build(); - return fallbackOptions == null ? options : TransactionOptions.merge(options, fallbackOptions); - } else { - if (LOGGER.isDebugEnabled()) { - LOGGER.debug("%s cannot be casted to %s. Transaction labels won't be evaluated as options".formatted( - TransactionDefinition.class.getName(), TransactionAttribute.class.getName())); - } - return fallbackOptions; - } - } - - @Nullable - private static TransactionOptions.Builder enhanceWithProperty(@Nullable TransactionOptions.Builder builder, - String key, String value) { - return switch (key) { - case MAX_COMMIT_TIME -> nullSafe(builder).maxCommitTime(Duration.parse(value).toMillis(), TimeUnit.MILLISECONDS); - case READ_CONCERN_OPTION -> nullSafe(builder).readConcern(new ReadConcern(ReadConcernLevel.fromString(value))); - case READ_PREFERENCE_OPTION -> nullSafe(builder).readPreference(ReadPreference.valueOf(value)); - case WRITE_CONCERN_OPTION -> nullSafe(builder).writeConcern(getWriteConcern(value)); - default -> builder; - }; - } - - private static TransactionOptions.Builder nullSafe(@Nullable TransactionOptions.Builder builder) { - return builder == null ? TransactionOptions.builder() : builder; - } - - private static WriteConcern getWriteConcern(String writeConcernAsString) { - WriteConcern writeConcern = WriteConcern.valueOf(writeConcernAsString); - if (writeConcern == null) { - throw new IllegalArgumentException("'%s' is not a valid WriteConcern".formatted(writeConcernAsString)); - } - return writeConcern; - } - -} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java index c8c38a622..3907acbb7 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/ReactiveMongoTransactionManager.java @@ -67,7 +67,8 @@ import com.mongodb.reactivestreams.client.ClientSession; public class ReactiveMongoTransactionManager extends AbstractReactiveTransactionManager implements InitializingBean { private @Nullable ReactiveMongoDatabaseFactory databaseFactory; - private @Nullable TransactionOptions options; + private @Nullable MongoTransactionOptions options; + private MongoTransactionOptionsResolver transactionOptionsResolver; /** * Create a new {@link ReactiveMongoTransactionManager} for bean-style usage. @@ -103,11 +104,27 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction */ public ReactiveMongoTransactionManager(ReactiveMongoDatabaseFactory databaseFactory, @Nullable TransactionOptions options) { + this(databaseFactory, MongoTransactionOptionsResolver.defaultResolver(), MongoTransactionOptions.of(options)); + } + + /** + * Create a new {@link ReactiveMongoTransactionManager} obtaining sessions from the given + * {@link ReactiveMongoDatabaseFactory} applying the given {@link TransactionOptions options}, if present, when + * starting a new transaction. + * + * @param databaseFactory must not be {@literal null}. + * @param transactionOptionsResolver + * @param defaultTransactionOptions can be {@literal null}. + * + */ + public ReactiveMongoTransactionManager(ReactiveMongoDatabaseFactory databaseFactory, MongoTransactionOptionsResolver transactionOptionsResolver, + @Nullable MongoTransactionOptions defaultTransactionOptions) { Assert.notNull(databaseFactory, "DatabaseFactory must not be null"); this.databaseFactory = databaseFactory; - this.options = options; + this.transactionOptionsResolver = transactionOptionsResolver; + this.options = defaultTransactionOptions; } @Override @@ -146,7 +163,8 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction }).doOnNext(resourceHolder -> { - mongoTransactionObject.startTransaction(MongoTransactionUtils.extractOptions(definition, options)); + MongoTransactionOptions mongoTransactionOptions = transactionOptionsResolver.resolve(definition).mergeWith(options); + mongoTransactionObject.startTransaction(mongoTransactionOptions.toDriverOptions()); if (logger.isDebugEnabled()) { logger.debug(String.format("Started transaction for session %s.", debugString(resourceHolder.getSession()))); @@ -291,7 +309,7 @@ public class ReactiveMongoTransactionManager extends AbstractReactiveTransaction * @param options can be {@literal null}. */ public void setOptions(@Nullable TransactionOptions options) { - this.options = options; + this.options = MongoTransactionOptions.of(options); } /** diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/SimpleMongoTransactionOptions.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/SimpleMongoTransactionOptions.java new file mode 100644 index 000000000..9cd2146ce --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/SimpleMongoTransactionOptions.java @@ -0,0 +1,154 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import java.time.Duration; +import java.util.Arrays; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; + +import org.springframework.lang.Nullable; +import org.springframework.util.Assert; + +import com.mongodb.Function; +import com.mongodb.ReadConcern; +import com.mongodb.ReadConcernLevel; +import com.mongodb.ReadPreference; +import com.mongodb.WriteConcern; + +/** + * Trivial implementation of {@link MongoTransactionOptions}. + * + * @author Christoph Strobl + * @since 4.3 + */ +class SimpleMongoTransactionOptions implements MongoTransactionOptions { + + static final Set KNOWN_KEYS = Arrays.stream(OptionKey.values()).map(OptionKey::getKey) + .collect(Collectors.toSet()); + + private final Duration maxCommitTime; + private final ReadConcern readConcern; + private final ReadPreference readPreference; + private final WriteConcern writeConcern; + + static SimpleMongoTransactionOptions of(Map options) { + return new SimpleMongoTransactionOptions(options); + } + + private SimpleMongoTransactionOptions(Map options) { + + this.maxCommitTime = doGetMaxCommitTime(options); + this.readConcern = doGetReadConcern(options); + this.readPreference = doGetReadPreference(options); + this.writeConcern = doGetWriteConcern(options); + } + + @Nullable + @Override + public Duration getMaxCommitTime() { + return maxCommitTime; + } + + @Nullable + @Override + public ReadConcern getReadConcern() { + return readConcern; + } + + @Nullable + @Override + public ReadPreference getReadPreference() { + return readPreference; + } + + @Nullable + @Override + public WriteConcern getWriteConcern() { + return writeConcern; + } + + @Nullable + private static Duration doGetMaxCommitTime(Map options) { + + return getValue(options, OptionKey.MAX_COMMIT_TIME, value -> { + + Duration timeout = Duration.parse(value); + Assert.isTrue(!timeout.isNegative(), "%s cannot be negative".formatted(OptionKey.MAX_COMMIT_TIME)); + return timeout; + }); + } + + @Nullable + private static ReadConcern doGetReadConcern(Map options) { + return getValue(options, OptionKey.READ_CONCERN, value -> new ReadConcern(ReadConcernLevel.fromString(value))); + } + + @Nullable + private static ReadPreference doGetReadPreference(Map options) { + return getValue(options, OptionKey.READ_PREFERENCE, ReadPreference::valueOf); + } + + @Nullable + private static WriteConcern doGetWriteConcern(Map options) { + + return getValue(options, OptionKey.WRITE_CONCERN, value -> { + + WriteConcern writeConcern = WriteConcern.valueOf(value); + if (writeConcern == null) { + throw new IllegalArgumentException("'%s' is not a valid WriteConcern".formatted(options.get("writeConcern"))); + } + return writeConcern; + }); + } + + @Nullable + private static T getValue(Map options, OptionKey key, Function convertFunction) { + + String value = options.get(key.getKey()); + return value != null ? convertFunction.apply(value) : null; + } + + @Override + public String toString() { + + return "DefaultMongoTransactionOptions{" + "maxCommitTime=" + maxCommitTime + ", readConcern=" + readConcern + + ", readPreference=" + readPreference + ", writeConcern=" + writeConcern + '}'; + } + + enum OptionKey { + + MAX_COMMIT_TIME("maxCommitTime"), READ_CONCERN("readConcern"), READ_PREFERENCE("readPreference"), WRITE_CONCERN( + "writeConcern"); + + final String key; + + OptionKey(String key) { + this.key = key; + } + + public String getKey() { + return key; + } + + @Override + public String toString() { + return getKey(); + } + } + +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionMetadata.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionMetadata.java new file mode 100644 index 000000000..fd01d180a --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionMetadata.java @@ -0,0 +1,34 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import java.time.Duration; + +import org.springframework.lang.Nullable; + +/** + * @author Christoph Strobl + * @since 4.3 + */ +public interface TransactionMetadata { + + @Nullable + Duration getMaxCommitTime(); + + default boolean hasMaxCommitTime() { + return getMaxCommitTime() != null; + } +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionOptionResolver.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionOptionResolver.java new file mode 100644 index 000000000..fc8432690 --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/TransactionOptionResolver.java @@ -0,0 +1,29 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import org.springframework.lang.Nullable; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.interceptor.TransactionAttribute; + +/** + * @author Christoph Strobl + */ +interface TransactionOptionResolver { + + @Nullable + T resolve(TransactionDefinition attribute); +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/WriteConcernAware.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/WriteConcernAware.java new file mode 100644 index 000000000..18b2e4d4f --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/WriteConcernAware.java @@ -0,0 +1,41 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb.core; + +import org.springframework.lang.Nullable; + +import com.mongodb.ReadPreference; +import com.mongodb.WriteConcern; + +/** + * @author Christoph Strobl + * @since 4.3 + */ +public interface WriteConcernAware { + + /** + * @return {@literal true} if a {@link com.mongodb.WriteConcern} is set. + */ + default boolean hasWriteConcern() { + return getWriteConcern() != null; + } + + /** + * @return the {@link ReadPreference} to apply or {@literal null} if none set. + */ + @Nullable + WriteConcern getWriteConcern(); +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/CapturingTransactionOptionsResolver.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/CapturingTransactionOptionsResolver.java new file mode 100644 index 000000000..c3c80d632 --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/CapturingTransactionOptionsResolver.java @@ -0,0 +1,64 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; + +import org.assertj.core.api.Assertions; +import org.assertj.core.api.ListAssert; +import org.springframework.lang.Nullable; +import org.springframework.util.CollectionUtils; + +/** + * @author Christoph Strobl + */ +public class CapturingTransactionOptionsResolver implements MongoTransactionOptionsResolver { + + private final MongoTransactionOptionsResolver delegateResolver; + private final List capturedOptions = new ArrayList<>(10); + + public CapturingTransactionOptionsResolver(MongoTransactionOptionsResolver delegateResolver) { + this.delegateResolver = delegateResolver; + } + + @Nullable + @Override + public String getLabelPrefix() { + return delegateResolver.getLabelPrefix(); + } + + @Override + public MongoTransactionOptions convert(Map source) { + + MongoTransactionOptions options = delegateResolver.convert(source); + capturedOptions.add(options); + return options; + } + + public void clear() { + capturedOptions.clear(); + } + + public List getCapturedOptions() { + return capturedOptions; + } + + public MongoTransactionOptions getLastCapturedOption() { + return CollectionUtils.lastElement(capturedOptions); + } +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionOptionsUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionOptionsUnitTests.java new file mode 100644 index 000000000..688bfbbe4 --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionOptionsUnitTests.java @@ -0,0 +1,116 @@ +/* + * Copyright 2024 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import static org.assertj.core.api.Assertions.*; + +import java.time.Duration; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.Test; +import org.springframework.lang.Nullable; + +import com.mongodb.ReadConcern; +import com.mongodb.ReadPreference; +import com.mongodb.TransactionOptions; +import com.mongodb.WriteConcern; + +/** + * @author Christoph Strobl + */ +class MongoTransactionOptionsUnitTests { + + private static final TransactionOptions NATIVE_OPTIONS = TransactionOptions.builder() // + .maxCommitTime(1L, TimeUnit.SECONDS) // + .readConcern(ReadConcern.SNAPSHOT) // + .readPreference(ReadPreference.secondaryPreferred()) // + .writeConcern(WriteConcern.W3) // + .build(); + + @Test // GH-1628 + void wrapsNativeDriverTransactionOptions() { + + assertThat(MongoTransactionOptions.of(NATIVE_OPTIONS)) + .returns(NATIVE_OPTIONS.getMaxCommitTime(TimeUnit.SECONDS), options -> options.getMaxCommitTime().toSeconds()) + .returns(NATIVE_OPTIONS.getReadConcern(), MongoTransactionOptions::getReadConcern) + .returns(NATIVE_OPTIONS.getReadPreference(), MongoTransactionOptions::getReadPreference) + .returns(NATIVE_OPTIONS.getWriteConcern(), MongoTransactionOptions::getWriteConcern) + .returns(NATIVE_OPTIONS, MongoTransactionOptions::toDriverOptions); + } + + @Test // GH-1628 + void mergeNoneWithDefaultsUsesDefaults() { + + assertThat(MongoTransactionOptions.NONE.mergeWith(MongoTransactionOptions.of(NATIVE_OPTIONS))) + .returns(NATIVE_OPTIONS.getMaxCommitTime(TimeUnit.SECONDS), options -> options.getMaxCommitTime().toSeconds()) + .returns(NATIVE_OPTIONS.getReadConcern(), MongoTransactionOptions::getReadConcern) + .returns(NATIVE_OPTIONS.getReadPreference(), MongoTransactionOptions::getReadPreference) + .returns(NATIVE_OPTIONS.getWriteConcern(), MongoTransactionOptions::getWriteConcern) + .returns(NATIVE_OPTIONS, MongoTransactionOptions::toDriverOptions); + } + + @Test // GH-1628 + void mergeExistingOptionsWithNoneUsesOptions() { + + MongoTransactionOptions source = MongoTransactionOptions.of(NATIVE_OPTIONS); + assertThat(source.mergeWith(MongoTransactionOptions.NONE)).isSameAs(source); + } + + @Test // GH-1628 + void mergeExistingOptionsWithUsesFirstNonNullValue() { + + MongoTransactionOptions source = MongoTransactionOptions + .of(TransactionOptions.builder().writeConcern(WriteConcern.UNACKNOWLEDGED).build()); + + assertThat(source.mergeWith(MongoTransactionOptions.of(NATIVE_OPTIONS))) + .returns(NATIVE_OPTIONS.getMaxCommitTime(TimeUnit.SECONDS), options -> options.getMaxCommitTime().toSeconds()) + .returns(NATIVE_OPTIONS.getReadConcern(), MongoTransactionOptions::getReadConcern) + .returns(NATIVE_OPTIONS.getReadPreference(), MongoTransactionOptions::getReadPreference) + .returns(source.getWriteConcern(), MongoTransactionOptions::getWriteConcern); + } + + @Test // GH-1628 + void testEquals() { + + assertThat(MongoTransactionOptions.NONE) // + .isSameAs(MongoTransactionOptions.NONE) // + .isNotEqualTo(new MongoTransactionOptions() { + @Nullable + @Override + public Duration getMaxCommitTime() { + return null; + } + + @Nullable + @Override + public ReadConcern getReadConcern() { + return null; + } + + @Nullable + @Override + public ReadPreference getReadPreference() { + return null; + } + + @Nullable + @Override + public WriteConcern getWriteConcern() { + return null; + } + }); + } +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionUtilsUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionUtilsUnitTests.java deleted file mode 100644 index 1e2916d00..000000000 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/MongoTransactionUtilsUnitTests.java +++ /dev/null @@ -1,227 +0,0 @@ -/* - * Copyright 2023 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.data.mongodb; - -import static java.util.UUID.*; -import static org.assertj.core.api.Assertions.*; - -import java.util.Set; -import java.util.concurrent.TimeUnit; - -import org.junit.jupiter.api.Test; -import org.springframework.transaction.interceptor.DefaultTransactionAttribute; -import org.springframework.transaction.interceptor.TransactionAttribute; -import org.springframework.transaction.support.DefaultTransactionDefinition; - -import com.mongodb.ReadConcern; -import com.mongodb.ReadPreference; -import com.mongodb.TransactionOptions; -import com.mongodb.WriteConcern; - -/** - * @author Yan Kardziyaka - */ -class MongoTransactionUtilsUnitTests { - - @Test // GH-1628 - public void shouldThrowIllegalArgumentExceptionIfLabelsContainInvalidMaxCommitTime() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:maxCommitTime=-PT5S")); - - assertThatThrownBy(() -> MongoTransactionUtils.extractOptions(attribute, fallbackOptions)) // - .isInstanceOf(IllegalArgumentException.class); - } - - @Test // GH-1628 - public void shouldThrowIllegalArgumentExceptionIfLabelsContainInvalidReadConcern() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:readConcern=invalidValue")); - - assertThatThrownBy(() -> MongoTransactionUtils.extractOptions(attribute, fallbackOptions)) // - .isInstanceOf(IllegalArgumentException.class); - } - - @Test // GH-1628 - public void shouldThrowIllegalArgumentExceptionIfLabelsContainInvalidReadPreference() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:readPreference=invalidValue")); - - assertThatThrownBy(() -> MongoTransactionUtils.extractOptions(attribute, fallbackOptions)) // - .isInstanceOf(IllegalArgumentException.class); - } - - @Test // GH-1628 - public void shouldThrowIllegalArgumentExceptionIfLabelsContainInvalidWriteConcern() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:writeConcern=invalidValue")); - - assertThatThrownBy(() -> MongoTransactionUtils.extractOptions(attribute, fallbackOptions)) // - .isInstanceOf(IllegalArgumentException.class); - } - - @Test // GH-1628 - public void shouldReturnFallbackOptionsIfNotTransactionAttribute() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionDefinition definition = new DefaultTransactionDefinition(); - - TransactionOptions result = MongoTransactionUtils.extractOptions(definition, fallbackOptions); - - assertThat(result).isSameAs(fallbackOptions); - } - - @Test // GH-1628 - public void shouldReturnFallbackOptionsIfNoLabelsProvided() { - TransactionOptions fallbackOptions = getTransactionOptions(); - TransactionAttribute attribute = new DefaultTransactionAttribute(); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isSameAs(fallbackOptions); - } - - @Test // GH-1628 - public void shouldReturnFallbackOptionsIfLabelsDoesNotContainValidOptions() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - Set labels = Set.of("mongo:readConcern", "writeConcern", "readPreference=SECONDARY", - "mongo:maxCommitTime PT5M", randomUUID().toString()); - attribute.setLabels(labels); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isSameAs(fallbackOptions); - } - - @Test // GH-1628 - public void shouldReturnMergedOptionsIfLabelsContainMaxCommitTime() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:maxCommitTime=PT5S")); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isNotSameAs(fallbackOptions) // - .returns(5L, from(options -> options.getMaxCommitTime(TimeUnit.SECONDS))) // - .returns(ReadConcern.AVAILABLE, from(TransactionOptions::getReadConcern)) // - .returns(ReadPreference.secondaryPreferred(), from(TransactionOptions::getReadPreference)) // - .returns(WriteConcern.UNACKNOWLEDGED, from(TransactionOptions::getWriteConcern)); - } - - @Test // GH-1628 - public void shouldReturnMergedOptionsIfLabelsContainReadConcern() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:readConcern=majority")); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isNotSameAs(fallbackOptions) // - .returns(1L, from(options -> options.getMaxCommitTime(TimeUnit.MINUTES))) // - .returns(ReadConcern.MAJORITY, from(TransactionOptions::getReadConcern)) // - .returns(ReadPreference.secondaryPreferred(), from(TransactionOptions::getReadPreference)) // - .returns(WriteConcern.UNACKNOWLEDGED, from(TransactionOptions::getWriteConcern)); - } - - @Test // GH-1628 - public void shouldReturnMergedOptionsIfLabelsContainReadPreference() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:readPreference=primaryPreferred")); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isNotSameAs(fallbackOptions) // - .returns(1L, from(options -> options.getMaxCommitTime(TimeUnit.MINUTES))) // - .returns(ReadConcern.AVAILABLE, from(TransactionOptions::getReadConcern)) // - .returns(ReadPreference.primaryPreferred(), from(TransactionOptions::getReadPreference)) // - .returns(WriteConcern.UNACKNOWLEDGED, from(TransactionOptions::getWriteConcern)); - } - - @Test // GH-1628 - public void shouldReturnMergedOptionsIfLabelsContainWriteConcern() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - attribute.setLabels(Set.of("mongo:writeConcern=w3")); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isNotSameAs(fallbackOptions) // - .returns(1L, from(options -> options.getMaxCommitTime(TimeUnit.MINUTES))) // - .returns(ReadConcern.AVAILABLE, from(TransactionOptions::getReadConcern)) // - .returns(ReadPreference.secondaryPreferred(), from(TransactionOptions::getReadPreference)) // - .returns(WriteConcern.W3, from(TransactionOptions::getWriteConcern)); - } - - @Test // GH-1628 - public void shouldReturnNewOptionsIfLabelsContainAllOptions() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - Set labels = Set.of("mongo:maxCommitTime=PT5S", "mongo:readConcern=majority", - "mongo:readPreference=primaryPreferred", "mongo:writeConcern=w3"); - attribute.setLabels(labels); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isNotSameAs(fallbackOptions) // - .returns(5L, from(options -> options.getMaxCommitTime(TimeUnit.SECONDS))) // - .returns(ReadConcern.MAJORITY, from(TransactionOptions::getReadConcern)) // - .returns(ReadPreference.primaryPreferred(), from(TransactionOptions::getReadPreference)) // - .returns(WriteConcern.W3, from(TransactionOptions::getWriteConcern)); - } - - @Test // GH-1628 - public void shouldReturnMergedOptionsIfLabelsContainOptionsMixedWithOrdinaryStrings() { - TransactionOptions fallbackOptions = getTransactionOptions(); - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - Set labels = Set.of("mongo:maxCommitTime=PT5S", "mongo:nonExistentOption=value", "label", - "mongo:writeConcern=w3"); - attribute.setLabels(labels); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, fallbackOptions); - - assertThat(result).isNotSameAs(fallbackOptions) // - .returns(5L, from(options -> options.getMaxCommitTime(TimeUnit.SECONDS))) // - .returns(ReadConcern.AVAILABLE, from(TransactionOptions::getReadConcern)) // - .returns(ReadPreference.secondaryPreferred(), from(TransactionOptions::getReadPreference)) // - .returns(WriteConcern.W3, from(TransactionOptions::getWriteConcern)); - } - - @Test // GH-1628 - public void shouldReturnNewOptionsIFallbackIsNull() { - DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); - Set labels = Set.of("mongo:maxCommitTime=PT5S", "mongo:writeConcern=w3"); - attribute.setLabels(labels); - - TransactionOptions result = MongoTransactionUtils.extractOptions(attribute, null); - - assertThat(result).returns(5L, from(options -> options.getMaxCommitTime(TimeUnit.SECONDS))) // - .returns(null, from(TransactionOptions::getReadConcern)) // - .returns(null, from(TransactionOptions::getReadPreference)) // - .returns(WriteConcern.W3, from(TransactionOptions::getWriteConcern)); - } - - private TransactionOptions getTransactionOptions() { - return TransactionOptions.builder() // - .maxCommitTime(1L, TimeUnit.MINUTES) // - .readConcern(ReadConcern.AVAILABLE) // - .readPreference(ReadPreference.secondaryPreferred()) // - .writeConcern(WriteConcern.UNACKNOWLEDGED).build(); - } -} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java index e77c93d1e..545def16c 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionIntegrationTests.java @@ -16,6 +16,7 @@ package org.springframework.data.mongodb; import static java.util.UUID.*; +import static org.assertj.core.api.Assertions.*; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -33,6 +34,7 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.condition.DisabledIfSystemProperty; import org.junit.jupiter.api.extension.ExtendWith; +import org.junitpioneer.jupiter.SetSystemProperty; import org.springframework.context.annotation.AnnotationConfigApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -53,6 +55,9 @@ import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.reactive.TransactionalOperator; import org.springframework.transaction.support.DefaultTransactionDefinition; +import com.mongodb.ReadConcern; +import com.mongodb.ReadConcernLevel; +import com.mongodb.WriteConcern; import com.mongodb.reactivestreams.client.MongoClient; /** @@ -66,6 +71,7 @@ import com.mongodb.reactivestreams.client.MongoClient; @EnableIfMongoServerVersion(isGreaterThanEqual = "4.0") @EnableIfReplicaSetAvailable @DisabledIfSystemProperty(named = "user.name", matches = "jenkins") +@SetSystemProperty(key = "tx.read.concern", value = "local") public class ReactiveTransactionIntegrationTests { private static final String DATABASE = "rxtx-test"; @@ -76,6 +82,7 @@ public class ReactiveTransactionIntegrationTests { PersonService personService; ReactiveMongoOperations operations; ReactiveTransactionOptionsTestService transactionOptionsTestService; + CapturingTransactionOptionsResolver transactionOptionsResolver; @BeforeAll public static void init() { @@ -93,6 +100,8 @@ public class ReactiveTransactionIntegrationTests { personService = context.getBean(PersonService.class); operations = context.getBean(ReactiveMongoOperations.class); transactionOptionsTestService = context.getBean(ReactiveTransactionOptionsTestService.class); + transactionOptionsResolver = context.getBean(CapturingTransactionOptionsResolver.class); + transactionOptionsResolver.clear(); // clean out left overs from dirty context try (MongoClient client = MongoTestUtils.reactiveClient()) { @@ -251,6 +260,9 @@ public class ReactiveTransactionIntegrationTests { .expectNext(person) // .verifyComplete(); + assertThat(transactionOptionsResolver.getLastCapturedOption()).returns(Duration.ofMinutes(1), + MongoTransactionOptions::getMaxCommitTime); + operations.count(new Query(), Person.class) // .as(StepVerifier::create) // .expectNext(1L) // @@ -271,6 +283,17 @@ public class ReactiveTransactionIntegrationTests { .verifyError(TransactionSystemException.class); } + @Test // GH-1628 + public void shouldReadTransactionOptionFromSystemProperty() { + + transactionOptionsTestService.environmentReadConcernFind(randomUUID().toString()).then().as(StepVerifier::create) + .verifyComplete(); + + assertThat(transactionOptionsResolver.getLastCapturedOption()).returns( + new ReadConcern(ReadConcernLevel.fromString(System.getProperty("tx.read.concern"))), + MongoTransactionOptions::getReadConcern); + } + @Test // GH-1628 public void shouldNotThrowOnTransactionWithMajorityReadConcern() { transactionOptionsTestService.majorityReadConcernFind(randomUUID().toString()) // @@ -337,6 +360,9 @@ public class ReactiveTransactionIntegrationTests { .expectNext(person) // .verifyComplete(); + assertThat(transactionOptionsResolver.getLastCapturedOption()).returns(WriteConcern.ACKNOWLEDGED, + MongoTransactionOptions::getWriteConcern); + operations.count(new Query(), Person.class) // .as(StepVerifier::create) // .expectNext(1L) // @@ -358,8 +384,14 @@ public class ReactiveTransactionIntegrationTests { } @Bean - public ReactiveMongoTransactionManager txManager(ReactiveMongoDatabaseFactory factory) { - return new ReactiveMongoTransactionManager(factory); + CapturingTransactionOptionsResolver txOptionsResolver() { + return new CapturingTransactionOptionsResolver(MongoTransactionOptionsResolver.defaultResolver()); + } + + @Bean + public ReactiveMongoTransactionManager txManager(ReactiveMongoDatabaseFactory factory, + MongoTransactionOptionsResolver txOptionsResolver) { + return new ReactiveMongoTransactionManager(factory, txOptionsResolver, MongoTransactionOptions.NONE); } @Bean @@ -421,10 +453,10 @@ public class ReactiveTransactionIntegrationTests { new DefaultTransactionDefinition()); return Flux.merge(operations.save(new EventLog(new ObjectId(), "beforeConvert")), // - operations.save(new EventLog(new ObjectId(), "afterConvert")), // - operations.save(new EventLog(new ObjectId(), "beforeInsert")), // - operations.save(person), // - operations.save(new EventLog(new ObjectId(), "afterInsert"))) // + operations.save(new EventLog(new ObjectId(), "afterConvert")), // + operations.save(new EventLog(new ObjectId(), "beforeInsert")), // + operations.save(person), // + operations.save(new EventLog(new ObjectId(), "afterInsert"))) // .thenMany(operations.query(EventLog.class).all()) // .as(transactionalOperator::transactional); } @@ -435,10 +467,10 @@ public class ReactiveTransactionIntegrationTests { new DefaultTransactionDefinition()); return Flux.merge(operations.save(new EventLog(new ObjectId(), "beforeConvert")), // - operations.save(new EventLog(new ObjectId(), "afterConvert")), // - operations.save(new EventLog(new ObjectId(), "beforeInsert")), // - operations.save(person), // - operations.save(new EventLog(new ObjectId(), "afterInsert"))) // + operations.save(new EventLog(new ObjectId(), "afterConvert")), // + operations.save(new EventLog(new ObjectId(), "beforeInsert")), // + operations.save(person), // + operations.save(new EventLog(new ObjectId(), "afterInsert"))) // . flatMap(it -> Mono.error(new RuntimeException("poof"))) // .as(transactionalOperator::transactional); } @@ -514,8 +546,8 @@ public class ReactiveTransactionIntegrationTests { return false; } Person person = (Person) o; - return Objects.equals(id, person.id) && Objects.equals(firstname, person.firstname) && Objects.equals(lastname, - person.lastname); + return Objects.equals(id, person.id) && Objects.equals(firstname, person.firstname) + && Objects.equals(lastname, person.lastname); } @Override @@ -524,7 +556,8 @@ public class ReactiveTransactionIntegrationTests { } public String toString() { - return "ReactiveTransactionIntegrationTests.Person(id=" + this.getId() + ", firstname=" + this.getFirstname() + ", lastname=" + this.getLastname() + ")"; + return "ReactiveTransactionIntegrationTests.Person(id=" + this.getId() + ", firstname=" + this.getFirstname() + + ", lastname=" + this.getLastname() + ")"; } } diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionOptionsTestService.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionOptionsTestService.java index 15f105687..f36aeef3b 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionOptionsTestService.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/ReactiveTransactionOptionsTestService.java @@ -59,6 +59,11 @@ public class ReactiveTransactionOptionsTestService { return findByIdFunction.apply(id); } + @Transactional(transactionManager = "txManager", label = { "mongo:readConcern=${tx.read.concern}" }) + public Mono environmentReadConcernFind(Object id) { + return findByIdFunction.apply(id); + } + @Transactional(transactionManager = "txManager", label = { "mongo:readConcern=majority" }) public Mono majorityReadConcernFind(Object id) { return findByIdFunction.apply(id); diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/SimpleMongoTransactionOptionsResolverUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/SimpleMongoTransactionOptionsResolverUnitTests.java new file mode 100644 index 000000000..cfe18a579 --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/SimpleMongoTransactionOptionsResolverUnitTests.java @@ -0,0 +1,131 @@ +/* + * Copyright 2023 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.mongodb; + +import static org.assertj.core.api.Assertions.*; + +import java.util.Set; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; +import org.springframework.transaction.interceptor.DefaultTransactionAttribute; +import org.springframework.transaction.interceptor.TransactionAttribute; +import org.springframework.transaction.support.DefaultTransactionDefinition; + +import com.mongodb.ReadConcern; +import com.mongodb.ReadPreference; +import com.mongodb.WriteConcern; + +/** + * @author Yan Kardziyaka + * @author Christoph Strobl + */ +class SimpleMongoTransactionOptionsResolverUnitTests { + + @ParameterizedTest + @ValueSource(strings = { "mongo:maxCommitTime=-PT5S", "mongo:readConcern=invalidValue", + "mongo:readPreference=invalidValue", "mongo:writeConcern=invalidValue", "mongo:invalidPreference=jedi", "mongo:readConcern", "mongo:readConcern:local", "mongo:readConcern=" }) + void shouldThrowExceptionOnInvalidAttribute(String label) { + + TransactionAttribute attribute = transactionAttribute(label); + + assertThatThrownBy(() -> DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(attribute)) // + .isInstanceOf(IllegalArgumentException.class); + } + + @Test // GH-1628 + public void shouldReturnEmptyOptionsIfNotTransactionAttribute() { + + DefaultTransactionDefinition definition = new DefaultTransactionDefinition(); + assertThat(DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(definition)) + .isSameAs(MongoTransactionOptions.NONE); + } + + @Test // GH-1628 + public void shouldReturnEmptyOptionsIfNoLabelsProvided() { + + TransactionAttribute attribute = new DefaultTransactionAttribute(); + + assertThat(DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(attribute)) + .isSameAs(MongoTransactionOptions.NONE); + } + + @Test // GH-1628 + public void shouldIgnoreNonMongoOptions() { + + TransactionAttribute attribute = transactionAttribute("jpa:ignore"); + + assertThat(DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(attribute)) + .isSameAs(MongoTransactionOptions.NONE); + } + + @Test // GH-1628 + public void shouldReturnMergedOptionsIfLabelsContainMaxCommitTime() { + + TransactionAttribute attribute = transactionAttribute("mongo:maxCommitTime=PT5S"); + + assertThat(DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(attribute)) + .returns(5L, from(options -> options.getMaxCommitTime().toSeconds())) // + .returns(null, from(MongoTransactionOptions::getReadConcern)) // + .returns(null, from(MongoTransactionOptions::getReadPreference)) // + .returns(null, from(MongoTransactionOptions::getWriteConcern)); + } + + @Test // GH-1628 + public void shouldReturnReadConcernWhenPresent() { + + TransactionAttribute attribute = transactionAttribute("mongo:readConcern=majority"); + + assertThat(DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(attribute)) + .returns(null, from(TransactionMetadata::getMaxCommitTime)) // + .returns(ReadConcern.MAJORITY, from(MongoTransactionOptions::getReadConcern)) // + .returns(null, from(MongoTransactionOptions::getReadPreference)) // + .returns(null, from(MongoTransactionOptions::getWriteConcern)); + } + + @Test // GH-1628 + public void shouldReturnMergedOptionsIfLabelsContainReadPreference() { + + TransactionAttribute attribute = transactionAttribute("mongo:readPreference=primaryPreferred"); + + assertThat(DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(attribute)) + .returns(null, from(TransactionMetadata::getMaxCommitTime)) // + .returns(null, from(MongoTransactionOptions::getReadConcern)) // + .returns(ReadPreference.primaryPreferred(), from(MongoTransactionOptions::getReadPreference)) // + .returns(null, from(MongoTransactionOptions::getWriteConcern)); + } + + @Test // GH-1628 + public void shouldReturnMergedOptionsIfLabelsContainWriteConcern() { + + TransactionAttribute attribute = transactionAttribute("mongo:writeConcern=w3"); + + assertThat(DefaultMongoTransactionOptionsResolver.INSTANCE.get().resolve(attribute)) + .returns(null, from(TransactionMetadata::getMaxCommitTime)) // + .returns(null, from(MongoTransactionOptions::getReadConcern)) // + .returns(null, from(MongoTransactionOptions::getReadPreference)) // + .returns(WriteConcern.W3, from(MongoTransactionOptions::getWriteConcern)); + + } + + private static TransactionAttribute transactionAttribute(String... labels) { + + DefaultTransactionAttribute attribute = new DefaultTransactionAttribute(); + attribute.setLabels(Set.of(labels)); + return attribute; + } +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateTransactionTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateTransactionTests.java index a3afbd91e..8263f69b0 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateTransactionTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateTransactionTests.java @@ -21,6 +21,7 @@ import static org.springframework.data.mongodb.core.query.Criteria.*; import static org.springframework.data.mongodb.core.query.Query.*; import static org.springframework.data.mongodb.test.util.MongoTestUtils.*; +import java.time.Duration; import java.util.Collections; import java.util.List; import java.util.Objects; @@ -31,14 +32,18 @@ import org.bson.Document; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; +import org.junitpioneer.jupiter.SetSystemProperty; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.dao.InvalidDataAccessApiUsageException; import org.springframework.data.annotation.Id; import org.springframework.data.domain.Persistable; +import org.springframework.data.mongodb.CapturingTransactionOptionsResolver; import org.springframework.data.mongodb.MongoDatabaseFactory; import org.springframework.data.mongodb.MongoTransactionManager; +import org.springframework.data.mongodb.MongoTransactionOptions; +import org.springframework.data.mongodb.MongoTransactionOptionsResolver; import org.springframework.data.mongodb.UncategorizedMongoDbException; import org.springframework.data.mongodb.config.AbstractMongoClientConfiguration; import org.springframework.data.mongodb.test.util.AfterTransactionAssertion; @@ -56,7 +61,10 @@ import org.springframework.transaction.annotation.EnableTransactionManagement; import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; +import com.mongodb.ReadConcern; +import com.mongodb.ReadConcernLevel; import com.mongodb.ReadPreference; +import com.mongodb.WriteConcern; import com.mongodb.client.MongoClient; import com.mongodb.client.MongoCollection; import com.mongodb.client.model.Filters; @@ -71,6 +79,7 @@ import com.mongodb.client.model.Filters; @EnableIfMongoServerVersion(isGreaterThanEqual = "4.0") @ContextConfiguration @Transactional(transactionManager = "txManager") +@SetSystemProperty(key = "tx.read.concern", value = "local") public class MongoTemplateTransactionTests { static final String DB_NAME = "template-tx-tests"; @@ -98,8 +107,14 @@ public class MongoTemplateTransactionTests { } @Bean - MongoTransactionManager txManager(MongoDatabaseFactory dbFactory) { - return new MongoTransactionManager(dbFactory); + CapturingTransactionOptionsResolver txOptionsResolver() { + return new CapturingTransactionOptionsResolver(MongoTransactionOptionsResolver.defaultResolver()); + } + + @Bean + MongoTransactionManager txManager(MongoDatabaseFactory dbFactory, + MongoTransactionOptionsResolver txOptionsResolver) { + return new MongoTransactionManager(dbFactory, txOptionsResolver, MongoTransactionOptions.NONE); } @Override @@ -113,12 +128,10 @@ public class MongoTemplateTransactionTests { } } - @Autowired - MongoTemplate template; - @Autowired - MongoClient client; - @Autowired - TransactionOptionsTestService transactionOptionsTestService; + @Autowired MongoTemplate template; + @Autowired MongoClient client; + @Autowired TransactionOptionsTestService transactionOptionsTestService; + @Autowired CapturingTransactionOptionsResolver transactionOptionsResolver; List>> assertionList; @@ -127,6 +140,7 @@ public class MongoTemplateTransactionTests { template.setReadPreference(ReadPreference.primary()); assertionList = new CopyOnWriteArrayList<>(); + transactionOptionsResolver.clear(); // clean out left overs from dirty context } @BeforeTransaction @@ -144,8 +158,8 @@ public class MongoTemplateTransactionTests { boolean isPresent = collection.countDocuments(Filters.eq("_id", it.getId())) != 0; - assertThat(isPresent).isEqualTo(it.shouldBePresent()).withFailMessage( - String.format("After transaction entity %s should %s.", it.getPersistable(), + assertThat(isPresent).isEqualTo(it.shouldBePresent()) + .withFailMessage(String.format("After transaction entity %s should %s.", it.getPersistable(), it.shouldBePresent() ? "be present" : "NOT be present")); }); } @@ -205,6 +219,9 @@ public class MongoTemplateTransactionTests { transactionOptionsTestService.saveWithinMaxCommitTime(assassin); + assertThat(transactionOptionsResolver.getLastCapturedOption()).returns(Duration.ofMinutes(1), + MongoTransactionOptions::getMaxCommitTime); + assertAfterTransaction(assassin).isPresent(); } @@ -226,6 +243,18 @@ public class MongoTemplateTransactionTests { .isInstanceOf(IllegalArgumentException.class); } + @Rollback(false) + @Test // GH-1628 + @Transactional(transactionManager = "txManager", propagation = Propagation.NEVER) + public void shouldReadTransactionOptionFromSystemProperty() { + + transactionOptionsTestService.environmentReadConcernFind(randomUUID().toString()); + + assertThat(transactionOptionsResolver.getLastCapturedOption()).returns( + new ReadConcern(ReadConcernLevel.fromString(System.getProperty("tx.read.concern"))), + MongoTransactionOptions::getReadConcern); + } + @Rollback(false) @Test // GH-1628 @Transactional(transactionManager = "txManager", propagation = Propagation.NEVER) @@ -296,6 +325,9 @@ public class MongoTemplateTransactionTests { transactionOptionsTestService.acknowledgedWriteConcernSave(assassin); + assertThat(transactionOptionsResolver.getLastCapturedOption()).returns(WriteConcern.ACKNOWLEDGED, + MongoTransactionOptions::getWriteConcern); + assertAfterTransaction(assassin).isPresent(); } @@ -311,8 +343,7 @@ public class MongoTemplateTransactionTests { @org.springframework.data.mongodb.core.mapping.Document(COLLECTION_NAME) static class Assassin implements Persistable { - @Id - String id; + @Id String id; String name; public Assassin(String id, String name) { diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/TransactionOptionsTestService.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/TransactionOptionsTestService.java index f075b9134..8bfe3db38 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/TransactionOptionsTestService.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/TransactionOptionsTestService.java @@ -60,6 +60,12 @@ public class TransactionOptionsTestService { return findByIdFunction.apply(id); } + @Nullable + @Transactional(transactionManager = "txManager", label = { "mongo:readConcern=${tx.read.concern}" }) + public T environmentReadConcernFind(Object id) { + return findByIdFunction.apply(id); + } + @Nullable @Transactional(transactionManager = "txManager", label = { "mongo:readConcern=majority" }) public T majorityReadConcernFind(Object id) {