Compatibility with Flow-based SmallRye Mutiny 2 at runtime
Includes simple Flow.Publisher bridge without Reactor. Closes gh-31000
This commit is contained in:
@@ -16,6 +16,7 @@
|
||||
|
||||
package org.springframework.core;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
@@ -24,9 +25,9 @@ import java.util.concurrent.CompletionStage;
|
||||
import java.util.concurrent.Flow;
|
||||
import java.util.function.Function;
|
||||
|
||||
import kotlinx.coroutines.CompletableDeferredKt;
|
||||
import kotlinx.coroutines.Deferred;
|
||||
import org.reactivestreams.Publisher;
|
||||
import org.reactivestreams.Subscriber;
|
||||
import org.reactivestreams.Subscription;
|
||||
import reactor.adapter.JdkFlowAdapter;
|
||||
import reactor.blockhound.BlockHound;
|
||||
import reactor.blockhound.integration.BlockHoundIntegration;
|
||||
@@ -36,15 +37,18 @@ import reactor.core.publisher.Mono;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.util.ClassUtils;
|
||||
import org.springframework.util.ConcurrentReferenceHashMap;
|
||||
import org.springframework.util.ReflectionUtils;
|
||||
|
||||
/**
|
||||
* A registry of adapters to adapt Reactive Streams {@link Publisher} to/from
|
||||
* various async/reactive types such as {@code CompletableFuture}, RxJava
|
||||
* {@code Flowable}, and others.
|
||||
* A registry of adapters to adapt Reactive Streams {@link Publisher} to/from various
|
||||
* async/reactive types such as {@code CompletableFuture}, RxJava {@code Flowable}, etc.
|
||||
* This is designed to complement Spring's Reactor {@code Mono}/{@code Flux} support while
|
||||
* also being usable without Reactor, e.g. just for {@code org.reactivestreams} bridging.
|
||||
*
|
||||
* <p>By default, depending on classpath availability, adapters are registered
|
||||
* for Reactor, RxJava 3, {@link CompletableFuture}, {@code Flow.Publisher},
|
||||
* and Kotlin Coroutines' {@code Deferred} and {@code Flow}.
|
||||
* <p>By default, depending on classpath availability, adapters are registered for Reactor
|
||||
* (including {@code CompletableFuture} and {@code Flow.Publisher} adapters), RxJava 3,
|
||||
* Kotlin Coroutines' {@code Deferred} (bridged via Reactor) and SmallRye Mutiny 1.x/2.x.
|
||||
* If Reactor is not present, a simple {@code Flow.Publisher} bridge will be registered.
|
||||
*
|
||||
* @author Rossen Stoyanchev
|
||||
* @author Sebastien Deleuze
|
||||
@@ -79,6 +83,7 @@ public class ReactiveAdapterRegistry {
|
||||
* Create a registry and auto-register default adapters.
|
||||
* @see #getSharedInstance()
|
||||
*/
|
||||
@SuppressWarnings("unchecked")
|
||||
public ReactiveAdapterRegistry() {
|
||||
// Reactor
|
||||
if (reactorPresent) {
|
||||
@@ -99,6 +104,14 @@ public class ReactiveAdapterRegistry {
|
||||
if (mutinyPresent) {
|
||||
new MutinyRegistrar().registerAdapters(this);
|
||||
}
|
||||
|
||||
// Simple Flow.Publisher bridge if Reactor is not present
|
||||
if (!reactorPresent) {
|
||||
registerReactiveType(
|
||||
ReactiveTypeDescriptor.multiValue(Flow.Publisher.class, () -> PublisherToRS.EMPTY_FLOW),
|
||||
source -> new PublisherToRS<>((Flow.Publisher<Object>) source),
|
||||
source -> new PublisherToFlow<>((Publisher<Object>) source));
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -304,9 +317,9 @@ public class ReactiveAdapterRegistry {
|
||||
@SuppressWarnings("KotlinInternalInJava")
|
||||
void registerAdapters(ReactiveAdapterRegistry registry) {
|
||||
registry.registerReactiveType(
|
||||
ReactiveTypeDescriptor.singleOptionalValue(Deferred.class,
|
||||
() -> CompletableDeferredKt.CompletableDeferred(null)),
|
||||
source -> CoroutinesUtils.deferredToMono((Deferred<?>) source),
|
||||
ReactiveTypeDescriptor.singleOptionalValue(kotlinx.coroutines.Deferred.class,
|
||||
() -> kotlinx.coroutines.CompletableDeferredKt.CompletableDeferred(null)),
|
||||
source -> CoroutinesUtils.deferredToMono((kotlinx.coroutines.Deferred<?>) source),
|
||||
source -> CoroutinesUtils.monoToDeferred(Mono.from(source)));
|
||||
|
||||
registry.registerReactiveType(
|
||||
@@ -319,20 +332,182 @@ public class ReactiveAdapterRegistry {
|
||||
|
||||
private static class MutinyRegistrar {
|
||||
|
||||
void registerAdapters(ReactiveAdapterRegistry registry) {
|
||||
registry.registerReactiveType(
|
||||
ReactiveTypeDescriptor.singleOptionalValue(
|
||||
io.smallrye.mutiny.Uni.class,
|
||||
() -> io.smallrye.mutiny.Uni.createFrom().nothing()),
|
||||
uni -> ((io.smallrye.mutiny.Uni<?>) uni).convert().toPublisher(),
|
||||
publisher -> io.smallrye.mutiny.Uni.createFrom().publisher(publisher));
|
||||
private static final Method uniToPublisher = ClassUtils.getMethod(io.smallrye.mutiny.groups.UniConvert.class, "toPublisher");
|
||||
|
||||
registry.registerReactiveType(
|
||||
ReactiveTypeDescriptor.multiValue(
|
||||
io.smallrye.mutiny.Multi.class,
|
||||
() -> io.smallrye.mutiny.Multi.createFrom().empty()),
|
||||
multi -> (io.smallrye.mutiny.Multi<?>) multi,
|
||||
publisher -> io.smallrye.mutiny.Multi.createFrom().publisher(publisher));
|
||||
@SuppressWarnings("unchecked")
|
||||
void registerAdapters(ReactiveAdapterRegistry registry) {
|
||||
ReactiveTypeDescriptor uniDesc = ReactiveTypeDescriptor.singleOptionalValue(
|
||||
io.smallrye.mutiny.Uni.class,
|
||||
() -> io.smallrye.mutiny.Uni.createFrom().nothing());
|
||||
ReactiveTypeDescriptor multiDesc = ReactiveTypeDescriptor.multiValue(
|
||||
io.smallrye.mutiny.Multi.class,
|
||||
() -> io.smallrye.mutiny.Multi.createFrom().empty());
|
||||
|
||||
if (Flow.Publisher.class.isAssignableFrom(uniToPublisher.getReturnType())) {
|
||||
// Mutiny 2 based on Flow.Publisher
|
||||
Method uniPublisher = ClassUtils.getMethod(io.smallrye.mutiny.groups.UniCreate.class, "publisher", Flow.Publisher.class);
|
||||
Method multiPublisher = ClassUtils.getMethod(io.smallrye.mutiny.groups.MultiCreate.class, "publisher", Flow.Publisher.class);
|
||||
registry.registerReactiveType(uniDesc,
|
||||
uni -> new PublisherToRS<>((Flow.Publisher<Object>) ReflectionUtils.invokeMethod(uniToPublisher, ((io.smallrye.mutiny.Uni<?>) uni).convert())),
|
||||
publisher -> ReflectionUtils.invokeMethod(uniPublisher, io.smallrye.mutiny.Uni.createFrom(), new PublisherToFlow<>(publisher)));
|
||||
registry.registerReactiveType(multiDesc,
|
||||
multi -> new PublisherToRS<>((Flow.Publisher<Object>) multi),
|
||||
publisher -> ReflectionUtils.invokeMethod(multiPublisher, io.smallrye.mutiny.Multi.createFrom(), new PublisherToFlow<>(publisher)));
|
||||
}
|
||||
else {
|
||||
// Mutiny 1 based on Reactive Streams
|
||||
registry.registerReactiveType(uniDesc,
|
||||
uni -> ((io.smallrye.mutiny.Uni<?>) uni).convert().toPublisher(),
|
||||
publisher -> io.smallrye.mutiny.Uni.createFrom().publisher(publisher));
|
||||
registry.registerReactiveType(multiDesc,
|
||||
multi -> (io.smallrye.mutiny.Multi<?>) multi,
|
||||
publisher -> io.smallrye.mutiny.Multi.createFrom().publisher(publisher));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private static class PublisherToFlow<T> implements Flow.Publisher<T> {
|
||||
|
||||
private static final Flow.Subscription EMPTY_SUBSCRIPTION = new Flow.Subscription() {
|
||||
@Override
|
||||
public void request(long n) {
|
||||
}
|
||||
@Override
|
||||
public void cancel() {
|
||||
}
|
||||
};
|
||||
|
||||
@Nullable
|
||||
private final Publisher<T> publisher;
|
||||
|
||||
public PublisherToFlow(@Nullable Publisher<T> publisher) {
|
||||
this.publisher = publisher;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void subscribe(Flow.Subscriber<? super T> subscriber) {
|
||||
if (this.publisher != null) {
|
||||
this.publisher.subscribe(new SubscriberToFlow<>(subscriber));
|
||||
}
|
||||
else {
|
||||
subscriber.onSubscribe(EMPTY_SUBSCRIPTION);
|
||||
subscriber.onComplete();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private static class PublisherToRS<T> implements Publisher<T> {
|
||||
|
||||
private static final Flow.Publisher<Object> EMPTY_FLOW = new PublisherToFlow<>(null);
|
||||
|
||||
private final Flow.Publisher<T> publisher;
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public PublisherToRS(@Nullable Flow.Publisher<T> publisher) {
|
||||
this.publisher = (publisher != null ? publisher : (Flow.Publisher<T>) EMPTY_FLOW);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void subscribe(Subscriber<? super T> subscriber) {
|
||||
this.publisher.subscribe(new SubscriberToRS<>(subscriber));
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private static class SubscriberToFlow<T> implements Subscriber<T>, Flow.Subscription {
|
||||
|
||||
private final Flow.Subscriber<? super T> subscriber;
|
||||
|
||||
@Nullable
|
||||
private Subscription subscription;
|
||||
|
||||
public SubscriberToFlow(Flow.Subscriber<? super T> subscriber) {
|
||||
this.subscriber = subscriber;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onSubscribe(Subscription subscription) {
|
||||
this.subscription = subscription;
|
||||
this.subscriber.onSubscribe(this);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onNext(T o) {
|
||||
this.subscriber.onNext(o);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onError(Throwable t) {
|
||||
this.subscriber.onError(t);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onComplete() {
|
||||
this.subscriber.onComplete();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void request(long n) {
|
||||
if (this.subscription != null) {
|
||||
this.subscription.request(n);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void cancel() {
|
||||
if (this.subscription != null) {
|
||||
this.subscription.cancel();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private static class SubscriberToRS<T> implements Flow.Subscriber<T>, Subscription {
|
||||
|
||||
private final Subscriber<? super T> subscriber;
|
||||
|
||||
@Nullable
|
||||
private Flow.Subscription subscription;
|
||||
|
||||
public SubscriberToRS(Subscriber<? super T> subscriber) {
|
||||
this.subscriber = subscriber;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onSubscribe(Flow.Subscription subscription) {
|
||||
this.subscription = subscription;
|
||||
this.subscriber.onSubscribe(this);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onNext(T o) {
|
||||
this.subscriber.onNext(o);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onError(Throwable throwable) {
|
||||
this.subscriber.onError(throwable);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onComplete() {
|
||||
this.subscriber.onComplete();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void request(long n) {
|
||||
if (this.subscription != null) {
|
||||
this.subscription.request(n);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void cancel() {
|
||||
if (this.subscription != null) {
|
||||
this.subscription.cancel();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -20,6 +20,7 @@ import java.time.Duration;
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Flow;
|
||||
|
||||
import io.smallrye.mutiny.Multi;
|
||||
import io.smallrye.mutiny.Uni;
|
||||
@@ -27,6 +28,7 @@ import kotlinx.coroutines.Deferred;
|
||||
import org.junit.jupiter.api.Nested;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.reactivestreams.Publisher;
|
||||
import reactor.adapter.JdkFlowAdapter;
|
||||
import reactor.core.CoreSubscriber;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
@@ -112,6 +114,16 @@ class ReactiveAdapterRegistryTests {
|
||||
assertThat(((Mono<Integer>) target).block(ONE_SECOND)).isEqualTo(Integer.valueOf(1));
|
||||
}
|
||||
|
||||
@Test
|
||||
void toFlowPublisher() {
|
||||
List<Integer> sequence = Arrays.asList(1, 2, 3);
|
||||
Publisher<Integer> source = io.reactivex.rxjava3.core.Flowable.fromIterable(sequence);
|
||||
Object target = getAdapter(Flow.Publisher.class).fromPublisher(source);
|
||||
assertThat(target).isInstanceOf(Flow.Publisher.class);
|
||||
assertThat(JdkFlowAdapter.flowPublisherToFlux((Flow.Publisher<Integer>) target)
|
||||
.collectList().block(ONE_SECOND)).isEqualTo(sequence);
|
||||
}
|
||||
|
||||
@Test
|
||||
void toCompletableFuture() throws Exception {
|
||||
Publisher<Integer> source = Flux.fromArray(new Integer[] {1, 2, 3});
|
||||
|
||||
Reference in New Issue
Block a user