From 00d6a364b2b0bfa53c8be6c2b410800cb2f3bff8 Mon Sep 17 00:00:00 2001 From: spencergibb Date: Fri, 22 Jan 2021 11:27:03 -0500 Subject: [PATCH] Moves CloudFlux to spring cloud package. FluxFirstNonEmptyEmitting and Tests used package private classes/interfaces from reactor. These were replaced with the public interfaces they were extending and any default methods copied to the implementation. LambdaSubscriber was copied to the test where it was used and had similar treatment for its package protected interfaces. Fixes gh-887 --- .../ReactiveCompositeDiscoveryClient.java | 2 +- .../cloud/commons}/publisher/CloudFlux.java | 7 ++-- .../publisher/FluxFirstNonEmptyEmitting.java | 37 ++++++++++++------- .../FluxFirstNonEmptyEmittingTests.java | 28 +++++++++----- 4 files changed, 46 insertions(+), 28 deletions(-) rename spring-cloud-commons/src/main/java/{reactor/core => org/springframework/cloud/commons}/publisher/CloudFlux.java (93%) rename spring-cloud-commons/src/main/java/{reactor/core => org/springframework/cloud/commons}/publisher/FluxFirstNonEmptyEmitting.java (90%) rename spring-cloud-commons/src/test/java/{reactor/core => org/springframework/cloud/commons}/publisher/FluxFirstNonEmptyEmittingTests.java (91%) diff --git a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/discovery/composite/reactive/ReactiveCompositeDiscoveryClient.java b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/discovery/composite/reactive/ReactiveCompositeDiscoveryClient.java index 1e366140..e7f8f982 100644 --- a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/discovery/composite/reactive/ReactiveCompositeDiscoveryClient.java +++ b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/discovery/composite/reactive/ReactiveCompositeDiscoveryClient.java @@ -19,11 +19,11 @@ package org.springframework.cloud.client.discovery.composite.reactive; import java.util.ArrayList; import java.util.List; -import reactor.core.publisher.CloudFlux; import reactor.core.publisher.Flux; import org.springframework.cloud.client.ServiceInstance; import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; +import org.springframework.cloud.commons.publisher.CloudFlux; import org.springframework.core.annotation.AnnotationAwareOrderComparator; /** diff --git a/spring-cloud-commons/src/main/java/reactor/core/publisher/CloudFlux.java b/spring-cloud-commons/src/main/java/org/springframework/cloud/commons/publisher/CloudFlux.java similarity index 93% rename from spring-cloud-commons/src/main/java/reactor/core/publisher/CloudFlux.java rename to spring-cloud-commons/src/main/java/org/springframework/cloud/commons/publisher/CloudFlux.java index c9f4b39a..87922b72 100644 --- a/spring-cloud-commons/src/main/java/reactor/core/publisher/CloudFlux.java +++ b/spring-cloud-commons/src/main/java/org/springframework/cloud/commons/publisher/CloudFlux.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2020 the original author or authors. + * Copyright 2013-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -14,13 +14,14 @@ * limitations under the License. */ -package reactor.core.publisher; +package org.springframework.cloud.commons.publisher; import org.reactivestreams.Publisher; +import reactor.core.publisher.Flux; /** * INTERNAL USAGE ONLY. This functionality will be ported to reactor-core and will be - * removed in a next release. + * removed in a future release. * * @author Tim Ysewyn */ diff --git a/spring-cloud-commons/src/main/java/reactor/core/publisher/FluxFirstNonEmptyEmitting.java b/spring-cloud-commons/src/main/java/org/springframework/cloud/commons/publisher/FluxFirstNonEmptyEmitting.java similarity index 90% rename from spring-cloud-commons/src/main/java/reactor/core/publisher/FluxFirstNonEmptyEmitting.java rename to spring-cloud-commons/src/main/java/org/springframework/cloud/commons/publisher/FluxFirstNonEmptyEmitting.java index d94caac1..c798c7f3 100644 --- a/spring-cloud-commons/src/main/java/reactor/core/publisher/FluxFirstNonEmptyEmitting.java +++ b/spring-cloud-commons/src/main/java/org/springframework/cloud/commons/publisher/FluxFirstNonEmptyEmitting.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2020 the original author or authors. + * Copyright 2013-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -14,8 +14,9 @@ * limitations under the License. */ -package reactor.core.publisher; +package org.springframework.cloud.commons.publisher; +import java.lang.reflect.Field; import java.util.Iterator; import java.util.Objects; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; @@ -25,12 +26,16 @@ import org.reactivestreams.Publisher; import org.reactivestreams.Subscription; import reactor.core.CoreSubscriber; import reactor.core.Scannable; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Operators; import reactor.util.annotation.Nullable; +import org.springframework.util.ReflectionUtils; + /** * @author Tim Ysewyn */ -final class FluxFirstNonEmptyEmitting extends Flux implements SourceProducer { +final class FluxFirstNonEmptyEmitting extends Flux implements Scannable, Publisher { final Publisher[] array; @@ -137,6 +142,11 @@ final class FluxFirstNonEmptyEmitting extends Flux implements SourceProduc return null; // no particular key to be represented, still useful in hooks } + @Override + public String stepName() { + return "source(" + getClass().getSimpleName() + ")"; + } + static final class RaceCoordinator implements Subscription, Scannable { final FirstNonEmptyEmittingSubscriber[] subscribers; @@ -264,8 +274,8 @@ final class FluxFirstNonEmptyEmitting extends Flux implements SourceProduc } - static final class FirstNonEmptyEmittingSubscriber - extends Operators.DeferredSubscription implements InnerOperator { + static final class FirstNonEmptyEmittingSubscriber extends Operators.DeferredSubscription + implements CoreSubscriber, Scannable, Subscription { final RaceCoordinator parent; @@ -285,14 +295,13 @@ final class FluxFirstNonEmptyEmitting extends Flux implements SourceProduc @Override @Nullable public Object scanUnsafe(Attr key) { - if (key == Attr.PARENT) { - return s; + if (key == Attr.ACTUAL) { + return actual; } if (key == Attr.CANCELLED) { return parent.cancelled; } - - return InnerOperator.super.scanUnsafe(key); + return super.scanUnsafe(key); } @Override @@ -300,11 +309,6 @@ final class FluxFirstNonEmptyEmitting extends Flux implements SourceProduc set(s); } - @Override - public CoreSubscriber actual() { - return actual; - } - @Override public void onNext(T t) { if (won) { @@ -334,6 +338,11 @@ final class FluxFirstNonEmptyEmitting extends Flux implements SourceProduc } } + @Override + public String stepName() { + return "CloudFlux.firstNonEmpty"; + } + } } diff --git a/spring-cloud-commons/src/test/java/reactor/core/publisher/FluxFirstNonEmptyEmittingTests.java b/spring-cloud-commons/src/test/java/org/springframework/cloud/commons/publisher/FluxFirstNonEmptyEmittingTests.java similarity index 91% rename from spring-cloud-commons/src/test/java/reactor/core/publisher/FluxFirstNonEmptyEmittingTests.java rename to spring-cloud-commons/src/test/java/org/springframework/cloud/commons/publisher/FluxFirstNonEmptyEmittingTests.java index 2a32bc5e..4ac24655 100644 --- a/spring-cloud-commons/src/test/java/reactor/core/publisher/FluxFirstNonEmptyEmittingTests.java +++ b/spring-cloud-commons/src/test/java/org/springframework/cloud/commons/publisher/FluxFirstNonEmptyEmittingTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019-2020 the original author or authors. + * Copyright 2013-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -14,7 +14,7 @@ * limitations under the License. */ -package reactor.core.publisher; +package org.springframework.cloud.commons.publisher; import java.time.Duration; import java.util.Arrays; @@ -24,6 +24,9 @@ import org.reactivestreams.Publisher; import org.reactivestreams.Subscription; import reactor.core.CoreSubscriber; import reactor.core.Scannable; +import reactor.core.publisher.BaseSubscriber; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Operators; import reactor.test.StepVerifier; import static java.util.Collections.singletonList; @@ -156,10 +159,8 @@ public class FluxFirstNonEmptyEmittingTests { @Test public void scanSubscriber() { - CoreSubscriber actual = new LambdaSubscriber<>(null, e -> { - }, null, null); - FluxFirstNonEmptyEmitting.RaceCoordinator parent = new FluxFirstNonEmptyEmitting.RaceCoordinator<>( - 1); + CoreSubscriber actual = new TestSubscriber<>(); + FluxFirstNonEmptyEmitting.RaceCoordinator parent = new FluxFirstNonEmptyEmitting.RaceCoordinator<>(1); FluxFirstNonEmptyEmitting.FirstNonEmptyEmittingSubscriber test = new FluxFirstNonEmptyEmitting.FirstNonEmptyEmittingSubscriber<>( actual, parent, 1); Subscription sub = Operators.emptySubscription(); @@ -174,10 +175,8 @@ public class FluxFirstNonEmptyEmittingTests { @Test public void scanRaceCoordinator() { - CoreSubscriber actual = new LambdaSubscriber<>(null, e -> { - }, null, null); - FluxFirstNonEmptyEmitting.RaceCoordinator parent = new FluxFirstNonEmptyEmitting.RaceCoordinator<>( - 1); + CoreSubscriber actual = new TestSubscriber<>(); + FluxFirstNonEmptyEmitting.RaceCoordinator parent = new FluxFirstNonEmptyEmitting.RaceCoordinator<>(1); FluxFirstNonEmptyEmitting.FirstNonEmptyEmittingSubscriber test = new FluxFirstNonEmptyEmitting.FirstNonEmptyEmittingSubscriber<>( actual, parent, 1); Subscription sub = Operators.emptySubscription(); @@ -190,4 +189,13 @@ public class FluxFirstNonEmptyEmittingTests { assertThat(parent.scan(Scannable.Attr.CANCELLED)).isTrue(); } + static class TestSubscriber extends BaseSubscriber implements Scannable { + + @Override + public Object scanUnsafe(Attr key) { + return null; + } + + } + }