Brought back netty-client instrumentation; fixes gh-1080

This commit is contained in:
Marcin Grzejszczak
2018-11-02 11:34:55 +01:00
parent 423ea6d5a9
commit aa2a0209de
8 changed files with 178 additions and 110 deletions

View File

@@ -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

View File

@@ -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));

View File

@@ -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<Mono<? extends Connection>, Bootstrap, Mono<? extends Connection>> {
class TracingDoOnRequest implements BiConsumer<HttpClientRequest, Connection> {
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<HttpHeaders, String> SETTER = new Propagation.Setter<HttpHeaders, String>() {
@Override
public void put(HttpHeaders carrier, String key, String value) {
if (!carrier.contains(key)) {
carrier.add(key, value);
public Mono<? extends Connection> apply(Mono<? extends Connection> 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<HttpHeaders, String> GETTER = new Propagation.Getter<HttpHeaders, String>() {
@Override
public String get(HttpHeaders carrier, String key) {
return carrier.get(key);
}
@Override
public String toString() {
return "HttpHeaders::get";
}
};
final Tracer tracer;
final HttpClientHandler<HttpClientRequest, HttpClientResponse> handler;
final TraceContext.Injector<HttpHeaders> 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<HttpClientRequest, Connection> {
@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<HttpHeaders, String> SETTER = new Propagation.Setter<HttpHeaders, String>() {
@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<HttpHeaders, String> GETTER = new Propagation.Getter<HttpHeaders, String>() {
@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<HttpClientRequest, HttpClientResponse> handler;
final TraceContext.Injector<HttpHeaders> 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<Object> 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<HttpClientResponse, Connection> {
class TracingDoOnResponse implements BiConsumer<HttpClientResponse, Connection> {
private static final Logger log = LoggerFactory.getLogger(TracingDoOnResponse.class);
final Tracer tracer;
final HttpClientHandler<HttpClientRequest, HttpClientResponse> 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<Object> 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<HttpClientRequest, Throwable> {
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<HttpClientResponse, Throwable> {
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<HttpClientRequest, HttpClientResponse> 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<HttpClientRequest, HttpClientResponse> {
class HttpAdapter
extends brave.http.HttpClientAdapter<HttpClientRequest, HttpClientResponse> {
@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();
}
}

View File

@@ -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()

View File

@@ -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<Span> spans = new LinkedBlockingQueue<>();

View File

@@ -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

View File

@@ -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");

View File

@@ -68,7 +68,7 @@ public class RequestSendingRunnable implements Runnable {
ResponseEntity<String> 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<Void> requestWithTraceId() {