Fix regression introduced in #1126 (brave headers propagation) (#1206)

* Fix regression introduced in #1126 (brave headers propagation)

Before #1126, the headers were eagerly set in
`TraceExchangeFilterFunction#filter`. After it, the side effect was
moved to lazy `MonoWebClientTrace#subscribe`.

However, we have everything to instrument the request in `filter`,
and it can be done eagerly

fixes gh-1199
This commit is contained in:
Sergei Egorov
2019-02-05 17:28:57 +01:00
committed by Marcin Grzejszczak
parent 5fb9177dec
commit 4cd1218ae9
2 changed files with 40 additions and 17 deletions

View File

@@ -146,7 +146,17 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction {
@Override
public Mono<ClientResponse> filter(ClientRequest request, ExchangeFunction next) {
return new MonoWebClientTrace(next, request, this);
ClientRequest.Builder builder = ClientRequest.from(request);
if (log.isDebugEnabled()) {
log.debug("Instrumenting WebClient call");
}
Span span = handler().handleSend(injector(), builder, request,
tracer().nextSpan());
if (log.isDebugEnabled()) {
log.debug("Handled send of " + span);
}
return new MonoWebClientTrace(next, builder.build(), this, span);
}
private static final class MonoWebClientTrace extends Mono<ClientResponse> {
@@ -165,8 +175,10 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction {
final Function<? super Publisher<DataBuffer>, ? extends Publisher<DataBuffer>> scopePassingTransformer;
private final Span span;
MonoWebClientTrace(ExchangeFunction next, ClientRequest request,
TraceExchangeFilterFunction parent) {
TraceExchangeFilterFunction parent, Span span) {
this.next = next;
this.request = request;
this.tracer = parent.tracer();
@@ -174,28 +186,16 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction {
this.injector = parent.injector();
this.tracing = parent.httpTracing().tracing();
this.scopePassingTransformer = parent.scopePassingTransformer;
this.span = span;
}
@Override
public void subscribe(CoreSubscriber<? super ClientResponse> subscriber) {
final ClientRequest.Builder builder = ClientRequest.from(this.request);
Context context = subscriber.currentContext();
this.next.exchange(builder.build()).subscribe(new WebClientTracerSubscriber(
subscriber, context, findOrCreateSpan(builder), this));
}
private Span findOrCreateSpan(ClientRequest.Builder builder) {
if (log.isDebugEnabled()) {
log.debug("Instrumenting WebClient call");
}
Span clientSpan = this.handler.handleSend(this.injector, builder,
this.request, this.tracer.nextSpan());
if (log.isDebugEnabled()) {
log.debug("Handled send of " + clientSpan);
}
return clientSpan;
this.next.exchange(request).subscribe(
new WebClientTracerSubscriber(subscriber, context, span, this));
}
static final class WebClientTracerSubscriber

View File

@@ -26,6 +26,7 @@ import java.util.Map;
import java.util.Optional;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors;
import javax.servlet.http.HttpServletRequest;
@@ -491,6 +492,28 @@ public class WebClientTests {
then(this.reporter.getSpans()).extracting("kind.name").contains("CLIENT");
}
@Test
public void should_add_headers_eagerly() {
Span span = this.tracer.nextSpan().name("foo").start();
AtomicReference<String> traceId = new AtomicReference<>();
try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) {
this.webClientBuilder
.filter((request, exchange) -> {
traceId.set(request.headers().getFirst("X-B3-SpanId"));
return exchange.exchange(request);
})
.build()
.get().uri("http://localhost:" + this.port + "/traceid")
.retrieve().bodyToMono(String.class).block();
}
finally {
span.finish();
}
then(traceId).doesNotHaveValue(null);
}
private String getHeader(ResponseEntity<String> response, String name) {
List<String> headers = response.getHeaders().get(name);
return headers == null || headers.isEmpty() ? null : headers.get(0);