From f6a0e38c7195b9bd9c418195d0c5e22db4e46de6 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Wed, 28 Feb 2018 17:23:52 +0100 Subject: [PATCH] Introduced lazy span operators and HookRegisteringBeanDefinitionRegistryPostProcessor without this change the Reactor Span operators where eagerly started. That means that beans like Tracing were required for an operation. Also, Sleuth requires to be the first to set a traceable scheduler factory. Unfortunately it wasn't always the case. with this change the Reactor Span operators are lazy. They create span operators at runtime. Thanks to this we could create a HookRegisteringBeanDefinitionRegistryPostProcessor that ensures that the hooks are applied and the factory is set before other components get initialized. fixes gh-866 --- .../reactor/LazySpanSubscriber.java | 67 ++++++++++++++++ .../instrument/reactor/ReactorSleuth.java | 47 +++++++---- .../reactor/ScopePassingSpanSubscriber.java | 17 ++-- .../instrument/reactor/SpanSubscriber.java | 4 +- .../instrument/reactor/SpanSubscription.java | 26 ++++++ .../reactor/SpanSubscriptionProvider.java | 55 +++++++++++++ .../TraceReactorAutoConfiguration.java | 80 ++++++++++++------- .../opentracing/BraveTracerTest.java | 10 --- .../reactor/Issue866Configuration.java | 53 ++++++++++++ .../reactor/sample/FlatMapTests.java | 34 +++++++- 10 files changed, 322 insertions(+), 71 deletions(-) create mode 100644 spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/LazySpanSubscriber.java create mode 100644 spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscription.java create mode 100644 spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProvider.java create mode 100644 spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/Issue866Configuration.java diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/LazySpanSubscriber.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/LazySpanSubscriber.java new file mode 100644 index 000000000..ae58a933f --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/LazySpanSubscriber.java @@ -0,0 +1,67 @@ +/* + * Copyright 2013-2018 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.reactor; + +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Supplier; + +import org.reactivestreams.Subscription; +import reactor.util.context.Context; + +/** + * A lazy representation of the {@link SpanSubscription} + * + * @author Marcin Grzejszczak + * @since 2.0.0 + */ +final class LazySpanSubscriber extends AtomicBoolean implements SpanSubscription { + + private final Supplier> supplier; + + LazySpanSubscriber(Supplier> supplier) { + this.supplier = supplier; + } + + @Override public void onSubscribe(Subscription subscription) { + this.supplier.get().onSubscribe(subscription); + } + + @Override public void request(long n) { + this.supplier.get().request(n); + } + + @Override public void cancel() { + this.supplier.get().cancel(); + } + + @Override public void onNext(T o) { + this.supplier.get().onNext(o); + } + + @Override public void onError(Throwable throwable) { + this.supplier.get().onError(throwable); + } + + @Override public void onComplete() { + this.supplier.get().onComplete(); + } + + @Override public Context currentContext() { + return this.supplier.get().currentContext(); + } +} + diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java index f5cf7c65b..0bd1792f8 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java @@ -20,10 +20,11 @@ import java.util.function.Function; import java.util.function.Predicate; import brave.Tracing; +import org.reactivestreams.Publisher; +import org.springframework.beans.factory.BeanFactory; import reactor.core.Fuseable; import reactor.core.Scannable; import reactor.core.publisher.Operators; -import org.reactivestreams.Publisher; import reactor.util.context.Context; /** @@ -35,29 +36,31 @@ import reactor.util.context.Context; public abstract class ReactorSleuth { /** - * Return a span operator pointcut given a {@link Tracing}. This can be used in reactor + * Return a span operator pointcut given a {@link BeanFactory}. This can be used in reactor * via {@link reactor.core.publisher.Flux#transform(Function)}, {@link * reactor.core.publisher.Mono#transform(Function)}, {@link * reactor.core.publisher.Hooks#onEachOperator(Function)} or {@link * reactor.core.publisher.Hooks#onLastOperator(Function)}. * - * @param tracing the {@link Tracing} instance to use in this span operator + * @param beanFactory * @param an arbitrary type that is left unchanged by the span operator * - * @return a new Span operator pointcut + * @return a new lazy span operator pointcut */ public static Function, ? extends Publisher> spanOperator( - Tracing tracing) { + BeanFactory beanFactory) { return Operators.lift(POINTCUT_FILTER, ((scannable, sub) -> { //do not trace fused flows if(scannable instanceof Fuseable && sub instanceof Fuseable.QueueSubscription){ return sub; } - return new SpanSubscriber<>( - sub, - sub.currentContext(), - tracing, - scannable.name()); + return new LazySpanSubscriber( + new SpanSubscriptionProvider( + beanFactory, + sub, + sub.currentContext(), + scannable.name()) + ); })); } @@ -69,22 +72,32 @@ public abstract class ReactorSleuth { * reactor.core.publisher.Hooks#onLastOperator(Function)}. The Span operator * pointcut will pass the Scope of the Span without ever creating any new spans. * - * @param tracing the {@link Tracing} instance to use in this span operator + * @param beanFactory * @param an arbitrary type that is left unchanged by the span operator * - * @return a new Span operator pointcut + * @return a new lazy span operator pointcut */ public static Function, ? extends Publisher> scopePassingSpanOperator( - Tracing tracing) { + BeanFactory beanFactory) { return Operators.lift(POINTCUT_FILTER, ((scannable, sub) -> { //do not trace fused flows if(scannable instanceof Fuseable && sub instanceof Fuseable.QueueSubscription){ return sub; } - return new ScopePassingSpanSubscriber<>( - sub, - sub != null ? sub.currentContext() : Context.empty(), - tracing); + return new LazySpanSubscriber( + new SpanSubscriptionProvider( + beanFactory, + sub, + sub.currentContext(), + scannable.name()) { + @Override SpanSubscription newCoreSubscriber(Tracing tracing) { + return new ScopePassingSpanSubscriber( + sub, + sub != null ? sub.currentContext() : Context.empty(), + tracing); + } + } + ); })); } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriber.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriber.java index 9415479b8..ac56b67b0 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriber.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriber.java @@ -23,9 +23,6 @@ import brave.Tracer; import brave.Tracing; import org.reactivestreams.Subscriber; import org.reactivestreams.Subscription; -import reactor.core.CoreSubscriber; -import reactor.util.Logger; -import reactor.util.Loggers; import reactor.util.context.Context; /** @@ -35,11 +32,7 @@ import reactor.util.context.Context; * @author Marcin Grzejszczak * @since 2.0.0 */ -final class ScopePassingSpanSubscriber extends AtomicBoolean implements Subscription, - CoreSubscriber { - - private static final Logger log = Loggers.getLogger( - ScopePassingSpanSubscriber.class); +final class ScopePassingSpanSubscriber extends AtomicBoolean implements SpanSubscription { private final Span span; private final Subscriber subscriber; @@ -83,11 +76,15 @@ final class ScopePassingSpanSubscriber extends AtomicBoolean implements Subsc } @Override public void onError(Throwable throwable) { - this.subscriber.onError(throwable); + try (Tracer.SpanInScope inScope = this.tracer.withSpanInScope(this.span)) { + this.subscriber.onError(throwable); + } } @Override public void onComplete() { - this.subscriber.onComplete(); + try (Tracer.SpanInScope inScope = this.tracer.withSpanInScope(this.span)) { + this.subscriber.onComplete(); + } } @Override public Context currentContext() { diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriber.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriber.java index e64b62b2a..26ff2ea33 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriber.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriber.java @@ -24,7 +24,6 @@ import brave.Tracing; import brave.propagation.TraceContextOrSamplingFlags; import org.reactivestreams.Subscriber; import org.reactivestreams.Subscription; -import reactor.core.CoreSubscriber; import reactor.util.Logger; import reactor.util.Loggers; import reactor.util.context.Context; @@ -36,8 +35,7 @@ import reactor.util.context.Context; * @author Marcin Grzejszczak * @since 2.0.0 */ -final class SpanSubscriber extends AtomicBoolean implements Subscription, - CoreSubscriber { +final class SpanSubscriber extends AtomicBoolean implements SpanSubscription { private static final Logger log = Loggers.getLogger( SpanSubscriber.class); diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscription.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscription.java new file mode 100644 index 000000000..c1992e39e --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscription.java @@ -0,0 +1,26 @@ +/* + * Copyright 2013-2018 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.reactor; + +import org.reactivestreams.Subscription; +import reactor.core.CoreSubscriber; + +/** + * @author Marcin Grzejszczak + */ +interface SpanSubscription extends Subscription, CoreSubscriber { +} diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProvider.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProvider.java new file mode 100644 index 000000000..16971012d --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProvider.java @@ -0,0 +1,55 @@ +/* + * Copyright 2013-2018 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.reactor; + +import java.util.function.Supplier; + +import brave.Tracing; +import org.reactivestreams.Subscriber; +import org.springframework.beans.factory.BeanFactory; +import reactor.util.context.Context; + +/** + * Supplier to lazily start a {@link SpanSubscription} + * + * @author Marcin Grzejszczak + */ +class SpanSubscriptionProvider implements Supplier> { + + final BeanFactory beanFactory; + final Subscriber subscriber; + final Context context; + final String name; + + SpanSubscriptionProvider(BeanFactory beanFactory, + Subscriber subscriber, + Context context, String name) { + this.beanFactory = beanFactory; + this.subscriber = subscriber; + this.context = context; + this.name = name; + } + + @Override public SpanSubscription get() { + Tracing tracing = this.beanFactory.getBean(Tracing.class); + return newCoreSubscriber(tracing); + } + + SpanSubscription newCoreSubscriber(Tracing tracing) { + return new SpanSubscriber<>(this.subscriber, this.context, tracing, this.name); + } +} diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java index 7d6194658..6cc93d508 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java @@ -16,20 +16,20 @@ package org.springframework.cloud.sleuth.instrument.reactor; -import javax.annotation.PostConstruct; -import javax.annotation.PreDestroy; import java.util.concurrent.ScheduledExecutorService; import java.util.function.Supplier; +import javax.annotation.PreDestroy; import brave.Tracing; -import reactor.core.publisher.Hooks; -import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; +import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; -import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.beans.factory.support.BeanDefinitionRegistryPostProcessor; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnNotWebApplication; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.condition.ConditionalOnWebApplication; @@ -37,6 +37,9 @@ import org.springframework.cloud.sleuth.instrument.async.TraceableScheduledExecu import org.springframework.cloud.sleuth.instrument.web.TraceWebFluxAutoConfiguration; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import reactor.core.publisher.Hooks; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; /** * {@link org.springframework.boot.autoconfigure.EnableAutoConfiguration Auto-configuration} @@ -56,34 +59,16 @@ public class TraceReactorAutoConfiguration { @ConditionalOnBean(Tracing.class) static class TraceReactorConfiguration { - private static final String SLEUTH_TRACE_REACTOR_KEY = TraceReactorConfiguration.class.getName(); - - @Autowired Tracing tracing; - @Autowired BeanFactory beanFactory; - @Autowired LastOperatorWrapper lastOperatorWrapper; + static final String SLEUTH_TRACE_REACTOR_KEY = TraceReactorConfiguration.class.getName(); @Bean @ConditionalOnNotWebApplication LastOperatorWrapper spanOperator() { - return tracer -> Hooks.onLastOperator(SLEUTH_TRACE_REACTOR_KEY, ReactorSleuth.spanOperator(tracer)); + return beanFactory -> Hooks.onLastOperator(SLEUTH_TRACE_REACTOR_KEY, ReactorSleuth.spanOperator(beanFactory)); } @Bean @ConditionalOnWebApplication LastOperatorWrapper noOpLastOperatorWrapper() { - return tracer -> { }; - } - - @PostConstruct - public void setupHooks() { - this.lastOperatorWrapper.wrapLastOperator(this.tracing); - Hooks.onEachOperator(SLEUTH_TRACE_REACTOR_KEY, ReactorSleuth.scopePassingSpanOperator(this.tracing)); - Schedulers.setFactory(new Schedulers.Factory() { - @Override public ScheduledExecutorService decorateExecutorService(String schedulerType, - Supplier actual) { - return new TraceableScheduledExecutorService( - TraceReactorConfiguration.this.beanFactory, - actual.get()); - } - }); + return beanFactory -> { }; } @PreDestroy @@ -92,9 +77,48 @@ public class TraceReactorAutoConfiguration { Hooks.resetOnEachOperator(SLEUTH_TRACE_REACTOR_KEY); Schedulers.resetFactory(); } + + @Bean + // for tests + @ConditionalOnMissingBean + HookRegisteringBeanDefinitionRegistryPostProcessor traceHookRegisteringBeanDefinitionRegistryPostProcessor() { + return new HookRegisteringBeanDefinitionRegistryPostProcessor(); + } } } interface LastOperatorWrapper { - void wrapLastOperator(Tracing tracer); + void wrapLastOperator(BeanFactory beanFactory); +} + +class HookRegisteringBeanDefinitionRegistryPostProcessor implements + BeanDefinitionRegistryPostProcessor { + + @Override public void postProcessBeanDefinitionRegistry( + BeanDefinitionRegistry registry) throws BeansException { + } + + @Override public void postProcessBeanFactory( + ConfigurableListableBeanFactory beanFactory) throws BeansException { + LastOperatorWrapper wrapper = beanFactory.getBean(LastOperatorWrapper.class); + setupHooks(wrapper, beanFactory); + } + + void setupHooks(LastOperatorWrapper wrapper, BeanFactory beanFactory) { + wrapper.wrapLastOperator(beanFactory); + Hooks.onEachOperator( + TraceReactorAutoConfiguration.TraceReactorConfiguration.SLEUTH_TRACE_REACTOR_KEY, + ReactorSleuth.scopePassingSpanOperator(beanFactory)); + Schedulers.setFactory(factoryInstance(beanFactory)); + } + + private Schedulers.Factory factoryInstance(final BeanFactory beanFactory) { + return new Schedulers.Factory() { + @Override public ScheduledExecutorService decorateExecutorService(String schedulerType, + Supplier actual) { + return new TraceableScheduledExecutorService(beanFactory, + actual.get()); + } + }; + } } \ No newline at end of file diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/opentracing/BraveTracerTest.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/opentracing/BraveTracerTest.java index a2d071734..15d5ec39c 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/opentracing/BraveTracerTest.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/opentracing/BraveTracerTest.java @@ -15,32 +15,23 @@ */ package org.springframework.cloud.sleuth.instrument.opentracing; -import java.util.ArrayList; -import java.util.Arrays; import java.util.LinkedHashMap; -import java.util.List; import java.util.Map; import brave.Span; -import brave.Tracer; import brave.Tracer.SpanInScope; import brave.Tracing; import brave.opentracing.BraveSpan; import brave.opentracing.BraveSpanContext; import brave.opentracing.BraveTracer; -import brave.propagation.B3Propagation; import brave.propagation.CurrentTraceContext; -import brave.propagation.ExtraFieldPropagation; -import brave.propagation.Propagation; import brave.propagation.StrictCurrentTraceContext; import brave.propagation.TraceContext; import brave.sampler.Sampler; import io.opentracing.Scope; import io.opentracing.propagation.Format; -import io.opentracing.propagation.TextMap; import io.opentracing.propagation.TextMapExtractAdapter; import io.opentracing.propagation.TextMapInjectAdapter; -import org.junit.After; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -52,7 +43,6 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.test.context.junit4.SpringRunner; import zipkin2.Annotation; -import zipkin2.Endpoint; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.data.MapEntry.entry; diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/Issue866Configuration.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/Issue866Configuration.java new file mode 100644 index 000000000..78885fc39 --- /dev/null +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/Issue866Configuration.java @@ -0,0 +1,53 @@ +/* + * Copyright 2013-2018 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.sleuth.instrument.reactor; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +/** + * @author Marcin Grzejszczak + */ +@Configuration +public class Issue866Configuration { + + // we don't want to force direct dependencies between components + // because Spring might just properly setup the context + // we want to ensure that the HRBDRPP is always executed before + // any other object is started + public static TestHook hook; + + @Bean + HookRegisteringBeanDefinitionRegistryPostProcessor overridingProcessorForTests() { + TestHook hook = new TestHook(); + Issue866Configuration.hook = hook; + return hook; + } + + public static class TestHook extends HookRegisteringBeanDefinitionRegistryPostProcessor { + public boolean executed = false; + + @Override public void postProcessBeanFactory( + ConfigurableListableBeanFactory beanFactory) throws BeansException { + super.postProcessBeanFactory(beanFactory); + this.executed = true; + } + } +} + diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java index 6be16ff32..398f798c7 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java @@ -23,6 +23,7 @@ import java.util.stream.Collectors; import brave.Tracer; import brave.sampler.Sampler; import org.awaitility.Awaitility; +import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.Rule; import org.junit.Test; @@ -34,10 +35,12 @@ import org.springframework.boot.autoconfigure.security.reactive.ReactiveSecurity import org.springframework.boot.autoconfigure.security.reactive.ReactiveUserDetailsServiceAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.boot.test.rule.OutputCapture; +import org.springframework.cloud.sleuth.instrument.reactor.Issue866Configuration; import org.springframework.cloud.sleuth.util.ArrayListSpanReporter; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Import; import org.springframework.core.env.Environment; import org.springframework.web.reactive.function.client.ClientResponse; import org.springframework.web.reactive.function.client.WebClient; @@ -46,6 +49,7 @@ import org.springframework.web.reactive.function.server.ServerResponse; import reactor.core.publisher.Flux; import reactor.core.publisher.Hooks; import reactor.core.publisher.Mono; +import reactor.core.scheduler.Scheduler; import reactor.core.scheduler.Schedulers; import zipkin2.Span; @@ -62,6 +66,12 @@ public class FlatMapTests { public static void setup() { Hooks.resetOnLastOperator(); Schedulers.resetFactory(); + Issue866Configuration.hook = null; + } + + @AfterClass + public static void cleanup() { + Issue866Configuration.hook = null; } @Rule public OutputCapture capture = new OutputCapture(); @@ -69,7 +79,8 @@ public class FlatMapTests { @Test public void should_work_with_flat_maps() { //given ConfigurableApplicationContext context = new SpringApplicationBuilder( - FlatMapTests.TestConfiguration.class).web(WebApplicationType.REACTIVE) + FlatMapTests.TestConfiguration.class, Issue866Configuration.class) + .web(WebApplicationType.REACTIVE) .properties("server.port=0", "spring.jmx.enabled=false", "spring.application.name=TraceWebFluxTests", "security.basic.enabled=false", "management.security.enabled=false").run(); @@ -77,6 +88,7 @@ public class FlatMapTests { int port = context.getBean(Environment.class).getProperty("local.server.port", Integer.class); RequestSender sender = context.getBean(RequestSender.class); TestConfiguration config = context.getBean(TestConfiguration.class); + FactoryUser factoryUser = context.getBean(FactoryUser.class); sender.port = port; accumulator.clear(); @@ -103,6 +115,8 @@ public class FlatMapTests { .collect(Collectors.toList())) .as("TraceFilter should not have any trace when receiving a request") .containsOnly(""); + //and #866 + then(factoryUser.wasSchedulerWrapped).isTrue(); }); } @@ -176,10 +190,24 @@ public class FlatMapTests { return Sampler.ALWAYS_SAMPLE; } - @Bean - RequestSender sender(WebClient client, Tracer tracer) { + @Bean RequestSender sender(WebClient client, Tracer tracer) { return new RequestSender(client, tracer); } + // https://github.com/spring-cloud/spring-cloud-sleuth/issues/866 + @Bean + FactoryUser factoryUser() { + return new FactoryUser(); + } + } } + +class FactoryUser { + boolean wasSchedulerWrapped = false; + + FactoryUser() { + Issue866Configuration.TestHook hook = Issue866Configuration.hook; + this.wasSchedulerWrapped = hook != null && hook.executed; + } +} \ No newline at end of file