From ae0b4fe615fe20d3bd30eb64dacdaf5b71ff08bd Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Wed, 18 Sep 2019 11:52:54 +0200 Subject: [PATCH] Came back to previous impl for WebClient instrumentation; fixes gh-1442 --- .../TraceWebClientBeanPostProcessor.java | 410 ++++-------------- .../feign/issues/issue362/Issue362Tests.java | 2 +- 2 files changed, 82 insertions(+), 330 deletions(-) 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 74a134925..123ecaf95 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 @@ -19,43 +19,25 @@ package org.springframework.cloud.sleuth.instrument.web.client; import java.util.Collections; import java.util.List; import java.util.function.Consumer; -import java.util.function.Function; import brave.Span; import brave.Tracer; -import brave.Tracing; import brave.http.HttpClientHandler; import brave.http.HttpTracing; import brave.propagation.Propagation; import brave.propagation.TraceContext; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.reactivestreams.Publisher; -import org.reactivestreams.Subscription; -import reactor.core.CoreSubscriber; -import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import reactor.util.annotation.Nullable; -import reactor.util.context.Context; import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.config.BeanPostProcessor; -import org.springframework.cloud.sleuth.instrument.reactor.ReactorSleuth; -import org.springframework.core.ParameterizedTypeReference; -import org.springframework.core.io.buffer.DataBuffer; -import org.springframework.http.HttpStatus; -import org.springframework.http.ResponseCookie; -import org.springframework.http.ResponseEntity; -import org.springframework.http.client.reactive.ClientHttpResponse; -import org.springframework.util.MultiValueMap; import org.springframework.web.client.RestClientException; -import org.springframework.web.reactive.function.BodyExtractor; import org.springframework.web.reactive.function.client.ClientRequest; 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.ExchangeStrategies; import org.springframework.web.reactive.function.client.WebClient; /** @@ -108,6 +90,9 @@ final class TraceWebClientBeanPostProcessor implements BeanPostProcessor { final class TraceExchangeFilterFunction implements ExchangeFilterFunction { private static final Log log = LogFactory.getLog(TraceExchangeFilterFunction.class); + + private static final String CLIENT_SPAN_KEY = "sleuth.webclient.clientSpan"; + static final Propagation.Setter SETTER = new Propagation.Setter() { @Override public void put(ClientRequest.Builder carrier, String key, String value) { @@ -126,14 +111,12 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction { } }; - private static final String CLIENT_SPAN_KEY = "sleuth.webclient.clientSpan"; - - private static final String CANCELLED_SUBSCRIPTION_ERROR = "CANCELLED"; + public static ExchangeFilterFunction create(BeanFactory beanFactory) { + return new TraceExchangeFilterFunction(beanFactory); + } final BeanFactory beanFactory; - final Function, ? extends Publisher> scopePassingTransformer; - Tracer tracer; HttpTracing httpTracing; @@ -144,27 +127,86 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction { TraceExchangeFilterFunction(BeanFactory beanFactory) { this.beanFactory = beanFactory; - this.scopePassingTransformer = ReactorSleuth - .scopePassingSpanOperator(beanFactory); - } - - public static ExchangeFilterFunction create(BeanFactory beanFactory) { - return new TraceExchangeFilterFunction(beanFactory); } @Override public Mono filter(ClientRequest request, ExchangeFunction next) { - 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); - } + final ClientRequest.Builder builder = ClientRequest.from(request); + Mono exchange = Mono.defer(() -> next.exchange(builder.build())) + .cast(Object.class).onErrorResume(Mono::just) + .zipWith(Mono.subscriberContext()).flatMap(anyAndContext -> { + if (log.isDebugEnabled()) { + log.debug("Wrapping the context [" + anyAndContext + "]"); + } + Object any = anyAndContext.getT1(); + Span clientSpan = anyAndContext.getT2().get(CLIENT_SPAN_KEY); + Mono continuation; + final Tracer.SpanInScope ws = tracer().withSpanInScope(clientSpan); + if (any instanceof Throwable) { + continuation = Mono.error((Throwable) any); + } + else { + continuation = Mono.just((ClientResponse) any); + } + return continuation + .doAfterSuccessOrError((clientResponse, throwable1) -> { + Throwable throwable = throwable1; + if (clientResponse == null + || clientResponse.statusCode() == null) { + if (log.isDebugEnabled()) { + log.debug( + "No response was returned. Will close the span [" + + clientSpan + "]"); + } + handleReceive(clientSpan, ws, clientResponse, + throwable); + return; + } + boolean error = clientResponse.statusCode() + .is4xxClientError() + || clientResponse.statusCode().is5xxServerError(); + if (error) { + if (log.isDebugEnabled()) { + log.debug( + "Non positive status code was returned from the call. Will close the span [" + + clientSpan + "]"); + } + throwable = new RestClientException( + "Status code of the response is [" + + clientResponse.statusCode().value() + + "] and the reason is [" + + clientResponse.statusCode() + .getReasonPhrase() + + "]"); + } + handleReceive(clientSpan, ws, clientResponse, throwable); + }); + }).subscriberContext(c -> { + if (log.isDebugEnabled()) { + log.debug("Instrumenting WebClient call"); + } + Span parent = c.getOrDefault(Span.class, null); + Span clientSpan = handler().handleSend(injector(), builder, request, + tracer().nextSpan()); + if (log.isDebugEnabled()) { + log.debug("Handled send of " + 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; + } - return new MonoWebClientTrace(next, builder.build(), this, span); + private void handleReceive(Span clientSpan, Tracer.SpanInScope ws, + ClientResponse clientResponse, Throwable throwable) { + handler().handleReceive(clientResponse, throwable, clientSpan); + ws.close(); } @SuppressWarnings("unchecked") @@ -199,296 +241,6 @@ final class TraceExchangeFilterFunction implements ExchangeFilterFunction { return this.injector; } - private static final class MonoWebClientTrace extends Mono { - - final ExchangeFunction next; - - final ClientRequest request; - - final Tracer tracer; - - final HttpClientHandler handler; - - final TraceContext.Injector injector; - - final Tracing tracing; - - final Function, ? extends Publisher> scopePassingTransformer; - - private final Span span; - - MonoWebClientTrace(ExchangeFunction next, ClientRequest request, - TraceExchangeFilterFunction parent, Span span) { - this.next = next; - this.request = request; - this.tracer = parent.tracer(); - this.handler = parent.handler(); - this.injector = parent.injector(); - this.tracing = parent.httpTracing().tracing(); - this.scopePassingTransformer = parent.scopePassingTransformer; - this.span = span; - } - - @Override - public void subscribe(CoreSubscriber subscriber) { - - Context context = subscriber.currentContext(); - - this.next.exchange(request).subscribe( - new WebClientTracerSubscriber(subscriber, context, span, this)); - } - - static final class WebClientTracerSubscriber - implements CoreSubscriber { - - final CoreSubscriber actual; - - final Context context; - - final Span span; - - final Tracer.SpanInScope ws; - - final HttpClientHandler handler; - - final Function, ? extends Publisher> scopePassingTransformer; - - final Tracing tracing; - - boolean done; - - WebClientTracerSubscriber(CoreSubscriber actual, - Context context, Span span, MonoWebClientTrace parent) { - this.actual = actual; - this.span = span; - this.handler = parent.handler; - this.tracing = parent.tracing; - this.scopePassingTransformer = parent.scopePassingTransformer; - - if (!context.hasKey(Span.class)) { - context = context.put(Span.class, span); - if (log.isDebugEnabled()) { - log.debug("Reactor Context got injected with the client span " - + span); - } - } - - this.context = context.put(CLIENT_SPAN_KEY, span); - this.ws = parent.tracer.withSpanInScope(span); - - } - - @Override - public void onSubscribe(Subscription subscription) { - this.actual.onSubscribe(new Subscription() { - @Override - public void request(long n) { - subscription.request(n); - } - - @Override - public void cancel() { - terminateSpanOnCancel(); - subscription.cancel(); - } - }); - } - - @Override - public void onNext(ClientResponse response) { - this.done = true; - try { - // decorate response body - this.actual.onNext(wrapped(response)); - } - finally { - terminateSpan(response, null); - } - } - - // TODO: Remove once fixed - // https://github.com/spring-projects/spring-framework/issues/23366 - private ClientResponse wrapped(ClientResponse response) { - return new ClientResponse() { - @Override - public HttpStatus statusCode() { - try { - return response.statusCode(); - } - catch (IllegalArgumentException ex) { - return null; - } - } - - @Override - public int rawStatusCode() { - return response.rawStatusCode(); - } - - @Override - public Headers headers() { - return response.headers(); - } - - @Override - public MultiValueMap cookies() { - return response.cookies(); - } - - @Override - public ExchangeStrategies strategies() { - return response.strategies(); - } - - @Override - public T body( - BodyExtractor extractor) { - return response.body(extractor); - } - - @Override - public Mono bodyToMono(Class elementClass) { - return response.bodyToMono(elementClass); - } - - @Override - public Mono bodyToMono( - ParameterizedTypeReference typeReference) { - return response.bodyToMono(typeReference); - } - - @Override - public Flux bodyToFlux(Class elementClass) { - return (Flux) response.bodyToFlux(DataBuffer.class) - .transform(scopePassingTransformer); - } - - @Override - public Flux bodyToFlux( - ParameterizedTypeReference typeReference) { - return (Flux) response.bodyToFlux(DataBuffer.class) - .transform(scopePassingTransformer); - } - - @Override - public Mono> toEntity(Class bodyType) { - return response.toEntity(bodyType); - } - - @Override - public Mono> toEntity( - ParameterizedTypeReference typeReference) { - return response.toEntity(typeReference); - } - - @Override - public Mono>> toEntityList( - Class elementType) { - return response.toEntityList(elementType); - } - - @Override - public Mono>> toEntityList( - ParameterizedTypeReference typeReference) { - return response.toEntityList(typeReference); - } - }; - } - - @Override - public void onError(Throwable t) { - try { - this.actual.onError(t); - } - finally { - terminateSpan(null, t); - } - } - - @Override - public void onComplete() { - try { - this.actual.onComplete(); - } - finally { - if (!this.done) { - terminateSpan(null, null); - } - } - } - - @Override - public Context currentContext() { - return this.context; - } - - void handleReceive(Span clientSpan, Tracer.SpanInScope ws, - ClientResponse clientResponse, Throwable throwable) { - this.handler.handleReceive(clientResponse, throwable, clientSpan); - ws.close(); - } - - void terminateSpanOnCancel() { - if (log.isDebugEnabled()) { - log.debug("Subscription was cancelled. Will close the span [" - + this.span + "]"); - } - - this.span.tag("error", CANCELLED_SUBSCRIPTION_ERROR); - handleReceive(this.span, this.ws, null, null); - } - - void terminateSpan(@Nullable ClientResponse clientResponse, - @Nullable Throwable throwable) { - if (clientResponse == null) { - if (log.isDebugEnabled()) { - log.debug("No response was returned. Will close the span [" - + this.span + "]"); - } - handleReceive(this.span, this.ws, clientResponse, throwable); - return; - } - int statusCode = statusCodeAsInt(clientResponse); - boolean error = isError(statusCode); - if (error) { - if (log.isDebugEnabled()) { - log.debug( - "Non positive status code was returned from the call. Will close the span [" - + this.span + "]"); - } - throwable = new RestClientException("Status code of the response is [" - + statusCode + "] and the reason is [" - + reasonPhrase(clientResponse) + "]"); - } - handleReceive(this.span, this.ws, clientResponse, throwable); - } - - private String reasonPhrase(ClientResponse clientResponse) { - try { - return clientResponse.statusCode().getReasonPhrase(); - } - catch (IllegalArgumentException ex) { - return ""; - } - } - - private boolean isError(int code) { - return code >= 400; - } - - private int statusCodeAsInt(ClientResponse response) { - try { - return response.rawStatusCode(); - } - catch (Exception dontCare) { - return 0; - } - } - - } - - } - static final class HttpAdapter extends brave.http.HttpClientAdapter { diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/issues/issue362/Issue362Tests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/issues/issue362/Issue362Tests.java index ee49aa591..ae5423f59 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/issues/issue362/Issue362Tests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/issues/issue362/Issue362Tests.java @@ -209,7 +209,7 @@ class CustomConfig { this.feignComponentAsserter.executedComponents.put(ErrorDecoder.class, true); if (response.status() == 409) { return new RetryableException(response.status(), "Article not Ready", - Request.HttpMethod.GET, new Date()); + Request.HttpMethod.GET, new Date(), response.request()); } else { return super.decode(methodKey, response);