Use the official RxJava to Reactive Streams adapter
This commit removes the usage of Reactor adapters (about to be moved from Reactor Core to a new Reactor Adapter module). Instead, RxReactiveStreams is now used for adapting RxJava 1 and Flowable methods are used for RxJava 2. Issue: SPR-14824
This commit is contained in:
@@ -26,11 +26,10 @@ import io.reactivex.BackpressureStrategy;
|
||||
import io.reactivex.Flowable;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import reactor.adapter.RxJava1Adapter;
|
||||
import reactor.adapter.RxJava2Adapter;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import rx.Observable;
|
||||
import rx.RxReactiveStreams;
|
||||
import rx.Single;
|
||||
|
||||
import org.springframework.core.MethodParameter;
|
||||
@@ -151,7 +150,7 @@ public class HttpEntityArgumentResolverTests {
|
||||
ResolvableType type = httpEntityType(forClassWithGenerics(Single.class, String.class));
|
||||
HttpEntity<Single<String>> entity = resolveValueWithEmptyBody(type);
|
||||
|
||||
TestSubscriber.subscribe(RxJava1Adapter.singleToMono(entity.getBody()))
|
||||
TestSubscriber.subscribe(RxReactiveStreams.toPublisher(entity.getBody()))
|
||||
.assertNoValues()
|
||||
.assertError(ServerWebInputException.class);
|
||||
}
|
||||
@@ -161,7 +160,7 @@ public class HttpEntityArgumentResolverTests {
|
||||
ResolvableType type = httpEntityType(forClassWithGenerics(io.reactivex.Single.class, String.class));
|
||||
HttpEntity<io.reactivex.Single<String>> entity = resolveValueWithEmptyBody(type);
|
||||
|
||||
TestSubscriber.subscribe(RxJava2Adapter.singleToMono(entity.getBody()))
|
||||
TestSubscriber.subscribe(entity.getBody().toFlowable())
|
||||
.assertNoValues()
|
||||
.assertError(ServerWebInputException.class);
|
||||
}
|
||||
@@ -171,7 +170,7 @@ public class HttpEntityArgumentResolverTests {
|
||||
ResolvableType type = httpEntityType(forClassWithGenerics(Observable.class, String.class));
|
||||
HttpEntity<Observable<String>> entity = resolveValueWithEmptyBody(type);
|
||||
|
||||
TestSubscriber.subscribe(RxJava1Adapter.observableToFlux(entity.getBody()))
|
||||
TestSubscriber.subscribe(RxReactiveStreams.toPublisher(entity.getBody()))
|
||||
.assertNoError()
|
||||
.assertComplete()
|
||||
.assertNoValues();
|
||||
@@ -182,7 +181,7 @@ public class HttpEntityArgumentResolverTests {
|
||||
ResolvableType type = httpEntityType(forClassWithGenerics(io.reactivex.Observable.class, String.class));
|
||||
HttpEntity<io.reactivex.Observable<String>> entity = resolveValueWithEmptyBody(type);
|
||||
|
||||
TestSubscriber.subscribe(RxJava2Adapter.observableToFlux(entity.getBody(), BackpressureStrategy.BUFFER))
|
||||
TestSubscriber.subscribe(entity.getBody().toFlowable(BackpressureStrategy.BUFFER))
|
||||
.assertNoError()
|
||||
.assertComplete()
|
||||
.assertNoValues();
|
||||
@@ -193,7 +192,7 @@ public class HttpEntityArgumentResolverTests {
|
||||
ResolvableType type = httpEntityType(forClassWithGenerics(Flowable.class, String.class));
|
||||
HttpEntity<Flowable<String>> entity = resolveValueWithEmptyBody(type);
|
||||
|
||||
TestSubscriber.subscribe(RxJava2Adapter.flowableToFlux(entity.getBody()))
|
||||
TestSubscriber.subscribe(entity.getBody())
|
||||
.assertNoError()
|
||||
.assertComplete()
|
||||
.assertNoValues();
|
||||
|
||||
@@ -24,10 +24,10 @@ import java.util.function.Predicate;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import reactor.adapter.RxJava1Adapter;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
import rx.Observable;
|
||||
import rx.RxReactiveStreams;
|
||||
import rx.Single;
|
||||
|
||||
import org.springframework.core.MethodParameter;
|
||||
@@ -151,12 +151,12 @@ public class RequestBodyArgumentResolverTests {
|
||||
ResolvableType type = forClassWithGenerics(Single.class, String.class);
|
||||
|
||||
Single<String> single = resolveValueWithEmptyBody(type, true);
|
||||
TestSubscriber.subscribe(RxJava1Adapter.singleToMono(single))
|
||||
TestSubscriber.subscribe(RxReactiveStreams.toPublisher(single))
|
||||
.assertNoValues()
|
||||
.assertError(ServerWebInputException.class);
|
||||
|
||||
single = resolveValueWithEmptyBody(type, false);
|
||||
TestSubscriber.subscribe(RxJava1Adapter.singleToMono(single))
|
||||
TestSubscriber.subscribe(RxReactiveStreams.toPublisher(single))
|
||||
.assertNoValues()
|
||||
.assertError(ServerWebInputException.class);
|
||||
}
|
||||
@@ -166,12 +166,12 @@ public class RequestBodyArgumentResolverTests {
|
||||
ResolvableType type = forClassWithGenerics(Observable.class, String.class);
|
||||
|
||||
Observable<String> observable = resolveValueWithEmptyBody(type, true);
|
||||
TestSubscriber.subscribe(RxJava1Adapter.observableToFlux(observable))
|
||||
TestSubscriber.subscribe(RxReactiveStreams.toPublisher(observable))
|
||||
.assertNoValues()
|
||||
.assertError(ServerWebInputException.class);
|
||||
|
||||
observable = resolveValueWithEmptyBody(type, false);
|
||||
TestSubscriber.subscribe(RxJava1Adapter.observableToFlux(observable))
|
||||
TestSubscriber.subscribe(RxReactiveStreams.toPublisher(observable))
|
||||
.assertNoValues()
|
||||
.assertComplete();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user