diff --git a/pom.xml b/pom.xml
index 7ebac07ad..77bb151d1 100644
--- a/pom.xml
+++ b/pom.xml
@@ -258,6 +258,12 @@
3.8.0
test
+
+ net.jcip
+ jcip-annotations
+ 1.0
+ test
+
diff --git a/spring-cloud-sleuth-core/pom.xml b/spring-cloud-sleuth-core/pom.xml
index 60b86be35..fb83731d8 100644
--- a/spring-cloud-sleuth-core/pom.xml
+++ b/spring-cloud-sleuth-core/pom.xml
@@ -271,6 +271,11 @@
spring-cloud-starter-netflix-eureka-client
test
+
+ net.jcip
+ jcip-annotations
+ test
+
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 e8c0fbd1a..46e3f751e 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
@@ -41,6 +41,8 @@ public abstract class ReactorSleuth {
private static final Log log = LogFactory.getLog(ReactorSleuth.class);
+ private static volatile boolean CONTEXT_REFRESHED = false;
+
/**
* Return a span operator pointcut given a {@link BeanFactory}. This can be used in reactor
* via {@link reactor.core.publisher.Flux#transform(Function)}, {@link
@@ -91,7 +93,7 @@ public abstract class ReactorSleuth {
beanFactory,
sub,
sub.currentContext(),
- scannable.name());
+ scannable);
}
/**
@@ -141,8 +143,15 @@ public abstract class ReactorSleuth {
}
private static boolean contextRefreshed(BeanFactory beanFactory) {
+ if (CONTEXT_REFRESHED) {
+ return true;
+ }
try {
- return beanFactory.getBean(ApplicationContextRefreshedListener.class).isRefreshed();
+ boolean contextRefreshed = beanFactory.getBean(ApplicationContextRefreshedListener.class).isRefreshed();
+ if (contextRefreshed) {
+ CONTEXT_REFRESHED = true;
+ }
+ return contextRefreshed;
} catch (NoSuchBeanDefinitionException e) {
return false;
}
@@ -154,7 +163,7 @@ public abstract class ReactorSleuth {
beanFactory,
sub,
sub.currentContext(),
- scannable.name()) {
+ scannable) {
@Override SpanSubscription newCoreSubscriber(Tracing tracing) {
return new ScopePassingSpanSubscriber(
sub,
@@ -166,4 +175,4 @@ public abstract class ReactorSleuth {
private ReactorSleuth() {
}
-}
\ No newline at end of file
+}
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 e3fa48c1d..0cbc073c6 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
@@ -19,8 +19,9 @@ package org.springframework.cloud.sleuth.instrument.reactor;
import java.util.concurrent.atomic.AtomicBoolean;
import brave.Span;
-import brave.Tracer;
import brave.Tracing;
+import brave.propagation.CurrentTraceContext;
+import brave.propagation.TraceContext;
import reactor.util.context.Context;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
@@ -38,18 +39,18 @@ final class ScopePassingSpanSubscriber extends AtomicBoolean implements SpanS
private static final Log log = LogFactory.getLog(ScopePassingSpanSubscriber.class);
- private final Span span;
+ private final TraceContext spanTraceContext;
+ private final CurrentTraceContext currentTraceContext;
private final Subscriber super T> subscriber;
private final Context context;
- private final Tracer tracer;
private Subscription s;
ScopePassingSpanSubscriber(Subscriber super T> subscriber, Context ctx, Tracing tracing) {
this.subscriber = subscriber;
- this.tracer = tracing.tracer();
+ this.currentTraceContext = tracing.currentTraceContext();
Span root = ctx != null ?
- ctx.getOrDefault(Span.class, this.tracer.currentSpan()) : null;
- this.span = root;
+ ctx.getOrDefault(Span.class, tracing.tracer().currentSpan()) : null;
+ this.spanTraceContext = root != null ? root.context() : null;
this.context = ctx != null && root != null ? ctx.put(Span.class, root) :
ctx != null ? ctx : Context.empty();
if (log.isTraceEnabled()) {
@@ -59,25 +60,25 @@ final class ScopePassingSpanSubscriber extends AtomicBoolean implements SpanS
@Override public void onSubscribe(Subscription subscription) {
this.s = subscription;
- try (Tracer.SpanInScope inScope = this.tracer.withSpanInScope(this.span)) {
+ try (CurrentTraceContext.Scope inScope = this.currentTraceContext.maybeScope(this.spanTraceContext)) {
this.subscriber.onSubscribe(this);
}
}
@Override public void request(long n) {
- try (Tracer.SpanInScope inScope = this.tracer.withSpanInScope(this.span)) {
+ try (CurrentTraceContext.Scope inScope = this.currentTraceContext.maybeScope(this.spanTraceContext)) {
this.s.request(n);
}
}
@Override public void cancel() {
- try (Tracer.SpanInScope inScope = this.tracer.withSpanInScope(this.span)) {
+ try (CurrentTraceContext.Scope inScope = this.currentTraceContext.maybeScope(this.spanTraceContext)) {
this.s.cancel();
}
}
@Override public void onNext(T o) {
- try (Tracer.SpanInScope inScope = this.tracer.withSpanInScope(this.span)) {
+ try (CurrentTraceContext.Scope inScope = this.currentTraceContext.maybeScope(this.spanTraceContext)) {
this.subscriber.onNext(o);
}
}
@@ -93,4 +94,4 @@ final class ScopePassingSpanSubscriber extends AtomicBoolean implements SpanS
@Override public Context currentContext() {
return this.context;
}
-}
\ No newline at end of file
+}
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 baa84eba9..7892b17a7 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
@@ -16,9 +16,11 @@
package org.springframework.cloud.sleuth.instrument.reactor;
+import java.util.Objects;
import java.util.function.Supplier;
import brave.Tracing;
+import reactor.core.Scannable;
import reactor.util.context.Context;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
@@ -37,18 +39,19 @@ class SpanSubscriptionProvider implements Supplier> {
final BeanFactory beanFactory;
final Subscriber super T> subscriber;
final Context context;
- final String name;
+ final Scannable scannable;
private Tracing tracing;
SpanSubscriptionProvider(BeanFactory beanFactory,
Subscriber super T> subscriber,
- Context context, String name) {
+ Context context, Scannable scannable) {
this.beanFactory = beanFactory;
this.subscriber = subscriber;
this.context = context;
- this.name = name;
+ Objects.requireNonNull(scannable, "scannable must not be null");
+ this.scannable = scannable;
if (log.isTraceEnabled()) {
- log.trace("Context [" + context + "], name [" + name + "]");
+ log.trace("Context [" + context + "], name [" + scannable.name() + "]");
}
}
@@ -57,7 +60,7 @@ class SpanSubscriptionProvider implements Supplier> {
}
SpanSubscription newCoreSubscriber(Tracing tracing) {
- return new SpanSubscriber<>(this.subscriber, this.context, tracing, this.name);
+ return new SpanSubscriber<>(this.subscriber, this.context, tracing, this.scannable.name());
}
private Tracing tracing() {
diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/log/Slf4jCurrentTraceContext.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/log/Slf4jCurrentTraceContext.java
index ce8aa5333..8bf1cbbc8 100644
--- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/log/Slf4jCurrentTraceContext.java
+++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/log/Slf4jCurrentTraceContext.java
@@ -144,4 +144,4 @@ public final class Slf4jCurrentTraceContext extends CurrentTraceContext {
MDC.remove(key);
}
}
-}
\ No newline at end of file
+}
diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/scheduling/TracingOnScheduledTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/scheduling/TracingOnScheduledTests.java
index 48c7c281b..e609da803 100644
--- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/scheduling/TracingOnScheduledTests.java
+++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/scheduling/TracingOnScheduledTests.java
@@ -22,6 +22,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
import brave.Span;
import brave.Tracing;
import brave.sampler.Sampler;
+import net.jcip.annotations.NotThreadSafe;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.junit.Before;
@@ -47,6 +48,7 @@ import static org.awaitility.Awaitility.await;
@RunWith(SpringRunner.class)
@SpringBootTest(classes = {ScheduledTestConfiguration.class})
@DirtiesContext
+@NotThreadSafe
public class TracingOnScheduledTests {
@Autowired
diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/log/Slf4JSpanLoggerTest.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/log/Slf4JSpanLoggerTest.java
index 83d3847a0..2d9d4702c 100644
--- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/log/Slf4JSpanLoggerTest.java
+++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/log/Slf4JSpanLoggerTest.java
@@ -21,6 +21,7 @@ import brave.Tracing;
import brave.propagation.CurrentTraceContext.Scope;
import brave.propagation.StrictScopeDecorator;
import brave.propagation.ThreadLocalCurrentTraceContext;
+import brave.propagation.TraceContext;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
@@ -56,31 +57,104 @@ public class Slf4JSpanLoggerTest {
@Test
public void should_set_entries_to_mdc_from_span() throws Exception {
- Scope scope = this.slf4jCurrentTraceContext.newScope(this.span.context());
+ try (Scope scope = this.slf4jCurrentTraceContext.newScope(this.span.context())) {
+ assertMDCInfoEqualToSpanInfo(span.context());
+ }
- assertThat(MDC.get("X-B3-TraceId")).isEqualTo(span.context().traceIdString());
- assertThat(MDC.get("traceId")).isEqualTo(span.context().traceIdString());
+ assertMDCInfoNullOrEmpty();
+ }
- scope.close();
+ @Test
+ public void should_set_entries_to_mdc_from_two_spans() throws Exception {
+ try (Scope scope = this.slf4jCurrentTraceContext.newScope(this.span.context())) {
- assertThat(MDC.get("X-B3-TraceId")).isNullOrEmpty();
- assertThat(MDC.get("traceId")).isNullOrEmpty();
+ assertMDCInfoEqualToSpanInfo(span.context());
+
+ try (Scope scopeInner = this.slf4jCurrentTraceContext.newScope(this.span.context())) {
+ assertMDCInfoEqualToSpanInfo(span.context());
+ }
+
+ assertMDCInfoEqualToSpanInfo(span.context());
+ }
+
+ assertMDCInfoNullOrEmpty();
+ }
+
+ @Test
+ public void should_set_entries_to_mdc_from_two_spans1() throws Exception {
+ try (Scope scope = this.slf4jCurrentTraceContext.newScope(this.span.context())) {
+
+ assertMDCInfoEqualToSpanInfo(span.context());
+
+ try (Scope scopeInner = this.slf4jCurrentTraceContext.newScope(null)) {
+ assertMDCInfoNullOrEmpty();
+ }
+
+ assertMDCInfoEqualToSpanInfo(span.context());
+ }
+
+ assertMDCInfoNullOrEmpty();
+ }
+
+ @Test
+ public void should_set_entries_to_mdc_from_two_spans2() throws Exception {
+ try (Scope scope = this.slf4jCurrentTraceContext.newScope(this.span.context())) {
+
+ assertMDCInfoEqualToSpanInfo(span.context());
+
+ TraceContext nextSpan = this.tracing.tracer().nextSpan().start().context();
+ try (Scope scopeInner = this.slf4jCurrentTraceContext.newScope(nextSpan)) {
+ assertMDCInfoEqualToSpanInfo(nextSpan);
+ }
+
+ assertMDCInfoEqualToSpanInfo(span.context());
+ }
+
+ assertMDCInfoNullOrEmpty();
+ }
+
+ @Test
+ public void should_set_entries_to_mdc_from_two_spans3() throws Exception {
+ try (Scope scope = this.slf4jCurrentTraceContext.newScope(null)) {
+
+ assertMDCInfoNullOrEmpty();
+
+ try (Scope scopeInner = this.slf4jCurrentTraceContext.newScope(null)) {
+ assertMDCInfoNullOrEmpty();
+ }
+
+ assertMDCInfoNullOrEmpty();
+ }
+
+ assertMDCInfoNullOrEmpty();
}
@Test
public void should_remove_entries_from_mdc_from_null_span() throws Exception {
MDC.put("X-B3-TraceId", "A");
MDC.put("traceId", "A");
+ MDC.put("X-B3-SpanId", "A");
+ MDC.put("spanId", "A");
- Scope scope = this.slf4jCurrentTraceContext
- .newScope(null);
-
- assertThat(MDC.get("X-B3-TraceId")).isNullOrEmpty();
- assertThat(MDC.get("traceId")).isNullOrEmpty();
-
- scope.close();
+ try (Scope scope = this.slf4jCurrentTraceContext.newScope(null)) {
+ assertMDCInfoNullOrEmpty();
+ }
assertThat(MDC.get("X-B3-TraceId")).isEqualTo("A");
assertThat(MDC.get("traceId")).isEqualTo("A");
}
-}
\ No newline at end of file
+
+ private void assertMDCInfoEqualToSpanInfo(TraceContext span){
+ assertThat(MDC.get("X-B3-TraceId")).isEqualTo(span.traceIdString());
+ assertThat(MDC.get("traceId")).isEqualTo(span.traceIdString());
+ assertThat(MDC.get("X-B3-SpanId")).isEqualTo(span.spanIdString());
+ assertThat(MDC.get("spanId")).isEqualTo(span.spanIdString());
+ }
+
+ private void assertMDCInfoNullOrEmpty(){
+ assertThat(MDC.get("X-B3-TraceId")).isNullOrEmpty();
+ assertThat(MDC.get("traceId")).isNullOrEmpty();
+ assertThat(MDC.get("X-B3-SpanId")).isNullOrEmpty();
+ assertThat(MDC.get("spanId")).isNullOrEmpty();
+ }
+}