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 index 576c6ffc0..667080201 100644 --- 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 @@ -29,6 +29,7 @@ import reactor.util.context.Context; * @author Marcin Grzejszczak * @since 2.0.0 */ +// TODO: why are we extending AtomicBoolean and not actually using its methods? final class LazySpanSubscriber extends AtomicBoolean implements SpanSubscription { private final Supplier> supplier; 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 10df1a826..a1746e0ce 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 @@ -18,7 +18,6 @@ package org.springframework.cloud.sleuth.instrument.reactor; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; -import java.util.function.BooleanSupplier; import java.util.function.Function; import brave.Tracing; @@ -67,24 +66,21 @@ public abstract class ReactorSleuth { log.trace("Scope passing operator [" + beanFactory + "]"); } - // Adapt if lazy bean factory - BooleanSupplier isActive = beanFactory instanceof ConfigurableApplicationContext - ? ((ConfigurableApplicationContext) beanFactory)::isActive : () -> true; - return Operators.liftPublisher((p, sub) -> { - // if Flux/Mono #just, #empty, #error + // While supply of scalar types may be deferred, we don't currently scope + // production of values in a trace context. This prevents excessive overhead + // when using constant results such as Flux/Mono #just, #empty, #error if (p instanceof Fuseable.ScalarCallable) { return sub; } - Scannable scannable = Scannable.from(p); - // rest of the logic unchanged... - if (isActive.getAsBoolean()) { + + if (beanFactory instanceof ConfigurableApplicationContext + && ((ConfigurableApplicationContext) beanFactory).isActive()) { if (log.isTraceEnabled()) { log.trace("Spring Context [" + beanFactory + "] already refreshed. Creating a scope " + "passing span subscriber with Reactor Context " + "[" - + sub.currentContext() + "] and name [" + scannable.name() - + "]"); + + sub.currentContext() + "] and name [" + name(sub) + "]"); } return scopePassingSpanSubscription(beanFactory, sub); @@ -92,18 +88,15 @@ public abstract class ReactorSleuth { if (log.isTraceEnabled()) { log.trace("Spring Context [" + beanFactory + "] is not yet refreshed, falling back to lazy span subscriber. Reactor Context is [" - + sub.currentContext() + "] and name is [" + scannable.name() - + "]"); + + sub.currentContext() + "] and name is [" + name(sub) + "]"); } return new LazySpanSubscriber<>( - lazyScopePassingSpanSubscription(beanFactory, scannable, sub)); + new SpanSubscriptionProvider<>(beanFactory, sub)); }); } - static SpanSubscriptionProvider lazyScopePassingSpanSubscription( - BeanFactory beanFactory, Scannable scannable, CoreSubscriber sub) { - return new SpanSubscriptionProvider<>(beanFactory, sub, sub.currentContext(), - scannable.name()); + static String name(CoreSubscriber sub) { + return Scannable.from(sub).name(); } private static Map CACHE = new ConcurrentHashMap<>(); 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 index e90b34836..2c79e5975 100644 --- 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 @@ -22,11 +22,13 @@ import brave.propagation.CurrentTraceContext; import brave.propagation.TraceContext; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.reactivestreams.Subscriber; +import reactor.core.CoreSubscriber; import reactor.util.context.Context; import org.springframework.beans.factory.BeanFactory; +import static org.springframework.cloud.sleuth.instrument.reactor.ReactorSleuth.name; + /** * Supplier to lazily start a {@link SpanSubscription}. * @@ -39,23 +41,20 @@ final class SpanSubscriptionProvider implements Supplier> final BeanFactory beanFactory; - final Subscriber subscriber; + final CoreSubscriber subscriber; final Context context; - final String name; - private volatile CurrentTraceContext currentTraceContext; - SpanSubscriptionProvider(BeanFactory beanFactory, Subscriber subscriber, - Context context, String name) { + SpanSubscriptionProvider(BeanFactory beanFactory, + CoreSubscriber subscriber) { this.beanFactory = beanFactory; this.subscriber = subscriber; - this.context = context; - this.name = name; + this.context = subscriber.currentContext(); if (log.isTraceEnabled()) { log.trace("Spring context [" + beanFactory + "], Reactor context [" + context - + "], name [" + name + "]"); + + "], name [" + name(subscriber) + "]"); } } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProviderTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProviderTests.java index 6838e4f3a..2146941f5 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProviderTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriptionProviderTests.java @@ -19,6 +19,7 @@ package org.springframework.cloud.sleuth.instrument.reactor; import org.assertj.core.api.BDDAssertions; import org.junit.Test; import org.mockito.BDDMockito; +import reactor.core.CoreSubscriber; import reactor.util.context.Context; import org.springframework.beans.factory.BeanFactory; @@ -27,13 +28,18 @@ public class SpanSubscriptionProviderTests { @Test public void should_return_default_tracing_instance_when_exception_thrown_upon_bean_retrieval() { + CoreSubscriber subscriber = BDDMockito.mock(CoreSubscriber.class); + BDDMockito.when(subscriber.currentContext()).thenReturn(Context.empty()); + BeanFactory beanFactory = BDDMockito.mock(BeanFactory.class); + BDDMockito.when(beanFactory.getBean(BDDMockito.any(Class.class))) .thenThrow(new IllegalStateException()); - SpanSubscriptionProvider provider = new SpanSubscriptionProvider(beanFactory, - null, Context.empty(), "example"); - SpanSubscription spanSubscription = provider.get(); + SpanSubscriptionProvider provider = new SpanSubscriptionProvider<>( + beanFactory, subscriber); + + SpanSubscription spanSubscription = provider.get(); BDDAssertions.then(spanSubscription).isNotNull(); }