From 8afa3e8d37ebc34131465b1bc00a818eda305b2f Mon Sep 17 00:00:00 2001 From: Nick Date: Thu, 19 Sep 2019 10:49:46 +0300 Subject: [PATCH] Fix ThreadPoolTaskScheduler proxy mechanism (#1447) --- .../async/ExecutorBeanPostProcessor.java | 9 +++- .../async/issues/issue410/Issue410Tests.java | 53 ++++++++++++++++--- 2 files changed, 54 insertions(+), 8 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessor.java index 5316f0467..cb7b865f3 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/async/ExecutorBeanPostProcessor.java @@ -185,8 +185,13 @@ class ExecutorBeanPostProcessor implements BeanPostProcessor { Object createAsyncTaskExecutorProxy(Object bean, boolean cglibProxy, AsyncTaskExecutor executor) { - return getProxiedObject(bean, cglibProxy, executor, - () -> new LazyTraceAsyncTaskExecutor(this.beanFactory, executor)); + return getProxiedObject(bean, cglibProxy, executor, () -> { + if (bean instanceof ThreadPoolTaskScheduler) { + return new LazyTraceThreadPoolTaskScheduler(this.beanFactory, + (ThreadPoolTaskScheduler) executor); + } + return new LazyTraceAsyncTaskExecutor(this.beanFactory, executor); + }); } private Object getProxiedObject(Object bean, boolean cglibProxy, Executor executor, diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue410/Issue410Tests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue410/Issue410Tests.java index bfef80152..07a20407d 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue410/Issue410Tests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/async/issues/issue410/Issue410Tests.java @@ -17,6 +17,7 @@ package org.springframework.cloud.sleuth.instrument.async.issues.issue410; import java.lang.invoke.MethodHandles; +import java.util.Date; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import java.util.concurrent.Executor; @@ -182,12 +183,35 @@ public class Issue410Tests { * Related to issue #1232 */ @Test - public void should_pass_tracing_info_for_completable_futures_with_threadPoolTaskScheduler() { + public void should_pass_tracing_info_for_submitted_tasks_with_threadPoolTaskScheduler() { Span span = this.tracer.nextSpan().name("foo"); log.info("Starting test"); try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) { String response = this.restTemplate.getForObject( - "http://localhost:" + port() + "/threadPoolTaskScheduler", + "http://localhost:" + port() + "/threadPoolTaskScheduler_submit", + String.class); + + then(response).isEqualTo(span.context().traceIdString()); + Awaitility.await().untilAsserted(() -> { + then(this.asyncTask.getSpan().get()).isNotNull(); + then(this.asyncTask.getSpan().get().context().traceId()) + .isEqualTo(span.context().traceId()); + }); + } + finally { + span.finish(); + } + + then(this.tracer.currentSpan()).isNull(); + } + + @Test + public void should_pass_tracing_info_for_scheduled_tasks_with_threadPoolTaskScheduler() { + Span span = this.tracer.nextSpan().name("foo"); + log.info("Starting test"); + try (Tracer.SpanInScope ws = this.tracer.withSpanInScope(span)) { + String response = this.restTemplate.getForObject( + "http://localhost:" + port() + "/threadPoolTaskScheduler_schedule", String.class); then(response).isEqualTo(span.context().traceIdString()); @@ -372,7 +396,7 @@ class AsyncTask { return this.span.get(); } - public Span threadPoolTaskScheduler() + public Span threadPoolTaskSchedulerSubmit() throws ExecutionException, InterruptedException { log.info("This task is running with ThreadPoolTaskScheduler"); this.threadPoolTaskScheduler.submit(() -> { @@ -382,6 +406,16 @@ class AsyncTask { return this.span.get(); } + public Span threadPoolTaskSchedulerSchedule() + throws ExecutionException, InterruptedException { + log.info("This task is running with ThreadPoolTaskScheduler"); + this.threadPoolTaskScheduler.schedule(() -> { + log.info("Hello from runnable"); + AsyncTask.this.span.set(AsyncTask.this.tracer.currentSpan()); + }, new Date()).get(); + return this.span.get(); + } + public AtomicReference getSpan() { return this.span; } @@ -427,11 +461,18 @@ class Application { return this.asyncTask.taskScheduler().context().traceIdString(); } - @RequestMapping("/threadPoolTaskScheduler") - public String threadPoolTaskScheduler() + @RequestMapping("/threadPoolTaskScheduler_submit") + public String threadPoolTaskSchedulerSubmit() throws ExecutionException, InterruptedException { log.info("Executing completable via ThreadPoolTaskScheduler"); - return this.asyncTask.threadPoolTaskScheduler().context().traceIdString(); + return this.asyncTask.threadPoolTaskSchedulerSubmit().context().traceIdString(); + } + + @RequestMapping("/threadPoolTaskScheduler_schedule") + public String threadPoolTaskSchedulerSchedule() + throws ExecutionException, InterruptedException { + log.info("Executing completable via ThreadPoolTaskScheduler"); + return this.asyncTask.threadPoolTaskSchedulerSchedule().context().traceIdString(); } @RequestMapping("/scheduledThreadPoolExecutor")