diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java index 21c36a3e7..df67b008e 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfiguration.java @@ -84,7 +84,7 @@ public class TraceReactorAutoConfiguration { if (log.isTraceEnabled()) { log.trace("Cleaning up hooks"); } - Hooks.resetOnLastOperator(SLEUTH_TRACE_REACTOR_KEY); + Hooks.resetOnEachOperator(SLEUTH_TRACE_REACTOR_KEY); Schedulers .removeExecutorServiceDecorator(SLEUTH_REACTOR_EXECUTOR_SERVICE_KEY); } @@ -115,7 +115,7 @@ class HookRegisteringBeanDefinitionRegistryPostProcessor } void setupHooks(BeanFactory beanFactory) { - Hooks.onLastOperator( + Hooks.onEachOperator( TraceReactorAutoConfiguration.TraceReactorConfiguration.SLEUTH_TRACE_REACTOR_KEY, ReactorSleuth.scopePassingSpanOperator(this.context)); Schedulers.setExecutorServiceDecorator( diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue1212/GH1212Tests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue1212/GH1212Tests.java index 96194d3a6..7c390fbfb 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue1212/GH1212Tests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue1212/GH1212Tests.java @@ -24,6 +24,7 @@ import org.slf4j.LoggerFactory; import org.springframework.aop.interceptor.AsyncExecutionAspectSupport; import org.springframework.boot.SpringBootConfiguration; +import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.context.ApplicationContext; @@ -48,7 +49,7 @@ public class GH1212Tests { @Test public void defaultTaskExecutor() throws Exception { try (ConfigurableApplicationContext ctx = new SpringApplicationBuilder(App.class, - DefaultTaskExecutorConfig.class).run()) { + DefaultTaskExecutorConfig.class).web(WebApplicationType.NONE).run()) { String asyncThreadName = getAsyncThreadName(ctx); assertThat(asyncThreadName).startsWith("defaultTaskExecutor"); } @@ -57,7 +58,7 @@ public class GH1212Tests { @Test public void singleTaskExecutor() throws Exception { try (ConfigurableApplicationContext ctx = new SpringApplicationBuilder(App.class, - SingleTaskExecutorConfig.class).run()) { + SingleTaskExecutorConfig.class).web(WebApplicationType.NONE).run()) { String asyncThreadName = getAsyncThreadName(ctx); assertThat(asyncThreadName).startsWith("singleTaskExecutor"); } @@ -66,7 +67,7 @@ public class GH1212Tests { @Test public void multipleTaskExecutors() throws Exception { try (ConfigurableApplicationContext ctx = new SpringApplicationBuilder(App.class, - MultipleTaskExecutorConfig.class).run()) { + MultipleTaskExecutorConfig.class).web(WebApplicationType.NONE).run()) { String asyncThreadName = getAsyncThreadName(ctx); assertThat(asyncThreadName).doesNotStartWith("multipleTaskExecutor"); assertThat(asyncThreadName).startsWith("SimpleAsyncTaskExecutor"); // <-- @@ -81,7 +82,7 @@ public class GH1212Tests { @Test public void customAsyncConfigurer() throws Exception { try (ConfigurableApplicationContext ctx = new SpringApplicationBuilder(App.class, - CustomAsyncConfigurerConfig.class).run()) { + CustomAsyncConfigurerConfig.class).web(WebApplicationType.NONE).run()) { String asyncThreadName = getAsyncThreadName(ctx); assertThat(asyncThreadName).startsWith("customAsyncConfigurer"); } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriberTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriberTests.java index a53685190..c7c74e217 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriberTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/SpanSubscriberTests.java @@ -32,7 +32,6 @@ import org.reactivestreams.Subscriber; import org.reactivestreams.Subscription; import reactor.core.CoreSubscriber; import reactor.core.publisher.Flux; -import reactor.core.publisher.Hooks; import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; @@ -107,7 +106,7 @@ public class SpanSubscriberTests { log.info("Hello"); // Disable global hooks for local hook testing - Hooks.resetOnLastOperator(); + TraceReactorAutoConfigurationAccessorConfiguration.close(); try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) { @@ -182,6 +181,8 @@ public class SpanSubscriberTests { Awaitility.await().untilAsserted(() -> { then(this.tracer.currentSpan()).isNull(); }); + + TraceReactorAutoConfigurationAccessorConfiguration.setup(this.factory); } @Test diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfigurationAccessorConfiguration.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfigurationAccessorConfiguration.java index 11d6dc3cb..8d02c8f4d 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfigurationAccessorConfiguration.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/TraceReactorAutoConfigurationAccessorConfiguration.java @@ -18,8 +18,8 @@ package org.springframework.cloud.sleuth.instrument.reactor; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import reactor.core.publisher.Hooks; -import reactor.core.scheduler.Schedulers; + +import org.springframework.context.ConfigurableApplicationContext; /** * @author Marcin Grzejszczak @@ -38,9 +38,15 @@ public final class TraceReactorAutoConfigurationAccessorConfiguration { log.trace("Cleaning up hooks"); } new TraceReactorAutoConfiguration.TraceReactorConfiguration().cleanupHooks(); - Hooks.resetOnLastOperator(); - Hooks.resetOnLastOperator(); - Schedulers.resetFactory(); + } + + public static void setup(ConfigurableApplicationContext context) { + if (log.isTraceEnabled()) { + log.trace("Setting up hooks"); + } + TraceReactorAutoConfiguration.TraceReactorConfiguration + .traceHookRegisteringBeanDefinitionRegistryPostProcessor(context) + .setupHooks(context); } } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java index 3f0ccb4b7..dd81a6272 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/reactor/sample/FlatMapTests.java @@ -30,9 +30,7 @@ import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import reactor.core.publisher.Flux; -import reactor.core.publisher.Hooks; import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; import zipkin2.Span; import org.springframework.boot.WebApplicationType; @@ -41,6 +39,7 @@ import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.boot.test.rule.OutputCapture; import org.springframework.cloud.sleuth.DisableWebFluxSecurity; import org.springframework.cloud.sleuth.instrument.reactor.Issue866Configuration; +import org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfigurationAccessorConfiguration; import org.springframework.cloud.sleuth.util.ArrayListSpanReporter; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -65,8 +64,7 @@ public class FlatMapTests { @BeforeClass public static void setup() { - Hooks.resetOnLastOperator(); - Schedulers.resetFactory(); + TraceReactorAutoConfigurationAccessorConfiguration.close(); Issue866Configuration.hook = null; } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java index 65891cfb9..734821a99 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/client/integration/WebClientTests.java @@ -57,8 +57,6 @@ import org.junit.ClassRule; import org.junit.Rule; import org.junit.Test; import org.junit.runner.RunWith; -import reactor.core.publisher.Hooks; -import reactor.core.scheduler.Schedulers; import reactor.netty.http.client.HttpClient; import reactor.netty.http.client.HttpClientResponse; import zipkin2.Annotation; @@ -80,6 +78,7 @@ import org.springframework.cloud.gateway.config.GatewayClassPathWarningAutoConfi import org.springframework.cloud.netflix.ribbon.RibbonClient; import org.springframework.cloud.openfeign.EnableFeignClients; import org.springframework.cloud.openfeign.FeignClient; +import org.springframework.cloud.sleuth.instrument.reactor.TraceReactorAutoConfigurationAccessorConfiguration; import org.springframework.cloud.sleuth.instrument.web.TraceWebServletAutoConfiguration; import org.springframework.cloud.sleuth.util.ArrayListSpanReporter; import org.springframework.context.annotation.Bean; @@ -166,8 +165,7 @@ public class WebClientTests { @BeforeClass public static void cleanup() { - Hooks.resetOnLastOperator(); - Schedulers.resetFactory(); + TraceReactorAutoConfigurationAccessorConfiguration.close(); } @After diff --git a/src/checkstyle/checkstyle-suppressions.xml b/src/checkstyle/checkstyle-suppressions.xml index c526753f2..2b2f25d0b 100644 --- a/src/checkstyle/checkstyle-suppressions.xml +++ b/src/checkstyle/checkstyle-suppressions.xml @@ -27,6 +27,7 @@ +