From ca2e63e3b1489d86b36a5d3a5a0792ba910a55b5 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Wed, 9 May 2018 20:10:01 +0200 Subject: [PATCH] Adding a workaround that should fix gh-973 --- .../instrument/reactor/ReactorSleuth.java | 100 ++++++++++-------- .../FeignClientServerErrorTests.java | 2 - 2 files changed, 57 insertions(+), 45 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java index aaf22cd36..df5440731 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java @@ -20,17 +20,18 @@ import java.util.function.Function; import java.util.function.Predicate; import brave.Tracing; -import reactor.core.CoreSubscriber; -import reactor.core.Fuseable; -import reactor.core.Scannable; -import reactor.core.publisher.ConnectableFlux; -import reactor.core.publisher.Operators; -import reactor.util.context.Context; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.reactivestreams.Publisher; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.NoSuchBeanDefinitionException; +import reactor.core.CoreSubscriber; +import reactor.core.Fuseable; +import reactor.core.Scannable; +import reactor.core.publisher.ConnectableFlux; +import reactor.core.publisher.GroupedFlux; +import reactor.core.publisher.Operators; +import reactor.util.context.Context; /** * Reactive Span pointcuts factories @@ -56,23 +57,28 @@ public abstract class ReactorSleuth { */ public static Function, ? extends Publisher> spanOperator( BeanFactory beanFactory) { - return Operators.lift(POINTCUT_FILTER, ((scannable, sub) -> { - //do not trace fused flows - if(scannable instanceof Fuseable && sub instanceof Fuseable.QueueSubscription){ - return sub; + return sourcePub -> { + //do the checks directly on actual original Publisher + if (sourcePub instanceof Fuseable //Sleuth can't handle that + || sourcePub instanceof Fuseable.ScalarCallable //IIRC Sleuth can't handle that + || sourcePub instanceof ConnectableFlux //Operators.lift can't handle that + || sourcePub instanceof GroupedFlux //Operators.lift can't handle that + ) { + return sourcePub; } - if (log.isTraceEnabled()) { - log.trace("Creating a lazy span subscriber with context " - + "[" + sub.currentContext() + "] and name [" + scannable.name() + "]"); - } - return new LazySpanSubscriber( - new SpanSubscriptionProvider( - beanFactory, - sub, - sub.currentContext(), - scannable.name()) - ); - })); + //no more POINTCUT_FILTER since mecanism is broken + Function, ? extends Publisher> lift = Operators.lift((scannable, sub) -> { + //rest of the logic unchanged... + return new LazySpanSubscriber( + new SpanSubscriptionProvider( + beanFactory, + sub, + sub.currentContext(), + scannable.name()) + ); + }); + return lift.apply(sourcePub); + }; } /** @@ -91,27 +97,35 @@ public abstract class ReactorSleuth { @SuppressWarnings("unchecked") public static Function, ? extends Publisher> scopePassingSpanOperator( BeanFactory beanFactory) { - return Operators.lift(POINTCUT_FILTER, ((scannable, sub) -> { - //do not trace fused flows - if(scannable instanceof Fuseable && sub instanceof Fuseable.QueueSubscription){ - return sub; + return sourcePub -> { + //do the checks directly on actual original Publisher + if (sourcePub instanceof Fuseable //Sleuth can't handle that + || sourcePub instanceof Fuseable.ScalarCallable //IIRC Sleuth can't handle that + || sourcePub instanceof ConnectableFlux //Operators.lift can't handle that + || sourcePub instanceof GroupedFlux //Operators.lift can't handle that + ) { + return sourcePub; } - if (contextRefreshed(beanFactory)) { - if (log.isTraceEnabled()) { - log.trace("Spring Context already refreshed. Creating a scope " - + "passing span subscriber with Reactor Context " - + "[" + sub.currentContext() + "] and name [" + scannable.name() + "]"); + //no more POINTCUT_FILTER since mecanism is broken + Function, ? extends Publisher> lift = Operators.lift((scannable, sub) -> { + //rest of the logic unchanged... + if (contextRefreshed(beanFactory)) { + if (log.isTraceEnabled()) { + log.trace("Spring Context already refreshed. Creating a scope " + "passing span subscriber with Reactor Context " + "[" + sub.currentContext() + "] and name [" + scannable.name() + "]"); + } + return scopePassingSpanSubscription(beanFactory, scannable, sub).get(); } - return scopePassingSpanSubscription(beanFactory, scannable, sub).get(); - } - if (log.isTraceEnabled()) { - log.trace("Spring Context is not yet refreshed, falling back to lazy span subscriber. " - + "Reactor Context is [" + sub.currentContext() + "] and name is [" + scannable.name() + "]"); - } - return new LazySpanSubscriber( - scopePassingSpanSubscription(beanFactory, scannable, sub) - ); - })); + if (log.isTraceEnabled()) { + log.trace( + "Spring Context is not yet refreshed, falling back to lazy span subscriber. " + "Reactor Context is [" + sub.currentContext() + "] and name is [" + scannable.name() + "]"); + } + return new LazySpanSubscriber( + scopePassingSpanSubscription(beanFactory, scannable, sub) + ); + }); + + return lift.apply(sourcePub); + }; } private static boolean contextRefreshed(BeanFactory beanFactory) { @@ -122,9 +136,9 @@ public abstract class ReactorSleuth { } } - private static SpanSubscriptionProvider scopePassingSpanSubscription( + private static SpanSubscriptionProvider scopePassingSpanSubscription( BeanFactory beanFactory, Scannable scannable, CoreSubscriber sub) { - return new SpanSubscriptionProvider( + return new SpanSubscriptionProvider( beanFactory, sub, sub.currentContext(), diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/servererrors/FeignClientServerErrorTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/servererrors/FeignClientServerErrorTests.java index 1420cd6c2..33074a3e0 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/servererrors/FeignClientServerErrorTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/feign/servererrors/FeignClientServerErrorTests.java @@ -77,8 +77,6 @@ import static org.assertj.core.api.BDDAssertions.then; "spring.sleuth.http.legacy.enabled=true", "hystrix.command.default.execution.isolation.thread.timeoutInMilliseconds=60000"}) @DirtiesContext -// TODO: Investigate why they fail -@Ignore("flakey") public class FeignClientServerErrorTests { private static final Log log = LogFactory.getLog(FeignClientServerErrorTests.class);