From b5fc46d0c0b19fd1d6cc72d34aa4e34bbabe7454 Mon Sep 17 00:00:00 2001 From: Vladimir Kulev Date: Wed, 4 Apr 2018 13:06:20 +0300 Subject: [PATCH] Improve reactive server span name generation --- .../sleuth/instrument/web/TraceWebFilter.java | 43 ++++++++++++++++--- .../instrument/web/TraceWebFluxTests.java | 30 +++++++++++++ 2 files changed, 67 insertions(+), 6 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 9d7cb4076..06aa7793f 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 @@ -31,6 +31,7 @@ import org.springframework.core.Ordered; import org.springframework.http.HttpHeaders; import org.springframework.http.server.reactive.ServerHttpRequest; import org.springframework.http.server.reactive.ServerHttpResponse; +import org.springframework.http.server.reactive.ServerHttpResponseDecorator; import org.springframework.web.method.HandlerMethod; import org.springframework.web.reactive.HandlerMapping; import org.springframework.web.server.ServerWebExchange; @@ -125,9 +126,7 @@ public final class TraceWebFilter implements WebFilter, Ordered { // clear any previous trace tracer().withSpanInScope(null); } - ServerHttpRequest request = exchange.getRequest(); - ServerHttpResponse response = exchange.getResponse(); - String uri = request.getPath().pathWithinApplication().value(); + String uri = exchange.getRequest().getPath().pathWithinApplication().value(); if (log.isDebugEnabled()) { log.debug("Received a request to uri [" + uri + "]"); } @@ -149,15 +148,22 @@ public final class TraceWebFilter implements WebFilter, Ordered { } else { continuation = Mono.empty(); } + String httpRoute = null; Object attribute = exchange .getAttribute(HandlerMapping.BEST_MATCHING_HANDLER_ATTRIBUTE); if (attribute instanceof HandlerMethod) { HandlerMethod handlerMethod = (HandlerMethod) attribute; addClassMethodTag(handlerMethod, span); addClassNameTag(handlerMethod, span); + Object pattern = exchange + .getAttribute(HandlerMapping.BEST_MATCHING_PATTERN_ATTRIBUTE); + httpRoute = pattern != null ? pattern.toString() : ""; } - addResponseTagsForSpanWithoutParent(exchange, response, span); - handler().handleSend(response, t, span); + addResponseTagsForSpanWithoutParent(exchange, exchange.getResponse(), span); + DecoratedServerHttpResponse delegate = new DecoratedServerHttpResponse( + exchange.getResponse(), exchange.getRequest().getMethodValue(), + httpRoute); + handler().handleSend(delegate, t, span); if (log.isDebugEnabled()) { log.debug("Handled send of " + span); } @@ -181,7 +187,7 @@ public final class TraceWebFilter implements WebFilter, Ordered { } } else { span = handler().handleReceive(extractor(), - request.getHeaders(), request); + exchange.getRequest().getHeaders(), exchange.getRequest()); if (log.isDebugEnabled()) { log.debug("Handled receive of span " + span); } @@ -255,6 +261,17 @@ public final class TraceWebFilter implements WebFilter, Ordered { return ORDER; } + static final class DecoratedServerHttpResponse extends ServerHttpResponseDecorator { + + final String method, httpRoute; + + DecoratedServerHttpResponse(ServerHttpResponse delegate, String method, String httpRoute) { + super(delegate); + this.method = method; + this.httpRoute = httpRoute; + } + } + static final class HttpAdapter extends brave.http.HttpServerAdapter { @@ -275,6 +292,20 @@ public final class TraceWebFilter implements WebFilter, Ordered { return response.getStatusCode() != null ? response.getStatusCode().value() : null; } + + @Override public String methodFromResponse(ServerHttpResponse response) { + if (response instanceof DecoratedServerHttpResponse) { + return ((DecoratedServerHttpResponse) response).method; + } + return null; + } + + @Override public String route(ServerHttpResponse response) { + if (response instanceof DecoratedServerHttpResponse) { + return ((DecoratedServerHttpResponse) response).httpRoute; + } + return null; + } } } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFluxTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFluxTests.java index 508beeb7b..8a3de36bb 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFluxTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceWebFluxTests.java @@ -41,6 +41,7 @@ import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.reactive.function.client.ClientResponse; import org.springframework.web.reactive.function.client.WebClient; +import org.springframework.web.reactive.function.server.*; import reactor.core.publisher.Flux; import reactor.core.publisher.Hooks; import reactor.core.publisher.Mono; @@ -76,6 +77,12 @@ public class TraceWebFluxTests { thenSpanWasReportedWithTags(accumulator, response); clean(accumulator, controller2); + // when + ClientResponse functionResponse = whenRequestIsSentToFunction(port); + // then + thenSpanWasReportedForFunction(accumulator, functionResponse); + accumulator.clear(); + // when ClientResponse nonSampledResponse = whenNonSampledRequestIsSent(port); // then @@ -101,11 +108,21 @@ public class TraceWebFluxTests { then(response.statusCode().value()).isEqualTo(200); then(accumulator.getSpans()).hasSize(1); }); + then(accumulator.getSpans().get(0).name()).isEqualTo("get /api/c2/{id}"); then(accumulator.getSpans().get(0).tags()) .containsEntry("mvc.controller.method", "successful") .containsEntry("mvc.controller.class", "Controller2"); } + private void thenSpanWasReportedForFunction(ArrayListSpanReporter accumulator, + ClientResponse response) { + Awaitility.await().untilAsserted(() -> { + then(response.statusCode().value()).isEqualTo(200); + then(accumulator.getSpans()).hasSize(1); + }); + then(accumulator.getSpans().get(0).name()).isEqualTo("get"); + } + private void thenNoSpanWasReported(ArrayListSpanReporter accumulator, ClientResponse response, Controller2 controller2) { Awaitility.await().untilAsserted(() -> { @@ -122,6 +139,12 @@ public class TraceWebFluxTests { return exchange.block(); } + private ClientResponse whenRequestIsSentToFunction(int port) { + Mono exchange = WebClient.create().get() + .uri("http://localhost:" + port + "/function").exchange(); + return exchange.block(); + } + private ClientResponse whenRequestIsSentToSkippedPattern(int port) { Mono exchange = WebClient.create().get() .uri("http://localhost:" + port + "/skipped").exchange(); @@ -160,6 +183,13 @@ public class TraceWebFluxTests { @Bean Controller2 controller2(Tracer tracer) { return new Controller2(tracer); } + + @Bean RouterFunction function() { + return RouterFunctions.route(RequestPredicates.GET("/function"), r -> { + then(MDC.get("X-B3-TraceId")).isNotEmpty(); + return ServerResponse.ok().syncBody("functionOk"); + }); + } } @RestController