From 4f6548e932e0b63cf60d78b32474831db23a78b6 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Wed, 8 Feb 2017 11:19:42 +0100 Subject: [PATCH] DATACASS-398 - Support RxJava 2 repositories. Add RxJava 2 dependency. Add test to verify RxJava 2 interoperability. Original Pull Request: #95 --- spring-data-cassandra/pom.xml | 7 ++ ...rtingReactiveCassandraRepositoryTests.java | 111 ++++++++++++++++-- 2 files changed, 111 insertions(+), 7 deletions(-) diff --git a/spring-data-cassandra/pom.xml b/spring-data-cassandra/pom.xml index 7d3dd7e8a..a7209cefd 100644 --- a/spring-data-cassandra/pom.xml +++ b/spring-data-cassandra/pom.xml @@ -92,6 +92,13 @@ true + + io.reactivex.rxjava2 + rxjava + ${rxjava2} + true + + javax.enterprise diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java index 48a8399d4..7ab7c2537 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/ConvertingReactiveCassandraRepositoryTests.java @@ -17,6 +17,10 @@ package org.springframework.data.cassandra.repository; import static org.assertj.core.api.Assertions.*; +import io.reactivex.Flowable; +import io.reactivex.Maybe; +import io.reactivex.observers.TestObserver; +import org.springframework.data.cassandra.core.ReactiveCassandraTemplate; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.test.StepVerifier; @@ -39,6 +43,7 @@ import org.springframework.data.cassandra.repository.config.EnableReactiveCassan import org.springframework.data.cassandra.test.integration.support.IntegrationTestConfig; import org.springframework.data.repository.reactive.ReactiveCrudRepository; import org.springframework.data.repository.reactive.RxJava1CrudRepository; +import org.springframework.data.repository.reactive.RxJava2CrudRepository; import org.springframework.stereotype.Repository; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -68,9 +73,11 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace } @Autowired Session session; - @Autowired MixedPersonRepostitory reactiveRepository; + @Autowired ReactiveCassandraTemplate template; + @Autowired MixedPersonRepository reactiveRepository; @Autowired PersonRepostitory reactivePersonRepostitory; @Autowired RxJava1PersonRepostitory rxJava1PersonRepostitory; + @Autowired RxJava2PersonRepostitory rxJava2PersonRepostitory; Person dave, oliver, carter, boyd; @@ -118,7 +125,8 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace } @Test // DATACASS-335 - public void simpleRxJavaMethodsShouldWork() { + public void simpleRxJava1MethodsShouldWork() { + rxJava1PersonRepostitory.exists(dave.getId()) // .test() // .awaitTerminalEvent() // @@ -128,7 +136,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace } @Test // DATACASS-335 - public void existsWithSingleRxJavaIdMethodsShouldWork() { + public void existsWithSingleRxJava1IdMethodsShouldWork() { rxJava1PersonRepostitory.exists(Single.just(dave.getId())) // .test() // @@ -139,7 +147,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace } @Test // DATACASS-335 - public void singleRxJavaQueryMethodShouldWork() { + public void singleRxJava1QueryMethodShouldWork() { rxJava1PersonRepostitory.findManyByLastname(dave.getLastname()) // .test() // @@ -150,7 +158,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace } @Test // DATACASS-335 - public void singleProjectedRxJavaQueryMethodShouldWork() { + public void singleProjectedRxJava1QueryMethodShouldWork() { List values = rxJava1PersonRepostitory.findProjectedByLastname(carter.getLastname()) // .test() // @@ -165,7 +173,7 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace } @Test // DATACASS-335 - public void observableRxJavaQueryMethodShouldWork() { + public void observableRxJava1QueryMethodShouldWork() { rxJava1PersonRepostitory.findByLastname(boyd.getLastname()) // .test() // @@ -175,6 +183,83 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace .assertCompleted(); } + @Test // DATACASS-398 + public void simpleRxJava2MethodsShouldWork() { + + TestObserver testObserver = rxJava2PersonRepostitory.exists(dave.getId()).test(); + + testObserver.awaitTerminalEvent(); + testObserver.assertComplete(); + testObserver.assertNoErrors(); + testObserver.assertValue(true); + } + + @Test // DATACASS-398 + public void existsWithSingleRxJava2IdMethodsShouldWork() { + + TestObserver testObserver = rxJava2PersonRepostitory.exists(io.reactivex.Single.just(dave.getId())).test(); + + testObserver.awaitTerminalEvent(); + testObserver.assertComplete(); + testObserver.assertNoErrors(); + testObserver.assertValue(true); + } + + @Test // DATACASS-398 + public void flowableRxJava2QueryMethodShouldWork() { + + io.reactivex.subscribers.TestSubscriber testSubscriber = rxJava2PersonRepostitory + .findManyByLastname(dave.getLastname()).test(); + + testSubscriber.awaitTerminalEvent(); + testSubscriber.assertComplete(); + testSubscriber.assertNoErrors(); + testSubscriber.assertValueCount(2); + } + + @Test // DATACASS-398 + public void singleProjectedRxJava2QueryMethodShouldWork() { + + TestObserver testObserver = rxJava2PersonRepostitory + .findProjectedByLastname(Maybe.just(carter.getLastname())).test(); + + testObserver.awaitTerminalEvent(); + testObserver.assertComplete(); + testObserver.assertNoErrors(); + + testObserver.assertValue(actual -> { + assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname()); + return true; + }); + } + + @Test // DATACASS-398 + public void observableProjectedRxJava2QueryMethodShouldWork() { + + TestObserver testObserver = rxJava2PersonRepostitory + .findProjectedByLastname(Single.just(carter.getLastname())).test(); + + testObserver.awaitTerminalEvent(); + testObserver.assertComplete(); + testObserver.assertNoErrors(); + + testObserver.assertValue(actual -> { + assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname()); + return true; + }); + } + + @Test // DATACASS-398 + public void maybeRxJava2QueryMethodShouldWork() { + + TestObserver testObserver = rxJava2PersonRepostitory.findByLastname(boyd.getLastname()).test(); + + testObserver.awaitTerminalEvent(); + testObserver.assertComplete(); + testObserver.assertNoErrors(); + testObserver.assertValue(boyd); + } + @Test // DATACASS-335 public void mixedRepositoryShouldWork() { @@ -214,7 +299,19 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace } @Repository - interface MixedPersonRepostitory extends ReactiveCassandraRepository { + interface RxJava2PersonRepostitory extends RxJava2CrudRepository { + + Flowable findManyByLastname(String lastname); + + Maybe findByLastname(String lastname); + + io.reactivex.Single findProjectedByLastname(Maybe lastname); + + io.reactivex.Observable findProjectedByLastname(Single lastname); + } + + @Repository + interface MixedPersonRepository extends ReactiveCassandraRepository { Single findByLastname(String lastname);