diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java index 58d5e44dc..1c1dfc84e 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java @@ -1906,8 +1906,8 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati publisher = options.getCollation().map(Collation::toMongoCollation).map(publisher::collation) .orElse(publisher); publisher = options.getResumeBsonTimestamp().map(publisher::startAtOperationTime).orElse(publisher); - - if(options.getFullDocumentBeforeChangeLookup().isPresent()) { + + if (options.getFullDocumentBeforeChangeLookup().isPresent()) { publisher = publisher.fullDocumentBeforeChange(options.getFullDocumentBeforeChangeLookup().get()); } return publisher.fullDocument(options.getFullDocumentLookup().orElse(fullDocument)); @@ -2495,7 +2495,7 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati if (ObjectUtils.nullSafeEquals(WriteResultChecking.EXCEPTION, writeResultChecking)) { if (wc == null || wc.getWObject() == null - || (wc.getWObject() instanceof Number && ((Number) wc.getWObject()).intValue() < 1)) { + || (wc.getWObject()instanceof Number && ((Number) wc.getWObject()).intValue() < 1)) { return WriteConcern.ACKNOWLEDGED; } } diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTask.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTask.java index 2733c0f32..42e3bd86c 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTask.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTask.java @@ -91,9 +91,9 @@ class ChangeStreamTask extends CursorReadingTask, BsonTimestamp startAt = null; boolean resumeAfter = true; - if (options instanceof ChangeStreamRequest.ChangeStreamRequestOptions) { + if (options instanceof ChangeStreamRequest.ChangeStreamRequestOptions requestOptions) { - ChangeStreamOptions changeStreamOptions = ((ChangeStreamRequestOptions) options).getChangeStreamOptions(); + ChangeStreamOptions changeStreamOptions = requestOptions.getChangeStreamOptions(); filter = prepareFilter(template, changeStreamOptions); if (changeStreamOptions.getFilter().isPresent()) { @@ -115,9 +115,7 @@ class ChangeStreamTask extends CursorReadingTask, .orElseGet(() -> ClassUtils.isAssignable(Document.class, targetType) ? FullDocument.DEFAULT : FullDocument.UPDATE_LOOKUP); - if(changeStreamOptions.getFullDocumentBeforeChangeLookup().isPresent()) { - fullDocumentBeforeChange = changeStreamOptions.getFullDocumentBeforeChangeLookup().get(); - } + fullDocumentBeforeChange = changeStreamOptions.getFullDocumentBeforeChangeLookup().orElse(null); startAt = changeStreamOptions.getResumeBsonTimestamp().orElse(null); } @@ -158,7 +156,7 @@ class ChangeStreamTask extends CursorReadingTask, } iterable = iterable.fullDocument(fullDocument); - if(fullDocumentBeforeChange != null) { + if (fullDocumentBeforeChange != null) { iterable = iterable.fullDocumentBeforeChange(fullDocumentBeforeChange); } @@ -181,7 +179,8 @@ class ChangeStreamTask extends CursorReadingTask, template.getConverter().getMappingContext(), queryMapper) : Aggregation.DEFAULT_CONTEXT; - return agg.toPipeline(new PrefixingDelegatingAggregationOperationContext(context, "fullDocument", denylist)); + return agg + .toPipeline(new PrefixingDelegatingAggregationOperationContext(context, "fullDocument", denylist)); } if (filter instanceof List) { diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java index 9453adca5..474aa90cb 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java @@ -23,7 +23,6 @@ import static org.springframework.data.mongodb.test.util.Assertions.assertThat; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; -import org.springframework.data.mongodb.util.BsonUtils; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.test.StepVerifier; @@ -93,6 +92,7 @@ import org.springframework.data.mongodb.core.query.NearQuery; import org.springframework.data.mongodb.core.query.Query; import org.springframework.data.mongodb.core.query.Update; import org.springframework.data.mongodb.core.timeseries.Granularity; +import org.springframework.data.mongodb.util.BsonUtils; import org.springframework.data.projection.SpelAwareProxyProjectionFactory; import org.springframework.lang.Nullable; import org.springframework.test.util.ReflectionTestUtils; @@ -100,7 +100,6 @@ import org.springframework.util.CollectionUtils; import com.mongodb.MongoClientSettings; import com.mongodb.ReadPreference; -import com.mongodb.WriteConcern; import com.mongodb.client.model.CountOptions; import com.mongodb.client.model.CreateCollectionOptions; import com.mongodb.client.model.DeleteOptions; @@ -1541,11 +1540,11 @@ public class ReactiveMongoTemplateUnitTests { when(changeStreamPublisher.fullDocument(any())).thenReturn(changeStreamPublisher); when(changeStreamPublisher.fullDocumentBeforeChange(any())).thenReturn(changeStreamPublisher); - template - .changeStream("database", "collection", ChangeStreamOptions.builder().fullDocumentBeforeChangeLookup(FullDocumentBeforeChange.REQUIRED).build(), Object.class) - .subscribe(); + ChangeStreamOptions options = ChangeStreamOptions.builder() + .fullDocumentBeforeChangeLookup(FullDocumentBeforeChange.REQUIRED).build(); + template.changeStream("database", "collection", options, Object.class).subscribe(); - verify(changeStreamPublisher).fullDocumentBeforeChange(eq(FullDocumentBeforeChange.REQUIRED)); + verify(changeStreamPublisher).fullDocumentBeforeChange(FullDocumentBeforeChange.REQUIRED); } diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTaskUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTaskUnitTests.java index a629cc8de..3f56da00c 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTaskUnitTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/messaging/ChangeStreamTaskUnitTests.java @@ -41,18 +41,21 @@ import com.mongodb.client.model.changestream.ChangeStreamDocument; import com.mongodb.client.model.changestream.FullDocumentBeforeChange; /** + * Unit tests for {@link ChangeStreamTask}. + * * @author Christoph Strobl * @author Myroslav Kosinskyi */ @ExtendWith(MockitoExtension.class) +@SuppressWarnings({ "unchecked", "rawtypes" }) class ChangeStreamTaskUnitTests { - ChangeStreamTask task; @Mock MongoTemplate template; @Mock MongoDatabase mongoDatabase; @Mock MongoCollection mongoCollection; @Mock ChangeStreamIterable changeStreamIterable; - private MongoConverter converter; + + MongoConverter converter; @BeforeEach void setUp() { @@ -64,9 +67,7 @@ class ChangeStreamTaskUnitTests { when(template.getDb()).thenReturn(mongoDatabase); when(mongoDatabase.getCollection(any())).thenReturn(mongoCollection); - when(mongoCollection.watch(eq(Document.class))).thenReturn(changeStreamIterable); - when(changeStreamIterable.fullDocument(any())).thenReturn(changeStreamIterable); } @@ -84,7 +85,7 @@ class ChangeStreamTaskUnitTests { initTask(request, Document.class); - verify(changeStreamIterable).resumeAfter(eq(resumeToken)); + verify(changeStreamIterable).resumeAfter(resumeToken); } @Test // DATAMONGO-2258 @@ -102,7 +103,7 @@ class ChangeStreamTaskUnitTests { initTask(request, Document.class); - verify(changeStreamIterable).resumeAfter(eq(resumeToken)); + verify(changeStreamIterable).resumeAfter(resumeToken); } @Test // DATAMONGO-2258 @@ -120,7 +121,7 @@ class ChangeStreamTaskUnitTests { initTask(request, Document.class); - verify(changeStreamIterable).startAfter(eq(resumeToken)); + verify(changeStreamIterable).startAfter(resumeToken); } @Test // GH-4495 @@ -136,7 +137,7 @@ class ChangeStreamTaskUnitTests { initTask(request, Document.class); - verify(changeStreamIterable).fullDocumentBeforeChange(eq(FullDocumentBeforeChange.REQUIRED)); + verify(changeStreamIterable).fullDocumentBeforeChange(FullDocumentBeforeChange.REQUIRED); } private MongoCursor> initTask(ChangeStreamRequest request, Class targetType) {