Improve reactive server span name generation

This commit is contained in:
Vladimir Kulev
2018-04-04 13:06:20 +03:00
committed by Adrian Cole
parent c5654003f8
commit b5fc46d0c0
2 changed files with 67 additions and 6 deletions

View File

@@ -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<ServerHttpRequest, ServerHttpResponse> {
@@ -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;
}
}
}

View File

@@ -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<ClientResponse> exchange = WebClient.create().get()
.uri("http://localhost:" + port + "/function").exchange();
return exchange.block();
}
private ClientResponse whenRequestIsSentToSkippedPattern(int port) {
Mono<ClientResponse> 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<ServerResponse> function() {
return RouterFunctions.route(RequestPredicates.GET("/function"), r -> {
then(MDC.get("X-B3-TraceId")).isNotEmpty();
return ServerResponse.ok().syncBody("functionOk");
});
}
}
@RestController