From a6ae8ed4fd3bd86a0b53ca317b83e1ceaa3b77a8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Simon=20Basl=C3=A9?= Date: Tue, 31 Oct 2017 16:49:02 +0100 Subject: [PATCH] make use of Reactor Context in WebFlux related tracing (#764) hopefully fixes #679 #677 --- .../sleuth/instrument/web/TraceWebFilter.java | 66 ++++++++--- .../TraceWebClientBeanPostProcessor.java | 103 ++++++++++++------ 2 files changed, 115 insertions(+), 54 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFilter.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFilter.java index cbf466cd3..76c86d756 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFilter.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFilter.java @@ -5,6 +5,9 @@ import java.util.regex.Pattern; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; +import reactor.core.Scannable; +import reactor.core.publisher.Mono; + import org.springframework.beans.factory.BeanFactory; import org.springframework.cloud.sleuth.ErrorParser; import org.springframework.cloud.sleuth.Span; @@ -23,7 +26,6 @@ import org.springframework.web.reactive.HandlerMapping; import org.springframework.web.server.ServerWebExchange; import org.springframework.web.server.WebFilter; import org.springframework.web.server.WebFilterChain; -import reactor.core.publisher.Mono; /** * A {@link WebFilter} that creates / continues / closes and detaches spans @@ -81,23 +83,51 @@ public class TraceWebFilter implements WebFilter, Ordered { continueSpan(exchange, spanFromAttribute); } String name = HTTP_COMPONENT + ":" + uri; - Span span = createSpan(request, exchange, skip, spanFromAttribute, name); - return chain.filter(exchange).compose(f -> f.doOnSuccess(t -> { - addResponseTags(response, null); - }).doOnError(t -> { - errorParser().parseErrorTags(tracer().getCurrentSpan(), t); - addResponseTags(response, t); - }).doFinally(t -> { - Object attribute = exchange - .getAttribute(HandlerMapping.BEST_MATCHING_HANDLER_ATTRIBUTE); - if (attribute instanceof HandlerMethod) { - HandlerMethod handlerMethod = (HandlerMethod) attribute; - addClassMethodTag(handlerMethod, span); - addClassNameTag(handlerMethod, span); - } - addResponseTagsForSpanWithoutParent(exchange, response); - detachOrCloseSpans(span); - })); + final String CONTEXT_ERROR = "sleuth.webfilter.context.error"; + return chain + .filter(exchange) + .compose(f -> f.then(Mono.subscriberContext()) + .onErrorResume(t -> Mono.subscriberContext().map(c -> c.put(CONTEXT_ERROR, t))) + .flatMap(c -> { + //reactivate span from context + Span span = c.getOrDefault(Span.class, null); + if (span != null) { + tracer().continueSpan(span); + } + Mono continuation; + + if (c.hasKey(CONTEXT_ERROR)) { + Throwable t = c.get(CONTEXT_ERROR); + errorParser().parseErrorTags(tracer().getCurrentSpan(), t); + addResponseTags(response, t); + continuation = Mono.error(t); + } else { + addResponseTags(response, null); + continuation = Mono.empty(); + } + Object attribute = exchange + .getAttribute(HandlerMapping.BEST_MATCHING_HANDLER_ATTRIBUTE); + if (attribute instanceof HandlerMethod) { + HandlerMethod handlerMethod = (HandlerMethod) attribute; + addClassMethodTag(handlerMethod, span); + addClassNameTag(handlerMethod, span); + } + addResponseTagsForSpanWithoutParent(exchange, response); + detachOrCloseSpans(span); + + return continuation; + }) + .subscriberContext(c -> { + Span span; + if (c.hasKey(Span.class)) { + Span parent = c.get(Span.class); + span = createSpan(request, exchange, skip, parent, name); + } else { + span = createSpan(request, exchange, skip, spanFromAttribute, name); + } + + return c.put(Span.class, span); + })); } private void addResponseTagsForSpanWithoutParent(ServerWebExchange exchange, 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 09b8daa40..fd4228c97 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 @@ -5,9 +5,14 @@ import java.util.AbstractMap; import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.stream.Stream; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; +import reactor.core.Scannable; +import reactor.core.publisher.Mono; +import reactor.util.function.Tuple2; + import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.config.BeanPostProcessor; @@ -25,7 +30,6 @@ import org.springframework.web.reactive.function.client.ClientResponse; import org.springframework.web.reactive.function.client.ExchangeFilterFunction; import org.springframework.web.reactive.function.client.ExchangeFunction; import org.springframework.web.reactive.function.client.WebClient; -import reactor.core.publisher.Mono; /** * {@link BeanPostProcessor} to wrap a {@link WebClient} instance into @@ -62,6 +66,8 @@ class TraceWebClientBeanPostProcessor implements BeanPostProcessor { class TraceExchangeFilterFunction implements ExchangeFilterFunction { + private static final String CLIENT_SPAN_KEY = "sleuth.webclient.clientSpan"; + private static final Log log = LogFactory.getLog(TraceExchangeFilterFunction.class); private Tracer tracer; @@ -76,43 +82,63 @@ class TraceExchangeFilterFunction implements ExchangeFilterFunction { @Override public Mono filter(ClientRequest request, ExchangeFunction next) { - if (log.isDebugEnabled()) { - log.debug("Creating a client span for the RPC"); - } - final Span clientSpan = createNewSpan(request); - ClientRequest.Builder builder = ClientRequest.from(request); - httpSpanInjector().inject(clientSpan, new ClientRequestTextMap(request, builder)); - if (log.isDebugEnabled()) { - log.debug("Headers got injected to the client span " + clientSpan); - } + final ClientRequest.Builder builder = ClientRequest.from(request); + Mono exchange = next.exchange(builder.build()) - .doOnError(throwable -> { + .cast(Object.class) + .onErrorResume(Mono::just) + .zipWith(Mono.subscriberContext()) + .flatMap(anyAndContext -> { + Object any = anyAndContext.getT1(); + Span clientSpan = anyAndContext.getT2().get(CLIENT_SPAN_KEY); + tracer().continueSpan(clientSpan); - errorParser().parseErrorTags(clientSpan, throwable); - }).doOnSuccess(response -> { - tracer().continueSpan(clientSpan); - boolean error = response.statusCode().is4xxClientError() || response - .statusCode().is5xxServerError(); - if (error) { - if (log.isDebugEnabled()) { - log.debug( - "Non positive status code was returned from the call. Will close the span [" - + clientSpan + "]"); + + Mono continuation; + if (any instanceof Throwable) { + Throwable throwable = (Throwable) any; + errorParser().parseErrorTags(clientSpan, throwable); + continuation = Mono.error(throwable); + } else { + ClientResponse response = (ClientResponse) any; + boolean error = response.statusCode().is4xxClientError() || response + .statusCode().is5xxServerError(); + if (error) { + if (log.isDebugEnabled()) { + log.debug( + "Non positive status code was returned from the call. Will close the span [" + + clientSpan + "]"); + } + errorParser().parseErrorTags(clientSpan, new RestClientException( + "Status code of the response is [" + response.statusCode() + .value() + "] and the reason is [" + response + .statusCode().getReasonPhrase() + "]")); } - errorParser().parseErrorTags(clientSpan, new RestClientException( - "Status code of the response is [" + response.statusCode() - .value() + "] and the reason is [" + response - .statusCode().getReasonPhrase() + "]")); + continuation = Mono.just(response); } - }).doFinally(signalType -> finish(clientSpan)); - if (log.isDebugEnabled()) { - log.debug("Will detach the client span " + clientSpan); - } - Span detachedSpan = tracer().detach(clientSpan); - tracer().continueSpan(detachedSpan); - if (log.isDebugEnabled()) { - log.debug("Client span detached"); - } + finish(clientSpan); + return continuation; + }) + .subscriberContext(c -> { + if (log.isDebugEnabled()) { + log.debug("Creating a client span for the WebClient"); + } + Span parent = c.getOrDefault(Span.class, null); + Span clientSpan = createNewSpan(request, parent); + + httpSpanInjector().inject(clientSpan, new ClientRequestTextMap(request, builder)); + if (log.isDebugEnabled()) { + log.debug("Headers got injected to the client span " + clientSpan); + } + + if (parent == null) { + c = c.put(Span.class, clientSpan); + if (log.isDebugEnabled()) { + log.debug("Reactor Context got injected with the client span " + clientSpan); + } + } + return c.put(CLIENT_SPAN_KEY, clientSpan); + }); return exchange; } @@ -120,10 +146,15 @@ class TraceExchangeFilterFunction implements ExchangeFilterFunction { * Enriches the request with proper headers and publishes * the client sent event */ - private Span createNewSpan(ClientRequest request) { + private Span createNewSpan(ClientRequest request, Span optionalParent) { URI uri = request.url(); String spanName = getName(uri); - Span newSpan = tracer().createSpan(spanName); + Span newSpan; + if (optionalParent == null) { + newSpan = tracer().createSpan(spanName); + } else { + newSpan = tracer().createSpan(spanName, optionalParent); + } addRequestTags(request); newSpan.logEvent(Span.CLIENT_SEND); if (log.isDebugEnabled()) {