From c2cd689ec7aa0193d44354d8742e2fda82cf140b Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 6 Feb 2020 13:36:57 +0800 Subject: [PATCH 01/12] Reduces the amount of times getBean is invoked from reactor-netty (#1551) This consolidates access to the `HttpTracing` bean so that the overhead of netty `HttpClient` is reduced. It also pulls the reactor-netty test into its own file. --- .../instrument/reactor/ReactorSleuth.java | 5 +- .../client/HttpClientBeanPostProcessor.java | 101 ++++++------- .../TraceWebClientAutoConfiguration.java | 4 +- .../reactor => internal}/LazyBean.java | 14 +- ...ReactorNettyHttpClientSpringBootTests.java | 138 ++++++++++++++++++ .../client/integration/WebClientTests.java | 34 ----- .../reactor => internal}/LazyBeanTests.java | 2 +- 7 files changed, 198 insertions(+), 100 deletions(-) rename spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/{instrument/reactor => internal}/LazyBean.java (83%) create mode 100644 spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java rename spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/{instrument/reactor => internal}/LazyBeanTests.java (95%) 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 7f095b779..fa013b7c8 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 @@ -30,6 +30,7 @@ import reactor.core.Scannable; import reactor.core.publisher.Operators; import reactor.util.context.Context; +import org.springframework.cloud.sleuth.internal.LazyBean; import org.springframework.context.ConfigurableApplicationContext; /** @@ -70,8 +71,8 @@ public abstract class ReactorSleuth { // keep a reference outside the lambda so that any caching will be visible to // all publishers - LazyBean lazyCurrentTraceContext = new LazyBean<>( - springContext, CurrentTraceContext.class); + LazyBean lazyCurrentTraceContext = LazyBean + .create(springContext, CurrentTraceContext.class); return Operators.liftPublisher((p, sub) -> { // We don't scope scalar results as they happen in an instant. This prevents diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java index dcf21694e..5558b5e99 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java @@ -33,26 +33,29 @@ import reactor.netty.http.client.HttpClientRequest; import reactor.netty.http.client.HttpClientResponse; import org.springframework.beans.BeansException; -import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.config.BeanPostProcessor; +import org.springframework.cloud.sleuth.internal.LazyBean; +import org.springframework.context.ConfigurableApplicationContext; class HttpClientBeanPostProcessor implements BeanPostProcessor { - private final BeanFactory beanFactory; + final ConfigurableApplicationContext springContext; - HttpClientBeanPostProcessor(BeanFactory beanFactory) { - this.beanFactory = beanFactory; + HttpClientBeanPostProcessor(ConfigurableApplicationContext springContext) { + this.springContext = springContext; } @Override public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { + LazyBean httpTracing = LazyBean.create(this.springContext, + HttpTracing.class); if (bean instanceof HttpClient) { - return ((HttpClient) bean).mapConnect(new TracingMapConnect(this.beanFactory)) - .doOnRequest(TracingDoOnRequest.create(this.beanFactory)) - .doOnRequestError(TracingDoOnErrorRequest.create(this.beanFactory)) - .doOnResponse(TracingDoOnResponse.create(this.beanFactory)) - .doOnResponseError(TracingDoOnErrorResponse.create(this.beanFactory)); + return ((HttpClient) bean).mapConnect(new TracingMapConnect(httpTracing)) + .doOnRequest(TracingDoOnRequest.create(httpTracing)) + .doOnRequestError(TracingDoOnErrorRequest.create(httpTracing)) + .doOnResponse(TracingDoOnResponse.create(httpTracing)) + .doOnResponseError(TracingDoOnErrorResponse.create(httpTracing)); } return bean; } @@ -60,12 +63,12 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { private static class TracingMapConnect implements BiFunction, Bootstrap, Mono> { - private final BeanFactory beanFactory; + final LazyBean httpTracing; - private Tracer tracer; + Tracer tracer; - TracingMapConnect(BeanFactory beanFactory) { - this.beanFactory = beanFactory; + TracingMapConnect(LazyBean httpTracing) { + this.httpTracing = httpTracing; } @Override @@ -77,7 +80,7 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { private Tracer tracer() { if (this.tracer == null) { - this.tracer = this.beanFactory.getBean(Tracer.class); + this.tracer = this.httpTracing.get().tracing().tracer(); } return this.tracer; } @@ -87,39 +90,30 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { private static class TracingDoOnRequest implements BiConsumer { - final BeanFactory beanFactory; - - HttpTracing httpTracing; + final LazyBean httpTracing; List propagationKeys; HttpClientHandler handler; - TracingDoOnRequest(BeanFactory beanFactory) { - this.beanFactory = beanFactory; + TracingDoOnRequest(LazyBean httpTracing) { + this.httpTracing = httpTracing; } - static TracingDoOnRequest create(BeanFactory beanFactory) { - return new TracingDoOnRequest(beanFactory); + static TracingDoOnRequest create(LazyBean httpTracing) { + return new TracingDoOnRequest(httpTracing); } - private HttpTracing httpTracing() { - if (this.httpTracing == null) { - this.httpTracing = this.beanFactory.getBean(HttpTracing.class); - } - return this.httpTracing; - } - - private List propagationKeys() { + List propagationKeys() { if (this.propagationKeys == null) { - this.propagationKeys = httpTracing().tracing().propagation().keys(); + this.propagationKeys = httpTracing.get().tracing().propagation().keys(); } return this.propagationKeys; } - private HttpClientHandler handler() { + HttpClientHandler handler() { if (this.handler == null) { - this.handler = HttpClientHandler.create(httpTracing()); + this.handler = HttpClientHandler.create(httpTracing.get()); } return this.handler; } @@ -147,12 +141,12 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { private static class TracingDoOnResponse extends AbstractTracingDoOnHandler implements BiConsumer { - TracingDoOnResponse(BeanFactory beanFactory) { - super(beanFactory); + TracingDoOnResponse(LazyBean httpTracing) { + super(httpTracing); } - static TracingDoOnResponse create(BeanFactory beanFactory) { - return new TracingDoOnResponse(beanFactory); + static TracingDoOnResponse create(LazyBean httpTracing) { + return new TracingDoOnResponse(httpTracing); } @Override @@ -165,12 +159,12 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { private static class TracingDoOnErrorRequest extends AbstractTracingDoOnHandler implements BiConsumer { - TracingDoOnErrorRequest(BeanFactory beanFactory) { - super(beanFactory); + TracingDoOnErrorRequest(LazyBean httpTracing) { + super(httpTracing); } - static TracingDoOnErrorRequest create(BeanFactory beanFactory) { - return new TracingDoOnErrorRequest(beanFactory); + static TracingDoOnErrorRequest create(LazyBean httpTracing) { + return new TracingDoOnErrorRequest(httpTracing); } @Override @@ -183,12 +177,12 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { private static class TracingDoOnErrorResponse extends AbstractTracingDoOnHandler implements BiConsumer { - TracingDoOnErrorResponse(BeanFactory beanFactory) { - super(beanFactory); + TracingDoOnErrorResponse(LazyBean httpTracing) { + super(httpTracing); } - static TracingDoOnErrorResponse create(BeanFactory beanFactory) { - return new TracingDoOnErrorResponse(beanFactory); + static TracingDoOnErrorResponse create(LazyBean httpTracing) { + return new TracingDoOnErrorResponse(httpTracing); } @Override @@ -200,26 +194,17 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { private static abstract class AbstractTracingDoOnHandler { - final BeanFactory beanFactory; - - HttpTracing httpTracing; + final LazyBean httpTracing; HttpClientHandler handler; - AbstractTracingDoOnHandler(BeanFactory beanFactory) { - this.beanFactory = beanFactory; + AbstractTracingDoOnHandler(LazyBean httpTracing) { + this.httpTracing = httpTracing; } - private HttpTracing httpTracing() { - if (this.httpTracing == null) { - this.httpTracing = this.beanFactory.getBean(HttpTracing.class); - } - return this.httpTracing; - } - - private HttpClientHandler handler() { + HttpClientHandler handler() { if (this.handler == null) { - this.handler = HttpClientHandler.create(httpTracing()); + this.handler = HttpClientHandler.create(httpTracing.get()); } return this.handler; } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientAutoConfiguration.java index ef596dddf..bb59752b7 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientAutoConfiguration.java @@ -162,8 +162,8 @@ public class TraceWebClientAutoConfiguration { @Bean static HttpClientBeanPostProcessor httpClientBeanPostProcessor( - BeanFactory beanFactory) { - return new HttpClientBeanPostProcessor(beanFactory); + ConfigurableApplicationContext springContext) { + return new HttpClientBeanPostProcessor(springContext); } } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/LazyBean.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/internal/LazyBean.java similarity index 83% rename from spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/LazyBean.java rename to spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/internal/LazyBean.java index 364d18648..8a55e0b31 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/LazyBean.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/internal/LazyBean.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.sleuth.instrument.reactor; +package org.springframework.cloud.sleuth.internal; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -25,8 +25,16 @@ import org.springframework.lang.Nullable; /** * Avoids calling the expensive {@link ConfigurableApplicationContext#getBean(Class)} many * times or throwing an exception. + * + *

+ * Note: This is an internal class to sleuth and must not be used by external code. */ -final class LazyBean { +public final class LazyBean { + + public static LazyBean create(ConfigurableApplicationContext springContext, + Class requiredType) { + return new LazyBean<>(springContext, requiredType); + } // spring-jcl uses commons-logging, so do we. private static final Log log = LogFactory.getLog(LazyBean.class); @@ -47,7 +55,7 @@ final class LazyBean { * @return the bean value or null if there was an exception getting it. */ @Nullable - T get() { + public T get() { if (this.value != null) { return this.value; } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java new file mode 100644 index 000000000..30b1dc1c8 --- /dev/null +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java @@ -0,0 +1,138 @@ +/* + * Copyright 2013-2019 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 + * + * https://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.web.client; + +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; + +import brave.propagation.B3SinglePropagation; +import brave.propagation.Propagation; +import brave.sampler.Sampler; +import org.junit.After; +import org.junit.Test; +import org.junit.runner.RunWith; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.netty.DisposableServer; +import reactor.netty.http.client.HttpClient; +import reactor.netty.http.server.HttpServer; +import zipkin2.Span; +import zipkin2.reporter.Reporter; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.test.context.junit4.SpringRunner; +import org.springframework.web.reactive.function.client.WebClient; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * This tests {@link HttpClient} instrumentation performed by + * {@link HttpClientBeanPostProcessor}, as wired by auto-configuration. + * + *

+ * Note: {@link HttpClient} can be an implementation of {@link WebClient}, so + * care should be taken to also test that integration. For example, it would be easy to + * create duplicate client spans for the same request. + */ +@SpringBootTest(classes = ReactorNettyHttpClientSpringBootTests.TestConfiguration.class, + webEnvironment = SpringBootTest.WebEnvironment.NONE) +@RunWith(SpringRunner.class) +public class ReactorNettyHttpClientSpringBootTests { + + DisposableServer disposableServer; + + @Autowired + HttpClient httpClient; + + @Autowired + BlockingQueue spans; + + @After + public void tearDown() { + if (disposableServer != null) { + disposableServer.disposeNow(); + } + this.spans.clear(); + } + + @Test + public void shouldSendTraceContextToServer_rootSpan() throws Exception { + disposableServer = HttpServer.create().port(0) + // this reads the trace context header, b3, returning it in the response + .handle((in, out) -> out + .sendString(Flux.just(in.requestHeaders().get("b3")))) + .bindNow(); + + Mono request = httpClient.port(disposableServer.port()).get().uri("/") + .responseContent().aggregate().asString(); + + String b3SingleHeaderReadByServer = request.block(); + + Span clientSpan = takeClientSpan(); + + assertThat(b3SingleHeaderReadByServer) + .isEqualTo(clientSpan.traceId() + "-" + clientSpan.id() + "-1"); + } + + /** Call this to block until a span was reported */ + Span takeClientSpan() throws InterruptedException { + Span result = spans.poll(1, TimeUnit.SECONDS); + assertThat(result).withFailMessage("Span was not reported").isNotNull(); + assertThat(result.kind()).isEqualTo(Span.Kind.CLIENT); + return result; + } + + @Configuration + @EnableAutoConfiguration + static class TestConfiguration { + + @Bean + Propagation.Factory propagationFactory() { + return B3SinglePropagation.FACTORY; + } + + @Bean + Sampler sampler() { + return Sampler.ALWAYS_SAMPLE; + } + + /** + * Use a blocking queue as it is simpler than wrapping everything in awaitility + */ + @Bean + BlockingQueue spans() { + return new LinkedBlockingQueue<>(); + } + + @Bean + Reporter spanReporter(BlockingQueue spans) { + return spans::add; + } + + @Bean + HttpClient reactorHttpClient() { + return HttpClient.create(); + } + + } + +} diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java index 23a066368..26962aa66 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java @@ -55,8 +55,6 @@ import org.junit.ClassRule; import org.junit.Rule; import org.junit.Test; import org.junit.runner.RunWith; -import reactor.netty.http.client.HttpClient; -import reactor.netty.http.client.HttpClientResponse; import zipkin2.Annotation; import zipkin2.reporter.Reporter; @@ -136,9 +134,6 @@ public class WebClientTests { @Autowired HttpClientBuilder httpClientBuilder; // #845 - @Autowired - HttpClient nettyHttpClient; - @Autowired HttpAsyncClientBuilder httpAsyncClientBuilder; // #845 @@ -278,30 +273,6 @@ public class WebClientTests { then(this.reporter.getSpans()).isNotEmpty(); } - @Test - @SuppressWarnings("unchecked") - public void shouldAttachTraceIdWhenCallingAnotherServiceForNettyHttpClient() - throws Exception { - Span span = this.tracer.nextSpan().name("foo").start(); - - try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) { - HttpClientResponse response = this.nettyHttpClient.get() - .uri("http://localhost:" + this.port).response().block(); - - then(response).isNotNull(); - } - - Awaitility.await().untilAsserted(() -> { - then(this.tracer.currentSpan()).isNull(); - System.out.println("Collected span " + this.reporter.getSpans()); - then(this.reporter.getSpans()).isNotEmpty() - .extracting("traceId", String.class) - // we can have some bizarre spans popping up - .contains(span.context().traceIdString()); - then(this.reporter.getSpans()).extracting("kind.name").contains("CLIENT"); - }); - } - @Test @SuppressWarnings("unchecked") public void shouldAttachTraceIdWhenCallingAnotherServiceForHttpClient() @@ -601,11 +572,6 @@ public class WebClientTests { return new MyRestTemplateCustomizer(); } - @Bean - HttpClient reactorHttpClient() { - return HttpClient.create(); - } - } static class MyRestTemplateCustomizer implements RestTemplateCustomizer { diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/LazyBeanTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/internal/LazyBeanTests.java similarity index 95% rename from spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/LazyBeanTests.java rename to spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/internal/LazyBeanTests.java index 3f9694ba6..db03fb410 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/LazyBeanTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/internal/LazyBeanTests.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.sleuth.instrument.reactor; +package org.springframework.cloud.sleuth.internal; import brave.propagation.CurrentTraceContext; import org.junit.Test; From 3c0932bb321dff7f8876ab63e8d5b92887903d3e Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 6 Feb 2020 15:02:12 +0800 Subject: [PATCH 02/12] Fixes null bugs on WebClient canceled request (#1548) --- .../TraceWebClientBeanPostProcessor.java | 36 +++++++++++-------- .../client/integration/WebClientTests.java | 23 ++++++++++-- 2 files changed, 42 insertions(+), 17 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientBeanPostProcessor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientBeanPostProcessor.java index 558f384cf..e50f31fc0 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientBeanPostProcessor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/TraceWebClientBeanPostProcessor.java @@ -18,6 +18,7 @@ package org.springframework.cloud.sleuth.instrument.web.client; import java.util.Collections; import java.util.List; +import java.util.concurrent.CancellationException; import java.util.function.Consumer; import java.util.function.Function; @@ -131,9 +132,14 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction { } }; - private static final String CLIENT_SPAN_KEY = "sleuth.webclient.clientSpan"; + static final String CLIENT_SPAN_KEY = "sleuth.webclient.clientSpan"; - private static final String CANCELLED_SUBSCRIPTION_ERROR = "CANCELLED"; + static final Exception CANCELLED_ERROR = new CancellationException("CANCELLED") { + @Override + public Throwable fillInStackTrace() { + return this; // stack trace doesn't add value here + } + }; final ConfigurableApplicationContext springContext; @@ -170,6 +176,8 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction { } MonoWebClientTrace trace = new MonoWebClientTrace(next, wrapper.buildRequest(), this, span); + // TODO: investigate why this commit leaks a scope: + // 8f5bcdabd7af23df443e771432eb85597f3b3076 tracer().withSpanInScope(parentSpan); return trace; } @@ -356,13 +364,14 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction { return this.context; } - void handleReceive(Span clientSpan, ClientResponse clientResponse, - Throwable throwable) { + void handleReceive(Span clientSpan, @Nullable ClientResponse res, + @Nullable Throwable error) { if (log.isTraceEnabled()) { log.trace("Handling receive"); } - this.handler.handleReceive(new HttpClientResponse(clientResponse), - throwable, clientSpan); + HttpClientResponse response = res != null ? new HttpClientResponse(res) + : null; + this.handler.handleReceive(response, error, clientSpan); if (log.isTraceEnabled()) { log.trace("Closed scope"); } @@ -374,32 +383,31 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction { + this.span + "]"); } - this.span.tag("error", CANCELLED_SUBSCRIPTION_ERROR); - handleReceive(this.span, null, null); + handleReceive(this.span, null, CANCELLED_ERROR); } void terminateSpan(@Nullable ClientResponse clientResponse, - @Nullable Throwable throwable) { + @Nullable Throwable error) { if (clientResponse == null) { if (log.isDebugEnabled()) { log.debug("No response was returned. Will close the span [" + this.span + "]"); } - handleReceive(this.span, clientResponse, throwable); + handleReceive(this.span, null, error); return; } int statusCode = clientResponse.rawStatusCode(); - boolean error = statusCode >= 400; - if (error) { + boolean isHttpError = statusCode >= 400; + if (isHttpError) { if (log.isDebugEnabled()) { log.debug( "Non positive status code was returned from the call. Will close the span [" + this.span + "]"); } - throwable = new RestClientException( + error = new RestClientException( "Status code of the response is [" + statusCode + "]"); } - handleReceive(this.span, clientResponse, throwable); + handleReceive(this.span, clientResponse, error); } } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java index 26962aa66..3c47698aa 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java @@ -40,6 +40,7 @@ import com.netflix.loadbalancer.ILoadBalancer; import com.netflix.loadbalancer.Server; import junitparams.JUnitParamsRunner; import junitparams.Parameters; +import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.http.HttpResponse; import org.apache.http.client.methods.HttpGet; @@ -55,6 +56,8 @@ import org.junit.ClassRule; import org.junit.Rule; import org.junit.Test; import org.junit.runner.RunWith; +import org.reactivestreams.Subscription; +import reactor.core.publisher.BaseSubscriber; import zipkin2.Annotation; import zipkin2.reporter.Reporter; @@ -112,8 +115,7 @@ public class WebClientTests { static final String SAMPLED_NAME = "X-B3-Sampled"; static final String PARENT_ID_NAME = "X-B3-ParentSpanId"; - private static final org.apache.commons.logging.Log log = LogFactory - .getLog(WebClientTests.class); + private static final Log log = LogFactory.getLog(WebClientTests.class); @Rule public final SpringMethodRule springMethodRule = new SpringMethodRule(); @@ -351,7 +353,7 @@ public class WebClientTests { @Test @SuppressWarnings("unchecked") - public void shouldWorkWhenCustomStatusCodeIsReturned() throws InterruptedException { + public void shouldWorkWhenCustomStatusCodeIsReturned() { Span span = this.tracer.nextSpan().name("foo").start(); try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) { @@ -370,6 +372,21 @@ public class WebClientTests { .contains("CLIENT"); } + @Test + public void shouldTagOnCancel() { + this.webClient.get().uri("http://localhost:" + this.port + "/doNotSkip") + .retrieve().bodyToMono(String.class) + .subscribe(new BaseSubscriber() { + @Override + protected void hookOnSubscribe(Subscription subscription) { + cancel(); + } + }); + + then(this.reporter.getSpans()).isNotEmpty(); + then(this.reporter.getSpans().get(0).tags()).containsEntry("error", "CANCELLED"); + } + @Test public void shouldRespectSkipPattern() { this.webClient.get().uri("http://localhost:" + this.port + "/skip").retrieve() From 2f0e3ed2c73e69662682bb724cbeea7e1f01080e Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 6 Feb 2020 15:36:50 +0800 Subject: [PATCH 03/12] formatting --- .../messaging/SleuthMessagingProperties.java | 6 +- .../messaging/TracingChannelInterceptor.java | 85 +++++---- .../TracingChannelInterceptorTest.java | 165 +++++++++++------- 3 files changed, 158 insertions(+), 98 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/SleuthMessagingProperties.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/SleuthMessagingProperties.java index 08806b453..80655caa9 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/SleuthMessagingProperties.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/SleuthMessagingProperties.java @@ -57,9 +57,11 @@ public class SleuthMessagingProperties { /** * An array of patterns against which channel names will be matched. * @see org.springframework.integration.config.GlobalChannelInterceptor#patterns() - * Defaults to any channel name not matching the Hystrix Stream and functional Stream channel names. + * Defaults to any channel name not matching the Hystrix Stream and functional + * Stream channel names. */ - private String[] patterns = new String[] { "!hystrixStreamOutput*", "*", "!channel*"}; + private String[] patterns = new String[] { "!hystrixStreamOutput*", "*", + "!channel*" }; /** * Enable Spring Integration sleuth instrumentation. diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java index bef74c433..6a0330028 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java @@ -59,7 +59,8 @@ import org.springframework.util.ClassUtils; * * @author Marcin Grzejszczak */ -public final class TracingChannelInterceptor extends ChannelInterceptorAdapter implements ExecutorChannelInterceptor { +public final class TracingChannelInterceptor extends ChannelInterceptorAdapter + implements ExecutorChannelInterceptor { /** * Name of the class in Spring Cloud Stream that is a direct channel. @@ -108,22 +109,25 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i @Autowired TracingChannelInterceptor(Tracing tracing) { - this(tracing, MessageHeaderPropagation.INSTANCE, MessageHeaderPropagation.INSTANCE); + this(tracing, MessageHeaderPropagation.INSTANCE, + MessageHeaderPropagation.INSTANCE); } - TracingChannelInterceptor(Tracing tracing, Propagation.Setter setter, + TracingChannelInterceptor(Tracing tracing, + Propagation.Setter setter, Propagation.Getter getter) { this.tracing = tracing; this.tracer = tracing.tracer(); this.threadLocalSpan = ThreadLocalSpan.create(this.tracer); this.injector = tracing.propagation().injector(setter); this.extractor = tracing.propagation().extractor(getter); - this.integrationObjectSupportPresent = ClassUtils - .isPresent("org.springframework.integration.context.IntegrationObjectSupport", null); - this.hasDirectChannelClass = ClassUtils.isPresent("org.springframework.integration.channel.DirectChannel", - null); - this.directWithAttributesChannelClass = ClassUtils.isPresent(STREAM_DIRECT_CHANNEL, null) - ? ClassUtils.resolveClassName(STREAM_DIRECT_CHANNEL, null) : null; + this.integrationObjectSupportPresent = ClassUtils.isPresent( + "org.springframework.integration.context.IntegrationObjectSupport", null); + this.hasDirectChannelClass = ClassUtils + .isPresent("org.springframework.integration.channel.DirectChannel", null); + this.directWithAttributesChannelClass = ClassUtils + .isPresent(STREAM_DIRECT_CHANNEL, null) + ? ClassUtils.resolveClassName(STREAM_DIRECT_CHANNEL, null) : null; } public static TracingChannelInterceptor create(Tracing tracing) { @@ -166,7 +170,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i MessageHeaderAccessor headers = mutableHeaderAccessor(retrievedMessage); TraceContextOrSamplingFlags extracted = this.extractor.extract(headers); Span span = this.threadLocalSpan.next(extracted); - MessageHeaderPropagation.removeAnyTraceHeaders(headers, this.tracing.propagation().keys()); + MessageHeaderPropagation.removeAnyTraceHeaders(headers, + this.tracing.propagation().keys()); this.injector.inject(span.context(), headers); if (!span.isNoop()) { span.kind(Span.Kind.PRODUCER).name("send").start(); @@ -195,19 +200,24 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i return REMOTE_SERVICE_NAME; } - private Message outputMessage(Message originalMessage, Message retrievedMessage, - MessageHeaderAccessor additionalHeaders) { - MessageHeaderAccessor headers = MessageHeaderAccessor.getMutableAccessor(originalMessage); + private Message outputMessage(Message originalMessage, + Message retrievedMessage, MessageHeaderAccessor additionalHeaders) { + MessageHeaderAccessor headers = MessageHeaderAccessor + .getMutableAccessor(originalMessage); if (originalMessage instanceof ErrorMessage) { ErrorMessage errorMessage = (ErrorMessage) originalMessage; - headers.copyHeaders(MessageHeaderPropagation.propagationHeaders(additionalHeaders.getMessageHeaders(), + headers.copyHeaders(MessageHeaderPropagation.propagationHeaders( + additionalHeaders.getMessageHeaders(), this.tracing.propagation().keys())); - return new ErrorMessage(errorMessage.getPayload(), isWebSockets(headers) ? headers.getMessageHeaders() - : new MessageHeaders(headers.getMessageHeaders()), errorMessage.getOriginalMessage()); + return new ErrorMessage(errorMessage.getPayload(), + isWebSockets(headers) ? headers.getMessageHeaders() + : new MessageHeaders(headers.getMessageHeaders()), + errorMessage.getOriginalMessage()); } headers.copyHeaders(additionalHeaders.getMessageHeaders()); return new GenericMessage<>(retrievedMessage.getPayload(), - isWebSockets(headers) ? headers.getMessageHeaders() : new MessageHeaders(headers.getMessageHeaders())); + isWebSockets(headers) ? headers.getMessageHeaders() + : new MessageHeaders(headers.getMessageHeaders())); } private boolean isWebSockets(MessageHeaderAccessor headerAccessor) { @@ -217,7 +227,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i private boolean isDirectChannel(MessageChannel channel) { Class targetClass = AopUtils.getTargetClass(channel); - boolean directChannel = this.hasDirectChannelClass && DirectChannel.class.isAssignableFrom(targetClass); + boolean directChannel = this.hasDirectChannelClass + && DirectChannel.class.isAssignableFrom(targetClass); if (!directChannel) { return false; } @@ -232,7 +243,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i } @Override - public void afterSendCompletion(Message message, MessageChannel channel, boolean sent, Exception ex) { + public void afterSendCompletion(Message message, MessageChannel channel, + boolean sent, Exception ex) { if (emptyMessage(message)) { return; } @@ -240,7 +252,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i afterMessageHandled(message, channel, null, ex); } if (log.isDebugEnabled()) { - log.debug("Will finish the current span after completion " + this.tracer.currentSpan()); + log.debug("Will finish the current span after completion " + + this.tracer.currentSpan()); } finishSpan(ex); } @@ -257,7 +270,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i MessageHeaderAccessor headers = mutableHeaderAccessor(message); TraceContextOrSamplingFlags extracted = this.extractor.extract(headers); Span span = this.threadLocalSpan.next(extracted); - MessageHeaderPropagation.removeAnyTraceHeaders(headers, this.tracing.propagation().keys()); + MessageHeaderPropagation.removeAnyTraceHeaders(headers, + this.tracing.propagation().keys()); this.injector.inject(span.context(), headers); if (!span.isNoop()) { span.kind(Span.Kind.CONSUMER).name("receive").start(); @@ -270,19 +284,21 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i headers.setImmutable(); if (message instanceof ErrorMessage) { ErrorMessage errorMessage = (ErrorMessage) message; - return new ErrorMessage(errorMessage.getPayload(), headers.getMessageHeaders(), - errorMessage.getOriginalMessage()); + return new ErrorMessage(errorMessage.getPayload(), + headers.getMessageHeaders(), errorMessage.getOriginalMessage()); } return new GenericMessage<>(message.getPayload(), headers.getMessageHeaders()); } @Override - public void afterReceiveCompletion(Message message, MessageChannel channel, Exception ex) { + public void afterReceiveCompletion(Message message, MessageChannel channel, + Exception ex) { if (emptyMessage(message)) { return; } if (log.isDebugEnabled()) { - log.debug("Will finish the current span after receive completion " + this.tracer.currentSpan()); + log.debug("Will finish the current span after receive completion " + + this.tracer.currentSpan()); } finishSpan(ex); } @@ -292,7 +308,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i * context. It then creates a span for the handler, placing it in scope. */ @Override - public Message beforeHandle(Message message, MessageChannel channel, MessageHandler handler) { + public Message beforeHandle(Message message, MessageChannel channel, + MessageHandler handler) { if (emptyMessage(message)) { return message; } @@ -307,28 +324,34 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter i consumerSpan.finish(); } // create and scope a span for the message processor - this.threadLocalSpan.next(TraceContextOrSamplingFlags.create(consumerSpan.context())).name("handle").start(); + this.threadLocalSpan + .next(TraceContextOrSamplingFlags.create(consumerSpan.context())) + .name("handle").start(); // remove any trace headers, but don't re-inject as we are synchronously // processing the // message and can rely on scoping to access this span later. - MessageHeaderPropagation.removeAnyTraceHeaders(headers, this.tracing.propagation().keys()); + MessageHeaderPropagation.removeAnyTraceHeaders(headers, + this.tracing.propagation().keys()); if (log.isDebugEnabled()) { log.debug("Created a new span in before handle" + consumerSpan); } if (message instanceof ErrorMessage) { - return new ErrorMessage((Throwable) message.getPayload(), headers.getMessageHeaders()); + return new ErrorMessage((Throwable) message.getPayload(), + headers.getMessageHeaders()); } headers.setImmutable(); return new GenericMessage<>(message.getPayload(), headers.getMessageHeaders()); } @Override - public void afterMessageHandled(Message message, MessageChannel channel, MessageHandler handler, Exception ex) { + public void afterMessageHandled(Message message, MessageChannel channel, + MessageHandler handler, Exception ex) { if (emptyMessage(message)) { return; } if (log.isDebugEnabled()) { - log.debug("Will finish the current span after message handled " + this.tracer.currentSpan()); + log.debug("Will finish the current span after message handled " + + this.tracer.currentSpan()); } finishSpan(ex); } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java index 10b5c1c91..88f334ccc 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java @@ -53,11 +53,10 @@ public class TracingChannelInterceptorTest { List spans = new ArrayList<>(); - ChannelInterceptor interceptor = TracingChannelInterceptor - .create(Tracing.newBuilder() - .currentTraceContext(ThreadLocalCurrentTraceContext.newBuilder() - .addScopeDecorator(StrictScopeDecorator.create()).build()) - .spanReporter(this.spans::add).build()); + ChannelInterceptor interceptor = TracingChannelInterceptor.create(Tracing.newBuilder() + .currentTraceContext(ThreadLocalCurrentTraceContext.newBuilder() + .addScopeDecorator(StrictScopeDecorator.create()).build()) + .spanReporter(this.spans::add).build()); QueueChannel channel = new QueueChannel(); @@ -86,9 +85,10 @@ public class TracingChannelInterceptorTest { this.channel.send(MessageBuilder.withPayload("foo").build()); - assertThat(this.channel.receive().getHeaders()).containsKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled", - "nativeHeaders"); - assertThat(this.spans).hasSize(1).flatExtracting(Span::kind).containsExactly(Span.Kind.PRODUCER); + assertThat(this.channel.receive().getHeaders()).containsKeys("X-B3-TraceId", + "X-B3-SpanId", "X-B3-Sampled", "nativeHeaders"); + assertThat(this.spans).hasSize(1).flatExtracting(Span::kind) + .containsExactly(Span.Kind.PRODUCER); } @Test @@ -98,9 +98,10 @@ public class TracingChannelInterceptorTest { this.directChannel.send(MessageBuilder.withPayload("foo").build()); assertThat(this.message).isNotNull(); - assertThat(this.message.getHeaders()).containsKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled", - "nativeHeaders"); - assertThat(this.spans).flatExtracting(Span::kind).contains(Span.Kind.CONSUMER, Span.Kind.PRODUCER); + assertThat(this.message.getHeaders()).containsKeys("X-B3-TraceId", "X-B3-SpanId", + "X-B3-Sampled", "nativeHeaders"); + assertThat(this.spans).flatExtracting(Span::kind).contains(Span.Kind.CONSUMER, + Span.Kind.PRODUCER); } @Test @@ -109,8 +110,9 @@ public class TracingChannelInterceptorTest { this.channel.send(MessageBuilder.withPayload("foo").build()); - assertThat((Map) this.channel.receive().getHeaders().get(NATIVE_HEADERS)).containsOnlyKeys("X-B3-TraceId", - "X-B3-SpanId", "X-B3-Sampled", "spanTraceId", "spanId", "spanSampled"); + assertThat((Map) this.channel.receive().getHeaders().get(NATIVE_HEADERS)) + .containsOnlyKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled", + "spanTraceId", "spanId", "spanSampled"); } /** @@ -122,11 +124,13 @@ public class TracingChannelInterceptorTest { public void producerConsidersOldSpanIds() { this.channel.addInterceptor(producerSideOnly(this.interceptor)); - this.channel.send(MessageBuilder.withPayload("foo").setHeader("X-B3-TraceId", "000000000000000a") - .setHeader("X-B3-ParentSpanId", "000000000000000a").setHeader("X-B3-SpanId", "000000000000000b") - .build()); + this.channel.send(MessageBuilder.withPayload("foo") + .setHeader("X-B3-TraceId", "000000000000000a") + .setHeader("X-B3-ParentSpanId", "000000000000000a") + .setHeader("X-B3-SpanId", "000000000000000b").build()); - assertThat(this.channel.receive().getHeaders()).containsEntry("X-B3-ParentSpanId", "000000000000000b"); + assertThat(this.channel.receive().getHeaders()).containsEntry("X-B3-ParentSpanId", + "000000000000000b"); } @Test @@ -140,10 +144,12 @@ public class TracingChannelInterceptorTest { accessor.setNativeHeader("X-B3-ParentSpanId", "000000000000000a"); accessor.setNativeHeader("X-B3-SpanId", "000000000000000b"); - this.channel.send(MessageBuilder.withPayload("foo").copyHeaders(accessor.toMessageHeaders()).build()); + this.channel.send(MessageBuilder.withPayload("foo") + .copyHeaders(accessor.toMessageHeaders()).build()); - assertThat((Map) this.channel.receive().getHeaders().get(NATIVE_HEADERS)).containsEntry("X-B3-ParentSpanId", - Collections.singletonList("000000000000000b")); + assertThat((Map) this.channel.receive().getHeaders().get(NATIVE_HEADERS)) + .containsEntry("X-B3-ParentSpanId", + Collections.singletonList("000000000000000b")); } /** @@ -156,9 +162,10 @@ public class TracingChannelInterceptorTest { this.channel.send(MessageBuilder.withPayload("foo").build()); - assertThat(this.channel.receive().getHeaders()).containsKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled", - "nativeHeaders"); - assertThat(this.spans).hasSize(1).flatExtracting(Span::kind).containsExactly(Span.Kind.CONSUMER); + assertThat(this.channel.receive().getHeaders()).containsKeys("X-B3-TraceId", + "X-B3-SpanId", "X-B3-Sampled", "nativeHeaders"); + assertThat(this.spans).hasSize(1).flatExtracting(Span::kind) + .containsExactly(Span.Kind.CONSUMER); } @Test @@ -167,8 +174,9 @@ public class TracingChannelInterceptorTest { this.channel.send(MessageBuilder.withPayload("foo").build()); - assertThat((Map) this.channel.receive().getHeaders().get(NATIVE_HEADERS)).containsOnlyKeys("X-B3-TraceId", - "X-B3-SpanId", "X-B3-Sampled", "spanTraceId", "spanId", "spanSampled"); + assertThat((Map) this.channel.receive().getHeaders().get(NATIVE_HEADERS)) + .containsOnlyKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled", + "spanTraceId", "spanId", "spanSampled"); } @Test @@ -180,9 +188,10 @@ public class TracingChannelInterceptorTest { channel.send(MessageBuilder.withPayload("foo").build()); - assertThat(messages.get(0).getHeaders()).doesNotContainKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled", - "nativeHeaders"); - assertThat(this.spans).flatExtracting(Span::kind).containsExactly(Span.Kind.CONSUMER, null); + assertThat(messages.get(0).getHeaders()).doesNotContainKeys("X-B3-TraceId", + "X-B3-SpanId", "X-B3-Sampled", "nativeHeaders"); + assertThat(this.spans).flatExtracting(Span::kind) + .containsExactly(Span.Kind.CONSUMER, null); } /** @@ -199,7 +208,8 @@ public class TracingChannelInterceptorTest { channel.send(MessageBuilder.withPayload("foo").build()); - assertThat(messages.get(0).getHeaders()).doesNotContainKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled"); + assertThat(messages.get(0).getHeaders()).doesNotContainKeys("X-B3-TraceId", + "X-B3-SpanId", "X-B3-Sampled"); } @Test @@ -211,8 +221,8 @@ public class TracingChannelInterceptorTest { channel.send(MessageBuilder.withPayload("foo").build()); - assertThat((Map) messages.get(0).getHeaders().get(NATIVE_HEADERS)).doesNotContainKeys("X-B3-TraceId", - "X-B3-SpanId", "X-B3-Sampled"); + assertThat((Map) messages.get(0).getHeaders().get(NATIVE_HEADERS)) + .doesNotContainKeys("X-B3-TraceId", "X-B3-SpanId", "X-B3-Sampled"); } @Test @@ -222,8 +232,8 @@ public class TracingChannelInterceptorTest { this.channel.send(MessageBuilder.withPayload("foo").build()); this.channel.receive(); - assertThat(this.spans).flatExtracting(Span::kind).containsExactlyInAnyOrder(Span.Kind.CONSUMER, - Span.Kind.PRODUCER); + assertThat(this.spans).flatExtracting(Span::kind) + .containsExactlyInAnyOrder(Span.Kind.CONSUMER, Span.Kind.PRODUCER); } @Test @@ -235,7 +245,8 @@ public class TracingChannelInterceptorTest { channel.send(MessageBuilder.withPayload("foo").build()); - assertThat(this.spans).flatExtracting(Span::kind).containsExactly(Span.Kind.CONSUMER, null, Span.Kind.PRODUCER); + assertThat(this.spans).flatExtracting(Span::kind) + .containsExactly(Span.Kind.CONSUMER, null, Span.Kind.PRODUCER); } @Test @@ -246,41 +257,54 @@ public class TracingChannelInterceptorTest { Map errorChannelHeaders = new HashMap<>(); errorChannelHeaders.put(MessageHeaders.REPLY_CHANNEL, errorsReplyChannel); errorChannelHeaders.put(MessageHeaders.ERROR_CHANNEL, errorsReplyChannel); - this.channel.send(new ErrorMessage( - new MessagingException(MessageBuilder.withPayload("hi") - .setHeader(TraceMessageHeaders.TRACE_ID_NAME, "000000000000000a") - .setHeader(TraceMessageHeaders.SPAN_ID_NAME, "000000000000000a") - .setReplyChannel(deadReplyChannel).setErrorChannel(deadReplyChannel).build()), - errorChannelHeaders)); + this.channel + .send(new ErrorMessage( + new MessagingException(MessageBuilder.withPayload("hi") + .setHeader(TraceMessageHeaders.TRACE_ID_NAME, + "000000000000000a") + .setHeader(TraceMessageHeaders.SPAN_ID_NAME, + "000000000000000a") + .setReplyChannel(deadReplyChannel) + .setErrorChannel(deadReplyChannel).build()), + errorChannelHeaders)); this.message = this.channel.receive(); assertThat(this.message).isNotNull(); - String spanId = this.message.getHeaders().get(TraceMessageHeaders.SPAN_ID_NAME, String.class); + String spanId = this.message.getHeaders().get(TraceMessageHeaders.SPAN_ID_NAME, + String.class); assertThat(spanId).isNotNull(); - String traceId = this.message.getHeaders().get(TraceMessageHeaders.TRACE_ID_NAME, String.class); + String traceId = this.message.getHeaders().get(TraceMessageHeaders.TRACE_ID_NAME, + String.class); assertThat(traceId).isEqualTo("000000000000000a"); assertThat(spanId).isNotEqualTo("000000000000000a"); assertThat(this.spans).hasSize(2); - assertThat(this.message.getHeaders().getReplyChannel()).isSameAs(errorsReplyChannel); - assertThat(this.message.getHeaders().getErrorChannel()).isSameAs(errorsReplyChannel); + assertThat(this.message.getHeaders().getReplyChannel()) + .isSameAs(errorsReplyChannel); + assertThat(this.message.getHeaders().getErrorChannel()) + .isSameAs(errorsReplyChannel); } @Test public void errorMessageOriginalMessageRetained() { this.channel.addInterceptor(this.interceptor); - Message originalMessage = MessageBuilder.withPayload("Hello").setHeader("header", "value").build(); - Message failedMessage = MessageBuilder.fromMessage(originalMessage).removeHeader("header").build(); - this.channel.send( - new ErrorMessage(new MessagingException(failedMessage), originalMessage.getHeaders(), originalMessage)); + Message originalMessage = MessageBuilder.withPayload("Hello") + .setHeader("header", "value").build(); + Message failedMessage = MessageBuilder.fromMessage(originalMessage) + .removeHeader("header").build(); + this.channel.send(new ErrorMessage(new MessagingException(failedMessage), + originalMessage.getHeaders(), originalMessage)); this.message = this.channel.receive(); assertThat(this.message).isNotNull(); - assertThat(this.message).isInstanceOfSatisfying(ErrorMessage.class, errorMessage -> { - assertThat(errorMessage.getOriginalMessage()).isSameAs(originalMessage); - assertThat(errorMessage.getHeaders().get("header")).isEqualTo("value"); - }); + assertThat(this.message).isInstanceOfSatisfying(ErrorMessage.class, + errorMessage -> { + assertThat(errorMessage.getOriginalMessage()) + .isSameAs(originalMessage); + assertThat(errorMessage.getHeaders().get("header")) + .isEqualTo("value"); + }); } @Test @@ -289,14 +313,17 @@ public class TracingChannelInterceptorTest { Map errorChannelHeaders = new HashMap<>(); errorChannelHeaders.put(TraceMessageHeaders.TRACE_ID_NAME, "000000000000000a"); errorChannelHeaders.put(TraceMessageHeaders.SPAN_ID_NAME, "000000000000000a"); - this.channel.send(new ErrorMessage(new MessagingException("exception"), errorChannelHeaders)); + this.channel.send(new ErrorMessage(new MessagingException("exception"), + errorChannelHeaders)); this.message = this.channel.receive(); assertThat(this.message).isNotNull(); - String spanId = this.message.getHeaders().get(TraceMessageHeaders.SPAN_ID_NAME, String.class); + String spanId = this.message.getHeaders().get(TraceMessageHeaders.SPAN_ID_NAME, + String.class); assertThat(spanId).isNotNull(); - String traceId = this.message.getHeaders().get(TraceMessageHeaders.TRACE_ID_NAME, String.class); + String traceId = this.message.getHeaders().get(TraceMessageHeaders.TRACE_ID_NAME, + String.class); assertThat(traceId).isEqualTo("000000000000000a"); assertThat(spanId).isNotEqualTo("000000000000000a"); assertThat(this.spans).hasSize(2); @@ -327,7 +354,8 @@ public class TracingChannelInterceptorTest { headers.put(AmqpHeaders.RECEIVED_ROUTING_KEY, "hello"); channel.send(MessageBuilder.createMessage("foo", new MessageHeaders(headers))); - assertThat(this.spans).flatExtracting(Span::remoteServiceName).contains("rabbitmq"); + assertThat(this.spans).flatExtracting(Span::remoteServiceName) + .contains("rabbitmq"); } @Test @@ -340,7 +368,8 @@ public class TracingChannelInterceptorTest { Map headers = new HashMap<>(); channel.send(MessageBuilder.createMessage("foo", new MessageHeaders(headers))); - assertThat(this.spans).flatExtracting(Span::remoteServiceName).containsOnly("broker", null); + assertThat(this.spans).flatExtracting(Span::remoteServiceName) + .containsOnly("broker", null); } ChannelInterceptor producerSideOnly(ChannelInterceptor delegate) { @@ -351,7 +380,8 @@ public class TracingChannelInterceptorTest { } @Override - public void afterSendCompletion(Message message, MessageChannel channel, boolean sent, Exception ex) { + public void afterSendCompletion(Message message, MessageChannel channel, + boolean sent, Exception ex) { delegate.afterSendCompletion(message, channel, sent, ex); } }; @@ -365,24 +395,29 @@ public class TracingChannelInterceptorTest { } @Override - public void afterReceiveCompletion(Message message, MessageChannel channel, Exception ex) { + public void afterReceiveCompletion(Message message, MessageChannel channel, + Exception ex) { delegate.afterReceiveCompletion(message, channel, ex); } }; } ExecutorChannelInterceptor executorSideOnly(ChannelInterceptor delegate) { - class ExecutorSideOnly extends ChannelInterceptorAdapter implements ExecutorChannelInterceptor { + class ExecutorSideOnly extends ChannelInterceptorAdapter + implements ExecutorChannelInterceptor { @Override - public Message beforeHandle(Message message, MessageChannel channel, MessageHandler handler) { - return ((ExecutorChannelInterceptor) delegate).beforeHandle(message, channel, handler); + public Message beforeHandle(Message message, MessageChannel channel, + MessageHandler handler) { + return ((ExecutorChannelInterceptor) delegate).beforeHandle(message, + channel, handler); } @Override - public void afterMessageHandled(Message message, MessageChannel channel, MessageHandler handler, - Exception ex) { - ((ExecutorChannelInterceptor) delegate).afterMessageHandled(message, channel, handler, ex); + public void afterMessageHandled(Message message, MessageChannel channel, + MessageHandler handler, Exception ex) { + ((ExecutorChannelInterceptor) delegate).afterMessageHandled(message, + channel, handler, ex); } } From 4de484504e3e02002a65b8f3e9fba3797279008e Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 6 Feb 2020 15:36:57 +0800 Subject: [PATCH 04/12] refactoring flakey test --- .../ScopePassingSpanSubscriberTests.java | 136 ++++++++++-------- 1 file changed, 75 insertions(+), 61 deletions(-) diff --git a/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java b/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java index b809a743e..ae96ea780 100644 --- a/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java +++ b/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java @@ -21,9 +21,9 @@ import java.util.function.Function; import brave.propagation.CurrentTraceContext; import brave.propagation.CurrentTraceContext.Scope; +import brave.propagation.StrictScopeDecorator; import brave.propagation.TraceContext; import org.assertj.core.presentation.StandardRepresentation; -import org.awaitility.Awaitility; import org.junit.After; import org.junit.Test; import org.reactivestreams.Publisher; @@ -54,7 +54,8 @@ public class ScopePassingSpanSubscriberTests { Objects::toString); } - final CurrentTraceContext currentTraceContext = CurrentTraceContext.Default.create(); + final CurrentTraceContext currentTraceContext = CurrentTraceContext.Default + .newBuilder().addScopeDecorator(StrictScopeDecorator.create()).build(); TraceContext context = TraceContext.newBuilder().traceId(1).spanId(1).sampled(true) .build(); @@ -62,6 +63,52 @@ public class ScopePassingSpanSubscriberTests { TraceContext context2 = TraceContext.newBuilder().traceId(1).spanId(2).sampled(true) .build(); + Subscriber assertNotScopePassingSpanSubscriber = new CoreSubscriber() { + @Override + public void onSubscribe(Subscription s) { + s.request(Long.MAX_VALUE); + assertThat(s).isNotInstanceOf(ScopePassingSpanSubscriber.class); + } + + @Override + public void onNext(Object o) { + + } + + @Override + public void onError(Throwable t) { + + } + + @Override + public void onComplete() { + + } + }; + + Subscriber assertScopePassingSpanSubscriber = new CoreSubscriber() { + @Override + public void onSubscribe(Subscription s) { + s.request(Long.MAX_VALUE); + assertThat(s).isInstanceOf(ScopePassingSpanSubscriber.class); + } + + @Override + public void onNext(Object o) { + + } + + @Override + public void onError(Throwable t) { + + } + + @Override + public void onComplete() { + + } + }; + AnnotationConfigApplicationContext springContext = new AnnotationConfigApplicationContext(); @After @@ -97,7 +144,7 @@ public class ScopePassingSpanSubscriberTests { } @Test - public void should_not_trace_scalar_flows() { + public void should_not_scope_scalar_subscribe() { springContext.registerBean(CurrentTraceContext.class, () -> currentTraceContext); springContext.refresh(); @@ -105,71 +152,38 @@ public class ScopePassingSpanSubscriberTests { this.springContext); try (Scope ws = this.currentTraceContext.newScope(context)) { - Subscriber assertNoSpanSubscriber = new CoreSubscriber() { - @Override - public void onSubscribe(Subscription s) { - s.request(Long.MAX_VALUE); - assertThat(s).isNotInstanceOf(ScopePassingSpanSubscriber.class); - } - @Override - public void onNext(Object o) { - - } - - @Override - public void onError(Throwable t) { - - } - - @Override - public void onComplete() { - - } - }; - - Subscriber assertSpanSubscriber = new CoreSubscriber() { - @Override - public void onSubscribe(Subscription s) { - s.request(Long.MAX_VALUE); - assertThat(s).isInstanceOf(ScopePassingSpanSubscriber.class); - } - - @Override - public void onNext(Object o) { - - } - - @Override - public void onError(Throwable t) { - - } - - @Override - public void onComplete() { - - } - }; - transformer.apply(Mono.just(1).hide()).subscribe(assertSpanSubscriber); - - transformer.apply(Mono.just(1)).subscribe(assertNoSpanSubscriber); - - transformer.apply(Mono.error(new Exception()).hide()) - .subscribe(assertSpanSubscriber); + transformer.apply(Mono.just(1)) + .subscribe(assertNotScopePassingSpanSubscriber); transformer.apply(Mono.error(new Exception())) - .subscribe(assertNoSpanSubscriber); + .subscribe(assertNotScopePassingSpanSubscriber); - transformer.apply(Mono.empty().hide()) - .subscribe(assertSpanSubscriber); - - transformer.apply(Mono.empty()).subscribe(assertNoSpanSubscriber); + transformer.apply(Mono.empty()) + .subscribe(assertNotScopePassingSpanSubscriber); } + } - Awaitility.await().untilAsserted(() -> { - then(this.currentTraceContext.get()).isNull(); - }); + @Test + public void should_scope_scalar_hide_subscribe() { + springContext.registerBean(CurrentTraceContext.class, () -> currentTraceContext); + springContext.refresh(); + + Function, ? extends Publisher> transformer = scopePassingSpanOperator( + this.springContext); + + try (Scope ws = this.currentTraceContext.newScope(context)) { + + transformer.apply(Mono.just(1).hide()) + .subscribe(assertScopePassingSpanSubscriber); + + transformer.apply(Mono.error(new Exception()).hide()) + .subscribe(assertScopePassingSpanSubscriber); + + transformer.apply(Mono.empty().hide()) + .subscribe(assertScopePassingSpanSubscriber); + } } } From 905ef2cc6f38c3b64f6b16afd4ff44409b1c09bf Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Fri, 7 Feb 2020 16:14:39 +0800 Subject: [PATCH 05/12] yolo attempt to deflake build --- .../reactor/ScopePassingSpanSubscriberTests.java | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java b/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java index ae96ea780..13427635a 100644 --- a/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java +++ b/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java @@ -25,13 +25,16 @@ import brave.propagation.StrictScopeDecorator; import brave.propagation.TraceContext; import org.assertj.core.presentation.StandardRepresentation; import org.junit.After; +import org.junit.Before; import org.junit.Test; import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import org.reactivestreams.Subscription; import reactor.core.CoreSubscriber; import reactor.core.publisher.BaseSubscriber; +import reactor.core.publisher.Hooks; import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; import reactor.util.context.Context; import org.springframework.context.annotation.AnnotationConfigApplicationContext; @@ -39,6 +42,8 @@ import org.springframework.context.annotation.AnnotationConfigApplicationContext import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.BDDAssertions.then; import static org.springframework.cloud.sleuth.instrument.reactor.ReactorSleuth.scopePassingSpanOperator; +import static org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfiguration.SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY; +import static org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfiguration.TraceReactorConfiguration.SLEUTH_TRACE_REACTOR_KEY; /** * @author Marcin Grzejszczak @@ -111,6 +116,15 @@ public class ScopePassingSpanSubscriberTests { AnnotationConfigApplicationContext springContext = new AnnotationConfigApplicationContext(); + @Before + public void resetHooks() { + // There's an assumption some other test is leaking hooks, so we clear them all to + // prevent should_not_scope_scalar_subscribe from being interfered with. + Hooks.resetOnEachOperator(SLEUTH_TRACE_REACTOR_KEY); + Hooks.resetOnLastOperator(SLEUTH_TRACE_REACTOR_KEY); + Schedulers.removeExecutorServiceDecorator(SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY); + } + @After public void close() { springContext.close(); From f53101f5b87d9a195717f6fd5a82b41952fefad7 Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Sat, 8 Feb 2020 08:56:51 +0800 Subject: [PATCH 06/12] Fixes propagation break and error handling in reactor-netty HttpClient (#1552) --- .../client/HttpClientBeanPostProcessor.java | 87 +++++++++---------- ...ReactorNettyHttpClientSpringBootTests.java | 19 ++++ 2 files changed, 62 insertions(+), 44 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java index 5558b5e99..7f90d40ed 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java @@ -31,11 +31,13 @@ import reactor.netty.Connection; import reactor.netty.http.client.HttpClient; import reactor.netty.http.client.HttpClientRequest; import reactor.netty.http.client.HttpClientResponse; +import reactor.util.context.Context; import org.springframework.beans.BeansException; import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.cloud.sleuth.internal.LazyBean; import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.lang.Nullable; class HttpClientBeanPostProcessor implements BeanPostProcessor { @@ -51,11 +53,15 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { LazyBean httpTracing = LazyBean.create(this.springContext, HttpTracing.class); if (bean instanceof HttpClient) { - return ((HttpClient) bean).mapConnect(new TracingMapConnect(httpTracing)) - .doOnRequest(TracingDoOnRequest.create(httpTracing)) - .doOnRequestError(TracingDoOnErrorRequest.create(httpTracing)) - .doOnResponse(TracingDoOnResponse.create(httpTracing)) - .doOnResponseError(TracingDoOnErrorResponse.create(httpTracing)); + // This adds handlers to manage the span lifecycle. All require explicit + // propagation of the current span as a reactor context property. + // This done in mapConnect, added last so that it is setup first. + // https://projectreactor.io/docs/core/release/reference/#_simple_context_examples + return ((HttpClient) bean).doOnRequest(new TracingDoOnRequest(httpTracing)) + .doOnRequestError(new TracingDoOnErrorRequest(httpTracing)) + .doOnResponse(new TracingDoOnResponse(httpTracing)) + .doOnResponseError(new TracingDoOnErrorResponse(httpTracing)) + .mapConnect(new TracingMapConnect(httpTracing)); } return bean; } @@ -74,11 +80,12 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { @Override public Mono apply(Mono mono, Bootstrap bootstrap) { + // This is read in this class and also inside ScopePassingSpanSubscriber return mono.subscriberContext(context -> context.put(AtomicReference.class, new AtomicReference<>(tracer().currentSpan()))); } - private Tracer tracer() { + Tracer tracer() { if (this.tracer == null) { this.tracer = this.httpTracing.get().tracing().tracer(); } @@ -100,10 +107,6 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { this.httpTracing = httpTracing; } - static TracingDoOnRequest create(LazyBean httpTracing) { - return new TracingDoOnRequest(httpTracing); - } - List propagationKeys() { if (this.propagationKeys == null) { this.propagationKeys = httpTracing.get().tracing().propagation().keys(); @@ -128,12 +131,21 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { return; } } - AtomicReference reference = req.currentContext() - .getOrDefault(AtomicReference.class, new AtomicReference<>()); + + // Look for a parent propagated by TracingMapConnect + AtomicReference ref = req.currentContext() + .getOrDefault(AtomicReference.class, null); + Span parent = ref != null ? ref.get() : null; + + // Start a new client span with the appropriate parent WrappedHttpClientRequest request = new WrappedHttpClientRequest(req); - Span span = reference.get() == null ? handler().handleSend(request) - : handler().handleSend(request, reference.get()); - reference.set(span); + Span clientSpan = parent != null ? handler().handleSend(request, parent) + : handler().handleSend(request); + + // Swap the ref with the client span, so that other hooks can see it + if (ref != null) { + ref.set(clientSpan); + } } } @@ -145,13 +157,9 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { super(httpTracing); } - static TracingDoOnResponse create(LazyBean httpTracing) { - return new TracingDoOnResponse(httpTracing); - } - @Override - public void accept(HttpClientResponse httpClientResponse, Connection connection) { - handle(httpClientResponse, null); + public void accept(HttpClientResponse response, Connection connection) { + handle(response.currentContext(), response, null); } } @@ -163,13 +171,10 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { super(httpTracing); } - static TracingDoOnErrorRequest create(LazyBean httpTracing) { - return new TracingDoOnErrorRequest(httpTracing); - } - @Override - public void accept(HttpClientRequest request, Throwable throwable) { - handle(null, throwable); + public void accept(HttpClientRequest request, Throwable error) { + // TODO: the current context here does not have the AtomicReference + handle(request.currentContext(), null, error); } } @@ -181,13 +186,9 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { super(httpTracing); } - static TracingDoOnErrorResponse create(LazyBean httpTracing) { - return new TracingDoOnErrorResponse(httpTracing); - } - @Override - public void accept(HttpClientResponse httpClientResponse, Throwable throwable) { - handle(httpClientResponse, throwable); + public void accept(HttpClientResponse response, Throwable error) { + handle(response.currentContext(), response, error); } } @@ -209,18 +210,16 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { return this.handler; } - protected void handle(HttpClientResponse httpClientResponse, - Throwable throwable) { - if (httpClientResponse == null) { - return; + void handle(Context context, @Nullable HttpClientResponse resp, + @Nullable Throwable error) { + AtomicReference ref = context.getOrDefault(AtomicReference.class, null); + Span span = ref != null ? ref.get() : null; + if (span == null) { + return; // Unexpected. In the handle method, without a span to finish! } - AtomicReference reference = httpClientResponse.currentContext() - .getOrDefault(AtomicReference.class, null); - if (reference == null || reference.get() == null) { - return; - } - handler().handleReceive(new WrappedHttpClientResponse(httpClientResponse), - throwable, (Span) reference.get()); + WrappedHttpClientResponse response = resp != null + ? new WrappedHttpClientResponse(resp) : null; + handler().handleReceive(response, error, span); } } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java index 30b1dc1c8..518b61fb5 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java @@ -30,6 +30,7 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.netty.DisposableServer; import reactor.netty.http.client.HttpClient; +import reactor.netty.http.client.PrematureCloseException; import reactor.netty.http.server.HttpServer; import zipkin2.Span; import zipkin2.reporter.Reporter; @@ -43,6 +44,7 @@ import org.springframework.test.context.junit4.SpringRunner; import org.springframework.web.reactive.function.client.WebClient; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; /** * This tests {@link HttpClient} instrumentation performed by @@ -93,6 +95,23 @@ public class ReactorNettyHttpClientSpringBootTests { .isEqualTo(clientSpan.traceId() + "-" + clientSpan.id() + "-1"); } + @Test + public void shouldTagOnRequestError() throws InterruptedException { + disposableServer = HttpServer.create().port(0).handle((req, resp) -> { + throw new RuntimeException("test"); + }).bindNow(); + + Mono request = httpClient.port(disposableServer.port()).get().uri("/") + .responseContent().aggregate().asString(); + + assertThatThrownBy(request::block) + .hasCauseInstanceOf(PrematureCloseException.class); + + Span clientSpan = takeClientSpan(); + + assertThat(clientSpan.tags()).containsKey("error"); + } + /** Call this to block until a span was reported */ Span takeClientSpan() throws InterruptedException { Span result = spans.poll(1, TimeUnit.SECONDS); From 236af9c70fb158174056638d9a92866bf3c6f762 Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 13 Feb 2020 12:06:04 -0800 Subject: [PATCH 07/12] bumps to latest brave --- benchmarks/pom.xml | 2 +- pom.xml | 2 +- spring-cloud-sleuth-dependencies/pom.xml | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/benchmarks/pom.xml b/benchmarks/pom.xml index 3c1f204f4..d07959a23 100644 --- a/benchmarks/pom.xml +++ b/benchmarks/pom.xml @@ -33,7 +33,7 @@ 1.8 1.8 2.3.0.BUILD-SNAPSHOT - 5.9.3 + 5.9.5 3.14.6 diff --git a/pom.xml b/pom.xml index 38a06d5d7..207d75c98 100644 --- a/pom.xml +++ b/pom.xml @@ -257,7 +257,7 @@ Horsham.SR1 2.2.2.BUILD-SNAPSHOT 2.2.2.BUILD-SNAPSHOT - 5.9.3 + 5.9.5 2.1.7.RELEASE 2.2.2.BUILD-SNAPSHOT false diff --git a/spring-cloud-sleuth-dependencies/pom.xml b/spring-cloud-sleuth-dependencies/pom.xml index acb206410..11cd57f54 100644 --- a/spring-cloud-sleuth-dependencies/pom.xml +++ b/spring-cloud-sleuth-dependencies/pom.xml @@ -31,7 +31,7 @@ spring-cloud-sleuth-dependencies Spring Cloud Sleuth Dependencies - 5.9.3 + 5.9.5 0.35.1 3.4.1 From 707abda92a7ced961e4fcde56f23cb02c755f15e Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 13 Feb 2020 14:12:06 -0800 Subject: [PATCH 08/12] fixes curse of the feign mocks --- .../web/client/feign/TracingFeignClient.java | 6 ++- .../client/feign/TracingFeignClientTests.java | 42 +++++++------------ 2 files changed, 20 insertions(+), 28 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClient.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClient.java index 662314df6..49914142b 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClient.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClient.java @@ -95,9 +95,13 @@ final class TracingFeignClient implements Client { Throwable error = null; try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) { Response res = this.delegate.execute(request.build(), options); - if (res != null) { // possibly null on bad implementation or mocks + if (res != null) { response = new HttpClientResponse(res); } + else { // possibly null on bad implementation or mocks + response = new HttpClientResponse( + Response.builder().request(req).build()); + } return res; } catch (IOException | RuntimeException | Error e) { diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClientTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClientTests.java index d195b7bd0..53dc877cd 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClientTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/TracingFeignClientTests.java @@ -17,8 +17,9 @@ package org.springframework.cloud.sleuth.instrument.web.client.feign; import java.io.IOException; -import java.nio.charset.Charset; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import brave.Span; import brave.Tracer; @@ -36,9 +37,7 @@ import org.mockito.BDDMockito; import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; -import org.springframework.beans.factory.BeanFactory; import org.springframework.cloud.sleuth.instrument.web.SleuthHttpParserAccessor; -import org.springframework.cloud.sleuth.util.ArrayListSpanReporter; import static org.assertj.core.api.BDDAssertions.then; @@ -48,15 +47,16 @@ import static org.assertj.core.api.BDDAssertions.then; @RunWith(MockitoJUnitRunner.class) public class TracingFeignClientTests { - ArrayListSpanReporter reporter = new ArrayListSpanReporter(); + Request request = Request.create("GET", "https://foo", new HashMap<>(), null, null); - @Mock - BeanFactory beanFactory; + Request.Options options = new Request.Options(); + + List spans = new ArrayList<>(); Tracing tracing = Tracing.newBuilder() .currentTraceContext(ThreadLocalCurrentTraceContext.newBuilder() .addScopeDecorator(StrictScopeDecorator.create()).build()) - .spanReporter(this.reporter).build(); + .spanReporter(spans::add).build(); Tracer tracer = this.tracing.tracer(); @@ -78,17 +78,13 @@ public class TracingFeignClientTests { Span span = this.tracer.nextSpan().name("foo"); try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span.start())) { - this.traceFeignClient - .execute( - Request.create("GET", "https://foo", new HashMap<>(), - "".getBytes(), Charset.defaultCharset()), - new Request.Options()); + this.traceFeignClient.execute(this.request, this.options); } finally { span.finish(); } - then(this.reporter.getSpans().get(0)).extracting("kind.ordinal") + then(spans.get(0)).extracting("kind.ordinal") .isEqualTo(Span.Kind.CLIENT.ordinal()); } @@ -99,11 +95,7 @@ public class TracingFeignClientTests { .willThrow(new RuntimeException("exception has occurred")); try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span.start())) { - this.traceFeignClient - .execute( - Request.create("GET", "https://foo", new HashMap<>(), - "".getBytes(), Charset.defaultCharset()), - new Request.Options()); + this.traceFeignClient.execute(this.request, this.options); BDDAssertions.fail("Exception should have been thrown"); } catch (Exception e) { @@ -112,21 +104,17 @@ public class TracingFeignClientTests { span.finish(); } - then(this.reporter.getSpans().get(0)).extracting("kind.ordinal") + then(this.spans.get(0)).extracting("kind.ordinal") .isEqualTo(Span.Kind.CLIENT.ordinal()); - then(this.reporter.getSpans().get(0).tags()).containsEntry("error", - "exception has occurred"); + then(this.spans.get(0).tags()).containsEntry("error", "exception has occurred"); } @Test public void should_shorten_the_span_name() throws IOException { - this.traceFeignClient - .execute( - Request.create("GET", "https://foo/" + bigName(), new HashMap<>(), - "".getBytes(), Charset.defaultCharset()), - new Request.Options()); + this.traceFeignClient.execute(Request.create("GET", "https://foo/" + bigName(), + new HashMap<>(), null, null), this.options); - then(this.reporter.getSpans().get(0).name()).hasSize(50); + then(this.spans.get(0).name()).hasSize(50); } private String bigName() { From 47416a8eefb6f0f5586b0042e224bd4f952f8559 Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 13 Feb 2020 14:56:47 -0800 Subject: [PATCH 09/12] Reduces scope copies when desired context is already current (#1558) --- .../reactor/ScopePassingSpanSubscriber.java | 4 +++- .../reactor/ScopePassingSpanSubscriberTests.java | 12 ++++++++++++ 2 files changed, 15 insertions(+), 1 deletion(-) 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 2e24dc278..a1b587052 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 @@ -54,7 +54,9 @@ final class ScopePassingSpanSubscriber implements SpanSubscription, Scanna this.subscriber = subscriber; this.currentTraceContext = currentTraceContext; this.parent = parent; - this.context = parent != null ? ctx.put(TraceContext.class, parent) : ctx; + this.context = parent != null + && !parent.equals(ctx.getOrDefault(TraceContext.class, null)) + ? ctx.put(TraceContext.class, parent) : ctx; if (log.isTraceEnabled()) { log.trace("Parent span [" + parent + "], context [" + this.context + "]"); } diff --git a/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java b/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java index 13427635a..65023e5da 100644 --- a/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java +++ b/tests/spring-cloud-sleuth-instrumentation-reactor-tests/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/ScopePassingSpanSubscriberTests.java @@ -138,6 +138,18 @@ public class ScopePassingSpanSubscriberTests { then((String) subscriber.currentContext().get("foo")).isEqualTo("bar"); } + /** + * This ensures when the desired context is in the reactor context we don't copy it. + */ + @Test + public void should_not_redundantly_copy_context() { + Context initial = Context.of(TraceContext.class, context); + ScopePassingSpanSubscriber subscriber = new ScopePassingSpanSubscriber<>(null, + initial, this.currentTraceContext, context); + + then(initial).isSameAs(subscriber.currentContext()); + } + @Test public void should_set_empty_context_when_context_is_null() { ScopePassingSpanSubscriber subscriber = new ScopePassingSpanSubscriber<>(null, From 0ecb23ee88acc5dc85975a5f2f32e333a668c57f Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 13 Feb 2020 15:28:31 -0800 Subject: [PATCH 10/12] Fixes parse glitches and adds remote endpoint to reactor HttpClient (#1559) --- .../client/HttpClientBeanPostProcessor.java | 19 ++++++++++++++++-- ...ReactorNettyHttpClientSpringBootTests.java | 20 +++++++++++++++++++ 2 files changed, 37 insertions(+), 2 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java index 7f90d40ed..052cf5467 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java @@ -16,6 +16,7 @@ package org.springframework.cloud.sleuth.instrument.web.client; +import java.net.InetSocketAddress; import java.util.List; import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; @@ -141,6 +142,7 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { WrappedHttpClientRequest request = new WrappedHttpClientRequest(req); Span clientSpan = parent != null ? handler().handleSend(request, parent) : handler().handleSend(request); + parseConnectionAddress(connection, clientSpan); // Swap the ref with the client span, so that other hooks can see it if (ref != null) { @@ -148,6 +150,14 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { } } + static void parseConnectionAddress(Connection connection, Span span) { + if (span.isNoop()) { + return; + } + InetSocketAddress socketAddress = connection.address(); + span.remoteIpAndPort(socketAddress.getHostString(), socketAddress.getPort()); + } + } private static class TracingDoOnResponse extends AbstractTracingDoOnHandler @@ -244,12 +254,12 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { @Override public String path() { - return delegate.path(); + return "/" + delegate.path(); // TODO: reactor/reactor-netty#999 } @Override public String url() { - return delegate.uri(); + return delegate.resourceUrl(); } @Override @@ -272,6 +282,11 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { this.delegate = delegate; } + @Override + public String method() { + return delegate.method().name(); + } + @Override public Object unwrap() { return delegate; diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java index 518b61fb5..791146de6 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java @@ -23,6 +23,7 @@ import java.util.concurrent.TimeUnit; import brave.propagation.B3SinglePropagation; import brave.propagation.Propagation; import brave.sampler.Sampler; +import io.netty.handler.codec.http.HttpResponseStatus; import org.junit.After; import org.junit.Test; import org.junit.runner.RunWith; @@ -30,6 +31,7 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.netty.DisposableServer; import reactor.netty.http.client.HttpClient; +import reactor.netty.http.client.HttpClientResponse; import reactor.netty.http.client.PrematureCloseException; import reactor.netty.http.server.HttpServer; import zipkin2.Span; @@ -76,6 +78,24 @@ public class ReactorNettyHttpClientSpringBootTests { this.spans.clear(); } + @Test + public void shouldRecordRemoteEndpoint() throws Exception { + disposableServer = HttpServer.create().port(0) + .handle((in, out) -> out.sendString(Flux.just("foo"))).bindNow(); + + HttpClientResponse response = httpClient.port(disposableServer.port()).get() + .uri("/").response().block(); + + assertThat(response.status()).isEqualTo(HttpResponseStatus.OK); + + Span clientSpan = takeClientSpan(); + + assertThat(clientSpan.remoteEndpoint()).satisfiesAnyOf( + ep -> assertThat(ep.ipv4()).isNotNull(), + ep -> assertThat(ep.ipv6()).isNotNull()); + assertThat(clientSpan.remoteEndpoint().portAsInt()).isNotZero(); + } + @Test public void shouldSendTraceContextToServer_rootSpan() throws Exception { disposableServer = HttpServer.create().port(0) From 20ca616b69cc81fcd3e0ae6905215afef1868b45 Mon Sep 17 00:00:00 2001 From: Adrian Cole Date: Thu, 13 Feb 2020 16:15:31 -0800 Subject: [PATCH 11/12] Fixes broken propagation in reactor HttpClient (#1560) --- .../client/HttpClientBeanPostProcessor.java | 106 ++++++++++-------- ...ReactorNettyHttpClientSpringBootTests.java | 29 +++++ 2 files changed, 90 insertions(+), 45 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java index 052cf5467..a358784da 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/client/HttpClientBeanPostProcessor.java @@ -17,15 +17,16 @@ package org.springframework.cloud.sleuth.instrument.web.client; import java.net.InetSocketAddress; -import java.util.List; import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; import java.util.function.BiFunction; import brave.Span; -import brave.Tracer; import brave.http.HttpClientHandler; import brave.http.HttpTracing; +import brave.propagation.CurrentTraceContext; +import brave.propagation.CurrentTraceContext.Scope; +import brave.propagation.TraceContext; import io.netty.bootstrap.Bootstrap; import reactor.core.publisher.Mono; import reactor.netty.Connection; @@ -58,21 +59,27 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { // propagation of the current span as a reactor context property. // This done in mapConnect, added last so that it is setup first. // https://projectreactor.io/docs/core/release/reference/#_simple_context_examples - return ((HttpClient) bean).doOnRequest(new TracingDoOnRequest(httpTracing)) - .doOnRequestError(new TracingDoOnErrorRequest(httpTracing)) - .doOnResponse(new TracingDoOnResponse(httpTracing)) + return ((HttpClient) bean) .doOnResponseError(new TracingDoOnErrorResponse(httpTracing)) + .doOnResponse(new TracingDoOnResponse(httpTracing)) + .doOnRequestError(new TracingDoOnErrorRequest(httpTracing)) + .doOnRequest(new TracingDoOnRequest(httpTracing)) .mapConnect(new TracingMapConnect(httpTracing)); } return bean; } + /** current client span, cleared on completion. */ + private static final class CurrentClientSpan extends AtomicReference { + + } + private static class TracingMapConnect implements BiFunction, Bootstrap, Mono> { final LazyBean httpTracing; - Tracer tracer; + CurrentTraceContext currentTraceContext; TracingMapConnect(LazyBean httpTracing) { this.httpTracing = httpTracing; @@ -81,16 +88,22 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { @Override public Mono apply(Mono mono, Bootstrap bootstrap) { - // This is read in this class and also inside ScopePassingSpanSubscriber - return mono.subscriberContext(context -> context.put(AtomicReference.class, - new AtomicReference<>(tracer().currentSpan()))); + return mono.subscriberContext(context -> { + TraceContext invocationContext = currentTraceContext().get(); + if (invocationContext != null) { + // Read in this processor and also in ScopePassingSpanSubscriber + context = context.put(TraceContext.class, invocationContext); + } + return context.put(CurrentClientSpan.class, new CurrentClientSpan()); + }); } - Tracer tracer() { - if (this.tracer == null) { - this.tracer = this.httpTracing.get().tracing().tracer(); + CurrentTraceContext currentTraceContext() { + if (this.currentTraceContext == null) { + this.currentTraceContext = this.httpTracing.get().tracing() + .currentTraceContext(); } - return this.tracer; + return this.currentTraceContext; } } @@ -100,21 +113,12 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { final LazyBean httpTracing; - List propagationKeys; - HttpClientHandler handler; TracingDoOnRequest(LazyBean httpTracing) { this.httpTracing = httpTracing; } - List propagationKeys() { - if (this.propagationKeys == null) { - this.propagationKeys = httpTracing.get().tracing().propagation().keys(); - } - return this.propagationKeys; - } - HttpClientHandler handler() { if (this.handler == null) { this.handler = HttpClientHandler.create(httpTracing.get()); @@ -122,30 +126,39 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { return this.handler; } + CurrentTraceContext currentTraceContext() { + return httpTracing.get().tracing().currentTraceContext(); + } + @Override public void accept(HttpClientRequest req, Connection connection) { - // request already instrumented - // TODO: consider another, cheaper way, like flagging a context - // property. If not, comment why. - for (String key : propagationKeys()) { - if (req.requestHeaders().contains(key)) { - return; - } + CurrentClientSpan ref = req.currentContext() + .getOrDefault(CurrentClientSpan.class, null); + if (ref == null) { // Somehow TracingMapConnect was not invoked.. skip out + return; } - // Look for a parent propagated by TracingMapConnect - AtomicReference ref = req.currentContext() - .getOrDefault(AtomicReference.class, null); - Span parent = ref != null ? ref.get() : null; + // This might be re-entrant on auto-redirect or connection retry: + // See reactor/reactor-netty#1000 for follow-ups. + Span clientSpan = ref.getAndSet(null); + if (clientSpan != null) { + // Retry from a connect fail wouldn't have parsed the request, leading to + // an empty span with no data if we finished it. An auto-redirect would + // have parsed the request, but we have no idea which status code it + // finished with. Since we can't see the preceding request state, we + // abandon its span in favor of the next. + clientSpan.abandon(); + } // Start a new client span with the appropriate parent + TraceContext parent = req.currentContext().getOrDefault(TraceContext.class, + null); WrappedHttpClientRequest request = new WrappedHttpClientRequest(req); - Span clientSpan = parent != null ? handler().handleSend(request, parent) - : handler().handleSend(request); - parseConnectionAddress(connection, clientSpan); - // Swap the ref with the client span, so that other hooks can see it - if (ref != null) { + // Simplify after openzipkin/brave#1082 + try (Scope ws = currentTraceContext().maybeScope(parent)) { + clientSpan = handler().handleSend(request); + parseConnectionAddress(connection, clientSpan); ref.set(clientSpan); } } @@ -182,9 +195,8 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { } @Override - public void accept(HttpClientRequest request, Throwable error) { - // TODO: the current context here does not have the AtomicReference - handle(request.currentContext(), null, error); + public void accept(HttpClientRequest req, Throwable error) { + handle(req.currentContext(), null, error); } } @@ -222,14 +234,18 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { void handle(Context context, @Nullable HttpClientResponse resp, @Nullable Throwable error) { - AtomicReference ref = context.getOrDefault(AtomicReference.class, null); - Span span = ref != null ? ref.get() : null; - if (span == null) { + CurrentClientSpan ref = context.getOrDefault(CurrentClientSpan.class, null); + if (ref == null) { // Somehow TracingMapConnect was not invoked.. skip out + return; + } + + Span clientSpan = ref.getAndSet(null); + if (clientSpan == null) { return; // Unexpected. In the handle method, without a span to finish! } WrappedHttpClientResponse response = resp != null ? new WrappedHttpClientResponse(resp) : null; - handler().handleReceive(response, error, span); + handler().handleReceive(response, error, clientSpan); } } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java index 791146de6..efecdd11b 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/ReactorNettyHttpClientSpringBootTests.java @@ -21,7 +21,10 @@ import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import brave.propagation.B3SinglePropagation; +import brave.propagation.CurrentTraceContext; +import brave.propagation.CurrentTraceContext.Scope; import brave.propagation.Propagation; +import brave.propagation.TraceContext; import brave.sampler.Sampler; import io.netty.handler.codec.http.HttpResponseStatus; import org.junit.After; @@ -70,6 +73,12 @@ public class ReactorNettyHttpClientSpringBootTests { @Autowired BlockingQueue spans; + @Autowired + CurrentTraceContext currentTraceContext; + + TraceContext context = TraceContext.newBuilder().traceId(1).spanId(1).sampled(true) + .build(); + @After public void tearDown() { if (disposableServer != null) { @@ -96,6 +105,26 @@ public class ReactorNettyHttpClientSpringBootTests { assertThat(clientSpan.remoteEndpoint().portAsInt()).isNotZero(); } + @Test + public void shouldUseInvocationContext() throws Exception { + disposableServer = HttpServer.create().port(0) + // this reads the trace context header, b3, returning it in the response + .handle((in, out) -> out + .sendString(Flux.just(in.requestHeaders().get("b3")))) + .bindNow(); + + String b3SingleHeaderReadByServer; + try (Scope ws = currentTraceContext.newScope(context)) { + b3SingleHeaderReadByServer = httpClient.port(disposableServer.port()).get() + .uri("/").responseContent().aggregate().asString().block(); + } + + Span clientSpan = takeClientSpan(); + + assertThat(b3SingleHeaderReadByServer).isEqualTo(context.traceIdString() + "-" + + clientSpan.id() + "-1-" + context.spanIdString()); + } + @Test public void shouldSendTraceContextToServer_rootSpan() throws Exception { disposableServer = HttpServer.create().port(0) From eff71d4f2291f650ca8ff22d99b54e6587fc82e2 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Tue, 18 Feb 2020 12:22:13 +0100 Subject: [PATCH 12/12] Moved more beans to the conditional on backward compatibility autoconfig; fixes gh-1555 --- ...ckwardsCompatibilityAutoConfiguration.java | 80 ++++++++++--------- ...dsCompatibilityAutoConfigurationTests.java | 4 +- 2 files changed, 46 insertions(+), 38 deletions(-) diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfiguration.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfiguration.java index e58326656..38aaf8bc8 100644 --- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfiguration.java +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfiguration.java @@ -65,32 +65,7 @@ import org.springframework.util.Assert; public class ZipkinBackwardsCompatibilityAutoConfiguration { /** - * Reporter that is depending on a {@link Sender} bean which is created in another - * auto-configuration than {@link ZipkinAutoConfiguration}. - * @param reporterMetrics metrics - * @param zipkin zipkin properties - * @param spanBytesEncoder encoder - * @param beanFactory Spring's Bean Factory - * @return span reporter - * @deprecated left for backwards compatibility - */ - @Bean - @Conditional(BackwardsCompatibilityCondition.class) - @Deprecated - Reporter reporter(ReporterMetrics reporterMetrics, ZipkinProperties zipkin, - BytesEncoder spanBytesEncoder, DefaultListableBeanFactory beanFactory) { - List beanNames = new ArrayList<>( - Arrays.asList(beanFactory.getBeanNamesForType(Sender.class))); - beanNames.remove(ZipkinAutoConfiguration.SENDER_BEAN_NAME); - Sender sender = (Sender) beanFactory.getBean(beanNames.get(0)); - // historical constraint. Note: AsyncReporter supports memory bounds - return AsyncReporter.builder(sender).queuedMaxSpans(1000) - .messageTimeout(zipkin.getMessageTimeout(), TimeUnit.SECONDS) - .metrics(reporterMetrics).build(spanBytesEncoder); - } - - /** - * Only used for creating a reporter bean with the method above. + * Only used for creating a reporter bean with the method below. * @param zipkinProperties zipkin properties * @return bytes encoder * @deprecated left for backwards compatibility @@ -102,17 +77,48 @@ public class ZipkinBackwardsCompatibilityAutoConfiguration { return zipkinProperties.getEncoder(); } - /** - * Deprecated because this is moved to {@link TraceAutoConfiguration}. Left for - * backwards compatibility reasons. - * @return reporter metrics - * @deprecated left for backwards compatibility - */ - @Bean - @ConditionalOnMissingBean - @Deprecated - ReporterMetrics zipkinReporterMetrics() { - return new InMemoryReporterMetrics(); + @Configuration(proxyBeanMethods = false) + @Conditional(BackwardsCompatibilityCondition.class) + static class BackwardsCompatibilityConfiguration { + + /** + * Reporter that is depending on a {@link Sender} bean which is created in another + * auto-configuration than {@link ZipkinAutoConfiguration}. + * @param reporterMetrics metrics + * @param zipkin zipkin properties + * @param spanBytesEncoder encoder + * @param beanFactory Spring's Bean Factory + * @return span reporter + * @deprecated left for backwards compatibility + */ + @Bean + @Deprecated + Reporter reporter(ReporterMetrics reporterMetrics, ZipkinProperties zipkin, + BytesEncoder spanBytesEncoder, + DefaultListableBeanFactory beanFactory) { + List beanNames = new ArrayList<>( + Arrays.asList(beanFactory.getBeanNamesForType(Sender.class))); + beanNames.remove(ZipkinAutoConfiguration.SENDER_BEAN_NAME); + Sender sender = (Sender) beanFactory.getBean(beanNames.get(0)); + // historical constraint. Note: AsyncReporter supports memory bounds + return AsyncReporter.builder(sender).queuedMaxSpans(1000) + .messageTimeout(zipkin.getMessageTimeout(), TimeUnit.SECONDS) + .metrics(reporterMetrics).build(spanBytesEncoder); + } + + /** + * Deprecated because this is moved to {@link TraceAutoConfiguration}. Left for + * backwards compatibility reasons. + * @return reporter metrics + * @deprecated left for backwards compatibility + */ + @Bean + @ConditionalOnMissingBean + @Deprecated + ReporterMetrics zipkinReporterMetrics() { + return new InMemoryReporterMetrics(); + } + } /** diff --git a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfigurationTests.java b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfigurationTests.java index 934c9ff9b..d22d70d1a 100644 --- a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfigurationTests.java +++ b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin2/ZipkinBackwardsCompatibilityAutoConfigurationTests.java @@ -18,6 +18,7 @@ package org.springframework.cloud.sleuth.zipkin2; import org.junit.Test; import zipkin2.codec.BytesEncoder; +import zipkin2.reporter.InMemoryReporterMetrics; import zipkin2.reporter.Reporter; import zipkin2.reporter.ReporterMetrics; @@ -43,7 +44,8 @@ public class ZipkinBackwardsCompatibilityAutoConfigurationTests { assertThat(context.getBean(ZipkinProperties.class)).isNotNull(); assertThat(context.getBean(Reporter.class)).isNotNull(); assertThat(context.getBean(BytesEncoder.class)).isNotNull(); - assertThat(context.getBean(ReporterMetrics.class)).isNotNull(); + assertThat(context.getBean(ReporterMetrics.class)) + .isInstanceOf(InMemoryReporterMetrics.class); }); }