From 40b33bca59b7c98eb1b3f839adfd350a9a6c463b Mon Sep 17 00:00:00 2001 From: Juergen Hoeller Date: Sun, 6 Aug 2023 14:04:24 +0200 Subject: [PATCH] Compatibility with Flow-based SmallRye Mutiny 2 at runtime Includes simple Flow.Publisher bridge without Reactor. Closes gh-31000 --- .../core/ReactiveAdapterRegistry.java | 223 ++++++++++++++++-- .../core/ReactiveAdapterRegistryTests.java | 12 + 2 files changed, 211 insertions(+), 24 deletions(-) diff --git a/spring-core/src/main/java/org/springframework/core/ReactiveAdapterRegistry.java b/spring-core/src/main/java/org/springframework/core/ReactiveAdapterRegistry.java index 50d8043294..f6a23e2722 100644 --- a/spring-core/src/main/java/org/springframework/core/ReactiveAdapterRegistry.java +++ b/spring-core/src/main/java/org/springframework/core/ReactiveAdapterRegistry.java @@ -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. * - *

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}. + *

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) source), + source -> new PublisherToFlow<>((Publisher) 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) 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) 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 implements Flow.Publisher { + + private static final Flow.Subscription EMPTY_SUBSCRIPTION = new Flow.Subscription() { + @Override + public void request(long n) { + } + @Override + public void cancel() { + } + }; + + @Nullable + private final Publisher publisher; + + public PublisherToFlow(@Nullable Publisher publisher) { + this.publisher = publisher; + } + + @Override + public void subscribe(Flow.Subscriber subscriber) { + if (this.publisher != null) { + this.publisher.subscribe(new SubscriberToFlow<>(subscriber)); + } + else { + subscriber.onSubscribe(EMPTY_SUBSCRIPTION); + subscriber.onComplete(); + } + } + } + + + private static class PublisherToRS implements Publisher { + + private static final Flow.Publisher EMPTY_FLOW = new PublisherToFlow<>(null); + + private final Flow.Publisher publisher; + + @SuppressWarnings("unchecked") + public PublisherToRS(@Nullable Flow.Publisher publisher) { + this.publisher = (publisher != null ? publisher : (Flow.Publisher) EMPTY_FLOW); + } + + @Override + public void subscribe(Subscriber subscriber) { + this.publisher.subscribe(new SubscriberToRS<>(subscriber)); + } + } + + + private static class SubscriberToFlow implements Subscriber, Flow.Subscription { + + private final Flow.Subscriber subscriber; + + @Nullable + private Subscription subscription; + + public SubscriberToFlow(Flow.Subscriber 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 implements Flow.Subscriber, Subscription { + + private final Subscriber subscriber; + + @Nullable + private Flow.Subscription subscription; + + public SubscriberToRS(Subscriber 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(); + } } } diff --git a/spring-core/src/test/java/org/springframework/core/ReactiveAdapterRegistryTests.java b/spring-core/src/test/java/org/springframework/core/ReactiveAdapterRegistryTests.java index 3a7a5ab249..a969a482f6 100644 --- a/spring-core/src/test/java/org/springframework/core/ReactiveAdapterRegistryTests.java +++ b/spring-core/src/test/java/org/springframework/core/ReactiveAdapterRegistryTests.java @@ -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) target).block(ONE_SECOND)).isEqualTo(Integer.valueOf(1)); } + @Test + void toFlowPublisher() { + List sequence = Arrays.asList(1, 2, 3); + Publisher 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) target) + .collectList().block(ONE_SECOND)).isEqualTo(sequence); + } + @Test void toCompletableFuture() throws Exception { Publisher source = Flux.fromArray(new Integer[] {1, 2, 3});