@@ -86,23 +86,9 @@
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>io.reactivex</groupId>
|
||||
<groupId>io.reactivex.rxjava3</groupId>
|
||||
<artifactId>rxjava</artifactId>
|
||||
<version>${rxjava}</version>
|
||||
<optional>true</optional>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>io.reactivex</groupId>
|
||||
<artifactId>rxjava-reactive-streams</artifactId>
|
||||
<version>${rxjava-reactive-streams}</version>
|
||||
<optional>true</optional>
|
||||
</dependency>
|
||||
|
||||
<dependency>
|
||||
<groupId>io.reactivex.rxjava2</groupId>
|
||||
<artifactId>rxjava</artifactId>
|
||||
<version>${rxjava2}</version>
|
||||
<version>${rxjava3}</version>
|
||||
<optional>true</optional>
|
||||
</dependency>
|
||||
|
||||
|
||||
@@ -35,7 +35,6 @@ import org.springframework.util.Assert;
|
||||
* @author Mark Paluch
|
||||
* @since 2.0
|
||||
* @see org.springframework.data.repository.reactive.ReactiveSortingRepository
|
||||
* @see org.springframework.data.repository.reactive.RxJava2SortingRepository
|
||||
*/
|
||||
public class ReactiveCassandraRepositoryFactoryBean<T extends Repository<S, ID>, S, ID>
|
||||
extends RepositoryFactoryBeanSupport<T, S, ID> {
|
||||
|
||||
@@ -17,13 +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.rxjava3.core.Single;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import reactor.test.StepVerifier;
|
||||
import rx.Observable;
|
||||
import rx.Single;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
@@ -40,7 +37,6 @@ import org.springframework.data.cassandra.repository.config.EnableReactiveCassan
|
||||
import org.springframework.data.cassandra.repository.support.AbstractSpringDataEmbeddedCassandraIntegrationTest;
|
||||
import org.springframework.data.cassandra.repository.support.IntegrationTestConfig;
|
||||
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
|
||||
import org.springframework.data.repository.reactive.RxJava2CrudRepository;
|
||||
import org.springframework.stereotype.Repository;
|
||||
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
|
||||
|
||||
@@ -71,8 +67,7 @@ class ConvertingReactiveCassandraRepositoryTests extends AbstractSpringDataEmbed
|
||||
@Autowired CqlSession session;
|
||||
@Autowired MixedUserRepository reactiveRepository;
|
||||
@Autowired UserRepostitory reactiveUserRepostitory;
|
||||
@Autowired RxJava1UserRepository rxJava1UserRepository;
|
||||
@Autowired RxJava2UserRepository rxJava2UserRepository;
|
||||
@Autowired RxJava3UserRepository rxJava3UserRepository;
|
||||
|
||||
private User dave;
|
||||
private User oliver;
|
||||
@@ -123,143 +118,72 @@ class ConvertingReactiveCassandraRepositoryTests extends AbstractSpringDataEmbed
|
||||
}
|
||||
|
||||
@Test // DATACASS-335
|
||||
void simpleRxJava1MethodsShouldWork() {
|
||||
void simpleRxJava3MethodsShouldWork() throws InterruptedException {
|
||||
|
||||
rxJava1UserRepository.existsById(dave.getId()) //
|
||||
rxJava3UserRepository.existsById(dave.getId()) //
|
||||
.test() //
|
||||
.awaitTerminalEvent() //
|
||||
.await() //
|
||||
.assertResult(true) //
|
||||
.assertCompleted() //
|
||||
.assertComplete() //
|
||||
.assertNoErrors();
|
||||
}
|
||||
|
||||
@Test // DATACASS-335
|
||||
void existsWithSingleRxJava1IdMethodsShouldWork() {
|
||||
void existsWithSingleRxJava3IdMethodsShouldWork() throws InterruptedException {
|
||||
|
||||
rxJava1UserRepository.existsById(Single.just(dave.getId())) //
|
||||
rxJava3UserRepository.existsById(Single.just(dave.getId())) //
|
||||
.test() //
|
||||
.awaitTerminalEvent() //
|
||||
.await() //
|
||||
.assertResult(true) //
|
||||
.assertCompleted() //
|
||||
.assertComplete() //
|
||||
.assertNoErrors();
|
||||
}
|
||||
|
||||
@Test // DATACASS-335
|
||||
void singleRxJava1QueryMethodShouldWork() {
|
||||
void singleRxJava3QueryMethodShouldWork() throws InterruptedException {
|
||||
|
||||
rxJava1UserRepository.findManyByLastname(dave.getLastname()) //
|
||||
rxJava3UserRepository.findManyByLastname(dave.getLastname()) //
|
||||
.test() //
|
||||
.awaitTerminalEvent() //
|
||||
.await() //
|
||||
.assertValueCount(2) //
|
||||
.assertNoErrors() //
|
||||
.assertCompleted();
|
||||
.assertComplete();
|
||||
}
|
||||
|
||||
@Test // DATACASS-335
|
||||
void singleProjectedRxJava1QueryMethodShouldWork() {
|
||||
void singleProjectedRxJava3QueryMethodShouldWork() throws InterruptedException {
|
||||
|
||||
List<ProjectedUser> values = rxJava1UserRepository.findProjectedByLastname(carter.getLastname()) //
|
||||
List<ProjectedUser> values = rxJava3UserRepository.findProjectedByLastname(carter.getLastname()) //
|
||||
.test() //
|
||||
.awaitTerminalEvent() //
|
||||
.await() //
|
||||
.assertValueCount(1) //
|
||||
.assertCompleted() //
|
||||
.assertComplete() //
|
||||
.assertNoErrors() //
|
||||
.getOnNextEvents();
|
||||
.values();
|
||||
|
||||
ProjectedUser projectedUser = values.get(0);
|
||||
assertThat(projectedUser.getFirstname()).isEqualTo(carter.getFirstname());
|
||||
}
|
||||
|
||||
@Test // DATACASS-335
|
||||
void observableRxJava1QueryMethodShouldWork() {
|
||||
void observableRxJava3QueryMethodShouldWork() throws InterruptedException {
|
||||
|
||||
rxJava1UserRepository.findByLastname(boyd.getLastname()) //
|
||||
rxJava3UserRepository.findByLastname(boyd.getLastname()) //
|
||||
.test() //
|
||||
.awaitTerminalEvent() //
|
||||
.await() //
|
||||
.assertValue(boyd) //
|
||||
.assertNoErrors() //
|
||||
.assertCompleted();
|
||||
}
|
||||
|
||||
@Test // DATACASS-398
|
||||
void simpleRxJava2MethodsShouldWork() {
|
||||
|
||||
rxJava2UserRepository.existsById(dave.getId()) //
|
||||
.test()//
|
||||
.assertValue(true) //
|
||||
.assertNoErrors() //
|
||||
.assertComplete() //
|
||||
.awaitTerminalEvent();
|
||||
}
|
||||
|
||||
@Test // DATACASS-398
|
||||
void existsWithSingleRxJava2IdMethodsShouldWork() {
|
||||
|
||||
rxJava2UserRepository.existsById(io.reactivex.Single.just(dave.getId())).test() //
|
||||
.assertValue(true) //
|
||||
.assertNoErrors() //
|
||||
.assertComplete() //
|
||||
.awaitTerminalEvent();
|
||||
}
|
||||
|
||||
@Test // DATACASS-398
|
||||
void flowableRxJava2QueryMethodShouldWork() {
|
||||
|
||||
rxJava2UserRepository.findManyByLastname(dave.getLastname()) //
|
||||
.test() //
|
||||
.assertValueCount(2) //
|
||||
.assertNoErrors() //
|
||||
.assertComplete() //
|
||||
.awaitTerminalEvent();
|
||||
}
|
||||
|
||||
@Test // DATACASS-398
|
||||
void singleProjectedRxJava2QueryMethodShouldWork() {
|
||||
|
||||
rxJava2UserRepository.findProjectedByLastname(Maybe.just(carter.getLastname())) //
|
||||
.test() //
|
||||
.assertValue(actual -> {
|
||||
assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname());
|
||||
return true;
|
||||
}) //
|
||||
.assertComplete() //
|
||||
.assertNoErrors() //
|
||||
.awaitTerminalEvent();
|
||||
}
|
||||
|
||||
@Test // DATACASS-398
|
||||
void observableProjectedRxJava2QueryMethodShouldWork() {
|
||||
|
||||
rxJava2UserRepository.findProjectedByLastname(Single.just(carter.getLastname())) //
|
||||
.test() //
|
||||
.assertValue(actual -> {
|
||||
assertThat(actual.getFirstname()).isEqualTo(carter.getFirstname());
|
||||
return true;
|
||||
}) //
|
||||
.assertComplete() //
|
||||
.assertNoErrors() //
|
||||
.awaitTerminalEvent();
|
||||
}
|
||||
|
||||
@Test // DATACASS-398
|
||||
void maybeRxJava2QueryMethodShouldWork() {
|
||||
|
||||
rxJava2UserRepository.findByLastname(boyd.getLastname()) //
|
||||
.test() //
|
||||
.assertValue(boyd) //
|
||||
.assertNoErrors() //
|
||||
.assertComplete() //
|
||||
.awaitTerminalEvent();
|
||||
.assertComplete();
|
||||
}
|
||||
|
||||
@Test // DATACASS-335
|
||||
void mixedRepositoryShouldWork() {
|
||||
void mixedRepositoryShouldWork() throws InterruptedException {
|
||||
|
||||
reactiveRepository.findByLastname(boyd.getLastname()) //
|
||||
.test() //
|
||||
.awaitTerminalEvent() //
|
||||
.await() //
|
||||
.assertValue(boyd) //
|
||||
.assertCompleted() //
|
||||
.assertComplete() //
|
||||
.assertNoErrors();
|
||||
}
|
||||
|
||||
@@ -280,29 +204,17 @@ class ConvertingReactiveCassandraRepositoryTests extends AbstractSpringDataEmbed
|
||||
}
|
||||
|
||||
@Repository
|
||||
interface RxJava1UserRepository extends org.springframework.data.repository.Repository<User, String> {
|
||||
interface RxJava3UserRepository extends org.springframework.data.repository.Repository<User, String> {
|
||||
|
||||
Observable<User> findManyByLastname(String lastname);
|
||||
io.reactivex.rxjava3.core.Observable<User> findManyByLastname(String lastname);
|
||||
|
||||
Single<User> findByLastname(String lastname);
|
||||
io.reactivex.rxjava3.core.Single<User> findByLastname(String lastname);
|
||||
|
||||
Single<ProjectedUser> findProjectedByLastname(String lastname);
|
||||
io.reactivex.rxjava3.core.Single<ProjectedUser> findProjectedByLastname(String lastname);
|
||||
|
||||
Single<Boolean> existsById(String id);
|
||||
io.reactivex.rxjava3.core.Single<Boolean> existsById(String id);
|
||||
|
||||
Single<Boolean> existsById(Single<String> id);
|
||||
}
|
||||
|
||||
@Repository
|
||||
interface RxJava2UserRepository extends RxJava2CrudRepository<User, String> {
|
||||
|
||||
Flowable<User> findManyByLastname(String lastname);
|
||||
|
||||
Maybe<User> findByLastname(String lastname);
|
||||
|
||||
io.reactivex.Single<ProjectedUser> findProjectedByLastname(Maybe<String> lastname);
|
||||
|
||||
io.reactivex.Observable<ProjectedUser> findProjectedByLastname(Single<String> lastname);
|
||||
io.reactivex.rxjava3.core.Single<Boolean> existsById(io.reactivex.rxjava3.core.Single<String> id);
|
||||
}
|
||||
|
||||
@Repository
|
||||
|
||||
@@ -21,6 +21,7 @@ import java.util.Collection;
|
||||
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import org.springframework.beans.factory.support.BeanDefinitionRegistry;
|
||||
import org.springframework.beans.factory.support.DefaultListableBeanFactory;
|
||||
import org.springframework.core.env.Environment;
|
||||
@@ -34,7 +35,7 @@ import org.springframework.data.repository.config.AnnotationRepositoryConfigurat
|
||||
import org.springframework.data.repository.config.RepositoryConfiguration;
|
||||
import org.springframework.data.repository.config.RepositoryConfigurationSource;
|
||||
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
|
||||
import org.springframework.data.repository.reactive.RxJava2CrudRepository;
|
||||
import org.springframework.data.repository.reactive.RxJava3CrudRepository;
|
||||
|
||||
/**
|
||||
* Unit tests for {@link ReactiveCassandraRepositoryConfigurationExtension}.
|
||||
@@ -107,7 +108,7 @@ public class ReactiveCassandraRepositoryConfigurationExtensionUnitTests {
|
||||
@Table
|
||||
private static class Sample {}
|
||||
|
||||
interface SampleRepository extends RxJava2CrudRepository<Sample, Long> {}
|
||||
interface SampleRepository extends RxJava3CrudRepository<Sample, Long> {}
|
||||
|
||||
interface UnannotatedRepository extends ReactiveCrudRepository<Object, Long> {}
|
||||
|
||||
|
||||
@@ -15,15 +15,16 @@
|
||||
*/
|
||||
package org.springframework.data.cassandra.repository.query;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.springframework.data.cassandra.core.mapping.CassandraType.Name;
|
||||
import static org.assertj.core.api.Assertions.*;
|
||||
import static org.springframework.data.cassandra.core.mapping.CassandraType.*;
|
||||
|
||||
import io.reactivex.rxjava3.core.Single;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.extension.ExtendWith;
|
||||
import org.mockito.Mock;
|
||||
@@ -39,7 +40,6 @@ import org.springframework.data.repository.core.support.DefaultRepositoryMetadat
|
||||
|
||||
import com.datastax.oss.driver.api.core.type.DataTypes;
|
||||
|
||||
import rx.Single;
|
||||
|
||||
/**
|
||||
* Unit tests for {@link ReactiveCassandraParameterAccessor}.
|
||||
|
||||
@@ -17,9 +17,9 @@ package org.springframework.data.cassandra.repository.query;
|
||||
|
||||
import static org.assertj.core.api.Assertions.*;
|
||||
|
||||
import io.reactivex.rxjava3.core.Single;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import rx.Single;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
|
||||
|
||||
@@ -18,9 +18,9 @@ package org.springframework.data.cassandra.repository.query;
|
||||
import static org.assertj.core.api.Assertions.*;
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
import io.reactivex.rxjava3.core.Single;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import rx.Single;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.util.Arrays;
|
||||
|
||||
Reference in New Issue
Block a user