diff --git a/src/main/java/org/springframework/data/repository/core/support/QueryExecutionResultHandler.java b/src/main/java/org/springframework/data/repository/core/support/QueryExecutionResultHandler.java index f585c3386..8fc5e3f8e 100644 --- a/src/main/java/org/springframework/data/repository/core/support/QueryExecutionResultHandler.java +++ b/src/main/java/org/springframework/data/repository/core/support/QueryExecutionResultHandler.java @@ -23,6 +23,7 @@ import org.springframework.core.convert.support.DefaultConversionService; import org.springframework.core.convert.support.GenericConversionService; import org.springframework.data.repository.util.NullableWrapper; import org.springframework.data.repository.util.QueryExecutionConverters; +import org.springframework.data.repository.util.ReactiveWrapperConverters; /** * Simple domain service to convert query results into a dedicated type. @@ -43,6 +44,7 @@ class QueryExecutionResultHandler { GenericConversionService conversionService = new DefaultConversionService(); QueryExecutionConverters.registerConvertersIn(conversionService); + ReactiveWrapperConverters.registerConvertersIn(conversionService); this.conversionService = conversionService; } diff --git a/src/main/java/org/springframework/data/repository/util/QueryExecutionConverters.java b/src/main/java/org/springframework/data/repository/util/QueryExecutionConverters.java index 5229c9c06..601325e82 100644 --- a/src/main/java/org/springframework/data/repository/util/QueryExecutionConverters.java +++ b/src/main/java/org/springframework/data/repository/util/QueryExecutionConverters.java @@ -21,14 +21,11 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Future; -import org.reactivestreams.Publisher; -import org.springframework.core.ReactiveAdapterRegistry; import org.springframework.core.convert.ConversionService; import org.springframework.core.convert.TypeDescriptor; import org.springframework.core.convert.converter.Converter; import org.springframework.core.convert.converter.GenericConverter; import org.springframework.core.convert.support.ConfigurableConversionService; -import org.springframework.data.repository.util.ReactiveWrappers.ReactiveLibrary; import org.springframework.scheduling.annotation.AsyncResult; import org.springframework.util.Assert; import org.springframework.util.ClassUtils; @@ -36,14 +33,6 @@ import org.springframework.util.concurrent.ListenableFuture; import com.google.common.base.Optional; -import io.reactivex.BackpressureStrategy; -import io.reactivex.Flowable; -import io.reactivex.Maybe; -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; -import rx.Completable; -import rx.Observable; -import rx.Single; import scala.Function0; import scala.Option; import scala.runtime.AbstractFunction0; @@ -85,7 +74,6 @@ public abstract class QueryExecutionConverters { private static final Set> WRAPPER_TYPES = new HashSet>(); private static final Set> UNWRAPPER_TYPES = new HashSet>(); private static final Set> UNWRAPPERS = new HashSet>(); - private static final ReactiveAdapterRegistry REACTIVE_ADAPTER_REGISTRY = new ReactiveAdapterRegistry(); static { @@ -117,9 +105,11 @@ public abstract class QueryExecutionConverters { UNWRAPPERS.add(ScalOptionUnwrapper.INSTANCE); } - WRAPPER_TYPES.addAll(ReactiveWrappers.getNoValueTypes()); - WRAPPER_TYPES.addAll(ReactiveWrappers.getSingleValueTypes()); - WRAPPER_TYPES.addAll(ReactiveWrappers.getMultiValueTypes()); + if(ReactiveWrappers.isAvailable()) { + WRAPPER_TYPES.addAll(ReactiveWrappers.getNoValueTypes()); + WRAPPER_TYPES.addAll(ReactiveWrappers.getSingleValueTypes()); + WRAPPER_TYPES.addAll(ReactiveWrappers.getMultiValueTypes()); + } } private QueryExecutionConverters() {} @@ -187,67 +177,6 @@ public abstract class QueryExecutionConverters { if (ASYNC_RESULT_PRESENT) { conversionService.addConverter(new NullableWrapperToFutureConverter(conversionService)); } - - if (ReactiveWrappers.isAvailable()) { - - if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) { - - conversionService.addConverter(PublisherToRxJava1CompletableConverter.INSTANCE); - conversionService.addConverter(RxJava1CompletableToPublisherConverter.INSTANCE); - conversionService.addConverter(RxJava1CompletableToMonoConverter.INSTANCE); - - conversionService.addConverter(PublisherToRxJava1SingleConverter.INSTANCE); - conversionService.addConverter(RxJava1SingleToPublisherConverter.INSTANCE); - conversionService.addConverter(RxJava1SingleToMonoConverter.INSTANCE); - conversionService.addConverter(RxJava1SingleToFluxConverter.INSTANCE); - - conversionService.addConverter(PublisherToRxJava1ObservableConverter.INSTANCE); - conversionService.addConverter(RxJava1ObservableToPublisherConverter.INSTANCE); - conversionService.addConverter(RxJava1ObservableToMonoConverter.INSTANCE); - conversionService.addConverter(RxJava1ObservableToFluxConverter.INSTANCE); - } - - if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) { - - conversionService.addConverter(PublisherToRxJava2CompletableConverter.INSTANCE); - conversionService.addConverter(RxJava2CompletableToPublisherConverter.INSTANCE); - conversionService.addConverter(RxJava2CompletableToMonoConverter.INSTANCE); - - conversionService.addConverter(PublisherToRxJava2SingleConverter.INSTANCE); - conversionService.addConverter(RxJava2SingleToPublisherConverter.INSTANCE); - conversionService.addConverter(RxJava2SingleToMonoConverter.INSTANCE); - conversionService.addConverter(RxJava2SingleToFluxConverter.INSTANCE); - - conversionService.addConverter(PublisherToRxJava2ObservableConverter.INSTANCE); - conversionService.addConverter(RxJava2ObservableToPublisherConverter.INSTANCE); - conversionService.addConverter(RxJava2ObservableToMonoConverter.INSTANCE); - conversionService.addConverter(RxJava2ObservableToFluxConverter.INSTANCE); - - conversionService.addConverter(PublisherToRxJava2FlowableConverter.INSTANCE); - conversionService.addConverter(RxJava2FlowableToPublisherConverter.INSTANCE); - - conversionService.addConverter(PublisherToRxJava2MaybeConverter.INSTANCE); - conversionService.addConverter(RxJava2MaybeToPublisherConverter.INSTANCE); - conversionService.addConverter(RxJava2MaybeToMonoConverter.INSTANCE); - conversionService.addConverter(RxJava2MaybeToFluxConverter.INSTANCE); - } - - if (ReactiveWrappers.isAvailable(ReactiveLibrary.PROJECT_REACTOR)) { - conversionService.addConverter(PublisherToMonoConverter.INSTANCE); - conversionService.addConverter(PublisherToFluxConverter.INSTANCE); - } - - if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) { - conversionService.addConverter(RxJava1SingleToObservableConverter.INSTANCE); - conversionService.addConverter(RxJava1ObservableToSingleConverter.INSTANCE); - } - - if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) { - conversionService.addConverter(RxJava2SingleToObservableConverter.INSTANCE); - conversionService.addConverter(RxJava2ObservableToSingleConverter.INSTANCE); - conversionService.addConverter(RxJava2ObservableToMaybeConverter.INSTANCE); - } - } } /** @@ -557,570 +486,4 @@ public abstract class QueryExecutionConverters { return source instanceof Option ? ((Option) source).getOrElse(alternative) : source; } } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link Flux}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToFluxConverter implements Converter, Flux> { - - INSTANCE; - - @Override - public Flux convert(Publisher source) { - return Flux.from(source); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToMonoConverter implements Converter, Mono> { - - INSTANCE; - - @Override - public Mono convert(Publisher source) { - return Mono.from(source); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link Single}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava1SingleConverter implements Converter, Single> { - - INSTANCE; - - @Override - public Single convert(Publisher source) { - return (Single) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Single.class).fromPublisher(Mono.from(source)); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link Completable}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava1CompletableConverter implements Converter, Completable> { - - INSTANCE; - - @Override - public Completable convert(Publisher source) { - return (Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Completable.class).fromPublisher(source); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link Observable}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava1ObservableConverter implements Converter, Observable> { - - INSTANCE; - - @Override - public Observable convert(Publisher source) { - return (Observable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Observable.class).fromPublisher(Flux.from(source)); - } - } - - /** - * A {@link Converter} to convert a {@link Single} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1SingleToPublisherConverter implements Converter, Publisher> { - - INSTANCE; - - @Override - public Publisher convert(Single source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toPublisher(source); - } - } - - /** - * A {@link Converter} to convert a {@link Single} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1SingleToMonoConverter implements Converter, Mono> { - - INSTANCE; - - @Override - public Mono convert(Single source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toMono(source); - } - } - - /** - * A {@link Converter} to convert a {@link Single} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1SingleToFluxConverter implements Converter, Flux> { - - INSTANCE; - - @Override - public Flux convert(Single source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toFlux(source); - } - } - - /** - * A {@link Converter} to convert a {@link Completable} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1CompletableToPublisherConverter implements Converter> { - - INSTANCE; - - @Override - public Publisher convert(Completable source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Completable.class).toFlux(source); - } - } - - /** - * A {@link Converter} to convert a {@link Completable} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1CompletableToMonoConverter implements Converter> { - - INSTANCE; - - @Override - public Mono convert(Completable source) { - return Mono.from(RxJava1CompletableToPublisherConverter.INSTANCE.convert(source)); - } - } - - /** - * A {@link Converter} to convert an {@link Observable} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1ObservableToPublisherConverter implements Converter, Publisher> { - - INSTANCE; - - @Override - public Publisher convert(Observable source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source); - } - } - - /** - * A {@link Converter} to convert a {@link Observable} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1ObservableToMonoConverter implements Converter, Mono> { - - INSTANCE; - - @Override - public Mono convert(Observable source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toMono(source); - } - } - - /** - * A {@link Converter} to convert a {@link Observable} to {@link Flux}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1ObservableToFluxConverter implements Converter, Flux> { - - INSTANCE; - - @Override - public Flux convert(Observable source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Single}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava2SingleConverter implements Converter, io.reactivex.Single> { - - INSTANCE; - - @Override - public io.reactivex.Single convert(Publisher source) { - return (io.reactivex.Single) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class) - .fromPublisher(source); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Completable}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava2CompletableConverter implements Converter, io.reactivex.Completable> { - - INSTANCE; - - @Override - public io.reactivex.Completable convert(Publisher source) { - return (io.reactivex.Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Completable.class) - .fromPublisher(source); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Observable}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava2ObservableConverter implements Converter, io.reactivex.Observable> { - - INSTANCE; - - @Override - public io.reactivex.Observable convert(Publisher source) { - return (io.reactivex.Observable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class) - .fromPublisher(source); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2SingleToPublisherConverter implements Converter, Publisher> { - - INSTANCE; - - @Override - public Publisher convert(io.reactivex.Single source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2SingleToMonoConverter implements Converter, Mono> { - - INSTANCE; - - @Override - public Mono convert(io.reactivex.Single source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2SingleToFluxConverter implements Converter, Flux> { - - INSTANCE; - - @Override - public Flux convert(io.reactivex.Single source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toFlux(source); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2CompletableToPublisherConverter implements Converter> { - - INSTANCE; - - @Override - public Publisher convert(io.reactivex.Completable source) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Completable.class).toFlux(source); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2CompletableToMonoConverter implements Converter> { - - INSTANCE; - - @Override - public Mono convert(io.reactivex.Completable source) { - return Mono.from(RxJava2CompletableToPublisherConverter.INSTANCE.convert(source)); - } - } - - /** - * A {@link Converter} to convert an {@link io.reactivex.Observable} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2ObservableToPublisherConverter implements Converter, Publisher> { - - INSTANCE; - - @Override - public Publisher convert(io.reactivex.Observable source) { - return source.toFlowable(BackpressureStrategy.BUFFER); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2ObservableToMonoConverter implements Converter, Mono> { - - INSTANCE; - - @Override - public Mono convert(io.reactivex.Observable source) { - return Mono.from(source.toFlowable(BackpressureStrategy.BUFFER)); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Flux}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2ObservableToFluxConverter implements Converter, Flux> { - - INSTANCE; - - @Override - public Flux convert(io.reactivex.Observable source) { - return Flux.from(source.toFlowable(BackpressureStrategy.BUFFER)); - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava2FlowableConverter implements Converter, io.reactivex.Flowable> { - - INSTANCE; - - @Override - public io.reactivex.Flowable convert(Publisher source) { - return Flowable.fromPublisher(source); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Flowable} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2FlowableToPublisherConverter implements Converter, Publisher> { - - INSTANCE; - - @Override - public Publisher convert(io.reactivex.Flowable source) { - return source; - } - } - - /** - * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum PublisherToRxJava2MaybeConverter implements Converter, io.reactivex.Maybe> { - - INSTANCE; - - @Override - public io.reactivex.Maybe convert(Publisher source) { - return (io.reactivex.Maybe) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Maybe.class).fromPublisher(source); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Publisher}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2MaybeToPublisherConverter implements Converter, Publisher> { - - INSTANCE; - - @Override - public Publisher convert(io.reactivex.Maybe source) { - return source.toFlowable(); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Mono}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2MaybeToMonoConverter implements Converter, Mono> { - - INSTANCE; - - @Override - public Mono convert(io.reactivex.Maybe source) { - return Mono.from(source.toFlowable()); - } - } - - /** - * A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Flux}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2MaybeToFluxConverter implements Converter, Flux> { - - INSTANCE; - - @Override - public Flux convert(io.reactivex.Maybe source) { - return Flux.from(source.toFlowable()); - } - } - - /** - * A {@link Converter} to convert a {@link Observable} to {@link Single}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1ObservableToSingleConverter implements Converter, Single> { - - INSTANCE; - - @Override - public Single convert(Observable source) { - return source.toSingle(); - } - } - - /** - * A {@link Converter} to convert a {@link Single} to {@link Single}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava1SingleToObservableConverter implements Converter, Observable> { - - INSTANCE; - - @Override - public Observable convert(Single source) { - return source.toObservable(); - } - } - - /** - * A {@link Converter} to convert a {@link Observable} to {@link Single}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2ObservableToSingleConverter - implements Converter, io.reactivex.Single> { - - INSTANCE; - - @Override - public io.reactivex.Single convert(io.reactivex.Observable source) { - return source.singleOrError(); - } - } - - /** - * A {@link Converter} to convert a {@link Observable} to {@link Maybe}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2ObservableToMaybeConverter - implements Converter, io.reactivex.Maybe> { - - INSTANCE; - - @Override - public io.reactivex.Maybe convert(io.reactivex.Observable source) { - return source.singleElement(); - } - } - - /** - * A {@link Converter} to convert a {@link Single} to {@link Single}. - * - * @author Mark Paluch - * @author 2.0 - */ - public enum RxJava2SingleToObservableConverter - implements Converter, io.reactivex.Observable> { - - INSTANCE; - - @Override - public io.reactivex.Observable convert(io.reactivex.Single source) { - return source.toObservable(); - } - } } diff --git a/src/main/java/org/springframework/data/repository/util/ReactiveWrapperConverters.java b/src/main/java/org/springframework/data/repository/util/ReactiveWrapperConverters.java index 20d4c63d0..6e5ed26aa 100644 --- a/src/main/java/org/springframework/data/repository/util/ReactiveWrapperConverters.java +++ b/src/main/java/org/springframework/data/repository/util/ReactiveWrapperConverters.java @@ -15,26 +15,32 @@ */ package org.springframework.data.repository.util; +import static org.springframework.data.repository.util.ReactiveWrapperConverters.RegistryHolder.*; + import java.util.ArrayList; import java.util.List; import org.reactivestreams.Publisher; import org.springframework.core.ReactiveAdapterRegistry; import org.springframework.core.convert.converter.Converter; +import org.springframework.core.convert.support.ConfigurableConversionService; import org.springframework.core.convert.support.GenericConversionService; import org.springframework.data.repository.util.ReactiveWrappers.ReactiveLibrary; import org.springframework.util.Assert; import org.springframework.util.ClassUtils; +import io.reactivex.BackpressureStrategy; import io.reactivex.Flowable; +import io.reactivex.Maybe; import lombok.experimental.UtilityClass; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import rx.Completable; import rx.Observable; import rx.Single; /** - * Conversion support for reactive wrapper types. This class is a logical extension to {@link QueryExecutionConverters}. + * Conversion support for reactive wrapper types. This class is a reactive extension to {@link QueryExecutionConverters}. *

* This class discovers reactive wrapper availability and their conversion support based on the class path. Reactive * wrapper types might be supported/on the class path but conversion may require additional dependencies. @@ -49,7 +55,6 @@ public class ReactiveWrapperConverters { private static final List> REACTIVE_WRAPPERS = new ArrayList<>(); private static final GenericConversionService GENERIC_CONVERSION_SERVICE = new GenericConversionService(); - private static final ReactiveAdapterRegistry REACTIVE_ADAPTER_REGISTRY = new ReactiveAdapterRegistry(); static { @@ -74,7 +79,76 @@ public class ReactiveWrapperConverters { REACTIVE_WRAPPERS.add(PublisherWrapper.INSTANCE); } - QueryExecutionConverters.registerConvertersIn(GENERIC_CONVERSION_SERVICE); + registerConvertersIn(GENERIC_CONVERSION_SERVICE); + } + + /** + * Registers converters for wrapper types found on the classpath. + * + * @param conversionService must not be {@literal null}. + */ + public static void registerConvertersIn(ConfigurableConversionService conversionService) { + + Assert.notNull(conversionService, "ConversionService must not be null!"); + + if (ReactiveWrappers.isAvailable(ReactiveLibrary.PROJECT_REACTOR)) { + + if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) { + + conversionService.addConverter(PublisherToRxJava1CompletableConverter.INSTANCE); + conversionService.addConverter(RxJava1CompletableToPublisherConverter.INSTANCE); + conversionService.addConverter(RxJava1CompletableToMonoConverter.INSTANCE); + + conversionService.addConverter(PublisherToRxJava1SingleConverter.INSTANCE); + conversionService.addConverter(RxJava1SingleToPublisherConverter.INSTANCE); + conversionService.addConverter(RxJava1SingleToMonoConverter.INSTANCE); + conversionService.addConverter(RxJava1SingleToFluxConverter.INSTANCE); + + conversionService.addConverter(PublisherToRxJava1ObservableConverter.INSTANCE); + conversionService.addConverter(RxJava1ObservableToPublisherConverter.INSTANCE); + conversionService.addConverter(RxJava1ObservableToMonoConverter.INSTANCE); + conversionService.addConverter(RxJava1ObservableToFluxConverter.INSTANCE); + } + + if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) { + + conversionService.addConverter(PublisherToRxJava2CompletableConverter.INSTANCE); + conversionService.addConverter(RxJava2CompletableToPublisherConverter.INSTANCE); + conversionService.addConverter(RxJava2CompletableToMonoConverter.INSTANCE); + + conversionService.addConverter(PublisherToRxJava2SingleConverter.INSTANCE); + conversionService.addConverter(RxJava2SingleToPublisherConverter.INSTANCE); + conversionService.addConverter(RxJava2SingleToMonoConverter.INSTANCE); + conversionService.addConverter(RxJava2SingleToFluxConverter.INSTANCE); + + conversionService.addConverter(PublisherToRxJava2ObservableConverter.INSTANCE); + conversionService.addConverter(RxJava2ObservableToPublisherConverter.INSTANCE); + conversionService.addConverter(RxJava2ObservableToMonoConverter.INSTANCE); + conversionService.addConverter(RxJava2ObservableToFluxConverter.INSTANCE); + + conversionService.addConverter(PublisherToRxJava2FlowableConverter.INSTANCE); + conversionService.addConverter(RxJava2FlowableToPublisherConverter.INSTANCE); + + conversionService.addConverter(PublisherToRxJava2MaybeConverter.INSTANCE); + conversionService.addConverter(RxJava2MaybeToPublisherConverter.INSTANCE); + conversionService.addConverter(RxJava2MaybeToMonoConverter.INSTANCE); + conversionService.addConverter(RxJava2MaybeToFluxConverter.INSTANCE); + } + + conversionService.addConverter(PublisherToMonoConverter.INSTANCE); + conversionService.addConverter(PublisherToFluxConverter.INSTANCE); + + if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA1)) { + conversionService.addConverter(RxJava1SingleToObservableConverter.INSTANCE); + conversionService.addConverter(RxJava1ObservableToSingleConverter.INSTANCE); + } + + if (ReactiveWrappers.isAvailable(ReactiveLibrary.RXJAVA2)) { + conversionService.addConverter(RxJava2SingleToObservableConverter.INSTANCE); + conversionService.addConverter(RxJava2ObservableToSingleConverter.INSTANCE); + conversionService.addConverter(RxJava2ObservableToMaybeConverter.INSTANCE); + } + } } /** @@ -88,7 +162,7 @@ public class ReactiveWrapperConverters { * @return {@literal true} if the {@code type} is a supported reactive wrapper type. */ public static boolean supports(Class type) { - return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(type) != null; + return RegistryHolder.REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(type) != null; } /** @@ -283,4 +357,580 @@ public class ReactiveWrapperConverters { return ((io.reactivex.Flowable) wrapper).map(converter::convert); } } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link Flux}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToFluxConverter implements Converter, Flux> { + + INSTANCE; + + @Override + public Flux convert(Publisher source) { + return Flux.from(source); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToMonoConverter implements Converter, Mono> { + + INSTANCE; + + @Override + public Mono convert(Publisher source) { + return Mono.from(source); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link Single}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava1SingleConverter implements Converter, Single> { + + INSTANCE; + + @Override + public Single convert(Publisher source) { + return (Single) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Single.class).fromPublisher(Mono.from(source)); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link Completable}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava1CompletableConverter implements Converter, Completable> { + + INSTANCE; + + @Override + public Completable convert(Publisher source) { + return (Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Completable.class).fromPublisher(source); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link Observable}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava1ObservableConverter implements Converter, Observable> { + + INSTANCE; + + @Override + public Observable convert(Publisher source) { + return (Observable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Observable.class).fromPublisher(Flux.from(source)); + } + } + + /** + * A {@link Converter} to convert a {@link Single} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1SingleToPublisherConverter implements Converter, Publisher> { + + INSTANCE; + + @Override + public Publisher convert(Single source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toPublisher(source); + } + } + + /** + * A {@link Converter} to convert a {@link Single} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1SingleToMonoConverter implements Converter, Mono> { + + INSTANCE; + + @Override + public Mono convert(Single source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toMono(source); + } + } + + /** + * A {@link Converter} to convert a {@link Single} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1SingleToFluxConverter implements Converter, Flux> { + + INSTANCE; + + @Override + public Flux convert(Single source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Single.class).toFlux(source); + } + } + + /** + * A {@link Converter} to convert a {@link Completable} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1CompletableToPublisherConverter implements Converter> { + + INSTANCE; + + @Override + public Publisher convert(Completable source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Completable.class).toFlux(source); + } + } + + /** + * A {@link Converter} to convert a {@link Completable} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1CompletableToMonoConverter implements Converter> { + + INSTANCE; + + @Override + public Mono convert(Completable source) { + return Mono.from(RxJava1CompletableToPublisherConverter.INSTANCE.convert(source)); + } + } + + /** + * A {@link Converter} to convert an {@link Observable} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1ObservableToPublisherConverter implements Converter, Publisher> { + + INSTANCE; + + @Override + public Publisher convert(Observable source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source); + } + } + + /** + * A {@link Converter} to convert a {@link Observable} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1ObservableToMonoConverter implements Converter, Mono> { + + INSTANCE; + + @Override + public Mono convert(Observable source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toMono(source); + } + } + + /** + * A {@link Converter} to convert a {@link Observable} to {@link Flux}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1ObservableToFluxConverter implements Converter, Flux> { + + INSTANCE; + + @Override + public Flux convert(Observable source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(Observable.class).toFlux(source); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Single}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava2SingleConverter implements Converter, io.reactivex.Single> { + + INSTANCE; + + @Override + public io.reactivex.Single convert(Publisher source) { + return (io.reactivex.Single) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class) + .fromPublisher(source); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Completable}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava2CompletableConverter implements Converter, io.reactivex.Completable> { + + INSTANCE; + + @Override + public io.reactivex.Completable convert(Publisher source) { + return (io.reactivex.Completable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Completable.class) + .fromPublisher(source); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Observable}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava2ObservableConverter implements Converter, io.reactivex.Observable> { + + INSTANCE; + + @Override + public io.reactivex.Observable convert(Publisher source) { + return (io.reactivex.Observable) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(io.reactivex.Single.class) + .fromPublisher(source); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2SingleToPublisherConverter implements Converter, Publisher> { + + INSTANCE; + + @Override + public Publisher convert(io.reactivex.Single source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2SingleToMonoConverter implements Converter, Mono> { + + INSTANCE; + + @Override + public Mono convert(io.reactivex.Single source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toMono(source); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Single} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2SingleToFluxConverter implements Converter, Flux> { + + INSTANCE; + + @Override + public Flux convert(io.reactivex.Single source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Single.class).toFlux(source); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2CompletableToPublisherConverter implements Converter> { + + INSTANCE; + + @Override + public Publisher convert(io.reactivex.Completable source) { + return REACTIVE_ADAPTER_REGISTRY.getAdapterFrom(io.reactivex.Completable.class).toFlux(source); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Completable} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2CompletableToMonoConverter implements Converter> { + + INSTANCE; + + @Override + public Mono convert(io.reactivex.Completable source) { + return Mono.from(RxJava2CompletableToPublisherConverter.INSTANCE.convert(source)); + } + } + + /** + * A {@link Converter} to convert an {@link io.reactivex.Observable} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2ObservableToPublisherConverter implements Converter, Publisher> { + + INSTANCE; + + @Override + public Publisher convert(io.reactivex.Observable source) { + return source.toFlowable(BackpressureStrategy.BUFFER); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2ObservableToMonoConverter implements Converter, Mono> { + + INSTANCE; + + @Override + public Mono convert(io.reactivex.Observable source) { + return Mono.from(source.toFlowable(BackpressureStrategy.BUFFER)); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Observable} to {@link Flux}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2ObservableToFluxConverter implements Converter, Flux> { + + INSTANCE; + + @Override + public Flux convert(io.reactivex.Observable source) { + return Flux.from(source.toFlowable(BackpressureStrategy.BUFFER)); + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava2FlowableConverter implements Converter, io.reactivex.Flowable> { + + INSTANCE; + + @Override + public io.reactivex.Flowable convert(Publisher source) { + return Flowable.fromPublisher(source); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Flowable} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2FlowableToPublisherConverter implements Converter, Publisher> { + + INSTANCE; + + @Override + public Publisher convert(io.reactivex.Flowable source) { + return source; + } + } + + /** + * A {@link Converter} to convert a {@link Publisher} to {@link io.reactivex.Flowable}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum PublisherToRxJava2MaybeConverter implements Converter, io.reactivex.Maybe> { + + INSTANCE; + + @Override + public io.reactivex.Maybe convert(Publisher source) { + return (io.reactivex.Maybe) REACTIVE_ADAPTER_REGISTRY.getAdapterTo(Maybe.class).fromPublisher(source); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Publisher}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2MaybeToPublisherConverter implements Converter, Publisher> { + + INSTANCE; + + @Override + public Publisher convert(io.reactivex.Maybe source) { + return source.toFlowable(); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Mono}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2MaybeToMonoConverter implements Converter, Mono> { + + INSTANCE; + + @Override + public Mono convert(io.reactivex.Maybe source) { + return Mono.from(source.toFlowable()); + } + } + + /** + * A {@link Converter} to convert a {@link io.reactivex.Maybe} to {@link Flux}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2MaybeToFluxConverter implements Converter, Flux> { + + INSTANCE; + + @Override + public Flux convert(io.reactivex.Maybe source) { + return Flux.from(source.toFlowable()); + } + } + + /** + * A {@link Converter} to convert a {@link Observable} to {@link Single}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1ObservableToSingleConverter implements Converter, Single> { + + INSTANCE; + + @Override + public Single convert(Observable source) { + return source.toSingle(); + } + } + + /** + * A {@link Converter} to convert a {@link Single} to {@link Single}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava1SingleToObservableConverter implements Converter, Observable> { + + INSTANCE; + + @Override + public Observable convert(Single source) { + return source.toObservable(); + } + } + + /** + * A {@link Converter} to convert a {@link Observable} to {@link Single}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2ObservableToSingleConverter + implements Converter, io.reactivex.Single> { + + INSTANCE; + + @Override + public io.reactivex.Single convert(io.reactivex.Observable source) { + return source.singleOrError(); + } + } + + /** + * A {@link Converter} to convert a {@link Observable} to {@link Maybe}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2ObservableToMaybeConverter + implements Converter, io.reactivex.Maybe> { + + INSTANCE; + + @Override + public io.reactivex.Maybe convert(io.reactivex.Observable source) { + return source.singleElement(); + } + } + + /** + * A {@link Converter} to convert a {@link Single} to {@link Single}. + * + * @author Mark Paluch + * @author 2.0 + */ + public enum RxJava2SingleToObservableConverter + implements Converter, io.reactivex.Observable> { + + INSTANCE; + + @Override + public io.reactivex.Observable convert(io.reactivex.Single source) { + return source.toObservable(); + } + } + + /** + * Holder for delayed initialization of {@link ReactiveAdapterRegistry}. + * + * @author Mark Paluch + * @author 2.0 + */ + static class RegistryHolder { + static final ReactiveAdapterRegistry REACTIVE_ADAPTER_REGISTRY = new ReactiveAdapterRegistry(); + } } diff --git a/src/test/java/org/springframework/data/repository/core/support/ReactiveRepositoryInformationUnitTests.java b/src/test/java/org/springframework/data/repository/core/support/ReactiveRepositoryInformationUnitTests.java index 1d4a0eeac..63fb8d0d0 100644 --- a/src/test/java/org/springframework/data/repository/core/support/ReactiveRepositoryInformationUnitTests.java +++ b/src/test/java/org/springframework/data/repository/core/support/ReactiveRepositoryInformationUnitTests.java @@ -18,8 +18,6 @@ package org.springframework.data.repository.core.support; import static org.hamcrest.Matchers.*; import static org.junit.Assert.*; -import rx.Observable; - import java.io.Serializable; import java.lang.reflect.Method; @@ -32,7 +30,9 @@ import org.springframework.data.repository.core.RepositoryMetadata; import org.springframework.data.repository.reactive.ReactiveCrudRepository; import org.springframework.data.repository.reactive.ReactiveSortingRepository; import org.springframework.data.repository.reactive.RxJava1CrudRepository; -import org.springframework.data.repository.util.QueryExecutionConverters; +import org.springframework.data.repository.util.ReactiveWrapperConverters; + +import rx.Observable; /** * Unit tests for {@link ReactiveRepositoryInformation}. @@ -60,7 +60,7 @@ public class ReactiveRepositoryInformationUnitTests { public void discoversMethodWithConvertibleArguments() throws Exception { DefaultConversionService conversionService = new DefaultConversionService(); - QueryExecutionConverters.registerConvertersIn(conversionService); + ReactiveWrapperConverters.registerConvertersIn(conversionService); Method method = RxJava1InterfaceWithGenerics.class.getMethod("save", Observable.class); RepositoryMetadata metadata = new DefaultRepositoryMetadata(RxJava1InterfaceWithGenerics.class); @@ -77,7 +77,7 @@ public class ReactiveRepositoryInformationUnitTests { public void discoversMethodAssignableArguments() throws Exception { DefaultConversionService conversionService = new DefaultConversionService(); - QueryExecutionConverters.registerConvertersIn(conversionService); + ReactiveWrapperConverters.registerConvertersIn(conversionService); Method method = ReactiveSortingRepository.class.getMethod("save", Publisher.class); RepositoryMetadata metadata = new DefaultRepositoryMetadata(ReactiveJavaInterfaceWithGenerics.class); @@ -94,7 +94,7 @@ public class ReactiveRepositoryInformationUnitTests { public void discoversMethodExactIterableArguments() throws Exception { DefaultConversionService conversionService = new DefaultConversionService(); - QueryExecutionConverters.registerConvertersIn(conversionService); + ReactiveWrapperConverters.registerConvertersIn(conversionService); Method method = ReactiveJavaInterfaceWithGenerics.class.getMethod("save", Iterable.class); RepositoryMetadata metadata = new DefaultRepositoryMetadata(ReactiveJavaInterfaceWithGenerics.class); @@ -111,7 +111,7 @@ public class ReactiveRepositoryInformationUnitTests { public void discoversMethodExactObjectArguments() throws Exception { DefaultConversionService conversionService = new DefaultConversionService(); - QueryExecutionConverters.registerConvertersIn(conversionService); + ReactiveWrapperConverters.registerConvertersIn(conversionService); Method method = ReactiveJavaInterfaceWithGenerics.class.getMethod("save", Object.class); RepositoryMetadata metadata = new DefaultRepositoryMetadata(ReactiveJavaInterfaceWithGenerics.class); @@ -124,8 +124,7 @@ public class ReactiveRepositoryInformationUnitTests { assertThat(reference.getParameterTypes()[0], is(equalTo(Object.class))); } - interface RxJava1InterfaceWithGenerics extends RxJava1CrudRepository - {} + interface RxJava1InterfaceWithGenerics extends RxJava1CrudRepository {} interface ReactiveJavaInterfaceWithGenerics extends ReactiveCrudRepository {} diff --git a/src/test/java/org/springframework/data/repository/core/support/ReactiveWrapperRepositoryFactorySupportUnitTests.java b/src/test/java/org/springframework/data/repository/core/support/ReactiveWrapperRepositoryFactorySupportUnitTests.java index 0dd0fad95..e53b6871b 100644 --- a/src/test/java/org/springframework/data/repository/core/support/ReactiveWrapperRepositoryFactorySupportUnitTests.java +++ b/src/test/java/org/springframework/data/repository/core/support/ReactiveWrapperRepositoryFactorySupportUnitTests.java @@ -15,10 +15,7 @@ */ package org.springframework.data.repository.core.support; -import static org.mockito.Mockito.any; -import static org.mockito.Mockito.times; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.verifyZeroInteractions; +import static org.mockito.Mockito.*; import java.io.Serializable; @@ -32,7 +29,7 @@ import org.mockito.runners.MockitoJUnitRunner; import org.springframework.core.convert.support.DefaultConversionService; import org.springframework.data.repository.Repository; import org.springframework.data.repository.reactive.ReactiveSortingRepository; -import org.springframework.data.repository.util.QueryExecutionConverters; +import org.springframework.data.repository.util.ReactiveWrapperConverters; import reactor.core.publisher.Mono; import rx.Single; @@ -56,7 +53,7 @@ public class ReactiveWrapperRepositoryFactorySupportUnitTests { public void setUp() { DefaultConversionService defaultConversionService = new DefaultConversionService(); - QueryExecutionConverters.registerConvertersIn(defaultConversionService); + ReactiveWrapperConverters.registerConvertersIn(defaultConversionService); factory = new DummyRepositoryFactory(backingRepo); factory.setConversionService(defaultConversionService);