From aa2a0209dedd2edd063b3c2843ba1df2aa7fe04d Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Fri, 2 Nov 2018 11:34:55 +0100 Subject: [PATCH] Brought back netty-client instrumentation; fixes gh-1080 --- .../sleuth/instrument/web/ServletUtils.java | 2 +- .../web/SleuthHttpServerParser.java | 2 +- .../TraceWebClientAutoConfiguration.java | 272 +++++++++++------- .../TraceWebClientBeanPostProcessor.java | 4 +- .../JmsTracingConfigurationTest.java | 2 +- ...raceCustomFilterResponseInjectorTests.java | 2 +- .../client/integration/WebClientTests.java | 2 +- .../java/tools/RequestSendingRunnable.java | 2 +- 8 files changed, 178 insertions(+), 110 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/ServletUtils.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/ServletUtils.java index 0809e0704..5aa688695 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/ServletUtils.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/ServletUtils.java @@ -20,7 +20,7 @@ import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; /** - * Utility class to retrieve data from Servlet HTTP request and response. + * Utility class to retrieve data from Servlet HTTP request and handle. * * @author Marcin Grzejszczak * @since 1.0.0 diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/SleuthHttpServerParser.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/SleuthHttpServerParser.java index d283db9a6..acaa50357 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/SleuthHttpServerParser.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/SleuthHttpServerParser.java @@ -72,7 +72,7 @@ class SleuthHttpServerParser extends HttpServerParser { return; } if (httpStatus == HttpServletResponse.SC_OK && error != null) { - // Filter chain threw exception but the response status may not have been set + // Filter chain threw exception but the handle status may not have been set // yet, so we have to guess. customizer.tag(STATUS_CODE_KEY, String.valueOf(HttpServletResponse.SC_INTERNAL_SERVER_ERROR)); 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 d0da65939..b83bbec4a 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 @@ -20,6 +20,7 @@ import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.function.BiConsumer; +import java.util.function.BiFunction; import brave.Span; import brave.Tracer; @@ -30,13 +31,13 @@ import brave.httpclient.TracingHttpClientBuilder; import brave.propagation.Propagation; import brave.propagation.TraceContext; import brave.spring.web.TracingClientHttpRequestInterceptor; +import io.netty.bootstrap.Bootstrap; import io.netty.handler.codec.http.HttpHeaders; -import io.netty.util.Attribute; -import io.netty.util.AttributeKey; import org.apache.http.impl.client.HttpClientBuilder; import org.apache.http.impl.nio.client.HttpAsyncClientBuilder; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import reactor.core.publisher.Mono; import reactor.netty.Connection; import reactor.netty.http.client.HttpClient; import reactor.netty.http.client.HttpClientRequest; @@ -325,138 +326,205 @@ class HttpClientBeanPostProcessor implements BeanPostProcessor { public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { if (bean instanceof HttpClient) { - return ((HttpClient) bean) + return ((HttpClient) bean).mapConnect(new TracingMapConnect(this.beanFactory)) .doOnRequest(TracingDoOnRequest.create(this.beanFactory)) - .doOnResponse(TracingDoOnResponse.create(this.beanFactory)); + .doOnRequestError(TracingDoOnErrorRequest.create(this.beanFactory)) + .doOnResponse(TracingDoOnResponse.create(this.beanFactory)) + .doOnResponseError(TracingDoOnErrorResponse.create(this.beanFactory)); } return bean; } -} + private static class TracingMapConnect implements + BiFunction, Bootstrap, Mono> { -class TracingDoOnRequest implements BiConsumer { + private final BeanFactory beanFactory; - private static final Logger log = LoggerFactory.getLogger(TracingDoOnRequest.class); + private Tracer tracer; + + TracingMapConnect(BeanFactory beanFactory) { + this.beanFactory = beanFactory; + } - static final Propagation.Setter SETTER = new Propagation.Setter() { @Override - public void put(HttpHeaders carrier, String key, String value) { - if (!carrier.contains(key)) { - carrier.add(key, value); + public Mono apply(Mono mono, + Bootstrap bootstrap) { + return mono.subscriberContext( + context -> context.put(Span.class, tracer().nextSpan())); + } + + private Tracer tracer() { + if (this.tracer == null) { + this.tracer = this.beanFactory.getBean(Tracer.class); } + return this.tracer; } - @Override - public String toString() { - return "HttpHeaders::add"; - } - }; - static final Propagation.Getter GETTER = new Propagation.Getter() { - @Override - public String get(HttpHeaders carrier, String key) { - return carrier.get(key); - } - - @Override - public String toString() { - return "HttpHeaders::get"; - } - }; - - final Tracer tracer; - - final HttpClientHandler handler; - - final TraceContext.Injector injector; - - final HttpTracing httpTracing; - - TracingDoOnRequest(HttpTracing httpTracing) { - this.tracer = httpTracing.tracing().tracer(); - this.handler = HttpClientHandler.create(httpTracing, new HttpAdapter()); - this.injector = httpTracing.tracing().propagation().injector(SETTER); - this.httpTracing = httpTracing; } - static TracingDoOnRequest create(BeanFactory beanFactory) { - return new TracingDoOnRequest(beanFactory.getBean(HttpTracing.class)); - } + private static class TracingDoOnRequest + implements BiConsumer { - @Override - public void accept(HttpClientRequest req, Connection connection) { - final Span currentSpan = this.tracer.currentSpan(); - try (Tracer.SpanInScope spanInScope = this.tracer.withSpanInScope(currentSpan)) { + static final Propagation.Setter SETTER = new Propagation.Setter() { + @Override + public void put(HttpHeaders carrier, String key, String value) { + if (!carrier.contains(key)) { + carrier.add(key, value); + } + } + + @Override + public String toString() { + return "HttpHeaders::add"; + } + }; + static final Propagation.Getter GETTER = new Propagation.Getter() { + @Override + public String get(HttpHeaders carrier, String key) { + return carrier.get(key); + } + + @Override + public String toString() { + return "HttpHeaders::get"; + } + }; + + private static final Logger log = LoggerFactory + .getLogger(TracingDoOnRequest.class); + + final Tracer tracer; + + final HttpClientHandler handler; + + final TraceContext.Injector injector; + + final HttpTracing httpTracing; + + TracingDoOnRequest(HttpTracing httpTracing) { + this.tracer = httpTracing.tracing().tracer(); + this.handler = HttpClientHandler.create(httpTracing, new HttpAdapter()); + this.injector = httpTracing.tracing().propagation().injector(SETTER); + this.httpTracing = httpTracing; + } + + static TracingDoOnRequest create(BeanFactory beanFactory) { + return new TracingDoOnRequest(beanFactory.getBean(HttpTracing.class)); + } + + @Override + public void accept(HttpClientRequest req, Connection connection) { + Span span = req.currentContext().getOrDefault(Span.class, + this.tracer.nextSpan()); if (log.isDebugEnabled()) { log.debug("Wrapping do on request"); } - Span span = this.handler.handleSend(this.injector, req.requestHeaders(), req); - Attribute attribute = connection.channel() - .attr(AttributeKey.valueOf("span")); - attribute.set(span); + this.handler.handleSend(this.injector, req.requestHeaders(), req, span); } + } -} + private static class TracingDoOnResponse extends AbstractTracingDoOnHandler + implements BiConsumer { -class TracingDoOnResponse implements BiConsumer { - - private static final Logger log = LoggerFactory.getLogger(TracingDoOnResponse.class); - - final Tracer tracer; - - final HttpClientHandler handler; - - TracingDoOnResponse(HttpTracing httpTracing) { - this.tracer = httpTracing.tracing().tracer(); - this.handler = HttpClientHandler.create(httpTracing, new HttpAdapter()); - } - - static TracingDoOnResponse create(BeanFactory beanFactory) { - return new TracingDoOnResponse(beanFactory.getBean(HttpTracing.class)); - } - - @Override - public void accept(HttpClientResponse httpClientResponse, Connection connection) { - Attribute spanAttr = connection.channel() - .attr(AttributeKey.valueOf("span")); - Span span = (Span) spanAttr.get(); - if (span == null) { - return; + TracingDoOnResponse(HttpTracing httpTracing) { + super(httpTracing); } - try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) { - if (log.isDebugEnabled()) { - log.debug("Setting client sent spans"); + + static TracingDoOnResponse create(BeanFactory beanFactory) { + return new TracingDoOnResponse(beanFactory.getBean(HttpTracing.class)); + } + + @Override + public void accept(HttpClientResponse httpClientResponse, Connection connection) { + handle(httpClientResponse, null); + } + + } + + private static class TracingDoOnErrorRequest extends AbstractTracingDoOnHandler + implements BiConsumer { + + TracingDoOnErrorRequest(HttpTracing httpTracing) { + super(httpTracing); + } + + static TracingDoOnErrorRequest create(BeanFactory beanFactory) { + return new TracingDoOnErrorRequest(beanFactory.getBean(HttpTracing.class)); + } + + @Override + public void accept(HttpClientRequest request, Throwable throwable) { + handle(null, throwable); + } + + } + + private static class TracingDoOnErrorResponse extends AbstractTracingDoOnHandler + implements BiConsumer { + + TracingDoOnErrorResponse(HttpTracing httpTracing) { + super(httpTracing); + } + + static TracingDoOnErrorResponse create(BeanFactory beanFactory) { + return new TracingDoOnErrorResponse(beanFactory.getBean(HttpTracing.class)); + } + + @Override + public void accept(HttpClientResponse httpClientResponse, Throwable throwable) { + handle(httpClientResponse, throwable); + } + + } + + private static abstract class AbstractTracingDoOnHandler { + + final Tracer tracer; + + final HttpClientHandler handler; + + AbstractTracingDoOnHandler(HttpTracing httpTracing) { + this.tracer = httpTracing.tracing().tracer(); + this.handler = HttpClientHandler.create(httpTracing, new HttpAdapter()); + } + + protected void handle(HttpClientResponse httpClientResponse, + Throwable throwable) { + Span span = httpClientResponse.currentContext().getOrDefault(Span.class, + null); + if (span == null) { + return; } - // status codes and CR - // TODO: Add throwable - this.handler.handleReceive(httpClientResponse, null, span); + this.handler.handleReceive(httpClientResponse, throwable, span); } + } -} + private static class HttpAdapter + extends brave.http.HttpClientAdapter { -class HttpAdapter - extends brave.http.HttpClientAdapter { + @Override + public String method(HttpClientRequest request) { + return request.method().name(); + } - @Override - public String method(HttpClientRequest request) { - return request.method().name(); - } + @Override + public String url(HttpClientRequest request) { + return request.uri(); + } - @Override - public String url(HttpClientRequest request) { - return request.uri(); - } + @Override + public String requestHeader(HttpClientRequest request, String name) { + Object result = request.requestHeaders().get(name); + return result != null ? result.toString() : ""; + } - @Override - public String requestHeader(HttpClientRequest request, String name) { - Object result = request.requestHeaders().get(name); - return result != null ? result.toString() : ""; - } + @Override + public Integer statusCode(HttpClientResponse response) { + return response.status().code(); + } - @Override - public Integer statusCode(HttpClientResponse response) { - return response.status().code(); } } 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 c117661ae..225f31286 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 @@ -155,7 +155,7 @@ class TraceExchangeFilterFunction implements ExchangeFilterFunction { || clientResponse.statusCode() == null) { if (log.isDebugEnabled()) { log.debug( - "No response was returned. Will close the span [" + "No handle was returned. Will close the span [" + clientSpan + "]"); } handleReceive(clientSpan, ws, clientResponse, @@ -172,7 +172,7 @@ class TraceExchangeFilterFunction implements ExchangeFilterFunction { + clientSpan + "]"); } throwable = new RestClientException( - "Status code of the response is [" + "Status code of the handle is [" + clientResponse.statusCode().value() + "] and the reason is [" + clientResponse.statusCode() diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/JmsTracingConfigurationTest.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/JmsTracingConfigurationTest.java index fc21624f9..b2a697d0c 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/JmsTracingConfigurationTest.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/JmsTracingConfigurationTest.java @@ -281,7 +281,7 @@ class JmsTestTracingConfiguration { /** * When testing servers or asynchronous clients, spans are reported on a worker * thread. In order to read them on the main thread, we use a concurrent queue. As - * some implementations report after a response is sent, we use a blocking queue to + * some implementations report after a handle is sent, we use a blocking queue to * prevent race conditions in tests. */ BlockingQueue spans = new LinkedBlockingQueue<>(); diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceCustomFilterResponseInjectorTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceCustomFilterResponseInjectorTests.java index 48bd3e9b3..34127df9a 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceCustomFilterResponseInjectorTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceCustomFilterResponseInjectorTests.java @@ -80,7 +80,7 @@ public class TraceCustomFilterResponseInjectorTests { Map.class); then(responseEntity.getHeaders()).containsKeys(TRACE_ID_NAME, SPAN_ID_NAME) - .as("Trace headers must be present in response headers"); + .as("Trace headers must be present in handle headers"); } @Configuration 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 464e2ef8d..351d9c748 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 @@ -282,7 +282,6 @@ public class WebClientTests { @Test @SuppressWarnings("unchecked") - @Ignore public void shouldAttachTraceIdWhenCallingAnotherServiceForNettyHttpClient() throws Exception { Span span = this.tracer.nextSpan().name("foo").start(); @@ -295,6 +294,7 @@ public class WebClientTests { } then(this.tracer.currentSpan()).isNull(); + System.out.println("Collected span " + this.reporter.getSpans()); then(this.reporter.getSpans()).isNotEmpty().extracting("traceId", String.class) .containsOnly(span.context().traceIdString()); then(this.reporter.getSpans()).extracting("kind.name").contains("CLIENT"); diff --git a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-test-core/src/main/java/tools/RequestSendingRunnable.java b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-test-core/src/main/java/tools/RequestSendingRunnable.java index c358a1e64..0cfaf2730 100644 --- a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-test-core/src/main/java/tools/RequestSendingRunnable.java +++ b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-test-core/src/main/java/tools/RequestSendingRunnable.java @@ -68,7 +68,7 @@ public class RequestSendingRunnable implements Runnable { ResponseEntity responseEntity = this.restTemplate .exchange(requestWithTraceId(), String.class); then(responseEntity.getStatusCode()).isEqualTo(HttpStatus.OK); - log.info(String.format("Received the following response [%s]", responseEntity)); + log.info(String.format("Received the following handle [%s]", responseEntity)); } private RequestEntity requestWithTraceId() {